Skip to content

fix(telemetry): rebuild worker on identity refresh - #19821

Open
litianningdatadog wants to merge 1 commit into
tianning.li/3-3-trace-writer-identity-refreshfrom
tianning.li/3-4-telemetry-identity-refresh
Open

litianningdatadog wants to merge 1 commit into
tianning.li/3-3-trace-writer-identity-refreshfrom
tianning.li/3-4-telemetry-identity-refresh

Conversation

@litianningdatadog

@litianningdatadog litianningdatadog commented Aug 23, 2026 •

Copy link
Copy Markdown
Contributor

Stacked PRs:

Description

Telemetry has the same stale-identity problem as traces. The native telemetry worker is created with the runtime identity available at that time, so an explicit MicroVM identity refresh must replace the worker before later telemetry is sent.

TelemetryWriter registers with the explicit identity-refresh callback registry only when running in a MicroVM (in_aws_lambda_microvm()); elsewhere the callback is never installed and the writer behaves as before. On refresh, it discards the current worker (without flushing its queued telemetry when the native drop() API is available; see below), clears the worker binding, and builds a fresh worker. If the previous worker had reported app-started, startup is reported again under the refreshed identity.

A worker rebuild starts with empty native state. The writer therefore replays accepted configuration events in sequence order and restores the latest integration and product-activation state on the replacement worker. Dependency reporting uses preserved tracker state and forces a full re-report so previously collected dependency metadata and SCA metadata are not lost.

Refresh can race with reporting calls that read or write self._worker (metrics, integrations, endpoints, configuration, logs, lifecycle, and fork handling). Each worker-accessing method now has a lock-free _without_lock implementation behind a public conditional-lock wrapper. MicroVM writers use the existing re-entrant lock; non-MicroVM writers use None and call the helper directly, avoiding nullcontext() overhead. The four add_*_metric methods are the exception: outside MicroVMs they keep their original inline body so the hot path pays no extra call frame, and MicroVM writers route through one shared locked helper.

The refresh notifies each trace exporter of the rebuilt worker on the /run thread. Re-pointing an exporter takes its exporter lock, which a send holds across libdatadog's retries (about 21 s against an unreachable agent), so in MicroVM mode the exporter does not wait for it: if a send holds the lock, the new worker is kept as pending and the next send applies it under the lock before sending. Outside MicroVMs the exporter waits for the lock as before.

Production status / native dependency

The preferred refresh path discards the old worker with the native TelemetryWorker.drop() API, which does not flush its queue. This branch pins libdatadog v43.0.1, whose TelemetryWorker exposes stop() but not drop().

Until the native API is available, refresh falls back to stop() and logs at debug level, since this happens on every refresh. The worker is still rebuilt with the refreshed runtime and session IDs, but stop() flushes the old runtime's queued telemetry (under the old IDs) and emits app-closing; its send_app_closing argument is currently ineffective. Raising instead would not help: listener exceptions are swallowed by core.dispatch, so /run proceeds either way, and the old worker would keep reporting under the snapshot's shared runtime ID for the life of every restored clone. The fallback limits stale telemetry to one bounded flush, so this degraded behavior is preferred. Once TelemetryWorker.drop() ships, the existing getattr(worker, "drop", None) check uses it automatically.

Reference

Testing

Added focused coverage for:

  • rebuilding the worker with refreshed runtime and session IDs
  • discarding the old worker without calling stop() or flushing its queue when drop() is available
  • falling back to stop() with a debug log when the native worker has no drop()
  • propagating discard failures without mutating the old worker state
  • retrying after replacement-worker build failure and preserving app-started lifecycle state
  • restoring app-started state when needed
  • recovering a failed replacement start on the next periodic() heartbeat without rebuilding the worker
  • reporting bootstrap configurations once across a worker rebuild
  • keeping a recreated trace exporter on the refreshed worker when a refresh lands before the exporter is published
  • replaying configuration, integration, and product-activation state
  • registering the identity-refresh callback only for MicroVM writers
  • serializing metric recording against a concurrent identity refresh
  • re-reporting preserved dependency metadata after refresh(), with and without SCA enabled
  • no-op writer compatibility and explicit callback registration
  • re-pointing a trace exporter during a send stuck in retries: MicroVM defers it to the next send, which applies it before sending; non-MicroVM still waits for the send
  • re-pointing an idle trace exporter immediately in both modes
  • an end-to-end MicroVM refresh with telemetry enabled and the trace exporter stuck in send(), finishing in under 1 s; this test stalls without the fix

Validation:

  • scripts/lint checks: passed
  • tests/telemetry/test_writer.py, tests/telemetry/test_dependency.py: 139 passed, 1 skipped (telemetry venv, Python 3.12)
  • tests/tracer/runtime/test_runtime_id.py, tests/tracer/test_writer.py: 215 passed, 10 skipped, 1 xpassed (tracer venv, Python 3.12)
  • Manually checked with the real libdatadog exporter against an agent that accepts and never answers: /run refresh took 21.02 s before the exporter change and 0.00 s after; the agent received only the in-flight payload's 6 retry attempts

Risks

Low outside MicroVM environments: the identity-refresh callback is never registered there, and non-MicroVM worker access remains lock-free. Inside MicroVMs, worker replacement occurs only on explicit runtime identity refresh, and the re-entrant lock serializes refresh with worker access. With the current native pin, refresh succeeds through the stop() fallback; the cost is that stale telemetry is flushed under the old identity (bounded by the flush interval) and may be duplicated across clones restored from one snapshot. A worker change that arrives during a stuck send reaches that exporter on its next send rather than immediately; a writer replaced by the refresh never sends again, so it never applies it. Native/libdatadog changes are not included in this PR.

Files (13)

  • ddtrace/internal/_runtime_id.py [MODIFIED] (+8 -0)
  • ddtrace/internal/runtime/init.py [MODIFIED] (+2 -0)
  • ddtrace/internal/telemetry/dependency.py [MODIFIED] (+4 -0)
  • ddtrace/internal/telemetry/dependency_tracker.py [MODIFIED] (+12 -6)
  • ddtrace/internal/telemetry/noop_writer.py [MODIFIED] (+7 -1)
  • ddtrace/internal/telemetry/writer.py [MODIFIED] (+385 -81)
  • ddtrace/internal/writer/writer.py [MODIFIED] (+60 -4)
  • releasenotes/notes/fix-telemetry-worker-identity-refresh-19fe1775be48b696.yaml [ADDED] (+5 -0)
  • tests/appsec/sca/test_telemetry.py [MODIFIED] (+1 -0)
  • tests/telemetry/test_dependency.py [MODIFIED] (+49 -1)
  • tests/telemetry/test_writer.py [MODIFIED] (+532 -0)
  • tests/tracer/runtime/test_runtime_id.py [MODIFIED] (+64 -0)
  • tests/tracer/test_writer.py [MODIFIED] (+151 -0)

🤖 Generated with Claude Code

@litianningdatadog litianningdatadog added changelog/no-changelog A changelog entry is not required for this PR. aws-microvm Work related to AWS MicroVM onboarding labels Aug 23, 2026
@cit-pr-commenter-54b7da

cit-pr-commenter-54b7da Bot commented Aug 23, 2026 •

Copy link
Copy Markdown

Circular import analysis

⚠️ Existing circular imports

There are 1 circular imports that already exist on the base branch and have not been changed by this PR.

ddtrace.errortracking._handled_exceptions.bytecode_injector -> ddtrace.errortracking._handled_exceptions.callbacks -> ddtrace.errortracking._handled_exceptions.collector -> ddtrace.errortracking._handled_exceptions.bytecode_reporting -> ddtrace.errortracking._handled_exceptions.bytecode_injector

@cit-pr-commenter-54b7da

cit-pr-commenter-54b7da Bot commented Aug 23, 2026 •

Copy link
Copy Markdown

Dependency direction analysis

⚠️ Existing dependency direction violations

There are 201 dependency direction violations that already exist on the base branch and have not been changed by this PR.

Show existing violations (showing 5 of 201 highest severity)
ddtrace.internal.tracemethods -×-> ddtrace.trace  (internal-core -> product:tracing, score=132)
ddtrace.internal.opentelemetry.span -×-> ddtrace.trace  (product:opentelemetry -> product:tracing, score=130)
ddtrace.internal.opentelemetry.context -×-> ddtrace.trace  (product:opentelemetry -> product:tracing, score=130)
ddtrace.llmobs._integrations.pydantic_ai -×-> ddtrace.trace  (product:llmobs -> product:tracing, score=130)
ddtrace.debugging._debugger -×-> ddtrace.trace  (product:debugging -> product:tracing, score=130)

To see all violations, download the layers-base.json and layers-pr.json artifacts from this CI job and run:

uv run --script scripts/import-analysis/layers.py compare layers-base.json layers-pr.json

@cit-pr-commenter-54b7da

cit-pr-commenter-54b7da Bot commented Aug 23, 2026 •

Copy link
Copy Markdown

Codeowners resolved as

Resolved from the full PR diff against tianning.li/3-3-trace-writer-identity-refresh using the target branch CODEOWNERS file.
CODEOWNERS team requests not listed below are not required by the current file set.

ddtrace/internal/_runtime_id.py                                         @DataDog/apm-core-python
ddtrace/internal/runtime/__init__.py                                    @DataDog/apm-sdk-capabilities-python
ddtrace/internal/telemetry/dependency.py                                @DataDog/apm-python
ddtrace/internal/telemetry/dependency_tracker.py                        @DataDog/apm-python
ddtrace/internal/telemetry/noop_writer.py                               @DataDog/apm-python
ddtrace/internal/telemetry/writer.py                                    @DataDog/apm-python
ddtrace/internal/writer/writer.py                                       @DataDog/apm-core-python
releasenotes/notes/fix-telemetry-worker-identity-refresh-19fe1775be48b696.yaml  @DataDog/apm-python
tests/appsec/sca/test_telemetry.py                                      @DataDog/asm-python
tests/telemetry/test_dependency.py                                      @DataDog/apm-core-python @DataDog/apm-python
tests/telemetry/test_writer.py                                          @DataDog/apm-core-python @DataDog/apm-python
tests/tracer/runtime/test_runtime_id.py                                 @DataDog/apm-sdk-capabilities-python
tests/tracer/test_writer.py                                             @DataDog/apm-sdk-capabilities-python

@litianningdatadog litianningdatadog changed the title fix(telemetry): rebuild worker on identity refresh chore(telemetry): rebuild worker on identity refresh Aug 23, 2026
@datadog-datadog-prod-us1-2

datadog-datadog-prod-us1-2 Bot commented Aug 23, 2026 •

Copy link
Copy Markdown
Contributor

Pipelines  Tests

❌ Errors

Your PR has failed checks. Please review the issues below and take necessary action before merging.

🚦 1 Pipeline job failed

DataDog/apm-reliability/dd-trace-py | download win_arm64 wheels

View more details · View in GitLab

ℹ️ Info

No other issues found (see more)

🧪 All tests passed
❄️ No new flaky tests detected

🔄 Datadog retried 1 test - 1 passed on retry View in Datadog

🚧 23 tests that failed were ignored due to quarantine View in Datadog

Useful? React with 👍 / 👎

This comment will be updated automatically if new data arrives.
🔗 Commit SHA: 11d4256 | Docs | View more details | Give us feedback!

@pr-commenter

pr-commenter Bot commented Aug 23, 2026 •

Copy link
Copy Markdown

Benchmarks

Benchmark execution time: 2026-10-05 21:08:45

Comparing candidate commit 11d4256 in PR branch tianning.li/3-4-telemetry-identity-refresh with baseline commit 2096489 in branch tianning.li/3-3-trace-writer-identity-refresh.

📊 Benchmarking dashboard

Found 0 performance improvements and 3 performance regressions! Performance is the same for 361 metrics, 9 unstable metrics, 4 known flaky benchmarks, 4 flaky benchmarks without significant changes.

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 ----------------------------------'

scenario:httppropagationextract-empty_headers

  • 🟥 execution_time [+86.248ns; +106.278ns] or [+11.580%; +14.269%]

scenario:msgpackencoderscenario-simple_one_span

  • 🟥 execution_time [+553.306ns; +610.621ns] or [+13.580%; +14.987%]

scenario:otelspan-start

  • 🟥 execution_time [+1.827ms; +2.722ms] or [+7.229%; +10.768%]

Unstable benchmarks

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

scenario:coreapiscenario-context_with_data_listeners

  • unstable execution_time [-798.839ns; +652.870ns] or [-7.630%; +6.236%]

scenario:coreapiscenario-core_dispatch_1_listener

  • unstable execution_time [-32.376ns; +45.599ns] or [-4.936%; +6.952%]

scenario:coreapiscenario-core_dispatch_50_listeners

  • unstable execution_time [-1907.122ns; +1904.157ns] or [-9.608%; +9.593%]

scenario:coreapiscenario-core_dispatch_exception_listeners

  • unstable execution_time [-1754.802ns; +1860.390ns] or [-9.291%; +9.850%]

scenario:coreapiscenario-core_dispatch_listeners

  • unstable execution_time [-387.856ns; +368.480ns] or [-9.121%; +8.665%]

scenario:coreapiscenario-core_dispatch_no_args_listeners

  • unstable execution_time [-235.850ns; +221.029ns] or [-8.821%; +8.266%]

scenario:coreapiscenario-core_dispatch_with_results_1_listener

  • unstable execution_time [-78.029ns; +108.999ns] or [-5.796%; +8.097%]

scenario:coreapiscenario-core_dispatch_with_results_50_listeners

  • unstable execution_time [-3996.651ns; +5278.679ns] or [-8.334%; +11.007%]

scenario:coreapiscenario-core_dispatch_with_results_listeners

  • unstable execution_time [-910.525ns; +974.264ns] or [-8.959%; +9.586%]

Known flaky benchmarks

These benchmarks are marked as flaky and will not trigger a failure. Modify FLAKY_BENCHMARKS_REGEX to control which benchmarks are marked as flaky.

scenario:httppropagationinject-ids_only

  • 🟥 execution_time [+2.987µs; +3.101µs] or [+21.302%; +22.113%]

scenario:span-start

  • 🟥 execution_time [+1.268ms; +1.705ms] or [+9.420%; +12.662%]

scenario:telemetryaddmetric-1-count-metric-1-times

  • 🟥 execution_time [+288.249ns; +331.813ns] or [+15.207%; +17.505%]

scenario:tracer-small

  • 🟥 execution_time [+42.633µs; +43.961µs] or [+16.238%; +16.744%]

Known flaky benchmarks without significant changes:

  • scenario:errortrackingflasksqli-baseline
  • scenario:flasksimple-iast-get
  • scenario:sethttpmeta-all-enabled
  • scenario:telemetryaddmetric-record-100-metrics

@litianningdatadog
litianningdatadog force-pushed the tianning.li/2-flask-web-request-starting-event branch from d59e112 to 16a5332 Compare August 24, 2026 02:23
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-4-telemetry-identity-refresh branch from 2165d36 to a44a4b8 Compare August 24, 2026 02:24
@litianningdatadog
litianningdatadog force-pushed the tianning.li/2-flask-web-request-starting-event branch from 16a5332 to 8dd7e8e Compare August 24, 2026 02:31
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-4-telemetry-identity-refresh branch from a44a4b8 to 2b1bf99 Compare August 24, 2026 02:31
@litianningdatadog
litianningdatadog force-pushed the tianning.li/2-flask-web-request-starting-event branch 5 times, most recently from cde3045 to a0e3c42 Compare August 24, 2026 23:56
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-4-telemetry-identity-refresh branch from 2b1bf99 to b4e8c78 Compare August 25, 2026 13:23
@litianningdatadog
litianningdatadog requested a lite review from Copilot August 25, 2026 13:36

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Pull request overview

Ensure telemetry emitted after a runtime identity refresh uses the refreshed runtime ID by tearing down the existing native telemetry worker and allowing it to be rebuilt.

Changes:

  • Wire TelemetryWriter to runtime identity changes and rebuild (drop) its native worker on refresh.
  • Stop the live native worker during identity refresh to prevent continued heartbeats with stale identity.
  • Add tests covering worker teardown on identity refresh and wiring through runtime.refresh_identity().

Reviewed changes

Copilot reviewed 2 out of 2 changed files in this pull request and generated 2 comments.

File Description
ddtrace/internal/telemetry/writer.py Subscribes to runtime-id changes and stops/drops the native telemetry worker on identity refresh.
tests/telemetry/test_writer.py Adds identity-refresh tests for worker stop/drop behavior and wiring through runtime.refresh_identity().

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread tests/telemetry/test_writer.py Outdated
Comment thread ddtrace/internal/telemetry/writer.py Outdated
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-4-telemetry-identity-refresh branch 4 times, most recently from 24f368e to fb190e7 Compare September 22, 2026 05:38
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-3-trace-writer-identity-refresh branch from 7d904ad to f2664ed Compare September 22, 2026 20:21
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-4-telemetry-identity-refresh branch from fb190e7 to e96b84d Compare September 22, 2026 20:25
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-3-trace-writer-identity-refresh branch 2 times, most recently from 8376c12 to ba3f177 Compare September 23, 2026 19:18
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-4-telemetry-identity-refresh branch from e96b84d to b51f682 Compare September 23, 2026 19:27
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-3-trace-writer-identity-refresh branch from ba3f177 to 805d8d1 Compare September 23, 2026 19:41
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-4-telemetry-identity-refresh branch 3 times, most recently from f7876cb to b84f21b Compare September 24, 2026 15:37
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-3-trace-writer-identity-refresh branch from 805d8d1 to 8dbfa50 Compare September 24, 2026 17:40
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-4-telemetry-identity-refresh branch 2 times, most recently from 156be45 to c8b45b6 Compare September 24, 2026 17:57
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-3-trace-writer-identity-refresh branch from 8dbfa50 to 19be48b Compare September 24, 2026 18:25
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-4-telemetry-identity-refresh branch from c8b45b6 to c641a12 Compare September 24, 2026 18:29
@litianningdatadog
litianningdatadog force-pushed the tianning.li/3-3-trace-writer-identity-refresh branch from 19be48b to 7f7d597 Compare September 24, 2026 19:13
@chatgpt-codex-connector

chatgpt-codex-connector Bot commented Sep 29, 2026 •

Copy link
Copy Markdown

Codex Review Summary

This comment shows the latest Codex review activity on this pull request.

Review Status Commit Review trigger
📝 Code Review ✅ Completed 2026-10-05T20:45:42.938287Z 11d4256 New commits
🔒 Security Review ✅ Completed 2026-10-05T20:45:20.379123Z 11d4256 New commits
ℹ️ About Codex in GitHub

Your team has set up Codex to 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" or "@codex security review".

Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings.

@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: 359a1ce7d0

ℹ️ About Codex in GitHub

Your team has set up Codex to 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 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

Comment on lines +466 to +472
if discard is None:
log.warning(
"Native TelemetryWorker does not support discard; stopping the worker %s. "
"Upgrade the native ddtrace dependency to avoid flushing stale telemetry.",
reason,
)
self._stop_worker(False, reason)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P1 Badge Fail refresh when the worker lacks discard support

With the currently pinned libdatadog v43.0.1, TelemetryWorker has no drop() method—the wrapper in src/native/telemetry.rs:254-280 only exposes stop(), which explicitly ignores send_app_closing, drains queued data, and emits app-closing. Consequently every MicroVM identity refresh takes this fallback, flushes telemetry carrying the previous runtime identity, and then returns successfully so the identity coordinator will not retry. This preserves the exact cross-invocation misattribution being fixed; the callback should fail without stopping until a non-flushing discard operation is available.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Keeping the stop() fallback; see "Production status / native dependency" in the PR description. Raising wouldn't keep stale telemetry off the wire. The old worker would stay alive and keep reporting under the snapshot's shared runtime ID in every restored clone until some later /run retried. The fallback limits stale telemetry to one bounded flush, and the getattr(worker, "drop", None) check switches to drop() automatically once libdatadog exposes it.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Addressed, see the reply above. This can be resolved if there are no other objections.

Comment on lines +474 to +477
discard()
self._worker = None
self.started = False
_unbind_metric_recorders(self)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Synchronize MetricRecorder calls before discarding the worker

During a MicroVM refresh concurrent callers using get_metric_recorder() remain unsynchronized: MetricRecorder.add() in ddtrace/internal/telemetry/metrics.py:231-235 reads its worker and calls add_point() without _worker_access_lock, while this path discards the native worker before rebinding recorders. Such a caller can therefore submit to the old worker after the runtime ID has rotated or race its teardown, losing or misattributing the metric despite the new locking around TelemetryWriter.add_*_metric().

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Leaving recorders lock-free on purpose: they are the zero-lookup path for every IAST aspect and propagation inject. A point racing the refresh can at worst go to the old worker, where it's flushed under the old ID by the stop() fallback or dropped by drop(). It can't raise, because native add_point swallows errors from a torn-down worker (drop_on_err in src/native/telemetry.rs). Unbinding before discarding wouldn't close the race either, since a thread may already hold the old worker. It would also break the guarantee that a failed drop() leaves state untouched for the retry (test_microvm_identity_refresh_discard_failure_propagates).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Addressed, see the reply above. This can be resolved if there are no other objections.

Comment on lines +1208 to +1210
if was_started:
self.app_started()
self._identity_refresh_started = False

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Propagate replacement worker start failures

If the old worker had started but the replacement worker's native start() call fails, app_started() catches the exception and returns with self.started still false. This code nevertheless clears _identity_refresh_started and returns success, causing the /run identity coordinator to remove the callback from its retry queue and mark the refresh complete; telemetry then remains permanently unstarted for that logical runtime. Verify self.started after this call and raise so the existing refresh retry mechanism can run again.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Reverted in 7b5eedc. A failed replacement start isn't permanent: the dependency collector calls periodic() every heartbeat, and periodic() begins with app_started(), which retries the start. Raising here would only add a retry on the next /run, which may never arrive. If one did arrive, it would discard a worker the heartbeat had already started under the new runtime ID, and with the stop() fallback that sends app-closing and then app-started for the new ID. test_microvm_identity_refresh_failed_replacement_start_recovers_on_periodic covers the heartbeat recovery.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Addressed, see the reply above. This can be resolved if there are no other objections.

Comment on lines +368 to +370
if get_parent_runtime_id() is None:
if not self.started:
self.add_configurations(get_python_config_vars())

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Avoid re-recording bootstrap configurations after replay

On every root-process identity rebuild, _replay_worker_state() has already copied all accepted configurations—including the initial get_python_config_vars() entries—into the replacement worker, but _discard_worker() reset started to false, so this branch immediately records the Python configuration list a second time with new sequence IDs. The replacement therefore reports duplicate configuration changes on every refresh, and the duplicates are appended back into the bounded 5,000-entry replay deque, potentially evicting real earlier configuration events near the limit. Bootstrap configurations should only be added for the initial worker, not after a state replay.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Fixed in 7b5eedc: enable() skips add_configurations(get_python_config_vars()) when a MicroVM writer already has configurations to replay, since _replay_worker_state() has just copied them onto the replacement worker. Covered by test_microvm_identity_refresh_reports_bootstrap_configurations_once, which checks both workers and the replay deque.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Addressed, see the reply above. This can be resolved if there are no other objections.

Comment thread ddtrace/internal/telemetry/writer.py Outdated
Comment on lines +982 to +986
if self._worker_access_lock:
with self._worker_access_lock:
self._add_count_metric_without_lock(namespace, name, value, tags)
else:
self._add_count_metric_without_lock(namespace, name, value, tags)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Preserve the metric fast path outside MicroVMs

For every normal non-MicroVM process, _worker_access_lock is None, yet each of the four metric APIs now performs an extra Python helper call per point. The previous implementation explicitly kept these bodies inline because avoiding that call frame measurably lowers overhead, and these methods run for every IAST aspect and propagation injection; this therefore regresses the dominant production path even though it requires no refresh synchronization. Keep the existing body inline for the None branch or select a specialized implementation once at initialization.

AGENTS.md reference: AGENTS.md:L16-L16

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Fixed in 7b5eedc: outside MicroVMs the four add_*_metric methods keep their original inline body again. MicroVM writers route through a single _add_metric_with_lock() helper that holds the worker access lock across the lookup and the native call.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Addressed, see the reply above. This can be resolved if there are no other objections.

@emmettbutler emmettbutler left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Deferring review since the base branch is not main

"Upgrade the native ddtrace dependency to avoid flushing stale telemetry.",
reason,
)
self._stop_worker(False, reason)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

P1 Reject refresh when discard support is unavailable

When the pinned native worker lacks drop(), identity refresh calls stop(), flushing queued telemetry and app-closing under the stale runtime identity. The callback then succeeds instead of remaining retryable, defeating the refresh contract for MicroVM users; raise without mutating the worker and update the legacy-worker test accordingly.

Was this helpful? React 👍 or 👎
🤖 Bits Code Review · @DataDog review to ask questions · Open Bits AI session

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Keeping the stop() fallback; see "Production status / native dependency" in the PR description. Raising wouldn't keep stale telemetry off the wire. The old worker would stay alive and keep reporting under the snapshot's shared runtime ID in every restored clone until some later /run retried. The fallback limits stale telemetry to one bounded flush, and the getattr(worker, "drop", None) check switches to drop() automatically once libdatadog exposes it.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Addressed, see the reply above. This can be resolved if there are no other objections.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

P2 Remove the production test warning

Every expected MicroVM identity refresh emits a warning claiming a test regenerated the runtime ID. This creates misleading application-log noise and exposes both runtime IDs during normal production lifecycle handling.

Suggested change
_refresh_runtime_id()

Was this helpful? React 👍 or 👎
🤖 Bits Code Review · @DataDog review to ask questions · Open Bits AI session

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

There is no such warning in ddtrace/internal/_runtime_id.py on this branch: refresh_identity() calls _refresh_runtime_id() directly, and nothing in the module logs a test-regeneration warning. This finding looks anchored to stale content, so no change here.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Addressed, see the reply above. This can be resolved if there are no other objections.

@datadog-datadog-prod-us1-2 datadog-datadog-prod-us1-2 Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Bits Code Review: FAIL

With the pinned worker lacking drop(), MicroVM refresh flushes old-identity telemetry and reports success instead of remaining retryable. Expected refreshes also emit a misleading test warning.

Open Bits AI session

🤖 Bits Code Review · Commit 359a1ce · @DataDog review to ask questions

@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: ee6f68585c

ℹ️ About Codex in GitHub

Your team has set up Codex to 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 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

# all child hooks have completed, without unparking the inherited Tokio runtime.
forksafe.register(self._fork_writer)
if self._is_microvm:
on_runtime_identity_refresh(self._refresh_runtime_identity)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Order telemetry refresh before rebuilding the trace writer

When both the tracer and this writer register refresh callbacks, _ON_RUNTIME_IDENTITY_REFRESH is a set, so callback order is nondeterministic. If the tracer callback runs first, it publishes a replacement NativeWriter whose exporter is initially attached to the old telemetry worker; a concurrent trace flush can then record trace_api.* metrics against the previous runtime until this callback runs, at which point those points are either discarded or, with the current stop fallback, flushed under the stale identity. Make the worker replacement and trace-writer rebuild deterministically ordered or otherwise atomic with respect to trace sends.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Leaving the ordering as is in this PR. The stale window is narrow:

  • In a MicroVM both callbacks run back to back under _RUNTIME_IDENTITY_REFRESH_LOCK, and the tracer callback recreates the writer with drop_buffered_traces=True, so the replacement exporter starts with an empty buffer. Only traces finished by other threads between the two callbacks can record trace_api.* points against the old worker.
  • As soon as the telemetry callback runs, _notify_worker_changed re-points the exporter at the new worker.
  • If telemetry runs first, the old writer's exporter is re-pointed before the tracer rebuilds, so that order has no window.

The worst case is a few health-metric points from that window being lost (with drop()) or attributed to the previous runtime ID (with the current stop() fallback). That doesn't justify adding ordering to the callback registry.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Addressed, see the reply above. This can be resolved if there are no other objections.

Comment on lines +1011 to +1013
telemetry_writer._subscribe_worker_changes(
self._on_telemetry_worker_changed, shared_worker, late_callback
)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Keep the unpublished exporter synchronized after subscribing

During exporter recreation in set_test_session_token() or _downgrade(), _create_exporter() runs before the caller assigns its result to self._exporter. If an identity refresh occurs after this subscription returns but before that assignment, the stored callback updates the old self._exporter, not the newly built local exporter; the caller then publishes the new exporter still attached to the discarded worker, and no later notification corrects it. The subscription or publication must keep the local exporter targeted throughout this race window.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Fixed in 7b5eedc: in MicroVMs _create_exporter() publishes the new exporter to self._exporter before subscribing, so a refresh notification after subscription updates the exporter the caller is about to keep. Both callers capture old_exporter first, so publishing early doesn't affect their shutdown of the old one. Covered by test_microvm_recreated_exporter_follows_refresh_before_caller_publishes.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Addressed, see the reply above. This can be resolved if there are no other objections.


def remove_runtime_identity_refresh(cb: t.Callable[[str], None]) -> None:
"""Unregister a callback for explicit runtime identity refreshes."""
_ON_RUNTIME_IDENTITY_REFRESH.discard(cb)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Serialize callback removal with identity refresh

In a MicroVM, refresh_identity() iterates _ON_RUNTIME_IDENTITY_REFRESH while holding _RUNTIME_IDENTITY_REFRESH_LOCK, but telemetry shutdown can call this unsynchronized discard() concurrently. If removal occurs while line 175 converts the set into the pending callback list, Python can raise RuntimeError: Set changed size during iteration after the runtime ID has already rotated; the retry path then sees the new ID with an empty pending list and can mark the transition complete without rebuilding any registered component. Acquire the identity-refresh lock when mutating this registry.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Not changing this. Taking _RUNTIME_IDENTITY_REFRESH_LOCK here would create a lock-order inversion: refresh_identity() holds that lock and then each callback takes its component lock (telemetry's _worker_access_lock), while TelemetryWriter.disable() calls this function with _worker_access_lock already held. Two threads on those paths would deadlock.

The race can't produce the error on GIL builds: list(_ON_RUNTIME_IDENTITY_REFRESH) and callback_snapshot[:] = _ON_RUNTIME_IDENTITY_REFRESH each copy the set inside a single C call, so no other thread can mutate it mid-copy. If the removal lands after the copy, the telemetry callback still runs and returns immediately because disable() already set _enabled = False.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Addressed, see the reply above. This can be resolved if there are no other objections.

@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: bfdcca813b

ℹ️ About Codex in GitHub

Your team has set up Codex to 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 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

Comment thread tests/telemetry/test_writer.py Outdated
Comment on lines +1029 to +1030
telemetry_writer._is_microvm = True
telemetry_writer._worker_access_lock = telemetry_writer._enable_lock

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Use a reentrant lock in the MicroVM test

When the native worker gains drop() and this test is no longer skipped, it will hang in runtime.refresh_identity(): the fixture constructs the writer outside a MicroVM, so _enable_lock is a non-reentrant forksafe.Lock, and assigning that same lock to _worker_access_lock means the refresh callback acquires it and then enable() tries to acquire it again after discarding the worker. Construct the writer with MicroVM detection enabled or replace both lock attributes with the same RLock before invoking the refresh.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Fixed in 7b5eedc: the test now installs one threading.RLock() as both _enable_lock and _worker_access_lock, matching a MicroVM writer, so the refresh can re-enter it through enable(). The test still skips until the native worker exposes drop().

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Addressed, see the reply above. This can be resolved if there are no other objections.

@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: 0654d6e62a

ℹ️ About Codex in GitHub

Your team has set up Codex to 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 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

seq_id,
)
# Replay the latest integration state for the rebuilt worker.
for integration_name, state in self._integrations.items():

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Snapshot replay state before iterating

In a non-MicroVM process, a worker rebuild (for example, the LLMObs agentless switch) publishes self._worker before entering this loop, while add_integration() remains lock-free. A concurrent registration of a previously unseen integration therefore mutates _integrations during iteration and raises RuntimeError: dictionary changed size during iteration, leaving the replacement worker only partially initialized; _product_statuses has the same race. Iterate over snapshots or keep the replacement unpublished until replay completes.

Useful? React with 👍 / 👎.

Comment on lines +452 to +453
except Exception:
log.debug("Failed to stop the native telemetry worker %s", reason, exc_info=True)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Propagate fallback stop failures during refresh

With the currently pinned worker, every identity refresh uses the stop() fallback, but this handler swallows a stop failure after _stop_worker() has already cleared self._worker. If native stop() fails before unregistering its handle—for example, ensure_runtime_after_fork() returns an error—the refresh proceeds and is marked complete while the unreachable old worker can remain active under the stale runtime ID alongside the replacement. Preserve the worker and propagate failures on the refresh path so the coordinator can retry the teardown.

Useful? React with 👍 / 👎.

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

@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: 11d42568ec

ℹ️ About Codex in GitHub

Your team has set up Codex to 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 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

Comment on lines +1230 to +1231
pending, self._pending_telemetry_worker = self._pending_telemetry_worker, None
self._set_telemetry_handle(*pending)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Preserve worker updates that arrive while consuming a pending one

When a second telemetry-worker notification arrives just after this tuple assignment reads the existing pending update but before it stores None, _on_telemetry_worker_changed() cannot acquire _exporter_lock and writes the newer update to _pending_telemetry_worker; this assignment then overwrites that update and applies only the stale worker. This can occur when a MicroVM send starts while one worker change is pending and another telemetry rebuild or teardown runs concurrently, leaving the exporter permanently attached to the previous or stopped worker. Consume the pending value without clearing a concurrently published replacement, or serialize publication through a separate lock/queue.

Useful? React with 👍 / 👎.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

aws-microvm Work related to AWS MicroVM onboarding

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants