Skip to content
Merged
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 src/agentex/lib/core/clients/temporal/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@
from temporalio.converter import PayloadCodec, DataConverter
from temporalio.contrib.pydantic import pydantic_data_converter

from agentex.lib.core.tracing.temporal import temporal_tracing_interceptors

# class DateTimeJSONEncoder(AdvancedJSONEncoder):
# def default(self, o: Any) -> Any:
# if isinstance(o, datetime.datetime):
Expand Down Expand Up @@ -136,6 +138,9 @@ async def get_temporal_client(
connect_kwargs: dict[str, Any] = {
"target_host": temporal_address,
"plugins": plugins,
# Propagate OTel trace context on outbound start_workflow / execute_activity
# (enabled by default; AGENTEX_TEMPORAL_TRACE_INTERCEPTOR_ENABLED=false to disable).
"interceptors": temporal_tracing_interceptors(),
}

if data_converter is not None:
Expand Down
8 changes: 7 additions & 1 deletion src/agentex/lib/core/temporal/workers/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@

from agentex.lib.utils.logging import make_logger
from agentex.lib.utils.registration import register_agent
from agentex.lib.core.tracing.temporal import temporal_tracing_interceptors
from agentex.lib.environment_variables import EnvironmentVariables
from agentex.lib.core.compat.version_guard import assert_backend_compatible

Expand Down Expand Up @@ -126,6 +127,9 @@ async def get_temporal_client(
connect_kwargs: dict[str, Any] = {
"target_host": temporal_address,
"plugins": plugins,
# Propagate OTel trace context on outbound start_workflow / execute_activity
# (enabled by default; AGENTEX_TEMPORAL_TRACE_INTERCEPTOR_ENABLED=false to disable).
"interceptors": temporal_tracing_interceptors(),
}

if data_converter is not None:
Expand Down Expand Up @@ -229,7 +233,9 @@ async def run(
max_concurrent_activities=self.max_concurrent_activities,
build_id=str(uuid.uuid4()),
debug_mode=debug_enabled, # Disable deadlock detection in debug mode
interceptors=self.interceptors, # Pass interceptors to Worker
# Tracing interceptor OUTERMOST so business interceptors (and the spans
# they create) nest under the propagated workflow/activity span.
interceptors=[*temporal_tracing_interceptors(), *self.interceptors],

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This double installs the tracing interceptor on the worker. Worker.init prepends client interceptors that also implement temporalio.worker.Interceptor, and TracingInterceptor implements both, so the instance already on this worker's client (from get_temporal_client above) applies to the worker automatically. With this line the worker runs two instances. Verified with a local repro on temporalio 1.26.0: every worker side span (RunWorkflow, CompleteWorkflow, StartActivity, RunActivity) is emitted twice, while client only wiring emits each span once with propagation intact.

Suggest interceptors=self.interceptors. The inherited instance is prepended, so tracing still sits outermost ahead of the business interceptors.

)

logger.info(f"Starting workers for task queue: {self.task_queue}")
Expand Down
73 changes: 73 additions & 0 deletions src/agentex/lib/core/tracing/temporal.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
"""OpenTelemetry trace-context propagation across Temporal boundaries.

Temporal serializes ``start_workflow`` / ``execute_activity`` across (potentially
cross-process) boundaries, and does NOT carry the active W3C ``traceparent`` by
default. So any span created inside a workflow or activity becomes a **new
detached root** -- the trace shatters at every Temporal hop.

This bites agentex directly: ``adk.tracing.span`` runs span creation as a
Temporal activity when ``in_temporal_workflow()`` is true, so without propagation
those business spans detach from the turn's obs trace.

Wiring temporalio's first-party ``TracingInterceptor`` onto the Temporal client
and worker injects the active span context into Temporal headers on the caller
side and extracts + continues it on the workflow/activity side, using the global
OpenTelemetry propagator -- so ``client -> workflow -> activity`` is one trace.

Enabled by DEFAULT. Set ``AGENTEX_TEMPORAL_TRACE_INTERCEPTOR_ENABLED=false``
(also accepts ``0`` / ``no`` / ``off``) to turn it off. It also degrades to a
no-op -- and never raises -- if temporalio's OpenTelemetry contrib isn't
importable, so enabling it by default can't break a worker.
"""

from __future__ import annotations

import os
from typing import Any

from agentex.lib.utils.logging import make_logger

logger = make_logger(__name__)

_ENABLE_ENV = "AGENTEX_TEMPORAL_TRACE_INTERCEPTOR_ENABLED"
_FALSEY = {"0", "false", "no", "off"}


def temporal_trace_interceptor_enabled() -> bool:
"""Whether the Temporal OTel trace interceptor should be installed.

Defaults to True; disabled only when ``AGENTEX_TEMPORAL_TRACE_INTERCEPTOR_ENABLED``
is set to a falsy value (``0`` / ``false`` / ``no`` / ``off``)."""
return os.environ.get(_ENABLE_ENV, "true").strip().lower() not in _FALSEY


def temporal_tracing_interceptors() -> list[Any]:
"""Interceptors that propagate OpenTelemetry trace context across Temporal.

Returns ``[TracingInterceptor()]`` (enabled by default) so callers can splat
it into a client's / worker's ``interceptors=`` list. Returns ``[]`` when
disabled via env, or when temporalio's OpenTelemetry contrib is not
importable. Never raises -- observability wiring must not break a worker.

``TracingInterceptor`` implements both the client and worker interceptor
interfaces, so the same call is used on both sides:
- on the **client**, it injects context on outbound ``start_workflow`` /
``execute_activity`` calls;
- on the **worker**, it extracts context and roots the workflow / activity
execution spans under it.
"""
if not temporal_trace_interceptor_enabled():
logger.info("Temporal OTel trace interceptor disabled via %s", _ENABLE_ENV)
return []
try:
from temporalio.contrib.opentelemetry import TracingInterceptor

# Construct inside the try so a constructor failure (not just a missing
# contrib) also falls back to a no-op instead of aborting worker startup.
return [TracingInterceptor()]
except Exception as exc: # contrib unavailable OR constructor failure -> no-op, never raise
logger.warning(
"Temporal OTel trace interceptor unavailable (%s); traces will not propagate across Temporal boundaries.",
exc,
)
return []
40 changes: 40 additions & 0 deletions tests/lib/core/tracing/test_temporal_interceptor.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
"""Unit tests for the Temporal OTel trace-interceptor wiring.

Verifies the interceptor is on by default, the opt-out env flag, and the safe
no-op fallback when temporalio's OpenTelemetry contrib isn't importable.
"""

import sys

import pytest

from agentex.lib.core.tracing import temporal as temporal_tracing


class TestTemporalTraceInterceptor:
def test_enabled_by_default(self, monkeypatch):
monkeypatch.delenv("AGENTEX_TEMPORAL_TRACE_INTERCEPTOR_ENABLED", raising=False)
assert temporal_tracing.temporal_trace_interceptor_enabled() is True

interceptors = temporal_tracing.temporal_tracing_interceptors()
assert len(interceptors) == 1
# temporalio's first-party OTel interceptor
assert type(interceptors[0]).__name__ == "TracingInterceptor"

@pytest.mark.parametrize("value", ["false", "0", "no", "off", "FALSE", "Off"])
def test_disabled_via_env(self, monkeypatch, value):
monkeypatch.setenv("AGENTEX_TEMPORAL_TRACE_INTERCEPTOR_ENABLED", value)
assert temporal_tracing.temporal_trace_interceptor_enabled() is False
assert temporal_tracing.temporal_tracing_interceptors() == []

@pytest.mark.parametrize("value", ["true", "1", "yes", "TRUE", "anything"])
def test_enabled_for_non_falsy_values(self, monkeypatch, value):
monkeypatch.setenv("AGENTEX_TEMPORAL_TRACE_INTERCEPTOR_ENABLED", value)
assert temporal_tracing.temporal_trace_interceptor_enabled() is True

def test_no_op_when_contrib_unimportable(self, monkeypatch):
# Enabled, but temporalio's OTel contrib not importable -> [] (never raises),
# so default-on can't break a worker that lacks the contrib.
monkeypatch.delenv("AGENTEX_TEMPORAL_TRACE_INTERCEPTOR_ENABLED", raising=False)
monkeypatch.setitem(sys.modules, "temporalio.contrib.opentelemetry", None)
assert temporal_tracing.temporal_tracing_interceptors() == []
Loading