Repository navigation
fix(telemetry): rebuild worker on identity refresh #19821
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: tianning.li/3-3-trace-writer-identity-refresh
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -62,6 +62,11 @@ def on_runtime_identity_refresh(cb: t.Callable[[str], None]) -> None: | |
| _ON_RUNTIME_IDENTITY_REFRESH.add(cb) | ||
|
|
||
|
|
||
| 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. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
In a MicroVM, Useful? React with 👍 / 👎.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Not changing this. Taking The race can't produce the error on GIL builds:
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Addressed, see the reply above. This can be resolved if there are no other objections. |
||
|
|
||
|
|
||
| def _notify_runtime_id_callbacks(callbacks: t.Set[t.Callable[[str], None]]) -> None: # noqa: UP006 | ||
| for cb in list(callbacks): | ||
| try: | ||
|
|
||
Large diffs are not rendered by default.
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
With telemetry enabled in a MicroVM, an agent returning 404/415 for v0.5/traces triggers downgrade. _create_exporter() now publishes the replacement before this assignment captures old_exporter, so downgrade shuts down the replacement instead of the original. Subsequent trace sends fail with 'TraceExporter has already been consumed' rather than continuing through v0.4. Capture the original handle before construction, as set_test_session_token() already does.
Suggested change
Was this helpful? React 👍 or 👎
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Fixed in afe7423:
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
With MicroVM telemetry enabled, a 404/415 from v0.5 triggers API downgrade.
Suggested change
Was this helpful? React 👍 or 👎
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Fixed in afe7423: |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -2,6 +2,7 @@ | |
| import binascii | ||
| from collections import defaultdict | ||
| from collections.abc import Sequence | ||
| from functools import partial | ||
| import gzip | ||
| import os | ||
| import socket | ||
|
|
@@ -74,6 +75,7 @@ | |
|
|
||
| if TYPE_CHECKING: # pragma: no cover | ||
| from ddtrace.internal.http import HTTPConnection # noqa:F401 | ||
| from ddtrace.internal.telemetry.writer import TelemetryWriter | ||
| from ddtrace.vendor.dogstatsd import DogStatsd | ||
|
|
||
|
|
||
|
|
@@ -899,6 +901,12 @@ def __init__( | |
| self._accepting_writes = True | ||
| self._exporter_dropped = False | ||
| self._owner_pid = os.getpid() | ||
| self._telemetry_writer: Optional[TelemetryWriter] = None | ||
| self._telemetry_worker_subscribed = False | ||
| # A MicroVM worker change that arrived while a send held the exporter lock. | ||
| self._pending_telemetry_worker: Optional[tuple] = None | ||
| # Serializes publishing a pending worker with a send consuming one, so the newer is not lost. | ||
| self._pending_telemetry_worker_lock = forksafe.Lock() | ||
|
|
||
| # Native exporter methods require exclusive access because PyO3 rejects | ||
| # overlapping mutable borrows. | ||
|
|
@@ -911,6 +919,7 @@ def __del__(self) -> None: | |
| try: | ||
| if getattr(self, "_owner_pid", None) != os.getpid(): | ||
| return | ||
| self._unsubscribe_telemetry_worker() | ||
| exporter = getattr(self, "_exporter", None) | ||
| if exporter is not None and not getattr(self, "_exporter_dropped", False): | ||
| if self._discarded_by_refresh: | ||
|
|
@@ -1013,14 +1022,57 @@ def _create_exporter(self) -> native.TraceExporter: | |
| exporter = builder.build(get_native_runtime()) | ||
| if shared_worker is not None: | ||
| exporter.set_telemetry_handle(shared_worker) | ||
| telemetry_writer._subscribe_worker_changes(self._on_telemetry_worker_changed) | ||
| late_callback = ( | ||
| partial(self._on_telemetry_worker_changed, exporter=exporter) if telemetry_writer._is_microvm else None | ||
| ) | ||
| if telemetry_writer._is_microvm: | ||
| # Publish before subscribing: worker-change notifications re-point self._exporter, | ||
| # so a refresh before the caller assigns this result would otherwise miss it. | ||
| self._exporter = exporter | ||
|
Comment on lines
+1029
to
+1031
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
In a MicroVM downgrade after a 404/415 response, this assignment publishes Useful? React with 👍 / 👎.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Fixed in afe7423: |
||
| telemetry_writer._subscribe_worker_changes( | ||
| self._on_telemetry_worker_changed, shared_worker, late_callback | ||
| ) | ||
|
Comment on lines
+1032
to
+1034
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
During exporter recreation in Useful? React with 👍 / 👎.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Fixed in 7b5eedc: in MicroVMs Edit: this reply originally said both callers capture
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Addressed, see the reply above. This can be resolved if there are no other objections. |
||
| else: | ||
| telemetry_writer._subscribe_worker_changes(self._on_telemetry_worker_changed, shared_worker) | ||
| self._telemetry_writer = telemetry_writer | ||
| self._telemetry_worker_subscribed = True | ||
| return exporter | ||
|
|
||
| def _on_telemetry_worker_changed(self, worker: "Optional[native.TelemetryWorker]") -> None: | ||
| def _unsubscribe_telemetry_worker(self) -> None: | ||
| if not self._telemetry_worker_subscribed: | ||
| return | ||
| telemetry_writer = self._telemetry_writer | ||
| self._telemetry_writer = None | ||
| self._telemetry_worker_subscribed = False | ||
| if telemetry_writer is not None: | ||
| telemetry_writer._unsubscribe_worker_changes(self._on_telemetry_worker_changed) | ||
|
|
||
| def _on_telemetry_worker_changed( | ||
| self, | ||
| worker: "Optional[native.TelemetryWorker]", | ||
| exporter: "Optional[native.TraceExporter]" = None, | ||
| ) -> None: | ||
| """Follow the telemetry writer onto a rebuilt worker (or off a stopped one).""" | ||
| # A send holds the exporter lock across its retries, and MicroVM identity refresh rebuilds | ||
| # the telemetry worker on the /run thread: leave the worker for the next send to apply. | ||
| if not self._exporter_lock.acquire(blocking=self._writer_lock is None): | ||
| with self._pending_telemetry_worker_lock: | ||
| self._pending_telemetry_worker = (worker, exporter) | ||
| return | ||
| try: | ||
| with self._exporter_lock: | ||
| self._exporter.set_telemetry_handle(worker) | ||
| self._pending_telemetry_worker = None | ||
| self._set_telemetry_handle(worker, exporter) | ||
| finally: | ||
| self._exporter_lock.release() | ||
|
|
||
| def _set_telemetry_handle( | ||
| self, | ||
| worker: "Optional[native.TelemetryWorker]", | ||
| exporter: "Optional[native.TraceExporter]" = None, | ||
| ) -> None: | ||
| # Caller holds the exporter lock. | ||
| try: | ||
| (exporter if exporter is not None else self._exporter).set_telemetry_handle(worker) | ||
| except Exception: | ||
| log.debug("Failed to re-point the trace exporter at the telemetry worker", exc_info=True) | ||
|
|
||
|
|
@@ -1053,9 +1105,11 @@ def set_test_session_token(self, token: Optional[str]) -> None: | |
|
|
||
| def shutdown_exporter(self) -> None: | ||
| """Tear down the native exporter without going through ``stop()``.""" | ||
| self._unsubscribe_telemetry_worker() | ||
| self._shutdown_exporter(self._exporter) | ||
|
|
||
| def _drop_exporter(self) -> None: | ||
| self._unsubscribe_telemetry_worker() | ||
| with self._exporter_lock: | ||
| # The refresh, the finishing sender, and on_shutdown() may all get here; only one drops. | ||
| if not self._exporter_dropped: | ||
|
|
@@ -1116,13 +1170,15 @@ def _downgrade(self, status, client): | |
| self._api_version = "v0.4" | ||
| # Built outside the exporter lock: _create_exporter() can take the telemetry enable lock, | ||
| # whose holder notifies _on_telemetry_worker_changed(), which takes the exporter lock. | ||
| # Read the old exporter first: in MicroVMs _create_exporter() publishes its result. | ||
| old_exporter = self._exporter | ||
| new_exporter = self._create_exporter() | ||
| with self._exporter_lock: | ||
| # A refresh may have dropped the exporter since the check above; a replacement | ||
| # installed now would be skipped by on_shutdown() and leak. | ||
| discarded = self._discarded_by_refresh | ||
| if not discarded: | ||
| old_exporter, self._exporter = self._exporter, new_exporter | ||
| self._exporter = new_exporter | ||
| if discarded: | ||
| new_exporter.drop() | ||
| return | ||
|
|
@@ -1191,6 +1247,10 @@ def _send_payload(self, payload: bytes, count: int, client: WriterClientBase): | |
| # starts after a refresh and none reaches a dropped exporter. | ||
| if self._discarded_by_refresh: | ||
| return | ||
| if self._pending_telemetry_worker is not None: | ||
| with self._pending_telemetry_worker_lock: | ||
| pending, self._pending_telemetry_worker = self._pending_telemetry_worker, None | ||
| self._set_telemetry_handle(*pending) | ||
| response_body = self._exporter.send(payload) | ||
| except native.RequestError as e: | ||
| try: | ||
|
|
@@ -1396,6 +1456,7 @@ def _stop_service( | |
|
|
||
| def on_shutdown(self): | ||
| if self._exporter_dropped: | ||
| self._unsubscribe_telemetry_worker() | ||
| return | ||
| if self._discarded_by_refresh: | ||
| # Finish the refresh's discard without the final flush. drop() avoids shutdown()'s stats flush but | ||
|
|
@@ -1406,6 +1467,7 @@ def on_shutdown(self): | |
| try: | ||
| self.periodic() | ||
| finally: | ||
| self._unsubscribe_telemetry_worker() | ||
| self._shutdown_exporter(self._exporter) | ||
|
|
||
|
|
||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,5 @@ | ||
| --- | ||
| fixes: | ||
| - | | ||
| telemetry: Fixes an issue where telemetry events can be associated with a previous | ||
| AWS Lambda MicroVM invocation after the runtime identity is refreshed. |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
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.
Was this helpful? React 👍 or 👎
🤖 Bits Code Review · @DataDog review to ask questions · Open Bits AI session
There was a problem hiding this comment.
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.pyon 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.There was a problem hiding this comment.
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.