diff --git a/agent/replay_cleanup.py b/agent/replay_cleanup.py index 814cc48fbe..c0ef24122b 100644 --- a/agent/replay_cleanup.py +++ b/agent/replay_cleanup.py @@ -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 diff --git a/agent/turn_context.py b/agent/turn_context.py index 63208dd191..4365ed71b3 100644 --- a/agent/turn_context.py +++ b/agent/turn_context.py @@ -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): diff --git a/tests/agent/test_replay_cleanup.py b/tests/agent/test_replay_cleanup.py index 7e51080318..6e3c4c5ad8 100644 --- a/tests/agent/test_replay_cleanup.py +++ b/tests/agent/test_replay_cleanup.py @@ -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"