fix(sessions): token-accounting guard stamps the agent's real source; trim salvage
When every row create of a turn loses to the SQLite lock, the queued token delta's "ensure the row exists" guard becomes the session's first writer and minted the row as source='unknown'. That placeholder was permanent on the real path even with the upsert repair from #112045: the turn lease (turn_facade_lease.admit_durable_turn) treats an existing row as proof the create already happened and sets _session_db_created, so the creator never returns to repair it. Live probe: a platform="desktop" AIAgent whose create_session raised "database is locked" for the whole first turn ended with a source='unknown' row on base AND on the contributor head; with this change the row is minted 'desktop' by the guard itself. Producer fix: update_token_counts gains an optional source= that the two agent call sites (agent/turn_usage.py, agent/codex_runtime.py) fill from _session_source_for_agent(platform), the same value _ensure_db_session would stamp. record_auxiliary_usage has no surface and keeps the placeholder, which the creator's upsert now repairs. Salvage trims: the contributor's SimpleNamespace dispatch test is replaced by a real-AIAgent invariant test under tests/agent/ (the dispatch hunk in _run_prompt_submit is kept; the INSERT-OR-IGNORE is idempotent under prompt.submit's own persist); narration comments cut to the WHY; docs list 'unknown' among the startup-sweep sources. Refs #111999
This commit is contained in:
@@ -77,7 +77,8 @@ def _queue_token_counts(agent, fail_msg: str, *fail_extra: Any, counts: Callable
|
||||
try:
|
||||
if not agent._session_db_created:
|
||||
agent._ensure_db_session()
|
||||
agent._session_db.queue_token_counts(agent.session_id, **counts())
|
||||
from agent.turn_usage import _agent_session_source
|
||||
agent._session_db.queue_token_counts(agent.session_id, source=_agent_session_source(agent), **counts())
|
||||
except Exception as exc:
|
||||
logger.debug(fail_msg, agent.session_id, *fail_extra, exc)
|
||||
|
||||
|
||||
@@ -22,6 +22,13 @@ from agent.usage_pricing import estimate_usage_cost, normalize_usage
|
||||
logger = logging.getLogger("agent.conversation_loop")
|
||||
|
||||
|
||||
def _agent_session_source(agent: Any) -> str:
|
||||
"""The surface the agent's own row create would stamp (``_ensure_db_session``), so an
|
||||
accounting guard that wins the row-creation race never mints an anonymous session."""
|
||||
from run_agent import _session_source_for_agent # late: run_agent imports this module
|
||||
return _session_source_for_agent(getattr(agent, "platform", None))
|
||||
|
||||
|
||||
@dataclass
|
||||
class ResponseUsageOutcome:
|
||||
"""``compression_attempts`` is the (possibly rearmed-to-zero) budget counter;
|
||||
@@ -248,6 +255,7 @@ def record_response_usage(
|
||||
agent._ensure_db_session()
|
||||
agent._session_db.queue_token_counts(
|
||||
agent.session_id,
|
||||
source=_agent_session_source(agent),
|
||||
input_tokens=canonical_usage.input_tokens,
|
||||
output_tokens=canonical_usage.output_tokens,
|
||||
cache_read_tokens=canonical_usage.cache_read_tokens,
|
||||
|
||||
+1
-1
@@ -1486,7 +1486,7 @@ class SessionDB(
|
||||
_TOKEN_DELTA_COST_FIELDS = ("estimated_cost_usd", "actual_cost_usd")
|
||||
_TOKEN_DELTA_ROUTE_FIELDS = (
|
||||
"model", "cost_status", "cost_source", "pricing_version", "billing_provider", "billing_base_url",
|
||||
"billing_mode",
|
||||
"billing_mode", "source",
|
||||
)
|
||||
|
||||
MAX_TITLE_LENGTH = 100
|
||||
|
||||
+10
-9
@@ -278,18 +278,19 @@ class SessionUsageMixin:
|
||||
actual_cost_usd: Optional[float]=None, cost_status: Optional[str]=None, cost_source: Optional[str]=None,
|
||||
pricing_version: Optional[str]=None, billing_provider: Optional[str]=None, billing_base_url: Optional[str]=None,
|
||||
billing_mode: Optional[str]=None, api_call_count: int=0, absolute: bool=False,
|
||||
source: Optional[str]=None,
|
||||
) -> None:
|
||||
"""Update token counters and backfill model if unset. *absolute*=False increments
|
||||
(per-API-call deltas, CLI path); *absolute*=True sets directly (gateway path,
|
||||
where the cached agent holds cumulative totals)."""
|
||||
where the cached agent holds cumulative totals). ``source`` is the session's real surface
|
||||
for the row-existence guard; callers that don't know it leave the placeholder."""
|
||||
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. 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)
|
||||
# locking, and the UPDATE would silently affect 0 rows. When this guard is the first
|
||||
# writer it must carry the agent's real source: the turn lease treats an existing row as
|
||||
# proof the create already happened, so the creator never returns to repair an anonymous
|
||||
# ``unknown`` placeholder and the session stays a phantom for life (#111999).
|
||||
self._insert_session_row(session_id, source or "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
|
||||
or api_call_count or estimated_cost_usd)
|
||||
@@ -384,8 +385,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 (its placeholder
|
||||
# source is repairable by the session's real creator — see _insert_session_row).
|
||||
# FK to sessions.id: same guard as update_token_counts; the aux path carries no surface, so
|
||||
# the placeholder stays repairable by the creator's upsert (_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))
|
||||
|
||||
|
||||
@@ -0,0 +1,60 @@
|
||||
"""The token-accounting guard never mints an anonymous session row (#111999).
|
||||
|
||||
When every row create of a turn loses to the SQLite lock, the queued token delta's
|
||||
"ensure the row exists" guard becomes the session's first writer. It must stamp the agent's
|
||||
real surface: the turn lease treats an existing row as proof the create happened, so the
|
||||
creator never returns to repair a ``source='unknown'`` placeholder.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import sqlite3
|
||||
from types import SimpleNamespace
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
from hermes_state import SessionDB
|
||||
from run_agent import AIAgent
|
||||
|
||||
SID = "20260913_210721_c89ac8"
|
||||
|
||||
|
||||
def _response():
|
||||
msg = SimpleNamespace(content="partial answer", tool_calls=None, reasoning=None, reasoning_content=None)
|
||||
choice = SimpleNamespace(message=msg, finish_reason="stop")
|
||||
usage = SimpleNamespace(prompt_tokens=120, completion_tokens=30, total_tokens=150,
|
||||
prompt_tokens_details=None, completion_tokens_details=None)
|
||||
return SimpleNamespace(choices=[choice], usage=usage, model="grok-4.6", id="x")
|
||||
|
||||
|
||||
def test_guard_first_writer_carries_the_agent_source(tmp_path):
|
||||
db = SessionDB(tmp_path / "state.db")
|
||||
with (
|
||||
patch("model_tools.get_tool_definitions", return_value=[]),
|
||||
patch("model_tools.check_toolset_requirements", return_value={}),
|
||||
patch("agent.process_bootstrap.OpenAI"),
|
||||
):
|
||||
agent = AIAgent(
|
||||
api_key="test-key-1234567890", base_url="https://openrouter.ai/api/v1", quiet_mode=True,
|
||||
skip_context_files=True, skip_memory=True, session_id=SID, session_db=db, platform="desktop",
|
||||
)
|
||||
agent.client = MagicMock()
|
||||
agent.client.chat.completions.create.return_value = _response()
|
||||
agent._cached_system_prompt = "You are helpful."
|
||||
agent._use_prompt_caching = False
|
||||
agent.compression_enabled = False
|
||||
agent.save_trajectories = False
|
||||
|
||||
# The whole first turn runs under contention: every create_session attempt is refused.
|
||||
def locked(*_a, **_k):
|
||||
raise sqlite3.OperationalError("database is locked")
|
||||
db.create_session = locked
|
||||
try:
|
||||
agent.run_conversation("first prompt")
|
||||
assert db.flush_token_counts()
|
||||
row = db.get_session(SID)
|
||||
assert row is not None, "the queued token delta must materialize the row"
|
||||
assert row["source"] == "desktop"
|
||||
assert db.list_sessions_rich(source="unknown", limit=50) == []
|
||||
finally:
|
||||
agent.close()
|
||||
db.close()
|
||||
@@ -1,35 +1,18 @@
|
||||
"""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). Two contracts keep such a placeholder from becoming a permanent 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.
|
||||
* the session's real creator repairs the placeholder source on the same id (the upsert used to
|
||||
keep whatever the first writer set);
|
||||
* the startup orphan sweep collects an ``unknown`` row 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.
|
||||
@@ -37,90 +20,23 @@ 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)."""
|
||||
stamps the real surface on the same id, and a real source is never downgraded."""
|
||||
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) == []
|
||||
# Control: a later writer with a different real surface never overwrites the creator's.
|
||||
db.create_session(ORPHAN_SID, source="tui")
|
||||
assert db.get_session(ORPHAN_SID)["source"] == "desktop"
|
||||
|
||||
|
||||
def test_startup_sweep_collects_a_legacy_unknown_phantom(tmp_path):
|
||||
|
||||
@@ -829,14 +829,10 @@ 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).
|
||||
# Every dispatch binds the session's own row (session_key, real source) before the turn writes:
|
||||
# the synthesized turns that enter here directly (crash auto-continue, queued-prompt drain,
|
||||
# wake-ups) bypass prompt.submit's persist, and a row-less turn is otherwise materialized by
|
||||
# the token-accounting guard as an anonymous session (#111999).
|
||||
if _ensure_session_db_row(session) is False:
|
||||
logger.warning(
|
||||
"prompt dispatch: session store unavailable for %s — this turn may not persist",
|
||||
|
||||
@@ -2890,4 +2890,4 @@ dashboard:
|
||||
- `ssh_isolated_idle_grace_s` (default `900`) — a Desktop-owned `hermes serve --isolated` backend reached over SSH is detached from the SSH session on purpose, so a laptop that sleeps mid-connection cannot tear it down; each dark-wake reconnect used to leave another backend holding `state.db`. The backend now retires itself once no client WebSocket has been connected for this long and no agent turn is running (a turn keeps it alive; an unreadable turn state keeps it alive too). Set high if you rely on a detached backend finishing long work after the laptop sleeps. Such backends also send a slow WebSocket ping (60 s, 10 min timeout) so a half-open tunnel is noticed.
|
||||
- `ws_orphan_reap_grace_s` — how long a WS-detached session waits before the orphan reaper collects it. Raise alongside the keepalive values if clients reconnect slowly. Periodic session maintenance also completes cleanup for closed sockets and re-arms a missing orphan timer, so a detached chat cannot keep its ownership lease solely because its initial cleanup or timer was lost. Reconnecting cancels that timer; active delegated work and healthy running turns remain protected by the normal orphan-reaper checks. (`HERMES_TUI_WS_ORPHAN_REAP_GRACE_S` remains as an internal override.)
|
||||
- `ws_orphan_activity_stale_s` (default `600`) — how long a detached **running** turn's activity clock (the same clock the `agent.turn_liveness` watchdog samples: API waits, stream tokens, tool heartbeats) must be idle before the orphan reaper interrupts it. A client-absent turn that is still actively producing keeps running to completion detached — closing the laptop, backgrounding the mobile app, or a desktop update no longer cancels healthy long turns; only a genuinely wedged turn is interrupted. Set `0` to interrupt at the grace window regardless of activity (old behavior).
|
||||
- `startup_orphan_sweep` (default `true`) — the WS-orphan reap timer above is in-process, so a gateway restart (update, crash, systemd) before it fires leaves the session row open forever — phantom "active" work in `/resume` and dashboards. On every gateway boot — both the stdio TUI (`entry.main`) and the desktop/dashboard WebSocket sidecar (`handle_ws`) — rows with source `tui` / `desktop` / `subagent` whose start time **and** newest message are both older than the session TTL (`HERMES_TUI_SESSION_TTL_S`, default 6 hours) are closed with `end_reason: startup_orphan_reap`. Messaging-platform sessions (Telegram, Discord, …) are never touched, live in-memory sessions (a client that already resumed) are excluded, and swept sessions remain resumable.
|
||||
- `startup_orphan_sweep` (default `true`) — the WS-orphan reap timer above is in-process, so a gateway restart (update, crash, systemd) before it fires leaves the session row open forever — phantom "active" work in `/resume` and dashboards. On every gateway boot — both the stdio TUI (`entry.main`) and the desktop/dashboard WebSocket sidecar (`handle_ws`) — rows with source `tui` / `desktop` / `subagent` / `unknown` (a row the token-accounting guard had to materialize itself) whose start time **and** newest message are both older than the session TTL (`HERMES_TUI_SESSION_TTL_S`, default 6 hours) are closed with `end_reason: startup_orphan_reap`. Messaging-platform sessions (Telegram, Discord, …) are never touched, live in-memory sessions (a client that already resumed) are excluded, and swept sessions remain resumable.
|
||||
|
||||
Reference in New Issue
Block a user