diff --git a/src/agentex/lib/core/tracing/obs_ids.py b/src/agentex/lib/core/tracing/obs_ids.py index 1be803a84..1b308ddca 100644 --- a/src/agentex/lib/core/tracing/obs_ids.py +++ b/src/agentex/lib/core/tracing/obs_ids.py @@ -24,7 +24,7 @@ import os from typing import Dict, Optional, Tuple -__all__ = ("get_obs_mode", "obs_correlation") +__all__ = ("get_obs_mode", "obs_correlation", "sync_ddtrace_to_lgtm") DD_ONLY = "dd_only" DUAL = "dual" @@ -81,3 +81,31 @@ def obs_correlation() -> Dict[str, str]: if not ids: return {} return {"obs.trace_id": ids[0], "obs.span_id": ids[1]} + + +def sync_ddtrace_to_lgtm() -> None: + """In ``dual`` mode, make ddtrace adopt the active OpenTelemetry/LGTM trace + context so ddtrace-emitted spans (best-effort Datadog) share the SAME + trace_id as the OTel trace. Call at request ingress once the OTel span is + active, and after any context boundary (e.g. entering a Temporal activity). + + No-op outside ``dual`` mode, or when OTel/ddtrace is unavailable or no OTel + span is active. Never raises. + """ + if get_obs_mode() != DUAL: + return + try: + from opentelemetry import trace + + try: + from ddtrace.trace import Context # ddtrace 3.x + except ImportError: # pragma: no cover - older ddtrace layout + from ddtrace.context import Context # type: ignore[no-redef] + from ddtrace import tracer + except ImportError: + return + + sc = trace.get_current_span().get_span_context() + if not (sc and sc.is_valid): + return + tracer.context_provider.activate(Context(trace_id=sc.trace_id, span_id=sc.span_id)) diff --git a/src/agentex/lib/sdk/fastacp/base/base_acp_server.py b/src/agentex/lib/sdk/fastacp/base/base_acp_server.py index 7cb04af7a..02cd11532 100644 --- a/src/agentex/lib/sdk/fastacp/base/base_acp_server.py +++ b/src/agentex/lib/sdk/fastacp/base/base_acp_server.py @@ -33,6 +33,7 @@ from agentex.lib.environment_variables import EnvironmentVariables, refreshed_environment_variables from agentex.types.task_message_update import TaskMessageUpdate, StreamTaskMessageFull from agentex.types.task_message_content import TaskMessageContent +from agentex.lib.core.tracing.obs_ids import sync_ddtrace_to_lgtm as _sync_ddtrace_to_lgtm from agentex.lib.core.tracing.span_queue import shutdown_default_span_queue from agentex.lib.core.compat.version_guard import assert_backend_compatible from agentex.lib.sdk.fastacp.base.constants import ( @@ -101,6 +102,9 @@ async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None: # request so business spans created here (and any downstream Temporal # workflow started from this request) share the same observability trace. token = _attach_incoming_trace_context(scope) + # In dual mode, make ddtrace adopt this OTel trace id so its best-effort + # Datadog spans share the same trace_id. No-op outside dual mode. + _sync_ddtrace_to_lgtm() try: await self.app(scope, receive, send) finally: