diff --git a/agent/conversation_loop.py b/agent/conversation_loop.py index 0437d3b149..3b9320ff08 100644 --- a/agent/conversation_loop.py +++ b/agent/conversation_loop.py @@ -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 diff --git a/agent/tool_executor.py b/agent/tool_executor.py index 2bc34bcf2f..2463f5948b 100644 --- a/agent/tool_executor.py +++ b/agent/tool_executor.py @@ -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 diff --git a/agent/turn_finalizer.py b/agent/turn_finalizer.py index 79f4397fc6..d51b8b22cf 100644 --- a/agent/turn_finalizer.py +++ b/agent/turn_finalizer.py @@ -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:'. 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). diff --git a/cron/scheduler.py b/cron/scheduler.py index b6b89010ad..1e520bad22 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -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, diff --git a/run_agent.py b/run_agent.py index 7eca7b8ce3..c28d067a9e 100644 --- a/run_agent.py +++ b/run_agent.py @@ -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. diff --git a/tests/run_agent/test_tool_call_incremental_persistence.py b/tests/run_agent/test_tool_call_incremental_persistence.py index b0914b53c5..20a2c2de9d 100644 --- a/tests/run_agent/test_tool_call_incremental_persistence.py +++ b/tests/run_agent/test_tool_call_incremental_persistence.py @@ -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" # --------------------------------------------------------------------------- diff --git a/tests/run_agent/test_turn_completion_explainer.py b/tests/run_agent/test_turn_completion_explainer.py index cb5bca9cf8..30edc08966 100644 --- a/tests/run_agent/test_turn_completion_explainer.py +++ b/tests/run_agent/test_turn_completion_explainer.py @@ -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 # --------------------------------------------------------------------------