fix(tui): heartbeat ownership follows the gateway's live routing index, not the row source

_notif_gateway_owns_heartbeat decided by the immutable sessions.source column,
so a heartbeat on an ARCHIVED gateway-sourced row (Telegram /reset, idle/daily
auto-reset, compression rotation) was skipped by the Desktop poller and never
registered by the gateway either — restore_heartbeat_watches only claims a key
whose current session_id is that row. The tick belonged to nobody and stayed
due forever, where origin/main's Desktop fired it.

Ownership now uses the same predicate the gateway does: a gateway_routing entry
whose current session_id is this session, with an origin and not suspended
(SessionDB.gateway_routing_entry_for_session; both the session's profile store
and the launch store are consulted so multiplexed and per-profile gateways are
covered). No entry is fail-open, as on main. The check runs after the cheap
is_active/is_due gate so idle sessions never touch the DB per poll.

Probe (SessionStore telegram -> force_new, heartbeat on the archived sid):
before: desktop fired False / gateway watches [] / still due True;
after: desktop fired True / still due False; the current gateway sid is still
left to the gateway (desktop fired False) — the hijack fix stays intact.
This commit is contained in:
teknium1
2026-09-15 14:47:13 -07:00
committed by Teknium
parent 037771a692
commit 490ee1607a
3 changed files with 51 additions and 25 deletions
+12
View File
@@ -313,6 +313,18 @@ class SessionGatewayMixin:
)
return [dict(r) for r in rows]
def gateway_routing_entry_for_session(self, session_id: str) -> Optional[Dict[str, Any]]:
"""The routing entry (any scope) whose current owner is *session_id*, or None. The id lives
only inside ``entry_json``, so matching is done in Python; an archived/rotated row has none."""
for row in self._read_all("SELECT entry_json FROM gateway_routing"):
try:
entry = json.loads(row["entry_json"] or "{}")
except Exception:
continue
if isinstance(entry, dict) and entry.get("session_id") == session_id:
return entry
return None
def _delete_routing_entries_for_sessions(self, session_ids: Set[str]) -> int:
"""Drop ``gateway_routing`` rows pointing at any of *session_ids*; the target id
lives only inside ``entry_json``, so matching is done in Python over all scopes."""
+21 -12
View File
@@ -92,22 +92,31 @@ def test_notification_poller_fires_due_heartbeat_when_idle(server, session):
assert load_heartbeat(key).fire_count == 1 and not load_heartbeat(key).is_due()
def test_desktop_poller_leaves_gateway_owned_heartbeat_for_gateway(server, session):
"""A Desktop viewer must not consume a messaging session's routed heartbeat."""
sid, key, s = session
s["source"] = "desktop"
server._get_db().create_session(key, source="telegram")
_arm_due(key)
def test_desktop_poller_leaves_gateway_owned_heartbeat_for_gateway(server, session, hermes_home):
"""A Desktop viewer must not consume a messaging session's routed heartbeat, but ownership follows the
gateway's live routing index, not the row's immutable ``source``: once /reset archives the row the gateway
never registers a watch for it again, so the Desktop viewer must fire it (else nobody does)."""
from gateway.config import GatewayConfig, Platform
from gateway.session import SessionSource, SessionStore
from hermes_cli.heartbeat import load_heartbeat
sid, _, s = session
s["source"] = "desktop"
store = SessionStore(hermes_home / "sessions", GatewayConfig())
store._db = server._get_db() # the routing index lives in the store the Desktop poller reads
src = SessionSource(platform=Platform.TELEGRAM, chat_id="42")
archived_key = store.get_or_create_session(src).session_id
current_key = store.get_or_create_session(src, force_new=True).session_id
for key, desktop_fires in ((current_key, False), (archived_key, True)):
s["session_key"], s["running"] = key, False
_arm_due(key)
p_submit, p_emit = _submits(server, MagicMock())
with p_submit as submit, p_emit:
server._maybe_fire_tui_heartbeat_tick(sid, s)
from hermes_cli.heartbeat import load_heartbeat
submit.assert_not_called()
assert s["running"] is False
assert load_heartbeat(key).fire_count == 0 and load_heartbeat(key).is_due()
assert submit.called is desktop_fires, key
assert (load_heartbeat(key).fire_count == 1) is desktop_fires
assert load_heartbeat(key).is_due() is not desktop_fires
@pytest.mark.parametrize("running,due", [(True, True), (False, False)])
+16 -11
View File
@@ -196,17 +196,22 @@ def _notif_slash_loop_tick(rid: str, sid: str, session: dict, mgr, wakeup: str)
def _notif_gateway_owns_heartbeat(session: dict, session_key: str) -> bool:
"""Whether the persisted conversation route owns heartbeat delivery.
"""Whether the gateway's heartbeat poller owns this session's due tick.
Desktop/TUI can attach to a messaging conversation, but its session-owner
poller has no adapter route for the reply. The gateway's heartbeat poller
retains that route, so leave its due tick untouched. A missing row is
fail-open for a newly-created local session.
Desktop/TUI can attach to a messaging conversation, but its session-owner poller has no adapter
route for the reply. Ownership is the gateway's LIVE routing index, not the row's immutable
``source``: ``gateway/run_heartbeat_restore.py`` only registers watches for a non-suspended,
origin-bearing key whose current ``session_id`` is this one, so a row archived by /reset,
auto-reset or compression rotation belongs to nobody there and must keep firing here. The index
lives in the gateway's home store (the launch handle for a multiplexed gateway, the profile's
own store for a per-profile gateway), so both are consulted; no entry is fail-open.
"""
try:
with _session_db(session) as db:
row = db.get_session(session_key) if db is not None else None
return bool(row and _is_gateway_owned_source(str(row.get("source") or "")))
for store in {id(d): d for d in (db, _get_db()) if d is not None}.values():
if (entry := store.gateway_routing_entry_for_session(session_key)) is not None:
return bool(entry.get("origin")) and not entry.get("suspended")
return False
except Exception:
return False
@@ -227,11 +232,11 @@ def _maybe_fire_tui_heartbeat_tick(sid: str, session: dict) -> None:
return
if not (sid_key := session.get("session_key") or ""):
return
if _notif_gateway_owns_heartbeat(session, sid_key):
return # the gateway poller owns the persisted messaging route
mgr = HeartbeatManager(session_id=sid_key)
if not mgr.is_active() or not mgr.state.is_due() or not _notif_claim_turn(session):
return # not due, or busy — the tick coalesces to the next idle poll
if not mgr.is_active() or not mgr.state.is_due() or _notif_gateway_owns_heartbeat(session, sid_key):
return # not due, or the gateway poller owns the routed conversation — stays due there
if not _notif_claim_turn(session):
return # busy — the tick coalesces to the next idle poll
if not (prompt := mgr.due_prompt()):
_notif_release_turn(session)
return