From ffab36579ac1c6b37e4a7a772f5555c81bb93852 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 17 Sep 2026 20:34:54 +0530 Subject: [PATCH] refactor(gateway): idle gates live in run_idle_gates; row semantics live with loops/heartbeat MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The probe helpers had been appended to the gateway/run.py facade and the three watchers each carried their own copy of the "probe → fail open → offload" shell, with gateway code parsing loop/heartbeat rows itself. One sibling now owns the gate shape (_gate + off_loop_gate), and hermes_cli.loops.store_has_active_loop / hermes_cli.heartbeat.store_has_active_heartbeat own what an ACTIVE row is (heartbeat gains the _META_PREFIX loops already had). One fail-open layer per read: has_pending_handoffs is a bare bounded query, the gate catches. The loop watcher calls its executor hop directly — it already did so unconditionally for the scan, so the "runner without a hop" fallback there was dead. --- gateway/run.py | 32 --------- gateway/run_adapters.py | 31 ++------- gateway/run_goals.py | 34 ++-------- gateway/run_heartbeat_restore.py | 24 +------ gateway/run_idle_gates.py | 66 +++++++++++++++++++ hermes_cli/heartbeat.py | 20 +++++- hermes_cli/loops.py | 11 ++++ hermes_state_gateway.py | 9 +-- tests/gateway/test_heartbeat_watch_restore.py | 2 +- tests/gateway/test_loop_command.py | 2 +- 10 files changed, 111 insertions(+), 120 deletions(-) create mode 100644 gateway/run_idle_gates.py diff --git a/gateway/run.py b/gateway/run.py index a0fa7ae6c1..0eac13a0e5 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -1788,38 +1788,6 @@ async def _async_profile_runtime_scope(profile_home: "Path"): yield -def _profile_session_db_probe(profile_home: "Path"): - """The goals-cached SessionDB for *profile_home* with ONLY the HERMES_HOME contextvar - installed — no config parse, no secret hydration, no terminal policy. Idle-path gates - ("is there any work for this profile at all?") use it before paying for a full - ``_profile_runtime_scope`` entry on every watcher tick. None when unavailable.""" - from hermes_cli.goals import _get_session_db - from hermes_constants import reset_hermes_home_override, set_hermes_home_override - - token = set_hermes_home_override(str(profile_home)) - try: - return _get_session_db() - except Exception: - logger.debug("session-db probe failed for %s", profile_home, exc_info=True) - return None - finally: - reset_hermes_home_override(token) - - -def _profile_meta_rows(profile_home: "Path", prefix: str) -> Optional[list]: - """``list_meta_prefix(prefix)`` rows from the probe DB, or None when the store is unavailable or - the read fails. None means "cannot prove emptiness": every idle gate treats it as work present, so - a broken or migrating store can never suppress a heartbeat restore or a due loop (fail OPEN).""" - db = _profile_session_db_probe(profile_home) - if db is None: - return None - try: - return db.list_meta_prefix(prefix) - except Exception: - logger.debug("meta probe %r failed for %s", prefix, profile_home, exc_info=True) - return None - - def load_gateway_config_for_runner() -> "GatewayConfig": """Load gateway config for the process-level GatewayRunner. An UNSET ``multiplex_profiles`` is settled first by ``resolve_multiplex_mode`` (the default is on; the boot guard keeps a fleet that diff --git a/gateway/run_adapters.py b/gateway/run_adapters.py index 390ee51ed0..bb9aa3b9e9 100644 --- a/gateway/run_adapters.py +++ b/gateway/run_adapters.py @@ -35,27 +35,6 @@ if TYPE_CHECKING: # string annotations only; never imported at runtime (cycle) logger = logging.getLogger("gateway.run") -async def _watcher_has_pending_handoffs(runner: object, profile_home: "Path") -> bool: - """Idle gate for the handoff watcher's per-profile ticks: True when the profile's store - holds a pending handoff — or when the probe can't prove otherwise (fail OPEN: a probe - error must never suppress a dispatch). The SessionDB read goes through the runner's - executor hop (off the loop thread); runners without one (bare test stand-ins) skip the - gate entirely and keep the historical always-enter behavior.""" - - def _probe() -> bool: - try: - from gateway.run import _profile_session_db_probe - db = _profile_session_db_probe(profile_home) - return True if db is None else db.has_pending_handoffs() - except Exception: - return True - - offload = getattr(runner, "_run_in_executor_with_context", None) - if not callable(offload): - return True - return bool(await offload(_probe)) - - class GatewayAdapterLifecycleMixin: """Adapter lifecycle: connect/teardown, fatal recovery, reconnect watcher, multiplex profiles.""" @@ -478,6 +457,7 @@ class GatewayAdapterLifecycleMixin: → running), re-bind the home channel to the CLI session_id, dispatch a synthetic event, mark ``completed``/``failed``.""" from gateway.run import _async_profile_runtime_scope, _handoff_watch_scopes, _reclaim_stale + from gateway.run_idle_gates import off_loop_gate, profile_has_pending_handoff await asyncio.sleep(5) # let platforms connect before dispatching through them # Does _process_handoff accept the profile argument? Test stand-ins bind a one-arg callable. try: @@ -545,11 +525,10 @@ class GatewayAdapterLifecycleMixin: while self._running: try: for profile_name, profile_home in _handoff_watch_scopes(self): - # Idle gate: the scope entry re-parses the profile's config/secrets, - # so only pay it when the profile's store actually holds a pending - # handoff. The root poll (None) is unscoped and stays cheap. - if profile_home is not None and not await _watcher_has_pending_handoffs( - self, profile_home): + # Idle gate (run_idle_gates): skip the scope entry when the profile's store + # holds no pending handoff. The root poll (None) is unscoped and stays cheap. + if profile_home is not None and not await off_loop_gate( + self, lambda home=profile_home: profile_has_pending_handoff(home)): continue async with _scope(profile_home): await _tick(profile_name) diff --git a/gateway/run_goals.py b/gateway/run_goals.py index 80dbc16d10..b983db21de 100644 --- a/gateway/run_goals.py +++ b/gateway/run_goals.py @@ -24,30 +24,6 @@ if TYPE_CHECKING: # string annotations only; never imported at runtime (cycle) logger = logging.getLogger("gateway.run") -async def _watcher_has_active_loops(runner: object, profile_home) -> bool: - """Idle gate for the loop wakeup watcher's per-profile scans: True when the profile's - store holds an ACTIVE ``loop:*`` row — or when the probe can't prove otherwise (fail - OPEN: an unavailable store, a failing read or a corrupt row must never skip a due loop). The scan's SessionDB read goes through - the runner's executor hop (off the loop thread, #92413); runners without one (bare test - stand-ins) skip the gate entirely and keep the historical always-enter behavior.""" - - def _probe() -> bool: - from gateway.run import _profile_meta_rows - from hermes_cli.loops import _META_PREFIX, _parse_state - - rows = _profile_meta_rows(profile_home, _META_PREFIX) - if rows is None: - return True - # ``list_active_loops()`` collapses an unavailable store to ``[]`` (fail CLOSED), so parse - # the rows here; a corrupt row is "unknown" and keeps the scan. - return any((state := _parse_state(raw)) is None or state.status == "active" for _key, raw in rows) - - offload = getattr(runner, "_run_in_executor_with_context", None) - if not callable(offload): - return True - return bool(await offload(_probe)) - - class GatewayGoalsMixin: """Goal/heartbeat continuation, post-turn hooks and loop-wakeup watcher methods for GatewayRunner.""" @@ -485,6 +461,7 @@ class GatewayGoalsMixin: profile's store is scanned under its own runtime scope (same shape as ``_handoff_watcher``), and each hit is fired against that profile's adapters.""" from gateway.run import _async_profile_runtime_scope, _handoff_watch_scopes + from gateway.run_idle_gates import profile_has_active_loop await asyncio.sleep(5) # let platforms finish connecting warned_no_route: set = set() @@ -507,11 +484,10 @@ class GatewayGoalsMixin: while self._running: try: for profile_name, profile_home in _handoff_watch_scopes(self): - # Idle gate: the scope entry re-parses the profile's config/secrets, so - # only pay it when the profile's store actually holds an active loop. - # The root scan (None) is unscoped and stays cheap. - if profile_home is not None and not await _watcher_has_active_loops( - self, profile_home): + # Idle gate (run_idle_gates): skip the scope entry when the profile's store holds + # no active loop. The root scan (None) is unscoped and stays cheap. + if profile_home is not None and not await self._run_in_executor_with_context( + profile_has_active_loop, profile_home): continue async with _scope(profile_home): await _scan_one_store(profile_name) diff --git a/gateway/run_heartbeat_restore.py b/gateway/run_heartbeat_restore.py index b60f1239b8..b430ba676e 100644 --- a/gateway/run_heartbeat_restore.py +++ b/gateway/run_heartbeat_restore.py @@ -7,27 +7,6 @@ from pathlib import Path logger = logging.getLogger("gateway.run") -def _profile_has_active_heartbeat(profile_home) -> bool: - """One profile's SessionDB holds a ``heartbeat:*`` row still ACTIVE — the only rows the sweep can - restore (``HeartbeatManager.is_active``). Reads the goals-cached DB only: no config/secret parsing. - Fails OPEN: an unavailable store, a failing read or a corrupt row all answer True, so the gate can - never suppress a restore the full sweep would have made. ``clear``/``pause`` keep their rows (status - ``cleared``/``paused``), so key existence alone would re-enable the sweep forever after first use.""" - from gateway.run import _profile_meta_rows - from hermes_cli.heartbeat import HeartbeatState - - rows = _profile_meta_rows(profile_home, "heartbeat:") - if rows is None: - return True - for _key, raw in rows: - try: - if HeartbeatState.from_json(raw).status == "active": - return True - except Exception: - return True - return False - - def _watched_homes(runner, default_home) -> list: """Every home the sweep's ``_profile_scope_for_source`` can resolve an origin to: the gateway home plus, under multiplex, the whole served set INCLUDING ``default`` — a ``-p work`` multiplexer's own @@ -48,6 +27,7 @@ async def restore_heartbeat_watches(runner) -> None: Run all storage work off-loop so a cold profile DB cannot block adapters. """ from gateway.run import _profile_runtime_scope + from gateway.run_idle_gates import profile_has_active_heartbeat from hermes_cli.heartbeat import HeartbeatManager from hermes_constants import get_hermes_home @@ -60,7 +40,7 @@ async def restore_heartbeat_watches(runner) -> None: home = getattr(store, "_routing_home", None) or get_hermes_home() # Cheap gate: with no heartbeat persisted in any served profile there is nothing to # restore — skip the per-origin profile-scope re-parse over every routed session. - if not any(_profile_has_active_heartbeat(h) for h in _watched_homes(runner, home)): + if not any(profile_has_active_heartbeat(h) for h in _watched_homes(runner, home)): return restored with _profile_runtime_scope(home): entries = store.list_sessions() diff --git a/gateway/run_idle_gates.py b/gateway/run_idle_gates.py new file mode 100644 index 0000000000..66ce74cf81 --- /dev/null +++ b/gateway/run_idle_gates.py @@ -0,0 +1,66 @@ +"""Idle gates for the gateway's per-profile pollers (heartbeat restore, handoff watcher, loop wakeup). + +Entering ``_profile_runtime_scope`` costs a config.yaml load, a ``.env`` parse, secret hydration and a +terminal-policy build; on a multiplex gateway the pollers paid that per profile per tick with nothing +to do. Each gate answers "does this profile's store hold work?" from the goals-cached SessionDB with +ONLY the HERMES_HOME contextvar installed. Every gate fails OPEN: an unavailable store, a failing +read or a corrupt row is "cannot prove emptiness", never "idle". +""" +from __future__ import annotations + +import logging +from pathlib import Path +from typing import Any, Callable, Optional + +logger = logging.getLogger("gateway.run") + + +def _profile_session_db_probe(profile_home: Path) -> Optional[Any]: + """The goals-cached SessionDB for *profile_home*; None when unavailable.""" + from hermes_cli.goals import _get_session_db + from hermes_constants import reset_hermes_home_override, set_hermes_home_override + + token = set_hermes_home_override(str(profile_home)) + try: + return _get_session_db() + except Exception: + logger.debug("session-db probe failed for %s", profile_home, exc_info=True) + return None + finally: + reset_hermes_home_override(token) + + +def _gate(profile_home: Path, store_has_work: Callable[[Any], bool]) -> bool: + db = _profile_session_db_probe(profile_home) + if db is None: + return True + try: + return store_has_work(db) + except Exception: + logger.debug("idle probe failed for %s; keeping the full sweep", profile_home, exc_info=True) + return True + + +def profile_has_active_heartbeat(profile_home: Path) -> bool: + from hermes_cli.heartbeat import store_has_active_heartbeat + + return _gate(profile_home, store_has_active_heartbeat) + + +def profile_has_active_loop(profile_home: Path) -> bool: + from hermes_cli.loops import store_has_active_loop + + return _gate(profile_home, store_has_active_loop) + + +def profile_has_pending_handoff(profile_home: Path) -> bool: + return _gate(profile_home, lambda db: db.has_pending_handoffs()) + + +async def off_loop_gate(runner: object, probe: Callable[[], bool]) -> bool: + """Run a sync gate through the runner's executor hop. Runners without one (bare test stand-ins + for the handoff watcher) keep the historical always-enter behaviour.""" + offload = getattr(runner, "_run_in_executor_with_context", None) + if not callable(offload): + return True + return bool(await offload(probe)) diff --git a/hermes_cli/heartbeat.py b/hermes_cli/heartbeat.py index d3f440d9ed..244edbd1ac 100644 --- a/hermes_cli/heartbeat.py +++ b/hermes_cli/heartbeat.py @@ -99,12 +99,15 @@ def _get_session_db() -> Optional[Any]: return None +_META_PREFIX = "heartbeat:" + + def load_heartbeat(session_id: str) -> Optional[HeartbeatState]: db = _get_session_db() if session_id else None if db is None: return None try: - raw = db.get_meta(f"heartbeat:{session_id}") + raw = db.get_meta(_META_PREFIX + session_id) except Exception as exc: logger.debug("HeartbeatManager: get_meta failed: %s", exc) return None @@ -116,6 +119,19 @@ def load_heartbeat(session_id: str) -> Optional[HeartbeatState]: return None if state is None or state.status == "cleared" else state +def store_has_active_heartbeat(db: Any) -> bool: + """True when *db* holds an ACTIVE ``heartbeat:*`` row — or one that cannot be parsed (unknown, so + the caller keeps its full sweep). ``clear``/``pause`` keep their rows (status ``cleared``/``paused``), + so key existence alone is not "active". Read errors propagate: "unavailable" is the caller's call.""" + for _key, raw in db.list_meta_prefix(_META_PREFIX): + try: + if HeartbeatState.from_json(raw).status == "active": + return True + except Exception: + return True + return False + + def save_heartbeat(session_id: str, state: HeartbeatState) -> None: if not session_id: return @@ -125,7 +141,7 @@ def save_heartbeat(session_id: str, state: HeartbeatState) -> None: _warn_dropped_write("HeartbeatManager", "heartbeat", session_id) return try: - db.set_meta(f"heartbeat:{session_id}", state.to_json()) + db.set_meta(_META_PREFIX + session_id, state.to_json()) except Exception as exc: logger.debug("HeartbeatManager: set_meta failed: %s", exc) diff --git a/hermes_cli/loops.py b/hermes_cli/loops.py index 7242004d5a..3c377470f8 100644 --- a/hermes_cli/loops.py +++ b/hermes_cli/loops.py @@ -321,6 +321,17 @@ def list_active_loops() -> List[Tuple[str, LoopState]]: return out +def store_has_active_loop(db: Any) -> bool: + """True when *db* holds an ACTIVE ``loop:*`` row — or a row that cannot be parsed (unknown, so the + caller keeps its full scan). Unlike :func:`list_active_loops` this takes the store explicitly and + propagates read errors, so an idle gate can tell "empty" from "unavailable".""" + for key, raw in db.list_meta_prefix(_META_PREFIX): + state = _parse_state(raw, key[len(_META_PREFIX):]) if raw else None + if state is None or state.status == "active": + return True + return False + + def migrate_loop_to_session(old_session_id: str, new_session_id: str, *, reason: str = "") -> bool: """Carry a /loop from a parent session to its continuation. Best-effort, never raises. diff --git a/hermes_state_gateway.py b/hermes_state_gateway.py index 788003d54c..ac7bd5a1a1 100644 --- a/hermes_state_gateway.py +++ b/hermes_state_gateway.py @@ -837,13 +837,8 @@ class SessionGatewayMixin: ) > 0 def has_pending_handoffs(self) -> bool: - """Cheap existence probe for the handoff watcher's idle gate. Fails OPEN: a probe - error answers True so a broken store never silently disables handoff dispatch.""" - try: - return self._read_one( - "SELECT 1 AS found FROM sessions WHERE handoff_state = 'pending' LIMIT 1") is not None - except Exception: - return True + """Bounded existence probe for the handoff watcher's idle gate (the gate fails open on error).""" + return self._read_one("SELECT 1 FROM sessions WHERE handoff_state = 'pending' LIMIT 1") is not None def complete_handoff(self, session_id: str) -> None: """Mark a handoff as completed.""" diff --git a/tests/gateway/test_heartbeat_watch_restore.py b/tests/gateway/test_heartbeat_watch_restore.py index 1cd125a362..eb350a9c56 100644 --- a/tests/gateway/test_heartbeat_watch_restore.py +++ b/tests/gateway/test_heartbeat_watch_restore.py @@ -128,7 +128,7 @@ async def test_restore_skips_session_sweep_when_no_heartbeats_exist(tmp_path, mo # Fail OPEN: a store the probe cannot open must not suppress the sweep. dbs[str(named)].set_meta('heartbeat:live', HeartbeatState( prompt='p', interval_seconds=60, status='cleared').to_json()) - monkeypatch.setattr('gateway.run._profile_session_db_probe', lambda _home: None) + monkeypatch.setattr('gateway.run_idle_gates._profile_session_db_probe', lambda _home: None) await restore_heartbeat_watches(runner) assert sweeps == [1, 1] finally: diff --git a/tests/gateway/test_loop_command.py b/tests/gateway/test_loop_command.py index 4d1031aea2..56e98d8380 100644 --- a/tests/gateway/test_loop_command.py +++ b/tests/gateway/test_loop_command.py @@ -474,7 +474,7 @@ async def test_loop_wakeup_watcher_gates_profile_scope_on_active_loops(loop_env, cleared = loops.LoopState.from_json(raw) cleared.status = "cleared" work_db.set_meta("loop:sid-work-loop", cleared.to_json()) - monkeypatch.setattr("gateway.run._profile_session_db_probe", lambda _home: None) + monkeypatch.setattr("gateway.run_idle_gates._profile_session_db_probe", lambda _home: None) await _run_one_tick() assert entered == [work_home], ( f"unavailable store must fall back to the historical scan; got {entered}")