#
# 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]
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