refactor(agent/small-modules): compact battery, browser_registry/provider, command_token_source, bounded_response (554->438 LOC)

This commit is contained in:
Teknium
2026-09-02 18:18:18 -07:00
parent 113f04616b
commit a0ccec3f6d
5 changed files with 88 additions and 204 deletions

View File

@@ -1,9 +1,8 @@
"""System-battery read-out for the CLI/TUI status bar.
Reads the host battery through ``psutil`` and exposes a compact, colour-coded
label. Everything degrades to "unavailable" (no battery / read failure) so
callers can render unconditionally. The status bar repaints on every keystroke,
so :func:`read_battery` memoises the reading for a few seconds.
Reads the host battery through ``psutil`` and exposes a compact, colour-coded label. Everything
degrades to "unavailable" (no battery / read failure) so callers can render unconditionally. The
status bar repaints on every keystroke, so :func:`read_battery` memoises the reading for a few seconds.
"""
from __future__ import annotations
@@ -15,8 +14,7 @@ from typing import Optional
@dataclass(frozen=True)
class BatteryStatus:
"""One reading: ``percent`` clamped 0-100; ``plugged`` None when the
platform can't tell."""
"""One reading: ``percent`` clamped 0-100; ``plugged`` None when the platform can't tell."""
available: bool
percent: Optional[int] = None
@@ -29,8 +27,7 @@ class BatteryStatus:
UNAVAILABLE = BatteryStatus(available=False)
# Colour buckets, mirroring the status-bar context styles but inverted (a full
# battery is "good", an empty one is "critical").
# Colour buckets, mirroring the status-bar context styles but inverted (full battery = "good").
CATEGORY_GOOD = "good"
CATEGORY_WARN = "warn"
CATEGORY_BAD = "bad"
@@ -54,7 +51,6 @@ def _read_battery_uncached() -> BatteryStatus:
return UNAVAILABLE
if batt is None:
return UNAVAILABLE
percent: Optional[int] = None
raw_percent = getattr(batt, "percent", None)
if raw_percent is not None:
@@ -69,11 +65,8 @@ def _read_battery_uncached() -> BatteryStatus:
def read_battery(use_cache: bool = True) -> BatteryStatus:
"""Return the current battery status (cached for a few seconds)."""
global _cache
if use_cache and _cache is not None:
ts, cached = _cache
if time.monotonic() - ts < _CACHE_TTL_SECONDS:
return cached
if use_cache and _cache is not None and time.monotonic() - _cache[0] < _CACHE_TTL_SECONDS:
return _cache[1]
status = _read_battery_uncached()
_cache = (time.monotonic(), status)
return status
@@ -89,8 +82,7 @@ def battery_category(status: BatteryStatus) -> str:
"""Bucket a reading into a colour category: good/warn/bad/critical/dim."""
if not status.available or status.percent is None:
return CATEGORY_DIM
# On AC power the level isn't a concern — always read as healthy.
if status.charging:
if status.charging: # on AC power the level isn't a concern
return CATEGORY_GOOD
for bound, category in _LEVEL_CATEGORIES:
if status.percent <= bound:
@@ -99,12 +91,12 @@ def battery_category(status: BatteryStatus) -> str:
def battery_glyph(status: BatteryStatus) -> str:
"""Return the leading glyph: a bolt while charging, else a battery."""
"""Leading glyph: a bolt while charging, else a battery."""
return "\u26a1" if status.charging else "\U0001f50b" # ⚡ / 🔋
def format_battery(status: BatteryStatus) -> str:
"""Return a compact label like ``🔋 82%`` / ``⚡ 82%`` (empty if N/A)."""
"""Compact label like ``🔋 82%`` / ``⚡ 82%`` (empty if N/A)."""
if not status.available or status.percent is None:
return ""
return f"{battery_glyph(status)} {status.percent}%"

View File

@@ -1,21 +1,12 @@
"""Bounded reads of HTTP error response bodies.
On a non-OK *streaming* response Hermes reads the body for a diagnostic. A
bare ``response.read()`` is unbounded two ways: a server can stream an
arbitrarily large body (memory), or open the body and stall forever (hang).
The diagnostic is only ever shown truncated to a few hundred chars, so
``read_streaming_error_body`` caps bytes and enforces a hard wall-clock
deadline; callers feed the returned text to their error builders instead of
touching ``response.text`` (unbounded / raises after a partial stream read).
Subtlety: ``httpx.iter_bytes()`` blocks *inside* the socket read, so a
deadline checked only between chunks can't interrupt a mid-chunk stall until
httpx's own (30s+) read timeout fires. The read therefore runs on a daemon
thread; on timeout we close the response (unblocking the read) and return the
partial bytes collected so far.
Covers the three streaming error-body sites: native Gemini, Gemini Cloud
Code, Antigravity Cloud Code.
On a non-OK *streaming* response Hermes reads the body for a diagnostic (only ever shown truncated to
a few hundred chars). A bare ``response.read()`` is unbounded two ways: arbitrarily large body
(memory) or a body that stalls forever (hang). ``read_streaming_error_body`` caps bytes and enforces a
hard wall-clock deadline; callers use the returned text instead of ``response.text`` (unbounded /
raises after a partial stream read). ``httpx.iter_bytes()`` blocks *inside* the socket read, so the
read runs on a daemon thread; on timeout we close the response (unblocking the read) and return the
partial bytes. Used by the streaming error-body sites: native Gemini, Gemini Cloud Code, Antigravity.
"""
from __future__ import annotations
@@ -28,11 +19,9 @@ import httpx
logger = logging.getLogger(__name__)
# Comfortably holds any real provider error envelope (Google RPC / Anthropic
# error JSON) while rejecting pathological bodies.
# Comfortably holds any real provider error envelope while rejecting pathological bodies.
DEFAULT_ERROR_BODY_MAX_BYTES = 64 * 1024
# Hard deadline for the whole read; past it the connection is closed and the
# partial bytes are kept.
# Hard deadline for the whole read; past it the connection is closed and the partial bytes are kept.
DEFAULT_ERROR_BODY_TIMEOUT_S = 10.0
@@ -44,9 +33,9 @@ def read_streaming_error_body(
) -> str:
"""Read a non-OK streaming body with a byte cap and a hard deadline.
Returns UTF-8 text (errors replaced) truncated to ``max_bytes``. Never
raises: transport errors, stalls and oversize bodies yield best-effort
partial text (or ""), so a read error can't mask the original failure.
Returns UTF-8 text (errors replaced) truncated to ``max_bytes``. Never raises: transport errors,
stalls and oversize bodies yield best-effort partial text (or ""), so a read error can't mask the
original failure.
"""
chunks: List[bytes] = []
state = {"truncated": False}
@@ -75,11 +64,9 @@ def read_streaming_error_body(
if not done.wait(timeout=timeout_s):
logger.debug(
"bounded error-body read: hard timeout after %.1fs (%d bytes so far)",
timeout_s,
sum(len(c) for c in chunks),
timeout_s, sum(len(c) for c in chunks),
)
# Closing cancels any in-flight socket read so the worker unwinds. We do
# not join (it is a daemon and may be blocked in C).
# Closing cancels any in-flight socket read so the worker unwinds. No join (daemon, may be blocked in C).
try:
response.close()
except Exception: # noqa: BLE001
@@ -87,8 +74,6 @@ def read_streaming_error_body(
if state["truncated"]:
logger.debug(
"bounded error-body read: capped at %d bytes (max=%d)",
sum(len(c) for c in chunks),
max_bytes,
"bounded error-body read: capped at %d bytes (max=%d)", sum(len(c) for c in chunks), max_bytes,
)
return b"".join(chunks).decode("utf-8", errors="replace")

View File

@@ -1,16 +1,12 @@
"""
Browser Provider ABC
====================
"""Browser Provider ABC: pluggable cloud browser backends (Browserbase, Browser Use, Firecrawl, …).
Pluggable-backend interface for cloud browser providers (Browserbase, Browser
Use, Firecrawl, …). Providers register via
:meth:`PluginContext.register_browser_provider`; the active one (selected by
``browser.cloud_provider`` in ``config.yaml``) services every cloud-mode
``browser_*`` tool call. Providers live in ``<repo>/plugins/browser/<name>/``
(built-in) or ``~/.hermes/plugins/browser/<name>/`` (user, opt-in).
Providers register via :meth:`PluginContext.register_browser_provider`; the active one (selected by
``browser.cloud_provider``) services every cloud-mode ``browser_*`` tool call. They live in
``<repo>/plugins/browser/<name>/`` (built-in) or ``~/.hermes/plugins/browser/<name>/`` (user).
Session metadata contract (preserved from the legacy ``CloudBrowserProvider``
so :mod:`tools.browser_tool` needs no translation)::
Session metadata contract (legacy ``CloudBrowserProvider`` shape; ``tools.browser_tool`` needs no
translation). ``bb_session_id`` is a legacy key name kept verbatim — it holds the provider's session ID
regardless of provider::
{
"session_name": str, # unique name for agent-browser --session
@@ -20,9 +16,6 @@ so :mod:`tools.browser_tool` needs no translation)::
"features": dict, # feature flags that were enabled
"external_call_id": str, # optional, managed-gateway billing key
}
``bb_session_id`` is a legacy key name kept verbatim for backward compat — it
holds the provider's session ID regardless of which provider is in use.
"""
from __future__ import annotations
@@ -36,46 +29,34 @@ from agent.provider_base import ProviderBase
class BrowserProvider(ProviderBase):
"""Abstract base class for a cloud browser backend.
Subclasses implement :attr:`name` (the ``browser.cloud_provider`` value,
e.g. ``browserbase``, ``browser-use``, ``firecrawl``), :meth:`is_available`,
and the lifecycle trio :meth:`create_session` / :meth:`close_session` /
:meth:`emergency_cleanup`. ``get_setup_schema`` may add ``"post_setup"``
(e.g. ``"agent_browser"``) to trigger the install hook.
Subclasses implement :attr:`name` (the ``browser.cloud_provider`` value), :meth:`is_available`, and
the lifecycle trio :meth:`create_session` / :meth:`close_session` / :meth:`emergency_cleanup`.
``get_setup_schema`` may add ``"post_setup"`` (e.g. ``"agent_browser"``) to trigger the install hook.
"""
@abc.abstractmethod
def is_available(self) -> bool:
"""True when this provider can service calls.
Cheap check only (env var present, managed-gateway token readable, dep
importable) — must NOT make network calls; runs at tool-registration
time and on every ``hermes tools`` paint.
"""
"""True when this provider can service calls. Cheap check only (env var, token readable, dep
importable) — must NOT make network calls; runs at tool-registration time and on every
``hermes tools`` paint."""
@abc.abstractmethod
def create_session(self, task_id: str) -> Dict[str, object]:
"""Create a cloud browser session and return the session metadata dict
described in the module docstring.
May raise ``ValueError`` (missing credentials) or ``RuntimeError``
(network / API failure); the dispatcher surfaces these to the user.
"""
"""Create a cloud browser session and return the metadata dict from the module docstring.
May raise ``ValueError`` (missing credentials) or ``RuntimeError`` (network / API failure);
the dispatcher surfaces these to the user."""
@abc.abstractmethod
def close_session(self, session_id: str) -> bool:
"""Release a cloud session by provider session ID.
Returns True on success, False on failure. Should not raise — log and
return False so the dispatcher's cleanup loop keeps moving.
"""
"""Release a cloud session by provider session ID. Returns True on success, False on failure;
should not raise (log and return False so the dispatcher's cleanup loop keeps moving)."""
@abc.abstractmethod
def emergency_cleanup(self, session_id: str) -> None:
"""Best-effort teardown from atexit / signal handlers. Must tolerate
missing credentials and network errors; must not raise."""
"""Best-effort teardown from atexit / signal handlers. Must tolerate missing credentials and
network errors; must not raise."""
# Legacy ``CloudBrowserProvider`` names still used by ``tools.browser_tool``
# and out-of-tree subclasses; thin delegations to the current API.
# Legacy ``CloudBrowserProvider`` names still used by ``tools.browser_tool`` and out-of-tree subclasses.
def is_configured(self) -> bool:
"""Backward-compat alias for :meth:`is_available`."""

View File

@@ -1,37 +1,11 @@
"""
Browser Provider Registry
=========================
"""Browser provider registry: cloud browser backends registered by plugins via
:meth:`PluginContext.register_browser_provider`, consumed by ``tools.browser_tool._get_cloud_provider``.
Central map of registered cloud browser providers. Populated by plugins at
import-time via :meth:`PluginContext.register_browser_provider`; consumed by
:func:`tools.browser_tool._get_cloud_provider` to route each cloud-mode
``browser_*`` tool call to the active backend.
Active selection
----------------
The active provider is chosen by configuration with this precedence:
1. ``browser.cloud_provider`` in ``config.yaml`` (explicit override).
2. Legacy preference order — ``browser-use`` → ``browserbase`` — filtered by
availability. Matches the historic auto-detect order in
:func:`tools.browser_tool._get_cloud_provider` (Browser Use checked first
because it covers both the managed Nous gateway and direct API key path;
Browserbase as the older direct-credentials fallback). ``firecrawl`` is
intentionally NOT in the legacy walk — users only get Firecrawl as a
cloud browser when they explicitly set ``browser.cloud_provider:
firecrawl``, matching pre-migration behaviour where Firecrawl was never
auto-selected.
3. Otherwise ``None`` — the dispatcher falls back to local browser mode.
The explicit-config branch (rule 1) intentionally ignores ``is_available()``
so the dispatcher surfaces a typed "X_API_KEY is not set" error to the user
instead of silently switching backends. Matches the legacy
:func:`tools.browser_tool._get_cloud_provider` behaviour for configured names.
Note: there is no "capability" split here (unlike the web subsystem, which
has search/extract/crawl). Every browser provider implements the full
:class:`agent.browser_provider.BrowserProvider` lifecycle; the registry's
job is purely selection, not capability routing.
Active-provider precedence (see :func:`_resolve`): ``browser.cloud_provider`` in config.yaml wins
regardless of ``is_available()`` (so the dispatcher surfaces a typed "X_API_KEY is not set" error
instead of silently switching); else the legacy auto-detect walk ``browser-use`` → ``browserbase``
filtered by availability; else ``None`` (local browser mode). There is no capability split here —
every provider implements the full :class:`agent.browser_provider.BrowserProvider` lifecycle.
"""
from __future__ import annotations
@@ -51,47 +25,31 @@ _registry: ProviderRegistry[BrowserProvider] = ProviderRegistry(
_registry.export(globals())
# ---------------------------------------------------------------------------
# Active-provider resolution
# ---------------------------------------------------------------------------
# Auto-detect order when ``browser.cloud_provider`` is unset (pre-migration
# walk of :func:`tools.browser_tool._get_cloud_provider`); see :func:`_resolve`
# for why Firecrawl is absent.
_LEGACY_PREFERENCE = (
"browser-use",
"browserbase",
)
# Auto-detect order when ``browser.cloud_provider`` is unset (historic order: Browser Use first because
# it covers both the managed Nous gateway and the direct API key path; Browserbase as the older
# direct-credentials fallback). Firecrawl is deliberately absent — see :func:`_resolve`.
_LEGACY_PREFERENCE = ("browser-use", "browserbase")
def _resolve(configured: Optional[str]) -> Optional[BrowserProvider]:
"""Resolve the active browser provider (rules in the module docstring).
There is intentionally NO "single-eligible shortcut" (unlike
:func:`agent.web_search_registry._resolve`): only ``_LEGACY_PREFERENCE``
names are auto-eligible. Firecrawl shares its API key with the *web*
extract plugin, so a user with ``FIRECRAWL_API_KEY`` must never be routed
to a paid cloud browser without setting ``browser.cloud_provider``; the
same gate applies to third-party browser-provider plugins.
Intentionally NO "single-eligible shortcut" (unlike ``agent.web_search_registry._resolve``): only
``_LEGACY_PREFERENCE`` names are auto-eligible. Firecrawl shares its API key with the *web* extract
plugin, so a user with ``FIRECRAWL_API_KEY`` must never be routed to a paid cloud browser without
setting ``browser.cloud_provider``; the same gate applies to third-party browser-provider plugins.
"""
snapshot = _registry.merged()
if configured == "local":
return None
# Explicit config wins regardless of is_available(): the dispatcher then
# surfaces a precise "X_API_KEY is not set" error instead of a silent switch.
if configured:
provider = snapshot.get(configured)
if provider is not None:
return provider
logger.debug(
"browser cloud_provider '%s' configured but not registered; "
"falling back to auto-detect",
"browser cloud_provider '%s' configured but not registered; falling back to auto-detect",
configured,
)
for legacy in _LEGACY_PREFERENCE:
provider = snapshot.get(legacy)
if provider is not None and is_available_safe(
@@ -100,6 +58,4 @@ def _resolve(configured: Optional[str]) -> Optional[BrowserProvider]:
level=logging.WARNING, exc_info=True,
):
return provider
return None

View File

@@ -1,24 +1,12 @@
"""Mint a provider API key by running a command (``key_cmd``).
Enterprise gateways (SSO/OIDC brokers, cloud IAM, auth proxies) issue
SHORT-LIVED bearers; a key copied into ``.env`` goes stale within the hour.
``key_cmd`` names a command that PRINTS a token (the ``apiKeyHelper`` /
``gcloud auth print-access-token`` / ``databricks auth token`` idiom)::
providers:
my-gateway:
base_url: https://gateway.internal.example.com/v1
api_mode: chat_completions
key_cmd: my-auth-cli print-token --profile prod
Both wire clients already accept a callable API key and invoke it per request,
so the token is simply always fresh; it is cached until shortly before expiry.
Output contract: ONLY the token on stdout, bare or as JSON with an
``access_token`` field (``expires_in`` honoured when present).
Precedence: an explicit ``--api-key`` still wins (one-off recovery escape
hatch); otherwise ``key_cmd`` beats a static ``api_key`` / ``key_env``.
Enterprise gateways (SSO/OIDC brokers, cloud IAM, auth proxies) issue SHORT-LIVED bearers; a key
copied into ``.env`` goes stale within the hour. ``key_cmd`` names a command that PRINTS a token
(the ``apiKeyHelper`` / ``gcloud auth print-access-token`` idiom). Both wire clients accept a
callable API key and invoke it per request; the token is cached until shortly before expiry.
Output contract: ONLY the token on stdout, bare or as JSON with an ``access_token`` field
(``expires_in`` / ISO ``expiry`` honoured). Precedence: explicit ``--api-key`` wins (one-off
recovery escape hatch); otherwise ``key_cmd`` beats a static ``api_key`` / ``key_env``.
"""
from __future__ import annotations
@@ -32,15 +20,13 @@ from typing import Callable, Optional
logger = logging.getLogger(__name__)
# Treat a token as spent slightly before expiry so a request can't be signed
# with one that dies in flight (60s = usual OAuth cache leeway).
# Treat a token as spent slightly before expiry so a request can't be signed with one that dies in
# flight (60s = usual OAuth cache leeway).
_TOKEN_REFRESH_LEEWAY_SECONDS = 60.0
# Helpers answer from a local cache in milliseconds; this long means hung.
_MINT_TIMEOUT_SECONDS = 15
# No advertised expiry: nothing in the request path re-mints on 401 (the SDK
# retries 429/5xx only), so a process-lifetime cache would 401 forever once the
# token died. Re-mint on a bounded window instead; helpers wanting a longer
# cache advertise their real expiry.
# No advertised expiry: nothing in the request path re-mints on 401 (the SDK retries 429/5xx only), so
# a process-lifetime cache would 401 forever once the token died. Re-mint on a bounded window instead.
_NO_TTL_REFRESH_SECONDS = 900.0
@@ -52,25 +38,18 @@ def _mint(command: str, label: str) -> tuple[str, Optional[float]]:
"""Run *command*, returning ``(token, ttl_seconds_or_None)``."""
try:
completed = subprocess.run(
command,
shell=True,
capture_output=True,
text=True,
timeout=_MINT_TIMEOUT_SECONDS,
command, shell=True, capture_output=True, text=True, timeout=_MINT_TIMEOUT_SECONDS,
)
except subprocess.TimeoutExpired as exc:
raise CommandTokenError(
f"key_cmd for provider {label!r} timed out after "
f"{_MINT_TIMEOUT_SECONDS}s"
f"key_cmd for provider {label!r} timed out after {_MINT_TIMEOUT_SECONDS}s"
) from exc
except OSError as exc:
raise CommandTokenError(
f"key_cmd for provider {label!r} could not be executed: {exc}"
) from exc
raise CommandTokenError(f"key_cmd for provider {label!r} could not be executed: {exc}") from exc
if completed.returncode != 0:
# NEVER include stdout/stderr (may hold a token) or the command string
# (may embed `--client-secret=…`); name the provider instead.
# NEVER include stdout/stderr (may hold a token) or the command string (may embed
# `--client-secret=…`); name the provider instead.
raise CommandTokenError(
f"key_cmd for provider {label!r} exited {completed.returncode}. "
f"Run that provider's key_cmd manually to see why "
@@ -91,28 +70,24 @@ def _mint(command: str, label: str) -> tuple[str, Optional[float]]:
token = str(payload.get("access_token") or "").strip()
if not token:
raise CommandTokenError(
f"key_cmd for provider {label!r} returned JSON without an "
"'access_token' field"
f"key_cmd for provider {label!r} returned JSON without an 'access_token' field"
)
ttl = payload.get("expires_in")
if isinstance(ttl, (int, float)) and ttl > 0:
return token, float(ttl)
# CLI helpers often print an absolute ISO 8601 deadline instead of
# OAuth's relative lifetime; honour it or the token 401s once past.
# Lazy import: hermes_cli.auth imports agent.* at module level.
# CLI helpers often print an absolute ISO 8601 deadline instead of OAuth's relative
# lifetime; honour it or the token 401s once past. Lazy import: hermes_cli.auth imports agent.*.
from hermes_cli.auth import _parse_iso_timestamp
for field in ("expiry", "expiresOn"):
deadline = _parse_iso_timestamp(payload.get(field))
if deadline is not None:
remaining = deadline - time.time()
if remaining > 0:
return token, remaining
remaining = deadline - time.time() if deadline is not None else 0
if remaining > 0:
return token, remaining
return token, None
# Bare token: stdout carries the token and nothing else. Do NOT keep one
# line of several — that turns a misconfigured helper (banner, warning) into
# a corrupt-key 401 far harder to diagnose than an explicit refusal.
# Bare token: stdout carries the token and nothing else. Do NOT keep one line of several — that
# turns a misconfigured helper (banner, warning) into a corrupt-key 401 far harder to diagnose.
token = stdout.strip()
if "\n" in token:
raise CommandTokenError(
@@ -148,12 +123,7 @@ class CommandTokenSource:
return token
def build_command_token_provider(
key_cmd: str,
provider_label: str = "custom",
) -> Optional[Callable[[], str]]:
def build_command_token_provider(key_cmd: str, provider_label: str = "custom") -> Optional[Callable[[], str]]:
"""A per-request token provider for *key_cmd*, or ``None`` when unset."""
command = str(key_cmd or "").strip()
if not command:
return None
return CommandTokenSource(command, provider_label)
return CommandTokenSource(command, provider_label) if command else None