Source code for pipecat.services.google.llm

#
# Copyright (c) 2024-2026, Daily
#
# SPDX-License-Identifier: BSD 2-Clause License
#

"""Google Gemini integration for Pipecat.

This module provides Google Gemini integration for the Pipecat framework,
including LLM services, context management, and message aggregation.
"""

import asyncio
import io
import os
import re
import uuid
from collections.abc import AsyncGenerator, AsyncIterator
from contextlib import aclosing
from dataclasses import dataclass, field
from typing import Any, Literal, Optional, Union, cast

from loguru import logger
from PIL import Image
from pydantic import BaseModel, Field

from pipecat.adapters.base_llm_adapter import LLMContextConversionError
from pipecat.adapters.services.gemini_adapter import GeminiLLMAdapter
from pipecat.frames.frames import (
    AssistantImageRawFrame,
    Frame,
    LLMContextFrame,
    LLMFullResponseEndFrame,
    LLMFullResponseStartFrame,
    LLMMessagesAppendFrame,
    LLMThoughtEndFrame,
    LLMThoughtStartFrame,
    LLMThoughtTextFrame,
)
from pipecat.metrics.metrics import LLMTokenUsage
from pipecat.processors.aggregators.llm_context import LLMContext
from pipecat.processors.frame_processor import FrameDirection
from pipecat.services.google.frames import LLMSearchResponseFrame
from pipecat.services.google.utils import update_google_client_http_options
from pipecat.services.llm_service import FunctionCallFromLLM, LLMService
from pipecat.services.settings import LLMSettings
from pipecat.utils.deprecation import deprecated
from pipecat.utils.tracing.service_decorators import traced_llm
from pipecat.utils.types import NOT_GIVEN, NotGiven, assert_given, is_given

# Suppress gRPC fork warnings
os.environ["GRPC_ENABLE_FORK_SUPPORT"] = "false"

try:
    import google.genai as genai
    from google.api_core.exceptions import DeadlineExceeded
    from google.genai.types import (
        FinishReason,
        GenerateContentConfig,
        GenerateContentResponse,
        HttpOptions,
        SafetySetting,
    )

    # Temporary hack to be able to process Nano Banana returned images.
    genai._api_client.READ_BUFFER_SIZE = 5 * 1024 * 1024  # pyright: ignore[reportAttributeAccessIssue]
except ModuleNotFoundError as e:
    logger.error(f"Exception: {e}")
    logger.error('In order to use Google AI, you need to `uv add "pipecat-ai[google]"`.')
    raise ImportError(f"Missing module: {e}") from e


# The lowest thinking level each model accepts, keyed by model name prefix. A
# model that isn't listed is assumed to accept "minimal", the fastest setting.
_LOWEST_MODEL_THINKING_LEVELS = {
    "gemini-3.7-flash": "low",
    "gemini-3.8-flash": "low",
}

# Models that take their thinking configuration from thinking_level, keyed by
# model name prefix. Whether one of these honors a thinking_budget set alongside
# it varies by model and by backend, so such a budget can't be relied on.
_MODELS_SUPPORTING_THINKING_LEVEL = ("gemini-3",)


[docs] class GoogleThinkingConfig(BaseModel): """Configuration for controlling the model's internal "thinking" process used before generating a response. Gemini 2.5 and 3 series models have this thinking process. Set either ``thinking_level`` or ``thinking_budget``, never both. Parameters: thinking_level: Thinking level, for Gemini 3 models. Gemini 3 Flash accepts "minimal", "low", "medium", and "high", except Gemini 3.7 Flash and Gemini 3.8 Flash, which accept only "low", "medium", and "high". Gemini 3 Pro accepts "low" and "high". If not provided, the flash models default to "medium" and Pro defaults to "high". Note: Gemini 2.5 series must use thinking_budget instead. thinking_budget: Token budget for thinking, for the Gemini 2.5 series. -1 for dynamic thinking (model decides), 0 to disable thinking, or a specific token count (e.g., 128-32768 for 2.5 Pro). If not provided, most models today default to dynamic thinking. See https://ai.google.dev/gemini-api/docs/thinking#set-budget for default values and allowed ranges. Note: Gemini 3 models must use thinking_level instead. include_thoughts: Whether to include thought summaries in the response. Today's models default to not including thoughts (False). """ thinking_budget: int | None = Field(default=None) # Why `| str` here? To not break compatibility in case Google adds more # levels in the future. thinking_level: Literal["low", "high", "medium", "minimal"] | str | None = Field(default=None) include_thoughts: bool | None = Field(default=None)
[docs] @dataclass class GoogleLLMSettings(LLMSettings): """Settings for GoogleLLMService. Parameters: thinking: Thinking configuration. safety_settings: Content safety filters, as a list of :class:`~google.genai.types.SafetySetting`. Each entry pairs a harm category with the threshold at which content is blocked. Categories left unspecified keep the Gemini API defaults. """ thinking: Union["GoogleLLMService.ThinkingConfig", None, NotGiven] = field( default_factory=lambda: NOT_GIVEN ) safety_settings: list[SafetySetting] | None | NotGiven = field( default_factory=lambda: NOT_GIVEN )
[docs] @classmethod def from_mapping(cls, settings): """Convert a plain dict to settings, coercing thinking and safety dicts. For backward compatibility, a ``thinking`` value that is a plain dict is converted to a :class:`GoogleLLMService.ThinkingConfig`, and any ``safety_settings`` entry that is a plain dict is converted to a :class:`~google.genai.types.SafetySetting`. """ instance = super().from_mapping(settings) if is_given(instance.thinking) and isinstance(instance.thinking, dict): instance.thinking = GoogleLLMService.ThinkingConfig(**instance.thinking) if is_given(instance.safety_settings) and instance.safety_settings: instance.safety_settings = [ SafetySetting(**entry) if isinstance(entry, dict) else entry for entry in instance.safety_settings ] return instance
[docs] class GoogleLLMService(LLMService[GeminiLLMAdapter]): """Google AI (Gemini) LLM service implementation. This class implements inference with Google's AI models, translating internally from an LLMContext to the messages format expected by the Google AI model. """ Settings = GoogleLLMSettings _settings: Settings # Overriding the default adapter to use the Gemini one. adapter_class = GeminiLLMAdapter supports_response_schema: bool = True # Backward compatibility: ThinkingConfig used to be defined inline here. ThinkingConfig = GoogleThinkingConfig
[docs] @deprecated( "`GoogleLLMService.InputParams` is deprecated since 0.0.105 and will be removed in 2.0.0. " "Use `GoogleLLMService.Settings` instead." ) class InputParams(BaseModel): """Input parameters for Google AI models. .. deprecated:: 0.0.105 Use ``settings=GoogleLLMService.Settings(...)`` instead. Will be removed in 2.0.0. Parameters: max_tokens: Maximum number of tokens to generate. temperature: Sampling temperature between 0.0 and 2.0. top_k: Top-k sampling parameter. top_p: Top-p sampling parameter between 0.0 and 1.0. thinking: Thinking configuration with thinking_budget, thinking_level, and include_thoughts. Used to control the model's internal "thinking" process used before generating a response. Gemini 2.5 series models use thinking_budget; Gemini 3 models use thinking_level. If this is not provided, Pipecat disables thinking for all models where that's possible (the 2.5 series, except 2.5 Pro), to reduce latency. extra: Additional parameters as a dictionary. """ max_tokens: int | None = Field(default=4096, ge=1) temperature: float | None = Field(default=None, ge=0.0, le=2.0) top_k: int | None = Field(default=None, ge=0) top_p: float | None = Field(default=None, ge=0.0, le=1.0) thinking: Optional["GoogleLLMService.ThinkingConfig"] = Field(default=None) extra: dict[str, Any] | None = Field(default_factory=dict)
[docs] def __init__( self, *, api_key: str, model: str | None = None, params: InputParams | None = None, settings: Settings | None = None, system_instruction: str | None = None, tools: list[dict[str, Any]] | None = None, tool_config: dict[str, Any] | None = None, http_options: HttpOptions | None = None, stream_idle_timeout_secs: float | None = 20.0, retry_timeout_secs: float | None = 5.0, retry_on_timeout: bool | None = False, **kwargs, ): """Initialize the Google LLM service. Args: api_key: Google AI API key for authentication. model: Model name to use. .. deprecated:: 0.0.105 Use ``settings=GoogleLLMService.Settings(model=...)`` instead. Will be removed in 2.0.0. params: Optional model parameters for inference. .. deprecated:: 0.0.105 Use ``settings=GoogleLLMService.Settings(...)`` instead. Will be removed in 2.0.0. settings: Runtime-updatable settings for this service. When both deprecated parameters and *settings* are provided, *settings* values take precedence. system_instruction: System instruction/prompt for the model. .. deprecated:: 0.0.105 Use ``settings=GoogleLLMService.Settings(system_instruction=...)`` instead. Will be removed in 2.0.0. tools: List of available tools/functions. tool_config: Configuration for tool usage. http_options: HTTP options for the client. stream_idle_timeout_secs: How long to wait for the next chunk of a streamed response before giving up on it. Bounds the wait when the API accepts a request and then stops producing without closing the stream. This is a gap between chunks, not a limit on how long a response may take overall. The first chunk is the slowest, since its wait spans the whole round trip including any thinking the model does before it emits anything; raise this for models configured to think at length. Set to ``None`` to wait indefinitely. retry_timeout_secs: How long to wait for the first chunk before giving up on the request and re-issuing it, when ``retry_on_timeout`` is set. Like the first-chunk wait above, this window spans the whole round trip including any thinking, so it is only a good fit for models that start emitting quickly. retry_on_timeout: Whether to re-issue the request once if the first chunk doesn't arrive within ``retry_timeout_secs``. Only the first chunk is retried: once a chunk has been pushed downstream, re-issuing would duplicate the response. **kwargs: Additional arguments passed to parent class. """ # 1. Initialize default_settings with hardcoded defaults default_settings = self.Settings( model="gemini-3.6-flash", system_instruction=None, max_tokens=4096, temperature=None, top_k=None, top_p=None, frequency_penalty=None, presence_penalty=None, seed=None, filter_incomplete_user_turns=False, user_turn_completion_config=None, thinking=None, safety_settings=None, extra={}, ) # 2. Apply direct init arg overrides (deprecated) if model is not None: self._warn_init_param_moved_to_settings("model", "model") default_settings.model = model if system_instruction is not None: self._warn_init_param_moved_to_settings("system_instruction", "system_instruction") default_settings.system_instruction = system_instruction # 3. Apply params overrides — only if settings not provided if params is not None: self._warn_init_param_moved_to_settings("params") if not settings: default_settings.max_tokens = params.max_tokens default_settings.temperature = params.temperature default_settings.top_k = params.top_k default_settings.top_p = params.top_p default_settings.thinking = params.thinking if isinstance(params.extra, dict): default_settings.extra = params.extra # 4. Apply settings delta (canonical API, always wins) if settings is not None: default_settings.apply_update(settings) super().__init__(settings=default_settings, **kwargs) self._api_key = api_key self._http_options = update_google_client_http_options(http_options) self._tools = tools self._tool_config = tool_config self._stream_idle_timeout_secs = stream_idle_timeout_secs self._retry_timeout_secs = retry_timeout_secs self._retry_on_timeout = retry_on_timeout # Initialize the API client. Subclasses can override this if needed. self.create_client() self._warn_if_thinking_budget_ignored()
[docs] def can_generate_metrics(self) -> bool: """Check if the service can generate usage metrics. Returns: True, as Google AI provides token usage metrics. """ return True
[docs] def create_client(self): """Create the Gemini client instance. Subclasses can override this.""" self._client = genai.Client(api_key=self._api_key, http_options=self._http_options)
[docs] @staticmethod def model_supports_response_schema(model: str) -> bool: """Whether a model can enforce a response schema. Gemini takes a JSON schema from the 2.5 models on. A model id without a version is assumed to support it. Args: model: The model name. """ match = re.search(r"gemini(?:-[a-z]+)*-(\d+)(?:\.(\d+))?", model) if not match: return True return (int(match.group(1)), int(match.group(2) or 0)) >= (2, 5)
[docs] async def run_inference( self, context: LLMContext, max_tokens: int | None = None, system_instruction: str | None = None, response_schema: dict[str, Any] | None = None, ) -> str | None: """Run a one-shot, out-of-band (i.e. out-of-pipeline) inference with the given LLM context. Args: context: The LLM context containing conversation history. max_tokens: Optional maximum number of tokens to generate. If provided, overrides the service's default max_tokens setting. system_instruction: Optional system instruction to use for this inference. If provided, overrides any system instruction in the context. response_schema: Optional JSON schema the reply must follow. The service asks the provider to enforce it, so the reply is JSON text matching the schema. Returns: The LLM's response as a string, or None if no response is generated. """ messages = [] system = [] tools = [] effective_instruction = system_instruction or assert_given( self._settings.system_instruction ) adapter = self.get_llm_adapter() params = adapter.get_llm_invocation_params( context, system_instruction=effective_instruction, ensure_last_message_is_user=self._should_inject_trailing_user_message(), ) messages = params["messages"] system = params["system_instruction"] tools = params["tools"] # Build generation config using the same method as streaming generation_params = self._build_generation_params( system_instruction=system, tools=tools if tools else None ) # Override max_output_tokens if provided if max_tokens is not None: generation_params["max_output_tokens"] = max_tokens response_schema = self._check_response_schema(response_schema) if response_schema is not None: generation_params["response_mime_type"] = "application/json" generation_params["response_json_schema"] = response_schema generation_config = GenerateContentConfig(**generation_params) # Use the new google-genai client's async method assert self._client is not None model = assert_given(self._settings.model) assert model is not None response = await self._client.aio.models.generate_content( model=model, contents=cast(Any, messages), config=generation_config, ) # Extract text from response if response.candidates and response.candidates[0].content: for part in response.candidates[0].content.parts or []: if part.text: return part.text return None
def _build_generation_params( self, system_instruction: str | None = None, tools: list | None = None, tool_config: dict[str, Any] | None = None, ) -> dict[str, Any]: """Build generation parameters for Google AI API. Args: system_instruction: Optional system instruction to use. tools: Optional list of tools to include. tool_config: Optional tool configuration. Returns: Dictionary of generation parameters with None values filtered out, carrying the low-latency thinking default when none is configured. """ # Filter out None values and create GenerationContentConfig generation_params = { k: v for k, v in { "system_instruction": system_instruction, "temperature": self._settings.temperature, "top_p": self._settings.top_p, "top_k": self._settings.top_k, "max_output_tokens": self._settings.max_tokens, "seed": self._settings.seed, "safety_settings": assert_given(self._settings.safety_settings), "tools": tools, "tool_config": tool_config, }.items() if v is not None } # Add thinking parameters if configured thinking = assert_given(self._settings.thinking) if thinking: generation_params["thinking_config"] = thinking.model_dump(exclude_unset=True) if self._settings.extra: generation_params.update(self._settings.extra) # Applied last, so an explicit thinking config from the settings or from # extra wins over the low-latency default. self._maybe_unset_thinking_budget(generation_params) return generation_params async def _update_settings(self, delta: GoogleLLMSettings) -> dict[str, Any]: """Apply a settings delta, re-checking the thinking configuration. Args: delta: An LLM settings delta. Returns: Dict mapping changed field names to their previous values. """ changed = await super()._update_settings(delta) if "model" in changed or "thinking" in changed: self._warn_if_thinking_budget_ignored() return changed def _warn_if_thinking_budget_ignored(self): """Warn when a thinking budget is set on a model that takes a level instead. Gemini 3 models control thinking through ``thinking_level``. Whether a budget set alongside one is honored, quietly ignored, or rejected varies by model and by backend, and the rejection names no field, so this warning is the only signal the caller gets in the cases that don't work. """ thinking = assert_given(self._settings.thinking) if not thinking or thinking.thinking_budget is None: return model = assert_given(self._settings.model) if not model or not model.startswith(_MODELS_SUPPORTING_THINKING_LEVEL): return logger.warning( f"{self}: thinking_budget is the pre-Gemini 3 thinking control, and {model} may " "ignore it or reject the request outright. Use thinking_level instead." ) def _maybe_unset_thinking_budget(self, generation_params: dict[str, Any]): try: model = assert_given(self._settings.model) # If we have an image model, we don't apply a thinking default. if model is None or "image" in model: return # If thinking_config is already set, don't override it. if "thinking_config" in generation_params: return # Apply model-aware low-latency thinking defaults. # Gemini 2.5 Flash: disable thinking via thinking_budget. # Gemini 3+ Flash: use the lowest thinking_level the model accepts. if model.startswith("gemini-2.5-flash"): generation_params["thinking_config"] = {"thinking_budget": 0} elif model.startswith("gemini-3") and "flash" in model: level = next( ( lowest for prefix, lowest in _LOWEST_MODEL_THINKING_LEVELS.items() if model.startswith(prefix) ), "minimal", ) generation_params["thinking_config"] = {"thinking_level": level} except Exception as e: logger.error(f"Failed to unset thinking budget: {e}") # Models known to accept a request whose contents end with a model turn, # continuing that turn as the start of the response. Newer models reject # such requests, so this is a frozen legacy set: any model NOT matching is # assumed to reject them and gets a trailing user message injected when # needed. gemini-3.5-flash accepts them but shares a prefix with # gemini-3.5-flash-lite, which doesn't, so it's left out. _PREFILL_SUPPORTED_PATTERNS = ( "gemini-2.", "gemini-3-", "gemini-3.1-", "gemini-pro-latest", ) def _should_inject_trailing_user_message(self) -> bool: """Whether to fix up requests whose contents end with a model turn. Models without support for continuing a trailing model turn reject such requests, so injection is on for every model not known to support it. Subclasses with other model naming can override ``_PREFILL_SUPPORTED_PATTERNS``. """ model = self._settings.model or "" return not any(model.startswith(p) for p in self._PREFILL_SUPPORTED_PATTERNS) async def _stream_content(self, context: LLMContext) -> AsyncIterator[GenerateContentResponse]: adapter = self.get_llm_adapter() params = adapter.get_llm_invocation_params( context, system_instruction=assert_given(self._settings.system_instruction), ensure_last_message_is_user=self._should_inject_trailing_user_message(), ) logger.debug( f"{self}: Generating chat from context {adapter.get_messages_for_logging(context)}" ) messages = params["messages"] # The adapter already resolved system_instruction vs context system message. system_instruction = params["system_instruction"] tools = [] if params["tools"]: tools = params["tools"] elif self._tools: tools = self._tools tool_config = None if self._tool_config: tool_config = self._tool_config # Build generation parameters generation_params = self._build_generation_params( system_instruction=system_instruction, tools=tools, tool_config=tool_config, ) generation_config = GenerateContentConfig(**generation_params) assert self._client is not None model = assert_given(self._settings.model) assert model is not None return await self._client.aio.models.generate_content_stream( model=model, contents=cast(Any, messages), config=generation_config, ) async def _stream_response(self, context: LLMContext) -> AsyncIterator[GenerateContentResponse]: """Yield the model's streamed chunks, re-issuing the request once if it stalls. The API client sends the request lazily, when the first chunk is pulled, so a request that is accepted and then produces nothing shows up as a first chunk that never arrives. Re-issuing is safe only up to that point: once a chunk has been yielded, its content is already on its way downstream. Args: context: The context to generate a response for. Yields: Each chunk of the streamed response. Raises: TimeoutError: If the stream stalls past the timeout it is bounded by. """ if self._retry_on_timeout: async with aclosing( self._iter_stream( await self._stream_content(context), first_chunk_timeout=self._retry_timeout_secs, ) ) as stream: try: first_chunk = await stream.__anext__() except StopAsyncIteration: return except (TimeoutError, DeadlineExceeded): logger.debug(f"{self}: retrying content generation due to timeout") else: yield first_chunk async for chunk in stream: yield chunk return # Either retries are disabled, or the initial attempt timed out. async with aclosing( self._iter_stream( await self._stream_content(context), first_chunk_timeout=self._stream_idle_timeout_secs, ) ) as stream: async for chunk in stream: yield chunk async def _iter_stream( self, response: AsyncIterator[GenerateContentResponse], *, first_chunk_timeout: float | None, ) -> AsyncGenerator[GenerateContentResponse, None]: """Yield streamed chunks, giving up if the stream stalls between them. A stream that stops producing without closing would otherwise leave the response open indefinitely, since the API client applies no timeout of its own. The timeout covers the gap between chunks rather than the response as a whole, so a slow but healthy stream is never cut short. Args: response: The streamed response to consume. first_chunk_timeout: How long to wait for the first chunk. Every later chunk is bounded by ``stream_idle_timeout_secs``. Yields: Each chunk of the streamed response. Raises: TimeoutError: If the next chunk doesn't arrive in time. """ chunks = response.__aiter__() timeout = first_chunk_timeout try: while True: try: if timeout is None: chunk = await chunks.__anext__() else: chunk = await asyncio.wait_for(chunks.__anext__(), timeout=timeout) except StopAsyncIteration: return yield chunk timeout = self._stream_idle_timeout_secs finally: # Release the HTTP resources held by a stream we walk away from, whether # that's because it stalled, was interrupted, or is being re-issued. # Closing is best-effort: the client's stream type only promises # AsyncIterator, which has no aclose(). aclose = getattr(response, "aclose", None) if aclose is not None: await aclose() def _handle_finish_reason(self, finish_reason: FinishReason): """Log why Gemini stopped generating, when it stopped for a notable reason. Anything other than a normal stop leaves the turn short of what the model would otherwise have said: the response was withheld (safety, recitation, prohibited content), a tool call was rejected, or the output hit the token limit. Whatever text did arrive is still passed downstream. """ if finish_reason in (FinishReason.STOP, FinishReason.FINISH_REASON_UNSPECIFIED): return logger.warning(f"{self}: response incomplete, the model stopped for {finish_reason.name}") @traced_llm async def _process_context(self, context: LLMContext): await self.push_frame(LLMFullResponseStartFrame()) prompt_tokens = 0 completion_tokens = 0 total_tokens = 0 cache_read_input_tokens = 0 reasoning_tokens = 0 grounding_metadata = None accumulated_text = "" try: await self.start_ttfb_metrics() function_calls = [] async for chunk in self._stream_response(context): # Gemini may send usage_metadata in multiple chunks with varying behavior: # - Sometimes a single chunk, sometimes multiple chunks # - Token counts may be cumulative (growing) or may change between chunks # - Early chunks may include estimates/overhead that gets refined # We use assignment (not accumulation) because the final chunk always contains # the authoritative, billable token usage for the entire response. if chunk.usage_metadata: prompt_tokens = chunk.usage_metadata.prompt_token_count or 0 completion_tokens = chunk.usage_metadata.candidates_token_count or 0 total_tokens = chunk.usage_metadata.total_token_count or 0 cache_read_input_tokens = chunk.usage_metadata.cached_content_token_count or 0 reasoning_tokens = chunk.usage_metadata.thoughts_token_count or 0 if not chunk.candidates: continue # A leading chunk can carry usage metadata and no candidates, so # TTFB ends at the first chunk that holds model output. await self.stop_ttfb_metrics() for candidate in chunk.candidates: if candidate.content and candidate.content.parts: for part in candidate.content.parts: function_call_id = None if part.text: if part.thought: # Gemini emits fully-formed thoughts rather # than chunks so bracket each thought in # start/end await self.push_frame(LLMThoughtStartFrame()) await self.push_frame(LLMThoughtTextFrame(part.text)) await self.push_frame(LLMThoughtEndFrame()) else: accumulated_text += part.text await self._push_llm_text(part.text) elif part.function_call: # A turn that only calls tools produces no answer # text, so the call itself is what the caller gets # and TTFAT ends here rather than going unmeasured. await self.stop_ttfat_metrics() function_call = part.function_call function_call_id = function_call.id or str(uuid.uuid4()) logger.debug( f"Function call: {function_call.name}:{function_call_id}" ) function_calls.append( FunctionCallFromLLM( context=context, tool_call_id=function_call_id, function_name=function_call.name or "", arguments=function_call.args or {}, ) ) elif part.inline_data and part.inline_data.data: # Here we assume that inline_data is an image. image = Image.open(io.BytesIO(part.inline_data.data)) await self.push_frame( AssistantImageRawFrame( image=image.tobytes(), size=image.size, format="RGB", original_data=part.inline_data.data, original_mime_type=part.inline_data.mime_type, ) ) # Handle Gemini thought signatures. # # - Gemini 2.5: they appear on function_call Parts, # and then (surprisingly) on the last(*) Part of # model responses following the first function_call # in a conversation. # - Gemini 3 Pro: they appear on the last(*) Part # of model responses, regardless of Part type. # # (*) Since we're using the streaming API, though, # where text Parts may be split across multiple # chunks (each represented by a Part, confusingly), # signatures may actually appear with the first # chunk (Gemini 2.5) or in a trailing empty-text # chunk (Gemini 3 Pro). if part.thought_signature: # Save a "bookmark" for the signature, so we # can later be sure we've put it in the right # place in context when sending the context # back to the LLM to continue the conversation. bookmark = {} if part.function_call: bookmark["function_call"] = function_call_id elif part.inline_data and part.inline_data.data: bookmark["inline_data"] = part.inline_data elif part.text is not None: # Account for Gemini 3 Pro trailing # empty-text chunk by using all the text # seen so far in this response's chunks. bookmark["text"] = accumulated_text else: logger.warning("Thought signature found on unhandled Part type") if bookmark: await self.push_frame( LLMMessagesAppendFrame( [ self.get_llm_adapter().create_llm_specific_message( { "type": "thought_signature", "signature": part.thought_signature, "bookmark": bookmark, } ) ] ) ) if ( candidate.grounding_metadata and candidate.grounding_metadata.grounding_chunks ): m = candidate.grounding_metadata rendered_content = ( m.search_entry_point.rendered_content if m.search_entry_point else None ) origins = [ { "site_uri": grounding_chunk.web.uri if grounding_chunk.web else None, "site_title": grounding_chunk.web.title if grounding_chunk.web else None, "results": [ { "text": grounding_support.segment.text if grounding_support.segment else "", "confidence": grounding_support.confidence_scores, } for grounding_support in ( m.grounding_supports if m.grounding_supports else [] ) if grounding_support.grounding_chunk_indices and index in grounding_support.grounding_chunk_indices ], } for index, grounding_chunk in enumerate( m.grounding_chunks if m.grounding_chunks else [] ) ] grounding_metadata = { "rendered_content": rendered_content, "origins": origins, } if candidate.finish_reason: self._handle_finish_reason(candidate.finish_reason) await self.run_function_calls(function_calls) except (TimeoutError, DeadlineExceeded) as e: await self._call_event_handler("on_completion_timeout") await self.push_error(error_msg="LLM completion timeout", exception=e) except LLMContextConversionError as e: await self.push_error(error_msg=str(e), exception=e) except Exception as e: await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e) finally: if grounding_metadata and isinstance(grounding_metadata, dict): llm_search_frame = LLMSearchResponseFrame( search_result=accumulated_text, origins=grounding_metadata["origins"], rendered_content=grounding_metadata["rendered_content"], ) await self.push_frame(llm_search_frame) await self.start_llm_usage_metrics( LLMTokenUsage( prompt_tokens=prompt_tokens, completion_tokens=completion_tokens, total_tokens=total_tokens, cache_read_input_tokens=cache_read_input_tokens, reasoning_tokens=reasoning_tokens, ) ) await self.push_frame(LLMFullResponseEndFrame())
[docs] async def process_frame(self, frame: Frame, direction: FrameDirection): """Process incoming frames and handle different frame types. Args: frame: The frame to process. direction: Direction of frame processing. """ await super().process_frame(frame, direction) if isinstance(frame, LLMContextFrame): await self._process_context(frame.context) else: await self.push_frame(frame, direction)
[docs] async def stop(self, frame): """Override stop to gracefully close the client.""" await super().stop(frame) await self._close_client()
[docs] async def cancel(self, frame): """Override cancel to gracefully close the client.""" await super().cancel(frame) await self._close_client()
[docs] async def cleanup(self): """Release resources held by the service.""" await super().cleanup() await self._close_client()
async def _close_client(self): if not self._client: return try: await self._client.aio.aclose() except Exception: # Do nothing - we're shutting down anyway pass finally: self._client = None