diff --git a/agent/agent_runtime_helpers.py b/agent/agent_runtime_helpers.py index c4018d20af..6f94189bf0 100644 --- a/agent/agent_runtime_helpers.py +++ b/agent/agent_runtime_helpers.py @@ -3858,10 +3858,22 @@ _INTERRUPTED_PLACEHOLDER = "[response interrupted]" # Repeated heals of the same poisoned transcript used to WARNING on every # send (#96870). Escalate once per session window, then stay quiet. -_EMPTY_HEAL_ESCALATE_AFTER = 5 +# ``_EMPTY_HEAL_ESCALATE_AFTER`` is the built-in default; deployments tune it +# via ``agent.sanitizer_heal_escalation_threshold`` in config.yaml (<= 0 +# disables escalation entirely — WARNINGs still fire per window). +_EMPTY_HEAL_ESCALATE_AFTER = 3 _EMPTY_HEAL_WINDOW_S = 600.0 _empty_heal_log_state: Dict[str, Dict[str, Any]] = {} _empty_heal_log_lock = threading.Lock() +# Session keys that already received the one-time user notice. Separate from +# the windowed log state so a new 10-minute window never re-notifies: the +# user is told ONCE per session, ever (#96870 — out-of-band, delivery +# channel only, never injected into conversation context). +_empty_heal_user_notified: set = set() +# One-shot pending notices keyed by session, drained by the conversation +# loop through ``consume_pending_sanitizer_heal_notice`` and delivered via +# the status/warning callback (the normal delivery channel). +_empty_heal_pending_notice: Dict[str, str] = {} def _msg_has_payload(msg: Dict[str, Any]) -> bool: @@ -3941,25 +3953,106 @@ def _session_id_for_heal_log() -> str: return "" +def _heal_escalation_threshold() -> int: + """Resolve the escalation threshold: config override, else the default. + + ``agent.sanitizer_heal_escalation_threshold`` in config.yaml. Fail-safe: + any read error falls back to the module default so the sanitiser can + never be broken by a bad config file. + """ + try: + from hermes_cli.config import load_config_readonly + + raw = (load_config_readonly().get("agent", {}) or {}).get( + "sanitizer_heal_escalation_threshold" + ) + if raw is not None: + return int(raw) + except Exception: + pass + return _EMPTY_HEAL_ESCALATE_AFTER + + +def consume_pending_sanitizer_heal_notice() -> Optional[str]: + """Drain the one-time user notice for the current session, if any. + + Called by the conversation loop right after the pre-send sanitizer pass; + the returned text is delivered through the status/warning callback (the + normal out-of-band delivery channel: gateway status message, CLI stderr + print). It is NEVER appended to the conversation context, so prompt + caching and role alternation are untouched. Returns at most one notice + per session for its whole lifetime. + """ + key = _session_id_for_heal_log() or "-" + with _empty_heal_log_lock: + return _empty_heal_pending_notice.pop(key, None) + + +def get_sanitizer_heal_stats() -> Dict[str, Dict[str, Any]]: + """Read-only snapshot of per-session sanitiser heal counters. + + Surfaced by diagnostics (``hermes doctor`` / debug share callers) so + repeated silent repairs are visible outside errors.log. Keys are session + ids; values carry ``heal_events`` (sanitizer invocations that healed at + least one message), ``messages_healed`` (total substituted turns) and + ``escalated`` (whether the ERROR + user notice fired). + """ + with _empty_heal_log_lock: + return { + k: { + "heal_events": v.get("total_events", v.get("count", 0)), + "messages_healed": v.get("total_healed", 0), + "escalated": k in _empty_heal_user_notified, + } + for k, v in _empty_heal_log_state.items() + } + + def _log_empty_non_final_heal(healed: int) -> None: """WARNING on the first heals in a window; one ERROR at the threshold. Further heals in the same session window stay silent so a poisoned transcript cannot flood ``errors.log`` (dozens of identical WARNINGs - per hour with no user-visible signal — #96870). + per hour with no user-visible signal — #96870). At the threshold the + escalation also queues a ONE-TIME out-of-band user notice (drained by + ``consume_pending_sanitizer_heal_notice``) pointing at ``/debug share`` + / ``hermes doctor`` — once per session, never re-armed by a new window. """ key = _session_id_for_heal_log() or "-" + threshold = _heal_escalation_threshold() now = time.monotonic() with _empty_heal_log_lock: state = _empty_heal_log_state.get(key) if state is None or (now - state["window_start"]) > _EMPTY_HEAL_WINDOW_S: - state = {"count": 0, "window_start": now, "escalated": False} + prior_events = state.get("total_events", 0) if state else 0 + prior_healed = state.get("total_healed", 0) if state else 0 + state = { + "count": 0, + "window_start": now, + "escalated": False, + "total_events": prior_events, + "total_healed": prior_healed, + } _empty_heal_log_state[key] = state state["count"] += 1 + state["total_events"] = state.get("total_events", 0) + 1 + state["total_healed"] = state.get("total_healed", 0) + healed count = state["count"] - if count >= _EMPTY_HEAL_ESCALATE_AFTER and not state["escalated"]: + total_events = state["total_events"] + total_healed = state["total_healed"] + if threshold > 0 and count >= threshold and not state["escalated"]: state["escalated"] = True level = "error" + if key not in _empty_heal_user_notified: + _empty_heal_user_notified.add(key) + _empty_heal_pending_notice[key] = ( + "⚠️ Your session transcript required repeated repair " + f"({total_events} heal passes so far). Replies keep " + "working, but a corrupted turn is stuck in this " + "session's history — run /debug share or `hermes " + "doctor` to capture diagnostics, or /new to start a " + "clean session." + ) elif state["escalated"]: level = "silent" else: @@ -3969,11 +4062,17 @@ def _log_empty_non_final_heal(healed: int) -> None: return if level == "error": _ra().logger.error( - "Pre-call sanitizer: healed %d empty non-final message(s) " - "(%d heals in this session window). The transcript is being " - "repaired on every send; /new drops the poisoned turns.", + "Pre-call sanitizer: repeated-heal escalation for session %s — " + "healed %d empty non-final message(s) this send; heal pattern: " + "%d heal events / %d messages healed this session " + "(%d in the current session window, threshold %d). The transcript " + "is being repaired on every send; /new drops the poisoned turns.", + key, healed, + total_events, + total_healed, count, + threshold, ) return _ra().logger.warning( diff --git a/agent/conversation_loop.py b/agent/conversation_loop.py index 7cb1346f38..5d2ac2688e 100644 --- a/agent/conversation_loop.py +++ b/agent/conversation_loop.py @@ -2570,6 +2570,24 @@ def run_conversation( # manual message manipulation are always caught. api_messages = agent._sanitize_api_messages(api_messages) + # One-time repeated-heal escalation notice (#96870): if the sanitizer + # above just crossed the per-session heal threshold, deliver the + # queued notice through the status/warning callback — the normal + # out-of-band delivery channel (gateway status message / CLI print). + # NEVER appended to messages/api_messages: conversation context and + # the cached prompt prefix stay byte-identical. + try: + from agent.agent_runtime_helpers import ( + consume_pending_sanitizer_heal_notice, + ) + + _heal_notice = consume_pending_sanitizer_heal_notice() + if _heal_notice: + agent._emit_warning(_heal_notice) + except Exception: + # A notice hiccup must never break the send path. + logger.debug("sanitizer heal notice delivery failed", exc_info=True) + # Drop thinking-only assistant turns (reasoning but no visible # output and no tool_calls) and merge any adjacent user messages # left behind. Prevents Anthropic 400s ("The final block in an diff --git a/cli-config.yaml.example b/cli-config.yaml.example index 575a03e5c5..c4025dc2b3 100644 --- a/cli-config.yaml.example +++ b/cli-config.yaml.example @@ -1055,6 +1055,16 @@ agent: # the turn (see gateway_timeout). 0 = disable. Default 300. # session_stall_timeout: 300 + # Transcript-sanitiser repeated-heal escalation (#96870). The pre-send + # sanitiser silently repairs empty non-final turns on the wire copy; after + # this many heal passes within a 10-minute session window Hermes logs one + # ERROR (with session id + heal pattern) and sends a ONE-TIME out-of-band + # notice to the user suggesting /debug share or `hermes doctor`. The notice + # goes through the status channel only — conversation context and prompt + # caching are untouched. Set 0 to disable escalation (per-window WARNINGs + # still fire). Default 3. + # sanitizer_heal_escalation_threshold: 3 + # Related in-agent compression timeouts (they live under the top-level # compression: block, shown here for discoverability next to the stall # watchdog they complement — a hung compression is a common stall cause): diff --git a/hermes_cli/config_defaults.py b/hermes_cli/config_defaults.py index 7990d88ad2..8c1733d551 100644 --- a/hermes_cli/config_defaults.py +++ b/hermes_cli/config_defaults.py @@ -297,6 +297,14 @@ DEFAULT_CONFIG = { # from gateway_timeout (which kills the turn) and # gateway_notify_interval ("still working" heartbeats). 0 = disable. "session_stall_timeout": 300, + # Transcript-sanitiser repeated-heal escalation threshold (#96870). + # After this many pre-send heal passes within a 10-minute session + # window, log one ERROR (session id + heal pattern) and queue a + # ONE-TIME out-of-band user notice pointing at /debug share or + # `hermes doctor`. Delivered via the status channel only — + # conversation context / prompt caching untouched. 0 = disable + # escalation (per-window WARNINGs still fire). + "sanitizer_heal_escalation_threshold": 3, # Long-lived reconnect-loop escalation (seconds). A platform that has # been continuously failing/reconnecting for this long gets # needs_attention flagged in gateway runtime status (visible in diff --git a/hermes_cli/debug.py b/hermes_cli/debug.py index f633df62f2..947562ecff 100644 --- a/hermes_cli/debug.py +++ b/hermes_cli/debug.py @@ -607,6 +607,26 @@ def collect_debug_report( if log_snapshots is None: log_snapshots = _capture_default_log_snapshots(log_lines) + # ── Sanitiser heal counters (#96870) ───────────────────────────────── + # In-process, in-memory counters: populated when this report is built + # inside a process that ran agent turns (gateway /debug share); empty + # from a fresh CLI process, where the errors.log tail below carries the + # same escalation lines instead. + try: + from agent.agent_runtime_helpers import get_sanitizer_heal_stats + + heal_stats = get_sanitizer_heal_stats() + if heal_stats: + buf.write("\n\n--- transcript sanitiser heal counters ---\n") + for sess, st in sorted(heal_stats.items()): + buf.write( + f"session {sess}: {st['heal_events']} heal events, " + f"{st['messages_healed']} messages healed, " + f"escalated={st['escalated']}\n" + ) + except Exception: + pass + # ── Recent log tails (summary only) ────────────────────────────────── buf.write("\n\n") buf.write(f"--- agent.log (last {log_lines} lines) ---\n") diff --git a/hermes_cli/dump.py b/hermes_cli/dump.py index c00d55f12f..c7399f39f8 100644 --- a/hermes_cli/dump.py +++ b/hermes_cli/dump.py @@ -237,6 +237,7 @@ def _config_overrides(config: dict) -> dict[str, str]: ("agent", "max_turns"), ("agent", "gateway_timeout"), ("agent", "session_stall_timeout"), + ("agent", "sanitizer_heal_escalation_threshold"), ("agent", "tool_use_enforcement"), ("agent", "execution_guidance"), ("terminal", "backend"), diff --git a/tests/run_agent/test_sanitiser_escalation.py b/tests/run_agent/test_sanitiser_escalation.py index 2aab5e7c21..d75be6ece2 100644 --- a/tests/run_agent/test_sanitiser_escalation.py +++ b/tests/run_agent/test_sanitiser_escalation.py @@ -16,7 +16,11 @@ import pytest from agent.agent_runtime_helpers import ( _INTERRUPTED_PLACEHOLDER, _empty_heal_log_state, + _empty_heal_pending_notice, + _empty_heal_user_notified, + consume_pending_sanitizer_heal_notice, fill_empty_non_final_wire_payload, + get_sanitizer_heal_stats, repair_empty_non_final_messages, ) from hermes_logging import clear_session_context, set_session_context @@ -25,9 +29,13 @@ from hermes_logging import clear_session_context, set_session_context @pytest.fixture(autouse=True) def _reset_heal_log(): _empty_heal_log_state.clear() + _empty_heal_pending_notice.clear() + _empty_heal_user_notified.clear() clear_session_context() yield _empty_heal_log_state.clear() + _empty_heal_pending_notice.clear() + _empty_heal_user_notified.clear() clear_session_context() @@ -83,7 +91,7 @@ class TestHealLogEscalation: def test_warning_then_one_error_then_silence(self, monkeypatch, caplog): import agent.agent_runtime_helpers as arh - monkeypatch.setattr(arh, "_EMPTY_HEAL_ESCALATE_AFTER", 3) + monkeypatch.setattr(arh, "_heal_escalation_threshold", lambda: 3) set_session_context("sess-heal") durable = _poisoned_rows() @@ -106,7 +114,7 @@ class TestHealLogEscalation: def test_sessions_do_not_share_heal_counters(self, monkeypatch, caplog): import agent.agent_runtime_helpers as arh - monkeypatch.setattr(arh, "_EMPTY_HEAL_ESCALATE_AFTER", 3) + monkeypatch.setattr(arh, "_heal_escalation_threshold", lambda: 3) with caplog.at_level(logging.WARNING, logger="run_agent"): set_session_context("sess-a") repair_empty_non_final_messages([dict(m) for m in _poisoned_rows()]) @@ -126,6 +134,132 @@ class TestHealLogEscalation: assert out is not durable +class TestOneTimeUserNotice: + def _heal_n(self, n): + for _ in range(n): + repair_empty_non_final_messages( + [dict(m) for m in _poisoned_rows()] + ) + + def test_notice_queued_once_at_threshold(self, monkeypatch): + import agent.agent_runtime_helpers as arh + + monkeypatch.setattr(arh, "_heal_escalation_threshold", lambda: 3) + set_session_context("sess-notice") + + self._heal_n(2) + assert consume_pending_sanitizer_heal_notice() is None + + self._heal_n(1) # crosses threshold + notice = consume_pending_sanitizer_heal_notice() + assert notice is not None + assert "repeated repair" in notice + assert "/debug share" in notice + assert "hermes doctor" in notice + + # drained: never delivered twice + assert consume_pending_sanitizer_heal_notice() is None + + def test_notice_never_rearms_in_new_window(self, monkeypatch): + import agent.agent_runtime_helpers as arh + + monkeypatch.setattr(arh, "_heal_escalation_threshold", lambda: 2) + set_session_context("sess-rearm") + + self._heal_n(2) + assert consume_pending_sanitizer_heal_notice() is not None + + # Simulate window expiry: force a fresh window, cross threshold again. + arh._empty_heal_log_state["sess-rearm"]["window_start"] -= ( + arh._EMPTY_HEAL_WINDOW_S + 1 + ) + self._heal_n(2) + assert consume_pending_sanitizer_heal_notice() is None + + def test_notice_scoped_to_session(self, monkeypatch): + import agent.agent_runtime_helpers as arh + + monkeypatch.setattr(arh, "_heal_escalation_threshold", lambda: 2) + set_session_context("sess-x") + self._heal_n(2) + set_session_context("sess-y") + # sess-y never escalated; its consume must not steal sess-x's notice + assert consume_pending_sanitizer_heal_notice() is None + set_session_context("sess-x") + assert consume_pending_sanitizer_heal_notice() is not None + + def test_threshold_zero_disables_escalation(self, monkeypatch, caplog): + import agent.agent_runtime_helpers as arh + + monkeypatch.setattr(arh, "_heal_escalation_threshold", lambda: 0) + set_session_context("sess-off") + with caplog.at_level(logging.WARNING, logger="run_agent"): + self._heal_n(6) + assert consume_pending_sanitizer_heal_notice() is None + assert not [r for r in caplog.records if r.levelno == logging.ERROR] + + def test_threshold_read_from_config(self, monkeypatch): + import agent.agent_runtime_helpers as arh + + monkeypatch.setattr( + "hermes_cli.config.load_config_readonly", + lambda: {"agent": {"sanitizer_heal_escalation_threshold": 7}}, + ) + assert arh._heal_escalation_threshold() == 7 + + def test_threshold_defaults_when_config_unreadable(self, monkeypatch): + import agent.agent_runtime_helpers as arh + + def _boom(): + raise RuntimeError("no config") + + monkeypatch.setattr("hermes_cli.config.load_config_readonly", _boom) + assert ( + arh._heal_escalation_threshold() == arh._EMPTY_HEAL_ESCALATE_AFTER + ) + + +class TestHealStatsSurface: + def test_counters_visible_and_escalation_flagged(self, monkeypatch): + import agent.agent_runtime_helpers as arh + + monkeypatch.setattr(arh, "_heal_escalation_threshold", lambda: 3) + set_session_context("sess-stats") + for _ in range(4): + repair_empty_non_final_messages( + [dict(m) for m in _poisoned_rows()] + ) + + stats = get_sanitizer_heal_stats() + assert stats["sess-stats"]["heal_events"] == 4 + assert stats["sess-stats"]["messages_healed"] == 4 + assert stats["sess-stats"]["escalated"] is True + + def test_debug_report_includes_heal_counters(self, monkeypatch): + import agent.agent_runtime_helpers as arh + from hermes_cli.debug import collect_debug_report, LogSnapshot + + monkeypatch.setattr(arh, "_heal_escalation_threshold", lambda: 2) + set_session_context("sess-report") + for _ in range(2): + repair_empty_non_final_messages( + [dict(m) for m in _poisoned_rows()] + ) + + empty = LogSnapshot(path=None, tail_text="", full_text="") + report = collect_debug_report( + log_lines=5, + dump_text="dump", + log_snapshots={ + k: empty + for k in ("agent", "errors", "gateway", "gui", "desktop") + }, + ) + assert "transcript sanitiser heal counters" in report + assert "sess-report: 2 heal events" in report + assert "escalated=True" in report + + class TestProjectionStopsReheal: def _loop_agent(self): from unittest.mock import MagicMock, patch @@ -209,3 +343,53 @@ class TestProjectionStopsReheal: assert wire_assistants[0]["content"] == _INTERRUPTED_PLACEHOLDER assert history[1]["content"] == "" assert "api_content" not in history[1] + + def test_pending_notice_delivered_out_of_band_not_in_context(self): + """A queued escalation notice is emitted through _emit_warning (the + status/delivery channel) and NEVER appears in the wire messages or + the durable history — message-flow / caching invariants (#96870).""" + from unittest.mock import patch + + import agent.agent_runtime_helpers as _arh + from tests.run_agent.test_run_agent import _mock_response + + agent = self._loop_agent() + agent.client.chat.completions.create.side_effect = [ + _mock_response(content="ok", finish_reason="stop"), + ] + + set_session_context("sess-loop-notice") + # The turn re-binds the log session context to the agent's own + # session id at turn start, so queue the notice under that key — + # exactly where the escalation path would have put it mid-session. + _live_key = str(getattr(agent, "session_id", None) or "-") + with _arh._empty_heal_log_lock: + _arh._empty_heal_user_notified.add(_live_key) + _arh._empty_heal_pending_notice[_live_key] = ( + "⚠️ Your session transcript required repeated repair — " + "run /debug share or `hermes doctor`." + ) + + warned = [] + history = [ + {"role": "user", "content": "start"}, + {"role": "assistant", "content": "earlier", "finish_reason": "stop"}, + ] + with ( + patch.object(agent, "_flush_messages_to_session_db"), + patch.object(agent, "_persist_session"), + patch.object(agent, "_save_trajectory"), + patch.object(agent, "_cleanup_task_resources"), + patch.object(agent, "_emit_warning", side_effect=warned.append), + ): + agent.run_conversation("next", conversation_history=history) + + assert warned and "repeated repair" in warned[0] + # one-time: drained after delivery + assert _arh._empty_heal_pending_notice == {} + # never injected into the wire copy or durable history + wire = agent.client.chat.completions.create.call_args.kwargs["messages"] + assert all("repeated repair" not in str(m.get("content")) for m in wire) + assert all( + "repeated repair" not in str(m.get("content")) for m in history + )