Files
hermes-agent/agent/compression_facade.py
beardthelion 09b72bc6d2 fix(compression): track commit fences as a registration stack
_compress_context published the active commit fence with a
save/restore cell: registration order was serialized by the fence
lock, but completion order is not. When attempt B registered over A
and A finished first, A's finally popped the slot, deleting B's live
fence mid-attempt (hard_interrupt lost the handle serializing cancel
admission against B's begin_commit). B's finally then republished A's
dead fence, which lingered until the next compression. The same
clobber existed in _publish_new_fence, which overwrote the slot
unconditionally when minting the stall-fallback retry fence.

Replace the cell with a stack of per-attempt registrations. The
finally removes only its own registration and republishes the newest
live entry (or clears the slot), so a dead fence can never be
restored over a live newer attempt. The stall-fallback retry swaps
its fence inside the owning registration and publishes only while
that attempt still holds the top registration. Registration moved
inside the try so an early exception cannot strand an entry.

(cherry picked from commit 574e9945cf186071c3da23c4bb517c0cbdaf73e6)
2026-09-20 15:27:11 -07:00

326 lines
18 KiB
Python

"""Host-side ``AIAgent._compress_context`` wrapper.
Publishes the commit fence ``hard_interrupt()`` reads, runs the compressor on a snapshot under the progress
timeout, mirrors ``_DB_PERSISTED_MARKER`` stamps back onto the live lists and rebinds the session context.
Extracted from ``run_agent.py``; every method resolves through ``AIAgent``'s MRO unchanged.
"""
import contextlib
import copy
import functools
import logging
import threading
from agent.session_activity import ActivityProvenance
# Same logger name as the origin module so log records / caplog filters are unchanged.
logger = logging.getLogger("run_agent")
class _CommitFenceRegistration:
"""One ``_compress_context`` attempt's entry on the agent's commit-fence stack.
The fence slot is a stack, not a save/restore cell: overlapping attempts register in
order but may complete in any order, so each attempt removes only its own entry and the
published slot always tracks the newest *live* attempt. The stall-fallback retry swaps
``fence`` in place so the slot update stays tied to the owning registration.
"""
__slots__ = ("fence",)
def __init__(self, fence) -> None:
self.fence = fence
def _timeout_fallback_prompt(agent, system_message: str) -> str:
"""Cached prompt, else a fresh build, else the raw ``system_message`` (never raises).
Resolved lazily by the timeout wrapper: an eager rebuild would raise before compress_context runs when
``_cached_system_prompt`` is unset and the builder fails."""
if cached := getattr(agent, "_cached_system_prompt", None):
return cached
try:
return agent._build_system_prompt(system_message)
except Exception:
logger.debug("compress_context timeout fallback prompt rebuild failed; using raw system_message", exc_info=True)
return system_message or ""
def _report_compression_timeout(
agent, *, idle: float, waited: float, since_progress: float, total_ceiling: float, total_exhausted: bool,
progress_observed: bool,
) -> None:
"""Host-side timeout bookkeeping: log, activity stamp, cooldown ladder, user warning."""
from agent.conversation_compression import mark_context_compression_timed_out
mark_context_compression_timed_out(agent)
if total_exhausted:
logger.warning(
"Context compression reached its total ceiling after %.1fs (progress observed=%s); continuing without compression",
waited, progress_observed,
)
else:
logger.warning(
"Context compression made no progress for %.1fs (total wait %.1fs, ceiling %.1fs); continuing without compression",
since_progress, waited, total_ceiling,
)
touch = getattr(agent, "_touch_activity", None)
if callable(touch):
try:
touch("context compression timed out", provenance=ActivityProvenance.AGENT_COMPRESSION_TIMEOUT)
except Exception:
logger.debug("compress_context timeout activity touch failed", exc_info=True)
# Same timeout cooldown ladder as summary-LLM timeouts: avoid re-burning the full idle budget every turn.
record = getattr(getattr(agent, "context_compressor", None), "record_timeout_failure", None)
if callable(record):
try:
if total_exhausted:
record("host compress_context total ceiling exhausted", failure_kind="ceiling_exhausted")
else:
record("host compress_context timeout (no summary progress)", failure_kind="stalled")
except Exception:
logger.debug("failed to record compress_context timeout cooldown", exc_info=True)
emit = getattr(agent, "_emit_warning", None)
if not callable(emit):
return
if total_exhausted:
progress = " after summary output was observed" if progress_observed else ""
emit(
"⚠ Context compression reached its total ceiling "
f"after {waited:.1f}s{progress}. No messages were "
"dropped — continuing without compression. Run /compress to retry or /new for a clean session."
)
else:
emit(
f"⚠ Context compression timed out after {idle:.1f}s with no output from the summary "
"model. No messages were dropped — continuing without compression. Run /compress to retry, /new "
"for a clean session, or check auxiliary.compression."
)
def _warn_commit_overrun(agent, waited: float, ceiling: float) -> None:
"""Commit-phase ceiling breach: the SessionDB mutation must complete, so only surface it."""
emit = getattr(agent, "_emit_warning", None)
if callable(emit):
emit(
f"⚠ Context compression commit is taking unusually long ({waited:.0f}s, ceiling {ceiling:.0f}s). "
"Waiting for it to finish safely — if this persists, check SessionDB health (disk / lock contention)."
)
def _sync_persisted_markers(target_messages, source_messages) -> None:
"""Mirror ``_DB_PERSISTED_MARKER`` stamps from the worker's snapshot onto a live list.
Matched by scoped identity; timestamp-less repeated content is ambiguous, so every scoped match is
stamped. Imported UNCONDITIONALLY: a silent fallback literal would split the stamping key from the flush's
and resurrect the duplicate-row bug."""
from agent.context_compressor import _DB_PERSISTED_MARKER
from agent.conversation_compression import _stamp_scoped_twins
if not isinstance(target_messages, list) or not isinstance(source_messages, list):
return
for source_message in source_messages:
if isinstance(source_message, dict) and source_message.get(_DB_PERSISTED_MARKER):
_stamp_scoped_twins(target_messages, source_message)
def _run_under_progress_timeout(
agent, run, messages, system_message, *, active_fence, registration, fence_registration_lock,
idle_timeout, total_ceiling, approx_tokens=None,
):
"""Run ``run(fence, target_messages=snapshot)`` on the pool under the progress-aware timeout.
The pooled worker must NEVER share the caller's live transcript — a late engine after a host timeout could
rewrite it. It deep-snapshots on the worker and publishes only via an ADMITTED commit; a no-op/abort
returns the snapshot unchanged, so the ORIGINAL list is handed back to keep identity semantics."""
from agent.conversation_compression import (
CompressionCommitFence, request_exceeds_model_window, run_compress_context_with_progress_timeout,
)
def _snapshot_worker(fence=None, *, same_turn_fallback_recovery=False):
# #76354 review F3: the pooled worker must NEVER share the caller's live transcript. Plugin/legacy
# context engines are allowed to mutate their input list in place; after a host timeout the worker
# stays alive, so a shared list would let a late engine rewrite the live conversation (roles,
# ordering, persisted content) behind the caller's back. Deep-snapshot here, on the worker thread,
# so the caller's list object is never touched by pooled code. Results are published to
# caller-visible state only via the returned value of an ADMITTED commit (the host discards results
# on timeout/cancel); durable SessionDB mutation is already gated behind the commit fence inside
# compress_context.
snapshot = copy.deepcopy(messages)
result_msgs, result_prompt = run(
fence, target_messages=snapshot, same_turn_fallback_recovery=same_turn_fallback_recovery
)
return (messages if result_msgs is snapshot else result_msgs), result_prompt
# The stall-fallback retry is the same recovery attempt as the stalled primary, but the cancelled primary
# worker records its stall_interrupted cooldown while unwinding — racing the retry's automatic gate
# (#112387). The retry therefore bypasses ONLY the summary-failure cooldown (never clears it; the
# structural breakers stay in force), exactly like provider-proven overflow recovery.
_same_turn_fallback_worker = functools.partial(_snapshot_worker, same_turn_fallback_recovery=True)
timeout_cause = {"total_exhausted": False, "progress_observed": False}
def _on_timeout_cause(total_exhausted, progress_observed):
timeout_cause.update(total_exhausted=total_exhausted, progress_observed=progress_observed)
def _on_timeout(idle, waited, since_progress):
_report_compression_timeout(
agent, idle=idle, waited=waited, since_progress=since_progress, total_ceiling=total_ceiling, **timeout_cause
)
def _publish_new_fence():
# The stall-fallback retry needs a fence the aborted attempt cannot veto; publish
# it on the slot hard_interrupt() reads, but only while this attempt still owns the
# top registration. A later attempt's live fence must not be clobbered.
retry_fence = CompressionCommitFence()
with fence_registration_lock:
registration.fence = retry_fence
stack = vars(agent).get("_compression_commit_fence_stack") or ()
if stack and stack[-1] is registration:
agent._active_compression_commit_fence = retry_fence
return retry_fence
return run_compress_context_with_progress_timeout(
worker=_snapshot_worker, messages=messages,
system_prompt_fallback=lambda: _timeout_fallback_prompt(agent, system_message),
idle_timeout_seconds=idle_timeout, total_ceiling_seconds=total_ceiling, on_timeout=_on_timeout,
on_timeout_cause=_on_timeout_cause,
on_commit_overrun=lambda waited, ceiling: _warn_commit_overrun(agent, waited, ceiling), fence=active_fence,
telemetry_agent=agent, new_fence=_publish_new_fence, fallback_worker=_same_turn_fallback_worker,
request_exceeds_window=request_exceeds_model_window(agent, approx_tokens) is True,
)
def _mirror_result_onto_live_lists(agent, result, messages, *, direct_path: bool) -> None:
"""Mirror persisted-marker stamps from the result list onto the live list(s)."""
if not (isinstance(result, tuple) and result and isinstance(result[0], list)):
return
result_messages = result[0]
# Direct-path callers bypass the snapshot worker but still need the post-publish mirror.
if direct_path or result_messages is not messages:
_sync_persisted_markers(messages, result_messages)
session_messages = getattr(agent, "_session_messages", None)
if isinstance(session_messages, list) and session_messages is not messages:
# Durable-parent adoption can leave `_session_messages` on the pre-adoption list.
_sync_persisted_markers(session_messages, result_messages)
def _rebind_caller_session_context(agent) -> None:
"""Propagate a rotated session id to the CALLER's thread/ContextVar (idempotent otherwise).
The worker thread rotated hermes_logging's thread-local id; post-compression tools must resolve
HERMES_SESSION_ID to the child id."""
with contextlib.suppress(Exception):
from hermes_logging import set_session_context
set_session_context(agent.session_id)
try:
from gateway.session_context import set_current_session_id
if agent.session_id:
set_current_session_id(agent.session_id)
except Exception:
logger.debug("post-compression session ContextVar rebind failed", exc_info=True)
class CompressionFacadeMixin:
"""``_compress_context`` (see module docstring)."""
def _compress_context(
self, messages: list, system_message: str, *, approx_tokens: int = None, task_id: str = "default",
focus_topic: str = None, force: bool = False, bypass_cooldown: bool = False,
defer_context_engine_notification: bool = False, commit_fence=None,
) -> tuple:
"""Forwarder — see ``agent.conversation_compression.compress_context``.
``force=True`` (manual /compress) bypasses the summary-failure cooldown; ``bypass_cooldown=True``
(provider-proven overflow recovery) runs one real attempt while the cooldown stays armed.
``force=True`` is passed by the manual ``/compress`` slash command so users can bypass the
summary-failure cooldown after an auto-compress abort. Auto-compress callers use the default
``force=False``. See #100661.
"""
# Per-attempt timeout signal for turn-start preflight and in-loop consumers: a stalled
# compression must not be mistaken for a structural no-op. Thread-local + per-agent lock.
# A stalled compression must not be mistaken for a structural no-op and followed by the oversized
# provider request it was meant to prevent. The typed helper upgrades the simple attribute to
# thread-local state guarded by a per-agent lock so overlapping automatic/manual entrypoints cannot
# clobber each other's outcome (#98741).
from agent.conversation_compression import (
CompressionCommitFence, compress_context, reset_context_compression_timeout_outcome,
resolve_context_compression_timeouts,
)
reset_context_compression_timeout_outcome(self)
from agent.portal_tags import (
get_affinity_scope, get_conversation_context, reset_affinity_scope, reset_conversation_context,
set_affinity_scope, set_conversation_context,
)
from agent.prompt_cache_scope import declared_conversation_scope_safe
# Out-of-turn compaction (/compact, gateway /compress, partial head compression) runs outside
# run_conversation's ambient scope; publish the root as a fallback so the summarizer's call carries
# the conversation tag. No-op for in-turn callers. Same for the ROUTING scope when declared.
token = None
if get_conversation_context() is None:
root = self._conversation_root_id()
if root:
token = set_conversation_context(root)
# Initialized alongside `token`: the turn-lease timeout/interrupt early returns leave the try block
# before set_affinity_scope() runs, and the finally reads this name unconditionally
# (UnboundLocalError otherwise — the 4 red cross-process lease tests on PR #97158).
affinity_token = None
if get_affinity_scope() is None:
declared = declared_conversation_scope_safe(self)
if declared:
affinity_token = set_affinity_scope(declared)
# Every compression has a fence; hard_interrupt() uses this exact instance to serialize cancel
# admission against begin_commit(). Publication is serialized so overlapping automatic/manual
# entrypoints cannot replace the fence of the attempt currently committing.
active_fence = commit_fence or CompressionCommitFence()
fence_registration_lock = vars(self).setdefault("_compression_commit_fence_lock", threading.RLock())
registration = _CommitFenceRegistration(active_fence)
try:
with fence_registration_lock:
fence_stack = vars(self).setdefault("_compression_commit_fence_stack", [])
fence_stack.append(registration)
self._active_compression_commit_fence = active_fence
def _run(fence=None, target_messages=None, same_turn_fallback_recovery=False):
return compress_context(
self, target_messages if target_messages is not None else messages, system_message,
approx_tokens=approx_tokens, task_id=task_id, focus_topic=focus_topic, force=force,
bypass_cooldown=bypass_cooldown or same_turn_fallback_recovery,
defer_context_engine_notification=(defer_context_engine_notification), commit_fence=fence,
)
# Callers that already own a progress-aware wait (gateway session
# hygiene) pass commit_fence and must not be double-wrapped.
direct_path = commit_fence is not None
if not direct_path:
idle_timeout, total_ceiling = resolve_context_compression_timeouts()
direct_path = idle_timeout <= 0
if direct_path:
result = _run(active_fence)
else:
result = _run_under_progress_timeout(
self, _run, messages, system_message,
active_fence=active_fence, registration=registration,
fence_registration_lock=fence_registration_lock,
idle_timeout=idle_timeout, total_ceiling=total_ceiling, approx_tokens=approx_tokens,
)
_mirror_result_onto_live_lists(self, result, messages, direct_path=direct_path)
_rebind_caller_session_context(self)
return result
finally:
# Remove only THIS attempt's registration. Completion order is not registration
# order: restoring a saved "previous" fence would resurrect a dead fence over a
# live newer attempt (and popping the slot outright would delete that attempt's
# fence mid-run). The slot always tracks the newest live registration.
with fence_registration_lock:
# Re-read rather than reuse the try-block local: an early exception before
# the append would otherwise raise NameError here and mask the real failure.
fence_stack = vars(self).get("_compression_commit_fence_stack") or ()
for idx in range(len(fence_stack) - 1, -1, -1):
if fence_stack[idx] is registration:
del fence_stack[idx]
break
if fence_stack:
self._active_compression_commit_fence = fence_stack[-1].fence
else:
vars(self).pop("_active_compression_commit_fence", None)
# Restore whatever the caller had, so a compaction never leaks its tag into the surrounding scope.
if token is not None:
reset_conversation_context(token)
if affinity_token is not None:
reset_affinity_scope(affinity_token)