From 98b7ce7dedd5e7dfcaa8d17c80e193f95dbdd916 Mon Sep 17 00:00:00 2001 From: SungJin1212 Date: Tue, 25 Aug 2026 18:00:35 +0900 Subject: [PATCH] Support partial write on PRW2.0 Signed-off-by: SungJin1212 --- CHANGELOG.md | 1 + integration/remote_write_v2_test.go | 169 +++++++++++++++++ pkg/util/push/push.go | 127 ++++++++----- pkg/util/push/push_test.go | 277 +++++++++++++++++++++++++++- 4 files changed, 531 insertions(+), 43 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7fe4ff88590..79a06342682 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -57,6 +57,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. #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..94c86f17fb1 100644 --- a/integration/remote_write_v2_test.go +++ b/integration/remote_write_v2_test.go @@ -881,6 +881,175 @@ func TestExemplar(t *testing.T) { require.Equal(t, 1, len(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 or histogram") + 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) { // Test `X-Prometheus-Remote-Write-Samples-Written` header value // with the replication. diff --git a/pkg/util/push/push.go b/pkg/util/push/push.go index 50860635ecf..93670090bee 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,75 @@ 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()) + badRequestErrs.add(fmt.Errorf("TimeSeries must contain at least one sample or histogram 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] + } else if enableTypeAndUnitLabels { + // The unit is part of the series identity here, so keeping the series would + // store it under a label set the sender never sent. + badRequestErrs.add(fmt.Errorf("invalid UnitRef %d: exceeds symbols length %d", v2Ts.Metadata.UnitRef, len(symbols))) + continue } - unit := symbols[v2Ts.Metadata.UnitRef] metricType := v2Ts.Metadata.Type shouldAttachTypeAndUnitLabels := enableTypeAndUnitLabels && (metricType != cortexpb.METRIC_TYPE_UNSPECIFIED || unit != "") if shouldAttachTypeAndUnitLabels { @@ -287,12 +331,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 { @@ -312,10 +351,11 @@ func convertV2RequestToV1(req *cortexpb.PreallocWriteRequestV2, enableTypeAndUni }) 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 +365,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 +379,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 +421,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 +436,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..cf70a36002c 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" @@ -819,6 +821,9 @@ func TestHandler_remoteWrite(t *testing.T) { expectedStatus int expectedBody string verifyResponse func(resp *httptest.ResponseRecorder) + // pushFunc overrides the default handler for cases where the Distributor is not + // called with the standard single series. + pushFunc Func }{ { name: "remote write v1", @@ -865,6 +870,12 @@ func TestHandler_remoteWrite(t *testing.T) { }, expectedStatus: http.StatusBadRequest, expectedBody: "TimeSeries must contain at least one sample or histogram for series {__name__=\"foo\"}", + // The only series is dropped during conversion, so the Distributor is still + // called, with an empty request. + pushFunc: func(_ context.Context, request *cortexpb.WriteRequest) (*cortexpb.WriteResponse, error) { + assert.Empty(t, request.Timeseries) + return &cortexpb.WriteResponse{}, nil + }, }, { name: "remote write v1 with oversized histogram returns 400", @@ -895,7 +906,11 @@ func TestHandler_remoteWrite(t *testing.T) { t.Run(test.name, func(t *testing.T) { ctx := context.Background() ctx = user.InjectOrgID(ctx, "user-1") - handler := Handler(true, false, 100000, overrides, nil, verifyWriteRequestHandler(t, cortexpb.API), nil) + pushFunc := test.pushFunc + if pushFunc == nil { + pushFunc = verifyWriteRequestHandler(t, cortexpb.API) + } + handler := Handler(true, false, 100000, overrides, nil, pushFunc, nil) body, isV2 := test.createBody() req := createRequest(t, body, isV2) @@ -1638,3 +1653,263 @@ 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") } + +// metricNameOf returns the metric name of a series that reached the Distributor. +func metricNameOf(lbs []cortexpb.LabelAdapter) string { + for _, l := range lbs { + if l.Name == labels.MetricName { + return l.Value + } + } + return "" +} + +func TestHandler_remoteWriteV2_PartialWrite(t *testing.T) { + var limits validation.Limits + flagext.DefaultValues(&limits) + + symbols := []string{ + "", // 0 + "__name__", // 1 + "foo", // 2 + "bar", // 3 + "baz", // 4 + "trace_id", // 5 + "abc123", // 6 + "help text", // 7 + "seconds", // 8 + } + // Any ref greater than or equal to the symbols table length is out of range. + const invalidRef = 9 + + validMetadata := writev2.Metadata{Type: writev2.Metadata_METRIC_TYPE_COUNTER, HelpRef: 7, UnitRef: 8} + hist := writev2.FromIntHistogram(20, tsdbutil.GenerateTestHistogram(1)) + + tests := []struct { + name string + timeseries []writev2.TimeSeries + // enableTypeAndUnitLabels turns the metadata unit into part of the series identity. + enableTypeAndUnitLabels bool + // The data expected to reach the Distributor. Samples, histograms and exemplars are + // listed as "@", metadata as the metric family name. + pushedSamples []string + pushedHistograms []string + pushedExemplars []string + pushedMetadata []string + expectedErrs []string + }{ + { + name: "an invalid series is skipped, the samples of the valid ones are pushed", + timeseries: []writev2.TimeSeries{ + {LabelsRefs: []uint32{1, invalidRef}, Samples: []writev2.Sample{{Value: 1, Timestamp: 10}}}, + {LabelsRefs: []uint32{1, 2}, Samples: []writev2.Sample{{Value: 1, Timestamp: 11}, {Value: 2, Timestamp: 12}}}, + }, + pushedSamples: []string{"foo@11", "foo@12"}, + expectedErrs: []string{"outside of symbols table"}, + }, + { + name: "an invalid series is skipped, the histograms of the valid ones are pushed", + timeseries: []writev2.TimeSeries{ + {LabelsRefs: []uint32{1, invalidRef}, Histograms: []writev2.Histogram{hist}}, + {LabelsRefs: []uint32{1, 2}, Histograms: []writev2.Histogram{hist}}, + }, + pushedHistograms: []string{"foo@20"}, + expectedErrs: []string{"outside of symbols table"}, + }, + { + name: "a malformed exemplar is dropped, the rest of its series is pushed", + timeseries: []writev2.TimeSeries{ + { + LabelsRefs: []uint32{1, 2}, + Samples: []writev2.Sample{{Value: 1, Timestamp: 10}}, + Exemplars: []writev2.Exemplar{ + {LabelsRefs: []uint32{5, invalidRef}, Value: 1, Timestamp: 30}, + {LabelsRefs: []uint32{5, 6}, Value: 2, Timestamp: 31}, + }, + }, + }, + pushedSamples: []string{"foo@10"}, + pushedExemplars: []string{"foo@31"}, + expectedErrs: []string{`parsing exemplar for series {__name__="foo"}`}, + }, + { + name: "an exemplar only series is rejected, the rest of the batch is pushed", + timeseries: []writev2.TimeSeries{ + { + LabelsRefs: []uint32{1, 2}, + Exemplars: []writev2.Exemplar{{LabelsRefs: []uint32{5, 6}, Value: 1, Timestamp: 30}}, + }, + {LabelsRefs: []uint32{1, 3}, Samples: []writev2.Sample{{Value: 1, Timestamp: 10}}}, + }, + pushedSamples: []string{"bar@10"}, + expectedErrs: []string{`TimeSeries must contain at least one sample or histogram for series {__name__="foo"}`}, + }, + { + name: "an out of range help ref drops the metadata, the series is still pushed", + timeseries: []writev2.TimeSeries{ + { + LabelsRefs: []uint32{1, 2}, + Samples: []writev2.Sample{{Value: 1, Timestamp: 10}}, + Metadata: writev2.Metadata{Type: writev2.Metadata_METRIC_TYPE_COUNTER, HelpRef: invalidRef, UnitRef: 8}, + }, + { + LabelsRefs: []uint32{1, 3}, + Samples: []writev2.Sample{{Value: 1, Timestamp: 11}}, + Metadata: validMetadata, + }, + }, + pushedSamples: []string{"foo@10", "bar@11"}, + pushedMetadata: []string{"bar"}, + expectedErrs: []string{"invalid HelpRef 9: exceeds symbols length 9"}, + }, + { + name: "an out of range unit ref drops the metadata, the series is still pushed", + timeseries: []writev2.TimeSeries{ + { + LabelsRefs: []uint32{1, 2}, + Histograms: []writev2.Histogram{hist}, + Metadata: writev2.Metadata{Type: writev2.Metadata_METRIC_TYPE_COUNTER, HelpRef: 7, UnitRef: invalidRef}, + }, + }, + pushedHistograms: []string{"foo@20"}, + expectedErrs: []string{"invalid UnitRef 9: exceeds symbols length 9"}, + }, + { + name: "an out of range unit ref drops the series when the unit is part of its identity", + enableTypeAndUnitLabels: true, + timeseries: []writev2.TimeSeries{ + { + LabelsRefs: []uint32{1, 2}, + Samples: []writev2.Sample{{Value: 1, Timestamp: 10}}, + Metadata: writev2.Metadata{Type: writev2.Metadata_METRIC_TYPE_COUNTER, HelpRef: 7, UnitRef: invalidRef}, + }, + {LabelsRefs: []uint32{1, 3}, Samples: []writev2.Sample{{Value: 1, Timestamp: 11}}}, + }, + pushedSamples: []string{"bar@11"}, + expectedErrs: []string{"invalid UnitRef 9: exceeds symbols length 9"}, + }, + { + 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}}, Metadata: validMetadata}, + }, + pushedSamples: []string{"bar@10"}, + pushedMetadata: []string{"bar"}, + expectedErrs: []string{`TimeSeries must contain at least one sample or histogram 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}}, + }, + pushedSamples: []string{"bar@10"}, + expectedErrs: []string{ + `TimeSeries must contain at least one sample or histogram for series {__name__="foo"}`, + `TimeSeries must contain at least one sample or histogram for series {__name__="baz"}`, + }, + }, + { + name: "every series is invalid, push is still called with an empty request", + timeseries: []writev2.TimeSeries{ + {LabelsRefs: []uint32{1, 2}}, + }, + expectedErrs: []string{`TimeSeries must contain at least one sample or histogram for series {__name__="foo"}`}, + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + var ( + pushCalled bool + pushedSamples []string + pushedHistograms []string + pushedExemplars []string + pushedMetadata []string + ) + pushFunc := func(ctx context.Context, req *cortexpb.WriteRequest) (*cortexpb.WriteResponse, error) { + pushCalled = true + for _, ts := range req.Timeseries { + name := metricNameOf(ts.Labels) + for _, s := range ts.Samples { + pushedSamples = append(pushedSamples, fmt.Sprintf("%s@%d", name, s.TimestampMs)) + } + for _, h := range ts.Histograms { + pushedHistograms = append(pushedHistograms, fmt.Sprintf("%s@%d", name, h.TimestampMs)) + } + for _, e := range ts.Exemplars { + pushedExemplars = append(pushedExemplars, fmt.Sprintf("%s@%d", name, e.TimestampMs)) + } + } + for _, m := range req.Metadata { + pushedMetadata = append(pushedMetadata, m.MetricFamilyName) + } + return &cortexpb.WriteResponse{ + Samples: int64(len(pushedSamples)), + Histograms: int64(len(pushedHistograms)), + Exemplars: int64(len(pushedExemplars)), + }, nil + } + + reqProto := writev2.Request{ + Symbols: symbols, + Timeseries: test.timeseries, + } + reqBytes, err := reqProto.Marshal() + require.NoError(t, err) + + caseLimits := limits + caseLimits.EnableTypeAndUnitLabels = test.enableTypeAndUnitLabels + overrides := validation.NewOverrides(caseLimits, nil) + + 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, down to the individual + // sample, histogram, exemplar and metadata. + assert.True(t, pushCalled, "push function must always be called") + assert.Equal(t, test.pushedSamples, pushedSamples, "samples") + assert.Equal(t, test.pushedHistograms, pushedHistograms, "histograms") + assert.Equal(t, test.pushedExemplars, pushedExemplars, "exemplars") + assert.Equal(t, test.pushedMetadata, pushedMetadata, "metadata") + + // The written stats headers must be set even when responding with a 400. + respHeader := resp.Header() + assert.Equal(t, strconv.Itoa(len(test.pushedSamples)), respHeader.Get(rw20WrittenSamplesHeader)) + assert.Equal(t, strconv.Itoa(len(test.pushedHistograms)), respHeader.Get(rw20WrittenHistogramsHeader)) + assert.Equal(t, strconv.Itoa(len(test.pushedExemplars)), 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]) +}