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
12 changes: 12 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,18 @@ to include examples, links to docs, or any other relevant information.

### Added

- Added experimental `temporalio.contrib.opentelemetry.ReplaySafeMeterProvider` and
`ReplaySafeLoggerProvider`, wrapping OpenTelemetry providers that drop synchronous instrument
recordings and emitted log records made from workflow code during replay. Install them as the
process-global providers when libraries record OpenTelemetry metrics or emit log events from
workflow code (e.g. Google ADK) so that workflow replays (cache eviction, worker restarts,
redeploys) do not duplicate telemetry; recordings are first-execution-only, matching
`temporalio.workflow.metric_meter()`.
`temporalio.contrib.opentelemetry.ReplaySafeTracerProvider` is now also exported.
`GoogleAdkPlugin` now warns at worker and replayer configuration time when the global
OpenTelemetry meter, tracer, or logger provider is positively identified as not replay-safe
(an OpenTelemetry SDK provider used directly).

### Changed

### Deprecated
Expand Down
58 changes: 58 additions & 0 deletions temporalio/contrib/google_adk_agents/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -199,6 +199,64 @@ agent = Agent(
)
```

## Telemetry and Workflow Replay

ADK records OpenTelemetry metrics (scope `gcp.vertex.agent`, e.g.
`gen_ai.client.token.usage`), spans, and log events (e.g. `gen_ai.choice`)
through the process-global OpenTelemetry providers from code that runs
inside the workflow. Workflow code re-executes on every replay, so with a
plain global provider each replay re-records all of that telemetry even
though no model or tool actually ran again — for example, 1 real execution
followed by 3 replays yields 4x the observations on every instrument and 4
copies of every log event. Replays happen routinely in production: workflow
cache eviction, worker restarts, redeploys, or running with
`max_cached_workflows=0`.

To avoid this, install Temporal's replay-safe providers as the global
OpenTelemetry providers. They pass recordings through on first execution and
drop them during replay:

```python
import opentelemetry._logs
import opentelemetry.metrics
import opentelemetry.trace
from opentelemetry.sdk._logs import LoggerProvider
from opentelemetry.sdk._logs.export import BatchLogRecordProcessor
from opentelemetry.sdk.metrics import MeterProvider
from opentelemetry.sdk.metrics.export import PeriodicExportingMetricReader
from opentelemetry.sdk.trace.export import BatchSpanProcessor

from temporalio.contrib.opentelemetry import (
ReplaySafeLoggerProvider,
ReplaySafeMeterProvider,
create_tracer_provider,
)

# The global set_*_provider functions only take effect once per process, so
# these wrappers must be the first and only global providers set.
opentelemetry.metrics.set_meter_provider(
ReplaySafeMeterProvider(
MeterProvider(metric_readers=[PeriodicExportingMetricReader(my_exporter)])
)
)
tracer_provider = create_tracer_provider()
tracer_provider.add_span_processor(BatchSpanProcessor(my_span_exporter))
opentelemetry.trace.set_tracer_provider(tracer_provider)
logger_provider = LoggerProvider()
logger_provider.add_log_record_processor(BatchLogRecordProcessor(my_log_exporter))
opentelemetry._logs.set_logger_provider(ReplaySafeLoggerProvider(logger_provider))
```

`GoogleAdkPlugin` warns at worker and replayer configuration time when a
global provider is installed that it can positively identify as not
replay-safe (an OpenTelemetry SDK provider used directly).

Recordings are first-execution-only, matching
`temporalio.workflow.metric_meter()`: a retried workflow task re-executes
live and can record again, and tokens consumed by failed activity attempts
are not counted. Telemetry recorded from activities (worker-side) is
unaffected.

## Integration Points

This integration provides comprehensive support for running Google ADK Agents within Temporal workflows while maintaining:
Expand Down
143 changes: 143 additions & 0 deletions temporalio/contrib/google_adk_agents/_plugin.py
Original file line number Diff line number Diff line change
@@ -1,29 +1,157 @@
from __future__ import annotations

import dataclasses
import sys
import time
import uuid
import warnings
from collections.abc import AsyncIterator, Callable
from contextlib import asynccontextmanager
from types import FrameType
from typing import Any

import opentelemetry._logs
import opentelemetry.metrics
import opentelemetry.trace
from opentelemetry._logs import LoggerProvider, NoOpLoggerProvider
from opentelemetry.metrics import MeterProvider, NoOpMeterProvider
from opentelemetry.sdk._logs import LoggerProvider as SdkLoggerProvider
from opentelemetry.sdk.metrics import MeterProvider as SdkMeterProvider
from opentelemetry.sdk.trace import TracerProvider as SdkTracerProvider
from opentelemetry.trace import (
NoOpTracerProvider,
ProxyTracerProvider,
TracerProvider,
)

from temporalio import workflow
from temporalio.contrib.google_adk_agents._mcp import TemporalMcpToolSetProvider
from temporalio.contrib.google_adk_agents._model import (
invoke_model,
invoke_model_streaming,
)
from temporalio.contrib.opentelemetry import (
ReplaySafeLoggerProvider,
ReplaySafeMeterProvider,
ReplaySafeTracerProvider,
)
from temporalio.contrib.pydantic import (
PydanticPayloadConverter,
ToJsonOptions,
)
from temporalio.converter import DataConverter, DefaultPayloadConverter
from temporalio.plugin import SimplePlugin
from temporalio.worker import (
ReplayerConfig,
WorkerConfig,
WorkflowRunner,
)
from temporalio.worker.workflow_sandbox import SandboxedWorkflowRunner

# Each classifier below returns True when the provider is replay-safe
# (including providers that drop all recordings), False when positively
# identified as replay-unsafe (an OpenTelemetry SDK provider), and None when
# it cannot be classified. Unknown provider types (e.g. a custom provider
# delegating to a replay-safe one) must not warn: a false positive is worse
# than a missed warning.


def _meter_provider_replay_safe(provider: MeterProvider) -> bool | None:
if isinstance(provider, (ReplaySafeMeterProvider, NoOpMeterProvider)):
return True
try:
# Unlike tracing's public ProxyTracerProvider, the proxy (unset) meter
# provider has no public counterpart. Import the private class lazily
# so a moved or removed symbol cannot break module import; it is
# present in opentelemetry-api 1.12 through at least 1.42.
from opentelemetry.metrics._internal import _ProxyMeterProvider

if isinstance(provider, _ProxyMeterProvider):
return True
except ImportError:
pass
return False if isinstance(provider, SdkMeterProvider) else None


def _tracer_provider_replay_safe(provider: TracerProvider) -> bool | None:
if isinstance(
provider,
(ReplaySafeTracerProvider, NoOpTracerProvider, ProxyTracerProvider),
):
return True
return False if isinstance(provider, SdkTracerProvider) else None


def _logger_provider_replay_safe(provider: LoggerProvider) -> bool | None:
if isinstance(provider, (ReplaySafeLoggerProvider, NoOpLoggerProvider)):
return True
try:
# The proxy (unset) logger provider has no public counterpart either;
# present in opentelemetry-api 1.23 through at least 1.42.
from opentelemetry._logs._internal import ProxyLoggerProvider

if isinstance(provider, ProxyLoggerProvider):
return True
except ImportError:
pass
return False if isinstance(provider, SdkLoggerProvider) else None


def _stacklevel_outside_temporalio() -> int:
# Attribute provider warnings to the nearest frame outside temporalio,
# e.g. the user's Worker(...)/Replayer(...) call or a user plugin that
# delegates here, however many plugin frames sit in between.
level = 1
frame: FrameType | None = sys._getframe(1)
while frame is not None:
module = frame.f_globals.get("__name__", "")
if module != "temporalio" and not module.startswith("temporalio."):
return level
frame = frame.f_back
level += 1
return 1


def _warn_if_global_otel_providers_not_replay_safe() -> None:
# ADK records metrics, spans, and log events through the process-global
# OpenTelemetry providers from code that runs workflow-side, so a
# non-replay-safe global provider re-emits that telemetry on every
# workflow replay. Unset (proxy) and no-op providers drop recordings and
# are fine.
stacklevel = _stacklevel_outside_temporalio()
if _meter_provider_replay_safe(opentelemetry.metrics.get_meter_provider()) is False:
warnings.warn(
"The global OpenTelemetry MeterProvider is not replay-safe: Google ADK "
"records metrics from workflow code, so every workflow replay will "
"re-record them. Wrap your provider in "
"temporalio.contrib.opentelemetry.ReplaySafeMeterProvider and make it "
"the first and only global provider set: "
"opentelemetry.metrics.set_meter_provider(ReplaySafeMeterProvider(provider))",
UserWarning,
stacklevel=stacklevel,
)
if _tracer_provider_replay_safe(opentelemetry.trace.get_tracer_provider()) is False:
warnings.warn(
"The global OpenTelemetry TracerProvider is not replay-safe: Google ADK "
"creates spans from workflow code, so every workflow replay will "
"re-emit them. Install a replay-safe provider: "
"opentelemetry.trace.set_tracer_provider("
"temporalio.contrib.opentelemetry.create_tracer_provider())",
UserWarning,
stacklevel=stacklevel,
)
if _logger_provider_replay_safe(opentelemetry._logs.get_logger_provider()) is False:
warnings.warn(
"The global OpenTelemetry LoggerProvider is not replay-safe: Google ADK "
"emits log events (e.g. gen_ai.choice) from workflow code, so every "
"workflow replay will re-emit them. Wrap your provider in "
"temporalio.contrib.opentelemetry.ReplaySafeLoggerProvider and make it "
"the first and only global provider set: "
"opentelemetry._logs.set_logger_provider(ReplaySafeLoggerProvider(provider))",
UserWarning,
stacklevel=stacklevel,
)


def setup_deterministic_runtime():
"""Configures ADK runtime for Temporal determinism.
Expand Down Expand Up @@ -118,6 +246,21 @@ def workflow_runner(runner: WorkflowRunner | None) -> WorkflowRunner:
workflow_runner=workflow_runner,
)

def configure_worker(self, config: WorkerConfig) -> WorkerConfig:
"""See base class. Also warns when the global OpenTelemetry providers
are not replay-safe, since ADK telemetry would duplicate on replay.
"""
_warn_if_global_otel_providers_not_replay_safe()
return super().configure_worker(config)

def configure_replayer(self, config: ReplayerConfig) -> ReplayerConfig:
"""See base class. Also warns when the global OpenTelemetry providers
are not replay-safe, since every replayed workflow would re-emit ADK
telemetry.
"""
_warn_if_global_otel_providers_not_replay_safe()
return super().configure_replayer(config)

def _configure_data_converter(
self, converter: DataConverter | None
) -> DataConverter:
Expand Down
55 changes: 55 additions & 0 deletions temporalio/contrib/opentelemetry/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -236,6 +236,61 @@ with tracer.start_as_current_span("my-operation") as span:
})
```

## Replay-Safe Metrics

For Temporal SDK metrics inside workflows, use `temporalio.workflow.metric_meter()`,
which is already replay-safe. However, third-party libraries (e.g. Google ADK) may
record OpenTelemetry metrics through the process-global meter provider from code
that runs inside workflows. Workflow code re-executes on every replay (cache
eviction, worker restart, redeploy), so a plain global meter provider re-records
those metrics on each replay, inflating counts.

`ReplaySafeMeterProvider` wraps your meter provider so synchronous instrument
recordings made from workflow code are dropped during replay, mirroring what
`create_tracer_provider()` does for spans:

```python
import opentelemetry.metrics
from opentelemetry.sdk.metrics import MeterProvider
from temporalio.contrib.opentelemetry import ReplaySafeMeterProvider

# set_meter_provider only takes effect once per process, so this wrapper must
# be the first and only global meter provider set, installed before any
# library records metrics.
opentelemetry.metrics.set_meter_provider(
ReplaySafeMeterProvider(MeterProvider(metric_readers=[my_reader]))
)
```

Recordings are first-execution-only, matching `workflow.metric_meter()`: a
retried workflow task re-executes live and can record again. Observable
(asynchronous) instruments and recordings made outside workflows pass through
untouched.

## Replay-Safe Log Events

Libraries may also emit OpenTelemetry log records through the process-global
logger provider from workflow code (e.g. Google ADK's `gen_ai.*` events),
which duplicate on every replay the same way. `ReplaySafeLoggerProvider`
wraps your logger provider so records emitted from workflow code are dropped
during replay:

```python
import opentelemetry._logs
from opentelemetry.sdk._logs import LoggerProvider
from temporalio.contrib.opentelemetry import ReplaySafeLoggerProvider

# set_logger_provider only takes effect once per process, so this wrapper
# must be the first and only global logger provider set, installed before
# any library emits log records.
opentelemetry._logs.set_logger_provider(
ReplaySafeLoggerProvider(my_logger_provider)
)
```

Emissions are first-execution-only: a retried workflow task re-executes live
and can emit again. Emissions outside workflows pass through untouched.

## Best Practices

1. **Register on Client**: Always register plugins/interceptors on the client, not the worker, to ensure proper context propagation
Expand Down
53 changes: 52 additions & 1 deletion temporalio/contrib/opentelemetry/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,18 +5,69 @@
propagation for distributed tracing.
"""

from typing import Any

from temporalio.contrib.opentelemetry._interceptor import (
TracingInterceptor,
TracingWorkflowInboundInterceptor,
)

_meter_provider_import_error: ImportError | None = None
try:
from temporalio.contrib.opentelemetry._meter_provider import (
ReplaySafeMeterProvider,
)
except ImportError as err:
# opentelemetry-api < 1.12 has no opentelemetry.metrics module. Keep the
# tracing integration importable and raise a clear error only when
# ReplaySafeMeterProvider is actually accessed (see __getattr__ below).
_meter_provider_import_error = err

_logger_provider_import_error: ImportError | None = None
try:
from temporalio.contrib.opentelemetry._logger_provider import (
ReplaySafeLoggerProvider,
)
except ImportError as err:
# opentelemetry-api < 1.15 has no opentelemetry._logs module. Keep the
# tracing integration importable and raise a clear error only when
# ReplaySafeLoggerProvider is actually accessed (see __getattr__ below).
_logger_provider_import_error = err

from temporalio.contrib.opentelemetry._otel_interceptor import OpenTelemetryInterceptor
from temporalio.contrib.opentelemetry._plugin import OpenTelemetryPlugin
from temporalio.contrib.opentelemetry._tracer_provider import create_tracer_provider
from temporalio.contrib.opentelemetry._tracer_provider import (
ReplaySafeTracerProvider,
create_tracer_provider,
)

__all__ = [
"TracingInterceptor",
"TracingWorkflowInboundInterceptor",
"OpenTelemetryInterceptor",
"OpenTelemetryPlugin",
"ReplaySafeLoggerProvider",
"ReplaySafeMeterProvider",
"ReplaySafeTracerProvider",
"create_tracer_provider",
]


def __getattr__(name: str) -> Any:
# Only reachable for the replay-safe providers when their guarded imports
# above failed; otherwise the module attributes exist and this is never
# called.
if name == "ReplaySafeMeterProvider":
raise ImportError(
"ReplaySafeMeterProvider requires the OpenTelemetry metrics API "
"(opentelemetry.metrics), which the installed opentelemetry-api "
"version does not provide. Install opentelemetry-api >= 1.12 "
"(>= 1.23 for synchronous gauge support)."
) from _meter_provider_import_error
if name == "ReplaySafeLoggerProvider":
raise ImportError(
"ReplaySafeLoggerProvider requires the OpenTelemetry logs API "
"(opentelemetry._logs), which the installed opentelemetry-api "
"version does not provide. Install opentelemetry-api >= 1.15."
) from _logger_provider_import_error
raise AttributeError(f"module {__name__!r} has no attribute {name!r}")
Loading
Loading