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.go b/pkg/promtail.go index 122ff87a..3134fc30 100644 --- a/pkg/promtail.go +++ b/pkg/promtail.go @@ -63,12 +63,15 @@ 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 survive mutation downstream + 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..4f0b5b3d --- /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 = 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{ + 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) + } +}