From 9cd1453f2666734738cb446b714d2f603800a856 Mon Sep 17 00:00:00 2001 From: Totoro-qaq Date: Thu, 10 Sep 2026 01:15:20 -0700 Subject: [PATCH] fix(sessions): a stale explicit-close stamp no longer makes a session permanently uncompressible Consolidated from PR #106543 (5 commits, final tree d2c4d908) by @Totoro-qaq. publish_compression_child() fails closed on any non-automatic end stamp and end_session() is first-stamp-wins, so a stale tui_close on a session the TUI still routes turned every rotation into "compute the summary, then discard it" (#106459). The host that still routes the session clears the stale explicit close via SessionDB.reopen_if_explicitly_closed() before the turn starts; publication never heals explicit closes. Review probes by @ehz0ah. --- agent/conversation_compression.py | 14 +- hermes_state_common.py | 4 + hermes_state_compression.py | 116 +++- .../test_106459_stale_explicit_close_stamp.py | 552 ++++++++++++++++++ .../test_106459_routing_provenance_reopen.py | 324 ++++++++++ .../test_bot_live_owner_delivery.py | 3 + tools/session_search_tool.py | 6 +- tui_gateway/prompt_turn.py | 39 +- 8 files changed, 1044 insertions(+), 14 deletions(-) create mode 100644 tests/hermes_state/test_106459_stale_explicit_close_stamp.py create mode 100644 tests/tui_gateway/test_106459_routing_provenance_reopen.py diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 1f49aa9060..851d2733c8 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -2899,8 +2899,18 @@ def _salvage_or_refuse_grown_transcript( def _parent_deliberately_ended(session_db: Any, session_id: str) -> bool: - """True when the parent row was ended by a non-automatic reason. Fails OPEN: an - unreadable row must not turn a cheap guard into a new way to lose compression.""" + """True when publish_compression_child() would fail closed on the parent's end stamp, so the + durable pre-publish flush is skipped. The store owns the verdict (#106459): a stale explicit close is + healed by the host that still routes the session before the turn starts, so any explicit close seen + here is preserved. Fails OPEN: an unreadable row must not turn a cheap guard into a new way to lose + compression.""" + verdict = getattr(session_db, "compression_parent_deliberately_ended", None) + if callable(verdict): + try: + return bool(verdict(session_id)) + except Exception: + return False + # Stores without the verdict (test stand-ins): taxonomy-only fallback. reader = getattr(session_db, "get_session", None) if not callable(reader): return False diff --git a/hermes_state_common.py b/hermes_state_common.py index dcb2f77d49..53ac2d4b63 100644 --- a/hermes_state_common.py +++ b/hermes_state_common.py @@ -120,6 +120,10 @@ _COMPRESSION_CHILD_SQL = ("EXISTS (SELECT 1 FROM sessions p WHERE p.id = # ended that way. Must stay identical to the recovery fence in find_latest_gateway_session_for_peer. _RESET_END_REASONS = ("session_reset", "session_switch", "idle", "daily", "suspended", "resume_pending_expired") _RESET_END_REASONS_SQL = ", ".join(f"'{reason}'" for reason in _RESET_END_REASONS) +# Deliberate conversation boundaries: the reset set plus CLI /new, which ends the predecessor as +# 'new_session' (hermes_cli/cli_session_mixin.py) without a reset child row. A compression rotation must +# never heal one of these (#106459); tools/session_search_tool.py derives its fresh-reset set from it. +_BOUNDARY_END_REASONS = frozenset(_RESET_END_REASONS) | {"new_session"} # Accidental end reasons recovery treats as resumable (docs/session-lifecycle.md); single source of truth for # recovery SQL and SessionDB.RECOVERABLE_END_REASONS. superseded_by_resume = sentinel-parked runtime replaced diff --git a/hermes_state_compression.py b/hermes_state_compression.py index 167470ae29..0ddbe62013 100644 --- a/hermes_state_compression.py +++ b/hermes_state_compression.py @@ -12,8 +12,8 @@ import time from typing import Any, Dict, List, Optional, Tuple from hermes_state_common import ( - _COMPRESSION_LOCK_ROW_SQL as _LOCK_ROW_SQL, _ENDED_ROW_SQL, _ended_by_compression, _sql_session_last_active, - is_automatic_end_reason) + _BOUNDARY_END_REASONS, _COMPRESSION_LOCK_ROW_SQL as _LOCK_ROW_SQL, _ENDED_ROW_SQL, _RESET_END_REASONS, + _ended_by_compression, _sql_session_last_active, is_automatic_end_reason) # Log-record parity with the origin module (caplog tests pin "hermes_state"). logger = logging.getLogger("hermes_state") @@ -70,6 +70,105 @@ def _claim_lease_row(conn, table: str, key_col: str, key: str, holder: str, now: class SessionCompressionMixin: """Compression lineage, cooldown/streak counters, locks and turn leases.""" + def _end_stamp_class(self, conn, session_id: str, row) -> Optional[str]: + """Classify *row*'s end stamp (a mapping with ``ended_at`` / ``end_reason``): None when live, + ``'automatic'`` (cleanup stamp, stale by construction, #88197), ``'compression'`` (names a + continuation), ``'boundary'`` (reset reasons and CLI ``new_session``: a deliberate end of the + conversation), ``'superseded'`` (an explicit close whose continuation child was already published), + or ``'explicit'`` (``tui_close``, ``cli_close``, ``webhook_complete``, ...: a close with no + continuation). Single owner of the taxonomy for ``reopen_if_explicitly_closed()``, + publish_compression_child() and the agent's pre-flush guard (#106459, never-patch-predicates).""" + if row is None or row["ended_at"] is None: + return None + reason = row["end_reason"] + if is_automatic_end_reason(reason): + return "automatic" + if reason == "compression": + return "compression" + if reason in _BOUNDARY_END_REASONS: + return "boundary" + child = conn.execute( + "SELECT 1 FROM sessions WHERE parent_session_id = ?" + + self._NON_CONTINUATION_CHILD_FILTER_SQL.format(alias="") + + " LIMIT 1", + (session_id, session_id, session_id), + ).fetchone() + return "superseded" if child is not None else "explicit" + + def _compression_parent_obstacle(self, conn, parent_session_id: str, parent) -> Optional[str]: + """Why publishing a compression child of *parent* must fail closed, or None when the parent is live + or carries an automatic-cleanup stamp that publish may clear (#88197). + + Every other stamp fails closed here, explicit closes included. A stale explicit close is healed only + by the host that still routes the session (``reopen_if_explicitly_closed()``, called by the TUI as it + starts a turn), never by the store: healing at publication cannot be made safe because + ``end_session()`` is first-stamp-wins -- while a stale stamp occupies the row, a close made during + the turn is a no-op write -- and healing at turn-lease admission cannot either, because a turn + worker can take its lease after the host has already closed the session. An explicit stamp present + here is therefore a deliberate close, a foreign writer's mis-stamp that costs this one rotation, or + a stamp no host has vouched against, and publish must not resurrect it.""" + kind = self._end_stamp_class(conn, parent_session_id, parent) + if kind in (None, "automatic"): + return None + reason = parent["end_reason"] + if kind == "compression": + return "closed by compression" + if kind == "boundary": + return f"closed by a {reason} boundary" + if kind == "superseded": + return f"closed ({reason}) with a published continuation" + return f"closed ({reason}); an explicit close is healed only by the host that still routes the session" + + def reopen_if_explicitly_closed( + self, session_id: str, *, provenance: str, patience_s: Optional[float] = None, + ) -> Optional[str]: + """Clear an explicit-close stamp (``tui_close``, ``cli_close``, ``webhook_complete``, ...) from a + session a HOST has just proven is still routed to it, returning the reason cleared or None. + *provenance* names that proof and is logged; *patience_s* bounds the write for a caller holding a + hot lock. Narrow twin of ``reopen_session()``, which clears any stamp: automatic (left to publish, + #88197), ``'compression'``, boundary and superseded stamps are lineage owned elsewhere and are + never touched here. The UPDATE is conditional on the exact stamp that was read, so a close landing + between the read and the write survives. + + The store cannot make this call itself (#106459). Acquiring the session turn lease is not proof of + routing: the TUI starts a turn worker before that worker reaches ``run_conversation()``, so a + ``session.close`` can pop the session, wait out its grace and stamp ``tui_close`` in between, and a + late lease would then clear a deliberate close that nothing re-applies. Only a host that holds the + registry claim can say "this conversation is still mine and I am accepting a turn for it" -- the + gateway's stale-route self-heal (#54878) is the same rule on the routing table. Call it under + whatever lock makes the host's claim atomic with its teardown.""" + if not session_id: + return None + def _do(conn): + row = conn.execute(_ENDED_ROW_SQL, (session_id,)).fetchone() + if self._end_stamp_class(conn, session_id, row) != "explicit": + return None + conn.execute( + "UPDATE sessions SET ended_at = NULL, end_reason = NULL " + "WHERE id = ? AND ended_at = ? AND end_reason = ?", + (session_id, row["ended_at"], row["end_reason"])) + return str(row["end_reason"]) + reason = self._execute_write(_do, patience_s=patience_s) + if reason is not None: + logger.warning( + "Session %s carried a stale %r end stamp while %s; cleared so the conversation can " + "compress and a later close is recorded (#106459)", session_id, reason, provenance) + return reason + + def compression_parent_deliberately_ended(self, session_id: str) -> bool: + """Read-only twin of publish_compression_child()'s liveness verdict, for the agent's guard that + runs BEFORE the durable pre-publish flush: True only when publish would fail closed on + *session_id*'s end stamp. A missing row is not an obstacle (publish reports that itself).""" + if not session_id: + return False + # The public reader on purpose: an unreadable row raises to the guard, which fails open + # (tests/agent/test_compression_rotation_state.py pins that contract). + parent = self.get_session(session_id) + if not parent or parent.get("ended_at") is None: + return False + with self._read_ctx() as conn: + return self._compression_parent_obstacle(conn, session_id, parent) is not None + def find_live_compression_child(self, parent_session_id: str) -> Optional[Dict[str, Any]]: """The unique live direct child of a compression-ended session, else None. A stale agent whose parent was rotated elsewhere may recover only when the lineage names @@ -222,12 +321,13 @@ class SessionCompressionMixin: if parent is None: raise RuntimeError(f"Compression parent not found: {parent_session_id}") if parent["ended_at"] is not None: - # An AUTOMATIC end stamp (tui_shutdown, ws_disconnect, orphan reap, idle/LRU - # evict) is stale by construction — this lease holder is still continuing the - # conversation, and left alone it wedges rotation forever. Clear it; the closure - # UPDATE below re-stamps end_reason='compression'. Deliberate boundaries fail closed. - if not is_automatic_end_reason(parent["end_reason"]): - raise RuntimeError(f"Compression parent already ended: {parent_session_id}") + # An automatic-cleanup stamp (#88197) is cleared here: this lease holder is still + # continuing the conversation, and left alone the stamp wedges rotation forever. The + # closure UPDATE below re-stamps end_reason='compression'. Every other stamp fails closed + # (see the verdict); a stale explicit close is healed by the routing host, not here (#106459). + obstacle = self._compression_parent_obstacle(conn, parent_session_id, parent) + if obstacle is not None: + raise RuntimeError(f"Compression parent already ended: {parent_session_id} ({obstacle})") conn.execute( "UPDATE sessions SET ended_at = NULL, end_reason = NULL WHERE id = ?", (parent_session_id,)) diff --git a/tests/hermes_state/test_106459_stale_explicit_close_stamp.py b/tests/hermes_state/test_106459_stale_explicit_close_stamp.py new file mode 100644 index 0000000000..fb129fbf12 --- /dev/null +++ b/tests/hermes_state/test_106459_stale_explicit_close_stamp.py @@ -0,0 +1,552 @@ +"""Regression tests for #106459 — an over-limit session must not become permanently +uncompressible because its row carries a stale explicit-close stamp, and a close the user +makes must never be resurrected. + +``publish_compression_child`` used to fail closed on ANY non-automatic ``end_reason`` +(``tui_close``, ``cli_close``, ``webhook_complete``, gateway-recovery and cron reasons), and +``end_session()`` is first-stamp-wins, so nothing ever cleared it. Every turn then ran the +compression, discarded its result at publication, kept the oversized history and hit +``Context length exceeded … Cannot compress further`` again; manual ``/compress`` reported +"No changes". The agent's pre-flush guard re-implemented the same taxonomy and aborted even +earlier. + +Contract under test. The store never decides on its own that an explicit close is stale. +Publication (single verdict owner ``_compression_parent_obstacle``, shared with the agent's +pre-flush guard) clears only automatic-cleanup stamps (#88197) and fails closed on every other +stamp. Turn-lease admission clears nothing: a TUI turn worker can take its lease after the host +has already closed the session. A stale explicit close is cleared only through +``reopen_if_explicitly_closed()``, called by the host that still routes the session -- the TUI, +under ``_sessions_lock`` as it starts a turn for a session it still has registered (covered in +``tests/tui_gateway/test_106459_routing_provenance_reopen.py``). That method clears only explicit +closes with no published continuation and leaves automatic, ``'compression'``, boundary (reset +reasons and CLI ``new_session``) and superseded stamps alone. Here the host's call is stood in +for by ``_host_reopen``. +""" + +from __future__ import annotations + +import logging +import os +import time +from pathlib import Path +from unittest.mock import MagicMock, patch + +import pytest + +from hermes_state import SessionDB +from hermes_state_common import _BOUNDARY_END_REASONS, _RESET_END_REASONS + +STALE_EXPLICIT_CLOSES = ["tui_close", "cli_close", "webhook_complete"] +BOUNDARIES = sorted(_BOUNDARY_END_REASONS) +HOLDER = "compression-writer" +# Holders carry this process's pid: a dead pid makes a lease reclaimable, which is not what these probe. +TURN = f"pid={os.getpid()}:turn=t1:platform=tui" +TURN_2 = f"pid={os.getpid()}:turn=t2:platform=tui" +NOT_HEALED = "healed only by the host that still routes the session" +PROVENANCE = "the test host still routes the session" + + +@pytest.fixture +def db(tmp_path: Path): + handle = SessionDB(db_path=tmp_path / "state.db") + try: + yield handle + finally: + handle.close() + + +def _publish(db: SessionDB, parent: str, child: str, *, holder: str | None = None) -> None: + db.publish_compression_child( + parent_session_id=parent, + child_session_id=child, + source="tui", + messages=[{"role": "user", "content": "[CONTEXT COMPACTION] summary"}], + require_compression_lease=holder is not None, + compression_lock_holder=holder, + ) + + +def _stamp(db: SessionDB, session_id: str, reason: str, *, age: float = 0.0) -> None: + """End the row with *reason*; ``age`` backdates the stamp so it is unambiguously older than + anything that follows.""" + db.end_session(session_id, reason) + if age: + db._write_sql("UPDATE sessions SET ended_at = ? WHERE id = ?", (time.time() - age, session_id)) + row = db.get_session(session_id) + assert row["ended_at"] is not None and row["end_reason"] == reason + + +def _lease(db: SessionDB, session_id: str, holder: str = HOLDER) -> str: + """The compression lease: publication ownership only.""" + assert db.try_acquire_compression_lock(session_id, holder, ttl_seconds=300.0) + return holder + + +def _admit(db: SessionDB, session_id: str, holder: str = TURN) -> str: + """Admit a turn on the conversation, as ``admit_durable_turn_lease`` does inside the turn worker.""" + assert db.try_acquire_session_turn_lease(session_id, holder, ttl_seconds=300.0) + return holder + + +def _host_reopen(db: SessionDB, session_id: str) -> str | None: + """What the TUI does under ``_sessions_lock`` for a session it still has registered.""" + return db.reopen_if_explicitly_closed(session_id, provenance=PROVENANCE) + + +def _live(db: SessionDB, session_id: str) -> bool: + row = db.get_session(session_id) + return row["ended_at"] is None and row["end_reason"] is None + + +class TestReopenIfExplicitlyClosed: + @pytest.mark.parametrize("reason", STALE_EXPLICIT_CLOSES) + def test_clears_each_explicit_close_and_logs_the_provenance(self, db: SessionDB, reason: str, caplog) -> None: + sid = f"S_{reason}" + db.create_session(sid, source="tui") + _stamp(db, sid, reason, age=60.0) + + with caplog.at_level(logging.WARNING, logger="hermes_state"): + assert _host_reopen(db, sid) == reason + + assert _live(db, sid) + assert any(reason in r.getMessage() and PROVENANCE in r.getMessage() and "stale" in r.getMessage() + for r in caplog.records), "clearing must be logged with the reason and the host's provenance" + + @pytest.mark.parametrize("reason", ["ws_disconnect", "compression", *BOUNDARIES]) + def test_leaves_every_other_stamp_alone(self, db: SessionDB, reason: str) -> None: + sid = f"S_{reason}" + db.create_session(sid, source="tui") + _stamp(db, sid, reason, age=60.0) + assert _host_reopen(db, sid) is None + row = db.get_session(sid) + assert row["end_reason"] == reason and row["ended_at"] is not None + + def test_leaves_a_superseded_close_alone(self, db: SessionDB) -> None: + parent = "P_superseded_reopen" + db.create_session(parent, source="tui") + _publish(db, parent, "C_first", holder=_lease(db, parent)) # a continuation now exists + db.release_compression_lock(parent, HOLDER) + db.reopen_session(parent) + _stamp(db, parent, "tui_close", age=60.0) + assert _host_reopen(db, parent) is None + assert db.get_session(parent)["end_reason"] == "tui_close" + + def test_live_missing_and_empty_ids_are_no_ops(self, db: SessionDB) -> None: + db.create_session("live", source="tui") + assert _host_reopen(db, "live") is None + assert _live(db, "live") + assert _host_reopen(db, "missing") is None + assert _host_reopen(db, "") is None + + +class TestTurnAdmissionHealsNothing: + @pytest.mark.parametrize("reason", STALE_EXPLICIT_CLOSES) + def test_admitting_a_turn_leaves_an_explicit_close(self, db: SessionDB, reason: str) -> None: + """A lease is not routing provenance: the store must not treat its acquisition as proof.""" + sid = f"S_admit_{reason}" + db.create_session(sid, source="tui") + _stamp(db, sid, reason, age=60.0) + _admit(db, sid) + assert db.get_session(sid)["end_reason"] == reason + + def test_a_late_turn_lease_after_the_host_closed_leaves_the_close(self, db: SessionDB) -> None: + """Fourth review probe, store level: the TUI accepted a prompt and started a worker, the user closed + the session (pop, grace, ``tui_close``), and only then did the worker reach ``run_conversation()`` + and take its turn lease. Nothing re-applies the close afterwards, so the lease must not clear it.""" + sid = "S_late_worker" + db.create_session(sid, source="tui") + _stamp(db, sid, "tui_close") # teardown finished before the worker took its lease + _admit(db, sid) + assert db.get_session(sid)["end_reason"] == "tui_close" + with pytest.raises(RuntimeError, match=f"already ended.*{NOT_HEALED}"): + _publish(db, sid, "C_never", holder=_lease(db, sid)) + assert db.get_session("C_never") is None + assert db.get_session(sid)["end_reason"] == "tui_close" + + +class TestHostReopenUnwedgesRotation: + @pytest.mark.parametrize("reason", STALE_EXPLICIT_CLOSES) + def test_a_stale_close_cleared_by_the_host_lets_the_rotation_publish(self, db: SessionDB, reason: str) -> None: + parent = f"P_{reason}" + db.create_session(parent, source="tui") + db.append_message(parent, "user", content="hello") + _stamp(db, parent, reason, age=60.0) # the #106459 shape: the mark predates this turn + assert _host_reopen(db, parent) == reason + _admit(db, parent) + _publish(db, parent, f"C_{reason}", holder=_lease(db, parent)) + + parent_row = db.get_session(parent) + assert parent_row["end_reason"] == "compression" and parent_row["ended_at"] is not None + child_row = db.get_session(f"C_{reason}") + assert child_row is not None and child_row["parent_session_id"] == parent + + def test_repeated_rotation_is_not_wedged(self, db: SessionDB) -> None: + """The field shape: the stamp lands once, and every later turn must still rotate.""" + parent = "P_repeat" + db.create_session(parent, source="tui") + _stamp(db, parent, "cli_close", age=60.0) + _host_reopen(db, parent) + _admit(db, parent) + _publish(db, parent, "C_1", holder=_lease(db, parent)) + db.release_session_turn_lease(parent, TURN) + _host_reopen(db, "C_1") # nothing to clear on the live continuation + _admit(db, "C_1", holder=TURN_2) + _publish(db, "C_1", "C_2", holder=_lease(db, "C_1")) + assert db.get_session("C_1")["end_reason"] == "compression" + assert db.get_session("C_2")["parent_session_id"] == "C_1" + + def test_a_stale_close_on_a_compression_child_is_cleared_on_that_row(self, db: SessionDB) -> None: + """The host reopens the row the turn writes to -- the compression tip -- not the lineage root.""" + root = "P_root" + db.create_session(root, source="tui") + _admit(db, root) + _publish(db, root, "C_1", holder=_lease(db, root)) + db.release_session_turn_lease(root, TURN) + _stamp(db, "C_1", "tui_close", age=60.0) + + assert _host_reopen(db, "C_1") == "tui_close" + assert db.get_session(root)["end_reason"] == "compression" # the root's stamp is lineage, untouched + _admit(db, "C_1", holder=TURN_2) + _publish(db, "C_1", "C_2", holder=_lease(db, "C_1")) + assert db.get_session("C_2")["parent_session_id"] == "C_1" + + def test_a_stamp_that_lands_mid_turn_costs_that_turn_only(self, db: SessionDB) -> None: + """A foreign writer's mis-stamp during the turn cannot be told from a deliberate close, so that + rotation fails closed; the host's reopen on the next prompt clears it and the rotation publishes.""" + parent = "P_mid_turn" + db.create_session(parent, source="tui") + _host_reopen(db, parent) + _admit(db, parent) + _stamp(db, parent, "webhook_complete") + with pytest.raises(RuntimeError, match=f"already ended.*{NOT_HEALED}"): + _publish(db, parent, "C_never", holder=_lease(db, parent)) + assert db.get_session("C_never") is None + db.release_compression_lock(parent, HOLDER) + db.release_session_turn_lease(parent, TURN) + + assert _host_reopen(db, parent) == "webhook_complete" + _admit(db, parent, holder=TURN_2) + _publish(db, parent, "C_next", holder=_lease(db, parent)) + assert db.get_session(parent)["end_reason"] == "compression" + assert db.get_session("C_next")["parent_session_id"] == parent + + +class TestADeliberateCloseIsNeverHealedAtPublication: + def test_close_during_compression_fails_closed(self, db: SessionDB) -> None: + """First review probe: lease, then the user closes, then publish.""" + parent = "P_closed_mid_compression" + db.create_session(parent, source="tui") + _admit(db, parent) + holder = _lease(db, parent) + _stamp(db, parent, "tui_close") + + with pytest.raises(RuntimeError, match=f"already ended.*{NOT_HEALED}"): + _publish(db, parent, "C_never", holder=holder) + assert db.get_session("C_never") is None + assert db.get_session(parent)["end_reason"] == "tui_close" # the user's close survives + + def test_close_after_the_turn_was_admitted_fails_closed(self, db: SessionDB) -> None: + """Second review probe: the turn is running, ``session.close`` stamps after its grace without + interrupting the turn thread, and that thread then takes the compression lease.""" + parent = "P_closed_during_turn" + db.create_session(parent, source="tui") + _admit(db, parent) + _stamp(db, parent, "tui_close") + holder = _lease(db, parent) + + with pytest.raises(RuntimeError, match=f"already ended.*{NOT_HEALED}"): + _publish(db, parent, "C_never", holder=holder) + assert db.get_session("C_never") is None + assert db.get_session(parent)["end_reason"] == "tui_close" + + def test_close_during_a_turn_that_began_on_a_stale_stamp_is_not_resurrected(self, db: SessionDB) -> None: + """Third review probe, store only: the turn began on the stale stamp this fix targets, the user's + close during it is a no-op write under first-stamp-wins, and publication must still refuse.""" + parent = "P_stale_then_closed" + db.create_session(parent, source="tui") + _stamp(db, parent, "tui_close", age=60.0) + _admit(db, parent) + db.end_session(parent, "tui_close") # no-op: the stale stamp already occupies the row + with pytest.raises(RuntimeError, match=f"already ended.*{NOT_HEALED}"): + _publish(db, parent, "C_never", holder=_lease(db, parent)) + assert db.get_session("C_never") is None + assert db.get_session(parent)["end_reason"] == "tui_close" + + def test_after_the_host_reopen_a_close_during_the_turn_is_recorded_and_wins(self, db: SessionDB) -> None: + """Third review probe on the host path: the TUI cleared the stale stamp as it started the turn, so + the user's close during it is an actual write with a fresh timestamp, and publication preserves it.""" + parent = "P_host_then_closed" + db.create_session(parent, source="tui") + _stamp(db, parent, "tui_close", age=60.0) + assert _host_reopen(db, parent) == "tui_close" + _admit(db, parent) + before = time.time() + db.end_session(parent, "tui_close") + with pytest.raises(RuntimeError, match=f"already ended.*{NOT_HEALED}"): + _publish(db, parent, "C_never", holder=_lease(db, parent)) + row = db.get_session(parent) + assert row["end_reason"] == "tui_close" and row["ended_at"] >= before, "the fresh close, not the stale one" + + @pytest.mark.parametrize("reason", STALE_EXPLICIT_CLOSES) + def test_without_the_host_a_stale_close_is_not_healed(self, db: SessionDB, reason: str) -> None: + parent = f"P_nohost_{reason}" + db.create_session(parent, source="tui") + _stamp(db, parent, reason, age=60.0) + _admit(db, parent) + holder = _lease(db, parent) # neither lease is evidence + with pytest.raises(RuntimeError, match=f"already ended.*{NOT_HEALED}"): + _publish(db, parent, f"C_{reason}", holder=holder) + with pytest.raises(RuntimeError, match=f"already ended.*{NOT_HEALED}"): + _publish(db, parent, f"C_{reason}") # lease-less publication: same contract + assert db.get_session(f"C_{reason}") is None + assert db.get_session(parent)["end_reason"] == reason + + +class TestLineageOwnedElsewhereStillFailsClosed: + def test_explicit_close_with_published_continuation_fails_closed(self, db: SessionDB) -> None: + parent = "P_superseded" + db.create_session(parent, source="tui") + _admit(db, parent) + _publish(db, parent, "C_first", holder=_lease(db, parent)) # a continuation now exists + db.release_compression_lock(parent, HOLDER) + db.release_session_turn_lease(parent, TURN) + db.reopen_session(parent) + _stamp(db, parent, "tui_close", age=60.0) + assert _host_reopen(db, parent) is None + _admit(db, parent, holder=TURN_2) + holder = _lease(db, parent, holder="second-writer") + + with pytest.raises(RuntimeError, match="already ended.*published continuation"): + _publish(db, parent, "C_second", holder=holder) + assert db.get_session("C_second") is None + assert db.get_session(parent)["end_reason"] == "tui_close" # untouched + + @pytest.mark.parametrize("reason", BOUNDARIES) + def test_boundary_fails_closed_even_after_a_host_reopen(self, db: SessionDB, reason: str) -> None: + assert reason in _RESET_END_REASONS or reason == "new_session" # CLI /new is a boundary too + parent = f"P_{reason}" + db.create_session(parent, source="tui") + _stamp(db, parent, reason, age=60.0) + assert _host_reopen(db, parent) is None + _admit(db, parent) + with pytest.raises(RuntimeError, match=f"already ended.*{reason} boundary"): + _publish(db, parent, f"C_{reason}", holder=_lease(db, parent)) + assert db.get_session(f"C_{reason}") is None + assert db.get_session(parent)["end_reason"] == reason + + def test_compression_stamp_fails_closed(self, db: SessionDB) -> None: + parent = "P_compression" + db.create_session(parent, source="tui") + _stamp(db, parent, "compression", age=60.0) + assert _host_reopen(db, parent) is None + _admit(db, parent) + with pytest.raises(RuntimeError, match="already ended.*closed by compression"): + _publish(db, parent, "C_x", holder=_lease(db, parent)) + + +class TestVerdictIsSharedWithTheAgentGuard: + """The pre-flush guard must agree with publish, or a durable flush is skipped (or written for nothing).""" + + def test_read_only_verdict_matches_publish(self, db: SessionDB) -> None: + for sid, reason in [("live", None), ("auto", "ws_disconnect")]: + db.create_session(sid, source="tui") + if reason: + _stamp(db, sid, reason, age=60.0) + assert db.compression_parent_deliberately_ended(sid) is False, (sid, reason) + + db.create_session("stale", source="tui") + _stamp(db, "stale", "tui_close", age=60.0) + _admit(db, "stale") + assert db.compression_parent_deliberately_ended("stale") is True # admission is not provenance + _host_reopen(db, "stale") + assert db.compression_parent_deliberately_ended("stale") is False + + db.create_session("closed_mid", source="tui") + _admit(db, "closed_mid") + _stamp(db, "closed_mid", "tui_close") + assert db.compression_parent_deliberately_ended("closed_mid") is True + + for sid, reason in [("reset", "session_reset"), ("new", "new_session"), ("comp", "compression")]: + db.create_session(sid, source="tui") + _stamp(db, sid, reason, age=60.0) + _host_reopen(db, sid) + assert db.compression_parent_deliberately_ended(sid) is True, (sid, reason) + + db.create_session("superseded", source="tui") + _publish(db, "superseded", "superseded_child", holder=_lease(db, "superseded")) + db.release_compression_lock("superseded", HOLDER) + db.reopen_session("superseded") + _stamp(db, "superseded", "cli_close", age=60.0) + _host_reopen(db, "superseded") + assert db.compression_parent_deliberately_ended("superseded") is True + assert db.compression_parent_deliberately_ended("missing") is False + assert db.compression_parent_deliberately_ended("") is False + + def test_agent_guard_delegates_to_the_store(self, db: SessionDB) -> None: + from agent.conversation_compression import _parent_deliberately_ended + + db.create_session("stale", source="tui") + _stamp(db, "stale", "tui_close", age=60.0) + _admit(db, "stale") + assert _parent_deliberately_ended(db, "stale") is True + _host_reopen(db, "stale") + assert _parent_deliberately_ended(db, "stale") is False + db.create_session("reset", source="tui") + _stamp(db, "reset", "session_reset", age=60.0) + assert _parent_deliberately_ended(db, "reset") is True + + def test_agent_guard_fails_open_and_keeps_the_taxonomy_fallback(self) -> None: + from agent.conversation_compression import _parent_deliberately_ended + + broken = MagicMock() + broken.compression_parent_deliberately_ended.side_effect = RuntimeError("db unavailable") + assert _parent_deliberately_ended(broken, "x") is False + + class TaxonomyOnlyStore: # a stand-in without the verdict: old behaviour is retained + def get_session(self, session_id): + return {"ended_at": 1.0, "end_reason": "tui_close"} + + assert _parent_deliberately_ended(TaxonomyOnlyStore(), "x") is True + + +class TestRotationEndToEnd: + def _build_agent(self, db: SessionDB, session_id: str): + with patch.dict(os.environ, {"OPENROUTER_API_KEY": "test-key"}): + from run_agent import AIAgent + + agent = AIAgent( + api_key="test-key", + base_url="https://openrouter.ai/api/v1", + model="test/model", + platform="tui", + quiet_mode=True, + session_db=db, + session_id=session_id, + skip_context_files=True, + skip_memory=True, + ) + compressor = MagicMock() + compressor.compress.return_value = [ + {"role": "user", "content": "[CONTEXT COMPACTION] summary"}, + {"role": "user", "content": "tail"}, + ] + compressor.compression_count = 1 + compressor.last_prompt_tokens = 0 + compressor.last_completion_tokens = 0 + compressor._last_summary_error = None + compressor._last_compress_aborted = False + compressor._last_summary_auth_failure = False + compressor._last_aux_model_failure_model = None + compressor._last_aux_model_failure_error = None + agent.context_compressor = compressor + agent.compression_in_place = False # rotation path + return agent + + @staticmethod + def _admit_turn(db: SessionDB, agent, session_id: str) -> None: + """What ``admit_durable_turn_lease`` does at the start of ``run_conversation``.""" + _admit(db, session_id) + agent._active_session_turn_lease_holder = TURN + agent._active_session_turn_lease_ttl_seconds = 300.0 + + @staticmethod + def _compress(agent) -> None: + msgs = [{"role": "user", "content": f"m{i}"} for i in range(20)] + agent._compress_context(list(msgs), "sys", approx_tokens=120_000) + + def test_a_stale_close_cleared_by_the_host_does_not_wedge_rotation(self, db: SessionDB) -> None: + """#106459 end-to-end: the live session's row carries an explicit-close stamp that predates the + turn; the routing host clears it as the turn starts, and auto-compaction rotates.""" + parent = "PARENT_106459_E2E" + db.create_session(parent, source="tui") + agent = self._build_agent(db, parent) + _stamp(db, parent, "tui_close", age=60.0) + _host_reopen(db, parent) + self._admit_turn(db, agent, parent) + + self._compress(agent) + + assert agent.session_id != parent, "rotation aborted on the stale explicit-close stamp" + assert db.get_session(parent)["end_reason"] == "compression" + child_row = db.get_session(agent.session_id) + assert child_row is not None and child_row["parent_session_id"] == parent + + def test_without_the_host_an_admitted_turn_does_not_rotate_over_a_close(self, db: SessionDB) -> None: + """The inverse of the previous revision's contract: a turn lease alone never clears the stamp.""" + parent = "PARENT_106459_ADMIT_ONLY" + db.create_session(parent, source="tui") + agent = self._build_agent(db, parent) + _stamp(db, parent, "tui_close", age=60.0) + self._admit_turn(db, agent, parent) + + self._compress(agent) + + assert agent.session_id == parent + assert db.get_session(parent)["end_reason"] == "tui_close" + assert db._read_one("SELECT COUNT(*) FROM sessions WHERE parent_session_id = ?", (parent,))[0] == 0 + + def test_close_during_a_turn_the_host_reopened_aborts_rotation(self, db: SessionDB) -> None: + """Third probe end-to-end on the host path: stale stamp, host reopen, admit, close, compress.""" + parent = "PARENT_106459_HOST_THEN_CLOSED" + db.create_session(parent, source="tui") + agent = self._build_agent(db, parent) + _stamp(db, parent, "tui_close", age=60.0) + _host_reopen(db, parent) + self._admit_turn(db, agent, parent) + before = time.time() + db.end_session(parent, "tui_close") # session.close during the turn + + self._compress(agent) + + assert agent.session_id == parent, "a close made during the turn must not be resurrected" + row = db.get_session(parent) + assert row["end_reason"] == "tui_close" and row["ended_at"] >= before + assert db._read_one("SELECT COUNT(*) FROM sessions WHERE parent_session_id = ?", (parent,))[0] == 0 + + def test_close_after_the_turn_was_admitted_aborts_rotation(self, db: SessionDB) -> None: + """Second probe end-to-end: close first, the still-running turn then compresses and takes the + lease second; publish creates no child.""" + parent = "PARENT_106459_CLOSE_FIRST" + db.create_session(parent, source="tui") + agent = self._build_agent(db, parent) + self._admit_turn(db, agent, parent) + db.end_session(parent, "tui_close") + + self._compress(agent) + + assert agent.session_id == parent + assert db.get_session(parent)["end_reason"] == "tui_close" + assert db._read_one("SELECT COUNT(*) FROM sessions WHERE parent_session_id = ?", (parent,))[0] == 0 + + def test_close_during_compression_still_aborts_rotation(self, db: SessionDB) -> None: + """First probe end-to-end: the user closes while the summary is being produced.""" + parent = "PARENT_106459_CLOSE_MID" + db.create_session(parent, source="tui") + agent = self._build_agent(db, parent) + self._admit_turn(db, agent, parent) + compressor = agent.context_compressor + summary = compressor.compress.return_value + + def close_while_summarizing(*args, **kwargs): + db.end_session(parent, "tui_close") # session.close lands after the lease was acquired + return summary + + compressor.compress.side_effect = close_while_summarizing + self._compress(agent) + + assert agent.session_id == parent, "a close made mid-compression must not be resurrected" + assert db.get_session(parent)["end_reason"] == "tui_close" + + @pytest.mark.parametrize("reason", ["session_reset", "new_session"]) + def test_boundary_still_aborts_rotation(self, db: SessionDB, reason: str) -> None: + parent = f"PARENT_106459_{reason}" + db.create_session(parent, source="tui") + agent = self._build_agent(db, parent) + _stamp(db, parent, reason, age=60.0) + _host_reopen(db, parent) + self._admit_turn(db, agent, parent) + + self._compress(agent) + + assert agent.session_id == parent + assert db.get_session(parent)["end_reason"] == reason diff --git a/tests/tui_gateway/test_106459_routing_provenance_reopen.py b/tests/tui_gateway/test_106459_routing_provenance_reopen.py new file mode 100644 index 0000000000..366f0bbb3a --- /dev/null +++ b/tests/tui_gateway/test_106459_routing_provenance_reopen.py @@ -0,0 +1,324 @@ +"""#106459 routing provenance: the TUI clears a stale explicit-close stamp only for a session it +still has registered, under ``_sessions_lock`` as it starts a turn, and never once the session has +been claimed for teardown. + +Fourth review probe on #106543: the TUI accepts a prompt and starts a worker before that worker +reaches ``run_conversation()`` and takes its turn lease, so ``session.close`` can pop the session, +wait out its grace and stamp ``tui_close`` in between. Clearing the stamp on lease acquisition +turned that deliberate close back into a live row. The clear now happens in ``_run_prompt_submit`` +while ``_sessions_lock`` is held and the session is still registered -- the lock +``_pop_session_by_id`` claims teardown under -- and a late lease clears nothing. +""" + +from __future__ import annotations + +import contextlib +import os +import sys +import threading +import time +from pathlib import Path + +import pytest + +from hermes_state import SessionDB +from tui_gateway import server + +TURN = f"pid={os.getpid()}:turn=tui:platform=tui" +# Captured at import: several tests replace ``threading.Thread`` (module-global) with a synchronous stand-in, +# and a lock probe must run on a genuinely different thread or an RLock simply re-enters. +_RealThread = threading.Thread + + +class _ImmediateThread: + """Run the turn inside ``start()`` so tests observe its final state synchronously.""" + + def __init__(self, target=None, daemon=None, **_kwargs): + self._target = target + + def start(self): + if self._target is not None: + self._target() + + def is_alive(self): + return False + + def join(self, timeout=None): + return None + + +class _Agent: + model = "test-model" + provider = "test-provider" + + def __init__(self, session_id: str, db: SessionDB, *, before_lease=None): + self.session_id = session_id + self._db = db + self._before_lease = before_lease + self.turns: list = [] + self.stamp_seen_by_turn: list = [] + + def clear_interrupt(self): + return None + + def run_conversation(self, prompt, conversation_history=None, stream_callback=None, **_kwargs): + if self._before_lease is not None: + self._before_lease() + # What admit_durable_turn_lease does at the top of the real run_conversation. + self._db.try_acquire_session_turn_lease(self.session_id, TURN, ttl_seconds=300.0) + self.stamp_seen_by_turn.append(self._db.get_session(self.session_id)["end_reason"]) + self.turns.append(prompt) + return {"final_response": "", "messages": []} + + +def _session(agent: _Agent, **extra) -> dict: + return { + "agent": agent, + "session_key": agent.session_id, + "history": [], + "history_lock": threading.Lock(), + "history_version": 0, + "running": True, + "attached_images": [], + "image_counter": 0, + "cols": 80, + "slash_worker": None, + "show_reasoning": False, + "tool_progress_mode": "all", + "inflight_turn": None, + **extra, + } + + +@pytest.fixture +def db(tmp_path: Path): + handle = SessionDB(db_path=tmp_path / "state.db") + try: + yield handle + finally: + handle.close() + + +@pytest.fixture +def turn_env(monkeypatch, tmp_path, db): + """The immediate-prompt harness of tests/test_tui_gateway_server.py, with a real SessionDB.""" + monkeypatch.setattr(server, "_emit", lambda *_a, **_k: None) + monkeypatch.setattr(server, "make_stream_renderer", lambda _cols: None) + monkeypatch.setattr(server, "render_message", lambda _raw, _cols: None) + monkeypatch.setattr(server, "_wire_callbacks", lambda _sid: None) + monkeypatch.setattr(server, "_sync_agent_model_with_config", lambda *_a: None) + monkeypatch.setattr(server, "_session_cwd", lambda _session: str(tmp_path)) + monkeypatch.setattr(server, "_register_session_cwd", lambda _session: None) + monkeypatch.setattr(server, "_set_session_context", lambda *_a, **_k: []) + monkeypatch.setattr(server, "_clear_session_context", lambda _tokens: None) + monkeypatch.setattr(server, "_session_info", lambda *_a: {}) + monkeypatch.setattr(server, "_get_usage", lambda _agent: {}) + monkeypatch.setattr(server, "_sync_session_key_after_compress", lambda *_a, **_k: None) + monkeypatch.setattr(server, "_drain_queued_prompt", lambda *_a: False) + monkeypatch.setattr(server, "_voice_tts_enabled", lambda: False) + monkeypatch.setattr(server, "_get_db", lambda: db) + return monkeypatch + + +@pytest.fixture +def registry(): + """Sessions registered by a test, removed afterwards.""" + added: list[str] = [] + + def _add(sid: str, session: dict) -> dict: + server._sessions[sid] = session + added.append(sid) + return session + + yield _add + for sid in added: + server._sessions.pop(sid, None) + + +def _stamp(db: SessionDB, row_id: str, reason: str = "tui_close", *, age: float = 60.0) -> None: + db.end_session(row_id, reason) + if age: + db._write_sql("UPDATE sessions SET ended_at = ? WHERE id = ?", (time.time() - age, row_id)) + + +def test_a_registered_session_is_reopened_before_its_turn_starts(turn_env, db, registry): + """The #106459 field shape: the TUI keeps accepting prompts on a session whose row carries a + stale ``tui_close``. The row must already be clear when the worker runs.""" + turn_env.setattr(server.threading, "Thread", _ImmediateThread) + db.create_session("row-1", source="tui") + _stamp(db, "row-1") + agent = _Agent("row-1", db) + session = registry("ui-1", _session(agent)) + + assert server._run_prompt_submit("rid", "ui-1", session, "go") is True + + assert agent.turns == ["go"] + assert agent.stamp_seen_by_turn == [None], "the stamp must be cleared before the worker starts" + row = db.get_session("row-1") + assert row["ended_at"] is None and row["end_reason"] is None + + +@pytest.mark.parametrize("reason", ["session_reset", "new_session", "compression", "ws_disconnect"]) +def test_only_explicit_closes_are_reopened(turn_env, db, registry, reason): + turn_env.setattr(server.threading, "Thread", _ImmediateThread) + db.create_session("row-2", source="tui") + _stamp(db, "row-2", reason) + session = registry("ui-2", _session(_Agent("row-2", db))) + + server._run_prompt_submit("rid", "ui-2", session, "go") + + assert db.get_session("row-2")["end_reason"] == reason + + +def test_an_unregistered_session_runs_its_turn_but_keeps_its_stamp(turn_env, db): + """``can_start`` also admits a session that is not in the registry; nothing proves it is routed + here, so its turn runs as before but its stamp is not cleared.""" + turn_env.setattr(server.threading, "Thread", _ImmediateThread) + db.create_session("row-3", source="tui") + _stamp(db, "row-3") + agent = _Agent("row-3", db) + + assert server._run_prompt_submit("rid", "ui-3", _session(agent), "go") is True + + assert agent.turns == ["go"] + assert agent.stamp_seen_by_turn == ["tui_close"] + assert db.get_session("row-3")["end_reason"] == "tui_close" + + +def test_a_session_already_claimed_for_teardown_is_not_reopened(turn_env, db, registry): + turn_env.setattr(server.threading, "Thread", _ImmediateThread) + db.create_session("row-4", source="tui") + _stamp(db, "row-4") + agent = _Agent("row-4", db) + session = registry("ui-4", _session(agent)) + assert server._pop_session_by_id("ui-4") is session # session.close claimed it + + assert server._run_prompt_submit("rid", "ui-4", session, "go") is False + + assert agent.turns == [] + assert db.get_session("row-4")["end_reason"] == "tui_close" + + +def test_a_close_that_wins_the_race_after_admission_is_not_reopened(turn_env, db, registry): + """``_admit_prompt_turn`` checks ``_closing`` under ``history_lock`` only, which does not exclude + ``_pop_session_by_id``. A close claimed between admission and the start gate must keep its row's + stamp -- the heal belongs under ``_sessions_lock`` with the gate, not in admission.""" + db.create_session("row-5", source="tui") + _stamp(db, "row-5") + agent = _Agent("row-5", db) + session = registry("ui-5", _session(agent)) + emit_entered, release_emit = threading.Event(), threading.Event() + + def _blocking_emit(event, *_a, **_k): + if event == "message.start": + emit_entered.set() + assert release_emit.wait(timeout=2.0) + + turn_env.setattr(server, "_emit", _blocking_emit) + results: list = [] + dispatch = threading.Thread(target=lambda: results.append( + server._run_prompt_submit("rid", "ui-5", session, "go"))) + try: + dispatch.start() + assert emit_entered.wait(timeout=1.0) + assert server._pop_session_by_id("ui-5") is session + finally: + release_emit.set() + dispatch.join(timeout=2.0) + + assert results == [False] + assert agent.turns == [] + assert db.get_session("row-5")["end_reason"] == "tui_close" + + +def test_a_late_worker_lease_after_session_close_leaves_the_close(turn_env, db, registry): + """Fourth review probe: the prompt is accepted and the worker started, the user closes the session + (pop, bounded grace, ``tui_close``) while the worker is still in pre-turn work, and only then does + the worker take its turn lease. Nothing re-applies the close afterwards, so it must survive.""" + db.create_session("row-6", source="tui") + worker_waiting, release_worker = threading.Event(), threading.Event() + + def _hold_before_lease(): + worker_waiting.set() + assert release_worker.wait(timeout=5.0) + + agent = _Agent("row-6", db, before_lease=_hold_before_lease) + session = registry("ui-6", _session(agent)) + stamped: list = [] + + def _teardown(popped, *, end_reason="tui_close"): + db.end_session(popped["agent"].session_id, end_reason) # what _finalize_session writes + stamped.append(end_reason) + + turn_env.setattr(server, "_teardown_session", _teardown) + turn_env.setattr(server, "_TURN_SETTLE_BEFORE_CLOSE_SECONDS", 0.2, raising=False) + try: + assert server._run_prompt_submit("rid", "ui-6", session, "go") is True + assert worker_waiting.wait(timeout=2.0) + assert server._close_session_by_id("ui-6") is True + assert stamped == ["tui_close"] + finally: + release_worker.set() + worker = session.get("_run_thread") + if worker is not None: + worker.join(timeout=5.0) + + assert agent.turns == ["go"] # the late worker still runs its turn, as at the merge base + assert agent.stamp_seen_by_turn == ["tui_close"] + assert db.get_session("row-6")["end_reason"] == "tui_close" + + +def test_the_session_db_is_resolved_before_the_sessions_lock_is_taken(turn_env, db, registry): + """Resolving a profile session's handle goes through the state registry; it must not run under the + lock that gates every create/close/prompt on this backend.""" + turn_env.setattr(server.threading, "Thread", _ImmediateThread) + db.create_session("row-7", source="tui") + _stamp(db, "row-7") + session = registry("ui-7", _session(_Agent("row-7", db))) + real_session_db = server._session_db + lock_free_at_resolution: list[bool] = [] + + @contextlib.contextmanager + def _probing_session_db(sess): + callers = {sys._getframe(depth).f_code.co_name for depth in range(1, 5)} + if "_run_prompt_submit" in callers: + probe: list[bool] = [] + + def _try_lock_from_another_thread(): + # Acquire and release on the SAME thread: an RLock left owned by a finished thread stays held. + got = server._sessions_lock.acquire(timeout=0.5) + probe.append(got) + if got: + server._sessions_lock.release() + + prober = _RealThread(target=_try_lock_from_another_thread) + prober.start() + prober.join() + lock_free_at_resolution.append(bool(probe and probe[0])) + with real_session_db(sess) as handle: + yield handle + + turn_env.setattr(server, "_session_db", _probing_session_db) + + assert server._run_prompt_submit("rid", "ui-7", session, "go") is True + + assert lock_free_at_resolution == [True] + assert db.get_session("row-7")["end_reason"] is None + + +def test_a_failed_routing_reopen_never_blocks_the_turn(turn_env, db, registry): + turn_env.setattr(server.threading, "Thread", _ImmediateThread) + db.create_session("row-8", source="tui") + _stamp(db, "row-8") + agent = _Agent("row-8", db) + session = registry("ui-8", _session(agent)) + + def _unavailable(*_a, **_k): + raise RuntimeError("database is locked") + + turn_env.setattr(db, "reopen_if_explicitly_closed", _unavailable) + + assert server._run_prompt_submit("rid", "ui-8", session, "go") is True + + assert agent.turns == ["go"] + assert db.get_session("row-8")["end_reason"] == "tui_close" diff --git a/tests/tui_gateway/test_bot_live_owner_delivery.py b/tests/tui_gateway/test_bot_live_owner_delivery.py index 9484d49dc5..34370495eb 100644 --- a/tests/tui_gateway/test_bot_live_owner_delivery.py +++ b/tests/tui_gateway/test_bot_live_owner_delivery.py @@ -8,6 +8,7 @@ from tui_gateway.turn_marker import record_turn_start, read_turn_marker def test_refused_input_commits_failed_mailbox_receipt(tmp_path): + import contextlib import contextvars import logging import time @@ -34,6 +35,8 @@ def test_refused_input_commits_failed_mailbox_receipt(tmp_path): "_finish_turn": noop, "_clear_inflight_turn": noop, "_retire_turn_marker": lambda *args: retired.append(args), "_emit_settled_session_info": noop, + "_routing_provenance_db": lambda _session: contextlib.nullcontext(None), + "_reopen_routed_session_row": noop, }) def terminal(outcome): mailbox.complete_delivery(tmp_path, queued["id"], status=outcome["status"], diff --git a/tools/session_search_tool.py b/tools/session_search_tool.py index b8a34be6e8..aa307649e3 100644 --- a/tools/session_search_tool.py +++ b/tools/session_search_tool.py @@ -13,7 +13,7 @@ import logging from datetime import datetime from typing import Any, Dict, List, Optional, Union -from hermes_state_common import _RESET_END_REASONS +from hermes_state_common import _BOUNDARY_END_REASONS # Hidden from browsing/searching — integrations (HERMES_SESSION_SOURCE=tool), delegate # subagent runs, kanban workers are not the user's history. @@ -36,8 +36,8 @@ _DISCOVER_SEARCH_FIELDS = ("id", "session_id", "role", "snippet", "source", "mod _COMPACTION_PREFIXES = ("[CONTEXT COMPACTION", "[CONTEXT SUMMARY]:") # /new, /reset, idle/daily expiry and CLI /new ("new_session") end the predecessor WITHOUT # carrying its transcript forward — unlike compression continuations and live delegation -# children. Derived from the gateway set so the two cannot drift. -_FRESH_RESET_END_REASONS = frozenset(_RESET_END_REASONS) | {"new_session"} +# children. The store's boundary set, so the two cannot drift. +_FRESH_RESET_END_REASONS = _BOUNDARY_END_REASONS def _quiet(fn, default, msg, *log_args, with_exc: bool = False): diff --git a/tui_gateway/prompt_turn.py b/tui_gateway/prompt_turn.py index 987aa66ebd..95c981d3fe 100644 --- a/tui_gateway/prompt_turn.py +++ b/tui_gateway/prompt_turn.py @@ -747,6 +747,37 @@ def _finish_turn(sid: str, session: dict, st: _TurnRun) -> None: _clear_session_context(scopes.session_tokens) +# Bounded so a contended state.db cannot hold ``_sessions_lock``; a skipped heal is retried on the next prompt. +_ROUTING_REOPEN_PATIENCE_S = 0.5 + + +def _routing_provenance_db(session: dict): + """The session's own SessionDB for :func:`_reopen_routed_session_row`, or a ``None`` context.""" + try: + return _session_db(session) + except Exception: + logger.debug("could not resolve the session db for routing provenance", exc_info=True) + return contextlib.nullcontext(None) + + +def _reopen_routed_session_row(db, sid: str, session: dict) -> None: + """Routing provenance for #106459, called under ``_sessions_lock`` while this backend still has the + session registered and is about to start a turn for it: an explicit-close stamp on its row missed a + live conversation, so it is cleared (the gateway's #54878 stale-route self-heal, which this backend + lacked). ``_pop_session_by_id`` claims teardown under the same lock and sets ``_closing`` before + ``_finalize_session`` stamps, so a close of THIS session is either not written yet or already stops + the turn -- it is never cleared here. Best-effort: a failure never blocks the turn.""" + session_id = getattr(session.get("agent"), "session_id", None) # the compression tip this turn writes + if db is None or not session_id: + return + try: + db.reopen_if_explicitly_closed( + session_id, provenance=f"TUI session {sid} is still registered and accepting a turn", + patience_s=_ROUTING_REOPEN_PATIENCE_S) + except Exception: + logger.debug("routing-provenance reopen failed for %s", session_id, exc_info=True) + + def _run_prompt_submit( rid, sid: str, session: dict, text: Any, *, display_kind: str | None = None, display_metadata: dict | None = None, image_paths: list[str] | None = None, @@ -842,10 +873,16 @@ def _run_prompt_submit( _emit_settled_session_info(sid, session, st.agent) _run_post_turn_followups(rid, sid, session, st.result, goal_followup) run_thread = threading.Thread(target=run, daemon=True) - with _sessions_lock: + # The handle is resolved BEFORE _sessions_lock: a profile session opens its own SessionDB through the + # state registry, and _sessions_lock gates every create/close/prompt on this backend. + with _routing_provenance_db(session) as routing_db, _sessions_lock: registered = _sessions.get(sid) can_start = not session.get("_closing") and (registered is None or registered is session) if can_start: + # Only a registered session is proof the conversation is routed here; an unregistered one may + # still run its turn, but its stamp stays (#106459). + if registered is session: + _reopen_routed_session_row(routing_db, sid, session) session["_run_thread"] = run_thread run_thread.start() if not can_start: