MCPcopy Create free account
hub / github.com/cortexproject/cortex / TestConcurrentGrpcCalls

Function TestConcurrentGrpcCalls

integration/grpc_server_test.go:143–243  ·  view source on GitHub ↗
(t *testing.T)

Source from the content-addressed store, hash-verified

141}
142
143func TestConcurrentGrpcCalls(t *testing.T) {
144 cfg := server.Config{}
145 (&cfg).RegisterFlags(flag.NewFlagSet("fake", flag.ContinueOnError))
146
147 tc := map[string]struct {
148 cfg server.Config
149 register func(s *grpc.Server)
150 validate func(t testing.TB, con *grpc.ClientConn)
151 }{
152 "distributor": {
153 cfg: cfg,
154 register: func(s *grpc.Server) {
155 d := &mockGprcServer{}
156 distributorpb.RegisterDistributorServer(s, d)
157 },
158 validate: func(t testing.TB, conn *grpc.ClientConn) {
159 client := distributorpb.NewDistributorClient(conn)
160 wg := sync.WaitGroup{}
161 n := 10000
162 wg.Add(n)
163 for i := 0; i < n; i++ {
164 go func(i int) {
165 defer wg.Done()
166 ctx := context.Background()
167 ctx = metadata.NewOutgoingContext(ctx, metadata.MD{"i": []string{strconv.Itoa(i)}})
168 _, err := client.Push(ctx, createRequest(i))
169 require.NoError(t, err)
170 }(i)
171 }
172
173 wg.Wait()
174 },
175 },
176 "distributor push stream": {
177 cfg: cfg,
178 register: func(s *grpc.Server) {
179 d := &mockGprcServer{}
180 ingester_client.RegisterIngesterServer(s, d)
181 },
182 validate: func(t testing.TB, conn *grpc.ClientConn) {
183 ctx := context.Background()
184 client := ingester_client.NewIngesterClient(conn)
185 wg := sync.WaitGroup{}
186 n := 10000
187 wg.Add(n)
188 for i := 0; i < n; i++ {
189 go func(i int) {
190 defer wg.Done()
191 stream, err := client.PushStream(ctx)
192 require.NoError(t, err)
193
194 ctx = metadata.NewOutgoingContext(ctx, metadata.MD{"i": []string{strconv.Itoa(i)}})
195 err = stream.Send(&cortexpb.StreamWriteRequest{TenantID: strconv.Itoa(i), Request: createRequest(i)})
196 require.NoError(t, err)
197 _, err = stream.Recv()
198 require.NoError(t, err)
199 //err = stream.Send(&cortexpb.StreamWriteRequest{"i", createRequest(i + 1)})
200 //require.NoError(t, err)

Callers

nothing calls this directly

Calls 15

PushMethod · 0.95
PushStreamMethod · 0.95
QueryStreamMethod · 0.95
NewDistributorClientFunction · 0.92
createStreamResponseFunction · 0.85
runFunction · 0.85
DoneMethod · 0.80
createRequestFunction · 0.70
RegisterFlagsMethod · 0.65
SendMethod · 0.65
RecvMethod · 0.65

Tested by

no test coverage detected