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:
@@ -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
|
||||
|
||||
@@ -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"]
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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."""
|
||||
|
||||
211
tests/acp_adapter/test_failed_turn_closure.py
Normal file
211
tests/acp_adapter/test_failed_turn_closure.py
Normal 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()
|
||||
Reference in New Issue
Block a user