Repository navigation
Clone labels before pipeline processing to fix batch structured_metadata bug #218
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| 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 | ||
| // 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 | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Why would this be true if every entry{} in the batch shares the same reference? Wouldn't they all see the mutated map?
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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) | ||
| } | ||
| } | ||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Should
parseCWEvent/processLogEventsbe making the map copy forentry{}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 IMOAlternatively, 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?
There was a problem hiding this comment.
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.