From 495f7c5f27d92689f9c6c2670680fd9b73bb765a Mon Sep 17 00:00:00 2001 From: SungJin1212 Date: Thu, 20 Aug 2026 16:10:09 +0900 Subject: [PATCH] Support partial write on PRW2.0 Signed-off-by: SungJin1212 --- CHANGELOG.md | 1 + integration/remote_write_v2_test.go | 186 +++++++++++++++++++++++++++ pkg/util/push/push.go | 135 +++++++++++++------- pkg/util/push/push_test.go | 187 +++++++++++++++++++++++++++- 4 files changed, 460 insertions(+), 49 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 16148eeae63..75ed45c319b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -55,6 +55,7 @@ * [ENHANCEMENT] Upgrade Thanos and promql-engine to latest. #7740 * [ENHANCEMENT] Ruler: Adjust ruler frontend decoder to not wrap query error messages with execution prefix, this makes error responses consistent between internal and external ruler paths. #7741 * [ENHANCEMENT] Distributor: Deduplicate metric metadata when converting PRW 2.0 requests. PRW 2.0 attaches metadata to every series, so a metric family was previously expanded into one `MetricMetadata` per series. #7760 +* [ENHANCEMENT] Distributor: Support partial write for Prometheus Remote Write 2.0 requests. Invalid series are now skipped and reported together in the `400` response instead of rejecting the whole batch, the valid ones are written, and the `X-Prometheus-Remote-Write-*-Written` response headers are set even when a `400` is returned. Exemplar only `TimeSeries`, which the Prometheus sender emits are also accepted now, consistently with the remote write 1.0 path. #7761 * [BUGFIX] Querier: Fix queryWithRetry and labelsWithRetry returning (nil, nil) on cancelled context by propagating ctx.Err(). #7370 * [BUGFIX] Metrics Helper: Fix non-deterministic bucket order in merged histograms by sorting buckets after map iteration, matching Prometheus client library behavior. #7380 * [BUGFIX] Distributor: Return HTTP 401 Unauthorized when tenant ID resolution fails in the Prometheus Remote Write 2.0 path. #7389 diff --git a/integration/remote_write_v2_test.go b/integration/remote_write_v2_test.go index f1455848b8b..9379dcc28d2 100644 --- a/integration/remote_write_v2_test.go +++ b/integration/remote_write_v2_test.go @@ -879,6 +879,192 @@ func TestExemplar(t *testing.T) { exemplars, err := c.QueryExemplars("test_metric", start, end) require.NoError(t, err) require.Equal(t, 1, len(exemplars)) + + // The Prometheus sender emits exemplar only TimeSeries, see + // https://github.com/prometheus/prometheus/issues/17857. + exemplarOnly := []writev2.TimeSeries{ + { + LabelsRefs: []uint32{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}, + Exemplars: []writev2.Exemplar{{LabelsRefs: []uint32{13, 14}, Value: 2, Timestamp: tsMillis + 1}}, + }, + } + writeStats, err = c.PushV2(symbols, exemplarOnly) + require.NoError(t, err) + testPushHeader(t, writeStats, 0, 0, 1) + + exemplars, err = c.QueryExemplars("test_metric", start, end) + require.NoError(t, err) + require.Equal(t, 1, len(exemplars)) + require.Equal(t, 2, len(exemplars[0].Exemplars)) +} + +func TestPRW2PartialWrite(t *testing.T) { + s, err := e2e.NewScenario(networkName) + require.NoError(t, err) + defer s.Close() + + // Start dependencies. + consul := e2edb.NewConsulWithName("consul") + require.NoError(t, s.StartAndWaitReady(consul)) + + flags := mergeFlags( + AlertmanagerLocalFlags(), + map[string]string{ + "-store.engine": blocksStorageEngine, + "-blocks-storage.backend": "filesystem", + "-blocks-storage.tsdb.head-compaction-interval": "4m", + "-blocks-storage.bucket-store.sync-interval": "15m", + "-blocks-storage.bucket-store.index-cache.backend": tsdb.IndexCacheBackendInMemory, + "-blocks-storage.bucket-store.bucket-index.enabled": "true", + "-blocks-storage.tsdb.ship-interval": "1s", + "-blocks-storage.tsdb.enable-native-histograms": "true", + // Ingester. + "-ring.store": "consul", + "-consul.hostname": consul.NetworkHTTPEndpoint(), + "-ingester.max-exemplars": "100", + // Distributor. + "-distributor.replication-factor": "1", + "-distributor.remote-writev2-enabled": "true", + // Store-gateway. + "-store-gateway.sharding-enabled": "false", + // alert manager + "-alertmanager.web.external-url": "http://localhost/alertmanager", + }, + ) + + // make alert manager config dir + require.NoError(t, writeFileToSharedDir(s, "alertmanager_configs", []byte{})) + + path := path.Join(s.SharedDir(), "cortex-1") + + flags = mergeFlags(flags, map[string]string{"-blocks-storage.filesystem.dir": path}) + // Start Cortex replicas. + cortex := e2ecortex.NewSingleBinary("cortex", flags, "") + require.NoError(t, s.StartAndWaitReady(cortex)) + + // Wait until Cortex replicas have updated the ring state. + require.NoError(t, cortex.WaitSumMetrics(e2e.Equals(float64(512)), "cortex_ring_tokens_total")) + + c, err := e2ecortex.NewClient(cortex.HTTPEndpoint(), cortex.HTTPEndpoint(), "", "", "user-1") + require.NoError(t, err) + + now := time.Now() + tsMillis := e2e.TimeToMilliseconds(now) + start := now.Add(-time.Minute) + end := now.Add(time.Minute) + + symbols := []string{ + "", // 0 + "__name__", // 1 + "good_sample", // 2 + "good_histogram", // 3 + "dropped_labels", // 4 + "dropped_exemplar", // 5 + "dropped_metadata", // 6 + "empty_series", // 7 + "trace_id", // 8 + "abc123", // 9 + } + // Any ref greater than or equal to the symbols table length is out of range. + const invalidRef = 10 + + t.Run("every series is dropped during conversion", func(t *testing.T) { + timeseries := []writev2.TimeSeries{ + {LabelsRefs: []uint32{1, 7}}, + {LabelsRefs: []uint32{1, invalidRef}, Samples: []writev2.Sample{{Value: 1, Timestamp: tsMillis}}}, + } + + writeStats, err := c.PushV2(symbols, timeseries) + require.Error(t, err) + require.Contains(t, err.Error(), "400") + require.Contains(t, err.Error(), "TimeSeries must contain at least one sample, histogram or exemplar") + require.Contains(t, err.Error(), "outside of symbols table") + + // Nothing was written, so the response headers must report zero of everything. + testPushHeader(t, writeStats, 0, 0, 0) + + result, err := c.Query("empty_series", now) + require.NoError(t, err) + require.Empty(t, result.(model.Vector)) + }) + + t.Run("dropped data is excluded from the written stats headers", func(t *testing.T) { + h := writev2.FromIntHistogram(tsMillis, tsdbutil.GenerateTestHistogram(1)) + + timeseries := []writev2.TimeSeries{ + // Written: 1 sample and 1 exemplar. + { + LabelsRefs: []uint32{1, 2}, + Samples: []writev2.Sample{{Value: 1, Timestamp: tsMillis}}, + Exemplars: []writev2.Exemplar{{LabelsRefs: []uint32{8, 9}, Value: 1, Timestamp: tsMillis}}, + }, + // Written: 1 histogram. + { + LabelsRefs: []uint32{1, 3}, + Histograms: []writev2.Histogram{h}, + }, + // Dropped on an out of range label ref, along with its 3 samples and 2 exemplars. + { + LabelsRefs: []uint32{1, invalidRef}, + Samples: []writev2.Sample{ + {Value: 1, Timestamp: tsMillis}, + {Value: 2, Timestamp: tsMillis + 1}, + {Value: 3, Timestamp: tsMillis + 2}, + }, + Exemplars: []writev2.Exemplar{ + {LabelsRefs: []uint32{8, 9}, Value: 1, Timestamp: tsMillis}, + {LabelsRefs: []uint32{8, 9}, Value: 2, Timestamp: tsMillis + 1}, + }, + }, + // Only the exemplar is dropped on an out of range label ref, the samples are kept. + { + LabelsRefs: []uint32{1, 5}, + Samples: []writev2.Sample{ + {Value: 1, Timestamp: tsMillis}, + {Value: 2, Timestamp: tsMillis + 1}, + }, + Exemplars: []writev2.Exemplar{{LabelsRefs: []uint32{8, invalidRef}, Value: 1, Timestamp: tsMillis}}, + }, + // The metadata is dropped on an out of range unit ref, but the series itself is kept. + { + LabelsRefs: []uint32{1, 6}, + Metadata: writev2.Metadata{UnitRef: invalidRef}, + Histograms: []writev2.Histogram{h}, + }, + // Dropped for holding no data at all. + {LabelsRefs: []uint32{1, 7}}, + } + + writeStats, err := c.PushV2(symbols, timeseries) + + // The client is told the request was a bad one, even though it was partially written. + require.Error(t, err) + require.Contains(t, err.Error(), "400") + + // Only the series that survived conversion are reported as written. + testPushHeader(t, writeStats, 3, 2, 1) + + // And they are the only ones actually stored. + for _, name := range []string{"good_sample", "good_histogram", "dropped_metadata", "dropped_exemplar"} { + result, err := c.Query(name, now) + require.NoError(t, err) + require.Len(t, result.(model.Vector), 1, "%s must be ingested", name) + } + for _, name := range []string{"empty_series"} { + result, err := c.Query(name, now) + require.NoError(t, err) + require.Empty(t, result.(model.Vector), "%s must not be ingested", name) + } + + exemplars, err := c.QueryExemplars("good_sample", start, end) + require.NoError(t, err) + require.Len(t, exemplars, 1) + require.Len(t, exemplars[0].Exemplars, 1) + + exemplars, err = c.QueryExemplars("dropped_exemplar", start, end) + require.NoError(t, err) + require.Empty(t, exemplars) + }) } func Test_WriteStatWithReplication(t *testing.T) { diff --git a/pkg/util/push/push.go b/pkg/util/push/push.go index 50860635ecf..26b1d821b6e 100644 --- a/pkg/util/push/push.go +++ b/pkg/util/push/push.go @@ -115,11 +115,11 @@ func Handler(remoteWrite2Enabled bool, acceptUnknownRemoteWriteContentType bool, req.Source = cortexpb.API } - v1Req, err := convertV2RequestToV1(req, overrides.EnableTypeAndUnitLabels(userID), overrides.EnableStartTimestamp(userID)) - if err != nil { - level.Error(logger).Log("err", err.Error()) - http.Error(w, err.Error(), http.StatusBadRequest) - return + // convertErr is a non-retriable bad request error. The series that converted + // successfully are still pushed, so the request is partially written. + v1Req, convertErr := convertV2RequestToV1(req, overrides.EnableTypeAndUnitLabels(userID), overrides.EnableStartTimestamp(userID)) + if convertErr != nil { + level.Warn(logger).Log("msg", "remote write v2 request partially converted", "err", convertErr) } v1Req.SkipLabelNameValidation = false @@ -127,7 +127,10 @@ func Handler(remoteWrite2Enabled bool, acceptUnknownRemoteWriteContentType bool, v1Req.Source = cortexpb.API } - if writeResp, err := push(ctx, &v1Req.WriteRequest); err != nil { + // The Distributor owns the pooled TimeSeries, so push must be called even when + // every series was skipped. + writeResp, err := push(ctx, &v1Req.WriteRequest) + if err != nil { if errors.Is(err, context.Canceled) { err = httpgrpc.Errorf(util_api.StatusClientClosedRequest, "%s", err.Error()) } @@ -148,11 +151,17 @@ func Handler(remoteWrite2Enabled bool, acceptUnknownRemoteWriteContentType bool, } else if resp.GetCode() != http.StatusAccepted && resp.GetCode() != http.StatusTooManyRequests && resp.GetCode() != util_api.StatusClientClosedRequest { level.Warn(logger).Log("msg", "push refused", "err", err) } + // The push error takes precedence over convertErr: a 5xx must be retried. http.Error(w, string(resp.Body), int(resp.Code)) - } else { - setPRW2RespHeader(w, writeResp.Samples, writeResp.Histograms, writeResp.Exemplars) - w.WriteHeader(http.StatusNoContent) + return } + + setPRW2RespHeader(w, writeResp.Samples, writeResp.Histograms, writeResp.Exemplars) + if convertErr != nil { + http.Error(w, convertErr.Error(), http.StatusBadRequest) + return + } + w.WriteHeader(http.StatusNoContent) } // follow Prometheus https://github.com/prometheus/prometheus/blob/v3.3.1/storage/remote/write_handler.go#L121 @@ -222,40 +231,73 @@ type v2MetadataKey struct { unitRef uint32 } -func convertV2RequestToV1(req *cortexpb.PreallocWriteRequestV2, enableTypeAndUnitLabels bool, enableStartTimestamp bool) (v1Req cortexpb.PreallocWriteRequest, err error) { +// maxConversionErrs is the maximum number of per-series conversion errors reported back +// to the client. +const maxConversionErrs = 10 + +// conversionErrs collects per-series conversion errors, keeping at most maxConversionErrs +// of them while still reporting the total count. +type conversionErrs struct { + errs []error + count int +} + +func (c *conversionErrs) add(err error) { + c.count++ + if len(c.errs) < maxConversionErrs { + c.errs = append(c.errs, err) + } +} + +// err joins the collected errors, summarizing the ones that were left out. +func (c *conversionErrs) err() error { + if c.count == 0 { + return nil + } + if omitted := c.count - len(c.errs); omitted > 0 { + return errors.Join(append(c.errs, fmt.Errorf("%d more errors omitted", omitted))...) + } + return errors.Join(c.errs...) +} + +// convertV2RequestToV1 converts a remote write v2 request into a v1 one. +// +// Malformed series are skipped and their errors are joined into the returned error, +// which is always a non-retriable bad request error. Series that converted successfully +// are still returned, so the request can be partially written, following the Prometheus +// receiver behavior. +func convertV2RequestToV1(req *cortexpb.PreallocWriteRequestV2, enableTypeAndUnitLabels bool, enableStartTimestamp bool) (cortexpb.PreallocWriteRequest, error) { + var v1Req cortexpb.PreallocWriteRequest v1Timeseries := make([]cortexpb.PreallocTimeseries, 0, len(req.Timeseries)) var v1Metadata []*cortexpb.MetricMetadata // v2 attaches metadata to every series, so a metric family repeats once per series. seenMetadata := make(map[v2MetadataKey]struct{}) - - // Release any pulled TimeSeries back to the pool to prevent memory leaks in case of an error. - defer func() { - if err != nil { - for _, pts := range v1Timeseries { - if pts.TimeSeries != nil { - cortexpb.ReuseTimeseries(pts.TimeSeries) - } - } - } - }() + var badRequestErrs conversionErrs b := labels.NewScratchBuilder(0) symbols := req.Symbols for _, v2Ts := range req.Timeseries { lbs, err := v2Ts.ToLabels(&b, symbols) if err != nil { - return v1Req, err + badRequestErrs.add(err) + continue } - if len(v2Ts.Samples) == 0 && len(v2Ts.Histograms) == 0 { - return v1Req, fmt.Errorf("TimeSeries must contain at least one sample or histogram for series %v", lbs.String()) + // The remote write 2.0 spec requires a TimeSeries to hold at least one sample or + // histogram, but the Prometheus sender emits exemplar only TimeSeries, see + // https://github.com/prometheus/prometheus/issues/17857. + if len(v2Ts.Samples) == 0 && len(v2Ts.Histograms) == 0 && len(v2Ts.Exemplars) == 0 { + badRequestErrs.add(fmt.Errorf("TimeSeries must contain at least one sample, histogram or exemplar for series %v", lbs.String())) + continue } - if int(v2Ts.Metadata.UnitRef) >= len(symbols) { - return v1Req, fmt.Errorf("invalid UnitRef %d: exceeds symbols length %d", v2Ts.Metadata.UnitRef, len(symbols)) + // An out of range UnitRef only invalidates the metadata, so keep the series and + // let convertV2ToV1Metadata report it. + var unit string + if int(v2Ts.Metadata.UnitRef) < len(symbols) { + unit = symbols[v2Ts.Metadata.UnitRef] } - unit := symbols[v2Ts.Metadata.UnitRef] metricType := v2Ts.Metadata.Type shouldAttachTypeAndUnitLabels := enableTypeAndUnitLabels && (metricType != cortexpb.METRIC_TYPE_UNSPECIFIED || unit != "") if shouldAttachTypeAndUnitLabels { @@ -287,12 +329,7 @@ func convertV2RequestToV1(req *cortexpb.PreallocWriteRequestV2, enableTypeAndUni ts.Samples = append(ts.Samples, sample) } - ts.Exemplars, err = convertV2ToV1Exemplars(&b, symbols, v2Ts.Exemplars, ts.Exemplars[:0]) - if err != nil { - // Current ts is not appended to the v1Timeseries, so we should call reuse here. - cortexpb.ReuseTimeseries(ts) - return v1Req, err - } + ts.Exemplars = convertV2ToV1Exemplars(&b, symbols, lbs, v2Ts.Exemplars, ts.Exemplars[:0], &badRequestErrs) ts.Histograms = ts.Histograms[:0] for _, histogram := range v2Ts.Histograms { @@ -307,15 +344,24 @@ func convertV2RequestToV1(req *cortexpb.PreallocWriteRequestV2, enableTypeAndUni ts.Histograms = append(ts.Histograms, histogram) } + // An exemplar only series ends up empty when every one of its exemplars was dropped + // above, and there is nothing left to write. + if len(ts.Samples) == 0 && len(ts.Histograms) == 0 && len(ts.Exemplars) == 0 { + // Current ts is not appended to the v1Timeseries, so we should call reuse here. + cortexpb.ReuseTimeseries(ts) + continue + } + v1Timeseries = append(v1Timeseries, cortexpb.PreallocTimeseries{ TimeSeries: ts, }) if shouldConvertV2Metadata(v2Ts.Metadata) { - var metricName string - metricName, err = extract.MetricNameFromLabels(lbs) + // The series has already been appended above, so only its metadata is dropped here. + metricName, err := extract.MetricNameFromLabels(lbs) if err != nil { - return v1Req, err + badRequestErrs.add(err) + continue } key := v2MetadataKey{ @@ -325,10 +371,10 @@ func convertV2RequestToV1(req *cortexpb.PreallocWriteRequestV2, enableTypeAndUni unitRef: v2Ts.Metadata.UnitRef, } if _, ok := seenMetadata[key]; !ok { - var metadata *cortexpb.MetricMetadata - metadata, err = convertV2ToV1Metadata(metricName, symbols, v2Ts.Metadata) + metadata, err := convertV2ToV1Metadata(metricName, symbols, v2Ts.Metadata) if err != nil { - return v1Req, err + badRequestErrs.add(err) + continue } seenMetadata[key] = struct{}{} v1Metadata = append(v1Metadata, metadata) @@ -339,7 +385,7 @@ func convertV2RequestToV1(req *cortexpb.PreallocWriteRequestV2, enableTypeAndUni v1Req.Timeseries = v1Timeseries v1Req.Metadata = v1Metadata - return v1Req, nil + return v1Req, badRequestErrs.err() } func shouldConvertV2Metadata(metadata cortexpb.MetadataV2) bool { @@ -381,11 +427,14 @@ func convertV2ToV1Metadata(name string, symbols []string, metadata cortexpb.Meta }, nil } -func convertV2ToV1Exemplars(b *labels.ScratchBuilder, symbols []string, v2Exemplars []cortexpb.ExemplarV2, v1Exemplars []cortexpb.Exemplar) ([]cortexpb.Exemplar, error) { +// convertV2ToV1Exemplars converts the exemplars of a single series. A malformed exemplar is +// skipped and reported, following the Prometheus receiver behavior. +func convertV2ToV1Exemplars(b *labels.ScratchBuilder, symbols []string, seriesLabels labels.Labels, v2Exemplars []cortexpb.ExemplarV2, v1Exemplars []cortexpb.Exemplar, badRequestErrs *conversionErrs) []cortexpb.Exemplar { for _, e := range v2Exemplars { lbs, err := e.ToLabels(b, symbols) if err != nil { - return v1Exemplars, err + badRequestErrs.add(fmt.Errorf("parsing exemplar for series %v: %w", seriesLabels.String(), err)) + continue } v1Exemplars = append(v1Exemplars, cortexpb.Exemplar{ Labels: cortexpb.FromLabelsToLabelAdapters(lbs), @@ -393,7 +442,7 @@ func convertV2ToV1Exemplars(b *labels.ScratchBuilder, symbols []string, v2Exempl TimestampMs: e.Timestamp, }) } - return v1Exemplars, nil + return v1Exemplars } func getTypeLabel(msgType remote.WriteMessageType, unknownOrInvalidContentType bool) string { diff --git a/pkg/util/push/push_test.go b/pkg/util/push/push_test.go index 4cbdfe7fcec..683e4e082c7 100644 --- a/pkg/util/push/push_test.go +++ b/pkg/util/push/push_test.go @@ -6,6 +6,8 @@ import ( "fmt" "net/http" "net/http/httptest" + "strconv" + "strings" "testing" "time" @@ -643,6 +645,14 @@ func Test_convertV2RequestToV1(t *testing.T) { Exemplars: []cortexpb.ExemplarV2{{LabelsRefs: []uint32{11, 12}, Value: 1, Timestamp: 1}}, }, }, + { + // The Prometheus sender emits exemplar only TimeSeries, see + // https://github.com/prometheus/prometheus/issues/17857 + TimeSeriesV2: &cortexpb.TimeSeriesV2{ + LabelsRefs: []uint32{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}, + Exemplars: []cortexpb.ExemplarV2{{LabelsRefs: []uint32{11, 12}, Value: 1, Timestamp: 1}}, + }, + }, } v2Req.Symbols = symbols @@ -650,7 +660,7 @@ func Test_convertV2RequestToV1(t *testing.T) { v1Req, err := convertV2RequestToV1(&v2Req, false, false) assert.NoError(t, err) expectedSamples := 3 - expectedExemplars := 2 + expectedExemplars := 3 expectedHistograms := 2 countSamples := 0 countExemplars := 0 @@ -665,7 +675,7 @@ func Test_convertV2RequestToV1(t *testing.T) { assert.Equal(t, expectedSamples, countSamples) assert.Equal(t, expectedExemplars, countExemplars) assert.Equal(t, expectedHistograms, countHistograms) - assert.Equal(t, 4, len(v1Req.Timeseries)) + assert.Equal(t, 5, len(v1Req.Timeseries)) assert.Equal(t, 1, len(v1Req.Metadata)) } @@ -841,9 +851,8 @@ func TestHandler_remoteWrite(t *testing.T) { }, }, { - name: "remote write v2 with empty samples and histograms should return 400", + name: "remote write v2 with an exemplar only series is accepted", createBody: func() ([]byte, bool) { - // Create a request with a TimeSeries that has no samples and no histograms reqProto := writev2.Request{ Symbols: []string{"", "__name__", "foo"}, Timeseries: []writev2.TimeSeries{ @@ -863,8 +872,13 @@ func TestHandler_remoteWrite(t *testing.T) { require.NoError(t, err) return reqBytes, true }, - expectedStatus: http.StatusBadRequest, - expectedBody: "TimeSeries must contain at least one sample or histogram for series {__name__=\"foo\"}", + expectedStatus: http.StatusNoContent, + verifyResponse: func(resp *httptest.ResponseRecorder) { + respHeader := resp.Header() + assert.Equal(t, "1", respHeader[rw20WrittenSamplesHeader][0]) + assert.Equal(t, "1", respHeader[rw20WrittenHistogramsHeader][0]) + assert.Equal(t, "1", respHeader[rw20WrittenExemplarsHeader][0]) + }, }, { name: "remote write v1 with oversized histogram returns 400", @@ -1638,3 +1652,164 @@ func TestHandler_remoteWriteV2_UnauthorizedWithoutTenantID(t *testing.T) { assert.Contains(t, resp.Body.String(), user.ErrNoOrgID.Error()) assert.False(t, pushCalled, "push function must not be called when tenant ID is missing") } + +func Test_convertV2RequestToV1_ExemplarOnlySeries(t *testing.T) { + var v2Req cortexpb.PreallocWriteRequestV2 + v2Req.Symbols = []string{"", "__name__", "test_metric", "trace_id", "abc123"} + v2Req.Timeseries = []cortexpb.PreallocTimeseriesV2{ + { + TimeSeriesV2: &cortexpb.TimeSeriesV2{ + LabelsRefs: []uint32{1, 2}, + Exemplars: []cortexpb.ExemplarV2{{LabelsRefs: []uint32{3, 4}, Value: 1.5, Timestamp: 20}}, + }, + }, + } + + v1Req, err := convertV2RequestToV1(&v2Req, false, false) + require.NoError(t, err) + require.Len(t, v1Req.Timeseries, 1) + + ts := v1Req.Timeseries[0] + assert.Empty(t, ts.Samples) + assert.Empty(t, ts.Histograms) + assert.Equal(t, []cortexpb.LabelAdapter{{Name: "__name__", Value: "test_metric"}}, ts.Labels) + require.Len(t, ts.Exemplars, 1) + assert.Equal(t, []cortexpb.LabelAdapter{{Name: "trace_id", Value: "abc123"}}, ts.Exemplars[0].Labels) + assert.Equal(t, 1.5, ts.Exemplars[0].Value) + assert.Equal(t, int64(20), ts.Exemplars[0].TimestampMs) +} + +func TestHandler_remoteWriteV2_PartialWrite(t *testing.T) { + var limits validation.Limits + flagext.DefaultValues(&limits) + overrides := validation.NewOverrides(limits, nil) + + // Symbol 5 is out of range for a symbols table of length 5. + const invalidRef = 5 + + tests := []struct { + name string + // pushedSeries lists the metric names expected to reach the Distributor, in order. + pushedSeries []string + timeseries []writev2.TimeSeries + expectedErrs []string + expectedStats []string + }{ + { + name: "a series holding no data at all is skipped, the rest of the batch is pushed", + timeseries: []writev2.TimeSeries{ + {LabelsRefs: []uint32{1, 2}}, + {LabelsRefs: []uint32{1, 3}, Samples: []writev2.Sample{{Value: 1, Timestamp: 10}}}, + }, + pushedSeries: []string{"bar"}, + expectedErrs: []string{`TimeSeries must contain at least one sample, histogram or exemplar for series {__name__="foo"}`}, + }, + { + name: "multiple invalid series produce a joined error", + timeseries: []writev2.TimeSeries{ + {LabelsRefs: []uint32{1, 2}}, + {LabelsRefs: []uint32{1, 3}, Samples: []writev2.Sample{{Value: 1, Timestamp: 10}}}, + {LabelsRefs: []uint32{1, 4}}, + }, + pushedSeries: []string{"bar"}, + expectedErrs: []string{ + `TimeSeries must contain at least one sample, histogram or exemplar for series {__name__="foo"}`, + `TimeSeries must contain at least one sample, histogram or exemplar for series {__name__="baz"}`, + }, + }, + { + name: "every series is invalid, push is still called with an empty request", + timeseries: []writev2.TimeSeries{ + {LabelsRefs: []uint32{1, 2}}, + }, + pushedSeries: nil, + expectedErrs: []string{`TimeSeries must contain at least one sample, histogram or exemplar for series {__name__="foo"}`}, + }, + { + name: "a series with an out of range label ref is skipped", + timeseries: []writev2.TimeSeries{ + {LabelsRefs: []uint32{1, invalidRef}, Samples: []writev2.Sample{{Value: 1, Timestamp: 10}}}, + {LabelsRefs: []uint32{1, 3}, Samples: []writev2.Sample{{Value: 1, Timestamp: 10}}}, + }, + pushedSeries: []string{"bar"}, + expectedErrs: []string{"outside of symbols table"}, + }, + { + name: "an out of range exemplar label ref drops the exemplar, not the series", + timeseries: []writev2.TimeSeries{ + { + LabelsRefs: []uint32{1, 2}, + Samples: []writev2.Sample{{Value: 1, Timestamp: 10}}, + Exemplars: []writev2.Exemplar{{LabelsRefs: []uint32{1, invalidRef}, Value: 1, Timestamp: 10}}, + }, + {LabelsRefs: []uint32{1, 3}, Samples: []writev2.Sample{{Value: 1, Timestamp: 10}}}, + }, + pushedSeries: []string{"foo", "bar"}, + expectedErrs: []string{`parsing exemplar for series {__name__="foo"}`, "outside of symbols table"}, + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + var pushed []string + pushCalled := false + pushFunc := func(ctx context.Context, req *cortexpb.WriteRequest) (*cortexpb.WriteResponse, error) { + pushCalled = true + for _, ts := range req.Timeseries { + pushed = append(pushed, ts.Labels[0].Value) + } + return &cortexpb.WriteResponse{Samples: int64(len(req.Timeseries))}, nil + } + + reqProto := writev2.Request{ + Symbols: []string{"", "__name__", "foo", "bar", "baz"}, + Timeseries: test.timeseries, + } + reqBytes, err := reqProto.Marshal() + require.NoError(t, err) + + handler := Handler(true, false, 100000, overrides, nil, pushFunc, nil) + httpReq := createRequest(t, reqBytes, true).WithContext(user.InjectOrgID(context.Background(), "user-1")) + + resp := httptest.NewRecorder() + handler.ServeHTTP(resp, httpReq) + + assert.Equal(t, http.StatusBadRequest, resp.Code) + for _, expectedErr := range test.expectedErrs { + assert.Contains(t, resp.Body.String(), expectedErr) + } + + // The valid part of the request must still be written. + assert.True(t, pushCalled, "push function must always be called") + assert.Equal(t, test.pushedSeries, pushed) + + // The written stats headers must be set even when responding with a 400. + respHeader := resp.Header() + assert.Equal(t, strconv.Itoa(len(test.pushedSeries)), respHeader.Get(rw20WrittenSamplesHeader)) + assert.Equal(t, "0", respHeader.Get(rw20WrittenHistogramsHeader)) + assert.Equal(t, "0", respHeader.Get(rw20WrittenExemplarsHeader)) + }) + } +} + +func Test_convertV2RequestToV1_BoundsReportedErrors(t *testing.T) { + const badSeries = maxConversionErrs + 5 + + var v2Req cortexpb.PreallocWriteRequestV2 + v2Req.Symbols = []string{"", "__name__", "test_metric"} + for range badSeries { + // A series holding no data at all. + v2Req.Timeseries = append(v2Req.Timeseries, cortexpb.PreallocTimeseriesV2{ + TimeSeriesV2: &cortexpb.TimeSeriesV2{LabelsRefs: []uint32{1, 2}}, + }) + } + + v1Req, err := convertV2RequestToV1(&v2Req, false, false) + require.Error(t, err) + assert.Empty(t, v1Req.Timeseries) + + // Only maxConversionErrs errors are reported, the rest are summarized. + lines := strings.Split(err.Error(), "\n") + assert.Len(t, lines, maxConversionErrs+1) + assert.Equal(t, fmt.Sprintf("%d more errors omitted", badSeries-maxConversionErrs), lines[len(lines)-1]) +}