Skip to content

Commit 1381744

Browse files
RKestcopybara-github
authored andcommitted
refactor(telemetry): extract invocation span into _instrumentation.record_invocation
Co-authored-by: Max Ind <maxind@google.com> PiperOrigin-RevId: 939850648
1 parent 17d5f38 commit 1381744

2 files changed

Lines changed: 33 additions & 2 deletions

File tree

src/google/adk/runners.py

Lines changed: 12 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -56,16 +56,22 @@
5656
from .sessions.base_session_service import BaseSessionService
5757
from .sessions.base_session_service import GetSessionConfig
5858
from .sessions.session import Session
59+
from .telemetry import _instrumentation
5960
from .telemetry.tracing import tracer
6061
from .tools.base_toolset import BaseToolset
6162
from .utils._debug_output import print_event
6263

6364
if TYPE_CHECKING:
6465
from .apps.app import App
6566
from .apps.app import ResumabilityConfig
67+
from .workflow._base_node import BaseNode
6668

6769
logger = logging.getLogger('google_adk.' + __name__)
6870

71+
# Silence unused warning.
72+
# tracer is imported for backwards compatibility, to avoid breaking change in the API.
73+
_ = tracer
74+
6975

7076
def _find_active_task_isolation_scope(session) -> Optional[str]:
7177
"""Walk session backwards; find the active paused task agent's scope.
@@ -454,7 +460,9 @@ async def _run_node_async(
454460
Events flow through ic._event_queue via NodeRunner.
455461
"""
456462

457-
with tracer.start_as_current_span('invocation'):
463+
with _instrumentation.record_invocation(
464+
entrypoint_node=node or self.agent, conversation_id=session_id
465+
):
458466
# 1. Setup
459467
session = await self._get_or_create_session(
460468
user_id=user_id, session_id=session_id
@@ -1040,7 +1048,9 @@ async def _run_with_trace(
10401048
new_message: Optional[types.Content] = None,
10411049
invocation_id: Optional[str] = None,
10421050
) -> AsyncGenerator[Event, None]:
1043-
with tracer.start_as_current_span('invocation'):
1051+
with _instrumentation.record_invocation(
1052+
entrypoint_node=self.agent, conversation_id=session_id
1053+
):
10441054
session = await self._get_or_create_session(
10451055
user_id=user_id,
10461056
session_id=session_id,

src/google/adk/telemetry/_instrumentation.py

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
import sys
2121
import time
2222
from typing import AsyncIterator
23+
from typing import Iterator
2324
from typing import TYPE_CHECKING
2425

2526
from opentelemetry import trace
@@ -35,6 +36,7 @@
3536
from ..models.llm_request import LlmRequest
3637
from ..models.llm_response import LlmResponse
3738
from ..tools.base_tool import BaseTool
39+
from ..workflow._base_node import BaseNode
3840

3941
logger = logging.getLogger("google_adk." + __name__)
4042

@@ -69,6 +71,25 @@ def _get_elapsed_s(
6971
return time.monotonic() - fallback_start
7072

7173

74+
@contextlib.contextmanager
75+
def record_invocation(
76+
entrypoint_node: BaseNode | None,
77+
conversation_id: str,
78+
) -> Iterator[None]:
79+
"""Top-level ``invocation`` span for a runner invocation.
80+
81+
Args:
82+
entrypoint_node: The runner's root agent/node.
83+
conversation_id: Session/conversation id.
84+
85+
Yields:
86+
Nothing; the span is active for the duration of the block.
87+
"""
88+
del entrypoint_node, conversation_id # Unused until schema v2 lands.
89+
with tracing.tracer.start_as_current_span("invocation"):
90+
yield
91+
92+
7293
@dataclasses.dataclass
7394
class TelemetryContext:
7495
"""Stores all telemetry related state."""

0 commit comments

Comments
 (0)