From 3b427418bf05d0250c38f574dc79e99f70e1b47a Mon Sep 17 00:00:00 2001 From: Florentin Labelle Date: Mon, 13 Jul 2026 17:00:04 +0200 Subject: [PATCH 1/3] fix(asm): avoid telemetry import race --- ddtrace/internal/telemetry/writer.py | 9 ++++++--- .../asm-telemetry-import-race-5be8c9d93f3a7261.yaml | 4 ++++ tests/telemetry/test_telemetry.py | 10 ++++++++++ 3 files changed, 20 insertions(+), 3 deletions(-) create mode 100644 releasenotes/notes/asm-telemetry-import-race-5be8c9d93f3a7261.yaml diff --git a/ddtrace/internal/telemetry/writer.py b/ddtrace/internal/telemetry/writer.py index 40367091bd4..efc33a1070d 100644 --- a/ddtrace/internal/telemetry/writer.py +++ b/ddtrace/internal/telemetry/writer.py @@ -374,16 +374,19 @@ def enable_sca_metadata(self) -> None: def _report_endpoints(self) -> Optional[dict[str, Any]]: """Adds a Telemetry event which sends the list of HTTP endpoints found at startup to the agent""" - import ddtrace.internal.settings.asm as asm_config_module + # AIDEV-NOTE: Do not import ASM from the telemetry worker. The worker starts while the + # ddtrace package is still initializing, so importing ASM here can deadlock with the + # main thread and expose a partially initialized module. A later heartbeat will retry. + asm_config = getattr(sys.modules.get("ddtrace.internal.settings.asm"), "config", None) - if not asm_config_module.config._api_security_endpoint_collection or not self._enabled: + if asm_config is None or not asm_config._api_security_endpoint_collection or not self._enabled: return None if not endpoint_collection.endpoints: return None with self._service_lock: - return endpoint_collection.flush(asm_config_module.config._api_security_endpoint_collection_limit) + return endpoint_collection.flush(asm_config._api_security_endpoint_collection_limit) def _report_products(self) -> dict[str, Any]: """Adds a Telemetry event which reports the enablement of an APM product""" diff --git a/releasenotes/notes/asm-telemetry-import-race-5be8c9d93f3a7261.yaml b/releasenotes/notes/asm-telemetry-import-race-5be8c9d93f3a7261.yaml new file mode 100644 index 00000000000..29d03431033 --- /dev/null +++ b/releasenotes/notes/asm-telemetry-import-race-5be8c9d93f3a7261.yaml @@ -0,0 +1,4 @@ +--- +fixes: + - | + ASM: Fixes an intermittent startup failure when instrumentation telemetry uses a short heartbeat interval. diff --git a/tests/telemetry/test_telemetry.py b/tests/telemetry/test_telemetry.py index c62199218c9..6c5a6807b1c 100644 --- a/tests/telemetry/test_telemetry.py +++ b/tests/telemetry/test_telemetry.py @@ -20,6 +20,16 @@ def test_enable(test_agent_session, run_python_code_in_subprocess): assert stderr == b"" +def test_enable_with_short_heartbeat_does_not_race_asm_import(test_agent_session, run_python_code_in_subprocess): + env = os.environ.copy() + env["DD_TELEMETRY_HEARTBEAT_INTERVAL"] = "0.00001" + env["DD_TELEMETRY_DEPENDENCY_COLLECTION_ENABLED"] = "false" + + _, stderr, status, _ = run_python_code_in_subprocess("import ddtrace", env=env) + + assert status == 0, stderr + + def test_enable_fork(test_agent_session, run_python_code_in_subprocess): """assert app-started/app-closing events are only sent in parent process""" code = """ From 59caec04b300f913b2d55802adfaa23d44b8a038 Mon Sep 17 00:00:00 2001 From: Florentin Labelle Date: Mon, 13 Jul 2026 17:09:05 +0200 Subject: [PATCH 2/3] chore(asm): remove implementation note --- ddtrace/internal/telemetry/writer.py | 3 --- 1 file changed, 3 deletions(-) diff --git a/ddtrace/internal/telemetry/writer.py b/ddtrace/internal/telemetry/writer.py index efc33a1070d..b644448d097 100644 --- a/ddtrace/internal/telemetry/writer.py +++ b/ddtrace/internal/telemetry/writer.py @@ -374,9 +374,6 @@ def enable_sca_metadata(self) -> None: def _report_endpoints(self) -> Optional[dict[str, Any]]: """Adds a Telemetry event which sends the list of HTTP endpoints found at startup to the agent""" - # AIDEV-NOTE: Do not import ASM from the telemetry worker. The worker starts while the - # ddtrace package is still initializing, so importing ASM here can deadlock with the - # main thread and expose a partially initialized module. A later heartbeat will retry. asm_config = getattr(sys.modules.get("ddtrace.internal.settings.asm"), "config", None) if asm_config is None or not asm_config._api_security_endpoint_collection or not self._enabled: From 759db1acdeb95ea956778bd09c01f243f8f9b0fd Mon Sep 17 00:00:00 2001 From: Florentin Labelle Date: Mon, 13 Jul 2026 17:22:59 +0200 Subject: [PATCH 3/3] fix(telemetry): avoid dependency import race --- .../internal/telemetry/dependency_tracker.py | 84 +++++++++++-------- ...elemetry-import-race-5be8c9d93f3a7261.yaml | 4 - ...-worker-import-races-5be8c9d93f3a7261.yaml | 4 + tests/telemetry/test_dependency.py | 24 +++++- tests/telemetry/test_telemetry.py | 7 +- 5 files changed, 79 insertions(+), 44 deletions(-) delete mode 100644 releasenotes/notes/asm-telemetry-import-race-5be8c9d93f3a7261.yaml create mode 100644 releasenotes/notes/telemetry-worker-import-races-5be8c9d93f3a7261.yaml diff --git a/ddtrace/internal/telemetry/dependency_tracker.py b/ddtrace/internal/telemetry/dependency_tracker.py index 9368ad70739..99d13cca42c 100644 --- a/ddtrace/internal/telemetry/dependency_tracker.py +++ b/ddtrace/internal/telemetry/dependency_tracker.py @@ -11,6 +11,7 @@ from importlib.metadata import PackageNotFoundError import re +import sys from threading import Lock from typing import Any from typing import Iterable @@ -41,12 +42,21 @@ def _normalize_dep_name(name: str) -> str: return _NORMALIZE_RE.sub("-", name).lower() +def _get_sca_enabled() -> Optional[bool]: + try: + tracer_config = getattr(sys.modules.get("ddtrace.internal.settings._config"), "config", None) + return None if tracer_config is None else bool(tracer_config._sca_enabled) + except Exception: + log.debug("Failed to read tracer config", exc_info=True) + return None + + class DependencyTracker: """Thread-safe tracker for imported dependencies and SCA metadata. All mutable access is protected by an internal lock. - SCA-enabled state is read from ``tracer_config._sca_enabled`` so + SCA-enabled state is read from the loaded ``tracer_config._sca_enabled`` so it reacts dynamically to Remote Configuration changes instead of relying on a one-time snapshot. The DependencyEntry.metadata field state drives the wire format: @@ -70,9 +80,13 @@ def collect_report(self) -> Optional[list[dict[str, Any]]]: if not config.DEPENDENCY_COLLECTION: return None + sca_enabled = _get_sca_enabled() + if sca_enabled is None: + return None + with self._lock: newly_imported_deps = modules.get_newly_imported_modules(self._modules_already_imported) - new_deps = update_imported_dependencies(self._imported_dependencies, newly_imported_deps) + new_deps = _update_imported_dependencies(self._imported_dependencies, newly_imported_deps, sca_enabled) # Normalize once; reuse the set for sent-marking and re-report dedup. new_keys = {_normalize_dep_name(d["name"]) for d in new_deps} @@ -83,9 +97,7 @@ def collect_report(self) -> Optional[list[dict[str, Any]]]: # 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. - from ddtrace.internal.settings._config import config as tracer_config - - if not tracer_config._sca_enabled: + if not sca_enabled: return new_deps if new_deps else None re_report_deps = self._collect_rereports(new_keys) @@ -140,10 +152,8 @@ def _ensure_entry(self, package_name: str) -> None: Caller must hold self._lock. """ - from ddtrace.internal.settings._config import config as tracer_config - key = _normalize_dep_name(package_name) - if key not in self._imported_dependencies and tracer_config._sca_enabled: + if key not in self._imported_dependencies and _get_sca_enabled(): try: from importlib.metadata import version as importlib_metadata_version @@ -198,7 +208,7 @@ def enable_sca_metadata(self) -> None: Called by the SCA product on start. Sets metadata from None to [] on all existing entries so the wire format includes the "metadata" - key. Future entries pick up the flag from tracer_config._sca_enabled. + key. Future entries pick up the flag from the loaded tracer configuration. """ with self._lock: for entry in self._imported_dependencies.values(): @@ -212,36 +222,11 @@ def reset(self) -> None: self._modules_already_imported = set() -def update_imported_dependencies( +def _update_imported_dependencies( already_imported: dict[str, DependencyEntry], new_modules: Iterable[str], + sca_enabled: bool, ) -> list[dict]: - """Standalone version of dependency discovery for backward compatibility. - - Mutates *already_imported* in place, adding a DependencyEntry for each - newly discovered package. Returns the list of serialized dependency - dicts ready for the ``app-dependencies-loaded`` telemetry payload. - - SCA-enabled state is read from ``tracer_config._sca_enabled`` so it - reacts dynamically to Remote Configuration changes. - - NOTE: This function is kept for backward compatibility with - tests and benchmarks that call it directly. Production code should use - DependencyTracker instead. - - Defensive: on interpreter shutdown or partial teardown, ``importlib.metadata`` - and ``sys.path`` resolution can fail, and even the ``tracer_config`` import - itself may raise once ``sys.modules`` starts being torn down. Any exception - is swallowed so the telemetry path never propagates to ``sys.excepthook``. - """ - try: - from ddtrace.internal.settings._config import config as tracer_config - - sca_enabled = tracer_config._sca_enabled - except Exception: - log.debug("update_imported_dependencies: failed to read tracer config", exc_info=True) - return [] - deps: list[dict] = [] for module_name in new_modules: try: @@ -262,3 +247,30 @@ def update_imported_dependencies( log.debug("update_imported_dependencies: failed for %r", module_name, exc_info=True) continue return deps + + +def update_imported_dependencies( + already_imported: dict[str, DependencyEntry], + new_modules: Iterable[str], +) -> list[dict]: + """Standalone version of dependency discovery for backward compatibility. + + Mutates *already_imported* in place, adding a DependencyEntry for each + newly discovered package. Returns the list of serialized dependency + dicts ready for the ``app-dependencies-loaded`` telemetry payload. + + SCA-enabled state is read from the loaded ``tracer_config._sca_enabled`` so it + reacts dynamically to Remote Configuration changes. + + NOTE: This function is kept for backward compatibility with + tests and benchmarks that call it directly. Production code should use + DependencyTracker instead. + + Defensive: on interpreter shutdown or partial teardown, tracer configuration, + ``importlib.metadata``, and ``sys.path`` resolution can fail. Any exception is + swallowed so the telemetry path never propagates to ``sys.excepthook``. + """ + sca_enabled = _get_sca_enabled() + if sca_enabled is None: + return [] + return _update_imported_dependencies(already_imported, new_modules, sca_enabled) diff --git a/releasenotes/notes/asm-telemetry-import-race-5be8c9d93f3a7261.yaml b/releasenotes/notes/asm-telemetry-import-race-5be8c9d93f3a7261.yaml deleted file mode 100644 index 29d03431033..00000000000 --- a/releasenotes/notes/asm-telemetry-import-race-5be8c9d93f3a7261.yaml +++ /dev/null @@ -1,4 +0,0 @@ ---- -fixes: - - | - ASM: Fixes an intermittent startup failure when instrumentation telemetry uses a short heartbeat interval. diff --git a/releasenotes/notes/telemetry-worker-import-races-5be8c9d93f3a7261.yaml b/releasenotes/notes/telemetry-worker-import-races-5be8c9d93f3a7261.yaml new file mode 100644 index 00000000000..a59fa5bd9ac --- /dev/null +++ b/releasenotes/notes/telemetry-worker-import-races-5be8c9d93f3a7261.yaml @@ -0,0 +1,4 @@ +--- +fixes: + - | + Telemetry: Fixes intermittent startup failures caused by the telemetry worker importing tracer and ASM settings before initialization completes. diff --git a/tests/telemetry/test_dependency.py b/tests/telemetry/test_dependency.py index 6d0761d2539..258621e1ac6 100644 --- a/tests/telemetry/test_dependency.py +++ b/tests/telemetry/test_dependency.py @@ -751,8 +751,8 @@ def _sca_enabled(self): assert result == [] - def test_update_imported_dependencies_swallows_tracer_config_import_failure(self): - """If the tracer_config import itself fails (e.g. sys.modules torn down), return [] without raising.""" + def test_update_imported_dependencies_swallows_unavailable_tracer_config(self): + """If tracer_config is unavailable (e.g. sys.modules torn down), return [] without raising.""" import sys from unittest.mock import patch @@ -765,6 +765,26 @@ def test_update_imported_dependencies_swallows_tracer_config_import_failure(self assert result == [] + def test_collect_report_waits_for_tracer_config(self): + import sys + from unittest.mock import patch + + from ddtrace.internal.telemetry.dependency_tracker import DependencyTracker + + tracker = DependencyTracker() + incomplete_module = type(sys)("ddtrace.internal.settings._config") + + with ( + patch.dict(sys.modules, {"ddtrace.internal.settings._config": incomplete_module}), + patch("ddtrace.internal.telemetry.dependency_tracker.config") as mock_config, + patch("ddtrace.internal.telemetry.dependency_tracker.modules") as mock_modules, + ): + mock_config.DEPENDENCY_COLLECTION = True + result = tracker.collect_report() + + assert result is None + mock_modules.get_newly_imported_modules.assert_not_called() + def test_collect_report_marks_sent_with_normalized_lookup(self): """collect_report should find entries by normalized key when marking as sent.""" from unittest.mock import patch diff --git a/tests/telemetry/test_telemetry.py b/tests/telemetry/test_telemetry.py index 6c5a6807b1c..c2419bbefd8 100644 --- a/tests/telemetry/test_telemetry.py +++ b/tests/telemetry/test_telemetry.py @@ -20,10 +20,13 @@ def test_enable(test_agent_session, run_python_code_in_subprocess): assert stderr == b"" -def test_enable_with_short_heartbeat_does_not_race_asm_import(test_agent_session, run_python_code_in_subprocess): +@pytest.mark.parametrize("dependency_collection_enabled", [True, False]) +def test_enable_with_short_heartbeat_does_not_race_imports( + dependency_collection_enabled, test_agent_session, run_python_code_in_subprocess +): env = os.environ.copy() env["DD_TELEMETRY_HEARTBEAT_INTERVAL"] = "0.00001" - env["DD_TELEMETRY_DEPENDENCY_COLLECTION_ENABLED"] = "false" + env["DD_TELEMETRY_DEPENDENCY_COLLECTION_ENABLED"] = str(dependency_collection_enabled) _, stderr, status, _ = run_python_code_in_subprocess("import ddtrace", env=env)