fix(agent): send-path canonicalization is prefix-only, clock is admission time, corrupt stamps fail closed

Follow-up to the two cherry-picked commits from #105308 (@JoaoMarcos44), closing the
three blockers raised on that thread plus one regression the salvage found:

- Prefix-only on the send path. build_api_messages now canonicalizes only
  messages[:current_turn_user_idx]; rows the current turn appended (its own tool
  calls/results) pass through verbatim. Canonicalizing the live tail rewrote a block
  the previous iteration had already sent whenever a tool result matched the
  interrupt heuristic, which is exactly the mid-turn prefix rewrite this fix exists
  to remove, and it also made the dangling-tail transform order-dependent on when
  the user row was appended.
- Exact interrupt marker. is_interrupted_tool_result matched
  "exit_code" + ("130" | "-1") + "interrupt" as substrings, so an ordinary
  `grep KeyboardInterrupt` result next to a diff hunk header rewrote a terminal
  result to an orphan notice (or dropped a read-only block). That heuristic was
  tolerable at resume time only; it now runs per request. Match the executors'
  bracketed markers ("[Command interrupted", "[execution interrupted") and nothing else.
- Admission-time clock. The frozen expiry clock was the input's platform-event
  stamp, so a message queued 70 s before the turn ran kept a 129 s-old confirmation
  live on the send path while replay expired it. _reset_per_turn_agent_state stamps
  time.time() once at admission; the three other writes (bind identity, stage
  message, build_api_messages write-back under suppress(Exception)) are gone.
- Fail closed on corrupt stamps. A present-but-unparseable timestamp (`"nan"`,
  `"not_a_number"`) made strip_stale_dangerous_confirmations keep the confirmation
  and its api_content sidecar. Coerce through hermes_cli.timefmt.coerce_epoch and
  treat an unknowable age as expired; missing stamps (legacy rows) are still left
  alone.
- Shape: drop the canonicalize_history_for_send alias (no consumer, never existed on
  main), the `now=` kwarg (no production caller), and the getattr/hasattr rewrite of
  _reset_per_turn_agent_state (only the test double needed it).
- Tests: 17 → 2 invariant tests. Real SessionDB round trip → canonicalize →
  ChatCompletionsTransport bytes, equal to the send path with sidecars applied and
  the durable list untouched, live tail preserved; admission-clock freeze across
  iterations + corrupt-stamp fail-closed. Each is red under the matching mutation
  (send path unpatched, whole-list canonicalization, per-request clock, fail-open,
  loose heuristic).
This commit is contained in:
kshitijk4poor
2026-09-12 22:55:16 +05:30
committed by kshitij
parent 401fef6e6f
commit 820d3ca65d
3 changed files with 138 additions and 395 deletions

View File

@@ -7,12 +7,14 @@ re-issues the unanswered call → endless "thinking"/reboot loop. These pure hel
from __future__ import annotations
import logging
import math
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
from agent.turn_context import drop_stale_api_content
from hermes_cli.timefmt import coerce_epoch
logger = logging.getLogger(__name__)
@@ -27,12 +29,20 @@ _DANGLING_NOTICES = (
)
# Every executor renders an interrupt as a bracketed marker: "[Command interrupted]" (terminal
# backends, tools/environments/), "[Command interrupted - Modal ...]" (managed_modal.py) and
# "[execution interrupted ...]" (code_execution_tool.py).
_INTERRUPT_MARKERS = ("[command interrupted", "[execution interrupted")
def is_interrupted_tool_result(content: Any) -> bool:
"""Return True if a tool result indicates the tool was interrupted."""
"""True only for an executor's interrupt marker. Nothing looser: this also runs on every
live request, where a substring heuristic rewrote an ordinary ``grep KeyboardInterrupt``
result mid-turn and broke the cached prefix."""
if not isinstance(content, str):
return False
lowered = content.lower()
return "[command interrupted]" in lowered or ("exit_code" in lowered and ("130" in lowered or "-1" in lowered) and "interrupt" in lowered)
return any(marker in lowered for marker in _INTERRUPT_MARKERS)
def _call_name(call: Dict[str, Any]) -> str:
@@ -132,15 +142,11 @@ def canonicalize_replay_history(
) -> 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.
Resume surfaces and the send path must serialize the same history bytes, or a
resumed request diverges in the middle of the cached prefix.
The input is never modified. ``now`` is injectable for deterministic tests;
production callers use the same wall clock as the existing expiry policy.
The input is never modified. ``now`` is the expiry clock; the send path passes the
turn's admission time so every request in one turn sees the same bytes.
"""
if not agent_history:
return agent_history
@@ -151,10 +157,6 @@ def canonicalize_replay_history(
return strip_stale_dangerous_confirmations(cleaned, now=now)
# Backward-compatible alias for the send-path name (2026-09-07 code).
canonicalize_history_for_send = canonicalize_replay_history
# --- Stale dangerous-confirmation text expiry ---
# Short on purpose: a dangerous confirmation must not survive any restart or resume gap.
@@ -203,21 +205,19 @@ def strip_stale_dangerous_confirmations(
cleaned: List[Dict[str, Any]] = []
for msg in agent_history:
ts = msg.get("timestamp") if isinstance(msg, dict) and msg.get("role") == "user" else None
try:
is_stale = (
ts is not None
and is_dangerous_confirmation(msg.get("content", ""))
and (float(now) - float(ts)) > expiry_seconds
)
except (ValueError, TypeError):
is_stale = False
if not is_stale:
if ts is None or not is_dangerous_confirmation(msg.get("content", "")):
cleaned.append(msg)
continue
# A present-but-corrupt stamp is treated as expired: the age is unknowable, and
# keeping the text (plus its api_content sidecar) would replay a live confirmation.
ts_f = coerce_epoch(ts, field="message timestamp")
age = math.inf if ts_f is None else now - ts_f
if age <= expiry_seconds:
cleaned.append(msg)
continue
logger.debug(
"Redacting stale dangerous-confirmation text in user message (age=%.1fs, expiry=%.1fs): %r",
float(now) - float(ts), expiry_seconds, (msg.get("content") or "")[:80],
age, expiry_seconds, (msg.get("content") or "")[:80],
)
redacted = dict(msg)
redacted["content"] = _EXPIRED_CONFIRMATION_SENTINEL

View File

@@ -466,7 +466,6 @@ def _bind_turn_identity(
agent._persist_user_message_override = persist_user_message
agent._persist_user_message_timestamp = persist_user_timestamp
agent._persist_user_message_platform_id = persist_user_platform_id
agent._current_turn_timestamp = persist_user_timestamp
# Unique task_id when not provided isolates VMs between tasks.
effective_task_id = task_id or str(uuid.uuid4())
agent._current_task_id = effective_task_id
@@ -495,7 +494,7 @@ _PER_TURN_RESET_STATE: Tuple[Tuple[str, Any], ...] = (
("_tool_guardrail_halt_decision", None), ("_vision_supported", True),
("_iteration_budget_warning_injected", False),
("_run_budget_wrapup_injected", False), ("_verification_stop_nudges", 0),
("_pre_verify_nudges", 0), ("_current_turn_timestamp", None),
("_pre_verify_nudges", 0),
)
@@ -505,13 +504,14 @@ def _reset_per_turn_agent_state(agent: Any) -> None:
setattr(agent, name, value)
agent._turn_failed_file_mutations = {}
agent._turn_file_mutation_paths = set()
_guardrails = getattr(agent, "_tool_guardrails", None)
if _guardrails is not None and hasattr(_guardrails, "reset_for_turn"):
_guardrails.reset_for_turn()
_mem = getattr(agent, "_memory_store", None)
_reset_consol = getattr(_mem, "reset_consolidation_failures", None) if _mem is not None else None
agent._tool_guardrails.reset_for_turn()
_reset_consol = getattr(agent._memory_store, "reset_consolidation_failures", None)
if callable(_reset_consol):
_reset_consol()
# Expiry clock for build_api_messages: admission time (not the input's platform-event
# stamp, which can predate admission by minutes), frozen so every request this turn
# sends identical bytes.
agent._current_turn_timestamp = time.time()
# Pre-turn connection health check: clean up dead TCP connections.
if agent.api_mode != "anthropic_messages":
@@ -567,8 +567,6 @@ def _stage_turn_user_message(
# CLI input is stamped when staged; gateway input may carry the platform event
# time. Preserve either value and cover any legacy unstamped handoff.
stamp_message_timestamp(user_msg, timestamp=persist_user_timestamp)
if agent is not None and getattr(agent, "_current_turn_timestamp", None) is None:
agent._current_turn_timestamp = user_msg.get("timestamp")
# Synthesized turns stamp their transcript type so the crash persist writes a typed
# row; the model still receives role/content unchanged (api_messages strips both).
@@ -1040,7 +1038,6 @@ def _sanitize_model_for(agent: Any, moa_config: Any) -> Any:
def build_api_messages(
agent: Any, messages: List[Dict[str, Any]], *, current_turn_user_idx: Any,
ext_prefetch_cache: Any, plugin_user_context: Any, moa_config: Any, active_system_prompt: Any,
now: Optional[float] = None,
) -> Tuple[List[Dict[str, Any]], str]:
"""Build the wire copy of ``messages`` for one API call plus the effective system
message. Returns ``(api_messages, effective_system)``.
@@ -1056,39 +1053,19 @@ def build_api_messages(
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
)
has_current = isinstance(current_turn_user_idx, int) and 0 <= current_turn_user_idx < len(messages)
current_turn_message = messages[current_turn_user_idx] if has_current else None
turn_now = now
if turn_now is None and agent is not None:
_agent_ts = getattr(agent, "_current_turn_timestamp", None)
if isinstance(_agent_ts, (int, float)):
turn_now = float(_agent_ts)
if turn_now is None and isinstance(current_turn_message, dict):
_msg_ts = current_turn_message.get("timestamp")
if isinstance(_msg_ts, (int, float)):
turn_now = float(_msg_ts)
elif isinstance(_msg_ts, str):
try:
turn_now = float(_msg_ts)
except ValueError:
pass
if turn_now is None:
turn_now = time.time()
if agent is not None:
with suppress(Exception):
agent._current_turn_timestamp = turn_now
# 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.
# The expiry evaluation is frozen for the active turn so tool-loop iterations
# cannot rewrite the prefix or withdraw confirmation mid-turn.
canonical_messages = canonicalize_replay_history(messages, now=turn_now)
# Replay consumers canonicalize the persisted prefix on read; the request copy must
# carry the same bytes or a resume diverges mid-prefix. Only the rows BEFORE this
# turn's user message are the replayed prefix — rows this turn appended (its tool
# calls/results) are live and must never be rewritten between iterations. The
# expiry clock is the turn's admission time, frozen in _reset_per_turn_agent_state.
# Without an anchor (compaction found no surviving user row) there is no provable
# persisted prefix, so nothing is canonicalized.
turn_now = getattr(agent, "_current_turn_timestamp", None) or time.time()
split = current_turn_user_idx if has_current else 0
canonical_messages = canonicalize_replay_history(messages[:split], now=turn_now) + messages[split:]
api_messages = []
for idx, msg in enumerate(canonical_messages):

View File

@@ -6,51 +6,12 @@ same way. Regression coverage for #29086 (WebUI session permanently stuck
because the dangling tool-call tail was replayed on every resume).
"""
import copy
import json
from pathlib import Path
import tempfile
import time
from agent.replay_cleanup import (
canonicalize_history_for_send,
canonicalize_replay_history,
is_interrupted_tool_result,
strip_dangling_tool_call_tail,
strip_interrupted_tool_tails,
strip_stale_dangerous_confirmations,
sanitize_replay_history,
)
from agent.transports.chat_completions import ChatCompletionsTransport
from agent.turn_context import build_api_messages
from hermes_state import SessionDB
def _wire(messages):
return ChatCompletionsTransport().convert_messages(list(messages))
def _canon(objs):
return json.dumps(objs, sort_keys=True, separators=(",", ":"))
class _Agent:
api_mode = "chat_completions"
ephemeral_system_prompt = None
_compression_warning = None
max_iterations = 10
@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
def _user(text):
@@ -134,319 +95,124 @@ 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"},
]
# --- Send/replay canonicalization parity (#105236 §6, salvage of #105308) ---
expected = strip_stale_dangerous_confirmations(
sanitize_replay_history(copy.deepcopy(history)), now=now
)
actual = canonicalize_replay_history(copy.deepcopy(history), now=now)
import copy
import json
assert actual == expected
from agent.replay_cleanup import canonicalize_replay_history
from agent.transports.chat_completions import ChatCompletionsTransport
from agent.turn_context import build_api_messages
from hermes_state import SessionDB
def test_canonicalize_history_for_send_alias():
assert canonicalize_history_for_send is canonicalize_replay_history
class _SendAgent:
api_mode = "chat_completions"
ephemeral_system_prompt = None
_compression_warning = None
_current_turn_timestamp = 10_000.0
@staticmethod
def _copy_reasoning_content_for_api(_source, _target):
return None
@staticmethod
def _should_sanitize_tool_calls():
return False
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
def _wire(messages):
return json.dumps(ChatCompletionsTransport().convert_messages(list(messages)), sort_keys=True)
now = 10_000.0
agent = _Agent()
agent._current_turn_timestamp = now
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)
def _send(agent, 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"
# Wire representation comparison: send wire matches replay wire with sidecar applied
wire_request = _wire(request)
expected_replay = canonicalize_replay_history(copy.deepcopy(history[:-1]), now=now) + [
{"role": "user", "content": "current-wire"}
]
assert _canon(wire_request) == _canon(_wire(expected_replay))
return request
def test_send_byte_identity_with_tui_replay_interrupted_block():
"""Interrupted read-only assistant->tool block: send path wire bytes match TUI replay."""
def test_send_wire_matches_replay_wire_after_db_round_trip(tmp_path):
"""The bytes a resumed session replays and the bytes the live send path emits for the
same persisted prefix are identical through the real transport, sidecars applied; the
durable transcript is untouched; rows appended by the CURRENT turn are never rewritten."""
now = 10_000.0
agent = _Agent()
agent._current_turn_timestamp = now
history = [
_user("u1"),
{"role": "assistant", "content": "a1"},
_assistant_tc("read_file"), _tool("[command interrupted]"),
_user("u2"),
db = SessionDB(db_path=tmp_path / "t.db")
db.create_session(session_id="s1", source="cli")
db.append_message("s1", role="user", content="hello", api_content="hello [with memory]", timestamp=now - 300)
db.append_message("s1", role="assistant", content="hi", timestamp=now - 299)
db.append_message("s1", role="user", content="confirm reboot", timestamp=now - 120)
db.append_message("s1", role="assistant", content="", tool_calls=[
{"id": "c1", "type": "function", "function": {"name": "read_file", "arguments": "{}"}}], timestamp=now - 119)
db.append_message("s1", role="tool", content="[execution interrupted — user stop]", tool_call_id="c1", tool_name="read_file", timestamp=now - 118)
# An ordinary result that merely mentions interrupts is NOT replay debris.
grep_hit = '{"output": "loop.py:12: except KeyboardInterrupt:\\n@@ -1,4 +1,5 @@", "exit_code": 0}'
db.append_message("s1", role="assistant", content="", tool_calls=[
{"id": "c2", "type": "function", "function": {"name": "search_files", "arguments": "{}"}}], timestamp=now - 110)
db.append_message("s1", role="tool", content=grep_hit, tool_call_id="c2", tool_name="search_files", timestamp=now - 109)
persisted = db.get_messages_as_conversation("s1")
db.close()
assert persisted[0].get("api_content") == "hello [with memory]" and persisted[2].get("timestamp")
# What every resume surface feeds the model, with the sidecar bytes the send path replays.
replay = [{**m, "content": m.get("api_content") or m.get("content")}
for m in canonicalize_replay_history(persisted, now=now)]
live = copy.deepcopy(persisted) + [{"role": "user", "content": "now", "timestamp": now}]
frozen = copy.deepcopy(live)
request = _send(_SendAgent(), live)
assert live == frozen
assert _wire(request) == _wire(replay + [{"role": "user", "content": "now"}])
assert "[with memory]" in request[0]["content"] and "EXPIRED" in request[2]["content"]
assert [m["role"] for m in request] == ["user", "assistant", "user", "assistant", "tool", "user"]
assert request[4]["content"] == grep_hit
# Rows this turn appended stay verbatim even when they look like replay debris.
live += [
{"role": "assistant", "content": "", "tool_calls": [
{"id": "c3", "type": "function", "function": {"name": "search_files", "arguments": "{}"}}]},
{"role": "tool", "tool_call_id": "c3", "content": "[Command interrupted]"},
]
send_request, _ = build_api_messages(
agent, history, current_turn_user_idx=len(history) - 1, ext_prefetch_cache="",
request2 = build_api_messages(
_SendAgent(), live, current_turn_user_idx=len(persisted), ext_prefetch_cache="",
plugin_user_context="", moa_config=None, active_system_prompt="",
)
wire_send = _wire(send_request)
wire_replay = _wire(sanitize_replay_history(copy.deepcopy(history)))
assert _canon(wire_send) == _canon(wire_replay)
)[0]
assert _wire(request2[: len(request)]) == _wire(request)
assert request2[-1]["content"] == "[Command interrupted]"
def test_send_byte_identity_with_dangling_tool_call_tail():
"""Trailing unanswered assistant(tool_calls): send path wire bytes match TUI replay."""
now = 10_000.0
history = [
_user("u1"),
{"role": "assistant", "content": "a1"},
{"role": "assistant", "content": "", "tool_calls": [{"id": "c9", "type": "function", "function": {"name": "read_file", "arguments": "{}"}}]},
]
wire_replay = _wire(sanitize_replay_history(copy.deepcopy(history)))
wire_canon = _wire(canonicalize_replay_history(copy.deepcopy(history), now=now))
assert _canon(wire_canon) == _canon(wire_replay)
def test_send_byte_identity_with_stale_dangerous_confirmation():
"""Stale confirmation (>60s): send path wire bytes match gateway replay (redacted to sentinel)."""
now = 10_000.0
agent = _Agent()
agent._current_turn_timestamp = now
history = [
{"role": "user", "content": "confirm forced restart", "timestamp": now - 120.0},
{"role": "assistant", "content": "a1"},
{"role": "user", "content": "u2", "timestamp": now},
]
send_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="",
)
wire_send = _wire(send_request)
wire_replay = _wire(strip_stale_dangerous_confirmations(copy.deepcopy(history), now=now))
assert _canon(wire_send) == _canon(wire_replay)
assert any("EXPIRED" in (m.get("content") or "") for m in wire_send)
def test_send_byte_identity_with_fresh_confirmation():
"""Fresh confirmation (<60s): send path wire bytes match gateway replay (preserved verbatim)."""
now = 10_000.0
agent = _Agent()
agent._current_turn_timestamp = now
history = [
{"role": "user", "content": "confirm reboot", "timestamp": now - 30.0},
{"role": "assistant", "content": "a1"},
{"role": "user", "content": "u2", "timestamp": now},
]
send_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="",
)
wire_send = _wire(send_request)
wire_replay = _wire(strip_stale_dangerous_confirmations(copy.deepcopy(history), now=now))
assert _canon(wire_send) == _canon(wire_replay)
assert wire_send[0]["content"] == "confirm reboot"
def test_send_byte_identity_clean_history():
"""Clean history (no interrupted blocks, no stale confirmations): send matches replay."""
now = 10_000.0
agent = _Agent()
agent._current_turn_timestamp = now
history = [
{"role": "user", "content": "u1"},
{"role": "assistant", "content": "a1"},
{"role": "user", "content": "u2"},
]
send_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="",
)
wire_send = _wire(send_request)
wire_replay = _wire(canonicalize_replay_history(copy.deepcopy(history), now=now))
assert _canon(wire_send) == _canon(wire_replay) == _canon(_wire(history))
def test_send_byte_identity_with_sidecar():
"""Historical user turn with api_content sidecar: wire representation reproduces sidecar bytes."""
now = 10_000.0
agent = _Agent()
agent._current_turn_timestamp = now
history = [
{"role": "user", "content": "hello", "api_content": "hello [with memory]", "timestamp": now - 100},
{"role": "assistant", "content": "hi"},
{"role": "user", "content": "current", "timestamp": now},
]
send_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="",
)
wire_send = _wire(send_request)
assert wire_send[0]["content"] == "hello [with memory]"
# Replay with sidecar intact matches
replay = canonicalize_replay_history(copy.deepcopy(history), now=now)
assert replay[0]["api_content"] == "hello [with memory]"
def test_active_turn_expiry_decision_frozen_across_tool_iterations(monkeypatch):
"""Deterministic wire-level probe (ehz0ah blocking defect):
Confirmation timestamp: 9941.0.
Active turn starts: 10000.0 (age = 59s <= 60s, fresh).
Request 1 assembled at 10000.0 sends 'confirm reboot' on wire.
Tool executes; request 2 assembled at 10002.0 (age = 61s > 60s).
Because both requests belong to the same active turn, the expiry decision
is frozen: request 2 does NOT rewrite the confirmation to EXPIRED, preserving
the prompt cache prefix across tool iterations.
Subsequent turn N+1 at 10070.0 DOES expire the confirmation and matches replay.
"""
def test_confirmation_expiry_uses_frozen_admission_clock_and_fails_closed(monkeypatch):
"""Expiry is judged once per turn at admission (not the input's event stamp, not
per-request wall time); a present-but-corrupt stamp is treated as expired."""
from agent.turn_context import _reset_per_turn_agent_state
agent = _Agent()
agent = _SendAgent()
agent._tool_guardrails = type("G", (), {"reset_for_turn": staticmethod(lambda: None)})()
agent._memory_store = None
agent.max_iterations = 4
history = [
{"role": "user", "content": "confirm reboot", "timestamp": 9941.0},
{"role": "assistant", "content": "Preparing reboot..."},
{"role": "user", "content": "proceed now", "timestamp": 10000.0},
{"role": "user", "content": "confirm reboot", "timestamp": 9_941.0},
{"role": "assistant", "content": "ok"},
{"role": "user", "content": "go", "timestamp": 10_000.0}, # platform event stamp: 59s old
]
turn_user_idx = 2
# Request 1 (first iteration of turn at t=10000.0):
monkeypatch.setattr(time, "time", lambda: 10000.0)
req1, _ = build_api_messages(
agent, history, current_turn_user_idx=turn_user_idx, ext_prefetch_cache="",
plugin_user_context="", moa_config=None, active_system_prompt="",
)
wire1 = _wire(req1)
assert wire1[0]["content"] == "confirm reboot"
# Tool executes during this turn; model response + tool result appended to history:
history.append({
"role": "assistant", "content": "",
"tool_calls": [{"id": "c1", "type": "function", "function": {"name": "reboot_check", "arguments": "{}"}}],
})
history.append({"role": "tool", "tool_call_id": "c1", "content": "ready"})
# Request 2 (second iteration of the SAME turn at t=10002.0 > 60s expiry threshold):
monkeypatch.setattr(time, "time", lambda: 10002.0)
req2, _ = build_api_messages(
agent, history, current_turn_user_idx=turn_user_idx, ext_prefetch_cache="",
plugin_user_context="", moa_config=None, active_system_prompt="",
)
wire2 = _wire(req2)
# The confirmation row MUST NOT have mutated to EXPIRED mid-turn:
assert wire2[0]["content"] == "confirm reboot"
# The prefix (all messages prior to the new tool call) remains byte-identical:
assert _canon(wire1) == _canon(wire2[:len(wire1)])
# Turn N+1: user sends a new message at t=10070.0 (well past expiry):
monkeypatch.setattr("agent.turn_context.time.time", lambda: 10_070.0) # admitted 129s later
_reset_per_turn_agent_state(agent)
history.append({"role": "user", "content": "system status", "timestamp": 10070.0})
monkeypatch.setattr(time, "time", lambda: 10070.0)
req3, _ = build_api_messages(
agent, history, current_turn_user_idx=len(history) - 1, ext_prefetch_cache="",
plugin_user_context="", moa_config=None, active_system_prompt="",
)
wire3 = _wire(req3)
# In the new turn, the confirmation row IS expired on wire:
assert "EXPIRED" in wire3[0]["content"]
# Wire matches replay canonicalization at this turn boundary:
wire_replay = _wire(canonicalize_replay_history(copy.deepcopy(history[:-1]), now=10070.0) + [history[-1]])
assert _canon(wire3) == _canon(wire_replay)
request = _send(agent, history)
assert "EXPIRED" in request[0]["content"]
assert _wire(request[:2]) == _wire(canonicalize_replay_history(history[:2], now=10_070.0))
agent._current_turn_timestamp = 9_990.0 # admitted at 49s: fresh, and stays fresh...
history += [{"role": "assistant", "content": "", "tool_calls": [
{"id": "c1", "type": "function", "function": {"name": "read_file", "arguments": "{}"}}]},
{"role": "tool", "tool_call_id": "c1", "content": "ready"}]
monkeypatch.setattr("agent.turn_context.time.time", lambda: 10_500.0) # ...however long the tools take
late = build_api_messages(agent, history, current_turn_user_idx=2, ext_prefetch_cache="",
plugin_user_context="", moa_config=None, active_system_prompt="")[0]
assert late[0]["content"] == "confirm reboot"
def test_canonicalize_is_idempotent_and_non_mutating():
"""canonicalize_replay_history is idempotent and does not mutate source."""
now = 10_000.0
history = [
_user("u1"),
{"role": "assistant", "content": "a1"},
_assistant_tc("read_file"), _tool("[command interrupted]"),
{"role": "user", "content": "confirm forced restart", "timestamp": now - 120},
{"role": "assistant", "content": "a2"},
_user("u3"),
]
original = copy.deepcopy(history)
out1 = canonicalize_replay_history(copy.deepcopy(history), now=now)
out2 = canonicalize_replay_history(copy.deepcopy(out1), now=now)
assert _canon(_wire(out1)) == _canon(_wire(out2))
assert history == original
def test_db_roundtrip_byte_identity():
"""SessionDB round-trip: stored messages read back and canonicalized match wire."""
messages = [
{"role": "user", "content": "check system"},
{"role": "assistant", "content": "all ok"},
{"role": "user", "content": "proceed"},
]
with tempfile.TemporaryDirectory() as tmp:
db = SessionDB(db_path=Path(tmp) / "t.db")
try:
db.create_session(session_id="s1", source="cli")
for m in messages:
db.append_message("s1", role=m["role"], content=m["content"])
conv = db.get_messages_as_conversation("s1")
read_back = [
{"role": m["role"], "content": m["content"]}
for m in conv
if m.get("content") is not None
]
canon_read = canonicalize_replay_history(read_back)
assert _canon(_wire(canon_read)) == _canon(_wire(messages))
finally:
db.close()
def test_canonicalize_replay_history_handles_malformed_timestamps():
"""Malformed or non-numeric timestamps must not raise exceptions and be safely preserved."""
now = 10_000.0
history = [
{"role": "user", "content": "confirm reboot", "timestamp": "2026-09-08T00:00:00Z"},
{"role": "user", "content": "confirm reboot", "timestamp": "not_a_number"},
{"role": "user", "content": "confirm reboot", "timestamp": None},
{"role": "user", "content": "confirm reboot", "timestamp": now - 120.0},
]
agent = _Agent()
# Should not raise TypeError or ValueError
canon = canonicalize_replay_history(history, now=now)
assert len(canon) == 4
# String / None timestamps are left untouched (not expired)
assert canon[0]["content"] == "confirm reboot"
assert canon[1]["content"] == "confirm reboot"
assert canon[2]["content"] == "confirm reboot"
# The valid numeric timestamp older than 60s is expired cleanly
assert "EXPIRED" in canon[3]["content"]
# Also verify build_api_messages with string timestamp on current_turn_message
res, _ = build_api_messages(
agent,
[{"role": "user", "content": "test", "timestamp": "2026-09-08T00:00:00Z"}],
current_turn_user_idx=0,
ext_prefetch_cache="",
plugin_user_context="",
moa_config=None,
active_system_prompt="",
)
assert len(res) == 1
assert res[0]["content"] == "test"
corrupt = [{"role": "user", "content": "confirm reboot", "timestamp": "nan", "api_content": "confirm reboot"}]
out = canonicalize_replay_history(corrupt, now=10_000.0)
assert "EXPIRED" in out[0]["content"] and "api_content" not in out[0]
assert canonicalize_replay_history([{"role": "user", "content": "confirm reboot"}], now=1e9)[0]["content"] == "confirm reboot"