diff --git a/hermes_state_sessions.py b/hermes_state_sessions.py index 7f65916994..2f2c6b2e9d 100644 --- a/hermes_state_sessions.py +++ b/hermes_state_sessions.py @@ -282,7 +282,10 @@ class SessionSessionsMixin: git_repo_root: str = None, origin_json: str = None, display_name: str = None, ) -> None: """Upsert a session row, never overwriting what an earlier writer set (the gateway creates a - bare row before create_session carries the real model/prompt). chat_id/thread_id scope gateway + bare row before create_session carries the real model/prompt) — the one exception is the + token-accounting guard's placeholder ``source='unknown'``, which a later writer's real surface + replaces (#111999): once minted, that placeholder otherwise labelled a real session anonymous + for life, because this upsert is the only writer that could correct it. chat_id/thread_id scope gateway /resume (IDOR). Children backfill from the parent; a missing profile_name is stamped with THIS store's own (NULL reads as unowned). @@ -320,6 +323,11 @@ class SessionSessionsMixin: ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, NULL, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(id) DO UPDATE SET + source = CASE + WHEN sessions.source = 'unknown' + THEN COALESCE(excluded.source, 'unknown') + ELSE sessions.source + END, model = COALESCE(sessions.model, excluded.model), model_config = CASE WHEN excluded.model_config IS NOT NULL diff --git a/hermes_state_usage.py b/hermes_state_usage.py index f272b76070..514f9937a2 100644 --- a/hermes_state_usage.py +++ b/hermes_state_usage.py @@ -284,7 +284,11 @@ class SessionUsageMixin: where the cached agent holds cumulative totals).""" usage = {k: v for k, v in locals().items() if k in _MODEL_USAGE_FIELDS} # Ensure the row exists: under concurrent load create_session() may have failed on - # locking, and the UPDATE would silently affect 0 rows. + # locking, and the UPDATE would silently affect 0 rows. The minted row carries the + # placeholder ``unknown`` source; a later writer's real surface replaces it in + # _insert_session_row's upsert, so the placeholder cannot outlive the session's creator + # (#111999). Until then the token guard is the only thing holding the row — never a + # session the user is shown as theirs. self._insert_session_row(session_id, "unknown", model=model) sql = _TOKEN_UPDATE_ABSOLUTE_SQL if absolute else _TOKEN_UPDATE_DELTA_SQL has_usage = bool(input_tokens or output_tokens or cache_read_tokens or cache_write_tokens or reasoning_tokens @@ -380,7 +384,8 @@ class SessionUsageMixin: if not session_id or not task: return usage["api_call_count"] = 1 if api_call_count is None else int(api_call_count) - # FK to sessions.id: same INSERT OR IGNORE guard as update_token_counts. + # FK to sessions.id: same INSERT OR IGNORE guard as update_token_counts (its placeholder + # source is repairable by the session's real creator — see _insert_session_row). self._insert_session_row(session_id, "unknown") self._execute_write(lambda conn: self._record_model_usage(conn, session_id, task=task, **usage)) diff --git a/tests/tui_gateway/test_bot_live_owner_delivery.py b/tests/tui_gateway/test_bot_live_owner_delivery.py index c79e9c6787..f16d7c212b 100644 --- a/tests/tui_gateway/test_bot_live_owner_delivery.py +++ b/tests/tui_gateway/test_bot_live_owner_delivery.py @@ -37,6 +37,8 @@ def test_refused_input_commits_failed_mailbox_receipt(tmp_path): "_emit_settled_session_info": noop, "_routing_provenance_db": lambda _session: contextlib.nullcontext(None), "_reopen_routed_session_row": noop, + # Every dispatch binds the session's own row before the turn writes (#111999). + "_ensure_session_db_row": noop, }) def terminal(outcome): mailbox.complete_delivery(tmp_path, queued["id"], status=outcome["status"], diff --git a/tests/tui_gateway/test_stream_interrupt_recovery_orphan.py b/tests/tui_gateway/test_stream_interrupt_recovery_orphan.py new file mode 100644 index 0000000000..ee84da6764 --- /dev/null +++ b/tests/tui_gateway/test_stream_interrupt_recovery_orphan.py @@ -0,0 +1,136 @@ +"""A stream that dies mid-answer must not leave an anonymous session behind (#111999). + +The two recovery injections — the gateway's crash auto-continue note and the loop's "Continue exactly +where you left off" stub — both write into the session that was interrupted. When that session's +durable row has not landed yet, three separate facts turned the recovery into a phantom: + +* ``update_token_counts`` is the only writer that mints ``source='unknown'`` (its "ensure the row + exists" guard, reached from the first API call's queued token delta); +* the row upsert keeps whatever the FIRST writer set, so the session's own creator could never repair + that placeholder — the record stayed anonymous for life; +* the startup orphan sweep skipped ``unknown``, so such a row stayed ``ended_at IS NULL`` forever. + +Contracts pinned here: + +* the recovery dispatch persists the session's OWN row (original session key, real source) before the + turn writes anything — so the recovery is written into the original session record; +* the session's real creator repairs the accounting guard's placeholder source on the same id; +* the startup sweep collects a phantom an older build already left on disk. +""" + +from __future__ import annotations + +import threading +import time +import types +from pathlib import Path + +import pytest + +from hermes_state import SessionDB +from hermes_state_registry import acquire, release_or_close +from tui_gateway import server +from tui_gateway.session_reaper import _ORPHAN_SWEEP_SOURCES + +# One of the four orphan ids from the report. +ORPHAN_SID = "20260913_210721_c89ac8" +IDLE_S = 6 * 3600 # mirror the TUI gateway's default session TTL + + +class _InlineThread: + """Run threads synchronously so tests observe final state.""" + + def __init__(self, target=None, daemon=None, args=(), kwargs=None): + self._target, self._args, self._kwargs = target, args, kwargs or {} + + def start(self): + if self._target is not None: + self._target(*self._args, **self._kwargs) + + def is_alive(self): + return False + + def join(self, timeout=None): + return None + + +def _session(agent, **extra): + return { + "agent": agent, "session_key": ORPHAN_SID, "history": [], "history_lock": threading.Lock(), + "history_version": 0, "running": False, "attached_images": [], "image_counter": 0, "cols": 80, + "slash_worker": None, "show_reasoning": False, "tool_progress_mode": "all", "inflight_turn": None, + "source": "desktop", **extra, + } + + +@pytest.fixture() +def recovery_env(monkeypatch, tmp_path): + """Neutralize the turn pipeline's environment-heavy side paths (same set the auto-continue suite uses).""" + monkeypatch.setattr(server, "threading", types.SimpleNamespace(Thread=_InlineThread)) + monkeypatch.setattr(server, "_hermes_home", tmp_path) + monkeypatch.setattr(server, "_wire_callbacks", lambda sid: None) + monkeypatch.setattr(server, "_sync_agent_model_with_config", lambda sid, session: None) + monkeypatch.setattr(server, "_session_cwd", lambda session: str(tmp_path)) + monkeypatch.setattr(server, "_register_session_cwd", lambda session: None) + monkeypatch.setattr(server, "_tts_stream_begin", lambda: None) + monkeypatch.setattr(server, "_sync_session_key_after_compress", lambda *a, **k: None) + monkeypatch.setattr(server, "_get_usage", lambda agent: {}) + + +def _recovery_dispatch(server_module, session, text): + """Dispatch a turn the way the crash auto-continue / queued-prompt drain do (straight into the turn + pipeline, bypassing prompt.submit's row persistence).""" + return server_module._run_prompt_submit("rid", "sid", session, text, display_kind="auto_continue") + + +def test_recovery_dispatch_binds_the_original_session_row(recovery_env, tmp_path): + """The interrupted turn's row never landed; the recovery must create it before writing, on the + SAME id and with the session's real source — not leave the store to materialize it later as an + anonymous ``unknown`` session holding only the half-finished assistant text.""" + agent = types.SimpleNamespace( + session_id=ORPHAN_SID, clear_interrupt=lambda: None, + run_conversation=lambda message, **kwargs: {"final_response": "resumed"}) + session = _session(agent, profile_home=str(tmp_path)) + + db = acquire(Path(tmp_path) / "state.db") + try: + assert db.get_session(ORPHAN_SID) is None # the state the recovery resumes from + + _recovery_dispatch(server, session, server._auto_continue_note("the original prompt")) + + row = db.get_session(ORPHAN_SID) + assert row is not None, "the recovery turn must own a durable row for its own session" + assert row["source"] == "desktop" + assert db.list_sessions_rich(source="unknown", limit=50) == [] + finally: + release_or_close(db) + + +def test_creator_repairs_the_accounting_placeholder_source(tmp_path): + """The accounting guard's mint is a placeholder, not an identity: the session's own creator + stamps the real surface on the same id (the upsert used to keep the first writer's value).""" + db = SessionDB(tmp_path / "state.db") + # First writer wins the INSERT because the row creation lost the race with the SQLite lock. + db.update_token_counts(ORPHAN_SID, input_tokens=1200, output_tokens=90, model="grok-4.6") + assert db.get_session(ORPHAN_SID)["source"] == "unknown" + + # The interrupted session's creator (prompt.submit / the recovery dispatch) claims it. + db.create_session(ORPHAN_SID, source="desktop") + + row = db.get_session(ORPHAN_SID) + assert row["source"] == "desktop" + assert row["model"] == "grok-4.6" # nothing else the first writer set is clobbered + assert db.list_sessions_rich(source="unknown", limit=50) == [] + + +def test_startup_sweep_collects_a_legacy_unknown_phantom(tmp_path): + """A phantom an older build already left on disk is collected, not left open forever.""" + db = SessionDB(tmp_path / "state.db") + db.update_token_counts(ORPHAN_SID, input_tokens=10, output_tokens=5, model="claude-sonnet-5") + db.append_message(ORPHAN_SID, role="assistant", content="review table, cut mid-stream") + stale = time.time() - 8 * 3600 + db._conn.execute("UPDATE sessions SET started_at = ? WHERE id = ?", (stale, ORPHAN_SID)) + db._conn.execute("UPDATE messages SET timestamp = ? WHERE session_id = ?", (stale, ORPHAN_SID)) + db._conn.commit() + + assert db.sweep_orphaned_sessions(max_idle_seconds=IDLE_S, sources=_ORPHAN_SWEEP_SOURCES) == [ORPHAN_SID] diff --git a/tui_gateway/prompt_turn.py b/tui_gateway/prompt_turn.py index 4152a1f863..4b8635a8b5 100644 --- a/tui_gateway/prompt_turn.py +++ b/tui_gateway/prompt_turn.py @@ -829,6 +829,18 @@ def _run_prompt_submit( queued_prompt_generation: int | None = None, terminal_callback: Callable[[dict[str, Any]], None] | None = None, turn_author: dict | None = None) -> bool: + # Every dispatch owns a durable row for THIS session before the turn writes anything. The + # prompt.submit handler persists it; the recovery dispatches that call straight in here (the + # crash auto-continue, the queued-prompt drain, the compute-host fallback) did not, so a turn whose + # row never landed — create deferred/failed under the SQLite lock, a record whose session_key had + # not been stamped yet — was materialized by the token-accounting guard instead: an anonymous + # source='unknown' session holding an assistant-first fragment, which the real creator's later + # upsert can never repair (the row insert keeps the first writer's source) and which the orphan + # sweep skips. Binding the row here keeps the recovery inside the ORIGINAL session (#111999). + if _ensure_session_db_row(session) is False: + logger.warning( + "prompt dispatch: session store unavailable for %s — this turn may not persist", + session.get("session_key") or sid) admitted = _admit_prompt_turn(sid, session, text, image_paths, queued_prompt_generation) if admitted is None: return False diff --git a/tui_gateway/session_reaper.py b/tui_gateway/session_reaper.py index 1611ab3cb1..f825b99641 100644 --- a/tui_gateway/session_reaper.py +++ b/tui_gateway/session_reaper.py @@ -290,7 +290,7 @@ def _schedule_session_cap_enforcement() -> None: # conservative. Disable via `dashboard.startup_orphan_sweep: false`. # This is the startup complement every other resource type already has (docker_orphan_reaper, compression # orphans). See #65194. -_ORPHAN_SWEEP_SOURCES = ("tui", "desktop", "subagent") +_ORPHAN_SWEEP_SOURCES = ("tui", "desktop", "subagent", "unknown") _startup_orphan_sweep_ran = False _startup_orphan_sweep_lock = threading.Lock()