fix(compression): repeated summary stall escalates to the deterministic fallback summary

A first stalled summary stream keeps today's behaviour: the transcript is left
alone, the stall-class cooldown (floored at the idle window) is armed and the
LLM route retries after it lapses. When the route stalls AGAIN while a
stall-class failure is still on the ladder, the stall retry ladder now ends
with a deterministic rung: the worker is re-run with the summary LLM skipped
(DETERMINISTIC_SUMMARY_ROUTE pin, consumed in _summarize_window) and
compress() commits its static fallback summary through the ordinary
lease/fence/watermark pipeline — the same degrade a failed summary call gets
(abort_on_summary_failure still aborts).

WHY: after "made no progress … continuing without compression" the context
stays oversized, so the next turn after the cooldown re-enters the same silent
stream and burns another full idle window; the reporter saw this every ~2 min
for hours (#112420). A route that has proven unhealthy twice must degrade once
instead of looping. The prune-on-stall hunk from #112504 was declined because
it committed outside the lease/fence; this rung reuses the same-turn fallback
worker (bypass_cooldown) and the commit path the fallback_chain retry already
uses, so no new commit surface is introduced.

Also: a pinned fallback_chain route whose summary call FAILS still commits the
static fallback summary (default abort_on_summary_failure=false); the host log
said "recovered on fallback_chain[0]" for that. It now logs "committed a
deterministic fallback summary on …" at WARNING (#112387 review caveat), keyed
on the post-commit fallback_compression_streak bump.

Docs: developer-guide failure-cooldown section + agent/AGENTS.md.
This commit is contained in:
teknium1
2026-09-16 22:08:12 -07:00
committed by Teknium
parent ba2bcfa6e7
commit 07c92d675a
6 changed files with 292 additions and 12 deletions

View File

@@ -82,7 +82,10 @@ Two layers: gateway session hygiene (85% threshold) and the agent `ContextCompre
configurable; per-model overrides; failure cooldown after provider-proven overflow). The algorithm
prunes old tool results first (no LLM call), then picks boundaries, then generates a structured
summary with the `auxiliary` compression model. In-place compaction keeps a single stable session
id; native Responses/Codex compaction paths are provider-specific. Compression is the sanctioned
id; native Responses/Codex compaction paths are provider-specific. A stalled summary stream retries
once on `auxiliary.compression.fallback_chain`, and a repeated stall (a stall-class failure already on
the cooldown ladder) ends with the deterministic fallback summary through the same pipeline — never a
prune committed outside the lease/fence. Compression is the sanctioned
cache break — keep it the only one. Full detail:
`website/docs/developer-guide/context-compression-and-caching.md`.

View File

@@ -89,6 +89,22 @@ def take_pinned_summary_route() -> Optional[Dict[str, Any]]:
return route
# Pinned route that names NO summary model: compress() skips the summary LLM and inserts its deterministic
# fallback summary instead (``abort_on_summary_failure`` still aborts). The host pins it when the summary
# route stalls again after a stall-class backoff already burned one idle window (#112420), so a provably
# unhealthy route degrades once instead of re-entering the same silent stream every turn.
DETERMINISTIC_SUMMARY_ROUTE: Dict[str, Any] = {"label": "deterministic fallback summary", "deterministic": True}
def take_deterministic_summary_pin() -> bool:
"""Consume the pin when it is the deterministic sentinel; a real route (or no pin) is left in place."""
route = _SUMMARY_ROUTE_PIN.get()
if not (isinstance(route, dict) and route.get("deterministic") is True):
return False
_SUMMARY_ROUTE_PIN.set(None)
return True
def _pinned_summary_call_kwargs() -> Dict[str, Any]:
"""Consume the pinned route as explicit ``call_llm`` keyword arguments."""
route = take_pinned_summary_route() or {}

View File

@@ -31,7 +31,19 @@ class SummaryDispatchMixin:
self, messages: List[Dict[str, Any]], turns_to_summarize: List[Dict[str, Any]], scan: "_HandoffScan",
focus_topic: Optional[str], memory_context: str, bypass_cooldown: bool,
) -> Optional[str]:
"""Run the summary LLM; a cancellation rolls back the handoff scan's self-heal mutation first."""
"""Run the summary LLM; a cancellation rolls back the handoff scan's self-heal mutation first.
A deterministic pin (repeated stall, #112420) skips the LLM: ``None`` lets Phase 3 insert the static
fallback summary, or abort under ``abort_on_summary_failure`` exactly like a failed summary call."""
from agent.context_compressor import take_deterministic_summary_pin
if take_deterministic_summary_pin():
# Surfaces through the fallback summary's reason line and the host's one-shot user warning.
self._last_summary_error = (
"summary model stalled again after a stall backoff; deterministic fallback summary inserted"
)
telemetry = getattr(self, "_active_compression_telemetry", None)
if isinstance(telemetry, dict):
telemetry["failure_class"] = "stall_deterministic_fallback"
return None
# Focus-topic derivation scans user turns; only pay when a summary is generated.
summary_kwargs: Dict[str, Any] = {
"focus_topic": focus_topic or self._derive_auto_focus_topic(messages),

View File

@@ -842,14 +842,32 @@ def resolve_compression_fallback_route() -> Optional[dict]:
return None
def _stall_retry_routes(escalate_deterministic: bool) -> list:
"""Pinned routes for the stall retry, in order: the configured chain entry, then (only once a
stall-class backoff has already burned a window in this session) the deterministic fallback summary."""
routes = [route for route in (resolve_compression_fallback_route(),) if route is not None]
if escalate_deterministic:
from agent.context_compressor import DETERMINISTIC_SUMMARY_ROUTE
routes.append(dict(DETERMINISTIC_SUMMARY_ROUTE))
return routes
def _prior_timeout_failures(agent: Any) -> int:
"""Timeout-class failures this session that no healthy summary has cleared yet (type-pinned)."""
count = getattr(getattr(agent, "context_compressor", None), "_consecutive_timeout_failures", 0)
return count if isinstance(count, int) and not isinstance(count, bool) else 0
def _retry_compression_on_fallback_chain(
*, worker: Callable[[CompressionCommitFence], Tuple[list, str]], messages: list,
system_prompt_fallback: Any, idle_timeout_seconds: float, total_ceiling_seconds: float,
on_commit_overrun: Optional[Callable[[float, float], None]] = None,
on_timeout_cause: Optional[Callable[[bool, bool], None]] = None, telemetry_agent: Any = None,
new_fence: Optional[Callable[[], CompressionCommitFence]] = None,
new_fence: Optional[Callable[[], CompressionCommitFence]] = None, escalate_deterministic: bool = False,
) -> Optional[Tuple[list, str]]:
"""Re-run an aborted compression once with the summary route pinned.
"""Re-run an aborted compression with the summary route pinned: once on the configured chain entry,
then — when ``escalate_deterministic`` (a stall backoff already burned one idle window this session,
#112420) — once with the summary LLM skipped so compress() commits its deterministic fallback summary.
Returns ``(messages, system_prompt)`` on real compression, else ``None`` and the caller degrades as
before. The entry's ``timeout`` sets the idle window. Re-runs the whole worker, so pre-compression
callbacks must be idempotent.
@@ -868,10 +886,25 @@ def _retry_compression_on_fallback_chain(
hard_cancel = getattr(telemetry_agent, "_hard_interrupt_requested", None)
if callable(getattr(hard_cancel, "is_set", None)) and hard_cancel.is_set():
return None
route = resolve_compression_fallback_route()
if route is None:
return None
for route in _stall_retry_routes(escalate_deterministic):
recovered = _run_pinned_compression_retry(
route, worker=worker, messages=messages, system_prompt_fallback=system_prompt_fallback,
idle_timeout_seconds=idle_timeout_seconds, total_ceiling_seconds=total_ceiling_seconds,
on_commit_overrun=on_commit_overrun, on_timeout_cause=on_timeout_cause, telemetry_agent=telemetry_agent,
new_fence=new_fence,
)
if recovered is not None:
return recovered
return None
def _run_pinned_compression_retry(
route: dict, *, worker: Callable[[CompressionCommitFence], Tuple[list, str]], messages: list,
system_prompt_fallback: Any, idle_timeout_seconds: float, total_ceiling_seconds: float,
on_commit_overrun: Optional[Callable[[float, float], None]], on_timeout_cause: Optional[Callable[[bool, bool], None]],
telemetry_agent: Any, new_fence: Optional[Callable[[], CompressionCommitFence]],
) -> Optional[Tuple[list, str]]:
"""One bounded re-run of ``worker`` with ``route`` pinned; ``None`` when it produced no compression."""
# The aborted fence refuses all commits; mint a fresh one via the host factory
# so a /stop during the retry serializes against THIS attempt's commit boundary.
retry_fence = None
@@ -892,10 +925,19 @@ def _retry_compression_on_fallback_chain(
retry_fence = CompressionCommitFence()
idle = float(route.get("timeout") or idle_timeout_seconds)
ceiling = max(float(total_ceiling_seconds), idle)
logger.warning(
"Context compression stalled on the configured summary route — "
"retrying once on %s (%s) before continuing without compression", route["label"], route["model"],
)
deterministic = route.get("deterministic") is True
if deterministic:
logger.warning(
"Context compression stalled again after a stall backoff — committing the %s (no summary model) "
"before continuing without compression", route["label"],
)
else:
logger.warning(
"Context compression stalled on the configured summary route — "
"retrying once on %s (%s) before continuing without compression", route["label"], route["model"],
)
compressor = getattr(telemetry_agent, "context_compressor", None)
streak_before = getattr(compressor, "_fallback_compression_streak", 0)
try:
from agent.context_compressor import pin_summary_route
with pin_summary_route(route):
@@ -917,7 +959,16 @@ def _retry_compression_on_fallback_chain(
route["label"],
)
return None
logger.info("Context compression recovered on %s after the primary summary route stalled", route["label"])
# A pinned summary call that failed still commits (static fallback summary under the default
# abort_on_summary_failure=false); the streak bump is the post-commit tell. Never call that "recovered".
streak_after = getattr(compressor, "_fallback_compression_streak", 0)
if deterministic or (isinstance(streak_after, int) and isinstance(streak_before, int) and streak_after > streak_before):
logger.warning(
"Context compression committed a deterministic fallback summary on %s after the primary summary route "
"stalled (no summary model produced output)", route["label"],
)
else:
logger.info("Context compression recovered on %s after the primary summary route stalled", route["label"])
return result_msgs, result_prompt
@@ -1046,6 +1097,10 @@ def run_compress_context_with_progress_timeout(
idle = float(idle_timeout_seconds)
fence = fence if fence is not None else CompressionCommitFence()
fence.set_total_ceiling_seconds(ceiling)
# Read BEFORE this attempt runs: the host's ``stalled`` record and the cancelled worker's
# ``stall_interrupted`` record both land during the unwind below, and this stall must not count as
# its own prior. One prior timeout-class failure = the route already burned a full idle window.
escalate_deterministic = stall_fallback and _prior_timeout_failures(telemetry_agent) >= 1
# Sync mirror of gateway hygiene's run_in_executor + wait_for loop: offload,
# poll idle budget + ceiling, fence-cancel on timeout so no late commit lands.
from tools.thread_context import propagate_context_to_thread
@@ -1143,6 +1198,7 @@ def run_compress_context_with_progress_timeout(
worker=fallback_worker or worker, messages=messages, system_prompt_fallback=system_prompt_fallback,
idle_timeout_seconds=idle, total_ceiling_seconds=ceiling, on_commit_overrun=on_commit_overrun,
on_timeout_cause=on_timeout_cause, telemetry_agent=telemetry_agent, new_fence=new_fence,
escalate_deterministic=escalate_deterministic,
)
if recovered is not None:
return recovered

View File

@@ -0,0 +1,176 @@
"""A repeated summary stall escalates to the deterministic fallback summary — #112420.
After one stall the host records a stall-class backoff (``_consecutive_timeout_failures`` >= 1). When the
next attempt stalls again, "continuing without compression" would re-enter the same silent route every
turn, so the stall retry ladder ends with a deterministic rung: the worker is re-run with the summary LLM
skipped and compress() commits its static fallback summary through the ordinary lease/fence pipeline.
A first stall keeps today's behaviour (backoff, LLM retry after it lapses). A pinned fallback route whose
summary call fails still commits under the default ``abort_on_summary_failure=false`` — that must not be
logged as a recovery (#112387 review caveat).
"""
from __future__ import annotations
import logging
import os
import time
from pathlib import Path
from unittest.mock import patch
import pytest
import agent.conversation_compression as cc
from agent.auxiliary_client import AuxiliaryExplicitCancellation
from agent.context_compressor import SUMMARY_PREFIX, pin_summary_route
from agent.conversation_compression import CompressionCommitFence, run_compress_context_with_progress_timeout
from hermes_state import SessionDB
CHAIN_ENTRY = {
"provider": "custom", "model": "backup-summarizer", "base_url": "https://fallback.invalid/v1",
"api_key": "sk-fallback", "timeout": 0.4,
}
def _make_agent(tmp_path, tag):
db = SessionDB(db_path=Path(tmp_path) / f"state-{tag}.db")
session_id = f"STALL_DETERMINISTIC_{tag}"
db.create_session(session_id, source="cli")
with patch.dict(os.environ, {"OPENROUTER_API_KEY": "test-key"}):
from run_agent import AIAgent
agent = AIAgent(
api_key="test-key", base_url="https://openrouter.ai/api/v1", model="test/model", quiet_mode=True,
session_db=db, session_id=session_id, skip_context_files=True, skip_memory=True,
)
agent._compression_feasibility_checked = True
agent.compression_in_place = True
agent._cached_system_prompt = "sys"
agent.context_compressor.threshold_tokens = 1_000
return agent
def _transcript():
return [
{"role": "user" if i % 2 == 0 else "assistant", "content": f"m{i} " + ("lorem ipsum " * 300)}
for i in range(40)
]
def _stalling_call_llm(compressor, calls, *, fail_when_pinned=False):
"""Summary call that never streams: hangs until the host cancels the fence, like a held-open socket."""
def _call(**kwargs):
calls.append(kwargs.get("provider") or "primary")
if fail_when_pinned and "provider" in kwargs:
raise RuntimeError("fallback route exploded")
cancelled = getattr(compressor, "_compression_cancelled_check", None)
deadline = time.monotonic() + 6
while time.monotonic() < deadline and not (callable(cancelled) and cancelled()):
time.sleep(0.001)
raise AuxiliaryExplicitCancellation()
return _call
def _summary_rows(messages):
return [m for m in messages if isinstance(m.get("content"), str) and m["content"].startswith(SUMMARY_PREFIX)]
@pytest.fixture
def fast_timeouts(monkeypatch):
monkeypatch.setattr(cc, "resolve_context_compression_timeouts", lambda compression_cfg=None: (0.4, 4.0))
def test_second_consecutive_stall_commits_the_deterministic_fallback_summary(tmp_path, fast_timeouts):
"""Real AIAgent + real ``compress_context`` + facade timeout wrap; only ``call_llm`` is stubbed.
Stall 1: backoff recorded, transcript untouched. Stall 2 (after the backoff lapsed): the deterministic
rung commits a fallback summary instead of continuing without compression."""
agent = _make_agent(tmp_path, "A")
compressor = agent.context_compressor
calls = []
live = _transcript()
with patch("agent.context_compressor.call_llm", side_effect=_stalling_call_llm(compressor, calls)), \
patch("agent.auxiliary_client._get_auxiliary_task_config", return_value={"fallback_chain": []}):
first, _ = agent._compress_context(live, "sys", approx_tokens=50_000)
assert first is live, "a first stall keeps the transcript and only arms the backoff"
assert compressor._consecutive_timeout_failures >= 1
assert compressor._summary_failure_cooldown_until > time.monotonic()
assert calls == ["primary"], "no deterministic rung on the FIRST stall: the LLM route gets its backoff retry"
# The backoff lapses; the still-oversized context re-triggers compression (the reporter's next turn).
compressor._summary_failure_cooldown_until = 0.0
compressor._session_db.clear_compression_failure_cooldown(compressor._session_id)
calls.clear()
second, _ = agent._compress_context(live, "sys", approx_tokens=50_000)
assert second is not live and len(second) < len(live)
assert len(_summary_rows(second)) == 1, "the deterministic fallback summary is committed as the handoff"
assert calls == ["primary"], "the deterministic rung makes no summary LLM call"
assert getattr(agent, "_last_compression_timed_out", None) is not True
# The stalled LLM route stays in its durable backoff even though the deterministic rung committed; the
# cancelled primary worker persists it while unwinding, so poll briefly for its row.
deadline = time.monotonic() + 3
remaining = 0.0
while time.monotonic() < deadline and remaining <= 0:
row = agent._session_db.get_compression_failure_cooldown(compressor._session_id) or {}
remaining = float(row.get("remaining_seconds") or 0.0)
time.sleep(0.01)
assert remaining > 0, "the stalled LLM route keeps its stall backoff after the deterministic commit"
def test_failing_pinned_fallback_route_is_not_logged_as_recovered(tmp_path, fast_timeouts, caplog):
"""Primary stalls, the fallback_chain route raises: compress() still commits its static fallback
summary (abort_on_summary_failure=false), and the host log must say so instead of 'recovered'."""
agent = _make_agent(tmp_path, "B")
compressor = agent.context_compressor
calls = []
live = _transcript()
caplog.set_level(logging.INFO, logger="agent.conversation_compression")
with patch(
"agent.context_compressor.call_llm", side_effect=_stalling_call_llm(compressor, calls, fail_when_pinned=True),
), patch("agent.auxiliary_client._get_auxiliary_task_config", return_value={"fallback_chain": [CHAIN_ENTRY]}):
out, _ = agent._compress_context(live, "sys", approx_tokens=50_000)
assert calls == ["primary", "custom"]
assert out is not live and len(_summary_rows(out)) == 1
records = [r for r in caplog.records if r.name == "agent.conversation_compression"]
assert not any("recovered on fallback_chain[0]" in r.getMessage() for r in records)
assert any(
"committed a deterministic fallback summary on fallback_chain[0]" in r.getMessage()
and r.levelno == logging.WARNING
for r in records
)
def test_deterministic_pin_is_consumed_and_a_real_route_is_left_alone():
from agent.context_compressor import (
DETERMINISTIC_SUMMARY_ROUTE, take_deterministic_summary_pin, take_pinned_summary_route,
)
with pin_summary_route(dict(CHAIN_ENTRY)):
assert take_deterministic_summary_pin() is False
assert take_pinned_summary_route()["model"] == "backup-summarizer"
with pin_summary_route(dict(DETERMINISTIC_SUMMARY_ROUTE)):
assert take_deterministic_summary_pin() is True
assert take_deterministic_summary_pin() is False, "single use"
assert take_deterministic_summary_pin() is False
def test_fence_level_retry_ladder_is_unchanged_without_a_prior_timeout():
"""No agent-level timeout history (fence-level callers, fresh compressors): a stall with no chain still
degrades in one attempt — the deterministic rung is escalation, not the default."""
attempts = []
def worker(fence: CompressionCommitFence):
attempts.append(fence)
time.sleep(0.3)
return ([{"role": "assistant", "content": "late"}], "late-prompt")
original = [{"role": "user", "content": "keep-me"}]
with patch("agent.auxiliary_client._get_auxiliary_task_config", return_value={"fallback_chain": []}):
msgs, prompt = run_compress_context_with_progress_timeout(
worker=worker, messages=original, system_prompt_fallback="degraded-prompt",
idle_timeout_seconds=0.05, total_ceiling_seconds=2.0, on_timeout=lambda *a: None,
)
assert msgs is original and prompt == "degraded-prompt"
assert len(attempts) == 1

View File

@@ -151,6 +151,23 @@ backend does not re-fire every turn. Three paths run a real attempt anyway:
- Manual `/compress` (`force=True`) — clears the cooldown and retries.
- The same-turn `fallback_chain` retry after a stalled primary route — the
cancelled primary's own stall cooldown must not suppress it (`bypass_cooldown`).
If that pinned route's summary call fails, compress() still commits its
deterministic fallback summary (default `abort_on_summary_failure: false`);
the log then says "committed a deterministic fallback summary", not
"recovered".
- **Repeated stall → deterministic fallback.** A first stall keeps the
transcript, arms the cooldown and lets the LLM route retry after it lapses.
When the route stalls *again* while a stall-class failure is still on the
ladder (`_consecutive_timeout_failures >= 1`), the retry ladder ends with a
deterministic rung: the worker is re-run with the summary LLM skipped
(`DETERMINISTIC_SUMMARY_ROUTE` pin) and commits the static fallback summary
through the ordinary lease/fence/watermark pipeline — the same degrade a
failed summary call gets — instead of "continuing without compression" and
re-entering the same silent stream every turn (#112420).
`abort_on_summary_failure: true` still aborts (nothing dropped). A committed
compaction rebinds the compressor and resets the ladder count, so each
compaction cycle grants the LLM route one stall before escalating; the
persisted cooldown row still paces attempts across turns and restarts.
- **Provider-proven overflow** — when the provider itself rejects the request
with a context-length error, the recovery pass ignores the cooldown for one
bounded attempt (`max_compression_attempts`) without clearing it. Deferring