fix(telemetry): rebuild worker on identity refresh - #19821
litianningdatadog wants to merge 1 commit into
Conversation
Circular import analysis
|
Dependency direction analysis
|
Codeowners resolved asResolved from the full PR diff against |
|
✅ All CI checks and tests passed. 🎉 All green!🧪 All tests passed 🔗 Commit SHA: 8529b92 | Docs | View more details | Give us feedback! |
BenchmarksBenchmark execution time: 2026-10-01 18:49:52 Comparing candidate commit 8529b92 in PR branch Found 0 performance improvements and 5 performance regressions! Performance is the same for 355 metrics, 9 unstable metrics, 4 known flaky benchmarks, 4 flaky benchmarks without significant changes.
|
d59e112 to
16a5332
Compare
2165d36 to
a44a4b8
Compare
16a5332 to
8dd7e8e
Compare
a44a4b8 to
2b1bf99
Compare
cde3045 to
a0e3c42
Compare
2b1bf99 to
b4e8c78
Compare
There was a problem hiding this comment.
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
TelemetryWriterto 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.
982f582 to
9c54e0e
Compare
24f368e to
fb190e7
Compare
7d904ad to
f2664ed
Compare
fb190e7 to
e96b84d
Compare
8376c12 to
ba3f177
Compare
e96b84d to
b51f682
Compare
ba3f177 to
805d8d1
Compare
f7876cb to
b84f21b
Compare
805d8d1 to
8dbfa50
Compare
156be45 to
c8b45b6
Compare
8dbfa50 to
19be48b
Compare
c8b45b6 to
c641a12
Compare
19be48b to
7f7d597
Compare
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. |
There was a problem hiding this comment.
💡 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".
| 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) |
There was a problem hiding this comment.
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 👍 / 👎.
| discard() | ||
| self._worker = None | ||
| self.started = False | ||
| _unbind_metric_recorders(self) |
There was a problem hiding this comment.
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 👍 / 👎.
| if was_started: | ||
| self.app_started() | ||
| self._identity_refresh_started = False |
There was a problem hiding this comment.
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 👍 / 👎.
| if get_parent_runtime_id() is None: | ||
| if not self.started: | ||
| self.add_configurations(get_python_config_vars()) |
There was a problem hiding this comment.
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 👍 / 👎.
| 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) |
There was a problem hiding this comment.
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 👍 / 👎.
emmettbutler
left a comment
There was a problem hiding this comment.
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) |
There was a problem hiding this comment.
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
There was a problem hiding this comment.
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.
| _refresh_runtime_id() |
Was this helpful? React 👍 or 👎
🤖 Bits Code Review · @DataDog review to ask questions · Open Bits AI session
There was a problem hiding this comment.
There was a problem hiding this comment.
💡 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) |
There was a problem hiding this comment.
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 👍 / 👎.
| telemetry_writer._subscribe_worker_changes( | ||
| self._on_telemetry_worker_changed, shared_worker, late_callback | ||
| ) |
There was a problem hiding this comment.
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 👍 / 👎.
|
|
||
| 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) |
There was a problem hiding this comment.
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 👍 / 👎.
There was a problem hiding this comment.
💡 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".
| telemetry_writer._is_microvm = True | ||
| telemetry_writer._worker_access_lock = telemetry_writer._enable_lock |
There was a problem hiding this comment.
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 👍 / 👎.
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.
TelemetryWriterregisters 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 nativedrop()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_lockimplementation behind a public conditional-lock wrapper. MicroVM writers use the existing re-entrant lock; non-MicroVM writers useNoneand call the helper directly, avoidingnullcontext()overhead on normal hot paths.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 libdatadogv43.0.1, whoseTelemetryWorkerexposesstop()but notdrop().Until the native API is available, refresh falls back to
stop()and logs a warning. The worker is still rebuilt with the refreshed runtime and session IDs, butstop()flushes the old runtime's queued telemetry (under the old IDs) and emitsapp-closing; itssend_app_closingargument is currently ineffective. Raising instead would fail the MicroVM/runlifecycle request, so this degraded behavior is preferred. OnceTelemetryWorker.drop()ships, the existinggetattr(worker, "drop", None)check uses it automatically.Reference
Testing
Added focused coverage for:
stop()or flushing its queue whendrop()is availablestop()with a warning when the native worker has nodrop()refresh(), with and without SCA enabledValidation:
scripts/lint fmt: passedscripts/lint style: passedtests/telemetry/test_writer.py: 71 passed, 1 skipped (Python 3.12)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. Native/libdatadog changes are not included in this PR.Files (12)
🤖 Generated with Claude Code