Repository navigation
feat: extract trace context from MSK events - #830
lucassarcanjo wants to merge 3 commits into
Conversation
BridgeAR
left a comment
There was a problem hiding this comment.
Thank you for the PR! I just left a few suggestions to make the code a bit faster :)
34ffa90 to
c552877
Compare
BridgeAR
left a comment
There was a problem hiding this comment.
Thank you for the quick follow-ups!
Code wise it seems fine to me.
|
Hi @lym953, could you take a look at this PR? |
|
@codex review |
|
To use Codex here, create a Codex account and connect to github. |
|
@codex review |
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
|
Codex Review: Didn't find any major issues. You're on a roll. Reviewed commit: ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
If Codex has suggestions, it will comment; otherwise it will react with 👍. Codex can also answer questions or update the PR. Try commenting "@codex address that feedback". |
|
LGTM 👍 |
|
Thanks @lucassarcanjo! This looks good to us but we need you to rebase onto the original main branch instead of a fork. If you'd rather us take care of it, please enable allow edits from maintainers so we can rebase and merge |
c552877 to
b6dd184
Compare
|
Hi @zarirhamza, done! |
- Guard event.records with an early return instead of an empty fallback - Wrap the record loop in a single try/catch so errors log once - Return null from getParsedRecordHeaders when nothing decodes, lazily creating the headers map - Drop per-byte validation and rely on Buffer.from
b6dd184 to
b98fe61
Compare
|
/remove |
|
View all feedbacks in Devflow UI.
|
| event.records !== null && | ||
| typeof event.records === "object" && | ||
| !Array.isArray(event.records) | ||
| ); |
There was a problem hiding this comment.
Lambda's self-managed Kafka trigger (any non-MSK cluster, e.g. on EC2 or Confluent Cloud) sends the same records and headers shape as MSK, but with eventSource: "SelfManagedKafka". Right now isMSKEvent only accepts "aws:kafka"`, so those Lambdas would still start a disconnected trace.
Could we accept both values?
static isKafkaEvent(event: any): event is MSKEvent | SelfManagedKafkaEvent {
return (
(event?.eventSource === "aws:kafka" || event?.eventSource === "SelfManagedKafka") &&
event.records !== null &&
typeof event.records === "object" &&
!Array.isArray(event.records)
);
}The header-decoding logic should work for both as-is. A test fixture with eventSource: "SelfManagedKafka" would cover it. If self-managed Kafka is intentionally out of scope for this PR, that's fine too. A follow-up issue would be enough.
| try { | ||
| // A Lambda span can have only one parent. Use the first record with valid | ||
| // trace context, without combining headers from different records. | ||
| for (const records of Object.values(event.records)) { |
There was a problem hiding this comment.
Every other batch extractor takes trace context from the first record only: sqs.ts, sns.ts, sns-sqs.ts, event-bridge-sqs.ts, and kinesis.ts all use Records[0]. Their loops over every record are only for DSM checkpoints. Can we do the same here? Read the first record of the first topic-partition, try to extract from it, and return.
This keeps MSK consistent with the other event sources, so the Lambda span's parent is always the first record. It also avoids decoding headers and calling tracer.extract on up to 10k records in every untraced batch.
The test cases that expect a later record's context to win would need to change to expect null.
purple4reina
left a comment
There was a problem hiding this comment.
Hi, just the two comments. Then I'll work to make sure the tests pass and we can get this merged. 😄
|
Sorry, there's too many datadog cooks in the kitchen, looks like we're all giving you different instructions. My apologies for any confusion. |
|
Oh wow, my bad. Looks like this was all already merged in #838. I'm gonna go ahead and close this PR. Thank you again for your help! |
What does this PR do?
Adds automatic trace context extraction for MSK-triggered Lambdas. Kafka header byte arrays are decoded as UTF-8 and passed to the existing tracer, preserving Datadog and W3C propagation context.
For batches, the Lambda span uses the first record with valid trace context across the topic-partition groups. Malformed headers are skipped, and headers from different records are never combined.
Motivation
Fixes #829. An instrumented producer's trace headers reach the Lambda event, but the existing dispatcher does not extract them, so the Lambda starts a separate trace.
Testing Guidelines
yarn test --runInBand: 691 tests and 3 snapshots passed on Node.js 22.yarn lintand formatting checks passed.dd-trace@5.118.0: the Lambda tracing lifecycle preserves the producer parent ID, sampling priority, and 128-bit trace ID with Datadog-only, W3C-only, and combined headers. The same reproduction fails with the original dispatcher.Types of Changes
Check all that apply