From d4c5c90ab6d89065ea64026613c191f58465e747 Mon Sep 17 00:00:00 2001 From: Nabil Dakkoune Date: Mon, 20 Jul 2026 15:01:13 +0200 Subject: [PATCH 1/5] feat(go-forwarder): add minimal lambda enrichment --- .../cloudwatch_lambda_coldstart.golden.json | 21 ++++-- ...dwatch_lambda_custom_log_group.golden.json | 21 ++++-- .../cloudwatch_lambda_timeout.golden.json | 28 +++++--- .../internal/handling/cloudwatch.go | 58 ++++++++++++++-- .../internal/handling/cloudwatch_test.go | 66 +++++++++++++++++-- .../internal/model/lambda.go | 10 +++ aws/logs_monitoring_go/internal/model/log.go | 19 +++--- 7 files changed, 186 insertions(+), 37 deletions(-) create mode 100644 aws/logs_monitoring_go/internal/model/lambda.go diff --git a/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_coldstart.golden.json b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_coldstart.golden.json index 3fb82499b..635e3b709 100644 --- a/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_coldstart.golden.json +++ b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_coldstart.golden.json @@ -10,14 +10,14 @@ }, "body": [ { - "host": "/aws/lambda/storms-cloudwatch-event", + "host": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event", "id": "35486831490800643125153606102923171443962457178576257024", "timestamp": 1591284559098, "message": "START RequestId: db275f87-a934-471a-8980-b63bf4dc1beb Version: $LATEST\\n", "service": "lambda", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "", + "ddtags": "functionname:storms-cloudwatch-event,env:none", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -25,17 +25,20 @@ "logStream": "2020/06/04/[$LATEST]af2b1e1843b84a2d80c67840ae3ffa72", "owner": "601427279990" } + }, + "lambda": { + "arn": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event" } }, { - "host": "/aws/lambda/storms-cloudwatch-event", + "host": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event", "id": "35486831490867545360749197972347778598780402263094198273", "timestamp": 1591284559101, "message": "END RequestId: db275f87-a934-471a-8980-b63bf4dc1beb\\n", "service": "lambda", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "", + "ddtags": "functionname:storms-cloudwatch-event,env:none", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -43,17 +46,20 @@ "logStream": "2020/06/04/[$LATEST]af2b1e1843b84a2d80c67840ae3ffa72", "owner": "601427279990" } + }, + "lambda": { + "arn": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event" } }, { - "host": "/aws/lambda/storms-cloudwatch-event", + "host": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event", "id": "35486831490867545360749197972347778598780402263094198274", "timestamp": 1591284559101, "message": "REPORT RequestId: db275f87-a934-471a-8980-b63bf4dc1beb\\tDuration: 1.76 ms\\tBilled Duration: 100 ms\\tMemory Size: 128 MB\\tMax Memory Used: 48 MB\\tInit Duration: 120.96 ms\\t\\n", "service": "lambda", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "", + "ddtags": "functionname:storms-cloudwatch-event,env:none", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -61,6 +67,9 @@ "logStream": "2020/06/04/[$LATEST]af2b1e1843b84a2d80c67840ae3ffa72", "owner": "601427279990" } + }, + "lambda": { + "arn": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event" } } ] diff --git a/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_custom_log_group.golden.json b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_custom_log_group.golden.json index 8890c817d..7a642cf6a 100644 --- a/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_custom_log_group.golden.json +++ b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_custom_log_group.golden.json @@ -10,14 +10,14 @@ }, "body": [ { - "host": "/aws/vendedlogs/states/anyLogGroupName", + "host": "arn:aws:lambda:us-east-1:123456789012:function:test-customized-loggroup", "id": "35311576111948622874033876462979853992919938886093242368", "timestamp": 1583425836114, "message": "2020-03-05T16:30:36.113Z\tf08bb4c8-d6b2-4f05-ac17-af7e2ba005fb\tDEBUG\t[dd.trace_id=3172564172058669914 dd.span_id=14292093692483532556] {\"status\":\"debug\",\"message\":\"datadog:Patched console output with trace context\"}\n", "service": "lambda", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "", + "ddtags": "functionname:test-customized-loggroup,env:none", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -25,17 +25,20 @@ "logStream": "2020/03/05/test-customized-loggroup[$LATEST]20bddfd5a2dc4c6b97ac02800eae90d0", "owner": "601427279990" } + }, + "lambda": { + "arn": "arn:aws:lambda:us-east-1:123456789012:function:test-customized-loggroup" } }, { - "host": "/aws/vendedlogs/states/anyLogGroupName", + "host": "arn:aws:lambda:us-east-1:123456789012:function:test-customized-loggroup", "id": "35311576111948622874033876462979853992919938886093242369", "timestamp": 1583425836114, "message": "2020-03-05T16:30:36.114Z\tf08bb4c8-d6b2-4f05-ac17-af7e2ba005fb\tDEBUG\t[dd.trace_id=3172564172058669914 dd.span_id=14292093692483532556] {\"autoPatchHTTP\":true,\"tracerInitialized\":true,\"status\":\"debug\",\"message\":\"datadog:Not patching HTTP libraries\"}\n", "service": "lambda", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "", + "ddtags": "functionname:test-customized-loggroup,env:none", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -43,17 +46,20 @@ "logStream": "2020/03/05/test-customized-loggroup[$LATEST]20bddfd5a2dc4c6b97ac02800eae90d0", "owner": "601427279990" } + }, + "lambda": { + "arn": "arn:aws:lambda:us-east-1:123456789012:function:test-customized-loggroup" } }, { - "host": "/aws/vendedlogs/states/anyLogGroupName", + "host": "arn:aws:lambda:us-east-1:123456789012:function:test-customized-loggroup", "id": "35311576111948622874033876462979853992919938886093242370", "timestamp": 1583425836114, "message": "2020-03-05T16:30:36.114Z\tf08bb4c8-d6b2-4f05-ac17-af7e2ba005fb\tDEBUG\t[dd.trace_id=3172564172058669914 dd.span_id=14292093692483532556] {\"status\":\"debug\",\"message\":\"datadog:Reading trace context from env var Root=1-5e61292c-cc1229a4dfbeae1043928548;Parent=c657b77d9514f70c;Sampled=1\"}\n", "service": "lambda", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "", + "ddtags": "functionname:test-customized-loggroup,env:none", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -61,6 +67,9 @@ "logStream": "2020/03/05/test-customized-loggroup[$LATEST]20bddfd5a2dc4c6b97ac02800eae90d0", "owner": "601427279990" } + }, + "lambda": { + "arn": "arn:aws:lambda:us-east-1:123456789012:function:test-customized-loggroup" } } ] diff --git a/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_timeout.golden.json b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_timeout.golden.json index 41d0544d2..0010bd278 100644 --- a/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_timeout.golden.json +++ b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_timeout.golden.json @@ -10,14 +10,14 @@ }, "body": [ { - "host": "/aws/lambda/storms-cloudwatch-event", + "host": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event", "id": "35496429375792603298393743017356257146982675867810398208", "timestamp": 1591714943146, "message": "START RequestId: 7c9567b5-107b-4a6c-8798-0157ac21db52 Version: $LATEST\\n", "service": "lambda", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "", + "ddtags": "functionname:storms-cloudwatch-event,env:none", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -25,17 +25,20 @@ "logStream": "2020/06/09/[$LATEST]b249865adaaf4fad80f95f8ad09725b8", "owner": "601427279990" } + }, + "lambda": { + "arn": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event" } }, { - "host": "/aws/lambda/storms-cloudwatch-event", + "host": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event", "id": "35496429442806342619978265557671090556291002193281548289", "timestamp": 1591714946151, "message": "END RequestId: 7c9567b5-107b-4a6c-8798-0157ac21db52\\n", "service": "lambda", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "", + "ddtags": "functionname:storms-cloudwatch-event,env:none", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -43,17 +46,20 @@ "logStream": "2020/06/09/[$LATEST]b249865adaaf4fad80f95f8ad09725b8", "owner": "601427279990" } + }, + "lambda": { + "arn": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event" } }, { - "host": "/aws/lambda/storms-cloudwatch-event", + "host": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event", "id": "35496429442806342619978265557671090556291002193281548290", "timestamp": 1591714946151, "message": "REPORT RequestId: 7c9567b5-107b-4a6c-8798-0157ac21db52\\tDuration: 3003.16 ms\\tBilled Duration: 3000 ms\\tMemory Size: 128 MB\\tMax Memory Used: 48 MB\\tInit Duration: 127.02 ms\\t\\n", "service": "lambda", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "", + "ddtags": "functionname:storms-cloudwatch-event,env:none", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -61,17 +67,20 @@ "logStream": "2020/06/09/[$LATEST]b249865adaaf4fad80f95f8ad09725b8", "owner": "601427279990" } + }, + "lambda": { + "arn": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event" } }, { - "host": "/aws/lambda/storms-cloudwatch-event", + "host": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event", "id": "35496429442806342619978265557671090556291002193281548291", "timestamp": 1591714946151, "message": "2020-06-09T15:02:26.150Z 7c9567b5-107b-4a6c-8798-0157ac21db52 Task timed out after 3.00 seconds\\n\\n", "service": "lambda", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "", + "ddtags": "functionname:storms-cloudwatch-event,env:none", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -79,6 +88,9 @@ "logStream": "2020/06/09/[$LATEST]b249865adaaf4fad80f95f8ad09725b8", "owner": "601427279990" } + }, + "lambda": { + "arn": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event" } } ] diff --git a/aws/logs_monitoring_go/internal/handling/cloudwatch.go b/aws/logs_monitoring_go/internal/handling/cloudwatch.go index 54a8a540b..bd6549607 100644 --- a/aws/logs_monitoring_go/internal/handling/cloudwatch.go +++ b/aws/logs_monitoring_go/internal/handling/cloudwatch.go @@ -34,9 +34,12 @@ const ( logStreamCloudtrail = "_CloudTrail_" ) +const envTag = "env" + // Custom log groups use the log stream format: YYYY/MM/DD/[][] +// The first capture group extracts the function name. var lambdaLogStreamRegex = regexp.MustCompile( - `^\d{4}/[01]\d/[0-3]\d/[\w.-]{1,75}\[(\$LATEST|[\w-]{1,129})\][0-9a-f]{32}$`, + `^\d{4}/[01]\d/[0-3]\d/([\w.-]{1,75})\[(?:\$LATEST|[\w-]{1,129})\][0-9a-f]{32}$`, ) type cloudwatchHandler struct { @@ -130,18 +133,24 @@ func (h cloudwatchHandler) newCloudwatchBaseEntry(data events.CloudwatchLogsData Owner: data.Owner, }, } + source := cloudwatchSource(strings.ToLower(logGroup), logStream) entry := model.NewLogEntry() - entry.Source = cmp.Or(h.cfg.Source, CloudwatchSource(strings.ToLower(logGroup), logStream)) + entry.Source = cmp.Or(h.cfg.Source, source) entry.Host = cmp.Or(h.cfg.Host, logGroup) entry.Metadata = metadata + + if entry.Source == sourceLambda { + enrichLambdaLog(&entry, lambdaOrigin.ARN, logGroup, logStream, h.cfg) + } + return entry } func (h cloudwatchHandler) newCloudwatchLogEntry(event events.CloudwatchLogsLogEvent, entry model.LogEntry) model.LogEntry { tags, service, message := extractFromMessage(event.Message) entry.Service = cmp.Or(h.cfg.Service, service, entry.Source) - entry.Tags = slices.Concat(tags, h.cfg.Tags) + entry.Tags = slices.Concat(tags, entry.Tags, h.cfg.Tags) entry.Message = message entry.ID = event.ID entry.Timestamp = event.Timestamp @@ -153,7 +162,7 @@ func (h cloudwatchHandler) newCloudwatchLogEntry(event events.CloudwatchLogsLogE return entry } -func CloudwatchSource(logGroup, logStream string) string { +func cloudwatchSource(logGroup, logStream string) string { if strings.HasPrefix(logStream, logStreamStepFunction) { return sourceStepFunction } @@ -177,3 +186,44 @@ func CloudwatchSource(logGroup, logStream string) string { } return sourceCloudwatch } + +func enrichLambdaLog(entry *model.LogEntry, forwarderARN, logGroup, logStream string, cfg *Config) { + name := lambdaName(strings.ToLower(logGroup), logStream) + if name != "" { + return + } + + prefix, _, found := strings.Cut(forwarderARN, "function:") + if !found { + return + } + + arn := prefix + "function:" + name + entry.Tags = append(entry.Tags, "functionname:"+name) + entry.Lambda = &model.LambdaLog{ARN: arn} + entry.Host = cmp.Or(cfg.Host, arn) + + if !hasTag(cfg.Tags, envTag) { + entry.Tags = append(entry.Tags, envTag+":none") + } +} + +func lambdaName(logGroup, logStream string) string { + if name := lambdaLogStreamRegex.FindString(logStream); name != "" { + return strings.ToLower(name) + } + + if _, name, found := strings.Cut(logGroup, logGroupLambda+"/"); found && name != "" { + return strings.ToLower(name) + } + return "" +} + +func hasTag(tags model.Tags, prefix string) bool { + for _, tag := range tags { + if strings.HasPrefix(tag, prefix) { + return true + } + } + return false +} diff --git a/aws/logs_monitoring_go/internal/handling/cloudwatch_test.go b/aws/logs_monitoring_go/internal/handling/cloudwatch_test.go index 6c5f31fa8..adb5705f4 100644 --- a/aws/logs_monitoring_go/internal/handling/cloudwatch_test.go +++ b/aws/logs_monitoring_go/internal/handling/cloudwatch_test.go @@ -70,9 +70,11 @@ func TestCloudwatchHandler_Handle(t *testing.T) { Source: "lambda", SourceCategory: "aws", Service: "lambda", - Host: "/aws/lambda/testing-datadog", + Host: "arn:aws:lambda:us-east-1:123456789012:function:testing-datadog", + Tags: model.Tags{"functionname:testing-datadog", "env:none"}, ID: "ev1", Timestamp: 1583425836114, + Lambda: &model.LambdaLog{ARN: "arn:aws:lambda:us-east-1:123456789012:function:testing-datadog"}, Metadata: model.CloudwatchMetadata{ LambdaOrigin: model.LambdaOrigin{ARN: "arn:aws:lambda:us-east-1:123456789012:function:forwarder"}, Origin: model.CloudwatchOrigin{ @@ -101,7 +103,10 @@ func TestCloudwatchHandler_Handle(t *testing.T) { { Message: "first", Source: "lambda", SourceCategory: "aws", Service: "lambda", - Host: "/aws/lambda/fn", ID: "a1", Timestamp: 1000, + Host: "arn:aws:lambda:us-east-1:123456789012:function:fn", + Tags: model.Tags{"functionname:fn", "env:none"}, + ID: "a1", Timestamp: 1000, + Lambda: &model.LambdaLog{ARN: "arn:aws:lambda:us-east-1:123456789012:function:fn"}, Metadata: model.CloudwatchMetadata{ LambdaOrigin: model.LambdaOrigin{ARN: "arn:aws:lambda:us-east-1:123456789012:function:forwarder"}, Origin: model.CloudwatchOrigin{LogGroup: "/aws/lambda/fn", LogStream: "stream", Owner: "111111111111"}, @@ -110,7 +115,10 @@ func TestCloudwatchHandler_Handle(t *testing.T) { { Message: "second", Source: "lambda", SourceCategory: "aws", Service: "lambda", - Host: "/aws/lambda/fn", ID: "a2", Timestamp: 2000, + Host: "arn:aws:lambda:us-east-1:123456789012:function:fn", + Tags: model.Tags{"functionname:fn", "env:none"}, + ID: "a2", Timestamp: 2000, + Lambda: &model.LambdaLog{ARN: "arn:aws:lambda:us-east-1:123456789012:function:fn"}, Metadata: model.CloudwatchMetadata{ LambdaOrigin: model.LambdaOrigin{ARN: "arn:aws:lambda:us-east-1:123456789012:function:forwarder"}, Origin: model.CloudwatchOrigin{LogGroup: "/aws/lambda/fn", LogStream: "stream", Owner: "111111111111"}, @@ -271,7 +279,57 @@ func TestCloudwatchSource(t *testing.T) { for name, tc := range tests { t.Run(name, func(t *testing.T) { t.Parallel() - assert.Equal(t, tc.want, CloudwatchSource(tc.logGroup, tc.logStream)) + assert.Equal(t, tc.want, cloudwatchSource(tc.logGroup, tc.logStream)) + }) + } +} + +func TestLambdaName(t *testing.T) { + t.Parallel() + + tests := map[string]struct { + logGroup string + logStream string + want string + }{ + "default log group": {logGroup: "/aws/lambda/my-function", logStream: "stream", want: "my-function"}, + "default log group lowercased": {logGroup: "/aws/lambda/My-Function", logStream: "stream", want: "my-function"}, + "custom log group name from stream": {logGroup: "/aws/vendedlogs/states/anyLogGroupName", logStream: "2020/03/05/Test-Customized-LogGroup[$LATEST]20bddfd5a2dc4c6b97ac02800eae90d0", want: "test-customized-loggroup"}, + "stream takes priority over log group": {logGroup: "/aws/lambda/from-group", logStream: "2020/03/05/from-stream[$LATEST]20bddfd5a2dc4c6b97ac02800eae90d0", want: "from-stream"}, + "stream without function name": {logGroup: "my-custom-group", logStream: "2023/11/04/[$LATEST]4426346c2cdf4c54a74d3bd2b929fc44", want: ""}, + "non-lambda log group": {logGroup: "/aws/rds/cluster", logStream: "stream", want: ""}, + "empty": {logGroup: "", logStream: "", want: ""}, + } + + for name, tc := range tests { + t.Run(name, func(t *testing.T) { + t.Parallel() + assert.Equal(t, tc.want, lambdaName(tc.logGroup, tc.logStream)) + }) + } +} + +func TestHasTag(t *testing.T) { + t.Parallel() + + tests := map[string]struct { + tags model.Tags + prefix string + want bool + }{ + "nil": {tags: nil, prefix: "env:", want: false}, + "no env": {tags: model.Tags{"team:infra", "service:api"}, prefix: "env:", want: false}, + "env present": {tags: model.Tags{"team:infra", "env:prod"}, prefix: "env:", want: true}, + "env none present": {tags: model.Tags{"env:none"}, prefix: "env:", want: true}, + "env prefix only": {tags: model.Tags{"environment:prod"}, prefix: "env:", want: false}, + "other prefix": {tags: model.Tags{"team:infra", "service:api"}, prefix: "service:", want: true}, + "empty prefix always": {tags: model.Tags{"team:infra"}, prefix: "", want: true}, + } + + for name, tc := range tests { + t.Run(name, func(t *testing.T) { + t.Parallel() + assert.Equal(t, tc.want, hasTag(tc.tags, tc.prefix)) }) } } diff --git a/aws/logs_monitoring_go/internal/model/lambda.go b/aws/logs_monitoring_go/internal/model/lambda.go new file mode 100644 index 000000000..c5e0f3bda --- /dev/null +++ b/aws/logs_monitoring_go/internal/model/lambda.go @@ -0,0 +1,10 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2026-Present Datadog, Inc. + +package model + +type LambdaLog struct { + ARN string `json:"arn"` +} diff --git a/aws/logs_monitoring_go/internal/model/log.go b/aws/logs_monitoring_go/internal/model/log.go index 3f669ffd9..c2f44f57d 100644 --- a/aws/logs_monitoring_go/internal/model/log.go +++ b/aws/logs_monitoring_go/internal/model/log.go @@ -13,15 +13,16 @@ import ( const sourceCategory = "aws" type LogEntry struct { - Host string `json:"host,omitempty"` - ID string `json:"id,omitempty"` - Timestamp int64 `json:"timestamp,omitempty"` - Message string `json:"message,omitempty"` - Service string `json:"service,omitempty"` - Source string `json:"ddsource"` - SourceCategory string `json:"ddsourcecategory"` - Tags Tags `json:"ddtags"` - Metadata any `json:"aws"` + Host string `json:"host,omitempty"` + ID string `json:"id,omitempty"` + Timestamp int64 `json:"timestamp,omitempty"` + Message string `json:"message,omitempty"` + Service string `json:"service,omitempty"` + Source string `json:"ddsource"` + SourceCategory string `json:"ddsourcecategory"` + Tags Tags `json:"ddtags"` + Metadata any `json:"aws"` + Lambda *LambdaLog `json:"lambda,omitempty"` } func NewLogEntry() LogEntry { From 293ec3eb3fb011570c70746a6a9d1e80bbe9165b Mon Sep 17 00:00:00 2001 From: Nabil Dakkoune Date: Tue, 21 Jul 2026 10:49:43 +0200 Subject: [PATCH 2/5] fix --- aws/logs_monitoring_go/internal/handling/cloudwatch.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/aws/logs_monitoring_go/internal/handling/cloudwatch.go b/aws/logs_monitoring_go/internal/handling/cloudwatch.go index bd6549607..c68fa2394 100644 --- a/aws/logs_monitoring_go/internal/handling/cloudwatch.go +++ b/aws/logs_monitoring_go/internal/handling/cloudwatch.go @@ -189,7 +189,7 @@ func cloudwatchSource(logGroup, logStream string) string { func enrichLambdaLog(entry *model.LogEntry, forwarderARN, logGroup, logStream string, cfg *Config) { name := lambdaName(strings.ToLower(logGroup), logStream) - if name != "" { + if name == "" { return } @@ -203,14 +203,14 @@ func enrichLambdaLog(entry *model.LogEntry, forwarderARN, logGroup, logStream st entry.Lambda = &model.LambdaLog{ARN: arn} entry.Host = cmp.Or(cfg.Host, arn) - if !hasTag(cfg.Tags, envTag) { + if !hasTag(cfg.Tags, envTag+":") { entry.Tags = append(entry.Tags, envTag+":none") } } func lambdaName(logGroup, logStream string) string { - if name := lambdaLogStreamRegex.FindString(logStream); name != "" { - return strings.ToLower(name) + if m := lambdaLogStreamRegex.FindStringSubmatch(logStream); m != nil { + return strings.ToLower(m[1]) } if _, name, found := strings.Cut(logGroup, logGroupLambda+"/"); found && name != "" { From 60cd6a1f5c99e25351443e8d778751c5bc4416bb Mon Sep 17 00:00:00 2001 From: Nabil Dakkoune Date: Tue, 21 Jul 2026 17:40:15 +0200 Subject: [PATCH 3/5] clean separation --- aws/logs_monitoring_go/internal/model/log.go | 26 ------- .../internal/model/log_test.go | 57 --------------- aws/logs_monitoring_go/internal/model/tags.go | 50 ++++++++++++++ .../internal/model/tags_test.go | 69 +++++++++++++++++++ 4 files changed, 119 insertions(+), 83 deletions(-) delete mode 100644 aws/logs_monitoring_go/internal/model/log_test.go create mode 100644 aws/logs_monitoring_go/internal/model/tags.go create mode 100644 aws/logs_monitoring_go/internal/model/tags_test.go diff --git a/aws/logs_monitoring_go/internal/model/log.go b/aws/logs_monitoring_go/internal/model/log.go index c2f44f57d..30fe7c9e2 100644 --- a/aws/logs_monitoring_go/internal/model/log.go +++ b/aws/logs_monitoring_go/internal/model/log.go @@ -5,11 +5,6 @@ package model -import ( - "encoding/json" - "strings" -) - const sourceCategory = "aws" type LogEntry struct { @@ -30,24 +25,3 @@ func NewLogEntry() LogEntry { SourceCategory: sourceCategory, } } - -type Tags []string - -func (t Tags) MarshalJSON() ([]byte, error) { - return json.Marshal(strings.Join(t, ",")) -} - -func (t *Tags) UnmarshalJSON(data []byte) error { - var s string - - if err := json.Unmarshal(data, &s); err != nil { - return err - } - if s == "" { - *t = nil - return nil - } - - *t = strings.Split(s, ",") - return nil -} diff --git a/aws/logs_monitoring_go/internal/model/log_test.go b/aws/logs_monitoring_go/internal/model/log_test.go deleted file mode 100644 index 811fc3f2c..000000000 --- a/aws/logs_monitoring_go/internal/model/log_test.go +++ /dev/null @@ -1,57 +0,0 @@ -// Unless explicitly stated otherwise all files in this repository are licensed -// under the Apache License Version 2.0. -// This product includes software developed at Datadog (https://www.datadoghq.com/). -// Copyright 2026-Present Datadog, Inc. - -package model - -import ( - "encoding/json" - "testing" - - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" -) - -func TestTags(t *testing.T) { - t.Parallel() - - tests := map[string]struct { - tags Tags - want string - }{ - "multiple_tags": { - tags: Tags{"env:prod", "team:aws"}, - want: `"env:prod,team:aws"`, - }, - "single_tag": { - tags: Tags{"env:prod"}, - want: `"env:prod"`, - }, - "empty": { - tags: Tags{}, - want: `""`, - }, - "nil": { - tags: nil, - want: `""`, - }, - } - - for name, tc := range tests { - t.Run(name, func(t *testing.T) { - t.Parallel() - got, err := json.Marshal(tc.tags) - require.NoError(t, err, "marshal") - assert.Equal(t, tc.want, string(got)) - - var tags Tags - require.NoError(t, json.Unmarshal(got, &tags), "unmarshal") - if len(tc.tags) == 0 { - assert.Empty(t, tags) - } else { - assert.Equal(t, tc.tags, tags) - } - }) - } -} diff --git a/aws/logs_monitoring_go/internal/model/tags.go b/aws/logs_monitoring_go/internal/model/tags.go new file mode 100644 index 000000000..598cd40fd --- /dev/null +++ b/aws/logs_monitoring_go/internal/model/tags.go @@ -0,0 +1,50 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2026-Present Datadog, Inc. + +package model + +import ( + "encoding/json" + "strings" +) + +const ( + KeyValueSeparator = ":" + TagSeparator = "," +) + +type Tags []string + +func (t Tags) MarshalJSON() ([]byte, error) { + return json.Marshal(strings.Join(t, TagSeparator)) +} + +func (t *Tags) UnmarshalJSON(data []byte) error { + var s string + + if err := json.Unmarshal(data, &s); err != nil { + return err + } + if s == "" { + *t = nil + return nil + } + + *t = strings.Split(s, ",") + return nil +} + +func (t *Tags) Add(key, value string) { + *t = append(*t, key+KeyValueSeparator+value) +} + +func (t *Tags) Has(key string) bool { + for _, tag := range *t { + if strings.HasPrefix(tag, key+KeyValueSeparator) { + return true + } + } + return false +} diff --git a/aws/logs_monitoring_go/internal/model/tags_test.go b/aws/logs_monitoring_go/internal/model/tags_test.go new file mode 100644 index 000000000..b5439d222 --- /dev/null +++ b/aws/logs_monitoring_go/internal/model/tags_test.go @@ -0,0 +1,69 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2026-Present Datadog, Inc. + +package model + +import ( + "encoding/json" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestTags_MarshalJSON_UnmarshalJSON(t *testing.T) { + t.Parallel() + + tests := map[string]struct { + tags Tags + want string + }{ + "nil": {tags: nil, want: `""`}, + "empty": {tags: Tags{}, want: `""`}, + "one": {tags: Tags{"env:prod"}, want: `"env:prod"`}, + "multiple": {tags: Tags{"env:prod", "team:aws"}, want: `"env:prod,team:aws"`}, + } + + for name, tc := range tests { + t.Run(name, func(t *testing.T) { + t.Parallel() + + got, err := json.Marshal(tc.tags) + + require.NoError(t, err) + assert.Equal(t, tc.want, string(got)) + + var tags Tags + require.NoError(t, json.Unmarshal(got, &tags)) + if len(tc.tags) == 0 { + assert.Empty(t, tags) + return + } + assert.Equal(t, tc.tags, tags) + }) + } +} + +func TestTags_Has(t *testing.T) { + t.Parallel() + + tests := map[string]struct { + tags Tags + key string + want bool + }{ + "nil": {tags: nil, key: "env", want: false}, + "present": {tags: Tags{"team:infra", "env:prod"}, key: "env", want: true}, + "not present": {tags: Tags{"team:infra", "service:api"}, key: "env", want: false}, + "not present when separator": {tags: Tags{"team:infra", "service:api"}, key: "service:", want: false}, + } + + for name, tc := range tests { + t.Run(name, func(t *testing.T) { + t.Parallel() + assert.Equal(t, tc.want, tc.tags.Has(tc.key)) + }) + } +} From 973d61e9c60ace5e2ca8619f82de76dfe3b0427e Mon Sep 17 00:00:00 2001 From: Nabil Dakkoune Date: Wed, 22 Jul 2026 10:12:56 +0200 Subject: [PATCH 4/5] add service and tag from lmbda name --- .../cloudwatch_lambda_coldstart.golden.json | 12 +-- ...dwatch_lambda_custom_log_group.golden.json | 12 +-- .../cloudwatch_lambda_timeout.golden.json | 16 ++-- .../internal/handling/cloudwatch.go | 54 +++-------- .../internal/handling/cloudwatch_test.go | 89 +++++++------------ .../internal/handling/lambda.go | 34 +++++++ .../internal/handling/lambda_test.go | 38 ++++++++ 7 files changed, 135 insertions(+), 120 deletions(-) create mode 100644 aws/logs_monitoring_go/internal/handling/lambda.go create mode 100644 aws/logs_monitoring_go/internal/handling/lambda_test.go diff --git a/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_coldstart.golden.json b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_coldstart.golden.json index 635e3b709..639c46713 100644 --- a/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_coldstart.golden.json +++ b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_coldstart.golden.json @@ -14,10 +14,10 @@ "id": "35486831490800643125153606102923171443962457178576257024", "timestamp": 1591284559098, "message": "START RequestId: db275f87-a934-471a-8980-b63bf4dc1beb Version: $LATEST\\n", - "service": "lambda", + "service": "storms-cloudwatch-event", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "functionname:storms-cloudwatch-event,env:none", + "ddtags": "functionname:storms-cloudwatch-event,env:none,service:storms-cloudwatch-event", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -35,10 +35,10 @@ "id": "35486831490867545360749197972347778598780402263094198273", "timestamp": 1591284559101, "message": "END RequestId: db275f87-a934-471a-8980-b63bf4dc1beb\\n", - "service": "lambda", + "service": "storms-cloudwatch-event", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "functionname:storms-cloudwatch-event,env:none", + "ddtags": "functionname:storms-cloudwatch-event,env:none,service:storms-cloudwatch-event", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -56,10 +56,10 @@ "id": "35486831490867545360749197972347778598780402263094198274", "timestamp": 1591284559101, "message": "REPORT RequestId: db275f87-a934-471a-8980-b63bf4dc1beb\\tDuration: 1.76 ms\\tBilled Duration: 100 ms\\tMemory Size: 128 MB\\tMax Memory Used: 48 MB\\tInit Duration: 120.96 ms\\t\\n", - "service": "lambda", + "service": "storms-cloudwatch-event", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "functionname:storms-cloudwatch-event,env:none", + "ddtags": "functionname:storms-cloudwatch-event,env:none,service:storms-cloudwatch-event", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { diff --git a/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_custom_log_group.golden.json b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_custom_log_group.golden.json index 7a642cf6a..3f59933df 100644 --- a/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_custom_log_group.golden.json +++ b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_custom_log_group.golden.json @@ -14,10 +14,10 @@ "id": "35311576111948622874033876462979853992919938886093242368", "timestamp": 1583425836114, "message": "2020-03-05T16:30:36.113Z\tf08bb4c8-d6b2-4f05-ac17-af7e2ba005fb\tDEBUG\t[dd.trace_id=3172564172058669914 dd.span_id=14292093692483532556] {\"status\":\"debug\",\"message\":\"datadog:Patched console output with trace context\"}\n", - "service": "lambda", + "service": "test-customized-loggroup", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "functionname:test-customized-loggroup,env:none", + "ddtags": "functionname:test-customized-loggroup,env:none,service:test-customized-loggroup", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -35,10 +35,10 @@ "id": "35311576111948622874033876462979853992919938886093242369", "timestamp": 1583425836114, "message": "2020-03-05T16:30:36.114Z\tf08bb4c8-d6b2-4f05-ac17-af7e2ba005fb\tDEBUG\t[dd.trace_id=3172564172058669914 dd.span_id=14292093692483532556] {\"autoPatchHTTP\":true,\"tracerInitialized\":true,\"status\":\"debug\",\"message\":\"datadog:Not patching HTTP libraries\"}\n", - "service": "lambda", + "service": "test-customized-loggroup", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "functionname:test-customized-loggroup,env:none", + "ddtags": "functionname:test-customized-loggroup,env:none,service:test-customized-loggroup", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -56,10 +56,10 @@ "id": "35311576111948622874033876462979853992919938886093242370", "timestamp": 1583425836114, "message": "2020-03-05T16:30:36.114Z\tf08bb4c8-d6b2-4f05-ac17-af7e2ba005fb\tDEBUG\t[dd.trace_id=3172564172058669914 dd.span_id=14292093692483532556] {\"status\":\"debug\",\"message\":\"datadog:Reading trace context from env var Root=1-5e61292c-cc1229a4dfbeae1043928548;Parent=c657b77d9514f70c;Sampled=1\"}\n", - "service": "lambda", + "service": "test-customized-loggroup", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "functionname:test-customized-loggroup,env:none", + "ddtags": "functionname:test-customized-loggroup,env:none,service:test-customized-loggroup", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { diff --git a/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_timeout.golden.json b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_timeout.golden.json index 0010bd278..abc81a45a 100644 --- a/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_timeout.golden.json +++ b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_timeout.golden.json @@ -14,10 +14,10 @@ "id": "35496429375792603298393743017356257146982675867810398208", "timestamp": 1591714943146, "message": "START RequestId: 7c9567b5-107b-4a6c-8798-0157ac21db52 Version: $LATEST\\n", - "service": "lambda", + "service": "storms-cloudwatch-event", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "functionname:storms-cloudwatch-event,env:none", + "ddtags": "functionname:storms-cloudwatch-event,env:none,service:storms-cloudwatch-event", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -35,10 +35,10 @@ "id": "35496429442806342619978265557671090556291002193281548289", "timestamp": 1591714946151, "message": "END RequestId: 7c9567b5-107b-4a6c-8798-0157ac21db52\\n", - "service": "lambda", + "service": "storms-cloudwatch-event", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "functionname:storms-cloudwatch-event,env:none", + "ddtags": "functionname:storms-cloudwatch-event,env:none,service:storms-cloudwatch-event", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -56,10 +56,10 @@ "id": "35496429442806342619978265557671090556291002193281548290", "timestamp": 1591714946151, "message": "REPORT RequestId: 7c9567b5-107b-4a6c-8798-0157ac21db52\\tDuration: 3003.16 ms\\tBilled Duration: 3000 ms\\tMemory Size: 128 MB\\tMax Memory Used: 48 MB\\tInit Duration: 127.02 ms\\t\\n", - "service": "lambda", + "service": "storms-cloudwatch-event", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "functionname:storms-cloudwatch-event,env:none", + "ddtags": "functionname:storms-cloudwatch-event,env:none,service:storms-cloudwatch-event", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { @@ -77,10 +77,10 @@ "id": "35496429442806342619978265557671090556291002193281548291", "timestamp": 1591714946151, "message": "2020-06-09T15:02:26.150Z 7c9567b5-107b-4a6c-8798-0157ac21db52 Task timed out after 3.00 seconds\\n\\n", - "service": "lambda", + "service": "storms-cloudwatch-event", "ddsource": "lambda", "ddsourcecategory": "aws", - "ddtags": "functionname:storms-cloudwatch-event,env:none", + "ddtags": "functionname:storms-cloudwatch-event,env:none,service:storms-cloudwatch-event", "aws": { "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", "awslogs": { diff --git a/aws/logs_monitoring_go/internal/handling/cloudwatch.go b/aws/logs_monitoring_go/internal/handling/cloudwatch.go index c68fa2394..58dd697ee 100644 --- a/aws/logs_monitoring_go/internal/handling/cloudwatch.go +++ b/aws/logs_monitoring_go/internal/handling/cloudwatch.go @@ -37,7 +37,6 @@ const ( const envTag = "env" // Custom log groups use the log stream format: YYYY/MM/DD/[][] -// The first capture group extracts the function name. var lambdaLogStreamRegex = regexp.MustCompile( `^\d{4}/[01]\d/[0-3]\d/([\w.-]{1,75})\[(?:\$LATEST|[\w-]{1,129})\][0-9a-f]{32}$`, ) @@ -141,7 +140,10 @@ func (h cloudwatchHandler) newCloudwatchBaseEntry(data events.CloudwatchLogsData entry.Metadata = metadata if entry.Source == sourceLambda { - enrichLambdaLog(&entry, lambdaOrigin.ARN, logGroup, logStream, h.cfg) + enrichLambdaLog(&entry, lambdaOrigin.ARN, logGroup, logStream) + if !h.cfg.Tags.Has(envTag) { + entry.Tags.Add(envTag, "none") + } } return entry @@ -149,7 +151,8 @@ func (h cloudwatchHandler) newCloudwatchBaseEntry(data events.CloudwatchLogsData func (h cloudwatchHandler) newCloudwatchLogEntry(event events.CloudwatchLogsLogEvent, entry model.LogEntry) model.LogEntry { tags, service, message := extractFromMessage(event.Message) - entry.Service = cmp.Or(h.cfg.Service, service, entry.Source) + + entry.Service = cmp.Or(h.cfg.Service, service, entry.Service, entry.Source) entry.Tags = slices.Concat(tags, entry.Tags, h.cfg.Tags) entry.Message = message entry.ID = event.ID @@ -159,6 +162,10 @@ func (h cloudwatchHandler) newCloudwatchLogEntry(event events.CloudwatchLogsLogE entry.Host = cloudtrailHost(event.Message) } + if entry.Lambda != nil { + entry.Tags.Add("service", entry.Service) + } + return entry } @@ -186,44 +193,3 @@ func cloudwatchSource(logGroup, logStream string) string { } return sourceCloudwatch } - -func enrichLambdaLog(entry *model.LogEntry, forwarderARN, logGroup, logStream string, cfg *Config) { - name := lambdaName(strings.ToLower(logGroup), logStream) - if name == "" { - return - } - - prefix, _, found := strings.Cut(forwarderARN, "function:") - if !found { - return - } - - arn := prefix + "function:" + name - entry.Tags = append(entry.Tags, "functionname:"+name) - entry.Lambda = &model.LambdaLog{ARN: arn} - entry.Host = cmp.Or(cfg.Host, arn) - - if !hasTag(cfg.Tags, envTag+":") { - entry.Tags = append(entry.Tags, envTag+":none") - } -} - -func lambdaName(logGroup, logStream string) string { - if m := lambdaLogStreamRegex.FindStringSubmatch(logStream); m != nil { - return strings.ToLower(m[1]) - } - - if _, name, found := strings.Cut(logGroup, logGroupLambda+"/"); found && name != "" { - return strings.ToLower(name) - } - return "" -} - -func hasTag(tags model.Tags, prefix string) bool { - for _, tag := range tags { - if strings.HasPrefix(tag, prefix) { - return true - } - } - return false -} diff --git a/aws/logs_monitoring_go/internal/handling/cloudwatch_test.go b/aws/logs_monitoring_go/internal/handling/cloudwatch_test.go index adb5705f4..20c3f447c 100644 --- a/aws/logs_monitoring_go/internal/handling/cloudwatch_test.go +++ b/aws/logs_monitoring_go/internal/handling/cloudwatch_test.go @@ -69,9 +69,9 @@ func TestCloudwatchHandler_Handle(t *testing.T) { Message: "hello", Source: "lambda", SourceCategory: "aws", - Service: "lambda", + Service: "testing-datadog", Host: "arn:aws:lambda:us-east-1:123456789012:function:testing-datadog", - Tags: model.Tags{"functionname:testing-datadog", "env:none"}, + Tags: model.Tags{"functionname:testing-datadog", "env:none", "service:testing-datadog"}, ID: "ev1", Timestamp: 1583425836114, Lambda: &model.LambdaLog{ARN: "arn:aws:lambda:us-east-1:123456789012:function:testing-datadog"}, @@ -102,9 +102,9 @@ func TestCloudwatchHandler_Handle(t *testing.T) { want: []model.LogEntry{ { Message: "first", Source: "lambda", SourceCategory: "aws", - Service: "lambda", + Service: "fn", Host: "arn:aws:lambda:us-east-1:123456789012:function:fn", - Tags: model.Tags{"functionname:fn", "env:none"}, + Tags: model.Tags{"functionname:fn", "env:none", "service:fn"}, ID: "a1", Timestamp: 1000, Lambda: &model.LambdaLog{ARN: "arn:aws:lambda:us-east-1:123456789012:function:fn"}, Metadata: model.CloudwatchMetadata{ @@ -114,9 +114,9 @@ func TestCloudwatchHandler_Handle(t *testing.T) { }, { Message: "second", Source: "lambda", SourceCategory: "aws", - Service: "lambda", + Service: "fn", Host: "arn:aws:lambda:us-east-1:123456789012:function:fn", - Tags: model.Tags{"functionname:fn", "env:none"}, + Tags: model.Tags{"functionname:fn", "env:none", "service:fn"}, ID: "a2", Timestamp: 2000, Lambda: &model.LambdaLog{ARN: "arn:aws:lambda:us-east-1:123456789012:function:fn"}, Metadata: model.CloudwatchMetadata{ @@ -228,6 +228,33 @@ func TestCloudwatchHandler_Handle(t *testing.T) { }, }, }, + "lambda keeps configured env and skips env:none": { + event: testutil.MustCloudwatchEvent(t, testutil.MustGzipJSON(t, map[string]any{ + "messageType": "DATA_MESSAGE", + "owner": "111111111111", + "logGroup": "/aws/lambda/fn", + "logStream": "stream", + "logEvents": []map[string]any{ + {"id": "ev1", "timestamp": 1000, "message": "hello"}, + }, + })), + config: &Config{Tags: model.Tags{"env:prod"}}, + chanSize: 1, + want: []model.LogEntry{ + { + Message: "hello", Source: "lambda", SourceCategory: "aws", + Service: "fn", + Host: "arn:aws:lambda:us-east-1:123456789012:function:fn", + Tags: model.Tags{"functionname:fn", "env:prod", "service:fn"}, + ID: "ev1", Timestamp: 1000, + Lambda: &model.LambdaLog{ARN: "arn:aws:lambda:us-east-1:123456789012:function:fn"}, + Metadata: model.CloudwatchMetadata{ + LambdaOrigin: model.LambdaOrigin{ARN: "arn:aws:lambda:us-east-1:123456789012:function:forwarder"}, + Origin: model.CloudwatchOrigin{LogGroup: "/aws/lambda/fn", LogStream: "stream", Owner: "111111111111"}, + }, + }, + }, + }, } for name, tc := range tests { @@ -283,53 +310,3 @@ func TestCloudwatchSource(t *testing.T) { }) } } - -func TestLambdaName(t *testing.T) { - t.Parallel() - - tests := map[string]struct { - logGroup string - logStream string - want string - }{ - "default log group": {logGroup: "/aws/lambda/my-function", logStream: "stream", want: "my-function"}, - "default log group lowercased": {logGroup: "/aws/lambda/My-Function", logStream: "stream", want: "my-function"}, - "custom log group name from stream": {logGroup: "/aws/vendedlogs/states/anyLogGroupName", logStream: "2020/03/05/Test-Customized-LogGroup[$LATEST]20bddfd5a2dc4c6b97ac02800eae90d0", want: "test-customized-loggroup"}, - "stream takes priority over log group": {logGroup: "/aws/lambda/from-group", logStream: "2020/03/05/from-stream[$LATEST]20bddfd5a2dc4c6b97ac02800eae90d0", want: "from-stream"}, - "stream without function name": {logGroup: "my-custom-group", logStream: "2023/11/04/[$LATEST]4426346c2cdf4c54a74d3bd2b929fc44", want: ""}, - "non-lambda log group": {logGroup: "/aws/rds/cluster", logStream: "stream", want: ""}, - "empty": {logGroup: "", logStream: "", want: ""}, - } - - for name, tc := range tests { - t.Run(name, func(t *testing.T) { - t.Parallel() - assert.Equal(t, tc.want, lambdaName(tc.logGroup, tc.logStream)) - }) - } -} - -func TestHasTag(t *testing.T) { - t.Parallel() - - tests := map[string]struct { - tags model.Tags - prefix string - want bool - }{ - "nil": {tags: nil, prefix: "env:", want: false}, - "no env": {tags: model.Tags{"team:infra", "service:api"}, prefix: "env:", want: false}, - "env present": {tags: model.Tags{"team:infra", "env:prod"}, prefix: "env:", want: true}, - "env none present": {tags: model.Tags{"env:none"}, prefix: "env:", want: true}, - "env prefix only": {tags: model.Tags{"environment:prod"}, prefix: "env:", want: false}, - "other prefix": {tags: model.Tags{"team:infra", "service:api"}, prefix: "service:", want: true}, - "empty prefix always": {tags: model.Tags{"team:infra"}, prefix: "", want: true}, - } - - for name, tc := range tests { - t.Run(name, func(t *testing.T) { - t.Parallel() - assert.Equal(t, tc.want, hasTag(tc.tags, tc.prefix)) - }) - } -} diff --git a/aws/logs_monitoring_go/internal/handling/lambda.go b/aws/logs_monitoring_go/internal/handling/lambda.go new file mode 100644 index 000000000..9b42ef217 --- /dev/null +++ b/aws/logs_monitoring_go/internal/handling/lambda.go @@ -0,0 +1,34 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2026-Present Datadog, Inc. + +package handling + +import ( + "strings" + + "github.com/DataDog/datadog-serverless-functions/aws/logs_monitoring_go/internal/model" +) + +func enrichLambdaLog(entry *model.LogEntry, forwarderARN, logGroup, logStream string) { + name := lambdaName(strings.ToLower(logGroup), logStream) + prefix, _, _ := strings.Cut(forwarderARN, "function:") + + arn := prefix + "function:" + name + entry.Tags.Add("functionname", name) + entry.Lambda = &model.LambdaLog{ARN: arn} + entry.Host = arn + entry.Service = name +} + +func lambdaName(logGroup, logStream string) string { + if m := lambdaLogStreamRegex.FindStringSubmatch(logStream); m != nil { + return strings.ToLower(m[1]) + } + + if _, name, found := strings.Cut(logGroup, logGroupLambda+"/"); found && name != "" { + return strings.ToLower(name) + } + return "" +} diff --git a/aws/logs_monitoring_go/internal/handling/lambda_test.go b/aws/logs_monitoring_go/internal/handling/lambda_test.go new file mode 100644 index 000000000..2ce262373 --- /dev/null +++ b/aws/logs_monitoring_go/internal/handling/lambda_test.go @@ -0,0 +1,38 @@ +// Unless explicitly stated otherwise all files in this repository are licensed +// under the Apache License Version 2.0. +// This product includes software developed at Datadog (https://www.datadoghq.com/). +// Copyright 2026-Present Datadog, Inc. + +package handling + +import ( + "testing" + + "github.com/stretchr/testify/assert" +) + +func TestLambdaName(t *testing.T) { + t.Parallel() + + tests := map[string]struct { + logGroup string + logStream string + want string + }{ + "default log group": {logGroup: "/aws/lambda/my-function", logStream: "stream", want: "my-function"}, + "default log group lowercased": {logGroup: "/aws/lambda/My-Function", logStream: "stream", want: "my-function"}, + "custom log group name from stream": {logGroup: "/aws/vendedlogs/states/anyLogGroupName", logStream: "2020/03/05/Test-Customized-LogGroup[$LATEST]20bddfd5a2dc4c6b97ac02800eae90d0", want: "test-customized-loggroup"}, + "stream takes priority over log group": {logGroup: "/aws/lambda/from-group", logStream: "2020/03/05/from-stream[$LATEST]20bddfd5a2dc4c6b97ac02800eae90d0", want: "from-stream"}, + "stream without function name": {logGroup: "my-custom-group", logStream: "2023/11/04/[$LATEST]4426346c2cdf4c54a74d3bd2b929fc44", want: ""}, + "non-lambda log group": {logGroup: "/aws/rds/cluster", logStream: "stream", want: ""}, + "empty": {logGroup: "", logStream: "", want: ""}, + } + + for name, tc := range tests { + t.Run(name, func(t *testing.T) { + t.Parallel() + + assert.Equal(t, tc.want, lambdaName(tc.logGroup, tc.logStream)) + }) + } +} From 8aeb533823026be3c6129dba9ff41cacd5f2741a Mon Sep 17 00:00:00 2001 From: Nabil Dakkoune Date: Wed, 22 Jul 2026 11:43:27 +0200 Subject: [PATCH 5/5] one service tag --- ...udwatch_lambda_with_dd_message.golden.json | 34 +++++++++++++++++++ ...oudwatch_lambda_with_dd_message.input.json | 15 ++++++++ .../internal/handling/cloudwatch.go | 4 +-- 3 files changed, 51 insertions(+), 2 deletions(-) create mode 100644 aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_with_dd_message.golden.json create mode 100644 aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_with_dd_message.input.json diff --git a/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_with_dd_message.golden.json b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_with_dd_message.golden.json new file mode 100644 index 000000000..0919790f7 --- /dev/null +++ b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_with_dd_message.golden.json @@ -0,0 +1,34 @@ +{ + "headers": { + "Content-Encoding": "gzip", + "Content-Type": "application/json", + "Dd-Api-Key": "abcdefghijklmnopqrstuvwxyz012345", + "Dd-Evp-Origin": "aws_forwarder", + "Dd-Evp-Origin-Version": "6.0", + "Dd-Storage-Tag": "cloudwatch", + "User-Agent": "Go-http-client/1.1" + }, + "body": [ + { + "host": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event", + "id": "35496429442806342619978265557671090556291002193281548291", + "timestamp": 1591714946151, + "message": "{\"message\":\"hello world\"}", + "service": "custom_service", + "ddsource": "lambda", + "ddsourcecategory": "aws", + "ddtags": "custom_tag1:value1,custom_tag2:value2,functionname:storms-cloudwatch-event,env:none,service:custom_service", + "aws": { + "invoked_function_arn": "arn:aws:lambda:us-east-1:123456789012:function:forwarder", + "awslogs": { + "logGroup": "/aws/lambda/storms-cloudwatch-event", + "logStream": "2020/06/09/[$LATEST]b249865adaaf4fad80f95f8ad09725b8", + "owner": "601427279990" + } + }, + "lambda": { + "arn": "arn:aws:lambda:us-east-1:123456789012:function:storms-cloudwatch-event" + } + } + ] +} diff --git a/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_with_dd_message.input.json b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_with_dd_message.input.json new file mode 100644 index 000000000..21e8f4350 --- /dev/null +++ b/aws/logs_monitoring_go/cmd/forwarder/testdata/cloudwatch_lambda_with_dd_message.input.json @@ -0,0 +1,15 @@ +{ + "messageType": "DATA_MESSAGE", + "owner": "601427279990", + "logGroup": "/aws/lambda/storms-cloudwatch-event", + "logStream": "2020/06/09/[$LATEST]b249865adaaf4fad80f95f8ad09725b8", + "subscriptionFilters": [ + "myevent" + ], + "logEvents": [{ + "id": "35496429442806342619978265557671090556291002193281548291", + "timestamp": 1591714946151, + "message": "{\"message\": \"hello world\", \"ddtags\": \"service:custom_service,custom_tag1:value1,custom_tag2:value2\"}\n" + } + ] +} \ No newline at end of file diff --git a/aws/logs_monitoring_go/internal/handling/cloudwatch.go b/aws/logs_monitoring_go/internal/handling/cloudwatch.go index 58dd697ee..a94c590f1 100644 --- a/aws/logs_monitoring_go/internal/handling/cloudwatch.go +++ b/aws/logs_monitoring_go/internal/handling/cloudwatch.go @@ -152,7 +152,7 @@ func (h cloudwatchHandler) newCloudwatchBaseEntry(data events.CloudwatchLogsData func (h cloudwatchHandler) newCloudwatchLogEntry(event events.CloudwatchLogsLogEvent, entry model.LogEntry) model.LogEntry { tags, service, message := extractFromMessage(event.Message) - entry.Service = cmp.Or(h.cfg.Service, service, entry.Service, entry.Source) + entry.Service = cmp.Or(service, entry.Service, h.cfg.Service, entry.Source) entry.Tags = slices.Concat(tags, entry.Tags, h.cfg.Tags) entry.Message = message entry.ID = event.ID @@ -162,7 +162,7 @@ func (h cloudwatchHandler) newCloudwatchLogEntry(event events.CloudwatchLogsLogE entry.Host = cloudtrailHost(event.Message) } - if entry.Lambda != nil { + if entry.Lambda != nil && !entry.Tags.Has("service") { entry.Tags.Add("service", entry.Service) }