Google ADK Agents SDK Integration for Temporal
Release stage: Pre-release.
This package provides the integration layer between the Google ADK and Temporal. It allows ADK Agents to run reliably within Temporal Workflows by ensuring determinism and correctly routing external calls (network I/O) through Temporal Activities.
Install
uv add temporalio-google-adk
Benefits of Temporal to the ADK
Temporal provides a holistic, unified solution that centralizes your orchestration needs in one Workflow abstraction. Rather than cobbling together separate servers, task queues, gateways, and databases, you get:
- Recovering from crashes and stalls automatically, rather than manually managing sessions and resuming them. (Google offers Vertex Agent Engine, which still leaves resumption to the user). No need to set up a separate database
- Along with Retries and mechanisms for handling backpressure and rate limits.
- Support for ambient/long-running agent patterns via blocking awaits and worker versioning.
- Automatic execution state persistence, not just for agent interactions but for any custom automations in your workflows, without setting up a separate database.
- For Human-in-the-Loop patterns, an api gateway to scalably route incoming messages (such as user chats) to awaken the correct workflow on your worker pool.
- **Long-running tools support** using Activities — no need to set up and maintain microservices.
- Manage and debug your agent workflow execution and pinpoint problems using Temporal UI.
Benefits of the ADK to Temporal
ADK provides: (from the ADK overview):
- Improved Agent development velocity with a first-class Agentic abstraction and integration with LLMs and an ecosystem of tools.
- Improved agent robustness using built-in evals
- Build complex agents using its Multi-agent architecture.
- Safety and security, via guardrails and integrations with sandboxing solutions like Vertex Agent Runtime.
What's Included
Core ADK Integration
TemporalModel: Intercepts model calls and executes them as Temporal activitiesGoogleAdkPlugin: Worker plugin that configures runtime determinism and Pydantic serializationinvoke_model: Activity for executing LLM model calls with proper error handling
MCP (Model Context Protocol) Integration
TemporalMcpToolSet: Executes MCP tools as Temporal activitiesTemporalMcpToolSetProvider: Manages toolset creation and activity registration- Full support for tool confirmation and event actions within workflows
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()withworkflow.now()when in workflow context - Replaces
uuid.uuid4()withworkflow.uuid4()for deterministic IDs - Automatic setup when using
GoogleAdkPlugin
2. Activity-Based Model Execution
Model calls are intercepted and executed as Temporal activities with configurable:
- Timeouts (schedule-to-close, start-to-close, heartbeat)
- Retry policies
- Task queues
- Cancellation behavior
- Priority levels
3. Sandbox Compatibility
- Automatic passthrough for
google.adk,google.genai,mcp, andopentelemetrymodules - Works with both sandboxed and unsandboxed workflow runners
4. Advanced Serialization
- Pydantic payload converter for ADK objects
- Proper handling of complex ADK data types
- Maintains type safety across workflow boundaries
Usage
Basic Setup
Agent (Workflow) Side:
from temporalio.google_adk import TemporalModel
from temporalio.workflow import ActivityConfig
from google.adk import Agent
# Add to agent
agent = Agent(
name="test_agent",
model=TemporalModel("gemini-2.5-pro", activity_config=ActivityConfig(summary="Researcher Agent")),
)
Worker Side:
from temporalio.client import Client
from temporalio.worker import Worker
from temporalio.google_adk import GoogleAdkPlugin
client = await Client.connect(
"localhost:7233",
plugins=[
GoogleAdkPlugin(),
],
)
worker = Worker(
client,
task_queue="my-queue",
)
Advanced Features
With MCP Tools:
Install the optional client-side MCP dependency:
uv add "temporalio-google-adk[mcp]"
The MCP toolsets require MCP Python SDK v1. The base package does not require MCP and can be installed alongside MCP v2 integrations.
import os
from google.adk import Agent
from google.adk.tools.mcp_tool import McpToolset
from google.adk.tools.mcp_tool.mcp_session_manager import StdioConnectionParams
from mcp import StdioServerParameters
from temporalio.client import Client
from temporalio.worker import Worker
from temporalio.google_adk import (
GoogleAdkPlugin,
TemporalMcpToolSetProvider,
TemporalMcpToolSet,
)
def toolset_factory(_):
return McpToolset(
connection_params=StdioConnectionParams(
server_params=StdioServerParameters(
command="npx",
args=[
"-y",
"@modelcontextprotocol/server-filesystem",
os.path.dirname(os.path.abspath(__file__)),
],
),
),
)
# Use in agent workflow
agent = Agent(
name="test_agent",
model="gemini-2.5-pro",
tools=[
TemporalMcpToolSet(
"my-tools",
not_in_workflow_toolset=toolset_factory,
)
],
)
client = await Client.connect(
"localhost:7233",
plugins=[
GoogleAdkPlugin(
toolset_providers=[
TemporalMcpToolSetProvider("my-tools", toolset_factory),
],
),
],
)
# Configure worker
worker = Worker(
client,
task_queue="task-queue"
)
TemporalMcpToolSet also accepts an optional factory_argument. It is sent to the toolset activities and passed to the registered toolset_factory when the McpToolset is created.
Do not pass secrets, credentials, or API keys through factory_argument. It is an activity argument, so it is recorded in workflow history and, without a payload codec, visible in the web UI. Resolve credentials worker-side inside the toolset factory instead.
Reading Session State in Activity Tools
ADK's live ToolContext holds non-serializable objects, so it cannot be an
activity argument. To read the serializable subset from an activity-backed
tool, declare a parameter named tool_context annotated with
ToolContextSnapshot:
from datetime import timedelta
from temporalio import activity
from temporalio.google_adk.workflow import (
ToolContextSnapshot,
activity_as_tool,
)
@activity.defn
async def get_weather(query: str, tool_context: ToolContextSnapshot) -> dict:
db_url = tool_context.state.get("url", "")
...
weather_tool = activity_as_tool(
get_weather, start_to_close_timeout=timedelta(seconds=30)
)
Exactly like a native ADK function tool's tool_context parameter, it is
excluded from the LLM-facing tool schema; at invocation the wrapper snapshots
the live ToolContext (session state as a plain dict, plus the function-call
id) and passes it to the activity. Annotating any parameter with a live ADK
context type raises ValueError at wrap time, since ADK would inject the
non-serializable context into it regardless of its name.
When running under Temporal, the entire session state crosses the activity boundary: every value in it must be serializable by the configured data converter and the total size must fit within payload limits, even for keys the tool never reads. Local ADK runs pass the snapshot in memory and have no such constraint.
The snapshot is one-way and should be treated as read-only: mutations inside the activity do not propagate back to the session, because the activity may run on a different worker (and in local ADK runs, where the snapshot is passed in memory, nested state values may alias the live session state, so mutating them can corrupt the session). To modify session state, return the needed information from the activity and apply it in workflow-side code (for example an ADK callback or a plain tool function).
Local ADK Runs
The same agent definitions can also be exercised outside Temporal with
adk run or adk web.
TemporalModelandactivity_as_tool(...)work in local ADK runs without additional configuration.- If the agent uses
TemporalMcpToolSet, define a shared toolset factory, register it withTemporalMcpToolSetProvider(...)for workflow runs, and reuse the same function fornot_in_workflow_toolset=...so the agent can fall back to the underlyingMcpToolsetwhen it is not running insideworkflow.in_workflow().
Example:
# Reuse the same toolset_factory registered in GoogleAdkPlugin above.
agent = Agent(
name="test_agent",
model=TemporalModel("gemini-2.5-pro"),
tools=[
TemporalMcpToolSet(
"my-tools",
not_in_workflow_toolset=toolset_factory,
)
],
)
Graph Workflows (ADK v2)
ADK v2's graph runtime (google.adk.workflow) runs inside Temporal workflows:
the scheduler is pure asyncio and executes deterministically on Temporal's
workflow event loop, while LLM calls (TemporalModel), MCP tools
(TemporalMcpToolSet), and activity-backed nodes leave the workflow as
activities.
Use activity_node(...) to run a graph node as a Temporal activity. The
previous node's output is passed to the activity — directly for a
single-parameter activity, bound by name (from a dict) for multi-parameter
activities:
from google.adk.workflow import JoinNode, Workflow
from temporalio.google_adk.workflow import activity_node
fetch = activity_node(fetch_data, start_to_close_timeout=timedelta(seconds=30))
def summarize(node_input): # plain nodes run in-workflow: keep deterministic
return f"{node_input} summarized"
graph = Workflow(name="pipeline", edges=[("START", fetch, summarize)])
Conditional routing ((router, {"KEY": handler, ...}), DEFAULT_ROUTE),
parallel fan-out with JoinNode, and LlmAgent nodes (with
mode="task"/"single_turn") all work — agent nodes route their model calls
through TemporalModel as usual.
Dynamic Workflows
Dynamic nodes (await ctx.run_node(...) with loops, branches, and
asyncio.gather) work in-workflow; child-run caching reads only the
in-memory session, so re-entry after a HITL resume replays deterministically.
from google.adk.workflow import node
@node(rerun_on_resume=True)
async def pipeline(ctx):
data = await ctx.run_node(fetch, "query") # activity_node child
results = await asyncio.gather(
*(ctx.run_node(worker, item) for item in data) # parallel children
)
return results
On a HITL resume, a rerun_on_resume=True dynamic node re-executes its body
while completed children are skipped from the session cache. Place activity
invocations in child nodes (activity_node, activity_as_tool) rather than
inline in the dynamic node body, or make them idempotent — inline calls run
again on re-entry (ADK's documented at-least-once semantics).
A Workflow with an input_schema can also be passed in an agent's
tools=[...] list (Workflow-as-Tool), letting the model invoke whole graphs
as tools.
Durable Human-in-the-Loop
ADK pauses a run for human input (a node yielding RequestInput) or tool
confirmation (FunctionTool(..., require_confirmation=True)); in a Temporal
workflow that pause becomes a durable wait. The
pending_hitl_requests / hitl_input_response / hitl_confirmation_response
helpers cover the wire format; the wait itself is ordinary workflow code:
from temporalio.google_adk import (
HitlRequest,
hitl_input_response,
pending_hitl_requests,
)
@workflow.defn
class ApprovalWorkflow:
def __init__(self) -> None:
self._pending: dict[str, HitlRequest] = {}
self._responses: dict[str, Any] = {}
@workflow.query
def pending_requests(self) -> list[HitlRequest]:
return list(self._pending.values())
@workflow.update
def respond(self, interrupt_id: str, response: Any) -> None:
self._responses[interrupt_id] = response
@workflow.run
async def run(self, prompt: str) -> str:
runner = Runner(
app_name="app", node=graph, session_service=InMemorySessionService()
)
session = await runner.session_service.create_session(
app_name="app", user_id="user"
)
message = types.Content(role="user", parts=[types.Part(text=prompt)])
result = ""
while True:
async for event in runner.run_async(
user_id="user", session_id=session.id, new_message=message
):
for request in pending_hitl_requests(event):
self._pending[request.interrupt_id] = request
if event.content and event.content.parts and event.content.parts[0].text:
result = event.content.parts[0].text
if not self._pending:
return result
await workflow.wait_condition(
lambda: any(i in self._responses for i in self._pending)
)
parts = [
hitl_input_response(i, self._responses.pop(i))
for i in list(self._pending)
if i in self._responses
]
for part in parts:
self._pending.pop(part.function_response.id)
message = types.Content(role="user", parts=parts)
Tool confirmation composes with activity_as_tool with no extra plumbing —
FunctionTool(func=activity_as_tool(risky_activity, ...), require_confirmation=True)
never schedules the activity until the human approves (answer with
hitl_confirmation_response(interrupt_id, confirmed=True)). MCP tools
requesting confirmation via tool_context.request_confirmation(...) flow
through the same loop. Partial responses are fine: unanswered requests stay
pending across run_async turns.
Auth requests (adk_request_credential) are not covered by these helpers: ADK
exchanges the credential with network I/O inside the flow, and the exchanged
secret would be recorded in workflow history. Resolve credentials worker-side
instead (for example inside an activity or an MCP toolset factory).
Replay-safety note: HITL resume matches recorded human responses against 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
RequestInputinterrupt ids) derive fromworkflow.uuid4()and replay identically.
Determinism Notes
- The plugin patches ADK's
google.adk.platformtime, uuid, and random providers toworkflow.now(),workflow.uuid4(), andworkflow.random()inside workflows. - ADK node
timeout=/RetryConfigmap onto durable timers (asyncio.wait_for/asyncio.sleep). For activity-backed nodes, prefer Temporal activity timeouts andretry_policyviaactivity_node(...)options; an ADKRetryConfigon top would retry on top of Temporal's own activity retries, and an ADK node timeout cancels the in-flight activity. - Never set
RunConfig.tool_thread_pool_configinside a workflow — it runs tools on threads, which breaks workflow determinism. Live/BIDI mode is likewise unsupported in workflows. - ADK resume is at-least-once: on a HITL resume, completed nodes fast-forward
from the in-memory session, but
rerun_on_resume=Truenode bodies re-execute. This is deterministic under Temporal replay; schedule side effects through activities (retried/tracked by Temporal) or make them idempotent. - Very long HITL conversations grow the workflow history with each turn;
consider
continue-as-newboundaries betweenrun_asyncturns for long-running chats.
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 while history events are replaying (queries and update validators are live operations and still record):
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:
- Determinism: All non-deterministic operations are routed through Temporal
- Observability: Full tracing and activity visibility
- Reliability: Proper retry handling and error propagation
- Extensibility: Support for custom tools via MCP protocol
Metadata
Release files for temporalio-google-adk 0.0.1
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| temporalio_google_adk-0.0.1.tar.gz | 28.8 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| temporalio_google_adk-0.0.1-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 60.2 kB
Release files / temporalio_google_adk-0.0.1.tar.gz
| Download URL | temporalio_google_adk-0.0.1.tar.gz |
|---|---|
| Size | 28.8 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
faad4a0a210e0788859b3548f5f14794de65a498b6bcb56b7457ea6e40286e9f
|
|
BLAKE2b-256 checksum How to use checksums |
e676560481d0fdf61f3644a7eeac329915e72d13cd842d4e734b46172a74ee6c
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Oct 7, 2026.
Transparency logRelease files / temporalio_google_adk-0.0.1-py3-none-any.whl
| Download URL | temporalio_google_adk-0.0.1-py3-none-any.whl |
|---|---|
| Size | 31.3 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
0c82b463d00f3c53eeb25c71e22ea4277502e8f9d9862dcd29e9770eca3d8f8e
|
|
BLAKE2b-256 checksum How to use checksums |
afe6f19f7f07268e7537b5028557426ec050c0645baabb94d81ee26707955d88
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Oct 7, 2026.
Transparency log