Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion pkg/cw_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
4 changes: 3 additions & 1 deletion pkg/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down Expand Up @@ -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)
}
Expand Down
7 changes: 5 additions & 2 deletions pkg/promtail.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
63 changes: 63 additions & 0 deletions pkg/promtail_test.go
Original file line number Diff line number Diff line change
@@ -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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

parseCWEvent/processLogEvents build one labels map and reuse it, by reference

Should parseCWEvent/processLogEvents be making the map copy for entry{}s in the first place then instead of setting up a reference only for it to be copied downstream? The semantics are a bit cleaner that way IMO

Alternatively, is it possible/would it make sense to have processing fix up the labels map before hand and have a separate collection for structured metadata entries, and pass both of those collections by reference?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

There are 4 places we'd have to do that — cw.go, kinesis.go, and s3.go (which has two). I think doing it once here is cleaner. I also think the user can define their own processing pipeline, and we can't guarantee they won't mutate the entries passed to it.

// 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

@tristanburgess tristanburgess Sep 18, 2026 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

only the first entry in the batch keeps the field

Why would this be true if every entry{} in the batch shares the same reference? Wouldn't they all see the mutated map?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is due to processing order. Entries go through the batch.add process one at a time, and the first one is correctly processed (label data extracted) and promoted to metadata and then that single label reference was deleted from the original e.labels labelset.

Entries 2 and 3 point at that exact same map — when they're processed next, the key's already gone, so the promotion check (e.Labels[source]) comes up empty and they never get the label added to StructuredMetadata.

// 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)
}
}
Loading