From 9be1750a67934db380c624b10a58a52bb05ce586 Mon Sep 17 00:00:00 2001 From: Matthew Sheppard Date: Fri, 18 Sep 2026 12:09:18 -0400 Subject: [PATCH 1/7] Clone labels before pipeline processing to fix batch structured_metadata bug structured_metadata mutates Entry.Labels in place, deleting a label once it's promoted to structured metadata. parseCWEvent and processLogEvents build one labels map per invocation and pass it by reference to every entry in a batch, so the delete from the first entry was visible to every later entry sharing that map -- only the first entry in a batch kept the field. --- pkg/promtail.go | 11 ++++++-- pkg/promtail_test.go | 63 ++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 72 insertions(+), 2 deletions(-) create mode 100644 pkg/promtail_test.go diff --git a/pkg/promtail.go b/pkg/promtail.go index 122ff87a..8179e37f 100644 --- a/pkg/promtail.go +++ b/pkg/promtail.go @@ -63,12 +63,19 @@ func newBatch(ctx context.Context, pClient Client, processingPipeline *LokiStage func (b *batch) add(ctx context.Context, e entry) error { if b.processor.Size() > 0 { + // Clone e.labels before handing it to the pipeline. Callers (e.g. parseCWEvent, + // processLogEvents) share one labels map across every entry in a batch, but stages + // like structured_metadata mutate Entry.Labels in place -- deleting a label once it's + // promoted to structured metadata. Without cloning, that delete is visible to every + // other entry sharing the map, so only the first entry in a batch keeps the field. + labels := e.labels.Clone() + // Apply pipeline stages to entry stageEntry := stages.Entry{ Extracted: map[string]interface{}{}, - Entry: api.Entry{Labels: e.labels, Entry: e.entry}, + Entry: api.Entry{Labels: labels, Entry: e.entry}, } - for labelName, labelValue := range e.labels { + for labelName, labelValue := range labels { stageEntry.Extracted[string(labelName)] = string(labelValue) } stageEntry = b.processor.Process(stageEntry) diff --git a/pkg/promtail_test.go b/pkg/promtail_test.go new file mode 100644 index 00000000..4ac95437 --- /dev/null +++ b/pkg/promtail_test.go @@ -0,0 +1,63 @@ +package main + +import ( + "context" + "testing" + "time" + + "github.com/grafana/loki/v3/pkg/logproto" + "github.com/prometheus/common/model" + "github.com/stretchr/testify/require" + + "github.com/grafana/loki/pkg/push" +) + +// Test_batch_add_SharedLabelsAcrossBatch reproduces +// https://github.com/grafana/support-escalations/issues/24182: when a Lambda invocation +// receives multiple log events in one batch (e.g. a CloudWatch put-log-events call with several +// events), parseCWEvent/processLogEvents build one labels map and reuse it, by reference, for +// every entry{} in the batch. The structured_metadata stage mutates Entry.Labels in place +// (deleting a label once it's promoted to structured metadata), so without batch.add cloning +// e.labels first, only the first entry in the batch keeps the field -- every later entry's +// Extracted copy is built from an already-mutated map and never sees the label at all. +func Test_batch_add_SharedLabelsAcrossBatch(t *testing.T) { + pipeline, err := ParsePipelineConfigs( + `[{"structured_metadata":{"log_stream":"__aws_cloudwatch_log_stream"}}]`, + nil, nil, + ) + require.NoError(t, err) + + sharedLabels := model.LabelSet{ + model.LabelName("__aws_log_type"): model.LabelValue("cloudwatch"), + model.LabelName("__aws_cloudwatch_log_group"): model.LabelValue("testLogGroup"), + model.LabelName("__aws_cloudwatch_log_stream"): model.LabelValue("testLogStream"), + } + + b := &batch{ + streams: map[string]*logproto.Stream{}, + processor: pipeline, + } + + batchSize = 131072 // large enough that add() never flushes mid-test + + for i, line := range []string{"first event", "second event", "third event"} { + err := b.add(context.Background(), entry{sharedLabels, logproto.Entry{ + Line: line, + Timestamp: time.Now(), + }}) + require.NoError(t, err, "event %d", i) + } + + require.Len(t, b.streams, 1, "all three entries should share the same remaining labels, hence one stream") + + var stream *logproto.Stream + for _, s := range b.streams { + stream = s + } + require.Len(t, stream.Entries, 3) + + for i, e := range stream.Entries { + require.Containsf(t, e.StructuredMetadata, push.LabelAdapter{Name: "log_stream", Value: "testLogStream"}, + "entry %d (%q) is missing log_stream structured metadata", i, e.Line) + } +} From 9d37e93520d928ec194aa29be5e7e2b917bf313c Mon Sep 17 00:00:00 2001 From: Matthew Sheppard Date: Fri, 18 Sep 2026 12:54:52 -0400 Subject: [PATCH 2/7] Shorten batch.add clone comment --- pkg/promtail.go | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/pkg/promtail.go b/pkg/promtail.go index 8179e37f..74d87a32 100644 --- a/pkg/promtail.go +++ b/pkg/promtail.go @@ -63,11 +63,7 @@ func newBatch(ctx context.Context, pClient Client, processingPipeline *LokiStage func (b *batch) add(ctx context.Context, e entry) error { if b.processor.Size() > 0 { - // Clone e.labels before handing it to the pipeline. Callers (e.g. parseCWEvent, - // processLogEvents) share one labels map across every entry in a batch, but stages - // like structured_metadata mutate Entry.Labels in place -- deleting a label once it's - // promoted to structured metadata. Without cloning, that delete is visible to every - // other entry sharing the map, so only the first entry in a batch keeps the field. + // Clone to sruvive mutation downstream labels := e.labels.Clone() // Apply pipeline stages to entry From d6e28e31e5b772278a9bdbfb074ada571a61560d Mon Sep 17 00:00:00 2001 From: Matthew Sheppard Date: Fri, 18 Sep 2026 12:55:27 -0400 Subject: [PATCH 3/7] Fix typo in comment --- pkg/promtail.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/promtail.go b/pkg/promtail.go index 74d87a32..3134fc30 100644 --- a/pkg/promtail.go +++ b/pkg/promtail.go @@ -63,7 +63,7 @@ func newBatch(ctx context.Context, pClient Client, processingPipeline *LokiStage func (b *batch) add(ctx context.Context, e entry) error { if b.processor.Size() > 0 { - // Clone to sruvive mutation downstream + // Clone to survive mutation downstream labels := e.labels.Clone() // Apply pipeline stages to entry From f5f1ebed86335d7ca4fc12e6b5c7cfec1c0c7b4e Mon Sep 17 00:00:00 2001 From: Matthew Sheppard Date: Thu, 8 Oct 2026 14:19:58 -0400 Subject: [PATCH 4/7] document batchsize usage --- pkg/cw_test.go | 2 +- pkg/main.go | 4 +++- pkg/promtail_test.go | 2 +- 3 files changed, 5 insertions(+), 3 deletions(-) diff --git a/pkg/cw_test.go b/pkg/cw_test.go index 208b9d24..4e06587f 100644 --- a/pkg/cw_test.go +++ b/pkg/cw_test.go @@ -49,7 +49,7 @@ func Test_parseCWEvent(t *testing.T) { } t.Run(tt.name, func(t *testing.T) { - batchSize = 131072 // Set large enough we don't send to promtail + batchSize = defaultBatchSize // Set large enough we don't send to promtail keepStream = tt.keepStream err := parseCWEvent(context.Background(), tt.b, cwevent) if err != nil { diff --git a/pkg/main.go b/pkg/main.go index e2cdabd3..7bb89421 100644 --- a/pkg/main.go +++ b/pkg/main.go @@ -29,6 +29,8 @@ const ( maxErrMsgLen = 1024 invalidExtraLabelsError = "invalid value for environment variable EXTRA_LABELS. Expected a comma separated list with an even number of entries. " + + defaultBatchSize = 131062 // see BATCH_SIZE in doc/sources/lambda-promtail-reference.md ) var ( @@ -109,7 +111,7 @@ func setupArguments(ctx context.Context, secretFetcher secretFetcher) { fmt.Println("keep stream: ", keepStream) batch := os.Getenv("BATCH_SIZE") - batchSize = 131072 + batchSize = defaultBatchSize if batch != "" { batchSize, _ = strconv.Atoi(batch) } diff --git a/pkg/promtail_test.go b/pkg/promtail_test.go index 4ac95437..4f0b5b3d 100644 --- a/pkg/promtail_test.go +++ b/pkg/promtail_test.go @@ -38,7 +38,7 @@ func Test_batch_add_SharedLabelsAcrossBatch(t *testing.T) { processor: pipeline, } - batchSize = 131072 // large enough that add() never flushes mid-test + batchSize = defaultBatchSize // large enough that add() never flushes mid-test for i, line := range []string{"first event", "second event", "third event"} { err := b.add(context.Background(), entry{sharedLabels, logproto.Entry{ From b3315829d55f84a948fd44cf3fd6adfa6d16ce1c Mon Sep 17 00:00:00 2001 From: Matthew Sheppard Date: Fri, 9 Oct 2026 13:10:50 -0400 Subject: [PATCH 5/7] remove external tool reference --- pkg/promtail_test.go | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/pkg/promtail_test.go b/pkg/promtail_test.go index 4f0b5b3d..036b4a8c 100644 --- a/pkg/promtail_test.go +++ b/pkg/promtail_test.go @@ -12,11 +12,10 @@ import ( "github.com/grafana/loki/pkg/push" ) -// Test_batch_add_SharedLabelsAcrossBatch reproduces -// https://github.com/grafana/support-escalations/issues/24182: when a Lambda invocation -// receives multiple log events in one batch (e.g. a CloudWatch put-log-events call with several -// events), parseCWEvent/processLogEvents build one labels map and reuse it, by reference, for -// every entry{} in the batch. The structured_metadata stage mutates Entry.Labels in place +// Test_batch_add_SharedLabelsAcrossBatch reproduces a bug where, when a Lambda invocation +// receives multiple log events in one batch (e.g. a CloudWatch put-log-events call with +// several events), parseCWEvent/processLogEvents build one labels map and reuse it, by +// reference, for every entry{} in the batch. The structured_metadata stage mutates Entry.Labels in place // (deleting a label once it's promoted to structured metadata), so without batch.add cloning // e.labels first, only the first entry in the batch keeps the field -- every later entry's // Extracted copy is built from an already-mutated map and never sees the label at all. From b5b402442482f8375d715947ff7743f0a3ebfbdf Mon Sep 17 00:00:00 2001 From: Matthew Sheppard Date: Fri, 9 Oct 2026 13:13:21 -0400 Subject: [PATCH 6/7] correct batch size --- pkg/main.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/main.go b/pkg/main.go index 7bb89421..eaf66344 100644 --- a/pkg/main.go +++ b/pkg/main.go @@ -30,7 +30,7 @@ const ( invalidExtraLabelsError = "invalid value for environment variable EXTRA_LABELS. Expected a comma separated list with an even number of entries. " - defaultBatchSize = 131062 // see BATCH_SIZE in doc/sources/lambda-promtail-reference.md + defaultBatchSize = 131072 // see BATCH_SIZE in doc/sources/lambda-promtail-reference.md ) var ( From 2ecd411dc76f89e85a0d0e7eab69a1b2ac015e28 Mon Sep 17 00:00:00 2001 From: Matthew Sheppard Date: Fri, 9 Oct 2026 13:35:52 -0400 Subject: [PATCH 7/7] check cloned labels are NOT mutated in test --- pkg/promtail_test.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/pkg/promtail_test.go b/pkg/promtail_test.go index 036b4a8c..3836437f 100644 --- a/pkg/promtail_test.go +++ b/pkg/promtail_test.go @@ -39,12 +39,15 @@ func Test_batch_add_SharedLabelsAcrossBatch(t *testing.T) { batchSize = defaultBatchSize // large enough that add() never flushes mid-test + want := sharedLabels.Clone() + for i, line := range []string{"first event", "second event", "third event"} { err := b.add(context.Background(), entry{sharedLabels, logproto.Entry{ Line: line, Timestamp: time.Now(), }}) require.NoError(t, err, "event %d", i) + require.Equal(t, want, sharedLabels, "batch.add mutated caller's labels on event %d", i) } require.Len(t, b.streams, 1, "all three entries should share the same remaining labels, hence one stream")