Skip to content

fix(google-cloud-pubsub): fix DSM context propagation for google-cloud-pubsub#7582

Open
robcarlan-datadog wants to merge 3 commits into
masterfrom
rob.carlan/dsm-context-leak-pubsub-test
Open

fix(google-cloud-pubsub): fix DSM context propagation for google-cloud-pubsub#7582
robcarlan-datadog wants to merge 3 commits into
masterfrom
rob.carlan/dsm-context-leak-pubsub-test

Conversation

@robcarlan-datadog

@robcarlan-datadog robcarlan-datadog commented Feb 19, 2026

Copy link
Copy Markdown
Contributor

Summary

  • Fixes flaky DSM behavior in google-cloud-pubsub by moving DSM checkpoint logic from bindStart to start in both the consumer and producer plugins, so DSM context is properly bound in the async store
  • Adds a regression test verifying that interleaved consume-produce flows maintain separate DSM contexts

Test plan

  • Verify regression test passes locally and in CI. Fails before adding fix

🤖 Generated with Claude Code

robcarlan-datadog added a commit that referenced this pull request Feb 19, 2026
Moved to #7582 to fix independently.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
@github-actions

github-actions Bot commented Feb 19, 2026

Copy link
Copy Markdown
Contributor

Overall package size

Self size: 6.77 MB
Deduped: 7.43 MB
No deduping: 7.43 MB

Dependency sizes | name | version | self size | total size | |------|---------|-----------|------------| | import-in-the-middle | 3.3.1 | 122.62 kB | 438.86 kB | | opentracing | 0.14.7 | 194.81 kB | 194.81 kB | | dc-polyfill | 0.1.11 | 25.74 kB | 25.74 kB |

🤖 This report was automatically generated by heaviest-objects-in-the-universe

@datadog-datadog-prod-us1

datadog-datadog-prod-us1 Bot commented Feb 19, 2026

Copy link
Copy Markdown

Tests

🎉 All green!

🧪 All tests passed
❄️ No new flaky tests detected

🎯 Code Coverage (details)
Patch Coverage: 100.00%
Overall Coverage: 98.34% (+0.01%)

This comment will be updated automatically if new data arrives.
🔗 Commit SHA: 52b7b61 | Docs | Datadog PR Page | Give us feedback!

robcarlan-datadog added a commit that referenced this pull request Feb 19, 2026
Moved DSM context fix and regression test for google-cloud-pubsub
to #7582 to fix independently.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
@pr-commenter

pr-commenter Bot commented Feb 19, 2026

Copy link
Copy Markdown

Benchmarks

Benchmark execution time: 2026-07-17 16:53:13

Comparing candidate commit 52b7b61 in PR branch rob.carlan/dsm-context-leak-pubsub-test with baseline commit 57d5f96 in branch master.

📊 Benchmarking dashboard

Found 0 performance improvements and 0 performance regressions! Performance is the same for 2316 metrics, 42 unstable metrics.

Explanation

This is an A/B test comparing a candidate commit's performance against that of a baseline commit. Performance changes are noted in the tables below as:

  • 🟩 = significantly better candidate vs. baseline
  • 🟥 = significantly worse candidate vs. baseline

We compute a confidence interval (CI) over the relative difference of means between metrics from the candidate and baseline commits, considering the baseline as the reference.

If the CI is entirely outside the configured SIGNIFICANT_IMPACT_THRESHOLD (or the deprecated UNCONFIDENCE_THRESHOLD), the change is considered significant.

Feel free to reach out to #apm-benchmarking-platform on Slack if you have any questions.

More details about the CI and significant changes

You can imagine this CI as a range of values that is likely to contain the true difference of means between the candidate and baseline commits.

CIs of the difference of means are often centered around 0%, because often changes are not that big:

---------------------------------(------|---^--------)-------------------------------->
                              -0.6%    0%  0.3%     +1.2%
                                 |          |        |
         lower bound of the CI --'          |        |
sample mean (center of the CI) -------------'        |
         upper bound of the CI ----------------------'

As described above, a change is considered significant if the CI is entirely outside the configured SIGNIFICANT_IMPACT_THRESHOLD (or the deprecated UNCONFIDENCE_THRESHOLD).

For instance, for an execution time metric, this confidence interval indicates a significantly worse performance:

----------------------------------------|---------|---(---------^---------)---------->
                                       0%        1%  1.3%      2.2%      3.1%
                                                  |   |         |         |
       significant impact threshold --------------'   |         |         |
                      lower bound of CI --------------'         |         |
       sample mean (center of the CI) --------------------------'         |
                      upper bound of CI ----------------------------------'

Unstable benchmarks

These benchmarks have a confidence interval too wide to call a change; treat them as noise rather than signal.

scenario:appsec-appsec-enabled-24

  • unstable execution_time [-210.592ms; +200.151ms] or [-7.946%; +7.553%]

scenario:appsec-appsec-enabled-26

  • unstable execution_time [-233.516ms; +237.793ms] or [-9.125%; +9.292%]

scenario:appsec-appsec-enabled-with-attacks-24

  • unstable execution_time [-171.041ms; +161.857ms] or [-5.537%; +5.240%]

scenario:appsec-appsec-enabled-with-attacks-26

  • unstable execution_time [-186.177ms; +190.156ms] or [-6.405%; +6.542%]

scenario:appsec-control-20

  • unstable execution_time [-111.650ms; +139.595ms] or [-6.766%; +8.460%]

scenario:appsec-control-24

  • unstable execution_time [-112233.162µs; +114006.429µs] or [-9.067%; +9.210%]

scenario:appsec-control-26

  • unstable execution_time [-124.725ms; +127.891ms] or [-10.056%; +10.311%]

scenario:appsec-iast-no-vulnerability-control-20

  • unstable execution_time [-18.257ms; +12.674ms] or [-6.922%; +4.805%]

scenario:appsec-iast-no-vulnerability-iast-enabled-always-active-20

  • unstable execution_time [-14.525ms; +18.098ms] or [-5.691%; +7.091%]

scenario:appsec-iast-with-vulnerability-iast-enabled-always-active-20

  • unstable execution_time [-28.449ms; +30.588ms] or [-5.171%; +5.560%]

scenario:appsec-iast-with-vulnerability-iast-enabled-default-config-20

  • unstable execution_time [-24.284ms; +34.225ms] or [-4.421%; +6.231%]

scenario:debugger-line-probe-with-snapshot-default-26

  • unstable cpu_user_time [-3513.610ms; +5099.658ms] or [-31.701%; +46.011%]
  • unstable execution_time [-3514.313ms; +5140.640ms] or [-29.770%; +43.546%]
  • unstable instructions [-31.9G instructions; +45.7G instructions] or [-34.205%; +48.916%]
  • unstable max_rss_usage [-11.134MB; +16.589MB] or [-6.780%; +10.102%]
  • unstable throughput [-1018.348op/s; +672.506op/s] or [-34.750%; +22.949%]

scenario:debugger-line-probe-with-snapshot-minimal-26

  • unstable cpu_user_time [-2632.336ms; +4217.706ms] or [-27.541%; +44.129%]
  • unstable execution_time [-2865.466ms; +4504.022ms] or [-27.806%; +43.706%]
  • unstable instructions [-23.2G instructions; +36.9G instructions] or [-29.086%; +46.407%]
  • unstable max_rss_usage [-6.993MB; +13.768MB] or [-4.382%; +8.628%]
  • unstable throughput [-853.184op/s; +547.992op/s] or [-26.412%; +16.964%]

scenario:debugger-line-probe-without-snapshot-24

  • unstable cpu_user_time [-1762.454ms; +563.368ms] or [-21.343%; +6.822%]
  • unstable execution_time [-1770.362ms; +576.217ms] or [-19.750%; +6.428%]
  • unstable instructions [-15.0G instructions; +4.8G instructions] or [-22.225%; +7.166%]
  • unstable throughput [-156.405op/s; +475.762op/s] or [-4.255%; +12.943%]

scenario:debugger-line-probe-without-snapshot-26

  • unstable cpu_user_time [-2269.880ms; +727.799ms] or [-23.857%; +7.650%]
  • unstable execution_time [-2299.662ms; +715.741ms] or [-22.410%; +6.975%]
  • unstable instructions [-20.2G instructions; +6.6G instructions] or [-25.434%; +8.268%]
  • unstable throughput [-141.416op/s; +460.752op/s] or [-4.381%; +14.273%]

scenario:dogstatsd-aggregated-24

  • unstable execution_time [-43.317ms; +72.971ms] or [-3.954%; +6.662%]

scenario:dogstatsd-with-tags-20

  • unstable cpu_user_time [-399.348ms; +250.260ms] or [-8.462%; +5.303%]
  • unstable execution_time [-395.531ms; +249.665ms] or [-8.238%; +5.200%]
  • unstable throughput [-88187.898op/s; +138782.333op/s] or [-5.046%; +7.941%]

scenario:plugin-aws-sdk-lambda-inject-with-context-24

  • unstable cpu_user_time [-180.714ms; +235.759ms] or [-4.748%; +6.194%]
  • unstable execution_time [-178.721ms; +236.326ms] or [-4.661%; +6.163%]

scenario:plugin-claude-agent-sdk-compact-stream-scan-24

  • unstable cpu_usage_percentage [-5.292%; +5.114%]

scenario:plugin-claude-agent-sdk-compact-stream-scan-26

  • unstable cpu_usage_percentage [-6.877%; +3.234%]

scenario:plugin-graphql-long-with-depth-and-collapse-off-20

  • unstable max_rss_usage [-35.559MB; +26.076MB] or [-9.193%; +6.741%]

scenario:plugin-graphql-long-with-depth-on-max-20

  • unstable cpu_user_time [-596.977ms; +605.541ms] or [-5.187%; +5.261%]
  • unstable execution_time [-607.107ms; +621.054ms] or [-5.170%; +5.289%]
  • unstable throughput [-3.626op/s; +3.553op/s] or [-5.291%; +5.185%]

scenario:test-optimization-large-suite-20

  • unstable max_rss_usage [-4259.668KB; +6080.668KB] or [-5.371%; +7.667%]

@robcarlan-datadog robcarlan-datadog changed the title test(dsm): add regression test for DSM context leak in google-cloud-pubsub fix(dsm): fix DSM context propagation for google-cloud-pubsub Feb 19, 2026
@codecov

codecov Bot commented Feb 19, 2026

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 98.34%. Comparing base (57d5f96) to head (52b7b61).

Additional details and impacted files
@@           Coverage Diff            @@
##           master    #7582    +/-   ##
========================================
  Coverage   98.34%   98.34%            
========================================
  Files         924      924            
  Lines      123056   123072    +16     
  Branches    10924    11757   +833     
========================================
+ Hits       121016   121040    +24     
+ Misses       2040     2032     -8     
Flag Coverage Δ
aiguard 58.22% <ø> (-0.03%) ⬇️
aiguard-integration 57.08% <ø> (ø)
apm-bucket-0 58.48% <ø> (-0.03%) ⬇️
apm-bucket-1 64.59% <ø> (-0.03%) ⬇️
apm-bucket-2 63.64% <ø> (-0.03%) ⬇️
apm-bucket-3 60.93% <ø> (-0.03%) ⬇️
apm-capabilities-tracing 62.30% <6.06%> (-0.01%) ⬇️
apm-integrations-aerospike 57.61% <ø> (-0.03%) ⬇️
apm-integrations-confluentinc-kafka-javascript 62.53% <ø> (+0.01%) ⬆️
apm-integrations-couchbase 57.89% <ø> (-0.02%) ⬇️
apm-integrations-http 63.63% <ø> (-0.03%) ⬇️
apm-integrations-kafkajs 63.09% <ø> (-0.03%) ⬇️
apm-integrations-next 60.02% <ø> (-0.03%) ⬇️
apm-integrations-prisma 59.45% <ø> (-0.03%) ⬇️
appsec 73.94% <ø> (ø)
appsec-express_fastify_graphql 71.53% <ø> (-0.02%) ⬇️
appsec-integration 51.90% <ø> (+0.14%) ⬆️
appsec-kafka_ldapjs_lodash 64.82% <ø> (-0.03%) ⬇️
appsec-mongodb-core_mongoose_mysql 68.55% <ø> (-0.02%) ⬇️
appsec-next 58.19% <ø> (-0.03%) ⬇️
appsec-node-serialize_passport_postgres 68.24% <ø> (-0.02%) ⬇️
appsec-sourcing_stripe_template 66.54% <ø> (-0.02%) ⬇️
debugger 65.88% <ø> (-0.01%) ⬇️
instrumentations-bucket-0 52.39% <ø> (-0.03%) ⬇️
instrumentations-bucket-1 61.12% <ø> (-0.03%) ⬇️
instrumentations-bucket-10 62.98% <ø> (-0.03%) ⬇️
instrumentations-bucket-11 52.40% <ø> (-0.03%) ⬇️
instrumentations-bucket-12 52.85% <ø> (-0.03%) ⬇️
instrumentations-bucket-13 52.35% <ø> (-0.03%) ⬇️
instrumentations-bucket-2 54.32% <ø> (-0.03%) ⬇️
instrumentations-bucket-3 60.09% <ø> (-0.03%) ⬇️
instrumentations-bucket-4 52.98% <ø> (-0.03%) ⬇️
instrumentations-bucket-5 58.32% <ø> (-0.03%) ⬇️
instrumentations-bucket-6 61.68% <ø> (-0.03%) ⬇️
instrumentations-bucket-7 59.22% <ø> (-0.03%) ⬇️
instrumentations-bucket-8 60.49% <ø> (-0.03%) ⬇️
instrumentations-bucket-9 62.35% <ø> (-0.03%) ⬇️
instrumentations-instrumentation-couchbase 51.91% <ø> (-0.03%) ⬇️
instrumentations-integration-esbuild 34.05% <ø> (ø)
llmobs-ai_anthropic_bedrock 63.32% <ø> (-0.03%) ⬇️
llmobs-bucket-1 62.62% <ø> (-0.02%) ⬇️
llmobs-openai 63.29% <ø> (-0.03%) ⬇️
llmobs-sdk 65.41% <ø> (-0.03%) ⬇️
llmobs-vertex-ai 59.97% <ø> (-0.03%) ⬇️
master-coverage 98.34% <100.00%> (?)
openfeature 54.70% <ø> (+<0.01%) ⬆️
openfeature-unit 53.43% <ø> (-0.03%) ⬇️
platform-core_esbuild_instrumentations-misc 40.67% <ø> (-0.02%) ⬇️
platform-integration 62.31% <ø> (ø)
platform-shimmer_unit-guardrails_webpack 39.34% <ø> (-0.02%) ⬇️
plugins-bucket-0 57.76% <ø> (-0.03%) ⬇️
plugins-bucket-1 55.17% <ø> (ø)
plugins-bucket-11 63.09% <ø> (+0.34%) ⬆️
plugins-bucket-17 ?
plugins-bucket-18 62.90% <ø> (-0.55%) ⬇️
plugins-bucket-19 60.96% <ø> (-1.60%) ⬇️
plugins-bucket-20 62.90% <ø> (-2.40%) ⬇️
plugins-bucket-4 59.39% <ø> (-0.03%) ⬇️
plugins-bullmq_cassandra_cookie 62.68% <ø> (-0.03%) ⬇️
plugins-cookie-parser_crypto_dd-trace-api 57.56% <ø> (-0.03%) ⬇️
plugins-fetch_fs_generic-pool 59.67% <ø> (+0.02%) ⬆️
plugins-google-cloud-pubsub_grpc_handlebars 65.68% <100.00%> (-0.01%) ⬇️
plugins-hapi_hono_ioredis 61.10% <ø> (-0.03%) ⬇️
plugins-jest_knex_langgraph 56.36% <ø> (?)
plugins-jest_langgraph_ldapjs ?
plugins-ldapjs_light-my-request_limitd-client 59.36% <ø> (?)
plugins-light-my-request_limitd-client_lodash ?
plugins-lodash_mariadb_memcached 58.85% <ø> (?)
plugins-mariadb_memcached_mercurius ?
plugins-moleculer_mongodb_mongodb-core 62.84% <ø> (?)
plugins-mongodb_mongodb-core_mongoose ?
plugins-mongoose_multer_mysql 59.85% <ø> (?)
plugins-multer_mysql_mysql2 ?
plugins-mysql2_nats_node-serialize 61.63% <ø> (?)
plugins-nats_node-serialize_opensearch ?
plugins-opensearch_passport-http_pino 60.41% <ø> (?)
plugins-passport-http_pino_postgres ?
plugins-postgres_process_pug 59.13% <ø> (?)
plugins-process_pug_redis ?
plugins-redis_router_sequelize 62.97% <ø> (?)
plugins-test-and-upstream-rhea_undici_url 62.54% <ø> (?)
plugins-undici_url_valkey ?
plugins-valkey_vm_winston 58.83% <ø> (?)
plugins-vm_winston_ws ?
plugins-ws 60.46% <ø> (?)
profiling 63.07% <ø> (-0.03%) ⬇️
serverless-aws-sdk-aws-sdk 55.71% <ø> (-0.03%) ⬇️
serverless-aws-sdk-bedrockruntime 55.43% <ø> (-0.03%) ⬇️
serverless-aws-sdk-client 57.17% <ø> (-0.03%) ⬇️
serverless-aws-sdk-dynamodb 56.37% <ø> (-0.03%) ⬇️
serverless-aws-sdk-eventbridge 49.93% <ø> (-0.03%) ⬇️
serverless-aws-sdk-kinesis 60.15% <ø> (-0.03%) ⬇️
serverless-aws-sdk-lambda 58.14% <ø> (-0.03%) ⬇️
serverless-aws-sdk-s3 56.30% <ø> (-0.03%) ⬇️
serverless-aws-sdk-serverless-peer-service 60.56% <ø> (-0.03%) ⬇️
serverless-aws-sdk-sns 60.96% <ø> (-0.03%) ⬇️
serverless-aws-sdk-sqs 61.40% <ø> (-0.03%) ⬇️
serverless-aws-sdk-stepfunctions 56.29% <ø> (-0.03%) ⬇️
serverless-aws-sdk-util 52.19% <ø> (-0.03%) ⬇️
serverless-bucket-0 55.22% <ø> (ø)
serverless-bucket-1 60.09% <ø> (-0.03%) ⬇️
test-optimization-cucumber 72.99% <ø> (-0.04%) ⬇️
test-optimization-cypress 66.42% <ø> (+0.08%) ⬆️
test-optimization-jest 74.39% <ø> (-0.05%) ⬇️
test-optimization-mocha 74.72% <ø> (+0.04%) ⬆️
test-optimization-playwright-playwright-atr 61.43% <ø> (-0.01%) ⬇️
test-optimization-playwright-playwright-efd 61.62% <ø> (ø)
test-optimization-playwright-playwright-final-status 61.59% <ø> (-0.01%) ⬇️
test-optimization-playwright-playwright-impacted-tests 61.31% <ø> (+0.16%) ⬆️
test-optimization-playwright-playwright-reporting 61.21% <ø> (-0.01%) ⬇️
test-optimization-playwright-playwright-test-management 62.15% <ø> (+<0.01%) ⬆️
test-optimization-playwright-playwright-test-span 61.33% <ø> (-0.07%) ⬇️
test-optimization-selenium 60.71% <ø> (-0.16%) ⬇️
test-optimization-testopt 59.21% <ø> (+0.08%) ⬆️
test-optimization-vitest 71.33% <ø> (+0.01%) ⬆️

Flags with carried forward coverage won't be shown. Click here to find out more.

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

robcarlan-datadog added a commit that referenced this pull request Feb 23, 2026
* fix(kafkajs): sync DSM context to currentStore to prevent context leaking

When processing concurrent Kafka messages, the DSM context was being set
via enterWith() on AsyncLocalStorage but not synced to ctx.currentStore.
Since ctx.currentStore is what gets returned from bindStart and bound to
async continuations via runStores, this caused DSM context to leak between
concurrent message handlers.

The fix syncs the DSM context to ctx.currentStore after DSM operations
complete, ensuring each handler's async continuations maintain the correct
DSM context.

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* test(dsm): add regression test for context propagation race condition

Add tests that verify DSM context is properly scoped to each handler's
async continuations when using diagnostic channels with runStores.

The tests verify that:
1. ctx.currentStore has dataStreamsContext after setDataStreamsContext
2. Concurrent handlers maintain isolated DSM contexts
3. DSM context persists through multiple async boundaries

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* fix(dsm): sync DSM context to currentStore to prevent context leaking

Add syncToStore helper in DSM context module that syncs DSM context
from AsyncLocalStorage to ctx.currentStore after DSM operations.

This fixes a race condition where DSM context was being set via
enterWith() but not synced to ctx.currentStore, which is what gets
bound to async continuations via store.run(). Without syncing, DSM
context would leak between concurrent message handlers.

Updated plugins:
- kafkajs (consumer, producer)
- amqplib (consumer, producer)
- bullmq (consumer, producer)
- rhea (consumer)
- google-cloud-pubsub (consumer, producer)
- aws-sdk (sqs, kinesis)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* remove comment

* chore: remove contrived regression test

The unit test was useful for validating the fix during development
but is contrived and doesn't add value as a permanent test.

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* test(dsm): add unit tests for syncToStore and integration spy tests

Add comprehensive tests for the new syncToStore helper:
- Unit tests for syncToStore in context.spec.js covering normal
  operation, edge cases, and integration with setDataStreamsContext
- Spy tests in 6 integration test files to verify syncToStore is
  called after produce and consume operations

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* fix(dsm): call syncToStore via module object for testability

Change all plugins to call DataStreamsContext.syncToStore(ctx)
instead of destructuring syncToStore at import time. This allows
sinon spies to intercept calls during testing.

When functions are destructured at require-time, they bind to the
original function reference. Spies set up later on the module
object don't affect these bindings. Calling via the module object
ensures spies work correctly.

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* fix(dsm): remove invalid producer-side syncToStore tests for AWS SDK

The syncToStore fix only applies to the consumer path where bindStart
is used. AWS SDK plugins (SQS, Kinesis) use requestInject for producers,
which doesn't need context synchronization since:

1. requestInject is called before the request is sent
2. DSM context is encoded directly into the message
3. There's no async continuation where context leaking would occur

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* chore: remove unnecessary comments from test files

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* fix(dsm): move DSM logic from bindStart to start for correct context binding

By moving DSM checkpoint/encoding logic into start() (which runs after
the child context is bound), the DSM context naturally lands in the
correct async store — eliminating the need for the syncToStore workaround
in kafkajs, amqplib, bullmq, and rhea consumer plugins.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix(dsm): remove syncToStore functionality

The actual fix for the context propagation race condition is moving DSM
logic from bindStart to start. The syncToStore workaround is unnecessary.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix(dsm): move DSM logic from bindStart to start for google-cloud-pubsub

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* chore: remove extra blank lines in google-cloud-pubsub

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* chore: fix comma placement in amqplib producer setCheckpoint call

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(kafkajs): add regression test for DSM context leak between concurrent consumers

When two KafkaJS consumers process messages concurrently and each
produces to a different topic, the DSM (Data Streams Monitoring) context
leaks between them. The first consumer to process loses its DSM context
entirely (null parent), while the second consumer picks up the first's
context instead of its own.

This test forces the interleaving by using promise gates to ensure both
eachMessage handlers have fired before either produces, reliably
reproducing the bug.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(kafkajs): remove redundant comments

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(dsm): add regression tests for DSM context leak in amqplib, bullmq, rhea, google-cloud-pubsub

Same root cause as the kafkajs test: ctx.currentStore is set by
startSpan before decodeDataStreamsContext/setCheckpoint call enterWith,
so the DSM context is never included in the bound store for async
continuations.

Verified deterministic: 10/10 failures across all plugins and versions.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(dsm): fix lint and remove redundant comments

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(rhea): use const for senderAOut

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(rhea): remove rhea DSM regression test

Rhea's producer DSM checkpoint happens in a separate encode hook, not
in bindStart/start, so the current fix doesn't cover it. Will be
addressed in a separate PR.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix(rhea): remove rhea consumer changes, tracked separately in #7581

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(dsm): move google-cloud-pubsub changes to separate PR

Moved DSM context fix and regression test for google-cloud-pubsub
to #7582 to fix independently.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.5 <noreply@anthropic.com>
dd-octo-sts Bot pushed a commit that referenced this pull request Feb 24, 2026
* fix(kafkajs): sync DSM context to currentStore to prevent context leaking

When processing concurrent Kafka messages, the DSM context was being set
via enterWith() on AsyncLocalStorage but not synced to ctx.currentStore.
Since ctx.currentStore is what gets returned from bindStart and bound to
async continuations via runStores, this caused DSM context to leak between
concurrent message handlers.

The fix syncs the DSM context to ctx.currentStore after DSM operations
complete, ensuring each handler's async continuations maintain the correct
DSM context.

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* test(dsm): add regression test for context propagation race condition

Add tests that verify DSM context is properly scoped to each handler's
async continuations when using diagnostic channels with runStores.

The tests verify that:
1. ctx.currentStore has dataStreamsContext after setDataStreamsContext
2. Concurrent handlers maintain isolated DSM contexts
3. DSM context persists through multiple async boundaries

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* fix(dsm): sync DSM context to currentStore to prevent context leaking

Add syncToStore helper in DSM context module that syncs DSM context
from AsyncLocalStorage to ctx.currentStore after DSM operations.

This fixes a race condition where DSM context was being set via
enterWith() but not synced to ctx.currentStore, which is what gets
bound to async continuations via store.run(). Without syncing, DSM
context would leak between concurrent message handlers.

Updated plugins:
- kafkajs (consumer, producer)
- amqplib (consumer, producer)
- bullmq (consumer, producer)
- rhea (consumer)
- google-cloud-pubsub (consumer, producer)
- aws-sdk (sqs, kinesis)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* remove comment

* chore: remove contrived regression test

The unit test was useful for validating the fix during development
but is contrived and doesn't add value as a permanent test.

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* test(dsm): add unit tests for syncToStore and integration spy tests

Add comprehensive tests for the new syncToStore helper:
- Unit tests for syncToStore in context.spec.js covering normal
  operation, edge cases, and integration with setDataStreamsContext
- Spy tests in 6 integration test files to verify syncToStore is
  called after produce and consume operations

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* fix(dsm): call syncToStore via module object for testability

Change all plugins to call DataStreamsContext.syncToStore(ctx)
instead of destructuring syncToStore at import time. This allows
sinon spies to intercept calls during testing.

When functions are destructured at require-time, they bind to the
original function reference. Spies set up later on the module
object don't affect these bindings. Calling via the module object
ensures spies work correctly.

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* fix(dsm): remove invalid producer-side syncToStore tests for AWS SDK

The syncToStore fix only applies to the consumer path where bindStart
is used. AWS SDK plugins (SQS, Kinesis) use requestInject for producers,
which doesn't need context synchronization since:

1. requestInject is called before the request is sent
2. DSM context is encoded directly into the message
3. There's no async continuation where context leaking would occur

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* chore: remove unnecessary comments from test files

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* fix(dsm): move DSM logic from bindStart to start for correct context binding

By moving DSM checkpoint/encoding logic into start() (which runs after
the child context is bound), the DSM context naturally lands in the
correct async store — eliminating the need for the syncToStore workaround
in kafkajs, amqplib, bullmq, and rhea consumer plugins.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix(dsm): remove syncToStore functionality

The actual fix for the context propagation race condition is moving DSM
logic from bindStart to start. The syncToStore workaround is unnecessary.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix(dsm): move DSM logic from bindStart to start for google-cloud-pubsub

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* chore: remove extra blank lines in google-cloud-pubsub

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* chore: fix comma placement in amqplib producer setCheckpoint call

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(kafkajs): add regression test for DSM context leak between concurrent consumers

When two KafkaJS consumers process messages concurrently and each
produces to a different topic, the DSM (Data Streams Monitoring) context
leaks between them. The first consumer to process loses its DSM context
entirely (null parent), while the second consumer picks up the first's
context instead of its own.

This test forces the interleaving by using promise gates to ensure both
eachMessage handlers have fired before either produces, reliably
reproducing the bug.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(kafkajs): remove redundant comments

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(dsm): add regression tests for DSM context leak in amqplib, bullmq, rhea, google-cloud-pubsub

Same root cause as the kafkajs test: ctx.currentStore is set by
startSpan before decodeDataStreamsContext/setCheckpoint call enterWith,
so the DSM context is never included in the bound store for async
continuations.

Verified deterministic: 10/10 failures across all plugins and versions.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(dsm): fix lint and remove redundant comments

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(rhea): use const for senderAOut

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(rhea): remove rhea DSM regression test

Rhea's producer DSM checkpoint happens in a separate encode hook, not
in bindStart/start, so the current fix doesn't cover it. Will be
addressed in a separate PR.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix(rhea): remove rhea consumer changes, tracked separately in #7581

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(dsm): move google-cloud-pubsub changes to separate PR

Moved DSM context fix and regression test for google-cloud-pubsub
to #7582 to fix independently.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.5 <noreply@anthropic.com>
juan-fernandez pushed a commit that referenced this pull request Mar 5, 2026
* fix(kafkajs): sync DSM context to currentStore to prevent context leaking

When processing concurrent Kafka messages, the DSM context was being set
via enterWith() on AsyncLocalStorage but not synced to ctx.currentStore.
Since ctx.currentStore is what gets returned from bindStart and bound to
async continuations via runStores, this caused DSM context to leak between
concurrent message handlers.

The fix syncs the DSM context to ctx.currentStore after DSM operations
complete, ensuring each handler's async continuations maintain the correct
DSM context.

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* test(dsm): add regression test for context propagation race condition

Add tests that verify DSM context is properly scoped to each handler's
async continuations when using diagnostic channels with runStores.

The tests verify that:
1. ctx.currentStore has dataStreamsContext after setDataStreamsContext
2. Concurrent handlers maintain isolated DSM contexts
3. DSM context persists through multiple async boundaries

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* fix(dsm): sync DSM context to currentStore to prevent context leaking

Add syncToStore helper in DSM context module that syncs DSM context
from AsyncLocalStorage to ctx.currentStore after DSM operations.

This fixes a race condition where DSM context was being set via
enterWith() but not synced to ctx.currentStore, which is what gets
bound to async continuations via store.run(). Without syncing, DSM
context would leak between concurrent message handlers.

Updated plugins:
- kafkajs (consumer, producer)
- amqplib (consumer, producer)
- bullmq (consumer, producer)
- rhea (consumer)
- google-cloud-pubsub (consumer, producer)
- aws-sdk (sqs, kinesis)

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* remove comment

* chore: remove contrived regression test

The unit test was useful for validating the fix during development
but is contrived and doesn't add value as a permanent test.

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* test(dsm): add unit tests for syncToStore and integration spy tests

Add comprehensive tests for the new syncToStore helper:
- Unit tests for syncToStore in context.spec.js covering normal
  operation, edge cases, and integration with setDataStreamsContext
- Spy tests in 6 integration test files to verify syncToStore is
  called after produce and consume operations

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* fix(dsm): call syncToStore via module object for testability

Change all plugins to call DataStreamsContext.syncToStore(ctx)
instead of destructuring syncToStore at import time. This allows
sinon spies to intercept calls during testing.

When functions are destructured at require-time, they bind to the
original function reference. Spies set up later on the module
object don't affect these bindings. Calling via the module object
ensures spies work correctly.

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* fix(dsm): remove invalid producer-side syncToStore tests for AWS SDK

The syncToStore fix only applies to the consumer path where bindStart
is used. AWS SDK plugins (SQS, Kinesis) use requestInject for producers,
which doesn't need context synchronization since:

1. requestInject is called before the request is sent
2. DSM context is encoded directly into the message
3. There's no async continuation where context leaking would occur

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* chore: remove unnecessary comments from test files

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

* fix(dsm): move DSM logic from bindStart to start for correct context binding

By moving DSM checkpoint/encoding logic into start() (which runs after
the child context is bound), the DSM context naturally lands in the
correct async store — eliminating the need for the syncToStore workaround
in kafkajs, amqplib, bullmq, and rhea consumer plugins.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix(dsm): remove syncToStore functionality

The actual fix for the context propagation race condition is moving DSM
logic from bindStart to start. The syncToStore workaround is unnecessary.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix(dsm): move DSM logic from bindStart to start for google-cloud-pubsub

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* chore: remove extra blank lines in google-cloud-pubsub

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* chore: fix comma placement in amqplib producer setCheckpoint call

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(kafkajs): add regression test for DSM context leak between concurrent consumers

When two KafkaJS consumers process messages concurrently and each
produces to a different topic, the DSM (Data Streams Monitoring) context
leaks between them. The first consumer to process loses its DSM context
entirely (null parent), while the second consumer picks up the first's
context instead of its own.

This test forces the interleaving by using promise gates to ensure both
eachMessage handlers have fired before either produces, reliably
reproducing the bug.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(kafkajs): remove redundant comments

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(dsm): add regression tests for DSM context leak in amqplib, bullmq, rhea, google-cloud-pubsub

Same root cause as the kafkajs test: ctx.currentStore is set by
startSpan before decodeDataStreamsContext/setCheckpoint call enterWith,
so the DSM context is never included in the bound store for async
continuations.

Verified deterministic: 10/10 failures across all plugins and versions.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(dsm): fix lint and remove redundant comments

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(rhea): use const for senderAOut

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(rhea): remove rhea DSM regression test

Rhea's producer DSM checkpoint happens in a separate encode hook, not
in bindStart/start, so the current fix doesn't cover it. Will be
addressed in a separate PR.

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* fix(rhea): remove rhea consumer changes, tracked separately in #7581

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>

* test(dsm): move google-cloud-pubsub changes to separate PR

Moved DSM context fix and regression test for google-cloud-pubsub
to #7582 to fix independently.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.5 <noreply@anthropic.com>
@robcarlan-datadog
robcarlan-datadog force-pushed the rob.carlan/dsm-context-leak-pubsub-test branch from e93f79d to 3f8fe3b Compare June 23, 2026 19:59
@robcarlan-datadog

Copy link
Copy Markdown
Contributor Author

@codex review

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 3f8fe3bc43

ℹ️ About Codex in GitHub

Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".

Comment thread packages/datadog-plugin-google-cloud-pubsub/src/producer.js
robcarlan-datadog and others added 3 commits July 17, 2026 12:37
…ubsub

Same root cause as the kafkajs test: ctx.currentStore is set by
startSpan before decodeDataStreamsContext/setCheckpoint call enterWith,
so the DSM context is never included in the bound store for async
continuations.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
Move DSM checkpoint logic from bindStart to start in both consumer and
producer plugins so DSM context is properly bound in the async store.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
@robcarlan-datadog
robcarlan-datadog force-pushed the rob.carlan/dsm-context-leak-pubsub-test branch from b839c52 to 52b7b61 Compare July 17, 2026 16:40
@robcarlan-datadog

Copy link
Copy Markdown
Contributor Author

@codex review

@chatgpt-codex-connector

Copy link
Copy Markdown

Codex Review: Didn't find any major issues. Keep it up!

Reviewed commit: 52b7b61b68

ℹ️ About Codex in GitHub

Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review".

If Codex has suggestions, it will comment; otherwise it will react with 👍.

When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".

@robcarlan-datadog
robcarlan-datadog marked this pull request as ready for review July 17, 2026 17:02
@robcarlan-datadog
robcarlan-datadog requested review from a team as code owners July 17, 2026 17:02
@robcarlan-datadog
robcarlan-datadog requested review from crysmags and removed request for a team July 17, 2026 17:02
@robcarlan-datadog robcarlan-datadog changed the title fix(dsm): fix DSM context propagation for google-cloud-pubsub fix(google-cloud-pubsub): fix DSM context propagation for google-cloud-pubsub Jul 20, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant