(t *testing.T)
| 3841 | } |
| 3842 | |
| 3843 | func TestIngester_QueryStream(t *testing.T) { |
| 3844 | // Create ingester. |
| 3845 | cfg := defaultIngesterTestConfig(t) |
| 3846 | |
| 3847 | for _, enc := range encodings { |
| 3848 | t.Run(enc.String(), func(t *testing.T) { |
| 3849 | i, err := prepareIngesterWithBlocksStorage(t, cfg, prometheus.NewRegistry()) |
| 3850 | require.NoError(t, err) |
| 3851 | require.NoError(t, services.StartAndAwaitRunning(context.Background(), i)) |
| 3852 | defer services.StopAndAwaitTerminated(context.Background(), i) //nolint:errcheck |
| 3853 | |
| 3854 | // Wait until it's ACTIVE. |
| 3855 | test.Poll(t, 1*time.Second, ring.ACTIVE, func() any { |
| 3856 | return i.lifecycler.GetState() |
| 3857 | }) |
| 3858 | |
| 3859 | // Push series. |
| 3860 | ctx := user.InjectOrgID(context.Background(), userID) |
| 3861 | lbls := labels.FromStrings(labels.MetricName, "foo") |
| 3862 | var ( |
| 3863 | req *cortexpb.WriteRequest |
| 3864 | expectedResponseChunks *client.QueryStreamResponse |
| 3865 | ) |
| 3866 | switch enc { |
| 3867 | case encoding.PrometheusXorChunk: |
| 3868 | req, expectedResponseChunks = mockWriteRequest(t, lbls, 123000, 456) |
| 3869 | case encoding.PrometheusHistogramChunk: |
| 3870 | req, expectedResponseChunks = mockHistogramWriteRequest(t, lbls, 123000, 456, false) |
| 3871 | case encoding.PrometheusFloatHistogramChunk: |
| 3872 | req, expectedResponseChunks = mockHistogramWriteRequest(t, lbls, 123000, 456, true) |
| 3873 | } |
| 3874 | _, err = i.Push(ctx, req) |
| 3875 | require.NoError(t, err) |
| 3876 | |
| 3877 | // Create a GRPC server used to query back the data. |
| 3878 | serv := grpc.NewServer(grpc.StreamInterceptor(middleware.StreamServerUserHeaderInterceptor)) |
| 3879 | defer serv.GracefulStop() |
| 3880 | client.RegisterIngesterServer(serv, i) |
| 3881 | |
| 3882 | listener, err := net.Listen("tcp", "localhost:0") |
| 3883 | require.NoError(t, err) |
| 3884 | |
| 3885 | go func() { |
| 3886 | require.NoError(t, serv.Serve(listener)) |
| 3887 | }() |
| 3888 | |
| 3889 | // Query back the series using GRPC streaming. |
| 3890 | c, err := client.MakeIngesterClient(listener.Addr().String(), defaultClientTestConfig(), false) |
| 3891 | require.NoError(t, err) |
| 3892 | defer c.Close() |
| 3893 | |
| 3894 | queryRequest := &client.QueryRequest{ |
| 3895 | StartTimestampMs: 0, |
| 3896 | EndTimestampMs: 200000, |
| 3897 | Matchers: []*client.LabelMatcher{{ |
| 3898 | Type: client.EQUAL, |
| 3899 | Name: model.MetricNameLabel, |
| 3900 | Value: "foo", |
nothing calls this directly
no test coverage detected