Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions ddtrace/internal/_runtime_id.py

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.

Original file line number Diff line number Diff line change
Expand Up @@ -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)

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.



def _notify_runtime_id_callbacks(callbacks: t.Set[t.Callable[[str], None]]) -> None: # noqa: UP006
for cb in list(callbacks):
try:
Expand Down
2 changes: 2 additions & 0 deletions ddtrace/internal/runtime/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
from ddtrace.internal._runtime_id import on_runtime_id_change
from ddtrace.internal._runtime_id import on_runtime_identity_refresh
from ddtrace.internal._runtime_id import refresh_identity
from ddtrace.internal._runtime_id import remove_runtime_identity_refresh


__all__ = [
Expand All @@ -17,6 +18,7 @@
"get_runtime_propagation_envs",
"on_runtime_id_change",
"on_runtime_identity_refresh",
"remove_runtime_identity_refresh",
"maybe_refresh_identity",
"refresh_identity",
]
4 changes: 4 additions & 0 deletions ddtrace/internal/telemetry/dependency.py
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,10 @@ def mark_all_metadata_sent(self) -> None:
for m in self.metadata:
m._mark_sent()

def reset_for_refresh(self) -> None:
"""Mark this dependency for reporting to a new worker."""
self._initial_report_sent = False

def add_metadata(self, cve_id: str, path: str = "", symbol: str = "", line: int = 0) -> bool:
"""Add or update reachability metadata for a CVE.

Expand Down
18 changes: 12 additions & 6 deletions ddtrace/internal/telemetry/dependency_tracker.py
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,7 @@ class DependencyTracker:

def __init__(self) -> None:
self._imported_dependencies: dict[str, DependencyEntry] = {}
self._report_all = False
self._modules_already_imported: set[str] = set()
self._lock = Lock()

Expand All @@ -77,15 +78,12 @@ def collect_report(self) -> Optional[list[dict[str, Any]]]:
new_keys = {_normalize_dep_name(d["name"]) for d in new_deps}
self._mark_sent(new_keys)

# Skip the re-report scan when SCA is disabled.
# Without SCA, no entry will ever have unsent metadata, so the
# scan over all _imported_dependencies is pure overhead (~887us
# at 10K deps). Only entries created by the SCA hook or with
# metadata attached can trigger needs_report() after initial send.
if not appsec_telemetry_config.SCA_ENABLED:
# Skip re-report scanning when SCA is disabled, except after identity refresh.
if not appsec_telemetry_config.SCA_ENABLED and not self._report_all:
return new_deps if new_deps else None

re_report_deps = self._collect_rereports(new_keys)
self._report_all = False
all_deps = new_deps + re_report_deps
return all_deps if all_deps else None

Expand Down Expand Up @@ -187,11 +185,19 @@ def enable_sca_metadata(self) -> None:
if entry.metadata is None:
entry.metadata = []

def refresh(self) -> None:
"""Preserve dependency metadata while scheduling a full report for a new worker."""
with self._lock:
for dependency in self._imported_dependencies.values():
dependency.reset_for_refresh()
self._report_all = True

def reset(self) -> None:
"""Reset all state (used on fork / queue reset)."""
with self._lock:
self._imported_dependencies = {}
self._modules_already_imported = set()
self._report_all = False


def update_imported_dependencies(
Expand Down
8 changes: 7 additions & 1 deletion ddtrace/internal/telemetry/noop_writer.py
Original file line number Diff line number Diff line change
Expand Up @@ -108,13 +108,19 @@ def set_test_session_token(self, token: Optional[str]) -> None:
def _restart_sequence(self) -> None:
pass

def _refresh_runtime_identity(self, _runtime_id: str) -> None:
pass

def _fork_writer(self) -> None:
pass

def _report_dependencies(self) -> Optional[list[dict[str, Any]]]:
return None

def _subscribe_worker_changes(self, callback: Any) -> None:
def _subscribe_worker_changes(self, callback: Any, expected_worker: Any) -> None:
pass

def _unsubscribe_worker_changes(self, callback: Any) -> None:
pass

def periodic(self, force_flush: bool = False) -> None:
Expand Down
499 changes: 417 additions & 82 deletions ddtrace/internal/telemetry/writer.py

Large diffs are not rendered by default.

72 changes: 67 additions & 5 deletions ddtrace/internal/writer/writer.py

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 Preserve the original exporter before building its replacement

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
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:
self._exporter = new_exporter

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

@litianningdatadog litianningdatadog Oct 10, 2026 •

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 afe7423: _downgrade() now reads old_exporter before calling _create_exporter(), matching set_test_session_token(), so it shuts down the original exporter and keeps the replacement that _create_exporter() has already published in MicroVMs. Covered by test_microvm_downgrade_shuts_down_original_exporter, which fails without the fix.

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 Capture the old exporter before constructing the downgrade replacement

With MicroVM telemetry enabled, a 404/415 from v0.5 triggers API downgrade. _create_exporter() publishes the replacement before this assignment captures old_exporter, so shutdown consumes the replacement instead of the original. Subsequent trace sends fail with TraceExporter has already been consumed. Capture the original exporter before construction, as the token-reconfiguration path already does.

Suggested change
with self._exporter_lock:
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:
self._exporter = new_exporter

Was this helpful? React 👍 or 👎
🤖 Bits Code Review · Open Bits AI session

@litianningdatadog litianningdatadog Oct 10, 2026 •

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 afe7423: _downgrade() now reads old_exporter before calling _create_exporter(), matching set_test_session_token(), so it shuts down the original exporter and keeps the replacement that _create_exporter() has already published in MicroVMs. Covered by test_microvm_downgrade_shuts_down_original_exporter, which fails without the fix.

Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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


Expand Down Expand Up @@ -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.
Expand All @@ -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:
Expand Down Expand Up @@ -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

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 Preserve the old exporter before publishing its replacement

In a MicroVM downgrade after a 404/415 response, this assignment publishes new_exporter before _downgrade() captures old_exporter at line 1170. Consequently both variables refer to the new exporter, _shutdown_exporter(old_exporter) immediately shuts down the exporter retained in self._exporter, and the actual old exporter is leaked, causing subsequent trace sends to use a shut-down exporter. Fresh evidence in the current code is that _downgrade() still captures self._exporter only after _create_exporter() returns; capture the old exporter before this early publication or otherwise preserve it through the locked swap.

Useful? React with 👍 / 👎.

@litianningdatadog litianningdatadog Oct 10, 2026 •

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 afe7423: _downgrade() now reads old_exporter before calling _create_exporter(), matching set_test_session_token(), so it shuts down the original exporter and keeps the replacement that _create_exporter() has already published in MicroVMs. Covered by test_microvm_downgrade_shuts_down_original_exporter, which fails without the fix.

telemetry_writer._subscribe_worker_changes(
self._on_telemetry_worker_changed, shared_worker, late_callback
)
Comment on lines +1032 to +1034

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 👍 / 👎.

@litianningdatadog litianningdatadog Oct 2, 2026 •

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. Covered by test_microvm_recreated_exporter_follows_refresh_before_caller_publishes.

Edit: this reply originally said both callers capture old_exporter first. That held for set_test_session_token() but not for _downgrade(), which read it after _create_exporter() had already published the replacement, and so shut down the exporter it kept. afe7423 makes _downgrade() read old_exporter before building the replacement too; see #19821 (comment).

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.

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)

Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand All @@ -1406,6 +1467,7 @@ def on_shutdown(self):
try:
self.periodic()
finally:
self._unsubscribe_telemetry_worker()
self._shutdown_exporter(self._exporter)


Expand Down
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.
1 change: 1 addition & 0 deletions tests/appsec/sca/test_telemetry.py
Original file line number Diff line number Diff line change
Expand Up @@ -282,6 +282,7 @@ def _make_writer_and_tracker(sca_enabled=False, deps=None, enabled=True):
appsec_telemetry_config.SCA_ENABLED = sca_enabled
writer = TelemetryWriter.__new__(TelemetryWriter)
writer._metric_lock = MagicMock()
writer._worker_access_lock = MagicMock()
writer._enabled = enabled
# The native worker is mocked: _report_dependencies() forwards to worker.add_dependency and
# returns the reported records, which is what these tests assert on.
Expand Down
50 changes: 49 additions & 1 deletion tests/telemetry/test_dependency.py
Original file line number Diff line number Diff line change
Expand Up @@ -452,9 +452,57 @@ def test_report_dependencies_no_rereport_without_new_metadata(self):

result = tracker.collect_report()

# No new deps, no new metadata -> None
assert result is None

def test_refresh_preserves_and_rereports_dependency_metadata(self):
from unittest.mock import patch

_, tracker = _make_writer_and_tracker(sca_enabled=True)
entry = DependencyEntry(name="requests", version="2.28.0", metadata=[])
entry.add_metadata("CVE-1", "requests.sessions", "send", 10)
entry.mark_initial_sent()
entry.mark_all_metadata_sent()
tracker._imported_dependencies["requests"] = entry

tracker.refresh()

with (
patch("ddtrace.internal.telemetry.dependency_tracker.modules") as mock_modules,
patch("ddtrace.internal.telemetry.dependency_tracker.telemetry_config") as mock_config,
):
mock_config.DEPENDENCY_COLLECTION = True
mock_modules.get_newly_imported_modules.return_value = set()

result = tracker.collect_report()

assert result is not None
assert result[0]["name"] == "requests"
assert len(result[0]["metadata"]) == 1
assert json.loads(result[0]["metadata"][0]["value"])["id"] == "CVE-1"

def test_refresh_rereports_dependencies_when_sca_disabled(self):
from unittest.mock import patch

_, tracker = _make_writer_and_tracker(sca_enabled=False)
entry = DependencyEntry(name="requests", version="2.28.0")
entry.mark_initial_sent()
tracker._imported_dependencies["requests"] = entry

tracker.refresh()

with (
patch("ddtrace.internal.telemetry.dependency_tracker.modules") as mock_modules,
patch("ddtrace.internal.telemetry.dependency_tracker.telemetry_config") as mock_config,
):
mock_config.DEPENDENCY_COLLECTION = True
mock_modules.get_newly_imported_modules.return_value = set()

result = tracker.collect_report()
second_result = tracker.collect_report()

assert result == [{"name": "requests", "version": "2.28.0"}]
assert second_result is None

def test_rereport_includes_all_metadata_per_rfc(self):
"""Re-report includes ALL metadata (sent + unsent) per RFC."""
from unittest.mock import patch
Expand Down
Loading
Loading