refactor(gateway): idle gates live in run_idle_gates; row semantics live with loops/heartbeat
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.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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()
|
||||
|
||||
66
gateway/run_idle_gates.py
Normal file
66
gateway/run_idle_gates.py
Normal file
@@ -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))
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
|
||||
@@ -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."""
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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}")
|
||||
|
||||
Reference in New Issue
Block a user