diff --git a/agent/battery.py b/agent/battery.py index 9343c69d27..043fa1c015 100644 --- a/agent/battery.py +++ b/agent/battery.py @@ -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}%" diff --git a/agent/bounded_response.py b/agent/bounded_response.py index 4ee32dfcb0..d0d5ac8ed8 100644 --- a/agent/bounded_response.py +++ b/agent/bounded_response.py @@ -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") diff --git a/agent/browser_provider.py b/agent/browser_provider.py index 1ea59a40a3..5b975543d9 100644 --- a/agent/browser_provider.py +++ b/agent/browser_provider.py @@ -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 ``/plugins/browser//`` -(built-in) or ``~/.hermes/plugins/browser//`` (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 +``/plugins/browser//`` (built-in) or ``~/.hermes/plugins/browser//`` (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`.""" diff --git a/agent/browser_registry.py b/agent/browser_registry.py index 93a88dd011..9c07aceb76 100644 --- a/agent/browser_registry.py +++ b/agent/browser_registry.py @@ -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 - diff --git a/agent/command_token_source.py b/agent/command_token_source.py index bebff150b6..6c7d5cdb5b 100644 --- a/agent/command_token_source.py +++ b/agent/command_token_source.py @@ -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