Skip to content

Commit 4d31fa9

Browse files
google-genai-botcopybara-github
authored andcommitted
feat(plugins): add on_schema_error and on_schema_ready callback hooks to BigQueryLoggerConfig
Allow hosts of `BigQueryAgentAnalyticsPlugin` to register optional `on_schema_error: Callable[[Exception], None]` and `on_schema_ready: Callable[[], None]` callbacks on `BigQueryLoggerConfig` that are invoked on a worker thread when the table readiness pass in `_ensure_schema_exists()` fails to read, create, or upgrade the BigQuery table, or completes with `True`, respectively. PiperOrigin-RevId: 996039246
1 parent 8e386a9 commit 4d31fa9

2 files changed

Lines changed: 640 additions & 3 deletions

File tree

‎src/google/adk/plugins/bigquery_agent_analytics_plugin.py‎

Lines changed: 113 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1709,6 +1709,13 @@ def _require_finite(name: str, value: Any, minimum_exclusive: float) -> float:
17091709
return float(value)
17101710

17111711

1712+
def _is_async_callable(fn: object) -> bool:
1713+
"""Returns whether ``fn`` is a coroutine function or async callable object."""
1714+
return inspect.iscoroutinefunction(fn) or inspect.iscoroutinefunction(
1715+
getattr(fn, "__call__", None)
1716+
)
1717+
1718+
17121719
def _validate_runtime_config(config: "BigQueryLoggerConfig") -> None:
17131720
"""Validates runtime settings at construction time.
17141721
@@ -1717,7 +1724,7 @@ def _validate_runtime_config(config: "BigQueryLoggerConfig") -> None:
17171724
every batch is dropped without a single attempt.
17181725
17191726
Raises:
1720-
ValueError: If any batch, queue, duration, or retry setting is
1727+
ValueError: If any batch, queue, duration, retry, or callback setting is
17211728
invalid.
17221729
"""
17231730
_require_count("batch_size", config.batch_size, 1)
@@ -1762,6 +1769,22 @@ def _validate_runtime_config(config: "BigQueryLoggerConfig") -> None:
17621769
"retry_config.max_delay must be >= initial_delay, got"
17631770
f" max_delay={retry.max_delay} initial_delay={retry.initial_delay}."
17641771
)
1772+
if config.on_schema_error is not None and (
1773+
not callable(config.on_schema_error)
1774+
or _is_async_callable(config.on_schema_error)
1775+
):
1776+
raise ValueError(
1777+
"on_schema_error must be a synchronous callable or None, got"
1778+
f" {config.on_schema_error!r}."
1779+
)
1780+
if config.on_schema_ready is not None and (
1781+
not callable(config.on_schema_ready)
1782+
or _is_async_callable(config.on_schema_ready)
1783+
):
1784+
raise ValueError(
1785+
"on_schema_ready must be a synchronous callable or None, got"
1786+
f" {config.on_schema_ready!r}."
1787+
)
17651788

17661789

17671790
# Cloud Platform OAuth scope. Assembled from parts so this module does not
@@ -2461,6 +2484,43 @@ class BigQueryLoggerConfig:
24612484
never written to BigQuery: the row's ``error_message`` still names
24622485
only the exception class. ``False`` (the default) logs a constant
24632486
message with no traceback.
2487+
on_schema_error: Optional synchronous callback invoked with the triggering
2488+
exception each time the table readiness pass fails to read, create, or
2489+
upgrade the BigQuery table. It runs on a worker thread, not the event
2490+
loop, and fires at most once per readiness pass, so it may be called
2491+
repeatedly while setup retries; keep it fast and non-blocking. Once a
2492+
pass completes with ``True`` (see ``on_schema_ready``), later setups of
2493+
the same plugin, including after ``close()``, after a fork, and in
2494+
copies pickled afterwards, skip the pass, so neither callback fires
2495+
again and later table outages surface only as append-time drops. Later
2496+
setups run the pass again only if the setup awaiting that pass was
2497+
cancelled, or after a ``create_analytics_views()`` call tries to create
2498+
the views and one fails; a call that creates every view makes later
2499+
setups skip the pass, and that call's own view creation never invokes
2500+
either callback. This callback is not invoked for analytics-view
2501+
creation failures or label-only ``update_table`` refresh failures (those
2502+
are logged and swallowed without raising). Exceptions it raises are
2503+
logged and suppressed so the original exception still propagates. If the
2504+
plugin is serialized with the standard ``pickle`` module, the callback
2505+
must be picklable (e.g. a module-level function; lambdas and closures
2506+
fail with standard ``pickle``, and bound methods pickle their instance).
2507+
on_schema_ready: Optional synchronous no-arg callback invoked each time
2508+
the table readiness pass completes with ``True`` (i.e. the table exists
2509+
or was created/upgraded and all requested views succeeded; not called
2510+
when view creation returns ``False`` or when an exception is raised).
2511+
After such a pass, later setups of the same plugin, including after
2512+
``close()``, after a fork, and in copies pickled afterwards, skip the
2513+
pass, so neither callback fires again and later table outages surface
2514+
only as append-time drops. Later setups run the pass again only if the
2515+
setup awaiting that pass was cancelled, or after a
2516+
``create_analytics_views()`` call tries to create the views and one
2517+
fails; a call that creates every view makes later setups skip the pass,
2518+
and that call's own view creation never invokes either callback. It runs
2519+
on a worker thread, not the event loop; keep it fast and non-blocking.
2520+
Exceptions it raises are logged and suppressed. If the plugin is
2521+
serialized with the standard ``pickle`` module, the callback must be
2522+
picklable (e.g. a module-level function; lambdas and closures fail with
2523+
standard ``pickle``, and bound methods pickle their instance).
24642524
"""
24652525

24662526
enabled: bool = True
@@ -2550,6 +2610,8 @@ class BigQueryLoggerConfig:
25502610
# local formatter-failure warning. The traceback can embed the unformatted
25512611
# content; see the class docstring before enabling it.
25522612
debug_content_formatter_errors: bool = False
2613+
on_schema_error: Optional[Callable[[Exception], None]] = None
2614+
on_schema_ready: Optional[Callable[[], None]] = None
25532615

25542616

25552617
# ==============================================================================
@@ -5934,6 +5996,44 @@ def _close() -> None:
59345996

59355997
await self._get_loop_state(claimed_generation=claimed_generation)
59365998

5999+
def _notify_schema_error(self, exc: Exception) -> None:
6000+
"""Invokes ``config.on_schema_error(exc)`` if configured."""
6001+
callback = self.config.on_schema_error
6002+
if callback is None:
6003+
return
6004+
try:
6005+
result = callback(exc)
6006+
if inspect.iscoroutine(result):
6007+
result.close()
6008+
logger.warning(
6009+
"BigQueryLoggerConfig.on_schema_error returned a coroutine;"
6010+
" callback must be synchronous."
6011+
)
6012+
except Exception: # pylint: disable=broad-exception-caught
6013+
logger.warning(
6014+
"Exception raised by BigQueryLoggerConfig.on_schema_error callback",
6015+
exc_info=True,
6016+
)
6017+
6018+
def _notify_schema_ready(self) -> None:
6019+
"""Invokes ``config.on_schema_ready()`` if configured."""
6020+
callback = self.config.on_schema_ready
6021+
if callback is None:
6022+
return
6023+
try:
6024+
result = callback()
6025+
if inspect.iscoroutine(result):
6026+
result.close()
6027+
logger.warning(
6028+
"BigQueryLoggerConfig.on_schema_ready returned a coroutine;"
6029+
" callback must be synchronous."
6030+
)
6031+
except Exception: # pylint: disable=broad-exception-caught
6032+
logger.warning(
6033+
"Exception raised by BigQueryLoggerConfig.on_schema_ready callback",
6034+
exc_info=True,
6035+
)
6036+
59376037
def _ensure_schema_exists(self) -> bool:
59386038
"""Ensures the BigQuery table exists with the correct schema.
59396039
@@ -5954,6 +6054,18 @@ def _ensure_schema_exists(self) -> bool:
59546054
or created, or a required schema upgrade failed.
59556055
"""
59566056
client = self._require_client()
6057+
self._require_schema()
6058+
try:
6059+
ready = self._ensure_schema_exists_impl(client)
6060+
except Exception as e:
6061+
self._notify_schema_error(e)
6062+
raise
6063+
if ready:
6064+
self._notify_schema_ready()
6065+
return ready
6066+
6067+
def _ensure_schema_exists_impl(self, client: bigquery.Client) -> bool:
6068+
"""Table readiness pass; see ``_ensure_schema_exists``."""
59576069
try:
59586070
existing_table = client.get_table(self.full_table_id)
59596071
if self.config.auto_schema_upgrade:

0 commit comments

Comments
 (0)