diff --git a/agent/session_persistence.py b/agent/session_persistence.py index c206adfe7b..23f2d155c9 100644 --- a/agent/session_persistence.py +++ b/agent/session_persistence.py @@ -129,6 +129,42 @@ def _persist_lock(agent): return nullcontext() if lock is None else lock +def adopt_unanswered_turn(history: List[Dict[str, Any]], query: Any, agent: Any) -> bool: + """Re-stage the transcript's unanswered tail row as THIS turn's user message; True when adopted. + + A dispatcher's re-run of a failed delivery turn resumes the DM its first attempt already persisted + instead of appending it again. Rows loaded from the store are born durable (``_rows_to_conversation``), + so handing the tail row back as ``agent._pending_cli_user_message`` makes ``_stage_turn_user_message`` + reuse it as this turn's user dict and the flush writes no second row. What differs per lane is only HOW + the dispatcher knows the DM is unanswered: + + * ``hermes_cli.quiet_single_query.adopt_unanswered_turn`` — the delivery lanes' re-run is a fresh CLI + process, told so through ``tools.bot_relay.RESUME_UNANSWERED_TURN_ENV``. + * ``gateway.platforms.api_server`` — the peer-DM lane re-runs the turn in-process and calls this + directly on the agent it just built for the re-run (#115325). + + The DM is not always the literal tail: a turn that died mid-way persisted its tool scaffolding — assistant + ``tool_calls`` rows and their ``tool`` results — behind the DM before the failure text was built, and the + dispatcher retries that too. The DM is still unanswered while nothing after it is a plain assistant reply, + so it is adopted and the failed attempt's scaffolding leaves the in-memory transcript: the re-run starts + the turn over from the DM (the rows stay in the DB as the record of the failed attempt; the re-run's + answer lands after them as a valid continuation). Anything else declines — no user row at the tail, or a + different text there — so a person's deliberate re-send of the same text is never swallowed. + """ + idx = next((i for i in range(len(history) - 1, -1, -1) + if isinstance(history[i], dict) and history[i].get("role") == "user"), None) + if idx is None or history[idx].get("content") != query: + return False + if not all(isinstance(row, dict) and (row.get("role") == "tool" or (row.get("role") == "assistant" and row.get("tool_calls"))) + for row in history[idx + 1:]): + return False + tail = history[idx] + del history[idx:] + tail[_DB_PERSISTED_MARKER] = True + agent._pending_cli_user_message = tail + return True + + # --- flush phases (module-level so the flush also works bound onto duck-typed agents) --- def _db_flush_seed_ids(agent) -> set: diff --git a/gateway/platforms/api_server.py b/gateway/platforms/api_server.py index 68221498f3..0a0975e8e5 100644 --- a/gateway/platforms/api_server.py +++ b/gateway/platforms/api_server.py @@ -3189,7 +3189,9 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): @_admit_api_agent_request async def _handle_session_chat(self, request: "web.Request") -> "web.Response": - """POST /api/sessions/{session_id}/chat — one synchronous agent turn.""" + """POST /api/sessions/{session_id}/chat — one synchronous agent turn (plus the delivery lanes' + one bounded re-run of a transient failure; ``hermes peer dm`` is the client).""" + from tools.bot_failure_reasons import RETRY_NONE, result_retry_action # This turn runs through _run_agent, so it already COUNTS toward the cap (#7483). # Spending the budget without checking it refused every other caller while never # refusing this route — and a fleet's cross-machine DMs all arrive here. @@ -3203,6 +3205,18 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): session_id = ctx["session_id"] history = await self._conversation_history_for_session(session_id) result, usage = await self._run_agent(conversation_history=history, **ctx["run_kwargs"]) + # One policy-gated re-run of a transiently failed turn — the peer-DM transport's half of the + # retry the local (``tools.bot_mode_dm``) and relayed (``tui_gateway.methods_bot_relay``) + # delivery lanes already apply (#93091 item 5, #115325). Same policy, same gate: transient + # classes (429 / 5xx) re-run the SAME session once, a context overflow lets the re-run's + # pre-API compaction shrink the transcript first, and auth/quota/config/model never re-run. The + # store is read again first: the failed attempt's turn-start persist left the DM as the + # transcript's unanswered tail row, and the re-run resumes that row instead of appending a + # second copy of it. A turn that fails again reaches the peer client exactly as before. + if result_retry_action(result) != RETRY_NONE: + history = await self._conversation_history_for_session(session_id) + result, usage = await self._run_agent( + conversation_history=history, resume_unanswered_turn=True, **ctx["run_kwargs"]) is_dict = isinstance(result, dict) effective_session_id = result.get("session_id") if is_dict else session_id final_response = _resolve_media_to_data_urls( @@ -3781,7 +3795,8 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): requested_runtime: Optional[Dict[str, Any]] = None, route_source: str = "global", confirmed_runtime_lock: bool = False, bind_declared_conversation: bool = False, session_history_delivery: str = "", turn_author: Optional[Dict[str, Any]] = None, - relay_metadata: Optional[Dict[str, Any]] = None, notification_category: str = "result") -> tuple: + relay_metadata: Optional[Dict[str, Any]] = None, notification_category: str = "result", + resume_unanswered_turn: bool = False) -> tuple: """Create an agent and run one turn in a thread executor -> ``(result, usage)``. ``agent_ref[0]`` receives the agent so SSE writers can interrupt it; ``active_run_id`` registers it in ``_active_run_agents``. Under a confirmed model lock the actual @@ -3789,7 +3804,11 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): ``session_history_delivery`` declares #98619 session-id provenance and default-denies: only audited producers whose client can address the id again pass "1" (see ``_bind_api_server_session``). - ``turn_author`` only labels the turn for memory attribution. It grants nothing.""" + ``turn_author`` only labels the turn for memory attribution. It grants nothing. + ``resume_unanswered_turn`` marks a policy-gated re-run of a turn whose user row the failed attempt + already persisted: the transcript's unanswered tail row is adopted from ``conversation_history`` + as THIS turn's user message instead of being appended a second time + (``agent.session_persistence.adopt_unanswered_turn``; #115325).""" loop = asyncio.get_running_loop() # ContextVars do not follow run_in_executor threads: capture here, re-enter in _run(). request_profile = _api_request_profile.get() @@ -3821,6 +3840,13 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): session_model=session_model, confirmed_runtime_lock=confirmed_runtime_lock) if agent_ref is not None: agent_ref[0] = agent + if resume_unanswered_turn: + # A dispatcher's re-run of a failed delivery turn: the DM's own row is already + # in the store (the failed attempt persisted it at turn start), so continue THAT + # row instead of appending a second copy of the same text (#115325). + from agent.session_persistence import adopt_unanswered_turn + + adopt_unanswered_turn(conversation_history, user_message, agent) if active_run_id: self._active_run_agents[active_run_id] = agent effective_task_id = session_id or str(uuid.uuid4()) diff --git a/hermes_cli/quiet_single_query.py b/hermes_cli/quiet_single_query.py index 8088a1d7d3..ad6e575021 100644 --- a/hermes_cli/quiet_single_query.py +++ b/hermes_cli/quiet_single_query.py @@ -175,33 +175,14 @@ def adopt_unanswered_turn(cli: Any, query: Any, environ: MutableMapping[str, str fresh process cannot know that by itself (``_DB_PERSISTED_MARKER`` is in-process only), and inferring it from an identical tail alone would swallow a person's deliberate re-send — so the dispatcher must say so with ``tools.bot_relay.RESUME_UNANSWERED_TURN_ENV``, consumed (popped) here - before the turn so tool subprocesses never inherit it. The row is re-staged as the pending CLI - dict already stamped durable: ``_stage_turn_user_message`` reuses it as this turn's user message and - the flush writes no second row. - - The DM is not always the literal tail: a turn that died mid-way (HTTP 503 on the call after a tool - round) persisted its tool scaffolding — assistant ``tool_calls`` rows and their ``tool`` results — - behind the DM before ``agent.turn_recovery`` built the failure text, and the dispatcher retries that - too. The DM is still unanswered while nothing after it is a plain assistant reply, so it is adopted - and the failed attempt's scaffolding leaves the in-memory transcript: the re-run starts the turn - over from the DM (the rows stay in the DB as the record of the failed attempt; the re-run's answer - lands after them as a valid continuation).""" + before the turn so tool subprocesses never inherit it. Which row counts as the unanswered DM, and + how it is re-staged as ``_pending_cli_user_message``, is shared with the in-process peer-DM lane + (``agent.session_persistence.adopt_unanswered_turn``, #115325). + """ from tools.bot_relay import RESUME_UNANSWERED_TURN_ENV if environ.pop(RESUME_UNANSWERED_TURN_ENV, None) != "1": return False - history = getattr(cli, "conversation_history", None) or [] - idx = next((i for i in range(len(history) - 1, -1, -1) - if isinstance(history[i], dict) and history[i].get("role") == "user"), None) - if idx is None or history[idx].get("content") != query: - return False - if not all(isinstance(row, dict) and (row.get("role") == "tool" or (row.get("role") == "assistant" and row.get("tool_calls"))) - for row in history[idx + 1:]): - return False - from agent.context_compressor import _DB_PERSISTED_MARKER + from agent.session_persistence import adopt_unanswered_turn as _adopt_tail - tail = history[idx] - del history[idx:] - tail[_DB_PERSISTED_MARKER] = True - cli.agent._pending_cli_user_message = tail - return True + return _adopt_tail(getattr(cli, "conversation_history", None) or [], query, cli.agent) diff --git a/tests/gateway/test_peer_dm_transient_retry.py b/tests/gateway/test_peer_dm_transient_retry.py new file mode 100644 index 0000000000..e5b781c81c --- /dev/null +++ b/tests/gateway/test_peer_dm_transient_retry.py @@ -0,0 +1,121 @@ +"""#115325: the peer-DM transport (``POST /api/sessions/{id}/chat``) retries a transiently failed +turn once, resuming the DM row the failed attempt already persisted. + +``hermes peer dm`` is the third Bot-DM transport. The local (``tools.bot_mode_dm``) and relayed +(``tui_gateway.methods_bot_relay``) lanes both re-run a transiently failed turn once and resume the +row the failed attempt left as the transcript's unanswered tail (``tools/bot_failure_reasons`` +``retry_action`` / ``RESUME_UNANSWERED_TURN``); the peer lane ran the turn once and handed whatever +came back — the provider's 429 paragraph — to the sender as the reply. + +Only the model turn is faked (``_create_agent`` hands out recording agents); the route, the store and +``_run_agent`` are the real ones, so both halves of the contract are asserted on what the retry was +actually handed: the persisted row adopted as this turn's user message, and no second attempt for a +failure the policy never retries. +""" + +from unittest.mock import MagicMock, patch + +import pytest +from aiohttp import web +from aiohttp.test_utils import TestClient, TestServer + +from gateway.config import PlatformConfig +from gateway.platforms.api_server import APIServerAdapter +from hermes_state import SessionDB + +SESSION_ID = "peer_dm_retry" +DM = "disk status?" + +# Transient: the policy's own class (429 / rate limit). +RATE_LIMIT = { + "final_response": "Rate limited by the provider. API call failed after 3 retries: 429 Too Many Requests", + "failed": True, "completed": False, "error": "429 Too Many Requests", + "failure_reason": "rate_limit", "messages": [], "api_calls": 3, +} +# Permanent: a wall no re-run can fix. +AUTH_WALL = { + "final_response": "Provider authentication failed: invalid api key", + "failed": True, "completed": False, "error": "Error code: 401 - invalid_api_key", + "failure_reason": "auth", "messages": [], "api_calls": 1, +} + + +def _ok(text: str) -> dict: + return {"final_response": text, "failed": False, "completed": True, + "messages": [], "api_calls": 1} + + +def _app(adapter: APIServerAdapter) -> web.Application: + app = web.Application() + app.router.add_post("/api/sessions/{session_id}/chat", adapter._handle_session_chat) + return app + + +def _adapter(tmp_path) -> tuple[APIServerAdapter, str]: + """A real adapter over a real store, holding the DM row a failed attempt's turn-start persist + would have left as the transcript's unanswered tail.""" + db = SessionDB(tmp_path / "state.db") + sid = db.create_session(SESSION_ID, "api_server") + db.append_message(sid, "user", content=DM) + adapter = APIServerAdapter(PlatformConfig(enabled=True)) + adapter._session_db = db + return adapter, sid + + +def _fake_agent(seen: list, outcome: dict) -> MagicMock: + """Records what each turn was handed; ``_pending_cli_user_message`` is the adopted row.""" + agent = MagicMock() + agent.session_id = SESSION_ID + agent.session_prompt_tokens = 0 + agent.session_completion_tokens = 0 + agent.session_total_tokens = 0 + + def _run(user_message=None, conversation_history=None, task_id=None, **_kwargs): + seen.append({ + "message": user_message, + "history": list(conversation_history or []), + "resumed": getattr(agent, "_pending_cli_user_message", None), + }) + return dict(outcome) + + agent.run_conversation.side_effect = _run + return agent + + +@pytest.mark.asyncio +async def test_peer_dm_retries_a_transient_turn_once_and_still_reports_a_permanent_one(tmp_path): + adapter, sid = _adapter(tmp_path) + app = _app(adapter) + + # ① Transient failure -> re-run once, resuming the persisted DM row, and the recovered reply is + # what the peer client receives. + seen: list = [] + agents = [_fake_agent(seen, RATE_LIMIT), _fake_agent(seen, _ok("2+2 is 4"))] + with patch.object(adapter, "_create_agent", side_effect=lambda **_kw: agents.pop(0)): + async with TestClient(TestServer(app)) as cli: + resp = await cli.post(f"/api/sessions/{sid}/chat", json={"message": DM}) + body = await resp.json() + + assert resp.status == 200 + assert len(seen) == 2, "a transiently failed peer DM must be re-run exactly once" + assert body["message"]["content"] == "2+2 is 4" + retry = seen[1] + assert retry["message"] == DM, "the re-run replays the same DM" + resumed = retry["resumed"] + assert isinstance(resumed, dict) and resumed.get("content") == DM, ( + "the re-run must resume the row the failed attempt persisted as the pending user message") + assert resumed.get("_db_persisted") is True, "the resumed row is the durable one, not a copy" + assert all(row.get("content") != DM for row in retry["history"]), ( + "the resumed row leaves the history, so the turn appends no second copy of the DM") + + # ② Permanent failure -> reported exactly as before: one attempt, its error copy as the reply. + seen2: list = [] + walls = [_fake_agent(seen2, AUTH_WALL), _fake_agent(seen2, AUTH_WALL)] + with patch.object(adapter, "_create_agent", side_effect=lambda **_kw: walls.pop(0)): + async with TestClient(TestServer(app)) as cli: + resp2 = await cli.post(f"/api/sessions/{sid}/chat", json={"message": DM}) + body2 = await resp2.json() + + assert resp2.status == 200 + assert len(seen2) == 1, "auth is never auto-retried: it cannot be fixed by a re-run" + assert body2["message"]["content"] == AUTH_WALL["final_response"] diff --git a/tools/bot_failure_reasons.py b/tools/bot_failure_reasons.py index aa462d928e..5e30b36829 100644 --- a/tools/bot_failure_reasons.py +++ b/tools/bot_failure_reasons.py @@ -11,6 +11,7 @@ provider 401 bodies (e.g. Anthropic) say "invalid, blocked or out of funds". from __future__ import annotations import re +from typing import Any # platform-side RUNTIME_OFFLINE = "runtime_offline" @@ -72,7 +73,10 @@ _RULES: tuple[tuple[re.Pattern[str], str], ...] = tuple( (rf"authentication_error|invalid api key|{_STATUS}(?:401|403)\b", PROVIDER_AUTH_OR_ACCESS), (rf"{_STATUS}402\b|out of funds|quota|balance", PROVIDER_QUOTA_LIMIT), (rf"{_STATUS}429\b|rate.?limit", PROVIDER_RATE_LIMIT), - (rf"{_STATUS}5\d{{2}}\b|server error|overloaded", PROVIDER_SERVER_ERROR), + # ``server[ _]?error`` / ``overloaded_error`` are the providers' own JSON spellings of the + # agent's typed ``server_error`` / ``overloaded`` verdicts (FailoverReason), which the + # in-process lanes classify from ``failure_reason`` rather than a status number. + (rf"{_STATUS}5\d{{2}}\b|server[ _]?error|overloaded", PROVIDER_SERVER_ERROR), (r"context length|context_overflow|maximum context", CONTEXT_OVERFLOW), (r"no llm provider configured|missing config|no access token", MISSING_CONFIG), (r"model .*(not found|does not exist)|model_not_found", MODEL_UNAVAILABLE), @@ -84,10 +88,29 @@ def turn_failure_text(stdout: str | None, stderr: str | None) -> str: """The error text of a failed ``hermes … -Q`` delivery turn: both streams, in the order the CLI writes them. The provider prose is the turn's final_response and lands on STDOUT; stderr carries session bookkeeping (``session_id: …``) on every run, so ``stderr or stdout`` only ever saw the - banner and every transient failure classified as ``unknown``.""" + banner and every transient failure classified as ``unknown``. + + The in-process lanes hand over the same two pieces of prose as an agent result dict + (``error`` + ``failure_reason``); ``result_retry_action`` joins them here too, so one classifier + sees every lane's failure text.""" return "\n".join(text.strip() for text in (stdout, stderr) if text and text.strip()) +def result_retry_action(result: Any) -> str: + """The retry action for a finished IN-PROCESS turn (an ``_run_agent`` result dict). + + The transport-agnostic twin of the child-process lanes' ``retry_action(classify_agent_error( + turn_failure_text(stdout, stderr)))``: same policy, the failure text just arrives as result fields + instead of two streams — ``error`` carries the raw provider summary (status codes included) and + ``failure_reason`` the turn loop's own typed verdict, which is what names an overflow the copy only + describes in prose. ``RETRY_NONE`` for anything that did not fail (and for a non-dict result), so a + successful turn whose text happens to mention 429 is never re-run.""" + if not isinstance(result, dict) or not result.get("failed"): + return RETRY_NONE + return retry_action(classify_agent_error( + turn_failure_text(result.get("error"), result.get("failure_reason")))) + + def classify_agent_error(text: str) -> str: """Map raw agent/provider error text to a closed reason code (``unknown`` when unmatched/empty).""" raw = str(text or "")