From e5ca5207de04e223dda7bc0ffac4c72777a9a4a8 Mon Sep 17 00:00:00 2001 From: joaomarcos Date: Mon, 7 Sep 2026 16:12:07 -0300 Subject: [PATCH] fix(agent): unify replay history canonicalization --- agent/replay_cleanup.py | 27 +++++++++++- agent/turn_context.py | 18 ++++++-- gateway/run.py | 22 ++++------ gateway/run_turn_runner.py | 8 ++-- tests/agent/test_replay_cleanup.py | 66 ++++++++++++++++++++++++++++++ tui_gateway/methods_session.py | 2 +- tui_gateway/server.py | 9 ++-- 7 files changed, 123 insertions(+), 29 deletions(-) diff --git a/agent/replay_cleanup.py b/agent/replay_cleanup.py index 7af23b143e..8940e0a47f 100644 --- a/agent/replay_cleanup.py +++ b/agent/replay_cleanup.py @@ -7,7 +7,8 @@ re-issues the unanswered call → endless "thinking"/reboot loop. These pure hel from __future__ import annotations import logging -from typing import Any, Dict, List +import time +from typing import Any, Dict, List, Optional from agent.tool_dispatch_helpers import make_tool_result_message from agent.tool_result_classification import tool_may_have_side_effect @@ -126,6 +127,30 @@ def sanitize_replay_history(agent_history: List[Dict[str, Any]]) -> List[Dict[st return strip_dangling_tool_call_tail(strip_interrupted_tool_tails(agent_history)) +def canonicalize_replay_history( + agent_history: List[Dict[str, Any]], *, now: Optional[float] = None +) -> List[Dict[str, Any]]: + """Apply every destructive replay transform in the shared, fixed order. + + Resume surfaces and the send path must serialize the same history bytes. The + older consumers each applied only a subset of these transforms: interrupted + blocks and dangling tails were handled by TUI replay, while stale dangerous + confirmations were handled by gateway replay. A request built from the + unmodified history could therefore diverge in the middle of the cached + prefix after a resume. + + The input is never modified. ``now`` is injectable for deterministic tests; + production callers use the same wall clock as the existing expiry policy. + """ + if not agent_history: + return agent_history + if now is None: + now = time.time() + cleaned = strip_interrupted_tool_tails(agent_history) + cleaned = strip_dangling_tool_call_tail(cleaned) + return strip_stale_dangerous_confirmations(cleaned, now=now) + + # --- Stale dangerous-confirmation text expiry --- # Short on purpose: a dangerous confirmation must not survive any restart or resume gap. diff --git a/agent/turn_context.py b/agent/turn_context.py index 8ff9048519..50b0100610 100644 --- a/agent/turn_context.py +++ b/agent/turn_context.py @@ -1047,9 +1047,21 @@ def build_api_messages( replayed verbatim.""" from agent.agent_runtime_helpers import fill_empty_non_final_wire_payload from agent.conversation_loop import _clone_message_for_send + from agent.replay_cleanup import canonicalize_replay_history + + current_turn_message = ( + messages[current_turn_user_idx] + if isinstance(current_turn_user_idx, int) + and 0 <= current_turn_user_idx < len(messages) + else None + ) + # Replay consumers rewrite interrupted blocks, dangling tails, and expired + # confirmations on read. Apply the exact same transform to this request-only + # copy before sidecars are substituted; the durable transcript remains intact. + canonical_messages = canonicalize_replay_history(messages) api_messages = [] - for idx, msg in enumerate(messages): + for idx, msg in enumerate(canonical_messages): # Structural clone, NOT msg.copy(): in-place transforms below must not reach # persisted history via nested containers; see _clone_message_for_send. api_msg = _clone_message_for_send(msg) @@ -1063,7 +1075,7 @@ def build_api_messages( # Inject ephemeral context (memory prefetch + pre_llm_call user hooks) # at API time only; `messages` is untouched beyond the api_content stamp. - if idx == current_turn_user_idx and msg.get("role") == "user": + if msg is current_turn_message and msg.get("role") == "user": if isinstance(_api_content, str) and _api_content: # Reuse the prologue's stamp so sidecar and wire cannot drift # and every pass this turn sends identical bytes. @@ -1094,7 +1106,7 @@ def build_api_messages( # Fill empty non-final user/assistant wire copies so the pre-call sanitizer # stops re-healing and flooding errors.log; durable history is untouched. # After the reasoning copy so thinking-only turns keep payload. - fill_empty_non_final_wire_payload(api_msg, is_final=(idx == len(messages) - 1)) + fill_empty_non_final_wire_payload(api_msg, is_final=(idx == len(canonical_messages) - 1)) # _thinking_prefill survives intentionally: the drop pass below needs it. # Strip length-continuation marks; some transports keep underscore keys. api_msg.pop("_length_continuation_fragment", None) diff --git a/gateway/run.py b/gateway/run.py index 1c7a68c8f5..fcaffe205a 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -1250,18 +1250,10 @@ def _build_gateway_agent_history( entry.pop("api_content", None) # prefix rewrite: the sidecar no longer matches agent_history.append(entry) - # Strip interrupted tool-call tails so the LLM doesn't re-execute tools killed mid-flight. - agent_history = strip_interrupted_tool_tails(agent_history) - - # Strip a dangling assistant(tool_calls) tail (SIGKILL-mid-tool-call); else the model re-issues it forever. - # Strip a dangling assistant(tool_calls) tail with no tool answers — the signature of a SIGKILL - # mid-tool-call (e.g. the tool itself ran `docker restart`/`kill` and took the gateway down before the - # result was persisted). Without this the model re-issues the unanswered call on resume and loops the - # restart forever (#49201). - agent_history = strip_dangling_tool_call_tail(agent_history) - - # Strip expired dangerous-confirmation phrases; replayed, a follow-up could read as a fresh confirmation. - agent_history = strip_stale_dangerous_confirmations(agent_history, now=time.time()) + # Keep gateway resume byte-identical to the TUI resume and send paths. The + # canonicalizer owns interrupted-block, dangling-tail, and stale-confirmation + # cleanup together so a middle-of-history rewrite cannot break the prefix cache. + agent_history = canonicalize_replay_history(agent_history) observed_context = "\n".join(observed_group_context).strip() or None return agent_history, observed_context @@ -1334,9 +1326,9 @@ def _last_transcript_timestamp(history: Optional[List[Dict[str, Any]]]) -> Any: # Tool output may hold literal MEDIA: examples (docs, logs); only deliberate media producers may auto-append. _AUTO_APPEND_MEDIA_TOOL_NAMES = {"text_to_speech", "text_to_speech_tool", "image_generate"} -# Replay-tail sanitization lives in agent/replay_cleanup.py so every resume surface shares one implementation. -from agent.replay_cleanup import ( # noqa: E402 - strip_interrupted_tool_tails, strip_dangling_tool_call_tail, strip_stale_dangerous_confirmations) +# Replay-history canonicalization lives in agent/replay_cleanup.py so every resume +# surface and the send path share one implementation. +from agent.replay_cleanup import canonicalize_replay_history # noqa: E402 _AUTO_CONTINUE_NOTE_PREFIX = "[System note: Your previous turn" diff --git a/gateway/run_turn_runner.py b/gateway/run_turn_runner.py index 3a147a5c1f..837492da3a 100644 --- a/gateway/run_turn_runner.py +++ b/gateway/run_turn_runner.py @@ -20,7 +20,7 @@ from datetime import datetime from typing import TYPE_CHECKING, Any, Dict, List, Optional from agent.interrupt_compat import _accepts_keyword -from agent.replay_cleanup import strip_stale_dangerous_confirmations +from agent.replay_cleanup import canonicalize_replay_history from gateway.config import Platform from gateway.media_repair import repair_explicit_computer_use_media_paths from gateway.platforms.base import BasePlatformAdapter @@ -1419,9 +1419,9 @@ class TurnRunner: "conversation context (possible FTS write corruption)", ctx.session_key, len(agent_history), len(selected), ) - # The live history bypassed _build_gateway_agent_history's cleanup — re-apply the - # stale-confirmation expiry so a dangerous confirmation can't slip through. - agent_history = strip_stale_dangerous_confirmations(selected, now=time.time()) + # The live history bypassed _build_gateway_agent_history's cleanup — re-apply + # the full canonicalization so no replay transform can slip through. + agent_history = canonicalize_replay_history(selected) # MEDIA paths already in history are excluded from this turn's extraction (compression-safe). return agent_history, observed_group_context, _collect_history_media_paths(agent_history) diff --git a/tests/agent/test_replay_cleanup.py b/tests/agent/test_replay_cleanup.py index afed30a407..5aeaf9448f 100644 --- a/tests/agent/test_replay_cleanup.py +++ b/tests/agent/test_replay_cleanup.py @@ -6,10 +6,14 @@ same way. Regression coverage for #29086 (WebUI session permanently stuck because the dangling tool-call tail was replayed on every resume). """ +import copy + from agent.replay_cleanup import ( + canonicalize_replay_history, is_interrupted_tool_result, strip_dangling_tool_call_tail, strip_interrupted_tool_tails, + strip_stale_dangerous_confirmations, sanitize_replay_history, ) @@ -93,3 +97,65 @@ def test_sanitize_replay_history_noop_on_clean_history(): def test_sanitize_replay_history_empty(): assert sanitize_replay_history([]) == [] + + +def test_canonicalize_replay_history_matches_all_resume_transforms(): + """Send and resume consumers must apply the same destructive transforms.""" + now = 10_000.0 + history = [ + _user("before"), + _assistant_tc("read_file"), _tool("[command interrupted]"), + {"role": "user", "content": "confirm forced restart", "timestamp": now - 120}, + {"role": "assistant", "content": "ack"}, + ] + + expected = strip_stale_dangerous_confirmations( + sanitize_replay_history(copy.deepcopy(history)), now=now + ) + actual = canonicalize_replay_history(copy.deepcopy(history), now=now) + + assert actual == expected + + +def test_send_builder_uses_canonical_history_without_mutating_source(): + """The request copy must match replay cleanup while durable history stays intact.""" + from agent.turn_context import build_api_messages + + class _Agent: + api_mode = "chat_completions" + ephemeral_system_prompt = None + + @staticmethod + def _copy_reasoning_content_for_api(_source, _target): + return None + + @staticmethod + def _should_sanitize_tool_calls(): + return False + + @staticmethod + def _sanitize_tool_calls_for_strict_api(*_args, **_kwargs): + return None + + now = 10_000.0 + history = [ + _user("before"), + {"role": "assistant", "content": "ack"}, + _assistant_tc("read_file"), _tool("[command interrupted]"), + {"role": "user", "content": "confirm reboot", "timestamp": now - 120}, + {"role": "assistant", "content": "ack2"}, + {"role": "user", "content": "current", "api_content": "current-wire"}, + ] + original = copy.deepcopy(history) + + request, _ = build_api_messages( + _Agent(), history, current_turn_user_idx=len(history) - 1, ext_prefetch_cache="", + plugin_user_context="", moa_config=None, active_system_prompt="", + ) + + assert history == original + assert [message["role"] for message in request] == [ + "user", "assistant", "user", "assistant", "user" + ] + assert "EXPIRED" in request[2]["content"] + assert request[-1]["content"] == "current-wire" diff --git a/tui_gateway/methods_session.py b/tui_gateway/methods_session.py index 885f0f2e99..47187312ca 100644 --- a/tui_gateway/methods_session.py +++ b/tui_gateway/methods_session.py @@ -518,7 +518,7 @@ class _Resume: def restore(self): """``(sanitized model history, display history, raw history)`` for a cold/eager resume.""" raw, display = self.read_history() - return sanitize_replay_history(raw), display, raw + return canonicalize_replay_history(raw), display, raw def info(self, cwd: str, overrides: dict) -> dict: return _lazy_resume_info(cwd, model=(overrides.get("model_override") or {}).get("model") or "", diff --git a/tui_gateway/server.py b/tui_gateway/server.py index 165afabd00..089c88a29f 100644 --- a/tui_gateway/server.py +++ b/tui_gateway/server.py @@ -29,7 +29,7 @@ from hermes_cli.env_loader import load_hermes_dotenv from utils import is_truthy_value from hermes_state_ids import new_session_id from tools.environments.local import hermes_subprocess_env -from agent.replay_cleanup import sanitize_replay_history +from agent.replay_cleanup import canonicalize_replay_history from agent.compaction_display import project_compaction_message_for_display # noqa: F401 from agent.skill_commands import describe_skill_invocation # noqa: F401 from agent.conversation_loop import INTERRUPT_WAITING_FOR_MODEL_PREFIX # noqa: F401 @@ -2558,10 +2558,9 @@ def _schedule_resume_hydration(sid: str, stored_id: str, db, *, close_db: bool = _emit("session.resume_progress", sid, {"phase": "history", "status": "loading"}) db.reopen_session(stored_id) raw_history, display_history, prefix = _load_resume_transcript(db, stored_id) - # Display keeps the full transcript; the model-fed history drops a dangling/interrupted - # tool-call tail so a session killed mid-loop does not replay the unanswered call forever - # (#29086). - history = sanitize_replay_history(raw_history) + # Display keeps the full transcript; the model-fed history uses the + # same canonicalization as gateway resume and the send path. + history = canonicalize_replay_history(raw_history) if _sessions.get(sid) is not session: return with session["history_lock"]: