From 888ab65a55ddbe7ea43ffbe7de2eebb7347c6f37 Mon Sep 17 00:00:00 2001 From: "james.eastham" Date: Fri, 19 Jun 2026 13:52:06 +0100 Subject: [PATCH 01/11] feat: add support for event bridge DSM context extraction --- datadog_lambda/config.py | 4 + datadog_lambda/tracing.py | 64 +++++++++++++++- tests/test_tracing.py | 153 ++++++++++++++++++++++++++++++++++++++ 3 files changed, 219 insertions(+), 2 deletions(-) diff --git a/datadog_lambda/config.py b/datadog_lambda/config.py index ce4924af8..a4e8a8179 100644 --- a/datadog_lambda/config.py +++ b/datadog_lambda/config.py @@ -89,6 +89,10 @@ def _resolve_env(self, key, default=None, cast=None, depends_on_tracing=False): data_streams_enabled = _get_env( "DD_DATA_STREAMS_ENABLED", "false", as_bool, depends_on_tracing=True ) + # EventBridge bus name used as the DSM `exchange` tag. The bus name is not + # present in the inbound event, so it must be provided explicitly to allow + # the consume checkpoint to pair with the EventBridge produce checkpoint. + dsm_exchange_name = _get_env("DD_DSM_EXCHANGE_NAME") appsec_enabled = _get_env("DD_APPSEC_ENABLED", "false", as_bool) sca_enabled = _get_env("DD_APPSEC_SCA_ENABLED", "false", as_bool) diff --git a/datadog_lambda/tracing.py b/datadog_lambda/tracing.py index b3f79a964..ae0994a8b 100644 --- a/datadog_lambda/tracing.py +++ b/datadog_lambda/tracing.py @@ -82,6 +82,43 @@ def _dsm_set_checkpoint(context_json, event_type, arn): ) +def _dsm_set_eventbridge_checkpoint(context_json, detail_type): + """Set a DSM consume checkpoint for an EventBridge event. + + Unlike the SQS/SNS/Kinesis helper, the EventBridge edge tags include an + `exchange` tag (the bus name) to mirror the produce-side tags so the + consume node pairs with the produce node. The bus name is not present in + the inbound event, so it is sourced from `DD_DSM_EXCHANGE_NAME` when set. + The public `set_consume_checkpoint` helper cannot emit an `exchange` tag, + so the lower-level processor API is used directly. + """ + if not config.data_streams_enabled: + return + + if not detail_type: + return + + try: + from ddtrace.internal.datastreams import data_streams_processor + from ddtrace.internal.datastreams.processor import PROPAGATION_KEY_BASE_64 + + processor = data_streams_processor() + if not processor: + return + + carrier_get = lambda k: context_json and context_json.get(k) # noqa: E731 + processor.decode_pathway_b64(carrier_get(PROPAGATION_KEY_BASE_64)) + + tags = ["direction:in", "topic:" + detail_type, "type:eventbridge"] + if config.dsm_exchange_name: + tags.append("exchange:" + config.dsm_exchange_name) + processor.set_checkpoint(tags) + except Exception as e: + logger.debug( + f"DSM:Failed to set consume checkpoint for eventbridge {detail_type}: {e}" + ) + + def _convert_xray_trace_id(xray_trace_id): """ Convert X-Ray trace id (hex)'s last 63 bits to a Datadog trace id (int). @@ -353,12 +390,32 @@ def _extract_context_from_eventbridge_sqs_event(event): This is only possible if first record in `Records` contains a `body` field which contains the EventBridge `detail` as a JSON string. """ - first_record = event.get("Records")[0] + records = event.get("Records") + first_record = records[0] body_str = first_record.get("body") body = json.loads(body_str) detail = body.get("detail") + # If `detail` is missing this is not an EventBridge -> SQS event; raising + # here lets the caller fall back to the regular SQS extraction path before + # any DSM checkpoint is set, avoiding double counting for plain SQS events. dd_context = detail.get("_datadog") + # The event has been confirmed as EventBridge -> SQS. Set a consume + # checkpoint for every record in the batch. The message is consumed from + # the SQS queue, so it follows SQS conventions (type:sqs, topic:queue ARN). + if config.data_streams_enabled: + for record in records: + try: + record_body = json.loads(record.get("body")) + record_context = (record_body.get("detail") or {}).get("_datadog") + _dsm_set_checkpoint( + record_context, "sqs", record.get("eventSourceARN", "") + ) + except Exception: + logger.debug( + "Failed to set DSM checkpoint for an EventBridge to SQS record." + ) + if is_step_function_event(dd_context): try: return extract_context_from_step_functions(dd_context, None) @@ -379,8 +436,11 @@ def extract_context_from_eventbridge_event(event, lambda_context): that header. """ try: - detail = event.get("detail") + detail = event.get("detail") or {} dd_context = detail.get("_datadog") + + _dsm_set_eventbridge_checkpoint(dd_context, event.get("detail-type")) + if not dd_context: return extract_context_from_lambda_context(lambda_context) diff --git a/tests/test_tracing.py b/tests/test_tracing.py index fc18f6e51..ae38702ea 100644 --- a/tests/test_tracing.py +++ b/tests/test_tracing.py @@ -46,6 +46,7 @@ _dsm_set_checkpoint, extract_context_from_kinesis_event, extract_context_from_sqs_or_sns_event_or_context, + extract_context_from_eventbridge_event, ) from datadog_lambda.trigger import parse_event_source @@ -3646,3 +3647,155 @@ def test_kinesis_data_streams_disabled(self): arn = "arn:aws:kinesis:us-east-1:123456789012:stream/test-stream" _dsm_set_checkpoint(context_json, event_type, arn) + + # EVENTBRIDGE -> SQS TESTS + + @staticmethod + def _eventbridge_sqs_record(queue_arn, pathway_ctx): + body = { + "detail-type": "MyDetailType", + "source": "my.event.source", + "detail": { + "_datadog": { + # Complete trace context so the extractor returns early and + # does not fall through to the regular SQS path. + "x-datadog-trace-id": "12345", + "x-datadog-parent-id": "67890", + "x-datadog-sampling-priority": "1", + "dd-pathway-ctx-base64": pathway_ctx, + } + }, + } + return { + "eventSourceARN": queue_arn, + "eventSource": "aws:sqs", + "body": json.dumps(body), + } + + def test_eventbridge_sqs_context_propagated(self): + queue_arn = "arn:aws:sqs:us-east-1:123456789012:eb-queue" + event = {"Records": [self._eventbridge_sqs_record(queue_arn, "12345")]} + + extract_context_from_sqs_or_sns_event_or_context( + event, self.lambda_context, parse_event_source(event) + ) + + # EventBridge -> SQS is consumed from the queue, so it uses SQS tags. + self.assertEqual(self.mock_checkpoint.call_count, 1) + args, _ = self.mock_checkpoint.call_args + self.assertEqual(args[0], "sqs") + self.assertEqual(args[1], queue_arn) + carrier_get = args[2] + self.assertEqual(carrier_get("dd-pathway-ctx-base64"), "12345") + + def test_eventbridge_sqs_checkpoints_all_records(self): + arn1 = "arn:aws:sqs:us-east-1:123456789012:eb-queue" + arn2 = "arn:aws:sqs:us-east-1:123456789012:eb-queue-2" + event = { + "Records": [ + self._eventbridge_sqs_record(arn1, "ctx-1"), + self._eventbridge_sqs_record(arn2, "ctx-2"), + ] + } + + extract_context_from_sqs_or_sns_event_or_context( + event, self.lambda_context, parse_event_source(event) + ) + + self.assertEqual(self.mock_checkpoint.call_count, 2) + first_args, _ = self.mock_checkpoint.call_args_list[0] + second_args, _ = self.mock_checkpoint.call_args_list[1] + self.assertEqual((first_args[0], first_args[1]), ("sqs", arn1)) + self.assertEqual(first_args[2]("dd-pathway-ctx-base64"), "ctx-1") + self.assertEqual((second_args[0], second_args[1]), ("sqs", arn2)) + self.assertEqual(second_args[2]("dd-pathway-ctx-base64"), "ctx-2") + + @patch("datadog_lambda.config.Config.data_streams_enabled", False) + def test_eventbridge_sqs_data_streams_disabled(self): + queue_arn = "arn:aws:sqs:us-east-1:123456789012:eb-queue" + event = {"Records": [self._eventbridge_sqs_record(queue_arn, "12345")]} + + extract_context_from_sqs_or_sns_event_or_context( + event, self.lambda_context, parse_event_source(event) + ) + + self.mock_checkpoint.assert_not_called() + + +class TestEventBridgeDSMLogic(unittest.TestCase): + def setUp(self): + self.lambda_context = get_mock_context() + self.mock_processor = Mock() + processor_patcher = patch( + "ddtrace.internal.datastreams.data_streams_processor", + return_value=self.mock_processor, + ) + processor_patcher.start() + self.addCleanup(processor_patcher.stop) + config_patcher = patch( + "datadog_lambda.config.Config.data_streams_enabled", True + ) + config_patcher.start() + self.addCleanup(config_patcher.stop) + + @staticmethod + def _eventbridge_event(detail_type="MyDetailType", pathway_ctx="12345"): + return { + "detail-type": detail_type, + "source": "my.event.source", + "detail": {"_datadog": {"dd-pathway-ctx-base64": pathway_ctx}}, + } + + def test_eventbridge_context_propagated(self): + event = self._eventbridge_event() + + extract_context_from_eventbridge_event(event, self.lambda_context) + + self.mock_processor.decode_pathway_b64.assert_called_once_with("12345") + self.mock_processor.set_checkpoint.assert_called_once() + (tags,), _ = self.mock_processor.set_checkpoint.call_args + self.assertIn("direction:in", tags) + self.assertIn("type:eventbridge", tags) + self.assertIn("topic:MyDetailType", tags) + self.assertFalse(any(t.startswith("exchange:") for t in tags)) + + @patch("datadog_lambda.config.Config.dsm_exchange_name", "my-event-bus") + def test_eventbridge_exchange_tag_from_env(self): + event = self._eventbridge_event() + + extract_context_from_eventbridge_event(event, self.lambda_context) + + (tags,), _ = self.mock_processor.set_checkpoint.call_args + self.assertIn("exchange:my-event-bus", tags) + self.assertIn("topic:MyDetailType", tags) + self.assertIn("type:eventbridge", tags) + + def test_eventbridge_no_detail_type_skips_checkpoint(self): + event = self._eventbridge_event(detail_type=None) + + extract_context_from_eventbridge_event(event, self.lambda_context) + + self.mock_processor.set_checkpoint.assert_not_called() + + def test_eventbridge_no_dd_context_still_checkpoints(self): + event = {"detail-type": "MyDetailType", "detail": {}} + + extract_context_from_eventbridge_event(event, self.lambda_context) + + self.mock_processor.decode_pathway_b64.assert_called_once_with(None) + self.mock_processor.set_checkpoint.assert_called_once() + + def test_eventbridge_missing_detail_still_checkpoints(self): + event = {"detail-type": "MyDetailType"} + + extract_context_from_eventbridge_event(event, self.lambda_context) + + self.mock_processor.set_checkpoint.assert_called_once() + + @patch("datadog_lambda.config.Config.data_streams_enabled", False) + def test_eventbridge_data_streams_disabled(self): + event = self._eventbridge_event() + + extract_context_from_eventbridge_event(event, self.lambda_context) + + self.mock_processor.set_checkpoint.assert_not_called() From 8fc0cd0295f7dad6d70d72344a6e48dc78bc2da9 Mon Sep 17 00:00:00 2001 From: "james.eastham" Date: Tue, 28 Jul 2026 08:43:44 +0100 Subject: [PATCH 02/11] chore: update processing of eventbridge DSM --- datadog_lambda/tracing.py | 46 ++++++++++++++++++++----------- tests/test_tracing.py | 57 ++++++++++++++++++++++++++++++--------- 2 files changed, 76 insertions(+), 27 deletions(-) diff --git a/datadog_lambda/tracing.py b/datadog_lambda/tracing.py index ae0994a8b..a1017e0d4 100644 --- a/datadog_lambda/tracing.py +++ b/datadog_lambda/tracing.py @@ -99,15 +99,16 @@ def _dsm_set_eventbridge_checkpoint(context_json, detail_type): return try: - from ddtrace.internal.datastreams import data_streams_processor - from ddtrace.internal.datastreams.processor import PROPAGATION_KEY_BASE_64 + from ddtrace.data_streams import PROPAGATION_KEY_BASE_64 + from ddtrace.data_streams import ddtrace as ddtrace_data_streams - processor = data_streams_processor() + processor = getattr(ddtrace_data_streams.tracer, "data_streams_processor", None) if not processor: return - carrier_get = lambda k: context_json and context_json.get(k) # noqa: E731 - processor.decode_pathway_b64(carrier_get(PROPAGATION_KEY_BASE_64)) + processor.decode_pathway_b64( + context_json.get(PROPAGATION_KEY_BASE_64) if context_json else None + ) tags = ["direction:in", "topic:" + detail_type, "type:eventbridge"] if config.dsm_exchange_name: @@ -289,9 +290,13 @@ def extract_context_from_sqs_or_sns_event_or_context( # EventBridge => SQS try: - context = _extract_context_from_eventbridge_sqs_event(event) - if _is_context_complete(context): - return context + context, is_eventbridge_sqs = _extract_context_from_eventbridge_sqs_event( + event + ) + if is_eventbridge_sqs: + if _is_context_complete(context): + return context + return extract_context_from_lambda_context(lambda_context) except Exception: logger.debug("Failed extracting context as EventBridge to SQS.") @@ -390,24 +395,35 @@ def _extract_context_from_eventbridge_sqs_event(event): This is only possible if first record in `Records` contains a `body` field which contains the EventBridge `detail` as a JSON string. """ - records = event.get("Records") + records = event.get("Records") or [] + if not records: + return None, False + first_record = records[0] body_str = first_record.get("body") body = json.loads(body_str) detail = body.get("detail") - # If `detail` is missing this is not an EventBridge -> SQS event; raising - # here lets the caller fall back to the regular SQS extraction path before - # any DSM checkpoint is set, avoiding double counting for plain SQS events. + if not isinstance(detail, dict): + return None, False + dd_context = detail.get("_datadog") # The event has been confirmed as EventBridge -> SQS. Set a consume # checkpoint for every record in the batch. The message is consumed from # the SQS queue, so it follows SQS conventions (type:sqs, topic:queue ARN). if config.data_streams_enabled: + _dsm_set_checkpoint(dd_context, "sqs", first_record.get("eventSourceARN", "")) for record in records: + if record is first_record: + continue try: record_body = json.loads(record.get("body")) - record_context = (record_body.get("detail") or {}).get("_datadog") + record_detail = record_body.get("detail") + record_context = ( + record_detail.get("_datadog") + if isinstance(record_detail, dict) + else None + ) _dsm_set_checkpoint( record_context, "sqs", record.get("eventSourceARN", "") ) @@ -418,13 +434,13 @@ def _extract_context_from_eventbridge_sqs_event(event): if is_step_function_event(dd_context): try: - return extract_context_from_step_functions(dd_context, None) + return extract_context_from_step_functions(dd_context, None), True except Exception: logger.debug( "Failed to extract Step Functions context from EventBridge to SQS event." ) - return propagator.extract(dd_context) + return propagator.extract(dd_context), True def extract_context_from_eventbridge_event(event, lambda_context): diff --git a/tests/test_tracing.py b/tests/test_tracing.py index ae38702ea..9cd01337b 100644 --- a/tests/test_tracing.py +++ b/tests/test_tracing.py @@ -3651,19 +3651,23 @@ def test_kinesis_data_streams_disabled(self): # EVENTBRIDGE -> SQS TESTS @staticmethod - def _eventbridge_sqs_record(queue_arn, pathway_ctx): - body = { - "detail-type": "MyDetailType", - "source": "my.event.source", - "detail": { - "_datadog": { + def _eventbridge_sqs_record(queue_arn, pathway_ctx, include_trace_headers=True): + dd_context = {"dd-pathway-ctx-base64": pathway_ctx} + if include_trace_headers: + dd_context.update( + { # Complete trace context so the extractor returns early and # does not fall through to the regular SQS path. "x-datadog-trace-id": "12345", "x-datadog-parent-id": "67890", "x-datadog-sampling-priority": "1", - "dd-pathway-ctx-base64": pathway_ctx, } + ) + body = { + "detail-type": "MyDetailType", + "source": "my.event.source", + "detail": { + "_datadog": dd_context }, } return { @@ -3710,6 +3714,35 @@ def test_eventbridge_sqs_checkpoints_all_records(self): self.assertEqual((second_args[0], second_args[1]), ("sqs", arn2)) self.assertEqual(second_args[2]("dd-pathway-ctx-base64"), "ctx-2") + @patch( + "datadog_lambda.tracing.extract_context_from_lambda_context", + return_value=Context(trace_id=111, span_id=222, sampling_priority=1), + ) + def test_eventbridge_sqs_incomplete_context_uses_single_checkpoint( + self, mock_extract_context + ): + queue_arn = "arn:aws:sqs:us-east-1:123456789012:eb-queue" + event = { + "Records": [ + self._eventbridge_sqs_record( + queue_arn, "ctx-only", include_trace_headers=False + ) + ] + } + + context = extract_context_from_sqs_or_sns_event_or_context( + event, self.lambda_context, parse_event_source(event) + ) + + self.assertEqual(context.trace_id, 111) + self.assertEqual(context.span_id, 222) + mock_extract_context.assert_called_once_with(self.lambda_context) + self.assertEqual(self.mock_checkpoint.call_count, 1) + args, _ = self.mock_checkpoint.call_args + self.assertEqual(args[0], "sqs") + self.assertEqual(args[1], queue_arn) + self.assertEqual(args[2]("dd-pathway-ctx-base64"), "ctx-only") + @patch("datadog_lambda.config.Config.data_streams_enabled", False) def test_eventbridge_sqs_data_streams_disabled(self): queue_arn = "arn:aws:sqs:us-east-1:123456789012:eb-queue" @@ -3726,12 +3759,12 @@ class TestEventBridgeDSMLogic(unittest.TestCase): def setUp(self): self.lambda_context = get_mock_context() self.mock_processor = Mock() - processor_patcher = patch( - "ddtrace.internal.datastreams.data_streams_processor", - return_value=self.mock_processor, + tracer_patcher = patch( + "ddtrace.data_streams.ddtrace.tracer", + new=SimpleNamespace(data_streams_processor=self.mock_processor), ) - processor_patcher.start() - self.addCleanup(processor_patcher.stop) + tracer_patcher.start() + self.addCleanup(tracer_patcher.stop) config_patcher = patch( "datadog_lambda.config.Config.data_streams_enabled", True ) From 9dcc6ecd02a36f755f03f93ab6362061520932a3 Mon Sep 17 00:00:00 2001 From: "james.eastham" Date: Thu, 30 Jul 2026 12:03:39 +0100 Subject: [PATCH 03/11] refactor: avoid extra allocation in eventbridge sqs parsing --- datadog_lambda/tracing.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/datadog_lambda/tracing.py b/datadog_lambda/tracing.py index a1017e0d4..991c2b2c9 100644 --- a/datadog_lambda/tracing.py +++ b/datadog_lambda/tracing.py @@ -395,7 +395,7 @@ def _extract_context_from_eventbridge_sqs_event(event): This is only possible if first record in `Records` contains a `body` field which contains the EventBridge `detail` as a JSON string. """ - records = event.get("Records") or [] + records = event.get("Records") if not records: return None, False From a45866f66e0babbd9fd5a2ecdae355616bce889e Mon Sep 17 00:00:00 2001 From: "james.eastham" Date: Thu, 30 Jul 2026 13:45:32 +0100 Subject: [PATCH 04/11] fix: require eventbridge envelope for sqs extraction --- datadog_lambda/tracing.py | 9 ++++++++- tests/test_tracing.py | 37 +++++++++++++++++++++++++++++++++++++ 2 files changed, 45 insertions(+), 1 deletion(-) diff --git a/datadog_lambda/tracing.py b/datadog_lambda/tracing.py index 991c2b2c9..acb7d3153 100644 --- a/datadog_lambda/tracing.py +++ b/datadog_lambda/tracing.py @@ -402,8 +402,15 @@ def _extract_context_from_eventbridge_sqs_event(event): first_record = records[0] body_str = first_record.get("body") body = json.loads(body_str) + if not isinstance(body, dict): + return None, False + detail = body.get("detail") - if not isinstance(detail, dict): + if not ( + isinstance(detail, dict) + and body.get("detail-type") + and body.get("source") + ): return None, False dd_context = detail.get("_datadog") diff --git a/tests/test_tracing.py b/tests/test_tracing.py index 9cd01337b..38828386d 100644 --- a/tests/test_tracing.py +++ b/tests/test_tracing.py @@ -2815,6 +2815,43 @@ def test_sqs_no_datadog_message_attribute(self): # None indicates no DSM context propagation self.assertEqual(carrier_get("dd-pathway-ctx-base64"), None) + @patch("datadog_lambda.tracing.extract_context_from_lambda_context") + def test_sqs_detail_body_still_uses_message_attributes(self, mock_extract_context): + dd_data = { + "x-datadog-trace-id": "12345", + "x-datadog-parent-id": "67890", + "x-datadog-sampling-priority": "1", + "dd-pathway-ctx-base64": "attr-ctx", + } + dd_json_data = json.dumps(dd_data) + + event = { + "Records": [ + { + "eventSourceARN": "arn:aws:sqs:us-east-1:123456789012:test-queue", + "messageAttributes": { + "_datadog": {"dataType": "String", "stringValue": dd_json_data} + }, + "eventSource": "aws:sqs", + "body": json.dumps({"detail": {"application": "payload"}}), + } + ] + } + + context = extract_context_from_sqs_or_sns_event_or_context( + event, self.lambda_context, parse_event_source(event) + ) + + mock_extract_context.assert_not_called() + self.assertEqual(context.trace_id, 12345) + self.assertEqual(context.span_id, 67890) + self.assertEqual(context.sampling_priority, 1) + self.assertEqual(self.mock_checkpoint.call_count, 1) + args, _ = self.mock_checkpoint.call_args + self.assertEqual(args[0], "sqs") + self.assertEqual(args[1], "arn:aws:sqs:us-east-1:123456789012:test-queue") + self.assertEqual(args[2]("dd-pathway-ctx-base64"), "attr-ctx") + def test_sqs_empty_datadog_message_attribute(self): event = { "Records": [ From ff42108b08082f98098c271856eb0fcaaed42ead Mon Sep 17 00:00:00 2001 From: "james.eastham" Date: Thu, 30 Jul 2026 13:52:13 +0100 Subject: [PATCH 05/11] fix: detect sqs carriers per record in eventbridge batches --- datadog_lambda/tracing.py | 76 ++++++++++++++++++++++++++------------- tests/test_tracing.py | 42 ++++++++++++++++++++-- 2 files changed, 90 insertions(+), 28 deletions(-) diff --git a/datadog_lambda/tracing.py b/datadog_lambda/tracing.py index acb7d3153..5508fb09d 100644 --- a/datadog_lambda/tracing.py +++ b/datadog_lambda/tracing.py @@ -290,9 +290,7 @@ def extract_context_from_sqs_or_sns_event_or_context( # EventBridge => SQS try: - context, is_eventbridge_sqs = _extract_context_from_eventbridge_sqs_event( - event - ) + context, is_eventbridge_sqs = _extract_context_from_eventbridge_sqs_event(event) if is_eventbridge_sqs: if _is_context_complete(context): return context @@ -400,37 +398,25 @@ def _extract_context_from_eventbridge_sqs_event(event): return None, False first_record = records[0] - body_str = first_record.get("body") - body = json.loads(body_str) - if not isinstance(body, dict): - return None, False - - detail = body.get("detail") - if not ( - isinstance(detail, dict) - and body.get("detail-type") - and body.get("source") - ): + dd_context, is_eventbridge_sqs = _extract_eventbridge_sqs_record_context( + first_record + ) + if not is_eventbridge_sqs: return None, False - dd_context = detail.get("_datadog") - # The event has been confirmed as EventBridge -> SQS. Set a consume # checkpoint for every record in the batch. The message is consumed from # the SQS queue, so it follows SQS conventions (type:sqs, topic:queue ARN). if config.data_streams_enabled: - _dsm_set_checkpoint(dd_context, "sqs", first_record.get("eventSourceARN", "")) for record in records: - if record is first_record: - continue try: - record_body = json.loads(record.get("body")) - record_detail = record_body.get("detail") - record_context = ( - record_detail.get("_datadog") - if isinstance(record_detail, dict) - else None + record_context, is_eventbridge_record = ( + _extract_eventbridge_sqs_record_context(record) ) + if not is_eventbridge_record: + record_context = _extract_sqs_record_message_attribute_context( + record + ) _dsm_set_checkpoint( record_context, "sqs", record.get("eventSourceARN", "") ) @@ -450,6 +436,46 @@ def _extract_context_from_eventbridge_sqs_event(event): return propagator.extract(dd_context), True +def _extract_eventbridge_sqs_record_context(record): + body_str = record.get("body") + body = json.loads(body_str) + if not isinstance(body, dict): + return None, False + + detail = body.get("detail") + if not ( + isinstance(detail, dict) and body.get("detail-type") and body.get("source") + ): + return None, False + + return detail.get("_datadog"), True + + +def _extract_sqs_record_message_attribute_context(record): + msg_attributes = record.get("messageAttributes") or {} + dd_payload = msg_attributes.get("_datadog") + if not dd_payload: + return None + + dd_json_data = None + dd_json_data_type = dd_payload.get("Type") or dd_payload.get("dataType") + if dd_json_data_type == "Binary": + import base64 + + dd_json_data = dd_payload.get("binaryValue") or dd_payload.get("Value") + if dd_json_data: + dd_json_data = base64.b64decode(dd_json_data) + elif dd_json_data_type == "String": + dd_json_data = dd_payload.get("stringValue") or dd_payload.get("Value") + else: + logger.debug( + "Datadog Lambda Python only supports extracting trace" + "context from String or Binary SQS/SNS message attributes" + ) + + return json.loads(dd_json_data) if dd_json_data else None + + def extract_context_from_eventbridge_event(event, lambda_context): """ Extract datadog trace context from an EventBridge message's Details. diff --git a/tests/test_tracing.py b/tests/test_tracing.py index 38828386d..bd4ac527a 100644 --- a/tests/test_tracing.py +++ b/tests/test_tracing.py @@ -3703,9 +3703,7 @@ def _eventbridge_sqs_record(queue_arn, pathway_ctx, include_trace_headers=True): body = { "detail-type": "MyDetailType", "source": "my.event.source", - "detail": { - "_datadog": dd_context - }, + "detail": {"_datadog": dd_context}, } return { "eventSourceARN": queue_arn, @@ -3751,6 +3749,44 @@ def test_eventbridge_sqs_checkpoints_all_records(self): self.assertEqual((second_args[0], second_args[1]), ("sqs", arn2)) self.assertEqual(second_args[2]("dd-pathway-ctx-base64"), "ctx-2") + def test_eventbridge_sqs_mixed_batch_uses_per_record_carriers(self): + arn1 = "arn:aws:sqs:us-east-1:123456789012:eb-queue" + arn2 = "arn:aws:sqs:us-east-1:123456789012:direct-queue" + second_dd_data = { + "x-datadog-trace-id": "12345", + "x-datadog-parent-id": "67890", + "x-datadog-sampling-priority": "1", + "dd-pathway-ctx-base64": "sqs-ctx", + } + event = { + "Records": [ + self._eventbridge_sqs_record(arn1, "eb-ctx"), + { + "eventSourceARN": arn2, + "eventSource": "aws:sqs", + "body": json.dumps({"message": "direct sqs payload"}), + "messageAttributes": { + "_datadog": { + "dataType": "String", + "stringValue": json.dumps(second_dd_data), + } + }, + }, + ] + } + + extract_context_from_sqs_or_sns_event_or_context( + event, self.lambda_context, parse_event_source(event) + ) + + self.assertEqual(self.mock_checkpoint.call_count, 2) + first_args, _ = self.mock_checkpoint.call_args_list[0] + second_args, _ = self.mock_checkpoint.call_args_list[1] + self.assertEqual((first_args[0], first_args[1]), ("sqs", arn1)) + self.assertEqual(first_args[2]("dd-pathway-ctx-base64"), "eb-ctx") + self.assertEqual((second_args[0], second_args[1]), ("sqs", arn2)) + self.assertEqual(second_args[2]("dd-pathway-ctx-base64"), "sqs-ctx") + @patch( "datadog_lambda.tracing.extract_context_from_lambda_context", return_value=Context(trace_id=111, span_id=222, sampling_priority=1), From 24da38968396f1092dadd90b3776c4aa4d071790 Mon Sep 17 00:00:00 2001 From: "james.eastham" Date: Thu, 30 Jul 2026 14:15:05 +0100 Subject: [PATCH 06/11] fix: use public dsm checkpoint for eventbridge --- datadog_lambda/tracing.py | 37 ++++----------------------- tests/test_tracing.py | 54 ++++++++++++++++++++++----------------- 2 files changed, 35 insertions(+), 56 deletions(-) diff --git a/datadog_lambda/tracing.py b/datadog_lambda/tracing.py index 5508fb09d..dded01811 100644 --- a/datadog_lambda/tracing.py +++ b/datadog_lambda/tracing.py @@ -85,39 +85,12 @@ def _dsm_set_checkpoint(context_json, event_type, arn): def _dsm_set_eventbridge_checkpoint(context_json, detail_type): """Set a DSM consume checkpoint for an EventBridge event. - Unlike the SQS/SNS/Kinesis helper, the EventBridge edge tags include an - `exchange` tag (the bus name) to mirror the produce-side tags so the - consume node pairs with the produce node. The bus name is not present in - the inbound event, so it is sourced from `DD_DSM_EXCHANGE_NAME` when set. - The public `set_consume_checkpoint` helper cannot emit an `exchange` tag, - so the lower-level processor API is used directly. + Temporary note: richer EventBridge consume checkpoint tagging is being + upstreamed into dd-trace-py. Until that lands, keep this on the same + public consume checkpoint API used by SQS/SNS/Kinesis and omit the + EventBridge exchange tag for now. """ - if not config.data_streams_enabled: - return - - if not detail_type: - return - - try: - from ddtrace.data_streams import PROPAGATION_KEY_BASE_64 - from ddtrace.data_streams import ddtrace as ddtrace_data_streams - - processor = getattr(ddtrace_data_streams.tracer, "data_streams_processor", None) - if not processor: - return - - processor.decode_pathway_b64( - context_json.get(PROPAGATION_KEY_BASE_64) if context_json else None - ) - - tags = ["direction:in", "topic:" + detail_type, "type:eventbridge"] - if config.dsm_exchange_name: - tags.append("exchange:" + config.dsm_exchange_name) - processor.set_checkpoint(tags) - except Exception as e: - logger.debug( - f"DSM:Failed to set consume checkpoint for eventbridge {detail_type}: {e}" - ) + _dsm_set_checkpoint(context_json, "eventbridge", detail_type) def _convert_xray_trace_id(xray_trace_id): diff --git a/tests/test_tracing.py b/tests/test_tracing.py index bd4ac527a..151d76031 100644 --- a/tests/test_tracing.py +++ b/tests/test_tracing.py @@ -3831,13 +3831,9 @@ def test_eventbridge_sqs_data_streams_disabled(self): class TestEventBridgeDSMLogic(unittest.TestCase): def setUp(self): self.lambda_context = get_mock_context() - self.mock_processor = Mock() - tracer_patcher = patch( - "ddtrace.data_streams.ddtrace.tracer", - new=SimpleNamespace(data_streams_processor=self.mock_processor), - ) - tracer_patcher.start() - self.addCleanup(tracer_patcher.stop) + checkpoint_patcher = patch("ddtrace.data_streams.set_consume_checkpoint") + self.mock_checkpoint = checkpoint_patcher.start() + self.addCleanup(checkpoint_patcher.stop) config_patcher = patch( "datadog_lambda.config.Config.data_streams_enabled", True ) @@ -3857,46 +3853,56 @@ def test_eventbridge_context_propagated(self): extract_context_from_eventbridge_event(event, self.lambda_context) - self.mock_processor.decode_pathway_b64.assert_called_once_with("12345") - self.mock_processor.set_checkpoint.assert_called_once() - (tags,), _ = self.mock_processor.set_checkpoint.call_args - self.assertIn("direction:in", tags) - self.assertIn("type:eventbridge", tags) - self.assertIn("topic:MyDetailType", tags) - self.assertFalse(any(t.startswith("exchange:") for t in tags)) + self.mock_checkpoint.assert_called_once() + args, kwargs = self.mock_checkpoint.call_args + self.assertEqual(args[0], "eventbridge") + self.assertEqual(args[1], "MyDetailType") + carrier_get = args[2] + self.assertEqual(carrier_get("dd-pathway-ctx-base64"), "12345") + self.assertEqual(kwargs, {"manual_checkpoint": False}) @patch("datadog_lambda.config.Config.dsm_exchange_name", "my-event-bus") - def test_eventbridge_exchange_tag_from_env(self): + def test_eventbridge_exchange_name_ignored_until_upstream_support_lands(self): event = self._eventbridge_event() extract_context_from_eventbridge_event(event, self.lambda_context) - (tags,), _ = self.mock_processor.set_checkpoint.call_args - self.assertIn("exchange:my-event-bus", tags) - self.assertIn("topic:MyDetailType", tags) - self.assertIn("type:eventbridge", tags) + self.mock_checkpoint.assert_called_once() + args, kwargs = self.mock_checkpoint.call_args + self.assertEqual(args[0], "eventbridge") + self.assertEqual(args[1], "MyDetailType") + carrier_get = args[2] + self.assertEqual(carrier_get("dd-pathway-ctx-base64"), "12345") + self.assertEqual(kwargs, {"manual_checkpoint": False}) def test_eventbridge_no_detail_type_skips_checkpoint(self): event = self._eventbridge_event(detail_type=None) extract_context_from_eventbridge_event(event, self.lambda_context) - self.mock_processor.set_checkpoint.assert_not_called() + self.mock_checkpoint.assert_not_called() def test_eventbridge_no_dd_context_still_checkpoints(self): event = {"detail-type": "MyDetailType", "detail": {}} extract_context_from_eventbridge_event(event, self.lambda_context) - self.mock_processor.decode_pathway_b64.assert_called_once_with(None) - self.mock_processor.set_checkpoint.assert_called_once() + self.mock_checkpoint.assert_called_once() + args, kwargs = self.mock_checkpoint.call_args + carrier_get = args[2] + self.assertIsNone(carrier_get("dd-pathway-ctx-base64")) + self.assertEqual(kwargs, {"manual_checkpoint": False}) def test_eventbridge_missing_detail_still_checkpoints(self): event = {"detail-type": "MyDetailType"} extract_context_from_eventbridge_event(event, self.lambda_context) - self.mock_processor.set_checkpoint.assert_called_once() + self.mock_checkpoint.assert_called_once() + args, kwargs = self.mock_checkpoint.call_args + carrier_get = args[2] + self.assertIsNone(carrier_get("dd-pathway-ctx-base64")) + self.assertEqual(kwargs, {"manual_checkpoint": False}) @patch("datadog_lambda.config.Config.data_streams_enabled", False) def test_eventbridge_data_streams_disabled(self): @@ -3904,4 +3910,4 @@ def test_eventbridge_data_streams_disabled(self): extract_context_from_eventbridge_event(event, self.lambda_context) - self.mock_processor.set_checkpoint.assert_not_called() + self.mock_checkpoint.assert_not_called() From fd813a34899c5210f75a1b47fe36eb0ec86cdb6b Mon Sep 17 00:00:00 2001 From: James Eastham Date: Thu, 3 Sep 2026 11:41:13 +0100 Subject: [PATCH 07/11] chore: update to handle new public DSM API --- datadog_lambda/tracing.py | 170 ++++++++++++++++++++++------ tests/test_tracing.py | 230 +++++++++++++++++++++++++++++++++++++- 2 files changed, 359 insertions(+), 41 deletions(-) diff --git a/datadog_lambda/tracing.py b/datadog_lambda/tracing.py index 7e9a5d576..d6b6d7075 100644 --- a/datadog_lambda/tracing.py +++ b/datadog_lambda/tracing.py @@ -86,12 +86,42 @@ def _dsm_set_checkpoint(context_json, event_type, arn): def _dsm_set_eventbridge_checkpoint(context_json, detail_type): """Set a DSM consume checkpoint for an EventBridge event. - Temporary note: richer EventBridge consume checkpoint tagging is being - upstreamed into dd-trace-py. Until that lands, keep this on the same - public consume checkpoint API used by SQS/SNS/Kinesis and omit the - EventBridge exchange tag for now. + When ddtrace >= 4.14 is present the public ``tags`` parameter is used to + attach the ``exchange:`` edge tag (sourced from ``DD_DSM_EXCHANGE_NAME``). + On older installs the checkpoint is still emitted, just without the + exchange tag. """ - _dsm_set_checkpoint(context_json, "eventbridge", detail_type) + if not config.data_streams_enabled or not detail_type: + return + + try: + from ddtrace.data_streams import set_consume_checkpoint + + carrier_get = lambda k: context_json and context_json.get(k) # noqa: E731 + try: + tags = ( + ["exchange:" + config.dsm_exchange_name] + if config.dsm_exchange_name + else None + ) + set_consume_checkpoint( + "eventbridge", + detail_type, + carrier_get, + manual_checkpoint=False, + tags=tags, + ) + except TypeError: + # ddtrace < 4.14 has no `tags` parameter. Retry without it, but + # keep `manual_checkpoint=False` so the checkpoint stays + # consistent with every other consume checkpoint. + set_consume_checkpoint( + "eventbridge", detail_type, carrier_get, manual_checkpoint=False + ) + except Exception as e: + logger.debug( + f"DSM:Failed to set consume checkpoint for eventbridge {detail_type}: {e}" + ) def _convert_xray_trace_id(xray_trace_id): @@ -262,9 +292,17 @@ def extract_context_from_sqs_or_sns_event_or_context( source_arn = "" event_type = "sqs" if event_source.equals(EventTypes.SQS) else "sns" - # EventBridge => SQS + # EventBridge => SQS. `dsm_handled` is True when the batch contained at + # least one EventBridge delivery and DSM checkpoints were already set + # per-record; in that case the regular SQS path below must not set another + # checkpoint for the first record or it would be double counted. + dsm_handled = False try: - context, is_eventbridge_sqs = _extract_context_from_eventbridge_sqs_event(event) + ( + context, + is_eventbridge_sqs, + dsm_handled, + ) = _extract_context_from_eventbridge_sqs_event(event) if is_eventbridge_sqs: if _is_context_complete(context): return context @@ -325,7 +363,8 @@ def extract_context_from_sqs_or_sns_event_or_context( "Failed to extract Step Functions context from SQS/SNS event." ) context = propagator.extract(dd_data) - _dsm_set_checkpoint(dd_data, event_type, source_arn) + if not dsm_handled: + _dsm_set_checkpoint(dd_data, event_type, source_arn) return context else: # Handle case where trace context is injected into attributes.AWSTraceHeader @@ -350,12 +389,14 @@ def extract_context_from_sqs_or_sns_event_or_context( sampling_priority=float(x_ray_context["sampled"]), ) # Still want to set a DSM checkpoint even if DSM context not propagated - _dsm_set_checkpoint(None, event_type, source_arn) + if not dsm_handled: + _dsm_set_checkpoint(None, event_type, source_arn) return extract_context_from_lambda_context(lambda_context) except Exception as e: logger.debug("The trace extractor returned with error %s", e) # Still want to set a DSM checkpoint even if DSM context not propagated - _dsm_set_checkpoint(None, event_type, source_arn) + if not dsm_handled: + _dsm_set_checkpoint(None, event_type, source_arn) return extract_context_from_lambda_context(lambda_context) @@ -366,53 +407,112 @@ def _extract_context_from_eventbridge_sqs_event(event): This is only possible if first record in `Records` contains a `body` field which contains the EventBridge `detail` as a JSON string. + + Returns a tuple ``(context, is_eventbridge_sqs, dsm_handled)``: + + * ``context`` / ``is_eventbridge_sqs`` describe the trace context and are + derived only from the first record, since that is the record whose trace + context becomes the Lambda's parent. + * ``dsm_handled`` reports whether this function already set DSM checkpoints + for the batch. It is ``True`` whenever *any* record in the batch is an + EventBridge delivery, so the caller must not set its own SQS checkpoint + (which would double count the first record). """ records = event.get("Records") if not records: - return None, False + return None, False, False first_record = records[0] dd_context, is_eventbridge_sqs = _extract_eventbridge_sqs_record_context( first_record ) + + # Set a consume checkpoint for every record in the batch whenever the batch + # contains at least one EventBridge delivery. Each record is classified + # independently so a mixed batch (EventBridge deliveries alongside direct + # SQS sends, in either order) uses the correct carrier per record. The + # message is consumed from the SQS queue, so it follows SQS conventions + # (type:sqs, topic:queue ARN). + dsm_handled = _dsm_set_eventbridge_sqs_batch_checkpoints(records) + if not is_eventbridge_sqs: - return None, False + return None, False, dsm_handled - # The event has been confirmed as EventBridge -> SQS. Set a consume - # checkpoint for every record in the batch. The message is consumed from - # the SQS queue, so it follows SQS conventions (type:sqs, topic:queue ARN). - if config.data_streams_enabled: - for record in records: + if is_step_function_event(dd_context): + try: + return ( + extract_context_from_step_functions(dd_context, None), + True, + dsm_handled, + ) + except Exception: + logger.debug( + "Failed to extract Step Functions context from EventBridge to SQS event." + ) + + return propagator.extract(dd_context), True, dsm_handled + + +def _dsm_set_eventbridge_sqs_batch_checkpoints(records): + """Set a per-record SQS DSM consume checkpoint for an EventBridge -> SQS + batch. + + Returns ``True`` when the batch contains at least one EventBridge delivery + (and checkpoints were therefore this function's responsibility), otherwise + ``False`` so the caller can fall back to its regular SQS checkpoint path. + Each record is classified independently: EventBridge records use the + carrier embedded in ``body.detail._datadog`` while other records fall back + to the SQS ``messageAttributes._datadog`` carrier. + """ + if not config.data_streams_enabled: + return False + + record_carriers = [] + batch_has_eventbridge = False + for record in records: + record_context, is_eventbridge_record = _extract_eventbridge_sqs_record_context( + record + ) + if is_eventbridge_record: + batch_has_eventbridge = True + else: try: - record_context, is_eventbridge_record = ( - _extract_eventbridge_sqs_record_context(record) - ) - if not is_eventbridge_record: - record_context = _extract_sqs_record_message_attribute_context( - record - ) - _dsm_set_checkpoint( - record_context, "sqs", record.get("eventSourceARN", "") - ) + record_context = _extract_sqs_record_message_attribute_context(record) except Exception: - logger.debug( - "Failed to set DSM checkpoint for an EventBridge to SQS record." - ) + record_context = None + record_carriers.append(record_context) - if is_step_function_event(dd_context): + if not batch_has_eventbridge: + return False + + for record, record_context in zip(records, record_carriers): try: - return extract_context_from_step_functions(dd_context, None), True + _dsm_set_checkpoint(record_context, "sqs", record.get("eventSourceARN", "")) except Exception: logger.debug( - "Failed to extract Step Functions context from EventBridge to SQS event." + "Failed to set DSM checkpoint for an EventBridge to SQS record." ) - return propagator.extract(dd_context), True + return True def _extract_eventbridge_sqs_record_context(record): + """Classify a single SQS record as an EventBridge delivery and return its + DSM carrier. + + Returns a tuple ``(dd_context, is_eventbridge)``. A record is only treated + as EventBridge when its ``body`` is a JSON object carrying the EventBridge + envelope fields (``detail`` object plus ``detail-type`` and ``source``). A + non-JSON or non-envelope body is not an error here: it simply means the + record is a regular SQS message, so the caller can fall back to the SQS + message attribute carrier instead. + """ body_str = record.get("body") - body = json.loads(body_str) + try: + body = json.loads(body_str) + except (ValueError, TypeError): + return None, False + if not isinstance(body, dict): return None, False diff --git a/tests/test_tracing.py b/tests/test_tracing.py index 2cce41a23..17c924a70 100644 --- a/tests/test_tracing.py +++ b/tests/test_tracing.py @@ -1,6 +1,8 @@ import base64 +import contextlib import copy import functools +import inspect import json import traceback import pytest @@ -44,6 +46,7 @@ propagator, emit_telemetry_on_exception_outside_of_handler, _dsm_set_checkpoint, + _dsm_set_eventbridge_checkpoint, extract_context_from_kinesis_event, extract_context_from_sqs_or_sns_event_or_context, extract_context_from_eventbridge_event, @@ -52,7 +55,6 @@ from datadog_lambda.trigger import parse_event_source from tests.utils import get_mock_context, ClientContext - function_arn = "arn:aws:lambda:us-west-1:123457598159:function:python-layer-test" fake_xray_header_value = ( @@ -4060,6 +4062,88 @@ def test_eventbridge_sqs_mixed_batch_uses_per_record_carriers(self): self.assertEqual((second_args[0], second_args[1]), ("sqs", arn2)) self.assertEqual(second_args[2]("dd-pathway-ctx-base64"), "sqs-ctx") + def test_eventbridge_sqs_direct_first_record_still_checkpoints_eventbridge(self): + # First record is a direct SQS send, a later record is an EventBridge + # delivery. Every record must still be classified and checkpointed + # independently, and the first record must not be double counted. + arn1 = "arn:aws:sqs:us-east-1:123456789012:direct-queue" + arn2 = "arn:aws:sqs:us-east-1:123456789012:eb-queue" + first_dd_data = { + "x-datadog-trace-id": "12345", + "x-datadog-parent-id": "67890", + "x-datadog-sampling-priority": "1", + "dd-pathway-ctx-base64": "sqs-ctx", + } + event = { + "Records": [ + { + "eventSourceARN": arn1, + "eventSource": "aws:sqs", + "body": json.dumps({"message": "direct sqs payload"}), + "messageAttributes": { + "_datadog": { + "dataType": "String", + "stringValue": json.dumps(first_dd_data), + } + }, + }, + self._eventbridge_sqs_record(arn2, "eb-ctx"), + ] + } + + extract_context_from_sqs_or_sns_event_or_context( + event, self.lambda_context, parse_event_source(event) + ) + + self.assertEqual(self.mock_checkpoint.call_count, 2) + first_args, _ = self.mock_checkpoint.call_args_list[0] + second_args, _ = self.mock_checkpoint.call_args_list[1] + self.assertEqual((first_args[0], first_args[1]), ("sqs", arn1)) + self.assertEqual(first_args[2]("dd-pathway-ctx-base64"), "sqs-ctx") + self.assertEqual((second_args[0], second_args[1]), ("sqs", arn2)) + self.assertEqual(second_args[2]("dd-pathway-ctx-base64"), "eb-ctx") + + def test_eventbridge_sqs_non_json_body_falls_back_to_sqs_attributes(self): + # A later direct SQS record can legitimately carry a non-JSON body + # while its DSM carrier lives in messageAttributes._datadog. The + # unparseable body must not prevent that record's checkpoint. + arn1 = "arn:aws:sqs:us-east-1:123456789012:eb-queue" + arn2 = "arn:aws:sqs:us-east-1:123456789012:direct-queue" + second_dd_data = { + "x-datadog-trace-id": "12345", + "x-datadog-parent-id": "67890", + "x-datadog-sampling-priority": "1", + "dd-pathway-ctx-base64": "sqs-ctx", + } + event = { + "Records": [ + self._eventbridge_sqs_record(arn1, "eb-ctx"), + { + "eventSourceARN": arn2, + "eventSource": "aws:sqs", + "body": "plain text, not json", + "messageAttributes": { + "_datadog": { + "dataType": "String", + "stringValue": json.dumps(second_dd_data), + } + }, + }, + ] + } + + extract_context_from_sqs_or_sns_event_or_context( + event, self.lambda_context, parse_event_source(event) + ) + + self.assertEqual(self.mock_checkpoint.call_count, 2) + first_args, _ = self.mock_checkpoint.call_args_list[0] + second_args, _ = self.mock_checkpoint.call_args_list[1] + self.assertEqual((first_args[0], first_args[1]), ("sqs", arn1)) + self.assertEqual(first_args[2]("dd-pathway-ctx-base64"), "eb-ctx") + self.assertEqual((second_args[0], second_args[1]), ("sqs", arn2)) + self.assertEqual(second_args[2]("dd-pathway-ctx-base64"), "sqs-ctx") + @patch( "datadog_lambda.tracing.extract_context_from_lambda_context", return_value=Context(trace_id=111, span_id=222, sampling_priority=1), @@ -4132,10 +4216,10 @@ def test_eventbridge_context_propagated(self): self.assertEqual(args[1], "MyDetailType") carrier_get = args[2] self.assertEqual(carrier_get("dd-pathway-ctx-base64"), "12345") - self.assertEqual(kwargs, {"manual_checkpoint": False}) + self.assertEqual(kwargs, {"manual_checkpoint": False, "tags": None}) @patch("datadog_lambda.config.Config.dsm_exchange_name", "my-event-bus") - def test_eventbridge_exchange_name_ignored_until_upstream_support_lands(self): + def test_eventbridge_exchange_name_used_when_upstream_support_present(self): event = self._eventbridge_event() extract_context_from_eventbridge_event(event, self.lambda_context) @@ -4146,7 +4230,9 @@ def test_eventbridge_exchange_name_ignored_until_upstream_support_lands(self): self.assertEqual(args[1], "MyDetailType") carrier_get = args[2] self.assertEqual(carrier_get("dd-pathway-ctx-base64"), "12345") - self.assertEqual(kwargs, {"manual_checkpoint": False}) + self.assertEqual( + kwargs, {"manual_checkpoint": False, "tags": ["exchange:my-event-bus"]} + ) def test_eventbridge_no_detail_type_skips_checkpoint(self): event = self._eventbridge_event(detail_type=None) @@ -4164,7 +4250,7 @@ def test_eventbridge_no_dd_context_still_checkpoints(self): args, kwargs = self.mock_checkpoint.call_args carrier_get = args[2] self.assertIsNone(carrier_get("dd-pathway-ctx-base64")) - self.assertEqual(kwargs, {"manual_checkpoint": False}) + self.assertEqual(kwargs, {"manual_checkpoint": False, "tags": None}) def test_eventbridge_missing_detail_still_checkpoints(self): event = {"detail-type": "MyDetailType"} @@ -4175,7 +4261,7 @@ def test_eventbridge_missing_detail_still_checkpoints(self): args, kwargs = self.mock_checkpoint.call_args carrier_get = args[2] self.assertIsNone(carrier_get("dd-pathway-ctx-base64")) - self.assertEqual(kwargs, {"manual_checkpoint": False}) + self.assertEqual(kwargs, {"manual_checkpoint": False, "tags": None}) @patch("datadog_lambda.config.Config.data_streams_enabled", False) def test_eventbridge_data_streams_disabled(self): @@ -4184,3 +4270,135 @@ def test_eventbridge_data_streams_disabled(self): extract_context_from_eventbridge_event(event, self.lambda_context) self.mock_checkpoint.assert_not_called() + + +def _real_set_consume_checkpoint_supports_tags(): + """True when the installed ddtrace `set_consume_checkpoint` exposes the + `tags` parameter (ddtrace >= 4.14).""" + from ddtrace.data_streams import set_consume_checkpoint + + return "tags" in inspect.signature(set_consume_checkpoint).parameters + + +@contextlib.contextmanager +def _stub_data_streams_processor(mock_processor): + """Drive the *real* `ddtrace.data_streams.set_consume_checkpoint` while + stubbing only the DSM processor via a public (non-`.internal`) seam. + + The public accessor differs by ddtrace version: + + * ddtrace >= 4.x exposes the module-level `ddtrace.data_streams. + data_streams_processor()` factory, which `set_consume_checkpoint` + calls directly. + * ddtrace 3.19.x reads `ddtrace.tracer.data_streams_processor` instead + (populated at startup when DSM is enabled). + + Patching whichever seam the running version uses keeps this test honest: + the real `set_consume_checkpoint` executes, so any upstream change to its + signature or behaviour surfaces here instead of being masked by a mock. + """ + import ddtrace + from ddtrace import data_streams + + ds_enabled = ddtrace.config._data_streams_enabled + ddtrace.config._data_streams_enabled = True + try: + if hasattr(data_streams, "data_streams_processor"): + with patch( + "ddtrace.data_streams.data_streams_processor", + return_value=mock_processor, + ): + yield + else: + with patch( + "ddtrace.tracer.data_streams_processor", + mock_processor, + create=True, + ): + yield + finally: + ddtrace.config._data_streams_enabled = ds_enabled + + +class TestDSMRealDdtraceApi(unittest.TestCase): + """Guard tests that exercise the real ddtrace `set_consume_checkpoint` + public API without mocking it, so a signature/behaviour change upstream + fails CI instead of passing silently against a permissive mock. + """ + + def setUp(self): + config_patcher = patch( + "datadog_lambda.config.Config.data_streams_enabled", True + ) + config_patcher.start() + self.addCleanup(config_patcher.stop) + + def _edge_tags(self, mock_processor): + self.assertTrue( + mock_processor.set_checkpoint.called, + "expected real set_consume_checkpoint to reach processor.set_checkpoint", + ) + return mock_processor.set_checkpoint.call_args.args[0] + + def test_real_api_sqs_checkpoint(self): + mock_processor = Mock() + context_json = {"dd-pathway-ctx-base64": "sqs-ctx"} + arn = "arn:aws:sqs:us-east-1:123456789012:test-queue" + + with _stub_data_streams_processor(mock_processor): + _dsm_set_checkpoint(context_json, "sqs", arn) + + edge_tags = self._edge_tags(mock_processor) + self.assertIn("type:sqs", edge_tags) + self.assertIn("topic:" + arn, edge_tags) + self.assertIn("direction:in", edge_tags) + # manual_checkpoint=False must not emit the manual_checkpoint tag. + self.assertNotIn("manual_checkpoint:true", edge_tags) + + def test_real_api_eventbridge_checkpoint(self): + mock_processor = Mock() + context_json = {"dd-pathway-ctx-base64": "eb-ctx"} + + with _stub_data_streams_processor(mock_processor): + _dsm_set_eventbridge_checkpoint(context_json, "MyDetailType") + + edge_tags = self._edge_tags(mock_processor) + self.assertIn("type:eventbridge", edge_tags) + self.assertIn("topic:MyDetailType", edge_tags) + self.assertIn("direction:in", edge_tags) + self.assertNotIn("manual_checkpoint:true", edge_tags) + + @unittest.skipUnless( + _real_set_consume_checkpoint_supports_tags(), + "installed ddtrace has no `tags` parameter (< 4.14)", + ) + @patch("datadog_lambda.config.Config.dsm_exchange_name", "my-event-bus") + def test_real_api_eventbridge_exchange_tag(self): + # Exercises the ddtrace >= 4.14 `tags` functionality end-to-end: the + # exchange edge tag must actually reach the processor. + mock_processor = Mock() + context_json = {"dd-pathway-ctx-base64": "eb-ctx"} + + with _stub_data_streams_processor(mock_processor): + _dsm_set_eventbridge_checkpoint(context_json, "MyDetailType") + + edge_tags = self._edge_tags(mock_processor) + self.assertIn("type:eventbridge", edge_tags) + self.assertIn("topic:MyDetailType", edge_tags) + self.assertIn("exchange:my-event-bus", edge_tags) + + @unittest.skipUnless( + _real_set_consume_checkpoint_supports_tags(), + "installed ddtrace has no `tags` parameter (< 4.14)", + ) + def test_real_api_eventbridge_no_exchange_tag_when_unset(self): + # With no DD_DSM_EXCHANGE_NAME configured, no exchange tag is emitted + # even on ddtrace versions that support `tags`. + mock_processor = Mock() + context_json = {"dd-pathway-ctx-base64": "eb-ctx"} + + with _stub_data_streams_processor(mock_processor): + _dsm_set_eventbridge_checkpoint(context_json, "MyDetailType") + + edge_tags = self._edge_tags(mock_processor) + self.assertFalse(any(t.startswith("exchange:") for t in edge_tags)) From 66fdc28cdf680d1d96730bb10f229b97185f8096 Mon Sep 17 00:00:00 2001 From: James Eastham Date: Wed, 23 Sep 2026 10:26:58 +0100 Subject: [PATCH 08/11] fix: extract SNS dsm context in mixed eventbridge/SNS sqs batches When EventBridge and SNS are both routed to the same SQS queue and their messages are batched into a single event payload, the SNS DSM context was dropped. _dsm_set_eventbridge_sqs_batch_checkpoints classifies every record in the batch and falls back to _extract_sqs_record_message_attribute_context for non-EventBridge records. That helper only read record.messageAttributes and never unwrapped an SNS notification carried in the SQS body (an SNS => SQS subscription without raw message delivery), so the SNS carrier resolved to None. Because the batch contained an EventBridge delivery, dsm_handled was True for the whole batch and the SNS-aware fallback path further down in extract_context_from_sqs_or_sns_event_or_context never ran. Teach the helper to unwrap the SNS envelope the same way the main SQS/SNS extraction path already does, so each record in a mixed batch is classified correctly regardless of whether it is a direct SQS send or an SNS delivery. The shared attribute decoding and envelope detection are factored out into _decode_dd_message_attribute and _parse_sns_notification_from_sqs_body. --- datadog_lambda/tracing.py | 47 ++++++++++++++++++++++++++++++++++-- tests/test_tracing.py | 51 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 96 insertions(+), 2 deletions(-) diff --git a/datadog_lambda/tracing.py b/datadog_lambda/tracing.py index d6b6d7075..4546511dd 100644 --- a/datadog_lambda/tracing.py +++ b/datadog_lambda/tracing.py @@ -526,8 +526,51 @@ def _extract_eventbridge_sqs_record_context(record): def _extract_sqs_record_message_attribute_context(record): - msg_attributes = record.get("messageAttributes") or {} - dd_payload = msg_attributes.get("_datadog") + """Return the ``_datadog`` carrier for a non-EventBridge SQS record. + + A record in an SQS batch may itself be an SNS notification (an SNS => SQS + subscription without raw message delivery), in which case the attributes + live in the SNS envelope inside ``body`` rather than in the record's own + ``messageAttributes``. Both shapes are handled here so that a mixed batch + (EventBridge deliveries alongside SNS deliveries on the same queue) does + not lose the SNS DSM context. + """ + msg_attributes = record.get("messageAttributes") + + if msg_attributes is None: + sns_record = _parse_sns_notification_from_sqs_body(record) or {} + msg_attributes = sns_record.get("MessageAttributes") or {} + + return _decode_dd_message_attribute(msg_attributes.get("_datadog")) + + +def _parse_sns_notification_from_sqs_body(record): + """Return the SNS envelope carried in an SQS record's ``body``, if any. + + Returns ``None`` when the body is not an SNS notification, which simply + means the record is a direct SQS send. + """ + try: + body = json.loads(record.get("body")) + except (ValueError, TypeError): + return None + + if ( + isinstance(body, dict) + and body.get("Type", "") == "Notification" + and "TopicArn" in body + ): + return body + + return None + + +def _decode_dd_message_attribute(dd_payload): + """Decode a ``_datadog`` SQS/SNS message attribute into a carrier dict. + + SQS uses ``dataType`` with ``stringValue``/``binaryValue`` while SNS uses + ``Type`` with ``Value``; both are supported. + """ if not dd_payload: return None diff --git a/tests/test_tracing.py b/tests/test_tracing.py index 17c924a70..74253916d 100644 --- a/tests/test_tracing.py +++ b/tests/test_tracing.py @@ -4103,6 +4103,57 @@ def test_eventbridge_sqs_direct_first_record_still_checkpoints_eventbridge(self) self.assertEqual((second_args[0], second_args[1]), ("sqs", arn2)) self.assertEqual(second_args[2]("dd-pathway-ctx-base64"), "eb-ctx") + def test_eventbridge_sqs_mixed_batch_with_sns_record_keeps_sns_context(self): + # EventBridge and SNS can both be routed to the same SQS queue and then + # batched into a single event payload. The SNS record's DSM context + # lives in the SNS envelope inside `body` (SNS => SQS subscription + # without raw message delivery), not in `messageAttributes`, and must + # still be extracted even though the batch contains an EventBridge + # delivery. + arn1 = "arn:aws:sqs:us-east-1:123456789012:eb-queue" + arn2 = "arn:aws:sqs:us-east-1:123456789012:sns-queue" + sns_dd_data = { + "x-datadog-trace-id": "12345", + "x-datadog-parent-id": "67890", + "x-datadog-sampling-priority": "1", + "dd-pathway-ctx-base64": "sns-ctx", + } + sns_body = { + "Type": "Notification", + "TopicArn": "arn:aws:sns:us-east-1:123456789012:my-topic", + "Message": "hello from sns", + "MessageAttributes": { + "_datadog": { + "Type": "String", + "Value": json.dumps(sns_dd_data), + } + }, + } + event = { + "Records": [ + self._eventbridge_sqs_record(arn1, "eb-ctx"), + { + "eventSourceARN": arn2, + "eventSource": "aws:sqs", + "body": json.dumps(sns_body), + }, + ] + } + + extract_context_from_sqs_or_sns_event_or_context( + event, self.lambda_context, parse_event_source(event) + ) + + self.assertEqual(self.mock_checkpoint.call_count, 2) + first_args, _ = self.mock_checkpoint.call_args_list[0] + second_args, _ = self.mock_checkpoint.call_args_list[1] + self.assertEqual((first_args[0], first_args[1]), ("sqs", arn1)) + self.assertEqual(first_args[2]("dd-pathway-ctx-base64"), "eb-ctx") + # The SNS record is consumed from the SQS queue, so it keeps SQS + # conventions, but its carrier must come from the SNS envelope. + self.assertEqual((second_args[0], second_args[1]), ("sqs", arn2)) + self.assertEqual(second_args[2]("dd-pathway-ctx-base64"), "sns-ctx") + def test_eventbridge_sqs_non_json_body_falls_back_to_sqs_attributes(self): # A later direct SQS record can legitimately carry a non-JSON body # while its DSM carrier lives in messageAttributes._datadog. The From d54041de8b910b2ded78905ee9f491a8e2a0c163 Mon Sep 17 00:00:00 2001 From: James Eastham Date: Wed, 23 Sep 2026 10:39:55 +0100 Subject: [PATCH 09/11] chore: bump lambda layer size limit for ddtrace 4.15.x The zipped layer size check started failing at 9273 kb against the 9231 kb limit. This is unrelated to any source change in this branch: the ddtrace cp311 manylinux x86_64 wheel grew ~470 kb in 4.15.0 (8.87 MB in 4.14.2 -> 9.33 MB in 4.15.0), and pyproject.toml pins ddtrace ">=4.1.1,<5,!=4.6.*" with no upper bound inside 4.x, so CI resolves to the newest 4.15.x. Raise MAX_LAYER_COMPRESSED_SIZE_KB from 9*1024+15 (9231 KB) to 9*1024+128 (9344 KB), which clears the observed 9273 kb with ~71 kb of headroom while keeping the guardrail tight enough to catch further growth. x86_64 is the larger of the two published arches (~245 kb above aarch64), so the failing cp311 x86_64 job is the worst case across the build matrix. The uncompressed limit is unchanged. --- scripts/check_layer_size.sh | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/scripts/check_layer_size.sh b/scripts/check_layer_size.sh index de9a8c4d4..a4b6677b9 100755 --- a/scripts/check_layer_size.sh +++ b/scripts/check_layer_size.sh @@ -8,7 +8,7 @@ # Compares layer size to threshold, and fails if below that threshold set -e -MAX_LAYER_COMPRESSED_SIZE_KB=$(expr 9 \* 1024 + 15) # 9231 KB +MAX_LAYER_COMPRESSED_SIZE_KB=$(expr 9 \* 1024 + 128) # 9344 KB MAX_LAYER_UNCOMPRESSED_SIZE_KB=$(expr 25 \* 1024) # 25600 KB From fa716719af8ae7e4bd5801a5ef7a575c34f3ab73 Mon Sep 17 00:00:00 2001 From: James Eastham Date: Fri, 25 Sep 2026 15:36:01 +0100 Subject: [PATCH 10/11] chore: update README with new flag --- README.md | 1 + 1 file changed, 1 insertion(+) diff --git a/README.md b/README.md index a1199e8eb..fc6c7f2c1 100644 --- a/README.md +++ b/README.md @@ -30,6 +30,7 @@ Besides the environment variables supported by dd-trace-py, the datadog-lambda-p | DD_CAPTURE_LAMBDA_PAYLOAD | [Captures incoming and outgoing AWS Lambda payloads][1] in the Datadog APM spans for Lambda invocations. | `false` | | DD_CAPTURE_LAMBDA_PAYLOAD_MAX_DEPTH | Determines the level of detail captured from AWS Lambda payloads, which are then assigned as tags for the `aws.lambda` span. It specifies the nesting depth of the JSON payload structure to process. Once the specified maximum depth is reached, the tag's value is set to the stringified value of any nested elements beyond this level.
For example, given the input payload:
{
"lv1" : {
"lv2": {
"lv3": "val"
}
}
}
If the depth is set to `2`, the resulting tag's key is set to `function.request.lv1.lv2` and the value is `{\"lv3\": \"val\"}`.
If the depth is set to `0`, the resulting tag's key is set to `function.request` and value is `{\"lv1\":{\"lv2\":{\"lv3\": \"val\"}}}` | `10` | | DD_EXCEPTION_REPLAY_ENABLED | When set to `true`, the Lambda will run with Error Tracking Exception Replay enabled, capturing local variables. | `false` | +| DD_DSM_EXCHANGE_NAME | If you are using Datadog Data Streams Monitoring and propagating context through Amazon Event Bridge, this variable specifies the name of the EventBridge EventBus. The name of the bus is *not* passed through in the event payload and therefore cannot be automatically inferred at consume time. | `false` | ## Opening Issues From 98c9ac1b285f8ce8437af318890fffe9a658d149 Mon Sep 17 00:00:00 2001 From: James Eastham Date: Fri, 25 Sep 2026 15:37:45 +0100 Subject: [PATCH 11/11] fix: handle msg_attributes being empty object --- datadog_lambda/tracing.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/datadog_lambda/tracing.py b/datadog_lambda/tracing.py index 4546511dd..741686494 100644 --- a/datadog_lambda/tracing.py +++ b/datadog_lambda/tracing.py @@ -537,7 +537,7 @@ def _extract_sqs_record_message_attribute_context(record): """ msg_attributes = record.get("messageAttributes") - if msg_attributes is None: + if not msg_attributes: sns_record = _parse_sns_notification_from_sqs_body(record) or {} msg_attributes = sns_record.get("MessageAttributes") or {}