From b133e9f001046b69fa235145108bd0025326ebe7 Mon Sep 17 00:00:00 2001 From: Sandy Chen Date: Mon, 3 Aug 2026 11:31:17 +0900 Subject: [PATCH 1/4] Querier: remove redundant chunk data copy in ingester streaming select detachChunksFromBuffer copied every chunk's data on the querier's ingester-read path to detach it from the gRPC receive buffer, but the chunk data never aliases that buffer in the first place: Recv() allocates a fresh QueryStreamResponse per message, TimeSeriesChunk.Unmarshal appends a zero Chunk value, and Chunk.Unmarshal's append therefore starts from a nil slice and allocates a private backing array. The copy cost an extra allocation and memcpy per chunk and kept both the original and the copy reachable for the duration of streamingSelect, raising peak heap. Labels are different: LabelAdapter.Unmarshal uses yoloString, so label names and values do alias the receive buffer, and the FromLabelAdaptersToLabelsWithCopy call is kept. The call-site comment now documents both facts, and the assumption the removal relies on: response messages are never reused across Recv calls. If QueryStreamResponse ever becomes pooled, a detach copy must be reinstated. Two tests pin the non-aliasing invariant, through a direct gogo round trip (with a wire-buffer mutation cross-check) and through the registered cortexCodec with a payload large enough to exceed the gRPC buffer-pooling threshold, so a future pooled or zero-copy decoder fails loudly. The original #7670 benchmark, which constructed chunk data as sub-slices of one shared buffer (a shape a real gRPC unmarshal never produces), is rebuilt with independently allocated chunk data: dropping the copy saves ~107 KB and 200 allocations per 100-series response batch. Fixes #7732 Signed-off-by: Sandy Chen --- pkg/querier/distributor_queryable.go | 41 ++-- pkg/querier/distributor_queryable_test.go | 260 +++++++++++++++++++--- 2 files changed, 257 insertions(+), 44 deletions(-) diff --git a/pkg/querier/distributor_queryable.go b/pkg/querier/distributor_queryable.go index 42f5fb59d4a..5369971a40b 100644 --- a/pkg/querier/distributor_queryable.go +++ b/pkg/querier/distributor_queryable.go @@ -186,10 +186,31 @@ func (q *distributorQuerier) streamingSelect(ctx context.Context, sortSeries, pa continue } - // Detach label and chunk data from gRPC unmarshal buffers so the Go GC - // can reclaim receive buffers and reduce heap usage. + // Copy labels out of the gRPC unmarshal buffer: LabelAdapter.Unmarshal uses + // yoloString, so label names/values alias the receive buffer. Today that + // buffer is never returned to a reusable pool while still referenced (see + // pkg/cortexpb/codec.go: QueryStreamResponse doesn't implement + // ReleasableMessage, so cortexCodec never hands it back to mem.BufferPool), + // so the immediate risk is heap retention rather than corruption -- but + // that is an implementation detail of the current codec/message types, not + // a documented guarantee. Keep this copy: it is both what lets the Go GC + // reclaim the receive buffer promptly and a safety net against a future + // change (e.g. QueryStreamResponse opting into ReleasableMessage for a + // zero-copy chunk path) silently turning that retention issue into a + // use-after-recycle correctness bug for label values. + // + // Chunk.Data does not need the same treatment: Recv() allocates a fresh + // QueryStreamResponse per message, TimeSeriesChunk.Unmarshal appends a + // zero client.Chunk value, and Chunk.Unmarshal's + // `m.Data = append(m.Data[:0], ...)` therefore starts from a nil slice + // and allocates a private backing array, never aliasing the receive + // buffer. This relies on response messages not being reused: if + // QueryStreamResponse ever becomes pooled or reused across Recv calls, + // that same append would recycle prior backing arrays and a detach copy + // must be reinstated here. TestChunkDataDoesNotAliasWireBuffer_* pin the + // non-aliasing invariant. ls := cortexpb.FromLabelAdaptersToLabelsWithCopy(result.Labels) - chunks, err := chunkcompat.FromChunks(ls, detachChunksFromBuffer(result.Chunks)) + chunks, err := chunkcompat.FromChunks(ls, result.Chunks) if err != nil { return storage.ErrSeriesSet(err) } @@ -450,17 +471,3 @@ func labelHintsToSelectHints(hints *storage.LabelHints) *storage.SelectHints { Limit: hints.Limit, } } - -// detachChunksFromBuffer returns a copy of the chunks slice with data byte -// slices re-allocated so that the series no longer references the gRPC -// unmarshal buffer, allowing the Go GC to reclaim receive buffers. -func detachChunksFromBuffer(chunks []client.Chunk) []client.Chunk { - copied := make([]client.Chunk, len(chunks)) - for i, c := range chunks { - copied[i] = c - if len(c.Data) > 0 { - copied[i].Data = append([]byte(nil), c.Data...) - } - } - return copied -} diff --git a/pkg/querier/distributor_queryable_test.go b/pkg/querier/distributor_queryable_test.go index f4dfa0d4dde..439166ab81f 100644 --- a/pkg/querier/distributor_queryable_test.go +++ b/pkg/querier/distributor_queryable_test.go @@ -1,10 +1,13 @@ package querier import ( + "bytes" "context" "fmt" + "runtime" "testing" "time" + "unsafe" "github.com/prometheus/common/model" "github.com/prometheus/prometheus/model/labels" @@ -14,6 +17,7 @@ import ( "github.com/stretchr/testify/mock" "github.com/stretchr/testify/require" "github.com/weaveworks/common/user" + "google.golang.org/grpc/encoding" "github.com/cortexproject/cortex/pkg/chunk" promchunk "github.com/cortexproject/cortex/pkg/chunk/encoding" @@ -643,39 +647,34 @@ func TestDistributorQuerier_QueryIngestersWithinBoundary(t *testing.T) { } } +// BenchmarkIngesterStreamingSelect quantifies the cost of the chunk detach +// copy removed for issue #7732, on an input with the shape a real gRPC +// Unmarshal produces for chunks: each chunk's Data is its own independently +// allocated slice (Chunk.Unmarshal appends into the nil Data slice of a +// fresh per-Recv message, so it never aliases the receive buffer or any +// other chunk). Label memory layout is irrelevant to what the two arms +// compare -- both run the identical label copy -- so labels are ordinary +// independent strings here. This replaces the original #7670 benchmark, +// which constructed Chunk.Data as sub-slices of one shared buffer, a shape +// a real gRPC unmarshal never produces. func BenchmarkIngesterStreamingSelect(b *testing.B) { - // Simulate a realistic ingester response: 100 series, each with labels and chunk data - // that reference a shared backing buffer (mimicking gRPC unmarshal behavior). - // This benchmark measures the allocations in the streamingSelect hot path. const numSeries = 100 const chunkDataSize = 1024 - // Build a single large buffer to simulate the gRPC receive buffer. - // All label strings and chunk data will be slices into this buffer. - bufSize := numSeries * (chunkDataSize + 200) // 200 bytes for label strings per series - buf := make([]byte, bufSize) - for i := range buf { - buf[i] = byte(i % 256) - } - buildResponse := func() *client.QueryStreamResponse { - offset := 0 series := make([]client.TimeSeriesChunk, numSeries) for i := range series { - // Labels that reference the shared buffer (simulating protobuf unmarshal) - nameStr := string(buf[offset : offset+10]) - offset += 10 - valueStr := string(buf[offset : offset+20]) - offset += 20 - - // Chunk data that references the shared buffer - chunkData := buf[offset : offset+chunkDataSize] - offset += chunkDataSize + // Chunk data independently allocated, as a real Chunk.Unmarshal + // produces -- never a sub-slice of a shared buffer. + chunkData := make([]byte, chunkDataSize) + for j := range chunkData { + chunkData[j] = byte(j % 256) + } series[i] = client.TimeSeriesChunk{ Labels: []cortexpb.LabelAdapter{ - {Name: "__name__", Value: nameStr}, - {Name: "instance", Value: valueStr}, + {Name: "__name__", Value: fmt.Sprintf("metric_%d", i)}, + {Name: "instance", Value: fmt.Sprintf("instance-%d", i)}, }, Chunks: []client.Chunk{ { @@ -689,24 +688,231 @@ func BenchmarkIngesterStreamingSelect(b *testing.B) { return &client.QueryStreamResponse{Chunkseries: series} } - b.Run("with_detach", func(b *testing.B) { + // after_fix mirrors the current streamingSelect body: copy labels out of + // the receive buffer, pass chunks through untouched. + b.Run("after_fix_no_redundant_copy", func(b *testing.B) { b.ReportAllocs() for i := 0; i < b.N; i++ { resp := buildResponse() for _, result := range resp.Chunkseries { _ = cortexpb.FromLabelAdaptersToLabelsWithCopy(result.Labels) - _ = detachChunksFromBuffer(result.Chunks) + _ = result.Chunks } } }) - b.Run("without_detach", func(b *testing.B) { + // before_fix reproduces the removed detachChunksFromBuffer behavior inline + // (the helper itself is gone from the production code) to quantify the + // allocation/byte cost of copying data that Unmarshal already allocated + // separately, even under this realistic input shape. + b.Run("before_fix_redundant_copy", func(b *testing.B) { b.ReportAllocs() for i := 0; i < b.N; i++ { resp := buildResponse() for _, result := range resp.Chunkseries { - _ = cortexpb.FromLabelAdaptersToLabels(result.Labels) + _ = cortexpb.FromLabelAdaptersToLabelsWithCopy(result.Labels) + copied := make([]client.Chunk, len(result.Chunks)) + for ci, c := range result.Chunks { + copied[ci] = c + if len(c.Data) > 0 { + copied[ci].Data = append([]byte(nil), c.Data...) + } + } + _ = copied } } }) } + +// newAliasingProbeQueryStreamResponse returns a QueryStreamResponse with one +// series carrying a label and a chunk whose payload bytes are individually +// distinguishable (0x00, 0x01, 0x02, ...), so that accidental aliasing can be +// detected either by pointer-range overlap or by mutating the source buffer +// and observing whether the decoded bytes change. It also returns an +// independent copy of the chunk payload to compare against after the wire +// buffer has been mutated, since wire-format field ordering is an +// implementation detail we don't want the test to depend on. +// +// chunkDataSize is a parameter (not a constant) because callers going through +// the registered cortexCodec need a marshaled message large enough to cross +// mem.IsBelowBufferPoolingThreshold (1KiB): below that threshold, Marshal +// returns a plain mem.SliceBuffer whose Ref/Free are no-ops, which would +// silently skip exercising the real pool-backed mem.Buffer code path. +func newAliasingProbeQueryStreamResponse(chunkDataSize int) (resp *client.QueryStreamResponse, wantChunkData []byte) { + chunkData := make([]byte, chunkDataSize) + for i := range chunkData { + chunkData[i] = byte(i) + } + resp = &client.QueryStreamResponse{ + Chunkseries: []client.TimeSeriesChunk{ + { + Labels: []cortexpb.LabelAdapter{ + {Name: "__name__", Value: "some_metric_name"}, + {Name: "instance", Value: "instance-aliasing-probe-01"}, + }, + Chunks: []client.Chunk{ + { + StartTimestampMs: 1000, + EndTimestampMs: 2000, + Encoding: 3, + Data: chunkData, + }, + }, + }, + }, + } + return resp, append([]byte(nil), chunkData...) +} + +// overlaps reports whether byte slices a and b share any backing memory. +// +// This compares raw addresses as uintptrs, which is only valid because Go's +// garbage collector does not move heap-allocated objects (unlike goroutine +// stacks, which can be copied). If a future Go runtime introduces a moving +// GC for the heap, this address-range comparison would need to be redone +// (e.g. by writing sentinel bytes and checking for their presence/absence +// instead of comparing pointers). runtime.KeepAlive calls below only guard +// against the compiler treating a/b as dead before the address is taken; +// they say nothing about GC movement. +func overlaps(a, b []byte) bool { + if len(a) == 0 || len(b) == 0 { + return false + } + aStart := uintptr(unsafe.Pointer(&a[0])) + aEnd := aStart + uintptr(len(a)) + bStart := uintptr(unsafe.Pointer(&b[0])) + bEnd := bStart + uintptr(len(b)) + result := aStart < bEnd && bStart < aEnd + runtime.KeepAlive(a) + runtime.KeepAlive(b) + return result +} + +// unsafeBytesFromString exposes a string's backing bytes without copying, so +// that yoloString-produced strings can be checked for aliasing with the +// buffer they were decoded from. +func unsafeBytesFromString(s string) []byte { + if len(s) == 0 { + return nil + } + return unsafe.Slice(unsafe.StringData(s), len(s)) +} + +// TestChunkDataDoesNotAliasWireBuffer_ButLabelValueDoes pins the invariant +// behind the fix for #7732: a gogo Marshal -> Unmarshal round trip of +// client.QueryStreamResponse must leave Chunk.Data as an independent +// allocation (never aliasing the wire buffer), while LabelAdapter.Value must +// still alias it via yoloString. If a future change to the gogo-generated +// Chunk.Unmarshal (or a replacement decoder) ever starts aliasing Data, this +// test must fail loudly -- both the pointer-overlap check AND the mutation +// check below have to be defeated for that regression to slip through. +func TestChunkDataDoesNotAliasWireBuffer_ButLabelValueDoes(t *testing.T) { + // Size is irrelevant here: this test calls gogo's Marshal/Unmarshal + // directly, never going through cortexCodec's buffer-pooling threshold. + orig, wantChunkData := newAliasingProbeQueryStreamResponse(256) + + wireBytes, err := orig.Marshal() + require.NoError(t, err) + + got := &client.QueryStreamResponse{} + require.NoError(t, got.Unmarshal(wireBytes)) + require.Len(t, got.Chunkseries, 1) + require.Len(t, got.Chunkseries[0].Chunks, 1) + require.Len(t, got.Chunkseries[0].Labels, 2) + + chunkData := got.Chunkseries[0].Chunks[0].Data + labelValue := got.Chunkseries[0].Labels[1].Value + require.Equal(t, "instance-aliasing-probe-01", labelValue) + + // --- Pointer-range checks --- + assert.False(t, overlaps(chunkData, wireBytes), + "invariant broken: Chunk.Data now aliases the wire buffer -- if Unmarshal changed "+ + "to make this true, a detach/copy step must be reinstated in streamingSelect") + assert.True(t, overlaps(unsafeBytesFromString(labelValue), wireBytes), + "expected LabelAdapter.Value to still alias the wire buffer via yoloString; if this "+ + "is now false, FromLabelAdaptersToLabelsWithCopy may no longer be load-bearing "+ + "(but should still be kept -- do not remove without re-verifying this test)") + + // --- Behavioral (mutation) checks: corrupt the wire buffer in place and + // confirm the chunk is unaffected while the label (known to alias) is. --- + for i := range wireBytes { + wireBytes[i] = 0xFF + } + assert.Equal(t, wantChunkData, chunkData, + "Chunk.Data must be a private copy: it should be unaffected by mutating the wire buffer") + assert.Equal(t, bytes.Repeat([]byte{0xFF}, len(labelValue)), unsafeBytesFromString(labelValue), + "sanity check failed: LabelAdapter.Value should have observably changed after mutating "+ + "the wire buffer it aliases -- if this assertion itself fails, the test's mutation "+ + "methodology is broken and the non-aliasing assertion above cannot be trusted") +} + +// TestChunkDataDoesNotAliasWireBuffer_ViaRegisteredCodec re-runs the same +// invariant through the actual gRPC codec registered for Cortex's ingester +// connections (pkg/cortexpb.codec.go's cortexCodec, registered under the +// "proto" content-subtype via encoding.RegisterCodecV2). That codec uses +// grpc's mem.BufferSlice/mem.Buffer machinery. The payload is sized above +// mem.IsBelowBufferPoolingThreshold (1KiB) so that cortexCodec.Marshal takes +// its real pool.Get/mem.NewBuffer branch -- a smaller payload would silently +// fall back to a plain mem.SliceBuffer, whose Ref/Free are no-ops, and this +// test would then never actually touch a pool-backed buffer. +// +// This test confirms that, even through that pool-backed buffer, decoding +// still routes to the classic gogo Marshal/Unmarshal methods (via the +// protobuf-go legacy-message shim) so Chunk.Data is still never an alias of +// the wire buffer -- this is the real path used by QueryStream, not just the +// pb.go method in isolation. +// +// What this does NOT cover: a message split across multiple transport reads +// (mem.BufferSlice with len > 1), which drives MaterializeToBuffer's other +// branch (pool.Get + CopyTo instead of the single-buffer Ref fast path). A +// same-process Marshal->Unmarshal round trip can't produce that shape -- +// reaching it would need an actual (or bufconn-based) gRPC transport with a +// large enough message to span multiple HTTP/2 DATA frames, which is +// integration-test territory and out of scope here. +func TestChunkDataDoesNotAliasWireBuffer_ViaRegisteredCodec(t *testing.T) { + codec := encoding.GetCodecV2(cortexpb.Name) + require.NotNil(t, codec, "expected a CodecV2 to be registered under name %q", cortexpb.Name) + + // 4096 bytes of chunk data comfortably clears the 1KiB pooling threshold + // once labels and protobuf framing overhead are added. + orig, _ := newAliasingProbeQueryStreamResponse(4096) + + wireData, err := codec.Marshal(orig) + require.NoError(t, err) + require.NotEmpty(t, wireData) + require.Greater(t, wireData.Len(), 1024, + "test payload must exceed the buffer-pooling threshold or this test silently stops "+ + "exercising the pool-backed mem.Buffer code path (see mem.IsBelowBufferPoolingThreshold)") + + got := &client.QueryStreamResponse{} + require.NoError(t, codec.Unmarshal(wireData, got)) + require.Len(t, got.Chunkseries, 1) + require.Len(t, got.Chunkseries[0].Chunks, 1) + require.Len(t, got.Chunkseries[0].Labels, 2) + + chunkData := got.Chunkseries[0].Chunks[0].Data + labelValue := got.Chunkseries[0].Labels[1].Value + require.Equal(t, "instance-aliasing-probe-01", labelValue) + + labelBytes := unsafeBytesFromString(labelValue) + chunkAliasesWire := false + labelAliasesWire := false + for _, b := range wireData { + wireSegment := b.ReadOnlyData() + if overlaps(chunkData, wireSegment) { + chunkAliasesWire = true + } + if overlaps(labelBytes, wireSegment) { + labelAliasesWire = true + } + } + + assert.False(t, chunkAliasesWire, + "invariant broken: through the registered gRPC codec, Chunk.Data aliases a wire "+ + "mem.Buffer segment -- this would mean a zero-copy/pooled decode path now exists "+ + "and detachChunksFromBuffer (or equivalent) must be reinstated") + assert.True(t, labelAliasesWire, + "expected LabelAdapter.Value to alias a wire mem.Buffer segment via yoloString even "+ + "through the registered codec; if false, re-verify FromLabelAdaptersToLabelsWithCopy "+ + "is still necessary before considering it dead weight") +} From 6747cb474d17766334c83f63650f0341783d8e89 Mon Sep 17 00:00:00 2001 From: Sandy Chen Date: Mon, 3 Aug 2026 11:33:47 +0900 Subject: [PATCH 2/4] Amend #7670 CHANGELOG entry for #7746 Signed-off-by: Sandy Chen --- CHANGELOG.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 6cc35782015..442727e6297 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -47,7 +47,7 @@ * [ENHANCEMENT] Ring: Cache `ShuffleShardWithLookback` subrings. The cached entry is invalidated on topology change or once `now` reaches the earliest `RegisteredTimestamp + lookbackPeriod` of any included instance. #7628 * [ENHANCEMENT] Query Frontend: Rename `time_taken` field to `time_taken_ms` and make it return millisecond count. #7649 * [ENHANCEMENT] Update prometheus alertmanager version to v0.33.0. #7647 -* [ENHANCEMENT] Querier/Ingester: Detach ingester series from gRPC buffers to reduce heap. #7670 +* [ENHANCEMENT] Querier/Ingester: Detach ingester series labels from gRPC buffers to reduce heap. Chunk data is already detached by the gRPC unmarshal itself, so it is passed through without an extra copy. #7670 #7746 * [ENHANCEMENT] Ingester: Add `cortex_ingester_tsdb_head_max_timestamp` metric that re-exports the TSDB head max timestamp (`prometheus_tsdb_head_max_time`) per user, to help investigate ingestion issues like out-of-bounds (too old sample) errors. #7694 * [ENHANCEMENT] Ingester: Include the TSDB head max time in the `out of bounds` and `too old sample` error messages, so that users can see how far behind the accepted time range a rejected sample is. #7695 * [ENHANCEMENT] Compactor: Reduce object storage GET calls when updating the bucket index by skipping re-reading parquet converter markers for blocks that already have a valid-version parquet entry in the previous index. #7669 From 28039dfc39c711aa31a173e7d74ec66b224579d8 Mon Sep 17 00:00:00 2001 From: Sandy Chen Date: Mon, 3 Aug 2026 13:55:14 +0900 Subject: [PATCH 3/4] Retrigger CI (flaky TestExpandedPostingsCacheFuzz / TestExperimentalPromQLFuncsWithPrometheus / compactor unit test) Signed-off-by: Sandy Chen From 116f76a2a9aa8ae88db4a3329fa2d192f341bc88 Mon Sep 17 00:00:00 2001 From: Sandy Chen Date: Mon, 3 Aug 2026 14:17:07 +0900 Subject: [PATCH 4/4] Retrigger CI (unrelated pkg/compactor flakes: DeleteLocalSyncFiles race, FailedWithHaltError) Signed-off-by: Sandy Chen