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


def get_runtime_identity_refresh_lock() -> t.ContextManager[None]:
"""Return the lock that serializes a MicroVM identity refresh with its consumers."""
return t.cast(t.ContextManager[None], _RUNTIME_IDENTITY_REFRESH_LOCK)
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 @@ -8,6 +8,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
from ddtrace.internal.serverless import MICROVM_RUN_HOOK_METHOD
from ddtrace.internal.serverless import MICROVM_RUN_HOOK_PATH
from ddtrace.internal.serverless import in_aws_lambda_microvm
Expand All @@ -21,6 +22,7 @@
"get_runtime_propagation_envs",
"on_runtime_id_change",
"on_runtime_identity_refresh",
"remove_runtime_identity_refresh",
"listen_for_identity_refresh_hooks",
"maybe_refresh_identity",
"refresh_identity",
Expand Down
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
Loading
Loading