refactor(sessions): trim #106543 to the host-side heal and two invariants
The publish-side taxonomy (_end_stamp_class, _compression_parent_obstacle, compression_parent_deliberately_ended) and the agent-guard delegation produced exactly main's verdict -- every non-automatic stamp still fails closed -- so they only changed an error string. Inline the "explicit close with no continuation" test into reopen_if_explicitly_closed(), restore main's publish branch and agent guard untouched, and keep two tests: the field shape rotates after the host clears the stamp; boundary/compression/automatic stamps and a session already claimed for teardown are never cleared.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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,))
|
||||
|
||||
@@ -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
|
||||
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user