diff --git a/google_adk_agents/metrics/README.md b/google_adk_agents/metrics/README.md new file mode 100644 index 00000000..14d92825 --- /dev/null +++ b/google_adk_agents/metrics/README.md @@ -0,0 +1,31 @@ +# Google ADK replay-safe metrics + +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. + +Start a local Temporal development server: + +```shell +temporal server start-dev +``` + +In another terminal, start the worker from the repository root: + +```shell +uv run --project google_adk_agents/metrics python -m google_adk_agents.metrics.run_worker +``` + +Then run the Workflow: + +```shell +uv run --project google_adk_agents/metrics python -m google_adk_agents.metrics.run_metrics_workflow +``` + +The starter prints `Replay-safe metrics are ready.` Inspect the metrics exposed by the worker: + +```shell +curl http://127.0.0.1:9464/metrics | rg 'gen_ai' +``` + +The output includes `gen_ai.invoke_agent` metrics, `gen_ai.client.operation.duration`, and `gen_ai.client.token.usage` from instrumentation scope `gcp.vertex.agent`. The worker sets `max_cached_workflows=0`, forcing replay between Workflow tasks. `ReplaySafeMeterProvider` drops observations made while replaying, so replay does not multiply the recorded counts. + +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. diff --git a/google_adk_agents/metrics/__init__.py b/google_adk_agents/metrics/__init__.py new file mode 100644 index 00000000..8b137891 --- /dev/null +++ b/google_adk_agents/metrics/__init__.py @@ -0,0 +1 @@ + diff --git a/google_adk_agents/metrics/models/__init__.py b/google_adk_agents/metrics/models/__init__.py new file mode 100644 index 00000000..8b137891 --- /dev/null +++ b/google_adk_agents/metrics/models/__init__.py @@ -0,0 +1 @@ + diff --git a/google_adk_agents/metrics/models/local_metrics_model.py b/google_adk_agents/metrics/models/local_metrics_model.py new file mode 100644 index 00000000..01e8d010 --- /dev/null +++ b/google_adk_agents/metrics/models/local_metrics_model.py @@ -0,0 +1,29 @@ +from collections.abc import AsyncGenerator + +from google.adk.models import BaseLlm +from google.adk.models.llm_request import LlmRequest +from google.adk.models.llm_response import LlmResponse +from google.genai import types + +MODEL_NAME = "local-metrics-model" + + +class LocalMetricsModel(BaseLlm): + @classmethod + def supported_models(cls) -> list[str]: + return [MODEL_NAME] + + async def generate_content_async( + self, llm_request: LlmRequest, stream: bool = False + ) -> AsyncGenerator[LlmResponse, None]: + yield LlmResponse( + content=types.Content( + role="model", + parts=[types.Part(text="Replay-safe metrics are ready.")], + ), + usage_metadata=types.GenerateContentResponseUsageMetadata( + prompt_token_count=8, + candidates_token_count=5, + total_token_count=13, + ), + ) diff --git a/google_adk_agents/metrics/pyproject.toml b/google_adk_agents/metrics/pyproject.toml new file mode 100644 index 00000000..2aa1dc0b --- /dev/null +++ b/google_adk_agents/metrics/pyproject.toml @@ -0,0 +1,25 @@ +[project] +name = "google-adk-agents-metrics-sample" +version = "0.1.0" +requires-python = ">=3.10" +dependencies = [ + "google-adk>=2.2.0,<3", + "opentelemetry-exporter-prometheus>=0.48b0", + "temporalio[google-adk,opentelemetry] @ git+https://github.com/temporalio/sdk-python.git@78159d7735b2defc6493669f6c14a0fae6eab985", +] + +[dependency-groups] +dev = [ + "pytest>=7.1.2,<9", + "pytest-asyncio>=0.23,<2", + "ruff>=0.5.0,<0.6", +] + +[tool.pytest.ini_options] +asyncio_mode = "auto" + +[tool.ruff] +target-version = "py310" + +[tool.uv] +package = false diff --git a/google_adk_agents/metrics/run_metrics_workflow.py b/google_adk_agents/metrics/run_metrics_workflow.py new file mode 100644 index 00000000..63c8d810 --- /dev/null +++ b/google_adk_agents/metrics/run_metrics_workflow.py @@ -0,0 +1,21 @@ +import asyncio + +from temporalio.client import Client +from temporalio.contrib.google_adk_agents import GoogleAdkPlugin + +from google_adk_agents.metrics.workflows.metrics_workflow import MetricsWorkflow + + +async def main() -> None: + client = await Client.connect("localhost:7233", plugins=[GoogleAdkPlugin()]) + result = await client.execute_workflow( + MetricsWorkflow.run, + "Explain replay-safe metrics.", + id="google-adk-agents-metrics-workflow-id", + task_queue="google-adk-agents-metrics", + ) + print(f"Result: {result}") + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/google_adk_agents/metrics/run_worker.py b/google_adk_agents/metrics/run_worker.py new file mode 100644 index 00000000..7405ac14 --- /dev/null +++ b/google_adk_agents/metrics/run_worker.py @@ -0,0 +1,34 @@ +import asyncio + +from opentelemetry.exporter.prometheus import PrometheusMetricReader + +from google_adk_agents.metrics.telemetry import install_meter_provider + + +async def main() -> None: + install_meter_provider(PrometheusMetricReader()) + + from google.adk.models import LLMRegistry + from prometheus_client import start_http_server + from temporalio.client import Client + from temporalio.contrib.google_adk_agents import GoogleAdkPlugin + from temporalio.worker import Worker + + from google_adk_agents.metrics.models.local_metrics_model import LocalMetricsModel + from google_adk_agents.metrics.workflows.metrics_workflow import MetricsWorkflow + + LLMRegistry.register(LocalMetricsModel) + start_http_server(port=9464, addr="127.0.0.1") + plugin = GoogleAdkPlugin() + client = await Client.connect("localhost:7233", plugins=[plugin]) + worker = Worker( + client, + task_queue="google-adk-agents-metrics", + workflows=[MetricsWorkflow], + max_cached_workflows=0, + ) + await worker.run() + + +if __name__ == "__main__": + asyncio.run(main()) diff --git a/google_adk_agents/metrics/telemetry.py b/google_adk_agents/metrics/telemetry.py new file mode 100644 index 00000000..e69f6f46 --- /dev/null +++ b/google_adk_agents/metrics/telemetry.py @@ -0,0 +1,12 @@ +import opentelemetry.metrics +from opentelemetry.sdk.metrics import MeterProvider +from opentelemetry.sdk.metrics.export import MetricReader +from temporalio.contrib.opentelemetry import ReplaySafeMeterProvider + + +def install_meter_provider(reader: MetricReader) -> ReplaySafeMeterProvider: + provider = ReplaySafeMeterProvider(MeterProvider(metric_readers=[reader])) + opentelemetry.metrics.set_meter_provider(provider) + if opentelemetry.metrics.get_meter_provider() is not provider: + raise RuntimeError("The global OpenTelemetry meter provider is already set") + return provider diff --git a/google_adk_agents/metrics/test_metrics.py b/google_adk_agents/metrics/test_metrics.py new file mode 100644 index 00000000..f1be6bd0 --- /dev/null +++ b/google_adk_agents/metrics/test_metrics.py @@ -0,0 +1,68 @@ +import uuid + +from opentelemetry.sdk.metrics.export import InMemoryMetricReader + +from google_adk_agents.metrics.telemetry import install_meter_provider + +ADK_METER_SCOPE = "gcp.vertex.agent" + + +async def test_metrics_are_not_inflated_by_replay() -> None: + reader = InMemoryMetricReader() + install_meter_provider(reader) + + from google.adk.models import LLMRegistry + from temporalio.contrib.google_adk_agents import GoogleAdkPlugin + from temporalio.testing import WorkflowEnvironment + from temporalio.worker import Replayer, Worker + + from google_adk_agents.metrics.models.local_metrics_model import LocalMetricsModel + from google_adk_agents.metrics.workflows.metrics_workflow import MetricsWorkflow + + LLMRegistry.register(LocalMetricsModel) + async with await WorkflowEnvironment.start_time_skipping() as environment: + plugin = GoogleAdkPlugin() + config = environment.client.config() + config["plugins"] = [*config["plugins"], plugin] + client = type(environment.client)(**config) + task_queue = f"google-adk-agents-metrics-{uuid.uuid4()}" + async with Worker( + client, + task_queue=task_queue, + workflows=[MetricsWorkflow], + ): + handle = await client.start_workflow( + MetricsWorkflow.run, + "Explain replay-safe metrics.", + id=f"google-adk-agents-metrics-{uuid.uuid4()}", + task_queue=task_queue, + ) + result = await handle.result() + + assert result == "Replay-safe metrics are ready." + counts_before_replay = metric_counts(reader) + assert counts_before_replay["gen_ai.invoke_agent.duration"] > 0 + assert counts_before_replay["gen_ai.invoke_agent.inference_calls"] > 0 + assert counts_before_replay["gen_ai.client.operation.duration"] > 0 + assert counts_before_replay["gen_ai.client.token.usage"] > 0 + + await Replayer(workflows=[MetricsWorkflow], plugins=[plugin]).replay_workflow( + await handle.fetch_history() + ) + + assert metric_counts(reader) == counts_before_replay + + +def metric_counts(reader: InMemoryMetricReader) -> dict[str, int]: + counts: dict[str, int] = {} + data = reader.get_metrics_data() + if data is not None: + for resource_metrics in data.resource_metrics: + for scope_metrics in resource_metrics.scope_metrics: + if scope_metrics.scope.name != ADK_METER_SCOPE: + continue + for metric in scope_metrics.metrics: + counts[metric.name] = sum( + getattr(point, "count", 1) for point in metric.data.data_points + ) + return counts diff --git a/google_adk_agents/metrics/workflows/__init__.py b/google_adk_agents/metrics/workflows/__init__.py new file mode 100644 index 00000000..8b137891 --- /dev/null +++ b/google_adk_agents/metrics/workflows/__init__.py @@ -0,0 +1 @@ + diff --git a/google_adk_agents/metrics/workflows/metrics_workflow.py b/google_adk_agents/metrics/workflows/metrics_workflow.py new file mode 100644 index 00000000..1cd7eb5a --- /dev/null +++ b/google_adk_agents/metrics/workflows/metrics_workflow.py @@ -0,0 +1,39 @@ +from datetime import timedelta + +from google.adk import Agent +from google.adk.runners import InMemoryRunner +from google.adk.utils.context_utils import Aclosing +from google.genai import types +from temporalio import workflow +from temporalio.contrib.google_adk_agents import TemporalModel + + +@workflow.defn +class MetricsWorkflow: + @workflow.run + async def run(self, prompt: str) -> str: + agent = Agent( + name="metrics_agent", + model=TemporalModel("local-metrics-model"), + instruction="Answer the user briefly.", + ) + runner = InMemoryRunner(agent=agent, app_name="metrics_app") + session = await runner.session_service.create_session( + app_name="metrics_app", user_id="sample-user" + ) + + final_text = "" + async with Aclosing( + runner.run_async( + user_id="sample-user", + session_id=session.id, + new_message=types.Content(role="user", parts=[types.Part(text=prompt)]), + ) + ) as events: + async for event in events: + if event.content and event.content.parts: + for part in event.content.parts: + if part.text: + final_text = part.text + await workflow.sleep(timedelta(milliseconds=1)) + return final_text