diff --git a/hermes_state_gateway.py b/hermes_state_gateway.py index 5133c5ed47..29905eddd6 100644 --- a/hermes_state_gateway.py +++ b/hermes_state_gateway.py @@ -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.""" diff --git a/tests/tui_gateway/test_heartbeat_tui_tick.py b/tests/tui_gateway/test_heartbeat_tui_tick.py index 2edf0e938d..7ee912d396 100644 --- a/tests/tui_gateway/test_heartbeat_tui_tick.py +++ b/tests/tui_gateway/test_heartbeat_tui_tick.py @@ -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) - - p_submit, p_emit = _submits(server, MagicMock()) - with p_submit as submit, p_emit: - server._maybe_fire_tui_heartbeat_tick(sid, s) - +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 - submit.assert_not_called() - assert s["running"] is False - assert load_heartbeat(key).fire_count == 0 and load_heartbeat(key).is_due() + 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) + 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)]) diff --git a/tui_gateway/session_notifications.py b/tui_gateway/session_notifications.py index 50c7c4024d..8b93074a96 100644 --- a/tui_gateway/session_notifications.py +++ b/tui_gateway/session_notifications.py @@ -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