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

### Added

- Added experimental `temporalio.contrib.opentelemetry.ReplaySafeMeterProvider` and
`ReplaySafeLoggerProvider`, which wrap a user-supplied OpenTelemetry provider and 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 or tracer provider is positively identified as not replay-safe
(an OpenTelemetry SDK provider used directly).

### Changed

- The `opentelemetry` and `lambda-worker-otel` extras now require
`opentelemetry-api`/`opentelemetry-sdk` `>= 1.24`, aligning the declared floor with what
`temporalio.contrib.opentelemetry` already required in practice (it has depended on an API
added in `opentelemetry-api` 1.24 since the tracing integration was introduced, and the
lambda worker builds on it).

### Deprecated

### :boom: Breaking Changes
Expand Down
6 changes: 3 additions & 3 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -26,15 +26,15 @@ classifiers = [

[project.optional-dependencies]
grpc = ["grpcio>=1.48.2,<2"]
opentelemetry = ["opentelemetry-api>=1.11.1,<2", "opentelemetry-sdk>=1.11.1,<2"]
opentelemetry = ["opentelemetry-api>=1.24,<2", "opentelemetry-sdk>=1.24,<2"]
pydantic = ["pydantic>=2.0.0,<3"]
openai-agents = ["openai-agents>=0.17.5", "mcp>=1.9.4, <2"]
google-adk = ["google-adk>=2.2.0,<3"]
langgraph = ["langgraph>=1.1.0"]
langsmith = ["langsmith>=0.7.34,<0.9"]
lambda-worker-otel = [
"opentelemetry-api>=1.11.1,<2",
"opentelemetry-sdk>=1.11.1,<2",
"opentelemetry-api>=1.24,<2",
"opentelemetry-sdk>=1.24,<2",
"opentelemetry-exporter-otlp-proto-grpc>=1.11.1,<2",
"opentelemetry-semantic-conventions>=0.40b0,<1",
"opentelemetry-sdk-extension-aws>=2.0.0,<3",
Expand Down
60 changes: 60 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,66 @@ 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 the
global meter or tracer provider is positively identified as not replay-safe
(an OpenTelemetry SDK provider used directly). The global logger provider is
not checked because the OpenTelemetry logs SDK has no public import path yet,
but the same replay duplication applies to it.

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
87 changes: 87 additions & 0 deletions temporalio/contrib/google_adk_agents/_plugin.py
Original file line number Diff line number Diff line change
@@ -1,12 +1,18 @@
from __future__ import annotations

import dataclasses
import inspect
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.metrics
import opentelemetry.trace

from temporalio import workflow
from temporalio.contrib.google_adk_agents._mcp import TemporalMcpToolSetProvider
from temporalio.contrib.google_adk_agents._model import (
Expand All @@ -20,11 +26,72 @@
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


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
own_frame: FrameType | None = inspect.currentframe()
frame = own_frame.f_back if own_frame is not None else None
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. Warn only on providers positively identified as
# replay-unsafe: an OpenTelemetry SDK provider used directly as the
# global. Anything else stays silent -- unset (proxy) and no-op providers
# drop recordings, and unknown provider types (e.g. a custom provider
# delegating to a replay-safe one) cannot be classified, where a false
# positive is worse than a missed warning. The SDK logger provider is not
# checked because its class is only importable from the underscore
# namespace opentelemetry.sdk._logs while OpenTelemetry logs are pre-GA.
try:
from opentelemetry.sdk.metrics import MeterProvider as SdkMeterProvider
from opentelemetry.sdk.trace import TracerProvider as SdkTracerProvider
except ImportError:
# Without the opentelemetry-sdk package installed no SDK provider can
# exist, so there is nothing replay-unsafe to warn about.
return
stacklevel = _stacklevel_outside_temporalio()
if isinstance(opentelemetry.metrics.get_meter_provider(), SdkMeterProvider):
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 isinstance(opentelemetry.trace.get_tracer_provider(), SdkTracerProvider):
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,
)


def setup_deterministic_runtime():
"""Configures ADK runtime for Temporal determinism.

Expand Down Expand Up @@ -68,6 +135,10 @@ class GoogleAdkPlugin(SimplePlugin):
This plugin configures:
- Pydantic Payload Converter (required for ADK objects).
- Sandbox Passthrough for google.adk and google.genai modules.

At worker and replayer configuration time it also warns when the global
OpenTelemetry meter or tracer provider is not replay-safe, since ADK
telemetry recorded from workflow code would duplicate on replay.
"""

def __init__(
Expand Down Expand Up @@ -118,6 +189,22 @@ 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 meter or
tracer provider is 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 meter or
tracer provider is 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
56 changes: 56 additions & 0 deletions temporalio/contrib/opentelemetry/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -236,6 +236,62 @@ 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 opentelemetry.sdk._logs.export import BatchLogRecordProcessor
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.
logger_provider = LoggerProvider()
logger_provider.add_log_record_processor(BatchLogRecordProcessor(my_log_exporter))
opentelemetry._logs.set_logger_provider(ReplaySafeLoggerProvider(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
17 changes: 15 additions & 2 deletions temporalio/contrib/opentelemetry/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,21 +2,34 @@

This package provides OpenTelemetry tracing integration for Temporal workflows,
activities, and other operations. It includes automatic span creation and
propagation for distributed tracing.
propagation for distributed tracing. It also provides replay-safe wrappers for
the global OpenTelemetry tracer, meter, and logger providers.
"""

from temporalio.contrib.opentelemetry._interceptor import (
TracingInterceptor,
TracingWorkflowInboundInterceptor,
)
from temporalio.contrib.opentelemetry._logger_provider import (
ReplaySafeLoggerProvider,
)
from temporalio.contrib.opentelemetry._meter_provider import (
ReplaySafeMeterProvider,
)
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",
]
Loading
Loading