From a4443445b05aba92cd52b521b87417308d6bea86 Mon Sep 17 00:00:00 2001 From: Tianning Li Date: Wed, 23 Sep 2026 11:52:00 -0400 Subject: [PATCH] chore(runtime): add explicit identity refresh Signed-off-by: Tianning Li --- ddtrace/internal/_runtime_id.py | 67 ++++- ddtrace/internal/runtime/__init__.py | 4 + tests/tracer/runtime/test_runtime_id.py | 328 ++++++++++++++++++++++++ 3 files changed, 396 insertions(+), 3 deletions(-) diff --git a/ddtrace/internal/_runtime_id.py b/ddtrace/internal/_runtime_id.py index ef745dda165..fe4472ee2b6 100644 --- a/ddtrace/internal/_runtime_id.py +++ b/ddtrace/internal/_runtime_id.py @@ -1,3 +1,4 @@ +import logging import typing as t import uuid @@ -6,12 +7,16 @@ from . import forksafe +log = logging.getLogger(__name__) + + __all__ = [ "get_ancestor_runtime_id", "get_process_role", "get_runtime_id", "get_parent_runtime_id", "get_runtime_propagation_envs", + "refresh_identity", ] @@ -31,6 +36,7 @@ def _generate_runtime_id() -> str: # IMPORTANT: Do not change t.Set to set until minimum Python version is 3.11+ # Module-level set[...] in Python 3.10 affects import timing. See packages.py for details. _ON_RUNTIME_ID_CHANGE: t.Set[t.Callable[[str], None]] = set() # noqa: UP006 +_ON_RUNTIME_IDENTITY_REFRESH: t.Set[t.Callable[[str], None]] = set() # noqa: UP006 def on_runtime_id_change(cb: t.Callable[[str], None]) -> None: @@ -42,6 +48,45 @@ def on_runtime_id_change(cb: t.Callable[[str], None]) -> None: _ON_RUNTIME_ID_CHANGE.add(cb) +def on_runtime_identity_refresh(cb: t.Callable[[str], None]) -> None: + """Register a callback to be called after an explicit runtime identity refresh. + + This is separate from the fork callback because a logical runtime replacement + must not be treated as a child process or update the fork lineage. + """ + global _ON_RUNTIME_IDENTITY_REFRESH + _ON_RUNTIME_IDENTITY_REFRESH.add(cb) + + +def _notify_runtime_id_callbacks(callbacks: t.Set[t.Callable[[str], None]]) -> None: # noqa: UP006 + for cb in list(callbacks): + try: + cb(_RUNTIME_ID) + except Exception: + log.exception("Exception ignored in runtime ID callback %r", cb) + + +def _notify_runtime_identity_refresh_callbacks(*, raise_on_error: bool = False) -> None: # noqa: UP006 + # Direct refresh callers keep subscriber failures isolated so every component gets + # a chance to rebuild. The MicroVM coordinator opts into propagation so a failed + # rebuild leaves its completion guard unset and the same identity can be retried. + for cb in list(_ON_RUNTIME_IDENTITY_REFRESH): + if raise_on_error: + cb(_RUNTIME_ID) + continue + try: + cb(_RUNTIME_ID) + except Exception: + log.exception("Exception ignored in runtime ID callback %r", cb) + + +def _refresh_runtime_id() -> None: + global _RUNTIME_ID + + _RUNTIME_ID = _generate_runtime_id() + _notify_runtime_id_callbacks(_ON_RUNTIME_ID_CHANGE) + + @forksafe.register def _set_runtime_id() -> None: global _RUNTIME_ID, _ANCESTOR_RUNTIME_ID, _PARENT_RUNTIME_ID @@ -51,9 +96,25 @@ def _set_runtime_id() -> None: _ANCESTOR_RUNTIME_ID = _RUNTIME_ID _PARENT_RUNTIME_ID = _RUNTIME_ID - _RUNTIME_ID = _generate_runtime_id() - for cb in _ON_RUNTIME_ID_CHANGE: - cb(_RUNTIME_ID) + _refresh_runtime_id() + + +def refresh_identity(raise_on_error: bool = False) -> None: + """Regenerate the runtime ID without recording fork lineage. + + Unlike a fork, this does not update _PARENT_RUNTIME_ID / _ANCESTOR_RUNTIME_ID: + the previous runtime ID was not a real parent process, so recording it there + would make get_process_role() and friends misreport a fork lineage that never + existed. Use this when a new logical process instance is created by a mechanism + other than fork(). + """ + # This is an opt-in lifecycle path: non-MicroVM processes do not call it from + # ordinary request handling, so their runtime-ID and fork behavior is unchanged. + # Notify consumers that only need the new ID first. The explicit refresh + # callbacks below are for components that must rebuild restore-sensitive + # state, which is different from the fork handling in _set_runtime_id(). + _refresh_runtime_id() + _notify_runtime_identity_refresh_callbacks(raise_on_error=raise_on_error) def get_runtime_id() -> str: diff --git a/ddtrace/internal/runtime/__init__.py b/ddtrace/internal/runtime/__init__.py index f29a8dbf11d..7657b2c46a5 100644 --- a/ddtrace/internal/runtime/__init__.py +++ b/ddtrace/internal/runtime/__init__.py @@ -4,6 +4,8 @@ from ddtrace.internal._runtime_id import get_runtime_id from ddtrace.internal._runtime_id import get_runtime_propagation_envs 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 __all__ = [ @@ -13,4 +15,6 @@ "get_parent_runtime_id", "get_runtime_propagation_envs", "on_runtime_id_change", + "on_runtime_identity_refresh", + "refresh_identity", ] diff --git a/tests/tracer/runtime/test_runtime_id.py b/tests/tracer/runtime/test_runtime_id.py index 0b44fccb902..4254cc4fa97 100644 --- a/tests/tracer/runtime/test_runtime_id.py +++ b/tests/tracer/runtime/test_runtime_id.py @@ -40,6 +40,85 @@ def test_get_runtime_id_fork(): assert exit_code == 42 +@pytest.mark.subprocess(env={"PYTHONWARNINGS": "ignore::DeprecationWarning"}) +def test_fork_notifies_runtime_id_subscribers(): + import os + + from ddtrace.internal import runtime + + seen = [] + + def on_change(new_id): + seen.append(new_id) + + runtime.on_runtime_id_change(on_change) + + child = os.fork() + if child == 0: + assert seen == [runtime.get_runtime_id()] + os._exit(42) + + _, status = os.waitpid(child, 0) + assert os.WEXITSTATUS(status) == 42 + + +@pytest.mark.subprocess(env={"PYTHONWARNINGS": "ignore::DeprecationWarning"}) +def test_fork_does_not_notify_runtime_identity_refresh_subscribers(): + import os + + from ddtrace.internal import runtime + + seen = [] + + def on_refresh(new_id): + seen.append(new_id) + + runtime.on_runtime_identity_refresh(on_refresh) + + child = os.fork() + if child == 0: + assert seen == [] + runtime.refresh_identity() + assert seen == [runtime.get_runtime_id()] + os._exit(42) + + _, status = os.waitpid(child, 0) + assert os.WEXITSTATUS(status) == 42 + + +@pytest.mark.subprocess( + env={"PYTHONWARNINGS": "ignore::DeprecationWarning"}, + err=lambda s: "Exception ignored in runtime ID callback" in s, +) +def test_runtime_id_callback_failure_does_not_block_other_callbacks(): + from ddtrace.internal import runtime + + seen = [] + + def on_change_raises(new_id): + seen.append(("raises", new_id)) + raise RuntimeError("callback failed") + + def on_change(new_id): + seen.append(("change", new_id)) + + def on_refresh(new_id): + seen.append(("refresh", new_id)) + + runtime.on_runtime_id_change(on_change_raises) + runtime.on_runtime_id_change(on_change) + runtime.on_runtime_identity_refresh(on_refresh) + + runtime.refresh_identity() + + runtime_id = runtime.get_runtime_id() + assert ("raises", runtime_id) in seen + assert [entry for entry in seen if entry[0] != "raises"] == [ + ("change", runtime_id), + ("refresh", runtime_id), + ] + + @pytest.mark.subprocess(env={"PYTHONWARNINGS": "ignore::DeprecationWarning"}) def test_get_runtime_id_double_fork(): import os @@ -210,3 +289,252 @@ def test_get_process_role_spawn_child() -> None: from ddtrace.internal.runtime import get_process_role assert get_process_role() == "worker", get_process_role() + + +def test_refresh_identity_changes_runtime_id(run_python_code_in_subprocess): + """refresh_identity() is the non-fork trigger for a new logical process instance.""" + code = """ +from ddtrace.internal import runtime + +runtime_id = runtime.get_runtime_id() +runtime.refresh_identity() +new_runtime_id = runtime.get_runtime_id() + +assert isinstance(new_runtime_id, str) +assert new_runtime_id != runtime_id +assert new_runtime_id == runtime.get_runtime_id() +""" + _, err, status, _ = run_python_code_in_subprocess(code) + assert status == 0, err + + +def test_refresh_identity_does_not_record_fork_lineage(run_python_code_in_subprocess): + """Unlike a fork, refresh_identity() must not make get_process_role() report a fake worker. + + The previous runtime ID was not a real parent process, so recording it as one would + corrupt process-lineage telemetry. + """ + import os + + env = os.environ.copy() + env.update( + { + "_DD_ROOT_PY_SESSION_ID": None, + "_DD_PARENT_PY_SESSION_ID": None, + "DD_TRACE_SUBPROCESS_ENABLED": "false", + } + ) + code = """ +from ddtrace.internal import runtime + +assert runtime.get_process_role() is None +assert runtime.get_parent_runtime_id() is None +assert runtime.get_ancestor_runtime_id() is None + +runtime.refresh_identity() + +assert runtime.get_process_role() is None +assert runtime.get_parent_runtime_id() is None +assert runtime.get_ancestor_runtime_id() is None +""" + _, err, status, _ = run_python_code_in_subprocess(code, env=env) + assert status == 0, err + + +def test_refresh_identity_preserves_spawned_lineage(run_python_code_in_subprocess): + import os + + env = os.environ.copy() + env.update( + { + "_DD_ROOT_PY_SESSION_ID": "ancestor-session-id", + "_DD_PARENT_PY_SESSION_ID": "parent-session-id", + "DD_TRACE_SUBPROCESS_ENABLED": "false", + } + ) + code = """ +from ddtrace.internal import runtime + +assert runtime.get_ancestor_runtime_id() == "ancestor-session-id" +assert runtime.get_parent_runtime_id() == "parent-session-id" +assert runtime.get_process_role() == "worker" + +runtime.refresh_identity() + +assert runtime.get_ancestor_runtime_id() == "ancestor-session-id" +assert runtime.get_parent_runtime_id() == "parent-session-id" +assert runtime.get_process_role() == "worker" +""" + _, err, status, _ = run_python_code_in_subprocess(code, env=env) + assert status == 0, err + + +def test_refresh_identity_preserves_fork_lineage(run_python_code_in_subprocess): + import os + + env = os.environ.copy() + env.update( + { + "_DD_ROOT_PY_SESSION_ID": None, + "_DD_PARENT_PY_SESSION_ID": None, + "DD_TRACE_SUBPROCESS_ENABLED": "false", + } + ) + code = """ +import os + +from ddtrace.internal import runtime + +root_id = runtime.get_runtime_id() +child = os.fork() + +if child == 0: + parent_id = runtime.get_parent_runtime_id() + ancestor_id = runtime.get_ancestor_runtime_id() + + assert parent_id == root_id + assert ancestor_id == root_id + assert runtime.get_process_role() == "worker" + + runtime.refresh_identity() + + assert runtime.get_parent_runtime_id() == parent_id + assert runtime.get_ancestor_runtime_id() == ancestor_id + assert runtime.get_process_role() == "worker" + os._exit(42) + +_, status = os.waitpid(child, 0) +assert os.WEXITSTATUS(status) == 42 +""" + _, err, status, _ = run_python_code_in_subprocess(code, env=env) + assert status == 0, err + + +def test_refresh_identity_notifies_subscribers(run_python_code_in_subprocess): + code = """ +from ddtrace.internal import runtime + +seen = [] + + +class _Subscriber: + def on_change(self, new_id): + seen.append(new_id) + + +subscriber = _Subscriber() +runtime.on_runtime_id_change(subscriber.on_change) + +runtime.refresh_identity() + +assert seen == [runtime.get_runtime_id()] +""" + _, err, status, _ = run_python_code_in_subprocess(code) + assert status == 0, err + + +def test_refresh_identity_notifies_refresh_subscribers(run_python_code_in_subprocess): + code = """ +from ddtrace.internal import runtime + +seen = [] + + +class _Subscriber: + def on_refresh(self, new_id): + seen.append(new_id) + + +subscriber = _Subscriber() +runtime.on_runtime_identity_refresh(subscriber.on_refresh) + +runtime.refresh_identity() + +assert seen == [runtime.get_runtime_id()] +""" + _, err, status, _ = run_python_code_in_subprocess(code) + assert status == 0, err + + +def test_refresh_identity_callback_failure_does_not_block_other_refresh_callbacks( + run_python_code_in_subprocess, +): + code = """ +from ddtrace.internal import runtime + +seen = [] + + +def on_refresh_raises(new_id): + seen.append(("raises", new_id)) + raise RuntimeError("refresh callback failed") + + +def on_refresh(new_id): + seen.append(("refresh", new_id)) + + +runtime.on_runtime_identity_refresh(on_refresh_raises) +runtime.on_runtime_identity_refresh(on_refresh) +runtime.refresh_identity() + +runtime_id = runtime.get_runtime_id() +assert ("raises", runtime_id) in seen +assert ("refresh", runtime_id) in seen +""" + _, err, status, _ = run_python_code_in_subprocess(code) + assert status == 0, err + + +def test_refresh_identity_raise_on_error_notifies_successful_callback(run_python_code_in_subprocess): + code = """ +from ddtrace.internal import runtime + +seen = [] + + +def on_refresh(new_id): + seen.append(new_id) + + +runtime.on_runtime_identity_refresh(on_refresh) +runtime.refresh_identity(raise_on_error=True) + +assert seen == [runtime.get_runtime_id()] +""" + _, err, status, _ = run_python_code_in_subprocess(code) + assert status == 0, err + + +def test_refresh_identity_propagates_refresh_callback_failure(run_python_code_in_subprocess): + code = """ +from ddtrace.internal import runtime + +old_runtime_id = runtime.get_runtime_id() +ordinary_callbacks = [] + + +def on_change(new_id): + ordinary_callbacks.append(new_id) + + +def on_refresh(new_id): + raise RuntimeError("refresh callback failed") + + +runtime.on_runtime_id_change(on_change) +runtime.on_runtime_identity_refresh(on_refresh) + +try: + runtime.refresh_identity(raise_on_error=True) +except RuntimeError as exc: + assert str(exc) == "refresh callback failed" +else: + raise AssertionError("refresh_identity() did not propagate the callback failure") + +new_runtime_id = runtime.get_runtime_id() +assert new_runtime_id != old_runtime_id +assert ordinary_callbacks == [new_runtime_id] +""" + _, err, status, _ = run_python_code_in_subprocess(code) + assert status == 0, err