From e3b29555f1d16ff58e256ab8d3cfe238ba6fbf31 Mon Sep 17 00:00:00 2001 From: Kevin Shi Date: Thu, 15 Jan 2026 16:10:20 -0800 Subject: [PATCH 1/6] Fix single large event handling for DynamoDB backend --- lib/events/dynamoevents/dynamoevents.go | 67 ++++++++++++++- lib/events/dynamoevents/dynamoevents_test.go | 85 ++++++++++++++++++++ 2 files changed, 150 insertions(+), 2 deletions(-) diff --git a/lib/events/dynamoevents/dynamoevents.go b/lib/events/dynamoevents/dynamoevents.go index b13789be6520e..0813af1c70fcd 100644 --- a/lib/events/dynamoevents/dynamoevents.go +++ b/lib/events/dynamoevents/dynamoevents.go @@ -42,6 +42,7 @@ import ( dynamodbtypes "github.com/aws/aws-sdk-go-v2/service/dynamodb/types" "github.com/aws/smithy-go" "github.com/aws/smithy-go/tracing/smithyoteltracing" + "github.com/dustin/go-humanize" "github.com/google/uuid" "github.com/gravitational/trace" "github.com/jonboulle/clockwork" @@ -1621,10 +1622,53 @@ func (l *eventsFetcher) processQueryOutput(output *dynamodb.QueryOutput) ([]even // Stop early when the fetcher's total size exceeds the response size limit. if l.totalSize+len(data) >= events.MaxEventBytesInResponse { - if err := l.saveCheckpointAtEvent(out[len(out)-1]); err != nil { + // Encountered an event that would push the total page over the size limit. + // Return all processed events, and the next event will be picked up on the next page. + if len(out) > 0 { + if err := l.saveCheckpointAtEvent(out[len(out)-1]); err != nil { + return nil, false, trace.Wrap(err) + } + return out, true, nil + } + + // A single event is larger than the max page size - the best we can + // do is try to trim it. + e.FieldsMap, err = trimToMaxSize(e.FieldsMap) + if err != nil { + return nil, false, trace.Wrap(err, "failed to trim event to max size") + } + trimmedData, err := json.Marshal(e.FieldsMap) + if err != nil { return nil, false, trace.Wrap(err) } - return out, true, nil + + if l.totalSize+len(trimmedData) <= events.MaxEventBytesInResponse { + events.MetricQueriedTrimmedEvents.Inc() + l.totalSize += len(trimmedData) + out = append(out, e) + l.left-- + + // Since we reach the response size limit, simply return the trimmed event. + if err := l.saveCheckpointAtEvent(out[len(out)-1]); err != nil { + return nil, false, trace.Wrap(err) + } + return out, true, nil + } + + // Failed to trim the event to size. + // If this condition is reached it should be considered a bug, any + // event that can possibly exceed the maximum size should implement + // TrimToMaxSize (until we can one day implement an API for storing + // and retrieving large events). + l.log.ErrorContext(context.Background(), "Failed to query event exceeding maximum response size.", + "event_type", e.FieldsMap.GetType(), + "event_id", e.FieldsMap.GetID(), + "event_size", len(data), + ) + return nil, false, trace.Errorf( + "%s event %s is %s and cannot be returned because it exceeds the maximum response size of %s", + e.FieldsMap.GetType(), e.FieldsMap.GetID(), humanize.IBytes(uint64(len(data))), humanize.IBytes(events.MaxEventBytesInResponse)) + } l.totalSize += len(data) out = append(out, e) @@ -1640,6 +1684,25 @@ func (l *eventsFetcher) processQueryOutput(output *dynamodb.QueryOutput) ([]even return out, false, nil } +// trimToMaxSize attempts to trim the event to fit into the maximum response size (MaxEventBytesInResponse). +// If the event is larger than the maximum response size, it will be trimmed +// to the maximum size, which may result in loss of data. +// Trimming requires unmarshalling the event to apievents.AuditEvent and then +// calling TrimToMaxSize on it. +// This is not an efficient operation, but it is executed once per page, +// so it should not be a problem in practice. +func trimToMaxSize(fields events.EventFields) (events.EventFields, error) { + event, err := events.FromEventFields(fields) + if err != nil { + return nil, trace.Wrap(err) + } + + event = event.TrimToMaxSize(events.MaxEventBytesInResponse) + + fields, err = events.ToEventFields(event) + return fields, trace.Wrap(err) +} + // saveCheckpointAtEvent updates the checkpoint iterator at the given event. // This overrides LastEvaluatedKey to resume future processing from this iterator. func (l *eventsFetcher) saveCheckpointAtEvent(e event) error { diff --git a/lib/events/dynamoevents/dynamoevents_test.go b/lib/events/dynamoevents/dynamoevents_test.go index 0d5110b4d8ed9..9881e1625da54 100644 --- a/lib/events/dynamoevents/dynamoevents_test.go +++ b/lib/events/dynamoevents/dynamoevents_test.go @@ -37,6 +37,7 @@ import ( "github.com/aws/aws-sdk-go-v2/feature/dynamodb/attributevalue" "github.com/aws/aws-sdk-go-v2/service/dynamodb" dynamodbtypes "github.com/aws/aws-sdk-go-v2/service/dynamodb/types" + "github.com/dustin/go-humanize" "github.com/google/go-cmp/cmp" "github.com/google/go-cmp/cmp/cmpopts" "github.com/google/uuid" @@ -1027,11 +1028,31 @@ func Test_eventsFetcher_QueryByDateIndex(t *testing.T) { AppName: "app-4", }, } + bigUntrimmableEvent := &apievents.AppCreate{ + Metadata: apievents.Metadata{ + ID: uuid.NewString(), + Time: time.Now().UTC(), + Type: events.AppCreateEvent, + }, + AppMetadata: apievents.AppMetadata{ + AppName: strings.Repeat("aaaaa", events.MaxEventBytesInResponse), + }, + } + bigTrimmableEvent := &apievents.DatabaseSessionQuery{ + Metadata: apievents.Metadata{ + ID: uuid.NewString(), + Time: time.Now().UTC(), + Type: events.DatabaseSessionQueryEvent, + }, + DatabaseQuery: strings.Repeat("aaaaa", events.MaxEventBytesInResponse), + } + bigTrimmedEvent := bigTrimmableEvent.TrimToMaxSize(events.MaxEventBytesInResponse) // have a deterministic session ID (UID) when used in test cases key1 := eventToKey(event1) key3 := eventToKey(event3) key4 := eventToKey(event4) + keyTrimmed := eventToKey(bigTrimmedEvent) tests := []struct { name string @@ -1039,6 +1060,7 @@ func Test_eventsFetcher_QueryByDateIndex(t *testing.T) { mockResponses map[EventKey]mockResponse wantEvents []apievents.AuditEvent wantKey *EventKey + wantErrorMsg string }{ { name: "no data returned from query, return empty results", @@ -1096,6 +1118,65 @@ func Test_eventsFetcher_QueryByDateIndex(t *testing.T) { wantEvents: []apievents.AuditEvent{event1, event2, event3}, wantKey: &key3, }, + { + name: "events with big untrimmable event exceeding > MaxEventBytesInResponse", + limit: 10, + mockResponses: map[EventKey]mockResponse{ + {}: { + events: []apievents.AuditEvent{event1}, + returnKey: &key1, + }, + key1: { + events: []apievents.AuditEvent{event2, event3, bigUntrimmableEvent}, + returnKey: nil, + }, + }, + // we don't expect bigUntrimmableEvent because it should go to next batch + wantEvents: []apievents.AuditEvent{event1, event2, event3}, + wantKey: &key3, + }, + { + name: "only 1 big untrimmable event", + limit: 10, + mockResponses: map[EventKey]mockResponse{ + {}: { + events: []apievents.AuditEvent{bigUntrimmableEvent}, + returnKey: nil, + }, + }, + wantErrorMsg: fmt.Sprintf( + "app.create event %s is 5.0 MiB and cannot be returned because it exceeds the maximum response size of %s", + bigUntrimmableEvent.Metadata.ID, humanize.IBytes(events.MaxEventBytesInResponse)), + }, + { + name: "events with big trimmable event exceeding > MaxEventBytesInResponse", + limit: 10, + mockResponses: map[EventKey]mockResponse{ + {}: { + events: []apievents.AuditEvent{event1}, + returnKey: &key1, + }, + key1: { + events: []apievents.AuditEvent{event2, event3, bigUntrimmableEvent}, + returnKey: nil, + }, + }, + // we don't expect bigTrimmedEvent because it should go to next batch + wantEvents: []apievents.AuditEvent{event1, event2, event3}, + wantKey: &key3, + }, + { + name: "only 1 big trimmable event", + limit: 10, + mockResponses: map[EventKey]mockResponse{ + {}: { + events: []apievents.AuditEvent{bigTrimmableEvent}, + returnKey: nil, + }, + }, + wantEvents: []apievents.AuditEvent{bigTrimmedEvent}, + wantKey: &keyTrimmed, + }, } for _, test := range tests { @@ -1116,6 +1197,10 @@ func Test_eventsFetcher_QueryByDateIndex(t *testing.T) { } gotRawEvents, err := ef.QueryByDateIndex(t.Context(), getExprFilter(ef.filter)) + if test.wantErrorMsg != "" { + require.ErrorContains(t, err, test.wantErrorMsg) + return + } require.NoError(t, err) if test.wantKey != nil { From 32bccf474c9cd13cbbaf68797c3b4cb40e669d51 Mon Sep 17 00:00:00 2001 From: Kevin Shi Date: Fri, 16 Jan 2026 14:28:27 -0800 Subject: [PATCH 2/6] Change inequality sign from >= to > --- lib/events/dynamoevents/dynamoevents.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/lib/events/dynamoevents/dynamoevents.go b/lib/events/dynamoevents/dynamoevents.go index 0813af1c70fcd..09184158a2010 100644 --- a/lib/events/dynamoevents/dynamoevents.go +++ b/lib/events/dynamoevents/dynamoevents.go @@ -1621,7 +1621,7 @@ func (l *eventsFetcher) processQueryOutput(output *dynamodb.QueryOutput) ([]even } // Stop early when the fetcher's total size exceeds the response size limit. - if l.totalSize+len(data) >= events.MaxEventBytesInResponse { + if l.totalSize+len(data) > events.MaxEventBytesInResponse { // Encountered an event that would push the total page over the size limit. // Return all processed events, and the next event will be picked up on the next page. if len(out) > 0 { From 6d5f9c1578c818d1bd794e6fa415b3cb749a0113 Mon Sep 17 00:00:00 2001 From: Kevin Shi Date: Tue, 20 Jan 2026 10:38:42 -0800 Subject: [PATCH 3/6] Fix tests, add comments --- lib/events/athena/querier.go | 3 ++- lib/events/dynamoevents/dynamoevents.go | 3 ++- lib/events/dynamoevents/dynamoevents_test.go | 6 +++--- 3 files changed, 7 insertions(+), 5 deletions(-) diff --git a/lib/events/athena/querier.go b/lib/events/athena/querier.go index fa28875d1c3e3..7024af6ba7fea 100644 --- a/lib/events/athena/querier.go +++ b/lib/events/athena/querier.go @@ -1030,7 +1030,8 @@ func (rb *responseBuilder) appendUntilSizeLimit(resultResp *athena.GetQueryResul // to the maximum size, which may result in loss of data. // Trimming requires unmarshalling the event to audit.Event and then // calling TrimToMaxSize on it. -// This is not an efficient operation, but it is executed once per page, +// This is not an efficient operation, but it is executed at most once per page, +// and only when a single event exceeds the limit, // so it should not be a problem in practice. func trimToMaxSize(fields events.EventFields) (events.EventFields, error) { event, err := events.FromEventFields(fields) diff --git a/lib/events/dynamoevents/dynamoevents.go b/lib/events/dynamoevents/dynamoevents.go index 09184158a2010..ee4cc30684e91 100644 --- a/lib/events/dynamoevents/dynamoevents.go +++ b/lib/events/dynamoevents/dynamoevents.go @@ -1689,7 +1689,8 @@ func (l *eventsFetcher) processQueryOutput(output *dynamodb.QueryOutput) ([]even // to the maximum size, which may result in loss of data. // Trimming requires unmarshalling the event to apievents.AuditEvent and then // calling TrimToMaxSize on it. -// This is not an efficient operation, but it is executed once per page, +// This is not an efficient operation, but it is executed at most once per page, +// and only when a single event exceeds the limit, // so it should not be a problem in practice. func trimToMaxSize(fields events.EventFields) (events.EventFields, error) { event, err := events.FromEventFields(fields) diff --git a/lib/events/dynamoevents/dynamoevents_test.go b/lib/events/dynamoevents/dynamoevents_test.go index 9881e1625da54..5686e216e4caa 100644 --- a/lib/events/dynamoevents/dynamoevents_test.go +++ b/lib/events/dynamoevents/dynamoevents_test.go @@ -1031,7 +1031,7 @@ func Test_eventsFetcher_QueryByDateIndex(t *testing.T) { bigUntrimmableEvent := &apievents.AppCreate{ Metadata: apievents.Metadata{ ID: uuid.NewString(), - Time: time.Now().UTC(), + Time: time.Date(2025, 2, 5, 0, 0, 0, 0, time.UTC), Type: events.AppCreateEvent, }, AppMetadata: apievents.AppMetadata{ @@ -1041,7 +1041,7 @@ func Test_eventsFetcher_QueryByDateIndex(t *testing.T) { bigTrimmableEvent := &apievents.DatabaseSessionQuery{ Metadata: apievents.Metadata{ ID: uuid.NewString(), - Time: time.Now().UTC(), + Time: time.Date(2025, 2, 5, 0, 0, 0, 0, time.UTC), Type: events.DatabaseSessionQueryEvent, }, DatabaseQuery: strings.Repeat("aaaaa", events.MaxEventBytesInResponse), @@ -1157,7 +1157,7 @@ func Test_eventsFetcher_QueryByDateIndex(t *testing.T) { returnKey: &key1, }, key1: { - events: []apievents.AuditEvent{event2, event3, bigUntrimmableEvent}, + events: []apievents.AuditEvent{event2, event3, bigTrimmableEvent}, returnKey: nil, }, }, From 6d99759592f8fd90991b737da8003d13f4e1aee7 Mon Sep 17 00:00:00 2001 From: Kevin Shi Date: Tue, 27 Jan 2026 16:09:54 -0800 Subject: [PATCH 4/6] Modify to serve the oversized event and bypass the limit --- lib/events/dynamoevents/dynamoevents.go | 41 ++++++++------------ lib/events/dynamoevents/dynamoevents_test.go | 13 ++----- 2 files changed, 21 insertions(+), 33 deletions(-) diff --git a/lib/events/dynamoevents/dynamoevents.go b/lib/events/dynamoevents/dynamoevents.go index ee4cc30684e91..db1c9f5c4f9a0 100644 --- a/lib/events/dynamoevents/dynamoevents.go +++ b/lib/events/dynamoevents/dynamoevents.go @@ -42,7 +42,6 @@ import ( dynamodbtypes "github.com/aws/aws-sdk-go-v2/service/dynamodb/types" "github.com/aws/smithy-go" "github.com/aws/smithy-go/tracing/smithyoteltracing" - "github.com/dustin/go-humanize" "github.com/google/uuid" "github.com/gravitational/trace" "github.com/jonboulle/clockwork" @@ -1644,31 +1643,25 @@ func (l *eventsFetcher) processQueryOutput(output *dynamodb.QueryOutput) ([]even if l.totalSize+len(trimmedData) <= events.MaxEventBytesInResponse { events.MetricQueriedTrimmedEvents.Inc() - l.totalSize += len(trimmedData) - out = append(out, e) - l.left-- - - // Since we reach the response size limit, simply return the trimmed event. - if err := l.saveCheckpointAtEvent(out[len(out)-1]); err != nil { - return nil, false, trace.Wrap(err) - } - return out, true, nil + } else { + // Failed to trim the event to size. + // Even if we fail to trim the event, we still try to return the oversized event. + l.log.ErrorContext(context.Background(), "Failed to trim event exceeding maximum response size.", + "event_type", e.FieldsMap.GetType(), + "event_id", e.FieldsMap.GetID(), + "event_size", len(data), + "event_size after trim", len(trimmedData), + ) } + l.totalSize += len(trimmedData) + out = append(out, e) + l.left-- - // Failed to trim the event to size. - // If this condition is reached it should be considered a bug, any - // event that can possibly exceed the maximum size should implement - // TrimToMaxSize (until we can one day implement an API for storing - // and retrieving large events). - l.log.ErrorContext(context.Background(), "Failed to query event exceeding maximum response size.", - "event_type", e.FieldsMap.GetType(), - "event_id", e.FieldsMap.GetID(), - "event_size", len(data), - ) - return nil, false, trace.Errorf( - "%s event %s is %s and cannot be returned because it exceeds the maximum response size of %s", - e.FieldsMap.GetType(), e.FieldsMap.GetID(), humanize.IBytes(uint64(len(data))), humanize.IBytes(events.MaxEventBytesInResponse)) - + // Since we reached the response size limit, simply return the event. + if err := l.saveCheckpointAtEvent(out[len(out)-1]); err != nil { + return nil, false, trace.Wrap(err) + } + return out, true, nil } l.totalSize += len(data) out = append(out, e) diff --git a/lib/events/dynamoevents/dynamoevents_test.go b/lib/events/dynamoevents/dynamoevents_test.go index 5686e216e4caa..ca46fc8d8d243 100644 --- a/lib/events/dynamoevents/dynamoevents_test.go +++ b/lib/events/dynamoevents/dynamoevents_test.go @@ -37,7 +37,6 @@ import ( "github.com/aws/aws-sdk-go-v2/feature/dynamodb/attributevalue" "github.com/aws/aws-sdk-go-v2/service/dynamodb" dynamodbtypes "github.com/aws/aws-sdk-go-v2/service/dynamodb/types" - "github.com/dustin/go-humanize" "github.com/google/go-cmp/cmp" "github.com/google/go-cmp/cmp/cmpopts" "github.com/google/uuid" @@ -1052,6 +1051,7 @@ func Test_eventsFetcher_QueryByDateIndex(t *testing.T) { key1 := eventToKey(event1) key3 := eventToKey(event3) key4 := eventToKey(event4) + keyUntrimmable := eventToKey(bigUntrimmableEvent) keyTrimmed := eventToKey(bigTrimmedEvent) tests := []struct { @@ -1060,7 +1060,6 @@ func Test_eventsFetcher_QueryByDateIndex(t *testing.T) { mockResponses map[EventKey]mockResponse wantEvents []apievents.AuditEvent wantKey *EventKey - wantErrorMsg string }{ { name: "no data returned from query, return empty results", @@ -1144,9 +1143,9 @@ func Test_eventsFetcher_QueryByDateIndex(t *testing.T) { returnKey: nil, }, }, - wantErrorMsg: fmt.Sprintf( - "app.create event %s is 5.0 MiB and cannot be returned because it exceeds the maximum response size of %s", - bigUntrimmableEvent.Metadata.ID, humanize.IBytes(events.MaxEventBytesInResponse)), + // we still want to receive the untrimmable event + wantEvents: []apievents.AuditEvent{bigUntrimmableEvent}, + wantKey: &keyUntrimmable, }, { name: "events with big trimmable event exceeding > MaxEventBytesInResponse", @@ -1197,10 +1196,6 @@ func Test_eventsFetcher_QueryByDateIndex(t *testing.T) { } gotRawEvents, err := ef.QueryByDateIndex(t.Context(), getExprFilter(ef.filter)) - if test.wantErrorMsg != "" { - require.ErrorContains(t, err, test.wantErrorMsg) - return - } require.NoError(t, err) if test.wantKey != nil { From c48f7ab1bb6ab0d71538c7b831ff612aa2335aed Mon Sep 17 00:00:00 2001 From: Kevin Shi Date: Thu, 29 Jan 2026 15:36:42 -0800 Subject: [PATCH 5/6] Fix minor typos --- lib/events/dynamoevents/dynamoevents.go | 2 +- lib/events/dynamoevents/dynamoevents_test.go | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/lib/events/dynamoevents/dynamoevents.go b/lib/events/dynamoevents/dynamoevents.go index db1c9f5c4f9a0..ba0dfbe18d0d3 100644 --- a/lib/events/dynamoevents/dynamoevents.go +++ b/lib/events/dynamoevents/dynamoevents.go @@ -1650,7 +1650,7 @@ func (l *eventsFetcher) processQueryOutput(output *dynamodb.QueryOutput) ([]even "event_type", e.FieldsMap.GetType(), "event_id", e.FieldsMap.GetID(), "event_size", len(data), - "event_size after trim", len(trimmedData), + "event_size_after_trim", len(trimmedData), ) } l.totalSize += len(trimmedData) diff --git a/lib/events/dynamoevents/dynamoevents_test.go b/lib/events/dynamoevents/dynamoevents_test.go index ca46fc8d8d243..6eec968abc110 100644 --- a/lib/events/dynamoevents/dynamoevents_test.go +++ b/lib/events/dynamoevents/dynamoevents_test.go @@ -1160,7 +1160,7 @@ func Test_eventsFetcher_QueryByDateIndex(t *testing.T) { returnKey: nil, }, }, - // we don't expect bigTrimmedEvent because it should go to next batch + // we don't expect bigTrimmableEvent because it should go to next batch wantEvents: []apievents.AuditEvent{event1, event2, event3}, wantKey: &key3, }, From 074274ac07bb2d0a493ee7c8de2cd8d2cf7539e5 Mon Sep 17 00:00:00 2001 From: Kevin Shi Date: Tue, 3 Feb 2026 13:57:30 -0800 Subject: [PATCH 6/6] Address feedback --- lib/events/dynamoevents/dynamoevents.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/lib/events/dynamoevents/dynamoevents.go b/lib/events/dynamoevents/dynamoevents.go index ba0dfbe18d0d3..2d3bca38679b3 100644 --- a/lib/events/dynamoevents/dynamoevents.go +++ b/lib/events/dynamoevents/dynamoevents.go @@ -1641,18 +1641,18 @@ func (l *eventsFetcher) processQueryOutput(output *dynamodb.QueryOutput) ([]even return nil, false, trace.Wrap(err) } - if l.totalSize+len(trimmedData) <= events.MaxEventBytesInResponse { - events.MetricQueriedTrimmedEvents.Inc() - } else { + if l.totalSize+len(trimmedData) > events.MaxEventBytesInResponse { // Failed to trim the event to size. // Even if we fail to trim the event, we still try to return the oversized event. - l.log.ErrorContext(context.Background(), "Failed to trim event exceeding maximum response size.", + l.log.WarnContext(context.Background(), "Failed to trim event exceeding maximum response size.", "event_type", e.FieldsMap.GetType(), "event_id", e.FieldsMap.GetID(), "event_size", len(data), "event_size_after_trim", len(trimmedData), ) } + events.MetricQueriedTrimmedEvents.Inc() + l.totalSize += len(trimmedData) out = append(out, e) l.left--