(t *testing.T)
| 141 | } |
| 142 | |
| 143 | func 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) |
nothing calls this directly
no test coverage detected