fix(agent): close a failed turn's durable user tail so the next prompt is not merged into it

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 <kendrick.kester@gmail.com>
Co-authored-by: Israel Lot <israel.lot@gmail.com>
This commit is contained in:
teknium1
2026-09-18 00:55:30 -07:00
committed by Teknium
parent d60280fde8
commit 8f0322da5b
5 changed files with 288 additions and 10 deletions

View File

@@ -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

View File

@@ -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"]

View File

@@ -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

View File

@@ -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."""

View File

@@ -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()