Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
35 changes: 35 additions & 0 deletions sdks/python/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,3 +24,38 @@ async def search_database(query: str) -> str:
```

Controls are defined centrally and enforced automatically at runtime. See [Python SDK Documentation](https://docs.agentcontrol.dev/sdk/python-sdk) for complete reference.

## Sharing an OpenTelemetry provider with Google ADK

When Google ADK and Agent Control should export through the same OpenTelemetry
pipeline, configure one SDK `TracerProvider` and pass it to both components:

```python
import agent_control
from opentelemetry import trace
from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter
from openinference.instrumentation.google_adk import GoogleADKInstrumentor
from opentelemetry.sdk.resources import Resource
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor

resource = Resource.create({"service.name": "my-google-adk-agent"})
provider = TracerProvider(resource=resource)
provider.add_span_processor(BatchSpanProcessor(OTLPSpanExporter()))
trace.set_tracer_provider(provider)

GoogleADKInstrumentor().instrument(tracer_provider=provider)
agent_control.init(
agent_name="my-google-adk-agent",
observability_enabled=True,
observability_sink_name="otel",
observability_sink_config={"enabled": True},
otel_tracer_provider=provider,
)
```

Agent Control reuses this provider without adding an exporter or span processor.
It force-flushes the provider during shutdown but leaves provider shutdown to the
application. If no provider is passed, the OTEL sink first checks for a globally
registered SDK provider, then falls back to its existing Agent Control-owned OTLP
pipeline.
13 changes: 12 additions & 1 deletion sdks/python/src/agent_control/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,8 @@ async def handle_input(user_message: str) -> str:
)
"""

from __future__ import annotations

from importlib.metadata import PackageNotFoundError, version

try:
Expand All @@ -51,7 +53,7 @@ async def handle_input(user_message: str) -> str:
from collections.abc import Callable, Coroutine
from dataclasses import dataclass
from datetime import UTC, datetime
from typing import Any, Literal, TypeVar
from typing import TYPE_CHECKING, Any, Literal, TypeVar

import httpx
from agent_control_models import (
Expand Down Expand Up @@ -123,6 +125,9 @@ async def handle_input(user_message: str) -> str:
)
from .validation import ensure_agent_name

if TYPE_CHECKING:
from opentelemetry.sdk.trace import TracerProvider

# Module logger
logger = get_logger(__name__)

Expand Down Expand Up @@ -460,6 +465,7 @@ def init(
target_type: str | None = None,
target_id: str | None = None,
runtime_token_header: str | None = None,
otel_tracer_provider: TracerProvider | None = None,
**kwargs: object
) -> Agent:
"""
Expand Down Expand Up @@ -514,6 +520,10 @@ def init(
X-Agent-Control-Runtime-Token) when the server runs behind a gateway
that reserves Authorization for its own identity JWT. The server must
be configured to read the same header.
otel_tracer_provider: Optional application-owned OpenTelemetry SDK provider.
Used only when observability_sink_name is ``"otel"``. If omitted,
the OTEL sink reuses a globally registered SDK provider when available,
otherwise it creates and owns a provider as before.
**kwargs: Additional metadata to store with the agent

Returns:
Expand Down Expand Up @@ -738,6 +748,7 @@ def run_in_thread() -> None:
enabled=observability_enabled,
sink_name=observability_sink_name,
sink_config=observability_sink_config,
otel_tracer_provider=otel_tracer_provider,
)
if batcher:
logger.info("Observability enabled")
Expand Down
63 changes: 59 additions & 4 deletions sdks/python/src/agent_control/observability.py
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,7 @@

if TYPE_CHECKING:
from agent_control_models import ControlExecutionEvent
from opentelemetry.sdk.trace import TracerProvider

from .otel_sink import (
OTEL_CONTROL_EVENT_SINK_NAME,
Expand Down Expand Up @@ -833,6 +834,8 @@ def get_stats(self) -> dict:
)
_configured_named_event_sink: ControlEventSink | None = None
_configured_named_event_sink_selection: ControlEventSinkSelection | None = None
_configured_otel_tracer_provider: TracerProvider | None = None
_configured_named_event_sink_provider: TracerProvider | None = None
_configured_named_event_sink_lock = threading.Lock()
_used_custom_event_sinks: list[ControlEventSink] = []
_used_custom_event_sinks_lock = threading.Lock()
Expand All @@ -842,7 +845,15 @@ def _register_builtin_control_event_sink_factories() -> None:
"""Ensure built-in named sink factories are available."""
_named_event_sink_factories.register(
OTEL_CONTROL_EVENT_SINK_NAME,
create_otel_control_event_sink,
_create_configured_otel_control_event_sink,
)


def _create_configured_otel_control_event_sink(config: JSONObject) -> BaseControlEventSink:
"""Create the OTEL sink with its runtime-only provider configuration."""
return create_otel_control_event_sink(
config,
tracer_provider=_configured_otel_tracer_provider,
)


Expand Down Expand Up @@ -976,14 +987,21 @@ def _get_or_create_named_control_event_sink(
selection: ControlEventSinkSelection,
) -> ControlEventSink | None:
"""Resolve and cache a named configured sink."""
global _configured_named_event_sink, _configured_named_event_sink_selection
global _configured_named_event_sink, _configured_named_event_sink_provider
global _configured_named_event_sink_selection

previous_sink: ControlEventSink | None = None
selected_provider = (
_configured_otel_tracer_provider
if selection.name == OTEL_CONTROL_EVENT_SINK_NAME
else None
)

with _configured_named_event_sink_lock:
if (
_configured_named_event_sink is not None
and _configured_named_event_sink_selection == selection
and _configured_named_event_sink_provider is selected_provider
):
return _configured_named_event_sink

Expand All @@ -1004,6 +1022,7 @@ def _get_or_create_named_control_event_sink(
previous_sink = _configured_named_event_sink
_configured_named_event_sink = sink
_configured_named_event_sink_selection = selection.model_copy(deep=True)
_configured_named_event_sink_provider = selected_provider

if previous_sink is not None:
_shutdown_custom_control_event_sink(previous_sink)
Expand Down Expand Up @@ -1051,18 +1070,41 @@ def _shutdown_built_in_event_sink() -> None:

def _shutdown_configured_named_event_sink() -> None:
"""Stop and clear the cached configured named sink if it is active."""
global _configured_named_event_sink, _configured_named_event_sink_selection
global _configured_named_event_sink, _configured_named_event_sink_provider
global _configured_named_event_sink_selection

configured_named_sink: ControlEventSink | None = None
with _configured_named_event_sink_lock:
configured_named_sink = _configured_named_event_sink
_configured_named_event_sink = None
_configured_named_event_sink_selection = None
_configured_named_event_sink_provider = None

if configured_named_sink is not None:
_shutdown_custom_control_event_sink(configured_named_sink)


def _invalidate_cached_otel_sink() -> None:
"""Clear and shut down only the cached built-in OTEL sink, if present."""
global _configured_named_event_sink, _configured_named_event_sink_provider
global _configured_named_event_sink_selection

configured_otel_sink: ControlEventSink | None = None
with _configured_named_event_sink_lock:
if (
_configured_named_event_sink_selection is None
or _configured_named_event_sink_selection.name != OTEL_CONTROL_EVENT_SINK_NAME
):
return
configured_otel_sink = _configured_named_event_sink
_configured_named_event_sink = None
_configured_named_event_sink_selection = None
_configured_named_event_sink_provider = None

if configured_otel_sink is not None:
_shutdown_custom_control_event_sink(configured_otel_sink)


def _shutdown_custom_control_event_sink(sink: ControlEventSink) -> None:
"""Flush and close a custom sink when it exposes lifecycle hooks."""
flush = getattr(sink, "flush", None)
Expand Down Expand Up @@ -1109,6 +1151,7 @@ def init_observability(
enabled: bool | None = None,
sink_name: str | None = None,
sink_config: JSONObject | None = None,
otel_tracer_provider: TracerProvider | None = None,
) -> EventBatcher | None:
"""
Initialize observability system.
Expand All @@ -1122,11 +1165,12 @@ def init_observability(
enabled: Override AGENT_CONTROL_OBSERVABILITY_ENABLED
sink_name: Override AGENT_CONTROL_OBSERVABILITY_SINK_NAME
sink_config: Override AGENT_CONTROL_OBSERVABILITY_SINK_CONFIG
otel_tracer_provider: Runtime-only provider for the built-in OTEL sink

Returns:
EventBatcher instance if enabled, None otherwise
"""
global _batcher, _event_sink
global _batcher, _configured_otel_tracer_provider, _event_sink

settings_updates: dict[str, object] = {}
current_settings = get_settings()
Expand All @@ -1144,6 +1188,17 @@ def init_observability(
if settings_updates:
configure_settings(**settings_updates)

selected_sink_name = get_settings().observability_sink_name
selected_otel_tracer_provider = (
otel_tracer_provider
if selected_sink_name == OTEL_CONTROL_EVENT_SINK_NAME
else None
)
provider_changed = selected_otel_tracer_provider is not _configured_otel_tracer_provider
_configured_otel_tracer_provider = selected_otel_tracer_provider
if provider_changed:
_invalidate_cached_otel_sink()

is_enabled = get_settings().observability_enabled

if not is_enabled:
Expand Down
Loading
Loading