From 38484b5be706ccbd4c87cd395ab847fed505f980 Mon Sep 17 00:00:00 2001 From: "ang-fleet-workers[bot]" <333956806+ang-fleet-workers[bot]@users.noreply.github.com> Date: Wed, 30 Sep 2026 12:39:01 -0700 Subject: [PATCH 1/6] wip: web search dead-backend circuit breaker (t_0f3cc21c) From 80fa78aa44eeed1cd7c577a3bedf31feab852f0e Mon Sep 17 00:00:00 2001 From: "ang-fleet-workers[bot]" <333956806+ang-fleet-workers[bot]@users.noreply.github.com> Date: Wed, 30 Sep 2026 12:44:37 -0700 Subject: [PATCH 2/6] feat(web): circuit-break 402/401-dead web backends, page once per episode A keyed web backend that answers HTTP 402 (out of credits) or 401 (bad key) was retried on every web_search/web_extract call: a doomed round-trip plus a WARNING per call before the fallback chain served it. Measured over 7 days: every one of 177 clanker + 330 athena + 30 other firecrawl searches was a 402 served by the keyless rescue. tools/web_backend_breaker.py: a 402/401-class failure opens a per-backend breaker for web.dead_backend_cooldown_seconds (default 3600, 0 = off). While open the backend is skipped with no network call and the existing fallback chain / keyless rescue serves the call. After the cooldown one call probes it; another 402 re-opens quietly, a success closes the episode. One WARNING at episode start, one INFO at recovery; the per-call rescue WARNING drops to DEBUG while the breaker is open. web.dead_backend_alert_command runs once per episode (env WEB_BACKEND/_STATUS/_ERROR), marked paged only on rc 0 so a failed page retries on the next trip. State is a flocked JSON file under the profile home state/, shared by every process of the profile. 429/5xx/timeouts do not trip. Wired into the search primary + keyed fallbacks and the extract primary + keyed fallbacks. Verified: tests/tools/test_web_backend_breaker.py 15 passed; the 3 wiring tests fail against the pre-change web_tools.py (RED proof). test_web_keyed_fallbacks.py + test_web_keyless_rescue.py still pass. Card: t_0f3cc21c --- hermes_cli/config_defaults.py | 10 + tests/tools/test_web_backend_breaker.py | 206 +++++++++++++++ tools/web_backend_breaker.py | 250 ++++++++++++++++++ tools/web_tools.py | 75 ++++-- website/docs/user-guide/configuration.md | 8 + .../docs/user-guide/features/web-search.md | 2 + 6 files changed, 535 insertions(+), 16 deletions(-) create mode 100644 tests/tools/test_web_backend_breaker.py create mode 100644 tools/web_backend_breaker.py diff --git a/hermes_cli/config_defaults.py b/hermes_cli/config_defaults.py index e37f9fe938242..863023f2b9d8f 100644 --- a/hermes_cli/config_defaults.py +++ b/hermes_cli/config_defaults.py @@ -621,6 +621,16 @@ # free-tier ring — the next call attempts the chosen backend again # (no sticky failover). Off when keyless_fallback is false. "keyless_rescue": True, + # Dead-backend circuit breaker: a keyed backend answering HTTP 402 + # (out of credits) or 401 (bad key) is skipped for this many seconds + # instead of being retried on every call; the fallback chain serves + # meanwhile. One WARNING per episode (first failure to next success). + # 0 disables the breaker. + "dead_backend_cooldown_seconds": 3600, + # Shell command run ONCE per dead-backend episode (e.g. a pager). + # Env: WEB_BACKEND, WEB_BACKEND_STATUS (402|401), WEB_BACKEND_ERROR. + # A nonzero exit is retried on the next trip. Empty = log only. + "dead_backend_alert_command": "", # Per-provider tier selection for ring vendors with both a keyless # free endpoint and a keyed paid path (exa, parallel, tavily, # firecrawl, keenable). Set by the `hermes tools` picker's diff --git a/tests/tools/test_web_backend_breaker.py b/tests/tools/test_web_backend_breaker.py new file mode 100644 index 0000000000000..f5be2c51fd72b --- /dev/null +++ b/tests/tools/test_web_backend_breaker.py @@ -0,0 +1,206 @@ +"""Dead-backend (402/401) circuit breaker for web_search / web_extract.""" +import json +import logging +import time + +import pytest + +from tools import web_backend_breaker as bb +from tools import web_tools + +PAYMENT = ("Payment Required: Failed to search. Insufficient credits to perform " + "this request.") + + +class Provider: + def __init__(self, name, calls, *, search=None, extract=None): + self.name = name + self.calls = calls + self._search = search + self._extract = extract + + def supports_search(self): + return self._search is not None + + def supports_extract(self): + return self._extract is not None + + def is_available(self): + return True + + def search(self, query, limit): + self.calls.append(self.name) + return dict(self._search) + + def extract(self, urls, **kwargs): + self.calls.append(self.name) + return [dict(r, url=u) for r, u in zip(self._extract * len(urls), urls)] + + +@pytest.fixture +def setup(monkeypatch, tmp_path): + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + cfg = { + "search_backend": "firecrawl", "extract_backend": "firecrawl", + "search_fallbacks": ["exa"], "extract_fallbacks": ["exa"], + "dead_backend_cooldown_seconds": 3600, + } + monkeypatch.setattr("agent.web_search_provider.get_provider_env", + lambda key: "test-key" if key == "EXA_API_KEY" else "") + monkeypatch.setattr(web_tools, "_ensure_web_plugins_loaded", lambda: None) + monkeypatch.setattr(web_tools, "_load_web_config", lambda: cfg) + monkeypatch.setattr(bb, "_web_config", lambda: cfg) + + async def safe(url): + return True + + monkeypatch.setattr(web_tools, "async_is_safe_url", safe) + calls, providers = [], {} + monkeypatch.setattr("agent.web_search_registry.get_provider", providers.get) + return calls, providers, cfg, tmp_path + + +def _search(q): + return json.loads(web_tools.web_search_tool(f"{q} {time.time_ns()}")) + + +@pytest.mark.parametrize("err,status", [ + (PAYMENT, 402), + ("HTTP 402", 402), + ("401 Unauthorized: invalid api key", 401), + ("Invalid API key provided", 401), + ("HTTP 429 Too Many Requests", None), + ("HTTP 500 upstream error", None), + ("timed out", None), + ("", None), +]) +def test_classify_dead(err, status): + assert bb.classify_dead(err) == status + + +def test_402_opens_breaker_skips_backend_and_logs_once(setup, caplog): + calls, providers, _cfg, _home = setup + ok = {"success": True, "data": {"web": [{"url": "https://ok.example"}]}} + providers.update(firecrawl=Provider("firecrawl", calls, search={"success": False, "error": PAYMENT}), + exa=Provider("exa", calls, search=ok)) + caplog.set_level(logging.DEBUG, logger="tools.web_backend_breaker") + for i in range(3): + assert _search(f"q{i}")["success"] is True + # Firecrawl was tried exactly once; the next two calls went straight to exa. + assert calls == ["firecrawl", "exa", "exa", "exa"] + warnings = [r for r in caplog.records + if r.levelno == logging.WARNING and "is dead" in r.getMessage()] + assert len(warnings) == 1 + + +def test_cooldown_expiry_reprobes_quietly_then_success_closes(setup, monkeypatch, caplog): + calls, providers, _cfg, home = setup + ok = {"success": True, "data": {"web": [{"url": "https://ok.example"}]}} + fc = Provider("firecrawl", calls, search={"success": False, "error": PAYMENT}) + providers.update(firecrawl=fc, exa=Provider("exa", calls, search=ok)) + caplog.set_level(logging.DEBUG, logger="tools.web_backend_breaker") + now = [1_000_000.0] + monkeypatch.setattr(bb.time, "time", lambda: now[0]) + + _search("a") + now[0] += 3601 # cooldown lapsed: next call probes firecrawl again + _search("b") + assert calls == ["firecrawl", "exa", "firecrawl", "exa"] + assert sum("is dead" in r.getMessage() for r in caplog.records) == 1 + + now[0] += 3601 + fc._search = ok # credits refilled + calls.clear() + assert _search("c")["success"] is True + assert calls == ["firecrawl"] + state = json.loads((home / "state" / "web_backend_breaker.json").read_text()) + assert "firecrawl" not in state + assert any("recovered" in r.getMessage() for r in caplog.records) + + +def test_transient_errors_do_not_trip(setup): + calls, providers, _cfg, _home = setup + ok = {"success": True, "data": {"web": [{"url": "https://ok.example"}]}} + providers.update(firecrawl=Provider("firecrawl", calls, search={"success": False, "error": "HTTP 429"}), + exa=Provider("exa", calls, search=ok)) + _search("x") + _search("y") + assert calls == ["firecrawl", "exa", "firecrawl", "exa"] + + +def test_cooldown_zero_disables_breaker(setup): + calls, providers, cfg, _home = setup + cfg["dead_backend_cooldown_seconds"] = 0 + ok = {"success": True, "data": {"web": [{"url": "https://ok.example"}]}} + providers.update(firecrawl=Provider("firecrawl", calls, search={"success": False, "error": PAYMENT}), + exa=Provider("exa", calls, search=ok)) + _search("x") + _search("y") + assert calls == ["firecrawl", "exa", "firecrawl", "exa"] + + +def _wait_lines(path, n, timeout=10.0): + deadline = time.time() + timeout + while time.time() < deadline: + if path.exists() and len(path.read_text().splitlines()) >= n: + break + time.sleep(0.05) + return path.read_text().splitlines() if path.exists() else [] + + +def _wait_paged(home, backend, timeout=10.0): + deadline = time.time() + timeout + while time.time() < deadline: + state = json.loads((home / "state" / "web_backend_breaker.json").read_text()) + if "paging" not in state.get(backend, {}): + return state[backend] + time.sleep(0.05) + raise AssertionError("alert thread never finished") + + +def test_alert_command_fires_once_per_episode(setup, tmp_path): + _calls, _providers, cfg, home = setup + out = tmp_path / "pages.txt" + cfg["dead_backend_alert_command"] = f'echo "$WEB_BACKEND $WEB_BACKEND_STATUS" >> {out}' + t0 = 2_000_000.0 + assert bb.record_failure("firecrawl", PAYMENT, now=t0) is True + assert _wait_lines(out, 1) == ["firecrawl 402"] + _wait_paged(home, "firecrawl") + # Re-trips in the same episode (e.g. hourly re-probes) do not page again. + assert bb.record_failure("firecrawl", PAYMENT, now=t0 + 3601) is False + assert bb.record_failure("firecrawl", PAYMENT, now=t0 + 7202) is False + time.sleep(0.3) + assert out.read_text().splitlines() == ["firecrawl 402"] + # Recovery closes the episode; a later death is a new episode and pages. + bb.record_success("firecrawl") + assert bb.record_failure("firecrawl", PAYMENT, now=t0 + 9000) is True + assert _wait_lines(out, 2) == ["firecrawl 402", "firecrawl 402"] + + +def test_failed_alert_is_retried_on_next_trip(setup, tmp_path): + _calls, _providers, cfg, home = setup + out = tmp_path / "attempts.txt" + cfg["dead_backend_alert_command"] = f"echo try >> {out}; exit 3" + t0 = 3_000_000.0 + bb.record_failure("firecrawl", PAYMENT, now=t0) + assert _wait_lines(out, 1) == ["try"] + entry = _wait_paged(home, "firecrawl") + assert not entry.get("paged") + cfg["dead_backend_alert_command"] = f"echo ok >> {out}" + bb.record_failure("firecrawl", PAYMENT, now=t0 + 3601) + assert _wait_lines(out, 2) == ["try", "ok"] + assert _wait_paged(home, "firecrawl").get("paged") is True + + +@pytest.mark.asyncio +async def test_extract_whole_batch_402_trips_breaker(setup): + calls, providers, _cfg, _home = setup + urls = ["https://a.example", "https://b.example"] + providers.update( + firecrawl=Provider("firecrawl", calls, extract=[{"content": "", "error": PAYMENT}]), + exa=Provider("exa", calls, extract=[{"content": "ok", "error": None}]), + ) + out1 = json.loads(await web_tools.web_extract_tool(urls)) + out2 = json.loads(await web_tools.web_extract_tool([u + "/2" for u in urls])) + assert calls == ["firecrawl", "exa", "exa"] + assert all(not r.get("error") for r in out1["results"] + out2["results"]) diff --git a/tools/web_backend_breaker.py b/tools/web_backend_breaker.py new file mode 100644 index 0000000000000..28b6889094a1f --- /dev/null +++ b/tools/web_backend_breaker.py @@ -0,0 +1,250 @@ +"""Circuit breaker for web backends that are dead for billing/auth reasons. + +A keyed web backend answering HTTP 402 (out of credits) or 401 (bad key) will +answer the same way on the next call, and the one after that. Without a +breaker every ``web_search`` pays a doomed round-trip plus a WARNING line +before the fallback chain serves the call. + +Behaviour: + +* A 402/401-class failure opens the breaker for ``web.dead_backend_cooldown_seconds`` + (default 3600). While open, the backend is skipped (no network call) and the + normal fallback chain / keyless rescue serves the call. +* When the cooldown lapses, the next call probes the backend once. Another + 402/401 re-opens it quietly; a success closes the episode. +* An EPISODE runs from the first dead response to the next success. It logs one + WARNING at the start and one INFO at recovery, and runs + ``web.dead_backend_alert_command`` at most once (it retries on a later trip + only if the command exited nonzero). + +State is a small JSON file under ``$HERMES_HOME/state`` so every process of a +profile (gateway, warm clients, subagents) shares one episode and one page. +""" + +from __future__ import annotations + +import json +import logging +import os +import re +import subprocess +import threading +import time +from contextlib import contextmanager +from pathlib import Path +from typing import Any, Dict, Optional + +logger = logging.getLogger(__name__) + +DEFAULT_COOLDOWN_SECONDS = 3600 +_ALERT_TIMEOUT_SECONDS = 30 +_STATE_FILENAME = "web_backend_breaker.json" + +# Billing (402) and auth (401) failures. Rate limits (429) and 5xx are +# transient and deliberately NOT matched. +_DEAD_PATTERNS = ( + (re.compile(r"\b402\b|payment required|insufficient credits|out of credits" + r"|credits? (?:exhausted|depleted)|quota exceeded for (?:this|your) (?:plan|account)", + re.IGNORECASE), 402), + (re.compile(r"\b401\b|unauthori[sz]ed|invalid api key|invalid_api_key" + r"|api key (?:is )?invalid", re.IGNORECASE), 401), +) + + +def classify_dead(error: Any) -> Optional[int]: + """Return 402 or 401 when *error* reads as a billing/auth death, else None.""" + text = str(error or "") + if not text: + return None + for pattern, status in _DEAD_PATTERNS: + if pattern.search(text): + return status + return None + + +def _web_config() -> Dict[str, Any]: + try: + from hermes_cli.config import load_config + + cfg = load_config().get("web", {}) + return cfg if isinstance(cfg, dict) else {} + except Exception: # noqa: BLE001 — config optional here + return {} + + +def cooldown_seconds(cfg: Optional[Dict[str, Any]] = None) -> int: + cfg = _web_config() if cfg is None else cfg + try: + return max(0, int(cfg.get("dead_backend_cooldown_seconds", DEFAULT_COOLDOWN_SECONDS))) + except (TypeError, ValueError): + return DEFAULT_COOLDOWN_SECONDS + + +def _state_path() -> Path: + from hermes_constants import get_hermes_home + + return get_hermes_home() / "state" / _STATE_FILENAME + + +_thread_lock = threading.Lock() + + +@contextmanager +def _locked_state(): + """Yield the mutable state dict under a thread + file lock; persist on exit.""" + path = _state_path() + with _thread_lock: + path.parent.mkdir(parents=True, exist_ok=True) + lock_path = path.with_suffix(".lock") + with open(lock_path, "a+", encoding="utf-8") as lock_fh: + try: + import fcntl + + fcntl.flock(lock_fh.fileno(), fcntl.LOCK_EX) + except (ImportError, OSError): # Windows / odd FS: thread lock only + pass + try: + state = json.loads(path.read_text(encoding="utf-8")) if path.exists() else {} + if not isinstance(state, dict): + state = {} + except (OSError, ValueError): + state = {} + before = json.dumps(state, sort_keys=True) + yield state + if json.dumps(state, sort_keys=True) != before: + tmp = path.with_suffix(f".tmp.{os.getpid()}") + tmp.write_text(json.dumps(state, indent=2, sort_keys=True), encoding="utf-8") + os.replace(tmp, path) + + +def open_until(backend: str, now: Optional[float] = None) -> Optional[float]: + """Return the epoch the breaker for *backend* reopens, or None when closed.""" + if not backend: + return None + now = time.time() if now is None else now + try: + with _locked_state() as state: + entry = state.get(backend) + until = float(entry.get("open_until", 0)) if isinstance(entry, dict) else 0.0 + except Exception as exc: # noqa: BLE001 — breaker must never break search + logger.debug("web backend breaker read failed: %s", exc) + return None + return until if until > now else None + + +def skip_error(backend: str, until: float) -> str: + """Error text for a call skipped because the breaker is open.""" + entry = _entry(backend) or {} + reopen = time.strftime("%H:%M", time.localtime(until)) + return (f"backend '{backend}' skipped: circuit open until {reopen} after HTTP " + f"{entry.get('status', '?')} ({str(entry.get('error', ''))[:160]})") + + +def _entry(backend: str) -> Optional[Dict[str, Any]]: + try: + with _locked_state() as state: + entry = state.get(backend) + return dict(entry) if isinstance(entry, dict) else None + except Exception: # noqa: BLE001 + return None + + +def record_failure(backend: str, error: Any, now: Optional[float] = None, + cfg: Optional[Dict[str, Any]] = None) -> bool: + """Trip the breaker when *error* is a 402/401 death. + + Returns True when this call started a NEW episode. + """ + status = classify_dead(error) + if status is None or not backend: + return False + cfg = _web_config() if cfg is None else cfg + cooldown = cooldown_seconds(cfg) + if cooldown <= 0: + return False + now = time.time() if now is None else now + new_episode = False + need_page = False + try: + with _locked_state() as state: + entry = state.get(backend) + if not isinstance(entry, dict): + entry = {"episode_start": now, "paged": False} + new_episode = True + entry.update(status=status, error=str(error)[:300], open_until=now + cooldown, + last_trip=now, trips=int(entry.get("trips", 0)) + 1) + state[backend] = entry + # "paging" marks an alert in flight; a stale mark (the process + # died mid-send) must not suppress the page forever. + in_flight = now - float(entry.get("paging") or 0) < 4 * _ALERT_TIMEOUT_SECONDS + need_page = not entry.get("paged") and not in_flight + if need_page: + entry["paging"] = now + except Exception as exc: # noqa: BLE001 + logger.debug("web backend breaker write failed: %s", exc) + return False + if new_episode: + logger.warning( + "web backend '%s' is dead (HTTP %s: %s); skipping it for %ds. " + "Further failures this episode are not logged per call.", + backend, status, str(error)[:200], cooldown, + ) + else: + logger.debug("web backend '%s' still dead (HTTP %s); breaker re-opened", backend, status) + if need_page: + _fire_alert(backend, status, str(error), cfg) + return new_episode + + +def record_success(backend: str) -> None: + """Close an open episode for *backend* (no-op when none is open).""" + if not backend: + return + try: + with _locked_state() as state: + entry = state.pop(backend, None) + except Exception as exc: # noqa: BLE001 + logger.debug("web backend breaker reset failed: %s", exc) + return + if isinstance(entry, dict): + logger.info("web backend '%s' recovered; dead-backend episode closed", backend) + + +def _fire_alert(backend: str, status: int, error: str, cfg: Dict[str, Any]) -> None: + command = str(cfg.get("dead_backend_alert_command") or "").strip() + if not command: + _mark_paged(backend, ok=True, note="no alert command configured") + return + + def _run() -> None: + env = dict(os.environ, WEB_BACKEND=backend, WEB_BACKEND_STATUS=str(status), + WEB_BACKEND_ERROR=error[:300]) + try: + proc = subprocess.run(command, shell=True, env=env, capture_output=True, + text=True, timeout=_ALERT_TIMEOUT_SECONDS) + ok = proc.returncode == 0 + note = f"rc={proc.returncode} {(proc.stderr or '').strip()[:200]}" + except Exception as exc: # noqa: BLE001 + ok, note = False, f"{type(exc).__name__}: {exc}" + if ok: + logger.info("web backend '%s' dead-backend alert sent", backend) + else: + logger.error("web backend '%s' dead-backend alert FAILED (%s); " + "will retry on the next trip", backend, note) + _mark_paged(backend, ok=ok, note=note) + + threading.Thread(target=_run, name="web-backend-alert", daemon=True).start() + + +def _mark_paged(backend: str, *, ok: bool, note: str) -> None: + try: + with _locked_state() as state: + entry = state.get(backend) + if not isinstance(entry, dict): + return + entry.pop("paging", None) + if ok: + entry["paged"] = True + entry["paged_note"] = note[:200] + except Exception as exc: # noqa: BLE001 + logger.debug("web backend breaker page mark failed: %s", exc) diff --git a/tools/web_tools.py b/tools/web_tools.py index 7cf9e20fc5f59..3f46effd2709d 100644 --- a/tools/web_tools.py +++ b/tools/web_tools.py @@ -475,6 +475,57 @@ def _failed_extract_batch(results: list, urls: list) -> bool: and not all(_policy_blocked_result(r) for r in results)) +# ─── Dead-backend circuit breaker (tools/web_backend_breaker.py) ───────────── + +def _breaker_search(provider, query: str, limit: int) -> dict: + """``provider.search`` behind the 402/401 breaker. + + An open breaker skips the network call and returns a failure, so the + normal fallback chain / keyless rescue serves the call without first + paying a doomed round-trip. + """ + from tools import web_backend_breaker as _bb + + until = _bb.open_until(provider.name) + if until: + return {"success": False, "error": _bb.skip_error(provider.name, until)} + try: + resp = provider.search(query, limit) + except Exception as exc: + _bb.record_failure(provider.name, exc, cfg=_load_web_config()) + raise + if resp.get("success"): + _bb.record_success(provider.name) + else: + _bb.record_failure(provider.name, resp.get("error"), cfg=_load_web_config()) + return resp + + +async def _breaker_extract(provider, urls: list, format) -> list: + """``provider.extract`` behind the 402/401 breaker (whole-batch semantics).""" + import inspect + + from tools import web_backend_breaker as _bb + + until = _bb.open_until(provider.name) + if until: + err = _bb.skip_error(provider.name, until) + return [{"url": u, "title": "", "content": "", "error": err} for u in urls] + try: + if inspect.iscoroutinefunction(provider.extract): + results = await provider.extract(urls, format=format) + else: + results = await asyncio.to_thread(provider.extract, urls, format=format) + except Exception as exc: + _bb.record_failure(provider.name, exc, cfg=_load_web_config()) + raise + if results and all(r.get("error") for r in results): + _bb.record_failure(provider.name, results[0].get("error"), cfg=_load_web_config()) + elif results: + _bb.record_success(provider.name) + return results + + # ─── One-shot keyless rescue (keyed/configured backend failed) ─────────────── def _keyless_rescue_enabled() -> bool: @@ -540,8 +591,11 @@ def _rescue_search(provider_name: str, original_error: str, query: str, limit: i user) can see the configured backend needs attention. """ from plugins.web.keyless_mcp import search_with_failover + from tools import web_backend_breaker as _bb - logger.warning( + # A dead (402/401) backend already logged once for its episode. + logger.log( + logging.DEBUG if _bb.open_until(provider_name) else logging.WARNING, "web_search backend '%s' failed (%s); one-shot keyless rescue", provider_name, (original_error or "")[:200], ) @@ -1031,7 +1085,7 @@ def _paid_search() -> tuple[dict, bool]: _rescued = False current = provider try: - _resp = current.search(query, _fetch_limit) + _resp = _breaker_search(current, query, _fetch_limit) except Exception as exc: # noqa: BLE001 — candidate for fallback _resp = {"success": False, "error": str(exc)} if not _load_web_config().get("search_fallbacks") and not _rescue_eligible(current): @@ -1045,7 +1099,7 @@ def _paid_search() -> tuple[dict, bool]: ) current = candidate try: - _resp = candidate.search(query, _fetch_limit) + _resp = _breaker_search(candidate, query, _fetch_limit) except Exception as exc: # noqa: BLE001 — next keyed candidate _resp = {"success": False, "error": str(exc)} if current is not provider: @@ -1374,16 +1428,10 @@ async def web_extract_tool( # Async-or-sync dispatch: parallel + firecrawl have async # extract(); exa + tavily are sync. - import inspect _extract_rescued = False current = provider try: - if inspect.iscoroutinefunction(current.extract): - results = await current.extract(fetch_urls, format=format) - else: - results = await asyncio.to_thread( - current.extract, fetch_urls, format=format - ) + results = await _breaker_extract(current, fetch_urls, format) except Exception as exc: # noqa: BLE001 — candidate for fallback if not _load_web_config().get("extract_fallbacks") and not _rescue_eligible(current): raise @@ -1402,12 +1450,7 @@ async def web_extract_tool( ) current = candidate try: - if inspect.iscoroutinefunction(current.extract): - results = await current.extract(fetch_urls, format=format) - else: - results = await asyncio.to_thread( - current.extract, fetch_urls, format=format - ) + results = await _breaker_extract(current, fetch_urls, format) except Exception as exc: # noqa: BLE001 — next keyed candidate results = [{"url": u, "title": "", "content": "", "error": str(exc)} for u in fetch_urls] diff --git a/website/docs/user-guide/configuration.md b/website/docs/user-guide/configuration.md index f4e0ed8445392..b88c1fc351328 100644 --- a/website/docs/user-guide/configuration.md +++ b/website/docs/user-guide/configuration.md @@ -2409,6 +2409,14 @@ web: # call attempts the chosen backend again (never sticky). keyless_rescue: true + # Dead-backend circuit breaker (default 3600 s). A backend answering + # HTTP 402 (out of credits) or 401 (bad key) is skipped for this long; + # the fallback chain serves meanwhile. One WARNING per episode. 0 = off. + dead_backend_cooldown_seconds: 3600 + # Optional command run once per dead-backend episode (e.g. a pager). + # Env: WEB_BACKEND, WEB_BACKEND_STATUS, WEB_BACKEND_ERROR. + dead_backend_alert_command: "" + # Pin Exa/Parallel to a tier (set by the hermes tools Free/Paid rows). # free = always the anonymous endpoint; paid = always the keyed SDK path; # unset = auto (key present -> paid, otherwise free). diff --git a/website/docs/user-guide/features/web-search.md b/website/docs/user-guide/features/web-search.md index ecc96ddbed202..e215f1ffc4749 100644 --- a/website/docs/user-guide/features/web-search.md +++ b/website/docs/user-guide/features/web-search.md @@ -425,6 +425,8 @@ If no backend has **ever** been selected (no `web.backend` / per-capability key **One-shot keyless rescue for keyed backends:** when your chosen/keyed backend fails a call (bad key, outage, upstream 5xx), that single call automatically retries on the keyless free-tier ring instead of erroring — the result notes which vendor served it and why (`rescued_from` / `backend_error`). The failover is never sticky: the very next `web_search`/`web_extract` call attempts your chosen backend again. Disable with `web.keyless_rescue: false` (also off whenever `keyless_fallback` is off). +**Dead-backend circuit breaker:** a keyed backend that answers HTTP 402 (out of credits) or 401 (bad key) is skipped for `web.dead_backend_cooldown_seconds` (default 3600) instead of being retried on every call, and logs one WARNING per episode rather than one per call. After the cooldown the next call probes it once; a success closes the episode. Set `web.dead_backend_alert_command` to run a pager command once per episode (env `WEB_BACKEND`, `WEB_BACKEND_STATUS`, `WEB_BACKEND_ERROR`). State lives in `$HERMES_HOME/state/web_backend_breaker.json`, shared by every process of the profile. + xAI Web Search is **not** in the auto-detection chain — having `XAI_API_KEY` set (or being signed in via xAI Grok OAuth) does not automatically route web traffic through xAI, since those credentials are also used for inference / TTS / image gen and the user may want a different backend for web. Opt in explicitly with `web.backend: "xai"`. --- From dccaef904d89862143cf3efede9adb62f665473a Mon Sep 17 00:00:00 2001 From: "ang-fleet-workers[bot]" <333956806+ang-fleet-workers[bot]@users.noreply.github.com> Date: Wed, 30 Sep 2026 13:14:48 -0700 Subject: [PATCH 3/6] fix(web): exclude exhausted vendors from keyless ring; scope alert episodes Clanker now uses searxng first and keyless only as the last resort; the profile sets web.keyless_exclude: [firecrawl] so even a ring rescue does not hit the exhausted vendor. The same per-profile opt-out applies to the registry's keyless resolution and the dispatch ring. Prism round 1 P1s: - Alert completion carries episode_start + alert_id; an old alert can no longer ack/clear a new episode or a new in-flight attempt. - Only vendor-boundary exceptions trip the extract breaker. Per-URL errors can be a paywalled target site's 401/402 and cannot disable a healthy search backend. Search error response failures still trip normally. - Capture the initiating ContextVar and Hermes home before spawning the alert thread; the pager receives the right profile's state/home and a minimal env rather than inheriting default-profile credentials. Verified: test_web_backend_breaker.py 18 passed; test_web_keyless_fallback.py 42 passed; ruff check clean. Card: t_0f3cc21c --- agent/web_search_registry.py | 9 +++- hermes_cli/config_defaults.py | 3 ++ plugins/web/keyless_mcp.py | 10 +++- tests/tools/test_web_backend_breaker.py | 63 +++++++++++++++++++++++- tests/tools/test_web_keyless_fallback.py | 7 +++ tools/web_backend_breaker.py | 43 ++++++++++++---- tools/web_tools.py | 7 +-- website/docs/user-guide/configuration.md | 3 ++ 8 files changed, 127 insertions(+), 18 deletions(-) diff --git a/agent/web_search_registry.py b/agent/web_search_registry.py index 87b2160b54ccf..37d0c531d2375 100644 --- a/agent/web_search_registry.py +++ b/agent/web_search_registry.py @@ -196,9 +196,14 @@ def _keyless_preference() -> tuple: from plugins.web.keyless_mcp import _KEYLESS_RING, _ring_cursor start = _ring_cursor % len(_KEYLESS_RING) + from tools.web_tools import _load_web_config + + excluded = _load_web_config().get("keyless_exclude", []) + if not isinstance(excluded, list): + excluded = [] return tuple( - _KEYLESS_RING[(start + i) % len(_KEYLESS_RING)] - for i in range(len(_KEYLESS_RING)) + vendor for i in range(len(_KEYLESS_RING)) + if (vendor := _KEYLESS_RING[(start + i) % len(_KEYLESS_RING)]) not in excluded ) except Exception as exc: # noqa: BLE001 — ring optional in stripped envs logger.debug("keyless ring order unavailable: %s", exc) diff --git a/hermes_cli/config_defaults.py b/hermes_cli/config_defaults.py index 611a19cce1533..cb44434a9be2f 100644 --- a/hermes_cli/config_defaults.py +++ b/hermes_cli/config_defaults.py @@ -616,6 +616,9 @@ # failing over to the next ring vendor on rate limits. Never # pre-empts a configured or keyed backend. Set false to disable. "keyless_fallback": True, + # Omit named vendors from the anonymous ring for this profile (e.g. + # an exhausted Firecrawl account); other free vendors still fail over. + "keyless_exclude": [], # One-shot keyless rescue: when the chosen/keyed backend fails a # web_search/web_extract call, THAT call retries once on the keyless # free-tier ring — the next call attempts the chosen backend again diff --git a/plugins/web/keyless_mcp.py b/plugins/web/keyless_mcp.py index c3c1947fd3360..cbcfa1752a236 100644 --- a/plugins/web/keyless_mcp.py +++ b/plugins/web/keyless_mcp.py @@ -799,7 +799,15 @@ def _ring_order(name: str) -> List[str]: _KEYLESS_RING[(start + i) % len(_KEYLESS_RING)] for i in range(len(_KEYLESS_RING)) ] - return [v for v in ordered if provider_tier(v) != "paid"] + try: + import tools.web_tools as _wt + + excluded = _wt._load_web_config().get("keyless_exclude", []) + if not isinstance(excluded, list): + excluded = [] + except Exception: # noqa: BLE001 — config optional during early discovery + excluded = [] + return [v for v in ordered if provider_tier(v) != "paid" and v not in excluded] def search_with_failover(name: str, query: str, limit: int = 5) -> Dict[str, Any]: diff --git a/tests/tools/test_web_backend_breaker.py b/tests/tools/test_web_backend_breaker.py index f5be2c51fd72b..c09df9806933b 100644 --- a/tests/tools/test_web_backend_breaker.py +++ b/tests/tools/test_web_backend_breaker.py @@ -193,7 +193,7 @@ def test_failed_alert_is_retried_on_next_trip(setup, tmp_path): @pytest.mark.asyncio -async def test_extract_whole_batch_402_trips_breaker(setup): +async def test_target_site_401_402_does_not_trip_shared_breaker(setup): calls, providers, _cfg, _home = setup urls = ["https://a.example", "https://b.example"] providers.update( @@ -202,5 +202,64 @@ async def test_extract_whole_batch_402_trips_breaker(setup): ) out1 = json.loads(await web_tools.web_extract_tool(urls)) out2 = json.loads(await web_tools.web_extract_tool([u + "/2" for u in urls])) - assert calls == ["firecrawl", "exa", "exa"] + assert calls == ["firecrawl", "exa", "firecrawl", "exa"] assert all(not r.get("error") for r in out1["results"] + out2["results"]) + assert bb.open_until("firecrawl") is None + + +def test_stale_page_completion_cannot_ack_new_episode(setup, monkeypatch): + _calls, _providers, cfg, home = setup + # Keep the alert worker queued so the first episode is removed before it + # finishes. Then a second episode begins. The old success must not ack it. + monkeypatch.setattr(bb, "_fire_alert", lambda *args: None) + cfg["dead_backend_alert_command"] = "true" + bb.record_failure("firecrawl", PAYMENT, now=1_000_000.0) + first = bb._entry("firecrawl") + bb.record_success("firecrawl") + bb.record_failure("firecrawl", PAYMENT, now=1_000_001.0) + second = bb._entry("firecrawl") + assert first["alert_id"] != second["alert_id"] + bb._mark_paged("firecrawl", first["episode_start"], first["alert_id"], + ok=True, note="old alert") + entry = bb._entry("firecrawl") + assert entry["alert_id"] == second["alert_id"] + assert not entry.get("paged") + assert "paging" in entry + + +def test_stale_alert_attempt_cannot_clear_new_attempt(setup, monkeypatch): + _calls, _providers, cfg, home = setup + monkeypatch.setattr(bb, "_fire_alert", lambda *args: None) + cfg["dead_backend_alert_command"] = "true" + bb.record_failure("firecrawl", PAYMENT, now=1_000_000.0) + first = bb._entry("firecrawl") + # The alert in flight has timed out and a new call attempts delivery. + bb.record_failure("firecrawl", PAYMENT, now=1_000_121.0) + second = bb._entry("firecrawl") + assert first["alert_id"] != second["alert_id"] + bb._mark_paged("firecrawl", first["episode_start"], first["alert_id"], + ok=False, note="old alert failed") + assert bb._entry("firecrawl")["alert_id"] == second["alert_id"] + assert "paging" in bb._entry("firecrawl") + + +def test_alert_thread_keeps_profile_home_and_scrubs_default_secrets(setup, monkeypatch): + _calls, _providers, cfg, home = setup + seen = [] + cfg["dead_backend_alert_command"] = "fake pager" + monkeypatch.setenv("UNRELATED_DEFAULT_PROFILE_SECRET", "do-not-copy") + + class Proc: + returncode = 0 + stderr = "" + + def fake_run(cmd, **kwargs): + seen.append(kwargs["env"]) + return Proc() + + monkeypatch.setattr(bb.subprocess, "run", fake_run) + bb.record_failure("firecrawl", PAYMENT) + _wait_paged(home, "firecrawl") + assert len(seen) == 1 + assert seen[0]["HERMES_HOME"] == str(home) + assert "UNRELATED_DEFAULT_PROFILE_SECRET" not in seen[0] diff --git a/tests/tools/test_web_keyless_fallback.py b/tests/tools/test_web_keyless_fallback.py index f2c12ec21b328..f46d3f54b256a 100644 --- a/tests/tools/test_web_keyless_fallback.py +++ b/tests/tools/test_web_keyless_fallback.py @@ -303,6 +303,13 @@ def test_keyless_ring_rotates_and_covers_all_vendors(self, fresh_registry, monke assert keyless_mcp._ring_order("tavily")[0] == "tavily" assert keyless_mcp._ring_order("tavily")[0] == "tavily" + def test_keyless_exclusion_removes_vendor_from_ring_and_resolution(self, fresh_registry, monkeypatch): + monkeypatch.setattr("tools.web_tools._load_web_config", + lambda: {"keyless_exclude": ["firecrawl"]}) + assert "firecrawl" not in keyless_mcp._ring_order("searxng") + assert "firecrawl" not in registry._keyless_preference() + assert all(keyless_mcp._ring_order("firecrawl")) # no empty ring + def test_registry_keyless_disabled_returns_none(self, fresh_registry, monkeypatch): monkeypatch.setattr(registry, "_read_config_key", lambda *p: None) monkeypatch.setattr(registry, "_keyless_tier_enabled", lambda: False) diff --git a/tools/web_backend_breaker.py b/tools/web_backend_breaker.py index 28b6889094a1f..6fea0140bd133 100644 --- a/tools/web_backend_breaker.py +++ b/tools/web_backend_breaker.py @@ -23,10 +23,12 @@ from __future__ import annotations +import contextvars import json import logging import os import re +import secrets import subprocess import threading import time @@ -180,6 +182,9 @@ def record_failure(backend: str, error: Any, now: Optional[float] = None, need_page = not entry.get("paged") and not in_flight if need_page: entry["paging"] = now + entry["alert_id"] = secrets.token_hex(8) + episode_start = entry["episode_start"] + alert_id = entry.get("alert_id", "") except Exception as exc: # noqa: BLE001 logger.debug("web backend breaker write failed: %s", exc) return False @@ -192,7 +197,7 @@ def record_failure(backend: str, error: Any, now: Optional[float] = None, else: logger.debug("web backend '%s' still dead (HTTP %s); breaker re-opened", backend, status) if need_page: - _fire_alert(backend, status, str(error), cfg) + _fire_alert(backend, status, str(error), cfg, episode_start, alert_id) return new_episode @@ -210,15 +215,29 @@ def record_success(backend: str) -> None: logger.info("web backend '%s' recovered; dead-backend episode closed", backend) -def _fire_alert(backend: str, status: int, error: str, cfg: Dict[str, Any]) -> None: +def _fire_alert(backend: str, status: int, error: str, cfg: Dict[str, Any], + episode_start: float, alert_id: str) -> None: command = str(cfg.get("dead_backend_alert_command") or "").strip() if not command: - _mark_paged(backend, ok=True, note="no alert command configured") + _mark_paged(backend, episode_start, alert_id, ok=True, + note="no alert command configured") return + from hermes_constants import get_hermes_home + + # The gateway can multiplex profiles with ContextVar-scoped homes. Capture + # the initiating context BEFORE spawning the alert thread and give the + # child only its profile home + path (never inherited default-profile + # credentials). A host-local pager can load its own credentials. + context = contextvars.copy_context() + profile_home = str(get_hermes_home()) + def _run() -> None: - env = dict(os.environ, WEB_BACKEND=backend, WEB_BACKEND_STATUS=str(status), - WEB_BACKEND_ERROR=error[:300]) + env = {"HOME": os.path.expanduser("~"), + "PATH": os.environ.get("PATH", "/usr/local/bin:/usr/bin:/bin"), + "HERMES_HOME": profile_home, + "WEB_BACKEND": backend, "WEB_BACKEND_STATUS": str(status), + "WEB_BACKEND_ERROR": error[:300]} try: proc = subprocess.run(command, shell=True, env=env, capture_output=True, text=True, timeout=_ALERT_TIMEOUT_SECONDS) @@ -231,17 +250,21 @@ def _run() -> None: else: logger.error("web backend '%s' dead-backend alert FAILED (%s); " "will retry on the next trip", backend, note) - _mark_paged(backend, ok=ok, note=note) + _mark_paged(backend, episode_start, alert_id, ok=ok, note=note) - threading.Thread(target=_run, name="web-backend-alert", daemon=True).start() + threading.Thread(target=lambda: context.run(_run), name="web-backend-alert", + daemon=True).start() -def _mark_paged(backend: str, *, ok: bool, note: str) -> None: +def _mark_paged(backend: str, episode_start: float, alert_id: str, + *, ok: bool, note: str) -> None: try: with _locked_state() as state: entry = state.get(backend) - if not isinstance(entry, dict): - return + if (not isinstance(entry, dict) + or entry.get("episode_start") != episode_start + or entry.get("alert_id") != alert_id): + return # stale alert from an earlier episode/attempt entry.pop("paging", None) if ok: entry["paged"] = True diff --git a/tools/web_tools.py b/tools/web_tools.py index 3f46effd2709d..7e569de48ae25 100644 --- a/tools/web_tools.py +++ b/tools/web_tools.py @@ -519,9 +519,10 @@ async def _breaker_extract(provider, urls: list, format) -> list: except Exception as exc: _bb.record_failure(provider.name, exc, cfg=_load_web_config()) raise - if results and all(r.get("error") for r in results): - _bb.record_failure(provider.name, results[0].get("error"), cfg=_load_web_config()) - elif results: + # Per-URL errors can be a target site's 401/402 (paywall/login), not + # the extract vendor's billing/auth status. Only an exception at the + # provider boundary above can trip the shared search/extract breaker. + if results and any(not r.get("error") for r in results): _bb.record_success(provider.name) return results diff --git a/website/docs/user-guide/configuration.md b/website/docs/user-guide/configuration.md index b88c1fc351328..c0980072a9381 100644 --- a/website/docs/user-guide/configuration.md +++ b/website/docs/user-guide/configuration.md @@ -2403,6 +2403,9 @@ web: # and no API keys present, web tools rotate across the Exa/Parallel/ # Tavily/Firecrawl/Keenable free tiers. Set false to disable. keyless_fallback: true + # Optional per-profile exclusions from the anonymous ring (not the + # explicit backend selection). E.g. skip an exhausted Firecrawl account: + keyless_exclude: [firecrawl] # One-shot keyless rescue (default: true). When the chosen/keyed backend # fails a call, that single call retries on the keyless ring; the next From fbd6217a3e8145d029f9967366d56517a0c6ac4e Mon Sep 17 00:00:00 2001 From: "ang-fleet-workers[bot]" <333956806+ang-fleet-workers[bot]@users.noreply.github.com> Date: Wed, 30 Sep 2026 13:18:58 -0700 Subject: [PATCH 4/6] fix(web): decode pager output explicitly as UTF-8 on Windows --- tools/web_backend_breaker.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/tools/web_backend_breaker.py b/tools/web_backend_breaker.py index 6fea0140bd133..c1d2ca423c4b9 100644 --- a/tools/web_backend_breaker.py +++ b/tools/web_backend_breaker.py @@ -240,7 +240,8 @@ def _run() -> None: "WEB_BACKEND_ERROR": error[:300]} try: proc = subprocess.run(command, shell=True, env=env, capture_output=True, - text=True, timeout=_ALERT_TIMEOUT_SECONDS) + text=True, encoding="utf-8", errors="replace", + timeout=_ALERT_TIMEOUT_SECONDS) ok = proc.returncode == 0 note = f"rc={proc.returncode} {(proc.stderr or '').strip()[:200]}" except Exception as exc: # noqa: BLE001 From 9984b14fc65f10ef4cafe858f7a9f5d0b36aca2d Mon Sep 17 00:00:00 2001 From: "ang-fleet-workers[bot]" <333956806+ang-fleet-workers[bot]@users.noreply.github.com> Date: Wed, 30 Sep 2026 13:28:17 -0700 Subject: [PATCH 5/6] test(web): cover both target-site auth and payment extract errors --- tests/tools/test_web_backend_breaker.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/tests/tools/test_web_backend_breaker.py b/tests/tools/test_web_backend_breaker.py index c09df9806933b..f85631820b938 100644 --- a/tests/tools/test_web_backend_breaker.py +++ b/tests/tools/test_web_backend_breaker.py @@ -192,12 +192,13 @@ def test_failed_alert_is_retried_on_next_trip(setup, tmp_path): assert _wait_paged(home, "firecrawl").get("paged") is True +@pytest.mark.parametrize("site_error", ["401 Unauthorized", "402 Payment Required"]) @pytest.mark.asyncio -async def test_target_site_401_402_does_not_trip_shared_breaker(setup): +async def test_target_site_401_402_does_not_trip_shared_breaker(setup, site_error): calls, providers, _cfg, _home = setup urls = ["https://a.example", "https://b.example"] providers.update( - firecrawl=Provider("firecrawl", calls, extract=[{"content": "", "error": PAYMENT}]), + firecrawl=Provider("firecrawl", calls, extract=[{"content": "", "error": site_error}]), exa=Provider("exa", calls, extract=[{"content": "ok", "error": None}]), ) out1 = json.loads(await web_tools.web_extract_tool(urls)) From 891a80e79b77c721271daa275800b459605ef14b Mon Sep 17 00:00:00 2001 From: "ang-fleet-workers[bot]" <333956806+ang-fleet-workers[bot]@users.noreply.github.com> Date: Wed, 30 Sep 2026 13:39:32 -0700 Subject: [PATCH 6/6] fix(web): disconnect pager child stdin to satisfy TUI process guard --- tools/web_backend_breaker.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/tools/web_backend_breaker.py b/tools/web_backend_breaker.py index c1d2ca423c4b9..5f7890c9f81cd 100644 --- a/tools/web_backend_breaker.py +++ b/tools/web_backend_breaker.py @@ -239,9 +239,9 @@ def _run() -> None: "WEB_BACKEND": backend, "WEB_BACKEND_STATUS": str(status), "WEB_BACKEND_ERROR": error[:300]} try: - proc = subprocess.run(command, shell=True, env=env, capture_output=True, - text=True, encoding="utf-8", errors="replace", - timeout=_ALERT_TIMEOUT_SECONDS) + proc = subprocess.run(command, shell=True, env=env, stdin=subprocess.DEVNULL, + capture_output=True, text=True, encoding="utf-8", + errors="replace", timeout=_ALERT_TIMEOUT_SECONDS) ok = proc.returncode == 0 note = f"rc={proc.returncode} {(proc.stderr or '').strip()[:200]}" except Exception as exc: # noqa: BLE001