Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
30 commits
Select commit Hold shift + click to select a range
2cd39b5
Apply Google ADK deterministic providers inside workflow tasks
DABH Sep 11, 2026
593912d
Address review on ADK provider install
DABH Sep 11, 2026
e977fd3
Survive ADK provider resets and give ADK a private random stream
DABH Sep 14, 2026
58dde24
Merge branch 'main' into fix/adk-providers-in-workflow-threads
DABH Sep 14, 2026
e22384e
Separate the plugin docstring's bullet list so pydoctor parses it
DABH Sep 14, 2026
4605343
Merge branch 'main' into fix/adk-providers-in-workflow-threads
DABH Sep 15, 2026
973f682
Fall back to entropy in read-only contexts and unify uuid-from-random
DABH Sep 15, 2026
f024398
Update the README's determinism bullets for the private-stream design
DABH Sep 15, 2026
646e303
Fall back to wall-clock time in read-only contexts
DABH Sep 15, 2026
b0ae875
Rename uuid4's generator parameter to rng
DABH Sep 16, 2026
1e89b41
Merge branch 'main' into fix/adk-providers-in-workflow-threads
DABH Sep 16, 2026
bae9c6e
Merge branch 'main' into fix/adk-providers-in-workflow-threads
DABH Sep 17, 2026
5be6770
Return fresh unseeded generators in the nondeterministic fallbacks
DABH Sep 17, 2026
1359ee4
Merge branch 'main' into fix/adk-providers-in-workflow-threads
DABH Sep 18, 2026
da8c932
Merge remote-tracking branch 'origin/main' into fix/adk-providers-in-…
DABH Sep 20, 2026
8392e61
clean up changelog
DABH Sep 20, 2026
7c3645b
Merge remote-tracking branch 'origin/main' into fix/adk-providers-in-…
DABH Oct 1, 2026
ba113f4
State what 1.34.0 did in read-only contexts in the changelog entry
DABH Oct 1, 2026
65caa28
clean up changelog
DABH Oct 1, 2026
14f29ca
Keep workflow time in read-only contexts and tighten the uuid4 rng docs
DABH Oct 1, 2026
cdc53b1
Merge branch 'main' into fix/adk-providers-in-workflow-threads
DABH Oct 1, 2026
3d13161
Build ADK ids from the plugin's stream instead of a public uuid4(rng=)
DABH Oct 1, 2026
e4533bd
Merge remote-tracking branch 'origin/main' into fix/adk-providers-in-…
DABH Oct 1, 2026
0387906
Describe the private stream in the ADK README
DABH Oct 1, 2026
75fb076
Merge main into ADK provider fix
brianstrauch Oct 1, 2026
295286d
Key the ADK stream by the run, not the workflow object, and keep read…
DABH Oct 2, 2026
84815a2
Merge remote-tracking branch 'origin/main' into fix/adk-providers-in-…
DABH Oct 2, 2026
0aabf5f
Name the ADK stream so its ids cannot coincide with workflow.uuid4()
DABH Oct 2, 2026
b022a6b
Spell out that an unnamed new_random() keeps its integer seed
DABH Oct 2, 2026
60cf7ba
Give every read-only context fresh entropy, as the OpenTelemetry id g…
DABH Oct 2, 2026
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
11 changes: 11 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,9 @@ to include examples, links to docs, or any other relevant information.

### Added

- `workflow.new_random()` accepts an optional `name` that is mixed into the seed, so differently
named generators, and `workflow.random()`, produce different sequences.

### Changed

- Payload converters exposed by data converters and workflow/activity accessors
Expand All @@ -30,6 +33,11 @@ to include examples, links to docs, or any other relevant information.

### :boom: Breaking Changes

- `temporalio.contrib.google_adk_agents`: ADK-generated ids and retry jitter now draw from a
workflow-private deterministic stream (a `workflow.new_random()` per run) instead of
`workflow.random()`. A workflow started under 1.34.0 that generated
ADK ids or jitter (for example one waiting on a HITL response) may not replay
deterministically across this upgrade; drain such workflows or use worker versioning.
- The OpenAI Agents integration has moved to the independently versioned
[`temporalio-openai-agents`](https://pypi.org/project/temporalio-openai-agents/)
package. The existing `temporalio[openai-agents]` extra now installs that
Expand All @@ -41,6 +49,9 @@ to include examples, links to docs, or any other relevant information.

### Fixed

- `GoogleAdkPlugin`'s deterministic providers now work in read-only contexts (query handlers,
update validators), returning the workflow's deterministic time and fresh entropy without
touching the workflow's random state.
- `contrib.google_adk_agents`: agents with an `output_schema` no longer fail every workflow task
when calling the model. The schema type is now sent to the model activity as its JSON schema.
Custom Pydantic schema generation is preserved.
Expand Down
13 changes: 7 additions & 6 deletions temporalio/contrib/google_adk_agents/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,13 +38,14 @@ ADK provides: (from the [ADK overview](https://google.github.io/adk-docs/#learn-
### OpenTelemetry Integration
- Automatic instrumentation for ADK components when exporters are provided
- Tracing integration that works within Temporal's execution context
- Support for custom span exporters

### Key Features

#### 1. Deterministic Runtime
- Replaces `time.time()` with `workflow.now()` when in workflow context
- Replaces `uuid.uuid4()` with `workflow.uuid4()` for deterministic IDs
- Installs ADK's `google.adk.platform` time, uuid, and random providers as process-wide defaults, so they apply inside workflow tasks (which run on worker threads with an empty `contextvars` context)
- Inside a workflow, time comes from `workflow.time()` and ids and randoms come from a workflow-private deterministic stream (a `workflow.new_random()` cached per run), so ADK-generated session, event, invocation, and function-call ids and retry jitter are reproducible on replay without shifting the sequences user code sees from `workflow.random()` and `workflow.uuid4()`. In read-only contexts (query handlers, update validators) time is still `workflow.time()` and ids and randoms come from fresh entropy that leaves the private stream untouched
- Outside a workflow in the same process (activities, client code) they fall back to the standard library
- Overrides through ADK's `set_*_provider` functions must be made after the Worker starts or from workflow code; one made earlier is replaced (with a warning) when the plugin installs its providers, and `reset_*_provider` restores the deterministic providers rather than the standard-library ones
- Automatic setup when using `GoogleAdkPlugin`

#### 2. Activity-Based Model Execution
Expand Down Expand Up @@ -385,13 +386,13 @@ instead (for example inside an activity or an MCP toolset factory).
> generated interrupt/function-call ids, so those ids must regenerate
> identically on replay. The plugin installs ADK's platform time/uuid/random
> providers as process-wide defaults, so the ids ADK generates (including
> default `RequestInput` interrupt ids) derive from `workflow.uuid4()` and
> replay identically.
> default `RequestInput` interrupt ids) derive from the workflow-private
> deterministic stream and replay identically.

## Determinism Notes

- The plugin patches ADK's `google.adk.platform` time, uuid, and random
providers to `workflow.now()`, `workflow.uuid4()`, and `workflow.random()`
providers to `workflow.time()` and a workflow-private deterministic stream
inside workflows.
- ADK node `timeout=`/`RetryConfig` map onto durable timers
(`asyncio.wait_for`/`asyncio.sleep`). For activity-backed nodes, prefer
Expand Down
182 changes: 143 additions & 39 deletions temporalio/contrib/google_adk_agents/_plugin.py
Original file line number Diff line number Diff line change
@@ -1,11 +1,14 @@
from __future__ import annotations

import contextvars
import dataclasses
import inspect
import random
import threading
import time
import uuid
import warnings
import weakref
from collections.abc import AsyncIterator, Callable
from contextlib import asynccontextmanager
from types import FrameType
Expand Down Expand Up @@ -37,14 +40,6 @@
from temporalio.worker.workflow_sandbox import SandboxedWorkflowRunner


def _install_provider(module: Any, var_name: str, provider: Callable[[], Any]) -> None:
"""Rebinds an ADK platform ContextVar so ``provider`` is its default in every context."""
from contextvars import ContextVar

context_var = getattr(module, var_name)
setattr(module, var_name, ContextVar(context_var.name, default=provider))


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
Expand Down Expand Up @@ -105,61 +100,165 @@ def _warn_if_global_otel_providers_not_replay_safe() -> None:
)


def setup_deterministic_runtime():
"""Configures ADK runtime for Temporal determinism.
def _deterministic_time_provider() -> float:
# workflow.time() in every in-workflow context, read-only ones included: a
# dynamic workflow's ``dynamic_config`` runs read-only and is replayed, so
# a wall-clock value there would not be replay-safe. In a query handler
# the value is the current activation's timestamp rather than the wall
# clock, which is harmless because nothing a query computes is persisted.
if workflow.in_workflow():
return workflow.time()
return time.time()

.. warning::
This function is experimental and may change in future versions.
Use with caution in production environments.

Installs Temporal-aware time, uuid, and random providers as the
process-wide defaults for ADK's ``google.adk.platform`` seams. Inside a
workflow they derive from ``workflow.now()`` / ``workflow.uuid4()`` /
``workflow.random()`` so replays are deterministic; outside a workflow
they fall back to the real primitives.
# Each run's private stream, keyed by the SDK's per-run runtime object: that
# exists during the workflow's __init__ (workflow.instance() does not yet) and
# leaves the user's class alone (it may use __slots__). Entries go away with
# the run.
_adk_randoms: weakref.WeakKeyDictionary[workflow._Runtime, random.Random] = (
weakref.WeakKeyDictionary()
)
_adk_randoms_lock = threading.Lock()


def _workflow_adk_random() -> random.Random:
# ADK draws from a private, named workflow.new_random() stream rather than
# sharing workflow.random(): ADK's draw count never shifts the sequence user
# code sees, and the name keeps the two from coinciding (an unnamed stream
# starts out identical to workflow.random(), so the Nth ADK id would equal
# the Nth workflow.uuid4()). Read-only code (query handlers, update
# validators) must not touch that stream, since a draw there would advance
# it and diverge later activations from replay; as in the opentelemetry id
# generator, it gets a fresh unseeded generator instead.
if workflow.unsafe.is_read_only():
return random.Random()
runtime = workflow._Runtime.current()
with _adk_randoms_lock:
rng = _adk_randoms.get(runtime)
if rng is None:
rng = workflow.new_random("temporalio.contrib.google_adk_agents")
_adk_randoms[runtime] = rng
return rng


def _uuid4_from(rng: random.Random) -> uuid.UUID:
# Same construction as workflow.uuid4(), drawn from the given stream.
return uuid.UUID(bytes=rng.getrandbits(16 * 8).to_bytes(16, "big"), version=4)


def _deterministic_id_provider() -> str:
if workflow.in_workflow():
return str(_uuid4_from(_workflow_adk_random()))
return str(uuid.uuid4())


def _deterministic_random_provider() -> random.Random:
# Outside a workflow, a fresh unseeded generator per call. ADK's
# set_random_provider docstring asks providers to return an existing
# instance so a seeded generator keeps its sequence across get_random()
# calls; an unseeded one draws fresh OS entropy either way, and ADK's only
# caller uses the result immediately (retry jitter).
if workflow.in_workflow():
return _workflow_adk_random()
return random.Random()


_install_provider_lock = threading.Lock()


def _install_provider(
module: Any, var_name: str, default_name: str, provider: Callable[[], Any]
) -> None:
"""Makes ``provider`` an ADK platform seam's default, everywhere.

ADK's ``set_*_provider`` functions set a value in the calling context only.
Workflow tasks run on worker threads, which start with an empty
contextvars context, so a value set from the worker's event loop never
reaches them and ADK falls back to its wall-clock and random defaults
there. A ContextVar's default, unlike a set value, is visible from every
context, so the module's variable is replaced with one that defaults to
``provider``. The module's ``_default_*`` binding is rebound too, because
``reset_*_provider`` restores that binding: without this, an override
followed by a reset would land on the standard-library provider rather
than back on ``provider``. ADK's ``set_*_provider`` and
``reset_*_provider`` operate on the new variable from then on; a value set
on the old one beforehand is orphaned, so it is warned about. A no-op when
``provider`` is already installed.
"""
current: contextvars.ContextVar[Callable[[], Any]] = getattr(module, var_name)
try:
import google.adk.platform._random
import google.adk.platform.time
import google.adk.platform.uuid
default = contextvars.Context().run(current.get)
except LookupError:
default = None
if default is provider and getattr(module, default_name) is provider:
return
if current.get(default) is not default:
warnings.warn(
f"Replacing the {module.__name__} provider set in this context before "
"GoogleAdkPlugin installed its deterministic providers; it will not "
"take effect. Set ADK provider overrides after the worker starts or "
"from workflow code.",
UserWarning,
stacklevel=_stacklevel_outside_temporalio(),
)
setattr(module, default_name, provider)
setattr(module, var_name, contextvars.ContextVar(current.name, default=provider))
Comment thread
DABH marked this conversation as resolved.

# Define safer, context-aware providers
def _deterministic_time_provider() -> float:
if workflow.in_workflow():
return workflow.now().timestamp()
return time.time()

def _deterministic_id_provider() -> str:
if workflow.in_workflow():
return str(workflow.uuid4())
return str(uuid.uuid4())
def setup_deterministic_runtime() -> None:
"""Installs Temporal's deterministic time, id, and random providers for ADK.

_local_random = random.Random()
.. warning::
This function is experimental and may change in future versions.
Use with caution in production environments.

def _deterministic_random_provider() -> random.Random:
if workflow.in_workflow():
return workflow.random()
return _local_random
The providers become the process-wide defaults of ADK's
``google.adk.platform`` time, uuid, and random seams, so they apply inside
workflow tasks (which run on worker threads with an empty contextvars
context) as well as in the calling context. Inside a workflow, time comes
from ``workflow.time()``, and ids and randoms come from a workflow-private
deterministic stream (a ``workflow.new_random()`` per run; ids are v4
UUIDs built from that stream), so ADK-generated ids and retry jitter are
reproducible on replay without shifting the sequence user code sees from
``workflow.random()`` and ``workflow.uuid4()``. In read-only contexts
(query handlers, update validators) time is still ``workflow.time()``,
while ids and randoms come from fresh entropy that leaves the private
stream untouched, since read-only results are never replayed. Outside a
workflow in the same process (activities, client code) they fall back to
``time.time()``,
``uuid.uuid4()``, and an unseeded ``random.Random()``.

Overrides through ADK's ``set_*_provider`` functions must be made after
this runs (after the worker starts, or from workflow code); one made
earlier is replaced, with a warning. ADK's ``reset_*_provider`` functions
restore these deterministic providers, not the standard-library ones.

:class:`GoogleAdkPlugin` calls this when a worker or replayer starts.
Calling it again is a no-op.
"""
import google.adk.platform._random
import google.adk.platform.time
import google.adk.platform.uuid

with _install_provider_lock:
_install_provider(
google.adk.platform.time,
"_time_provider_context_var",
"_default_time_provider",
_deterministic_time_provider,
)
_install_provider(
google.adk.platform.uuid,
"_id_provider_context_var",
"_default_id_provider",
_deterministic_id_provider,
)
_install_provider(
google.adk.platform._random,
"_random_provider_context_var",
"_default_random_provider",
_deterministic_random_provider,
)
except ImportError:
pass
except Exception as e:
print(f"Warning: Failed to set deterministic runtime providers: {e}")


class GoogleAdkPlugin(SimplePlugin):
Expand All @@ -170,8 +269,13 @@ class GoogleAdkPlugin(SimplePlugin):
Use with caution in production environments.

This plugin configures:

- Pydantic Payload Converter (required for ADK objects).
- Sandbox Passthrough for google.adk, google.genai, and OpenTelemetry modules.
- ADK's time, id, and random providers, so ADK-generated ids and retry
jitter come from the workflow's deterministic clock and a
workflow-private deterministic random stream
(see :func:`setup_deterministic_runtime`).

At worker and replayer configuration time it also warns when the global
OpenTelemetry meter or tracer provider is not replay-safe, since ADK
Expand Down
19 changes: 15 additions & 4 deletions temporalio/workflow/_context.py
Original file line number Diff line number Diff line change
Expand Up @@ -879,21 +879,32 @@ def register_random_seed_callback(callback: Callable[[int], None]) -> None:
return _Runtime.current().workflow_register_random_seed_callback(callback)


def new_random() -> Random:
def new_random(name: str | None = None) -> Random:
"""Create a Random instance that automatically reseeds when the workflow seed changes.

This creates a new Random instance that is initially seeded with the current
workflow seed, and automatically registers a callback to reseed itself
whenever the workflow receives a new seed from core.

Args:
name: Mixed into the seed when given, so differently named instances,
and :py:func:`random`, produce different sequences. Without it the
instance starts out identical to :py:func:`random`.

Returns:
A Random instance that stays synchronized with the workflow's randomness.
"""
current_seed = random_seed()
auto_random = Random(current_seed)

def seed_for(workflow_seed: int) -> int | str:
if name is None:
# Unchanged: the same integer seed as :py:func:`random`.
return workflow_seed
return f"{workflow_seed}:{name}"

auto_random = Random(seed_for(random_seed()))

def reseed_callback(new_seed: int) -> None:
auto_random.seed(new_seed)
auto_random.seed(seed_for(new_seed))

register_random_seed_callback(reseed_callback)
return auto_random
Expand Down
7 changes: 4 additions & 3 deletions tests/contrib/google_adk_agents/test_adk_graph_workflows.py
Original file line number Diff line number Diff line change
Expand Up @@ -279,8 +279,9 @@ class JitteredRetryGraphWorkflow:
"""A retried node with default-style jitter must replay deterministically.

Retry jitter feeds asyncio.sleep, i.e. a durable timer; unless the delay is
drawn from workflow.random() (via ADK's platform random seam), replays
compute a different timer duration and diverge.
drawn from the workflow's deterministic random stream (the plugin's
provider behind ADK's platform random seam), replays compute a different
timer duration and diverge.
"""

@workflow.run
Expand Down Expand Up @@ -504,7 +505,7 @@ async def test_graph_node_retry_jitter_replay_safe(client: Client):
assert result == "ok-after-2"
history = await handle.fetch_history()
# The jittered retry delay is a durable timer; replay must recompute the
# exact same duration from workflow.random().
# exact same duration from the plugin's deterministic random provider.
await Replayer(
workflows=[JitteredRetryGraphWorkflow], plugins=[GoogleAdkPlugin()]
).replay_workflow(history)
5 changes: 3 additions & 2 deletions tests/contrib/google_adk_agents/test_adk_hitl.py
Original file line number Diff line number Diff line change
Expand Up @@ -406,8 +406,9 @@ async def test_tool_confirmation_activity_as_tool(client: Client, confirmed: boo
# max_cached_workflows=0 forces a full history replay on every workflow
# task, proving the confirmation resume is replay-safe: the recorded human
# response references the confirmation function-call id, which must
# regenerate identically on replay (it derives from workflow.uuid4() via
# the platform uuid seam the plugin installs).
# regenerate identically on replay (it derives from the workflow's
# deterministic random stream via the platform uuid seam the plugin
# installs).
async with _worker(client):
LLMRegistry.register(ConfirmationModel)
handle = await client.start_workflow(
Expand Down
Loading
Loading