fix(agent): classify session-persistence failures so lock contention is not misdiagnosed as disk-full

An enterprise deployment hit sustained SQLite write-lock contention on a
shared multi-gigabyte state.db (gateway + CLI processes writing
concurrently). Turns correctly failed closed with
session_persistence_failed, but the only user-facing wording claimed the
disk was full and the gateway rendered a generic failure.

The fast-fail semantics are deliberate and unchanged. This adds a pure
classifier (locked / disk / unknown) applied where the SQLite error is
still visible, threads the cause through the turn-completion explainer,
and stamps a machine-readable failure_reason
(session_persistence_failed:<cause>) plus a guaranteed non-empty error on
the result for downstream surfaces. The cron scheduler's explainer-text
suppression now matches every cause variant so refined wording cannot
leak into scheduled-job deliveries.
This commit is contained in:
Victor Kyriazakos
2026-08-07 00:52:18 +00:00
committed by kshitij
parent 063f6941e5
commit 2a9f5b3476
7 changed files with 288 additions and 11 deletions
+13
View File
@@ -1472,6 +1472,11 @@ def run_conversation(
# A configured SessionDB append failure halts only the affected turn. A
# cached gateway agent must recover on the next message if storage did.
agent._incremental_persistence_failed = False
# Cause of the most recent persistence failure this turn ('locked',
# 'disk', or 'unknown' — see run_agent.classify_persistence_error).
# Reset alongside the failure flag so a lock-contention diagnosis from a
# previous turn can never leak into this turn's user-facing explanation.
agent._last_persistence_error_cause = None
# Main conversation loop counters (pure locals consumed by the loop below).
api_call_count = 0
@@ -6505,6 +6510,10 @@ def run_conversation(
)
except Exception as exc:
_tool_turn_persisted = False
from run_agent import classify_persistence_error
agent._last_persistence_error_cause = (
classify_persistence_error(exc)
)
logger.warning(
"Incremental tool-call persistence failed before execution "
"(session=%s): %s",
@@ -6517,6 +6526,10 @@ def run_conversation(
# run side-effecting tools from state that exists only in
# this process. Breaking also avoids retrying the same
# unpersisted turn until the iteration budget is exhausted.
# The flush may have classified the cause internally; if
# nothing was recorded, the cause is genuinely unknown.
if getattr(agent, "_last_persistence_error_cause", None) is None:
agent._last_persistence_error_cause = "unknown"
_turn_exit_reason = "session_persistence_failed"
final_response = ""
failed = True
+7
View File
@@ -192,9 +192,16 @@ def _flush_session_db_after_tool_progress(
persisted = agent._flush_messages_to_session_db(messages) is not False
if not persisted:
agent._incremental_persistence_failed = True
# The flush caught its own exception and returned False; the
# classified cause (if any) was captured at the catch site. Only
# fall back to 'unknown' when nothing more specific is recorded.
if getattr(agent, "_last_persistence_error_cause", None) is None:
agent._last_persistence_error_cause = "unknown"
return persisted
except Exception as exc:
agent._incremental_persistence_failed = True
from run_agent import classify_persistence_error
agent._last_persistence_error_cause = classify_persistence_error(exc)
logger.warning("Incremental tool-call persistence failed after %s: %s", stage, exc)
return False
+13 -3
View File
@@ -529,7 +529,8 @@ def finalize_turn(
or _is_partial_stream_recovery
):
_explanation = agent._format_turn_completion_explanation(
_turn_exit_reason
_turn_exit_reason,
getattr(agent, "_last_persistence_error_cause", None),
)
if _explanation:
if _is_empty_terminal:
@@ -678,11 +679,20 @@ def finalize_turn(
result["guardrail"] = agent._tool_guardrail_halt_decision.to_metadata()
# Persistence failures already set failed=True + an explanation in
# final_response; also stamp `error` so gateway surfaces status="error"
# (and desktop can toast disk-full) instead of a quiet complete frame.
# (and desktop can toast the cause) instead of a quiet complete frame.
if failed and str(_turn_exit_reason) == "session_persistence_failed":
result["error"] = final_response or (
"session storage could not be written — free disk space and try again"
"session storage could not be written — check the state database "
"health (`hermes doctor`), then send your message again"
)
# Machine-readable cause for the gateway/desktop: exactly
# 'session_persistence_failed:<locked|disk|unknown>'. Never clobber a
# failure_reason another path already stamped on this result.
if "failure_reason" not in result:
_cause = getattr(agent, "_last_persistence_error_cause", None)
result["failure_reason"] = (
"session_persistence_failed:" + (_cause or "unknown")
)
# Surface any post-loop cleanup failures so the caller can distinguish a
# clean turn from one whose trajectory/session/resource teardown raised
# (the response is still returned either way — #8049).
+23 -5
View File
@@ -4250,11 +4250,29 @@ def run_job(
# and soft-failure marking below apply — restoring pre-#34452 silence
# for scheduled jobs without disabling the explainer everywhere.
if final_response.strip() and turn_exit_reason:
try:
_explainer_text = AIAgent._format_turn_completion_explanation(turn_exit_reason)
except Exception:
_explainer_text = ""
if _explainer_text and final_response.strip() == _explainer_text.strip():
# The formatter's wording varies by persistence cause (locked /
# disk / unknown), so render every variant — matching only the
# one-argument render would let cause-refined explainer text slip
# through and be delivered as a cron warning.
_explainer_variants = []
for _cause in (None, "locked", "disk", "unknown"):
try:
_variant = AIAgent._format_turn_completion_explanation(
turn_exit_reason, _cause
)
except TypeError:
# Older single-argument formatter (or a test double).
try:
_variant = AIAgent._format_turn_completion_explanation(
turn_exit_reason
)
except Exception:
_variant = ""
except Exception:
_variant = ""
if _variant:
_explainer_variants.append(_variant.strip())
if final_response.strip() in _explainer_variants:
logger.info(
"Job '%s': abnormal empty turn (%s) — suppressing explainer for cron delivery",
job_id,
+66 -3
View File
@@ -279,6 +279,37 @@ _MAX_TOOL_WORKERS = 8
_DB_PERSISTED_MARKER = "_db_persisted"
def classify_persistence_error(exc_or_str) -> str:
"""Classify a session-persistence failure into a coarse cause bucket.
Fast-failing a turn on a SessionDB write error is deliberate (the
transcript would otherwise be lost on restart), but the *guidance* the
user gets must match the cause: sustained SQLite write-lock contention
("database is locked" on a shared state.db) needs "storage was busy,
send it again", while a full disk or read-only database needs the
disk-space/permissions advice. Returns one of:
* ``"locked"`` — lock/busy contention (another process holds the write
lock); transient, retry-later guidance applies.
* ``"disk"`` — disk full / read-only / permission-shaped failures.
* ``"unknown"`` — anything else (or no visible exception at all).
"""
if exc_or_str is None:
return "unknown"
text = str(exc_or_str).lower()
if "locked" in text or "busy" in text:
return "locked"
if (
"disk" in text
or "readonly" in text
or "read-only" in text
or "no space" in text
or "database or disk is full" in text
):
return "disk"
return "unknown"
# Guard so the OpenRouter metadata pre-warm thread is only spawned once per
# process, not once per AIAgent instantiation. Without this, long-running
# gateway processes leak one OS thread per incoming message and eventually
@@ -2302,6 +2333,11 @@ class AIAgent:
# Force a full re-scan on the next flush: an exception mid-loop
# leaves messages with mixed dispositions.
self._db_flush_scan_prefix = None
# This is the one place the underlying SQLite error is visible
# before it is swallowed into a bare ``False`` — classify it here
# so the turn-end explanation can distinguish lock contention
# ("storage was busy, send it again") from disk-full/read-only.
self._last_persistence_error_cause = classify_persistence_error(e)
logger.warning("Session DB append_message failed: %s", e)
return False
@@ -3579,7 +3615,9 @@ class AIAgent:
return True # safe default: explainer on
@staticmethod
def _format_turn_completion_explanation(turn_exit_reason: str) -> str:
def _format_turn_completion_explanation(
turn_exit_reason: str, persistence_cause: Optional[str] = None
) -> str:
"""Render a user-facing explanation for an abnormal turn ending.
Maps the internal ``turn_exit_reason`` to a short, actionable
@@ -3588,6 +3626,13 @@ class AIAgent:
tool result, or an iteration/budget limit) is never silent from
the UI's perspective — the symptom users report in #34452.
``persistence_cause`` refines the ``session_persistence_failed``
wording (see ``classify_persistence_error``): lock contention gets
"storage was busy, send it again" instead of the disk-space advice,
which was a misdiagnosis for that failure mode. It is optional and
ignored for every other reason, so one-argument callers keep the
exact behavior they had before.
Returns an empty string for reasons that are NOT abnormal (e.g.
a normal ``text_response(...)`` exit), so callers can concatenate
or substitute unconditionally without warning on healthy turns
@@ -3668,12 +3713,30 @@ class AIAgent:
"let it summarize."
)
if reason == "session_persistence_failed":
cause = persistence_cause or "unknown"
if cause == "locked":
return (
prefix
+ "the turn was stopped because session storage was busy "
"(another Hermes process was writing to the state "
"database). Your message was saved — please send it "
"again in a moment."
)
if cause == "disk":
return (
prefix
+ "the turn was stopped because session storage could not "
"be written (the transcript would have been lost on "
"restart). This is often a full disk — free some space "
"(or fix state.db permissions), then send your message "
"again."
)
return (
prefix
+ "the turn was stopped because session storage could not be "
"written (the transcript would have been lost on restart). "
"This is often a full disk — free some space (or fix state.db "
"permissions), then send your message again."
"Check the state database health (`hermes doctor`), then "
"send your message again."
)
# Unknown/diagnostic-only reasons (e.g. "unknown", guardrail_halt
# which already surfaces its own message) — don't second-guess.
@@ -241,6 +241,88 @@ def test_failed_assistant_persist_blocks_ui_projection_and_tool_side_effects():
assert result["failed"] is True
assert result["completed"] is False
assert result["turn_exit_reason"] == "session_persistence_failed"
# No exception was visible (flush returned False), so the cause is
# unknown — but the machine-readable contract fields must still be set.
assert result["failure_reason"] == "session_persistence_failed:unknown"
assert isinstance(result.get("error"), str) and result["error"].strip() != ""
def test_locked_flush_exception_surfaces_locked_cause_in_result_contract():
"""SQLite write-lock contention must surface as a 'locked' cause.
Gateway contract: result['failure_reason'] is exactly
'session_persistence_failed:locked' and result['error'] is a non-empty
string whose wording talks about busy storage, NOT disk space.
"""
import sqlite3
agent = _make_agent()
tool_call = _mock_tool_call(call_id="must-not-run")
agent.client.chat.completions.create.return_value = _mock_response(
content="I'll inspect the repository now.",
finish_reason="tool_calls",
tool_calls=[tool_call],
)
agent._flush_messages_to_session_db = MagicMock(
side_effect=sqlite3.OperationalError("database is locked")
)
agent.interim_assistant_callback = MagicMock()
agent._execute_tool_calls = MagicMock()
with (
patch.object(agent, "_persist_session"),
patch.object(agent, "_save_trajectory"),
patch.object(agent, "_cleanup_task_resources"),
):
result = agent.run_conversation("inspect the repository")
agent.interim_assistant_callback.assert_not_called()
agent._execute_tool_calls.assert_not_called()
assert result["failed"] is True
assert result["turn_exit_reason"] == "session_persistence_failed"
assert result["failure_reason"] == "session_persistence_failed:locked"
assert isinstance(result.get("error"), str) and result["error"].strip() != ""
assert "busy" in result["error"].lower()
assert "disk" not in result["error"].lower()
def test_persistence_cause_resets_between_turns():
"""A locked failure on turn 1 must not leak its cause into turn 2."""
import sqlite3
agent = _make_agent()
tool_call = _mock_tool_call(call_id="must-not-run")
agent.client.chat.completions.create.return_value = _mock_response(
content="I'll inspect the repository now.",
finish_reason="tool_calls",
tool_calls=[tool_call],
)
agent._flush_messages_to_session_db = MagicMock(
side_effect=sqlite3.OperationalError("database is locked")
)
agent._execute_tool_calls = MagicMock()
with (
patch.object(agent, "_persist_session"),
patch.object(agent, "_save_trajectory"),
patch.object(agent, "_cleanup_task_resources"),
):
first = agent.run_conversation("inspect the repository")
assert first["failure_reason"] == "session_persistence_failed:locked"
# Storage recovered but the flush function now reports a bare False
# (no exception): the stale 'locked' cause must not be reused.
agent.client.chat.completions.create.side_effect = None
agent.client.chat.completions.create.return_value = _mock_response(
content="I'll inspect the repository now.",
finish_reason="tool_calls",
tool_calls=[_mock_tool_call(call_id="must-not-run-2")],
)
agent._flush_messages_to_session_db = MagicMock(return_value=False)
second = agent.run_conversation("inspect the repository again")
assert second["turn_exit_reason"] == "session_persistence_failed"
assert second["failure_reason"] == "session_persistence_failed:unknown"
# ---------------------------------------------------------------------------
@@ -93,6 +93,90 @@ def test_explanation_for_max_iterations_reached_prefix_match():
# --------------------------------------------------------------------------
# 1b. Cause-aware session-persistence wording
# --------------------------------------------------------------------------
def test_explanation_persistence_locked_cause_says_busy_not_disk():
"""Write-lock contention must NOT be misdiagnosed as a disk problem."""
out = AIAgent._format_turn_completion_explanation(
"session_persistence_failed", "locked"
)
lower = out.lower()
assert "busy" in lower
assert "saved" in lower
assert "send it again" in lower
assert "disk" not in lower
assert "permission" not in lower
def test_explanation_persistence_disk_cause_keeps_disk_wording():
out = AIAgent._format_turn_completion_explanation(
"session_persistence_failed", "disk"
)
lower = out.lower()
assert "disk" in lower
assert "free some space" in lower or "disk space" in lower
def test_explanation_persistence_unknown_cause_is_neutral():
"""None/'unknown' cause must not claim disk-full — point at diagnostics."""
for cause in (None, "unknown"):
out = AIAgent._format_turn_completion_explanation(
"session_persistence_failed", cause
)
lower = out.lower()
assert out.strip() != ""
assert "disk space" not in lower
assert "full disk" not in lower
assert "hermes doctor" in lower
assert "again" in lower
def test_explanation_persistence_one_arg_backward_compat():
"""Existing one-arg callers must keep working (optional second param)."""
out = AIAgent._format_turn_completion_explanation("session_persistence_failed")
assert out.strip() != ""
assert "session storage" in out.lower()
def test_explanation_cause_ignored_for_other_reasons():
"""The cause parameter must not perturb non-persistence reasons."""
assert (
AIAgent._format_turn_completion_explanation(
"text_response(finish_reason=stop)", "locked"
)
== ""
)
out = AIAgent._format_turn_completion_explanation(
"max_iterations_reached(10/10)", "locked"
)
assert "iteration" in out.lower()
# --------------------------------------------------------------------------
# 1c. classify_persistence_error — the pure cause classifier
# --------------------------------------------------------------------------
def test_classify_persistence_error_categories():
import sqlite3
from run_agent import classify_persistence_error
assert classify_persistence_error(
sqlite3.OperationalError("database is locked")
) == "locked"
assert classify_persistence_error("SQLITE_BUSY: busy") == "locked"
assert classify_persistence_error(
sqlite3.OperationalError("database or disk is full")
) == "disk"
assert classify_persistence_error("attempt to write a readonly database") == "disk"
assert classify_persistence_error("read-only file system") == "disk"
assert classify_persistence_error("no space left on device") == "disk"
assert classify_persistence_error("disk I/O error") == "disk"
assert classify_persistence_error("something else entirely") == "unknown"
assert classify_persistence_error(None) == "unknown"
assert classify_persistence_error("") == "unknown"
# --------------------------------------------------------------------------
# 2. Enable/disable seam
# --------------------------------------------------------------------------