Skip to content

Commit 01a7e2b

Browse files
authored
AI-472 Add replay-safe Google ADK metrics sample (#355)
* AI-472 Add replay-safe Google ADK metrics sample * AI-472 Avoid duplicate Google ADK worker plugin registration * AI-472 Make replay metrics test deterministic * AI-472 Address Google ADK metrics review
1 parent 4b7f573 commit 01a7e2b

13 files changed

Lines changed: 319 additions & 31 deletions

File tree

google_adk_agents/README.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,7 @@ Each directory contains a complete example with its own README:
4747
| [agent_patterns](./agent_patterns/README.md) | A coordinator `LlmAgent` with `sub_agents`, each a `TemporalModel` with a per-agent activity summary. |
4848
| [mcp](./mcp/README.md) | A local echo MCP toolset via `TemporalMcpToolSet` / `TemporalMcpToolSetProvider`, running MCP tools as activities. Self-contained, no Node required. |
4949
| [streaming](./streaming/README.md) | Token streaming via `TemporalModel(streaming_topic=...)` + `WorkflowStream`, consumed by a starter with `WorkflowStreamClient`. |
50+
| [metrics](./metrics/README.md) | Google ADK OpenTelemetry metrics exported to a local Prometheus endpoint, with replay suppression through `ReplaySafeMeterProvider`. |
5051

5152
To run any scenario, start its worker in one terminal and its workflow starter
5253
in another:
Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,33 @@
1+
# Google ADK replay-safe metrics
2+
3+
This sample exports Google ADK's OpenTelemetry metrics to a local Prometheus endpoint while preventing Workflow replay from recording the same observations again. The default scripted model is deterministic and makes no network model calls, so no API key is needed.
4+
5+
Start a local Temporal development server:
6+
7+
```shell
8+
temporal server start-dev
9+
```
10+
11+
In another terminal, start the worker from the repository root:
12+
13+
```shell
14+
uv run python -m google_adk_agents.metrics.run_worker
15+
```
16+
17+
Then run the Workflow:
18+
19+
```shell
20+
uv run python -m google_adk_agents.metrics.run_metrics_workflow
21+
```
22+
23+
The starter prints `Replay-safe metrics are ready.` Inspect the metrics exposed by the worker:
24+
25+
```shell
26+
curl -s http://127.0.0.1:9464/metrics | grep gen_ai
27+
```
28+
29+
The output includes `gen_ai.invoke_agent`, `gen_ai.client.operation.duration`, and `gen_ai.client.token.usage` metrics. Prometheus replaces dots with underscores, so an exported line looks like `gen_ai_invoke_agent_duration_seconds_count{gen_ai_agent_name="metrics_agent"} 1.0`. `ReplaySafeMeterProvider` drops observations made while replaying, so replay does not multiply the recorded counts.
30+
31+
Recordings are first-execution-only rather than exactly-once. Replay is suppressed, but a Workflow Task retry re-executes live and can record again, so treat these metrics as at-least-once usage signals.
32+
33+
OpenTelemetry's global meter provider can be installed only once per process. `run_worker.py` installs the replay-safe provider before importing Google ADK or the Workflow. Applications embedding this setup must likewise make it the first and only global meter provider installation in that process.
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,33 @@
1+
from collections.abc import AsyncGenerator
2+
3+
from google.adk.models import BaseLlm
4+
from google.adk.models.llm_request import LlmRequest
5+
from google.adk.models.llm_response import LlmResponse
6+
from google.genai import types
7+
8+
MODEL_NAME = "local-metrics-model"
9+
10+
11+
class LocalMetricsModel(BaseLlm):
12+
@classmethod
13+
def supported_models(cls) -> list[str]:
14+
return [MODEL_NAME]
15+
16+
async def generate_content_async(
17+
self, llm_request: LlmRequest, stream: bool = False
18+
) -> AsyncGenerator[LlmResponse, None]:
19+
if stream:
20+
raise NotImplementedError(
21+
"LocalMetricsModel does not implement streaming responses."
22+
)
23+
yield LlmResponse(
24+
content=types.Content(
25+
role="model",
26+
parts=[types.Part(text="Replay-safe metrics are ready.")],
27+
),
28+
usage_metadata=types.GenerateContentResponseUsageMetadata(
29+
prompt_token_count=8,
30+
candidates_token_count=5,
31+
total_token_count=13,
32+
),
33+
)
Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,21 @@
1+
import asyncio
2+
3+
from temporalio.client import Client
4+
from temporalio.contrib.google_adk_agents import GoogleAdkPlugin
5+
6+
from google_adk_agents.metrics.workflows.metrics_workflow import MetricsWorkflow
7+
8+
9+
async def main() -> None:
10+
client = await Client.connect("localhost:7233", plugins=[GoogleAdkPlugin()])
11+
result = await client.execute_workflow(
12+
MetricsWorkflow.run,
13+
"Explain replay-safe metrics.",
14+
id="google-adk-agents-metrics-workflow-id",
15+
task_queue="google-adk-agents-metrics",
16+
)
17+
print(f"Result: {result}")
18+
19+
20+
if __name__ == "__main__":
21+
asyncio.run(main())
Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,33 @@
1+
import asyncio
2+
3+
from opentelemetry.exporter.prometheus import PrometheusMetricReader
4+
5+
from google_adk_agents.metrics.telemetry import install_meter_provider
6+
7+
8+
async def main() -> None:
9+
install_meter_provider(PrometheusMetricReader())
10+
11+
from google.adk.models import LLMRegistry
12+
from prometheus_client import start_http_server
13+
from temporalio.client import Client
14+
from temporalio.contrib.google_adk_agents import GoogleAdkPlugin
15+
from temporalio.worker import Worker
16+
17+
from google_adk_agents.metrics.models.local_metrics_model import LocalMetricsModel
18+
from google_adk_agents.metrics.workflows.metrics_workflow import MetricsWorkflow
19+
20+
LLMRegistry.register(LocalMetricsModel)
21+
start_http_server(port=9464, addr="127.0.0.1")
22+
plugin = GoogleAdkPlugin()
23+
client = await Client.connect("localhost:7233", plugins=[plugin])
24+
worker = Worker(
25+
client,
26+
task_queue="google-adk-agents-metrics",
27+
workflows=[MetricsWorkflow],
28+
)
29+
await worker.run()
30+
31+
32+
if __name__ == "__main__":
33+
asyncio.run(main())
Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,12 @@
1+
import opentelemetry.metrics
2+
from opentelemetry.sdk.metrics import MeterProvider
3+
from opentelemetry.sdk.metrics.export import MetricReader
4+
from temporalio.contrib.opentelemetry import ReplaySafeMeterProvider
5+
6+
7+
def install_meter_provider(reader: MetricReader) -> ReplaySafeMeterProvider:
8+
provider = ReplaySafeMeterProvider(MeterProvider(metric_readers=[reader]))
9+
opentelemetry.metrics.set_meter_provider(provider)
10+
if opentelemetry.metrics.get_meter_provider() is not provider:
11+
raise RuntimeError("The global OpenTelemetry meter provider is already set")
12+
return provider
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,38 @@
1+
from google.adk import Agent
2+
from google.adk.runners import InMemoryRunner
3+
from google.adk.utils.context_utils import Aclosing
4+
from google.genai import types
5+
from temporalio import workflow
6+
from temporalio.contrib.google_adk_agents import TemporalModel
7+
8+
from google_adk_agents.metrics.models.local_metrics_model import MODEL_NAME
9+
10+
11+
@workflow.defn
12+
class MetricsWorkflow:
13+
@workflow.run
14+
async def run(self, prompt: str) -> str:
15+
agent = Agent(
16+
name="metrics_agent",
17+
model=TemporalModel(MODEL_NAME),
18+
instruction="Answer the user briefly.",
19+
)
20+
runner = InMemoryRunner(agent=agent, app_name="metrics_app")
21+
session = await runner.session_service.create_session(
22+
app_name="metrics_app", user_id="sample-user"
23+
)
24+
25+
final_text = ""
26+
async with Aclosing(
27+
runner.run_async(
28+
user_id="sample-user",
29+
session_id=session.id,
30+
new_message=types.Content(role="user", parts=[types.Part(text=prompt)]),
31+
)
32+
) as events:
33+
async for event in events:
34+
if event.content and event.content.parts:
35+
for part in event.content.parts:
36+
if part.text:
37+
final_text = part.text
38+
return final_text

0 commit comments

Comments
 (0)