From 8f0322da5b82029f3bc4d16fbaa2c986299abfc6 Mon Sep 17 00:00:00 2001 From: teknium1 <127238744+teknium1@users.noreply.github.com> Date: Fri, 18 Sep 2026 00:55:30 -0700 Subject: [PATCH] fix(agent): close a failed turn's durable user tail so the next prompt is not merged into it MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Terminal-failure paths (HTTP-200 content-policy refusal, ``_Trunc.end_turn``, retry exhaustion, interrupt before any assistant text) persist the accepted user row and return before ``finalize_turn``, so ``user`` stays the durable conversation tail. The next prompt appends a second user row, ``repair_message_sequence`` merges the pair, and the provider is asked to act on the failed request again. The gateway compensates with ``_hmwa_close_failed_turn`` (#108033); standalone ACP, the CLI and the TUI/Desktop hand ``result["messages"]`` straight back as history and had no closer. Close it once, at ``agent/conversation_loop.py::run_conversation`` — the seam every envelope leaves through — with a Hermes-authored assistant boundary (``agent/turn_failure_copy.py::FAILED_TURN_NOTICE`` / ``PARTIAL_FAILED_TURN_NOTICE``, which the gateway now aliases instead of keeping its own copy). Idempotence is keyed on ``SessionDB.latest_conversation_role`` (durable state, not content), so a redelivery or a tail another writer already closed is a no-op and the gateway's closer no-ops in turn. The context-pressure classes (``compression_exhausted``, ``compression_deferred``, ``failure_reason == "context_overflow"``) are excluded: appending to an oversized session is the #1630 growth loop; their repair is rotation. Adjacent defect from the same report: ``acp_adapter/server.py::_finish_turn`` called ``final_response.startswith`` on ``None`` for an interrupted turn — the same one-line fix PR #64471 by @israellot filed first (its wider prompt()-restructure is superseded by the current ``_finish_turn`` shape). Slimmer redo of #114168 by @kendrickkester (same seam and invariants; the +1023-line PR carried a new copy module, an accepted-turn re-anchoring scan and an 859-line suite). Two invariant tests: the real ACP path (loopback provider, refusal then a new prompt) and the durable-tail idempotence / overflow exclusion. Co-authored-by: Kendrick Kester Co-authored-by: Israel Lot --- acp_adapter/server.py | 2 +- agent/conversation_loop.py | 47 +++- agent/turn_failure_copy.py | 26 +++ gateway/run_turn.py | 12 +- tests/acp_adapter/test_failed_turn_closure.py | 211 ++++++++++++++++++ 5 files changed, 288 insertions(+), 10 deletions(-) create mode 100644 tests/acp_adapter/test_failed_turn_closure.py diff --git a/acp_adapter/server.py b/acp_adapter/server.py index 406e44e6a2..b8dc358b73 100644 --- a/acp_adapter/server.py +++ b/acp_adapter/server.py @@ -913,7 +913,7 @@ class HermesACPAgent(SlashCommandsMixin, acp.Agent): except Exception: logger.debug("Could not emit ACP provenance update after rotation for %s", session_id, exc_info=True) - final_response = result.get("final_response", "") + final_response = result.get("final_response") or "" # None on an interrupted turn cancelled = bool(state.cancel_event and state.cancel_event.is_set()) # The local "waiting for model" interrupt status is metadata, not prose; stop_reason carries it. from agent.conversation_loop import INTERRUPT_WAITING_FOR_MODEL_PREFIX diff --git a/agent/conversation_loop.py b/agent/conversation_loop.py index 6065da4ac0..4ce7008935 100644 --- a/agent/conversation_loop.py +++ b/agent/conversation_loop.py @@ -40,7 +40,7 @@ from agent.turn_retry_state import TurnRetryState from agent.turn_api_call import handle_api_interrupt, nous_rate_limit_guard, perform_api_call from agent.turn_api_error import handle_api_error from agent.turn_api_request import build_api_request -from agent.turn_failure_copy import site_copy +from agent.turn_failure_copy import failed_turn_notice, site_copy from agent.turn_final_response import finish_text_response from agent.turn_finalizer import finalize_turn from agent.turn_iteration_prep import ( @@ -1615,7 +1615,50 @@ def run_conversation( moa_config=moa_config, turn_author=turn_author, ) - return export_current_turn_boundary(agent, result, user_message) + result = export_current_turn_boundary(agent, result, user_message) + _close_durable_failed_turn(agent, result) + return result + + +def _close_durable_failed_turn(agent, result: Any) -> None: + """Append a Hermes-authored assistant boundary when a failed turn left ``user`` as the + durable conversation tail (in place, on ``result["messages"]`` and in SessionDB). + + The terminal-failure paths (content-policy refusal, ``_Trunc.end_turn``, retry exhaustion, + interrupt before any assistant text) persist the accepted user row and return without + reaching ``finalize_turn``; the next prompt then appends a second user row and + ``repair_message_sequence`` merges the failed request into the new one. The gateway + compensates with ``_hmwa_close_failed_turn``; CLI, TUI/Desktop and ACP hosts hand + ``result["messages"]`` straight back as history, so the seam is here. + + Excluded: the context-pressure classes (``compression_exhausted``, ``compression_deferred``, + ``failure_reason == "context_overflow"``) — appending to an already-oversized session is the + #1630 growth loop; their repair is rotation or a retry. Idempotence is keyed on the DURABLE + tail (``SessionDB.latest_conversation_role``), so a redelivery or a tail already closed by + another writer is a no-op, and the gateway's own closer then no-ops in turn. + """ + try: + if not isinstance(result, dict) or result.get("completed") is True: + return + if ( + result.get("compression_exhausted") or result.get("compression_deferred") + or result.get("failure_reason") == "context_overflow" + ): + return + messages = result.get("messages") + db, session_id = getattr(agent, "_session_db", None), getattr(agent, "session_id", None) + if not isinstance(messages, list) or not messages or db is None or not session_id: + return + if getattr(agent, "_persist_disabled", False) or db.latest_conversation_role(session_id) != "user": + return + # Scope the "did a tool run" scan to this turn when its boundary is proven; otherwise + # hedge over the whole list rather than under-report a possible side effect. + start = result.get("current_turn_user_idx") + turn_messages = messages[start:] if isinstance(start, int) and 0 <= start < len(messages) else messages + append_message(messages, {"role": "assistant", "content": failed_turn_notice(turn_messages)}) + agent._flush_messages_to_session_db(messages) + except Exception: + logger.debug("failed-turn boundary not written", exc_info=True) __all__ = ["run_conversation"] diff --git a/agent/turn_failure_copy.py b/agent/turn_failure_copy.py index ea3e880328..ffd7f06bea 100644 --- a/agent/turn_failure_copy.py +++ b/agent/turn_failure_copy.py @@ -28,6 +28,32 @@ def stamp_failure(result: Dict[str, Any], reason: str, retryable: bool) -> Dict[ return result +# ---- failed-turn transcript boundary ---------------------------------------------------------- +# The Hermes-authored assistant row that closes a durable turn which ended without one. A +# transcript boundary, NOT the model's answer: no provider/model error or refusal detail is +# ever interpolated (that rides ``final_response``). Owned here so the core closer +# (``agent/conversation_loop.py::run_conversation``) and the gateway's own writer +# (``gateway/run_turn.py::_hmwa_close_failed_turn``) say the same thing. + +FAILED_TURN_NOTICE = ( + "Your request was not processed. Send it again if you still want me to carry it out." +) +PARTIAL_FAILED_TURN_NOTICE = ( + "This turn did not complete. Some actions may already have run; verify their effects " + "before resending." +) + + +def failed_turn_notice(turn_messages: Any) -> str: + """Boundary copy for a failed turn: never claim "not processed" when a tool may have run.""" + for row in turn_messages or (): + if isinstance(row, dict) and ( + row.get("role") == "tool" or (row.get("role") == "assistant" and row.get("tool_calls")) + ): + return PARTIAL_FAILED_TURN_NOTICE + return FAILED_TURN_NOTICE + + def provider_label_for(provider: Any) -> str: """Human-friendly provider name for chat copy (``"OpenRouter"``, ``"Nous Portal"``…).""" from hermes_cli.models import provider_label diff --git a/gateway/run_turn.py b/gateway/run_turn.py index f0afcd6bf2..819376a603 100644 --- a/gateway/run_turn.py +++ b/gateway/run_turn.py @@ -17,6 +17,7 @@ import threading import time from agent.i18n import t from agent.session_activity import format_iteration_progress +from agent.turn_failure_copy import FAILED_TURN_NOTICE, PARTIAL_FAILED_TURN_NOTICE from contextlib import nullcontext, suppress from contextvars import copy_context from gateway.config import Platform @@ -1606,13 +1607,10 @@ class GatewayTurnMixin: except Exception as e: logger.debug("Watch queue drain error: %s", e) - _FAILED_TURN_NOTICE = ( - "Your request was not processed. Send it again if you still want me to carry it out." - ) - _PARTIAL_FAILED_TURN_NOTICE = ( - "This turn did not complete. Some actions may already have run; verify their effects " - "before resending." - ) + # One owner for the boundary copy: the core closer (agent/conversation_loop.py) writes the + # same row on the paths that never reach this layer; alias rather than keep a second copy. + _FAILED_TURN_NOTICE = FAILED_TURN_NOTICE + _PARTIAL_FAILED_TURN_NOTICE = PARTIAL_FAILED_TURN_NOTICE def _hmwa_add_failed_turn_notice(self, response, notice): """Make failed-turn delivery explicit without replacing the provider-specific guidance.""" diff --git a/tests/acp_adapter/test_failed_turn_closure.py b/tests/acp_adapter/test_failed_turn_closure.py new file mode 100644 index 0000000000..010ea27d6a --- /dev/null +++ b/tests/acp_adapter/test_failed_turn_closure.py @@ -0,0 +1,211 @@ +"""A failed turn must not leave its accepted user row as the durable conversation tail. + +The gateway closes failed turns itself (``gateway/run_turn.py::_hmwa_close_failed_turn``); +standalone ACP hands ``result["messages"]`` straight back as history, so the terminal-failure +paths that persist the user row and return before ``finalize_turn`` used to leave ``user`` as +the tail. The next prompt then appended a second user row and ``repair_message_sequence`` +merged the failed request into the new one. Closed at the core seam +(``agent/conversation_loop.py::_close_durable_failed_turn``). + +Driven through the real ``HermesACPAgent.prompt()`` with a real ``AIAgent`` and ``SessionDB`` +against a loopback OpenAI-compatible endpoint; assertions read SQLite rows and the request +body the provider actually received. +""" + +from __future__ import annotations + +import asyncio +import json +import sqlite3 +import threading +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer + +import pytest + +_MODEL = "fixture-model" + + +class _LoopbackProvider: + """Scripted OpenAI-compatible endpoint: one ``{"finish_reason", "content"}`` per request.""" + + def __init__(self) -> None: + self.script: list[dict] = [] + self.requests: list[dict] = [] + provider = self + + class _Handler(BaseHTTPRequestHandler): + def log_message(self, *args): + pass + + def do_POST(self): # noqa: N802 - BaseHTTPRequestHandler API + body = json.loads(self.rfile.read(int(self.headers["Content-Length"]))) + provider.requests.append(body) + spec = provider.script.pop(0) + delta = {"role": "assistant", "content": spec["content"]} + usage = {"prompt_tokens": 10, "completion_tokens": 4, "total_tokens": 14} + if body.get("stream"): + chunks = [ + {"id": "c", "object": "chat.completion.chunk", "created": 1, "model": _MODEL, + "choices": [{"index": 0, "delta": delta, "finish_reason": None}]}, + {"id": "c", "object": "chat.completion.chunk", "created": 1, "model": _MODEL, + "choices": [{"index": 0, "delta": {}, "finish_reason": spec["finish_reason"]}], + "usage": usage}, + ] + mime, payload = "text/event-stream", ( + "".join("data: " + json.dumps(c) + "\n\n" for c in chunks) + "data: [DONE]\n\n" + ).encode() + else: + mime, payload = "application/json", json.dumps({ + "id": "c", "object": "chat.completion", "created": 1, "model": _MODEL, + "choices": [{"index": 0, "message": delta, "finish_reason": spec["finish_reason"]}], + "usage": usage, + }).encode() + self.send_response(200) + self.send_header("Content-Type", mime) + self.send_header("Content-Length", str(len(payload))) + self.end_headers() + self.wfile.write(payload) + + self._server = ThreadingHTTPServer(("127.0.0.1", 0), _Handler) + threading.Thread(target=self._server.serve_forever, daemon=True).start() + self.base_url = f"http://127.0.0.1:{self._server.server_port}/v1" + + def shutdown(self) -> None: + self._server.shutdown() + + +class _RecordingConn: + def __init__(self) -> None: + self.updates: list = [] + + async def session_update(self, session_id, update): + self.updates.append(update) + + async def request_permission(self, *args, **kwargs): + raise AssertionError("no tool approval expected") + + +@pytest.fixture +def acp(tmp_path, monkeypatch): + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + monkeypatch.setenv("HERMES_DISABLE_PLUGINS", "1") + monkeypatch.setenv("NO_PROXY", "127.0.0.1,localhost") + provider = _LoopbackProvider() + + import acp_adapter.session as acp_session + import hermes_cli.config as cli_config + import hermes_cli.mcp_startup as mcp_startup + import hermes_cli.runtime_provider as runtime_provider + + monkeypatch.setattr(cli_config, "load_config", lambda *a, **k: { + "model": {"provider": "openai-compat", "default": _MODEL, "context_length": 131072}, + "agent": {"max_iterations": 2}, "compression": {"enabled": False}, + }) + monkeypatch.setattr(runtime_provider, "resolve_runtime_provider", lambda *a, **k: { + "provider": "openai-compat", "api_mode": "chat_completions", + "base_url": provider.base_url, "api_key": "fixture-only", + }) + monkeypatch.setattr(mcp_startup, "ensure_mcp_discovery_before_agent_build", lambda **k: None) + monkeypatch.setattr(acp_session, "_expand_acp_enabled_toolsets", lambda *a, **k: []) + + from acp_adapter.server import HermesACPAgent, TextContentBlock + from acp_adapter.session import SessionManager + from hermes_state import SessionDB + + db_path = tmp_path / "state.db" + db = SessionDB(db_path) + server = HermesACPAgent(session_manager=SessionManager(db=db)) + conn = _RecordingConn() + server.on_connect(conn) + sid = server.session_manager.create_session(cwd=str(tmp_path)).session_id + + def prompt(text): + return asyncio.run(server.prompt(prompt=[TextContentBlock(type="text", text=text)], session_id=sid)) + + def conversation_rows(): + with sqlite3.connect(db_path) as c: + return [tuple(r) for r in c.execute( + "SELECT role, content FROM messages WHERE session_id = ? AND active = 1 " + "AND role NOT IN ('session_meta', 'system') ORDER BY id", (sid,))] + + yield provider, prompt, conversation_rows, db, sid, conn + provider.shutdown() + db.close() + + +_REFUSED = "Explain how to pick the lock on my neighbour's front door" +_REFUSAL_DETAIL = "I can't help with breaking into someone else's property." +_NEW_REQUEST = "What is the capital of France?" + + +def test_acp_refusal_closes_the_turn_and_is_not_replayed_into_the_next_prompt(acp): + """Turn 1: HTTP-200 ``content_filter`` refusal. Turn 2: unrelated request. + + Invariant: the durable tail after a failed turn is a Hermes-authored assistant row (never + provider text), and the next prompt reaches the provider as its own user row. + """ + from agent.turn_failure_copy import FAILED_TURN_NOTICE + + provider, prompt, conversation_rows, db, sid, conn = acp + + provider.script = [{"finish_reason": "content_filter", "content": _REFUSAL_DETAIL}] + prompt(_REFUSED) + + rows = conversation_rows() + assert [r[0] for r in rows] == ["user", "assistant"], rows + assert rows[1][1] == FAILED_TURN_NOTICE # no tool ran: no hedging, no provider detail + assert db.latest_conversation_role(sid) == "assistant" + # The refusal reaches the ACP client, but never becomes canonical assistant history. + texts = [getattr(getattr(u, "content", None), "text", None) or getattr(u, "text", None) for u in conn.updates] + assert any(_REFUSAL_DETAIL in (t or "") for t in texts) + + provider.script = [{"finish_reason": "stop", "content": "Paris."}] + prompt(_NEW_REQUEST) + + sent = [m for m in provider.requests[-1]["messages"] if m["role"] != "system"] + assert [m["role"] for m in sent] == ["user", "assistant", "user"], sent + assert [m["content"] for m in sent if m["role"] == "user"] == [_REFUSED, _NEW_REQUEST] + + +def test_failed_turn_boundary_is_idempotent_on_the_durable_tail_and_skips_context_overflow(tmp_path, monkeypatch): + """Keyed on ``SessionDB.latest_conversation_role``: a redelivery of an already-closed turn + writes no second row, an assistant tail is left alone, and the context-overflow class + (whose repair is rotation, not a row — #1630) is never closed.""" + from types import SimpleNamespace + + from agent.conversation_loop import _close_durable_failed_turn + from agent.turn_failure_copy import PARTIAL_FAILED_TURN_NOTICE + from hermes_state import SessionDB + + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + db = SessionDB(tmp_path / "state.db") + sid = "s1" + db.create_session(session_id=sid, source="acp", model="m") + db.append_message(sid, "user", "delete the build directory") + + written: list[list[dict]] = [] + + def flush(messages): + db.append_message(sid, messages[-1]["role"], messages[-1]["content"]) + written.append(list(messages)) + + agent = SimpleNamespace(_session_db=db, session_id=sid, _flush_messages_to_session_db=flush) + messages = [ + {"role": "user", "content": "delete the build directory"}, + {"role": "assistant", "content": None, "tool_calls": [{"id": "t1", "type": "function", + "function": {"name": "terminal", "arguments": "{}"}}]}, + {"role": "tool", "tool_call_id": "t1", "content": "removed"}, + ] + + overflow = {"completed": False, "failed": True, "failure_reason": "context_overflow", "messages": list(messages)} + _close_durable_failed_turn(agent, overflow) + assert written == [] and db.latest_conversation_role(sid) == "user" + + failed = {"completed": False, "failed": True, "failure_reason": "server_error", "messages": messages, + "current_turn_user_idx": 0} + _close_durable_failed_turn(agent, failed) + _close_durable_failed_turn(agent, failed) # redelivery: tail is already assistant + assert len(written) == 1 + assert (messages[-1]["role"], messages[-1]["content"]) == ("assistant", PARTIAL_FAILED_TURN_NOTICE) # a tool ran: hedge + assert db.latest_conversation_role(sid) == "assistant" + db.close()