Skip to content

Commit 4692854

Browse files
committed
Merge remote-tracking branch 'origin/main' into fix-mcp-dependency-declarations
# Conflicts: # pyproject.toml # uv.lock
2 parents 3a817fc + 98b6d63 commit 4692854

23 files changed

Lines changed: 491 additions & 72 deletions

File tree

.github/CODEOWNERS

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,10 @@
11
* @temporalio/sdk
22

3-
# SDK & Nexus own the README, pyproject.toml, and uv.lock
3+
# SDK & Nexus teams share ownership of these files
44
/README.md @temporalio/sdk @temporalio/nexus
55
/pyproject.toml @temporalio/sdk @temporalio/nexus
66
/uv.lock @temporalio/sdk @temporalio/nexus
7+
/tests/conftest.py @temporalio/sdk @temporalio/nexus
78

89
# The Nexus team owns any folder whose name starts or ends with "nexus",
910
# both at the repo root and under tests/

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+

0 commit comments

Comments
 (0)