feat(agent): escalate repeated transcript-sanitiser heals with a one-time user notice (#96870)

Builds the escalation layer on top of HexLab98's heal-log windowing
(salvaged from PR #96916):

- Per-session heal counters (heal events + messages healed) tracked by the
  repair path in agent_runtime_helpers.py, session totals preserved across
  10-minute log windows.
- Threshold escalation: after N heals in a session window (default 3,
  configurable via agent.sanitizer_heal_escalation_threshold in
  config.yaml, 0 = off) log ONE ERROR carrying session id + heal pattern
  (events/messages/window/threshold), then stay quiet.
- ONE-TIME out-of-band user notice queued at the threshold and delivered by
  the conversation loop through _emit_warning (status callback -> gateway
  status message / CLI print). Never injected into conversation context or
  the wire copy: prompt caching, role alternation, and durable history are
  untouched. Never re-arms on a new window; scoped per session.
- Counters visible in diagnostics: get_sanitizer_heal_stats() rendered in
  the /debug share // hermes debug report, and the config key surfaced in
  hermes dump overrides. errors.log carries the ERROR line for `hermes logs
  errors`.
This commit is contained in:
Teknium
2026-08-31 11:54:24 -07:00
parent a779f527fa
commit fb9b2c893f
7 changed files with 349 additions and 9 deletions
+106 -7
View File
@@ -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(
+18
View File
@@ -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
+10
View File
@@ -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):
+8
View File
@@ -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
+20
View File
@@ -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")
+1
View File
@@ -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"),
+186 -2
View File
@@ -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
)