diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 851d2733c8..1f49aa9060 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -2899,18 +2899,8 @@ def _salvage_or_refuse_grown_transcript( def _parent_deliberately_ended(session_db: Any, session_id: str) -> bool: - """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. + """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.""" reader = getattr(session_db, "get_session", None) if not callable(reader): return False diff --git a/hermes_state_compression.py b/hermes_state_compression.py index 0ddbe62013..c902bfed1f 100644 --- a/hermes_state_compression.py +++ b/hermes_state_compression.py @@ -12,8 +12,7 @@ import time from typing import Any, Dict, List, Optional, Tuple from hermes_state_common import ( - _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) + _BOUNDARY_END_REASONS, _COMPRESSION_LOCK_ROW_SQL as _LOCK_ROW_SQL, _ENDED_ROW_SQL, _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,84 +69,42 @@ 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. + session a HOST has just proven is still routed to it, returning the reason cleared or None (#106459). + Narrow twin of ``reopen_session()``: automatic stamps are left to publish (#88197); ``compression``, + boundary (reset reasons, CLI ``new_session``) and stamps with a published continuation own lineage + elsewhere and are never touched. The UPDATE is conditional on the exact stamp read, so a close + landing between read and 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.""" + Only the routing host can make this call. Publication cannot: ``end_session()`` is first-stamp-wins, + so a close made during a turn that began on a stale stamp is a no-op write. Turn-lease admission + cannot: the TUI starts its worker before it reaches ``run_conversation()``, so ``session.close`` can + stamp ``tui_close`` in between and a late lease would clear a deliberate close. Call it under the + lock that makes the host's registry claim atomic with its teardown (#54878 on the routing table).""" 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": + if row is None or row["ended_at"] is None: + return None + reason = row["end_reason"] + if is_automatic_end_reason(reason) or reason == "compression" or reason in _BOUNDARY_END_REASONS: + return None + superseded = 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() + if superseded is not None: 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"]) + (session_id, row["ended_at"], reason)) + return str(reason) reason = self._execute_write(_do, patience_s=patience_s) if reason is not None: logger.warning( @@ -155,20 +112,6 @@ class SessionCompressionMixin: "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 @@ -321,13 +264,12 @@ 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-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})") + # 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}") 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 deleted file mode 100644 index fb129fbf12..0000000000 --- a/tests/hermes_state/test_106459_stale_explicit_close_stamp.py +++ /dev/null @@ -1,552 +0,0 @@ -"""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 index 366f0bbb3a..2570c54ed2 100644 --- a/tests/tui_gateway/test_106459_routing_provenance_reopen.py +++ b/tests/tui_gateway/test_106459_routing_provenance_reopen.py @@ -1,20 +1,12 @@ -"""#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. +"""#106459: a stale explicit-close stamp (``tui_close`` on a session the TUI still routes) wedged +compression forever, because publish fails closed on it and ``end_session()`` is first-stamp-wins. +The TUI clears it under ``_sessions_lock`` as it starts a turn for a still-registered session -- +and only then: the store never decides on its own that an explicit close is stale. """ from __future__ import annotations -import contextlib import os -import sys import threading import time from pathlib import Path @@ -25,9 +17,6 @@ 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: @@ -51,10 +40,9 @@ class _Agent: model = "test-model" provider = "test-provider" - def __init__(self, session_id: str, db: SessionDB, *, before_lease=None): + def __init__(self, session_id: str, db: SessionDB): self.session_id = session_id self._db = db - self._before_lease = before_lease self.turns: list = [] self.stamp_seen_by_turn: list = [] @@ -62,8 +50,6 @@ class _Agent: 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"]) @@ -71,22 +57,11 @@ class _Agent: return {"final_response": "", "messages": []} -def _session(agent: _Agent, **extra) -> dict: +def _session(agent: _Agent) -> 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, + "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, } @@ -102,223 +77,61 @@ def db(tmp_path: Path): @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: + for name, stub in { + "_emit": lambda *_a, **_k: None, "make_stream_renderer": lambda _cols: None, + "render_message": lambda _raw, _cols: None, "_wire_callbacks": lambda _sid: None, + "_sync_agent_model_with_config": lambda *_a: None, "_session_cwd": lambda _session: str(tmp_path), + "_register_session_cwd": lambda _session: None, "_set_session_context": lambda *_a, **_k: [], + "_clear_session_context": lambda _tokens: None, "_session_info": lambda *_a: {}, + "_get_usage": lambda _agent: {}, "_sync_session_key_after_compress": lambda *_a, **_k: None, + "_drain_queued_prompt": lambda *_a: False, "_voice_tts_enabled": lambda: False, "_get_db": lambda: db, + }.items(): + monkeypatch.setattr(server, name, stub) + monkeypatch.setattr(server.threading, "Thread", _ImmediateThread) + yield monkeypatch + for sid in [s for s in server._sessions if s.startswith("ui-106459-")]: server._sessions.pop(sid, None) -def _stamp(db: SessionDB, row_id: str, reason: str = "tui_close", *, age: float = 60.0) -> None: +def _stamp(db: SessionDB, row_id: str, reason: str) -> None: db.end_session(row_id, reason) - if age: - db._write_sql("UPDATE sessions SET ended_at = ? WHERE id = ?", (time.time() - age, row_id)) + db._write_sql("UPDATE sessions SET ended_at = ? WHERE id = ?", (time.time() - 60.0, 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) +def test_a_registered_stale_close_is_cleared_before_the_turn_and_rotation_publishes(turn_env, db): + """The field shape: the TUI keeps accepting prompts on a row stamped ``tui_close`` an hour ago. + The row is clear when the worker runs, and the rotation that follows publishes instead of wedging.""" db.create_session("row-1", source="tui") - _stamp(db, "row-1") + _stamp(db, "row-1", "tui_close") agent = _Agent("row-1", db) - session = registry("ui-1", _session(agent)) + session = server._sessions["ui-106459-1"] = _session(agent) - assert server._run_prompt_submit("rid", "ui-1", session, "go") is True + assert server._run_prompt_submit("rid", "ui-106459-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 + assert db.try_acquire_compression_lock("row-1", "holder") + db.publish_compression_child( + parent_session_id="row-1", child_session_id="row-1-child", source="tui", + messages=[{"role": "user", "content": "go"}], model="m", compression_lock_holder="holder") + assert db.get_session("row-1")["end_reason"] == "compression" -@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) +@pytest.mark.parametrize("reason, claimed_for_teardown", [ + ("session_reset", False), ("new_session", False), ("compression", False), ("ws_disconnect", False), + ("tui_close", True), +]) +def test_only_an_explicit_close_on_a_still_registered_session_is_cleared(turn_env, db, reason, claimed_for_teardown): + """Boundaries, compression and automatic stamps own lineage elsewhere; a session ``session.close`` + already claimed under ``_sessions_lock`` is refused a turn and its deliberate close survives.""" db.create_session("row-2", source="tui") _stamp(db, "row-2", reason) - session = registry("ui-2", _session(_Agent("row-2", db))) + agent = _Agent("row-2", db) + session = server._sessions["ui-106459-2"] = _session(agent) + if claimed_for_teardown: + assert server._pop_session_by_id("ui-106459-2") is session - server._run_prompt_submit("rid", "ui-2", session, "go") + started = server._run_prompt_submit("rid", "ui-106459-2", session, "go") + assert started is not claimed_for_teardown 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"