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
24 changes: 24 additions & 0 deletions src/agentex/lib/core/clients/temporal/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,26 @@ def validate_worker_interceptors(interceptors: list[Any]) -> None:
)


def _build_otel_interceptors() -> list:
"""Return an OpenTelemetry Temporal interceptor when an OTel observability
mode is active (SGP_OBS_MODE in {dual, lgtm}), so trace context propagates
across the workflow -> activity boundary. This lets business spans created
inside the model-loop activity adopt the caller's observability trace_id.

Empty in the default dd_only mode, or when the temporalio OTel contrib is
unavailable (safe no-op).
"""
import os

if os.getenv("SGP_OBS_MODE", "").strip().lower() not in ("dual", "lgtm"):
return []
try:
from temporalio.contrib.opentelemetry import TracingInterceptor
except ImportError:
return []
return [TracingInterceptor()]


async def get_temporal_client(
temporal_address: str,
metrics_url: str | None = None,
Expand Down Expand Up @@ -146,6 +166,10 @@ async def get_temporal_client(
dc = dataclasses.replace(dc, payload_codec=payload_codec)
connect_kwargs["data_converter"] = dc

otel_interceptors = _build_otel_interceptors()
if otel_interceptors:
connect_kwargs["interceptors"] = otel_interceptors

if not metrics_url:
client = await Client.connect(**connect_kwargs)
else:
Expand Down
22 changes: 22 additions & 0 deletions src/agentex/lib/core/temporal/workers/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,24 @@ def _validate_interceptors(interceptors: list) -> None:
)


def _build_otel_interceptors() -> list:
"""Return an OpenTelemetry Temporal interceptor when an OTel observability
mode is active (SGP_OBS_MODE in {dual, lgtm}), so trace context propagates
across the workflow -> activity boundary. This lets business spans created
inside the model-loop activity adopt the caller's observability trace_id.

Empty in the default dd_only mode, or when the temporalio OTel contrib is
unavailable (safe no-op).
"""
if os.getenv("SGP_OBS_MODE", "").strip().lower() not in ("dual", "lgtm"):
return []
try:
from temporalio.contrib.opentelemetry import TracingInterceptor
except ImportError:
return []
return [TracingInterceptor()]


async def get_temporal_client(
temporal_address: str,
metrics_url: str | None = None,
Expand Down Expand Up @@ -136,6 +154,10 @@ async def get_temporal_client(
dc = dataclasses.replace(dc, payload_codec=payload_codec)
connect_kwargs["data_converter"] = dc

otel_interceptors = _build_otel_interceptors()
if otel_interceptors:
connect_kwargs["interceptors"] = otel_interceptors

if not metrics_url:
client = await Client.connect(**connect_kwargs)
else:
Expand Down
Loading