Repository navigation
feat: add support for event bridge DSM context extraction #836
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from 3 commits
888ab65
218a15b
8fc0cd0
9dcc6ec
a45866f
ff42108
24da389
f9ba2ed
fd813a3
66fdc28
d54041d
fa71671
98c9ac1
0af9e16
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -82,6 +82,44 @@ 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.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}" | ||
| ) | ||
|
|
||
|
|
||
| def _convert_xray_trace_id(xray_trace_id): | ||
| """ | ||
| Convert X-Ray trace id (hex)'s last 63 bits to a Datadog trace id (int). | ||
|
|
@@ -252,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: | ||
|
jeastham1993 marked this conversation as resolved.
|
||
| logger.debug("Failed extracting context as EventBridge to SQS.") | ||
|
|
||
|
|
@@ -353,21 +395,52 @@ 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") or [] | ||
|
jeastham1993 marked this conversation as resolved.
Outdated
|
||
| if not records: | ||
| return None, False | ||
|
|
||
| first_record = records[0] | ||
| body_str = first_record.get("body") | ||
| body = json.loads(body_str) | ||
|
jeastham1993 marked this conversation as resolved.
Outdated
|
||
| detail = body.get("detail") | ||
| 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_detail = record_body.get("detail") | ||
| record_context = ( | ||
| record_detail.get("_datadog") | ||
| if isinstance(record_detail, dict) | ||
|
jeastham1993 marked this conversation as resolved.
Outdated
|
||
| else None | ||
| ) | ||
| _dsm_set_checkpoint( | ||
| record_context, "sqs", record.get("eventSourceARN", "") | ||
| ) | ||
|
jeastham1993 marked this conversation as resolved.
Outdated
|
||
| 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) | ||
| 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): | ||
|
|
@@ -379,8 +452,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")) | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
For standard scheduled EventBridge events ( Useful? React with 👍 / 👎. |
||
|
|
||
| if not dd_context: | ||
| return extract_context_from_lambda_context(lambda_context) | ||
|
|
||
|
|
||
Uh oh!
There was an error while loading. Please reload this page.