Skip to content

Commit a0fe575

Browse files
authored
feat: add support for event bridge DSM context extraction (#836)
* feat: add support for event bridge DSM context extraction * chore: update processing of eventbridge DSM * refactor: avoid extra allocation in eventbridge sqs parsing * fix: require eventbridge envelope for sqs extraction * fix: detect sqs carriers per record in eventbridge batches * fix: use public dsm checkpoint for eventbridge * chore: update to handle new public DSM API * 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. * 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. * chore: update README with new flag * fix: handle msg_attributes being empty object
1 parent d28d13e commit a0fe575

4 files changed

Lines changed: 780 additions & 16 deletions

File tree

‎README.md‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ Besides the environment variables supported by dd-trace-py, the datadog-lambda-p
3030
| DD_CAPTURE_LAMBDA_PAYLOAD | [Captures incoming and outgoing AWS Lambda payloads][1] in the Datadog APM spans for Lambda invocations. | `false` |
3131
| 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. <br> For example, given the input payload: <pre>{<br> "lv1" : {<br> "lv2": {<br> "lv3": "val"<br> }<br> }<br>}</pre> 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\"}`. <br> 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` |
3232
| DD_EXCEPTION_REPLAY_ENABLED | When set to `true`, the Lambda will run with Error Tracking Exception Replay enabled, capturing local variables. | `false` |
33+
| 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` |
3334

3435

3536
## Opening Issues

‎datadog_lambda/config.py‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -96,6 +96,10 @@ def _resolve_env(self, key, default=None, cast=None, depends_on_tracing=False):
9696
data_streams_enabled = _get_env(
9797
"DD_DATA_STREAMS_ENABLED", "false", as_bool, depends_on_tracing=True
9898
)
99+
# EventBridge bus name used as the DSM `exchange` tag. The bus name is not
100+
# present in the inbound event, so it must be provided explicitly to allow
101+
# the consume checkpoint to pair with the EventBridge produce checkpoint.
102+
dsm_exchange_name = _get_env("DD_DSM_EXCHANGE_NAME")
99103
appsec_enabled = _get_env("DD_APPSEC_ENABLED", "false", as_bool)
100104
sca_enabled = _get_env("DD_APPSEC_SCA_ENABLED", "false", as_bool)
101105

‎datadog_lambda/tracing.py‎

Lines changed: 240 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -83,6 +83,47 @@ def _dsm_set_checkpoint(context_json, event_type, arn):
8383
)
8484

8585

86+
def _dsm_set_eventbridge_checkpoint(context_json, detail_type):
87+
"""Set a DSM consume checkpoint for an EventBridge event.
88+
89+
When ddtrace >= 4.14 is present the public ``tags`` parameter is used to
90+
attach the ``exchange:`` edge tag (sourced from ``DD_DSM_EXCHANGE_NAME``).
91+
On older installs the checkpoint is still emitted, just without the
92+
exchange tag.
93+
"""
94+
if not config.data_streams_enabled or not detail_type:
95+
return
96+
97+
try:
98+
from ddtrace.data_streams import set_consume_checkpoint
99+
100+
carrier_get = lambda k: context_json and context_json.get(k) # noqa: E731
101+
try:
102+
tags = (
103+
["exchange:" + config.dsm_exchange_name]
104+
if config.dsm_exchange_name
105+
else None
106+
)
107+
set_consume_checkpoint(
108+
"eventbridge",
109+
detail_type,
110+
carrier_get,
111+
manual_checkpoint=False,
112+
tags=tags,
113+
)
114+
except TypeError:
115+
# ddtrace < 4.14 has no `tags` parameter. Retry without it, but
116+
# keep `manual_checkpoint=False` so the checkpoint stays
117+
# consistent with every other consume checkpoint.
118+
set_consume_checkpoint(
119+
"eventbridge", detail_type, carrier_get, manual_checkpoint=False
120+
)
121+
except Exception as e:
122+
logger.debug(
123+
f"DSM:Failed to set consume checkpoint for eventbridge {detail_type}: {e}"
124+
)
125+
126+
86127
def _convert_xray_trace_id(xray_trace_id):
87128
"""
88129
Convert X-Ray trace id (hex)'s last 63 bits to a Datadog trace id (int).
@@ -251,11 +292,21 @@ def extract_context_from_sqs_or_sns_event_or_context(
251292
source_arn = ""
252293
event_type = "sqs" if event_source.equals(EventTypes.SQS) else "sns"
253294

254-
# EventBridge => SQS
295+
# EventBridge => SQS. `dsm_handled` is True when the batch contained at
296+
# least one EventBridge delivery and DSM checkpoints were already set
297+
# per-record; in that case the regular SQS path below must not set another
298+
# checkpoint for the first record or it would be double counted.
299+
dsm_handled = False
255300
try:
256-
context = _extract_context_from_eventbridge_sqs_event(event)
257-
if _is_context_complete(context):
258-
return context
301+
(
302+
context,
303+
is_eventbridge_sqs,
304+
dsm_handled,
305+
) = _extract_context_from_eventbridge_sqs_event(event)
306+
if is_eventbridge_sqs:
307+
if _is_context_complete(context):
308+
return context
309+
return extract_context_from_lambda_context(lambda_context)
259310
except Exception:
260311
logger.debug("Failed extracting context as EventBridge to SQS.")
261312

@@ -312,7 +363,8 @@ def extract_context_from_sqs_or_sns_event_or_context(
312363
"Failed to extract Step Functions context from SQS/SNS event."
313364
)
314365
context = propagator.extract(dd_data)
315-
_dsm_set_checkpoint(dd_data, event_type, source_arn)
366+
if not dsm_handled:
367+
_dsm_set_checkpoint(dd_data, event_type, source_arn)
316368
return context
317369
else:
318370
# Handle case where trace context is injected into attributes.AWSTraceHeader
@@ -337,12 +389,14 @@ def extract_context_from_sqs_or_sns_event_or_context(
337389
sampling_priority=float(x_ray_context["sampled"]),
338390
)
339391
# Still want to set a DSM checkpoint even if DSM context not propagated
340-
_dsm_set_checkpoint(None, event_type, source_arn)
392+
if not dsm_handled:
393+
_dsm_set_checkpoint(None, event_type, source_arn)
341394
return extract_context_from_lambda_context(lambda_context)
342395
except Exception as e:
343396
logger.debug("The trace extractor returned with error %s", e)
344397
# Still want to set a DSM checkpoint even if DSM context not propagated
345-
_dsm_set_checkpoint(None, event_type, source_arn)
398+
if not dsm_handled:
399+
_dsm_set_checkpoint(None, event_type, source_arn)
346400
return extract_context_from_lambda_context(lambda_context)
347401

348402

@@ -353,22 +407,190 @@ def _extract_context_from_eventbridge_sqs_event(event):
353407
354408
This is only possible if first record in `Records` contains a
355409
`body` field which contains the EventBridge `detail` as a JSON string.
410+
411+
Returns a tuple ``(context, is_eventbridge_sqs, dsm_handled)``:
412+
413+
* ``context`` / ``is_eventbridge_sqs`` describe the trace context and are
414+
derived only from the first record, since that is the record whose trace
415+
context becomes the Lambda's parent.
416+
* ``dsm_handled`` reports whether this function already set DSM checkpoints
417+
for the batch. It is ``True`` whenever *any* record in the batch is an
418+
EventBridge delivery, so the caller must not set its own SQS checkpoint
419+
(which would double count the first record).
356420
"""
357-
first_record = event.get("Records")[0]
358-
body_str = first_record.get("body")
359-
body = json.loads(body_str)
360-
detail = body.get("detail")
361-
dd_context = detail.get("_datadog")
421+
records = event.get("Records")
422+
if not records:
423+
return None, False, False
424+
425+
first_record = records[0]
426+
dd_context, is_eventbridge_sqs = _extract_eventbridge_sqs_record_context(
427+
first_record
428+
)
429+
430+
# Set a consume checkpoint for every record in the batch whenever the batch
431+
# contains at least one EventBridge delivery. Each record is classified
432+
# independently so a mixed batch (EventBridge deliveries alongside direct
433+
# SQS sends, in either order) uses the correct carrier per record. The
434+
# message is consumed from the SQS queue, so it follows SQS conventions
435+
# (type:sqs, topic:queue ARN).
436+
dsm_handled = _dsm_set_eventbridge_sqs_batch_checkpoints(records)
437+
438+
if not is_eventbridge_sqs:
439+
return None, False, dsm_handled
362440

363441
if is_step_function_event(dd_context):
364442
try:
365-
return extract_context_from_step_functions(dd_context, None)
443+
return (
444+
extract_context_from_step_functions(dd_context, None),
445+
True,
446+
dsm_handled,
447+
)
366448
except Exception:
367449
logger.debug(
368450
"Failed to extract Step Functions context from EventBridge to SQS event."
369451
)
370452

371-
return propagator.extract(dd_context)
453+
return propagator.extract(dd_context), True, dsm_handled
454+
455+
456+
def _dsm_set_eventbridge_sqs_batch_checkpoints(records):
457+
"""Set a per-record SQS DSM consume checkpoint for an EventBridge -> SQS
458+
batch.
459+
460+
Returns ``True`` when the batch contains at least one EventBridge delivery
461+
(and checkpoints were therefore this function's responsibility), otherwise
462+
``False`` so the caller can fall back to its regular SQS checkpoint path.
463+
Each record is classified independently: EventBridge records use the
464+
carrier embedded in ``body.detail._datadog`` while other records fall back
465+
to the SQS ``messageAttributes._datadog`` carrier.
466+
"""
467+
if not config.data_streams_enabled:
468+
return False
469+
470+
record_carriers = []
471+
batch_has_eventbridge = False
472+
for record in records:
473+
record_context, is_eventbridge_record = _extract_eventbridge_sqs_record_context(
474+
record
475+
)
476+
if is_eventbridge_record:
477+
batch_has_eventbridge = True
478+
else:
479+
try:
480+
record_context = _extract_sqs_record_message_attribute_context(record)
481+
except Exception:
482+
record_context = None
483+
record_carriers.append(record_context)
484+
485+
if not batch_has_eventbridge:
486+
return False
487+
488+
for record, record_context in zip(records, record_carriers):
489+
try:
490+
_dsm_set_checkpoint(record_context, "sqs", record.get("eventSourceARN", ""))
491+
except Exception:
492+
logger.debug(
493+
"Failed to set DSM checkpoint for an EventBridge to SQS record."
494+
)
495+
496+
return True
497+
498+
499+
def _extract_eventbridge_sqs_record_context(record):
500+
"""Classify a single SQS record as an EventBridge delivery and return its
501+
DSM carrier.
502+
503+
Returns a tuple ``(dd_context, is_eventbridge)``. A record is only treated
504+
as EventBridge when its ``body`` is a JSON object carrying the EventBridge
505+
envelope fields (``detail`` object plus ``detail-type`` and ``source``). A
506+
non-JSON or non-envelope body is not an error here: it simply means the
507+
record is a regular SQS message, so the caller can fall back to the SQS
508+
message attribute carrier instead.
509+
"""
510+
body_str = record.get("body")
511+
try:
512+
body = json.loads(body_str)
513+
except (ValueError, TypeError):
514+
return None, False
515+
516+
if not isinstance(body, dict):
517+
return None, False
518+
519+
detail = body.get("detail")
520+
if not (
521+
isinstance(detail, dict) and body.get("detail-type") and body.get("source")
522+
):
523+
return None, False
524+
525+
return detail.get("_datadog"), True
526+
527+
528+
def _extract_sqs_record_message_attribute_context(record):
529+
"""Return the ``_datadog`` carrier for a non-EventBridge SQS record.
530+
531+
A record in an SQS batch may itself be an SNS notification (an SNS => SQS
532+
subscription without raw message delivery), in which case the attributes
533+
live in the SNS envelope inside ``body`` rather than in the record's own
534+
``messageAttributes``. Both shapes are handled here so that a mixed batch
535+
(EventBridge deliveries alongside SNS deliveries on the same queue) does
536+
not lose the SNS DSM context.
537+
"""
538+
msg_attributes = record.get("messageAttributes")
539+
540+
if not msg_attributes:
541+
sns_record = _parse_sns_notification_from_sqs_body(record) or {}
542+
msg_attributes = sns_record.get("MessageAttributes") or {}
543+
544+
return _decode_dd_message_attribute(msg_attributes.get("_datadog"))
545+
546+
547+
def _parse_sns_notification_from_sqs_body(record):
548+
"""Return the SNS envelope carried in an SQS record's ``body``, if any.
549+
550+
Returns ``None`` when the body is not an SNS notification, which simply
551+
means the record is a direct SQS send.
552+
"""
553+
try:
554+
body = json.loads(record.get("body"))
555+
except (ValueError, TypeError):
556+
return None
557+
558+
if (
559+
isinstance(body, dict)
560+
and body.get("Type", "") == "Notification"
561+
and "TopicArn" in body
562+
):
563+
return body
564+
565+
return None
566+
567+
568+
def _decode_dd_message_attribute(dd_payload):
569+
"""Decode a ``_datadog`` SQS/SNS message attribute into a carrier dict.
570+
571+
SQS uses ``dataType`` with ``stringValue``/``binaryValue`` while SNS uses
572+
``Type`` with ``Value``; both are supported.
573+
"""
574+
if not dd_payload:
575+
return None
576+
577+
dd_json_data = None
578+
dd_json_data_type = dd_payload.get("Type") or dd_payload.get("dataType")
579+
if dd_json_data_type == "Binary":
580+
import base64
581+
582+
dd_json_data = dd_payload.get("binaryValue") or dd_payload.get("Value")
583+
if dd_json_data:
584+
dd_json_data = base64.b64decode(dd_json_data)
585+
elif dd_json_data_type == "String":
586+
dd_json_data = dd_payload.get("stringValue") or dd_payload.get("Value")
587+
else:
588+
logger.debug(
589+
"Datadog Lambda Python only supports extracting trace"
590+
"context from String or Binary SQS/SNS message attributes"
591+
)
592+
593+
return json.loads(dd_json_data) if dd_json_data else None
372594

373595

374596
def extract_context_from_eventbridge_event(event, lambda_context):
@@ -380,8 +602,11 @@ def extract_context_from_eventbridge_event(event, lambda_context):
380602
that header.
381603
"""
382604
try:
383-
detail = event.get("detail")
605+
detail = event.get("detail") or {}
384606
dd_context = detail.get("_datadog")
607+
608+
_dsm_set_eventbridge_checkpoint(dd_context, event.get("detail-type"))
609+
385610
if not dd_context:
386611
return extract_context_from_lambda_context(lambda_context)
387612

0 commit comments

Comments
 (0)