compression: emit periodic client-visible status during compaction
Context compression can stream for minutes with no deltas, tool events, or status lines reaching remote transports. Idle-progress watchdogs on those clients treat the silence as a dead turn and interrupt it — the Android relay app's 180s turn watchdog fires session.interrupt, killing a healthy compression mid-flight and rolling back its work. On sessions near the context ceiling this loops forever: every new prompt retriggers preflight compression, which dies at exactly +180s again (observed telemetry: attempts aborted at 180469ms/180233ms/180219ms with failure_class=explicit_interrupt). Fix: the existing _CompressionActivityHeartbeat (which today only refreshes the SessionDB activity tracker) now also emits a 'compacting' status through agent.status_callback — once at compression start and on every heartbeat tick. The gateway already routes status_callback to status.update events, and clients already reset their watchdogs on any received event, so each heartbeat re-arms them. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> (cherry picked from commit 6d2e0e860d6fdb44ce549d1fa11a06662b1a4b0b)
This commit is contained in:
@@ -2337,6 +2337,7 @@ class _CompressionActivityHeartbeat:
|
||||
# if a prior timeout/cooldown stamp is still on the agent.
|
||||
self._suppressed = False
|
||||
self._touch("context compression started", allow_terminal_overwrite=True)
|
||||
self._emit_progress_status()
|
||||
self._thread.start()
|
||||
return self
|
||||
|
||||
@@ -2397,11 +2398,38 @@ class _CompressionActivityHeartbeat:
|
||||
except Exception:
|
||||
logger.debug("compression activity heartbeat touch failed", exc_info=True)
|
||||
|
||||
def _emit_progress_status(self) -> None:
|
||||
"""Publish a client-visible ``compacting`` status line.
|
||||
|
||||
Compression can stream for minutes with no deltas, tool events, or
|
||||
status lines reaching remote transports. Idle-progress watchdogs on
|
||||
those clients (e.g. the Android relay app's 180s turn watchdog)
|
||||
treat the silence as a dead turn and fire ``session.interrupt`` —
|
||||
killing a healthy compression mid-flight and rolling back its work,
|
||||
which retriggers on the next prompt and loops forever on sessions
|
||||
near the context ceiling. A periodic ``status_callback`` heartbeat
|
||||
gives every transport a progress event to re-arm on.
|
||||
"""
|
||||
status_callback = getattr(self._agent, "status_callback", None)
|
||||
if not status_callback:
|
||||
return
|
||||
try:
|
||||
status_callback(
|
||||
"compacting",
|
||||
f"🗜️ {COMPACTION_STATUS_MARKER} — still summarizing "
|
||||
"earlier conversation so I can continue...",
|
||||
)
|
||||
except Exception:
|
||||
logger.debug(
|
||||
"status_callback error in compression heartbeat", exc_info=True
|
||||
)
|
||||
|
||||
def _run(self) -> None:
|
||||
while not self._stop.wait(self._interval_seconds):
|
||||
if self._should_suppress():
|
||||
return
|
||||
self._touch("context compression in progress")
|
||||
self._emit_progress_status()
|
||||
|
||||
def _direct_messages_for_pre_compress_memory(messages: Any) -> list[dict[str, Any]]:
|
||||
"""Return direct user/assistant evidence safe for memory checkpointing.
|
||||
|
||||
@@ -167,6 +167,48 @@ def test_compression_activity_heartbeat_touches_agent_during_long_compress(tmp_p
|
||||
assert db.get_compression_lock_holder(session_id) is None
|
||||
|
||||
|
||||
def test_compression_activity_heartbeat_emits_client_status_events(tmp_path: Path) -> None:
|
||||
"""The heartbeat must emit ``compacting`` status events, not just DB touches.
|
||||
|
||||
Remote transports (e.g. the Android relay app) run idle-progress turn
|
||||
watchdogs that ``session.interrupt`` a turn after ~180s with no gateway
|
||||
events. Compression is silent on the event stream, so without periodic
|
||||
``status_callback`` heartbeats a long compression is killed mid-flight
|
||||
and retriggers forever on sessions near the context ceiling.
|
||||
"""
|
||||
db = SessionDB(db_path=tmp_path / "state.db")
|
||||
session_id = "HEARTBEAT_STATUS_TEST"
|
||||
db.create_session(session_id, source="test")
|
||||
|
||||
agent = _build_agent_with_db(db, session_id)
|
||||
agent._compression_activity_heartbeat_interval = 0.1
|
||||
touch_calls: list[str] = []
|
||||
agent._touch_activity = lambda desc, **_kw: touch_calls.append(desc)
|
||||
status_events: list[tuple[str, str]] = []
|
||||
setattr(
|
||||
agent,
|
||||
"status_callback",
|
||||
lambda event, message: status_events.append((event, message)),
|
||||
)
|
||||
|
||||
def _slow_compress(*_a, **_kw):
|
||||
_wait_for_touch(touch_calls, "context compression in progress")
|
||||
return [
|
||||
{"role": "user", "content": "[CONTEXT COMPACTION] summary"},
|
||||
{"role": "user", "content": "tail"},
|
||||
]
|
||||
|
||||
agent.context_compressor.compress.side_effect = _slow_compress
|
||||
messages = [{"role": "user", "content": f"m{i}"} for i in range(20)]
|
||||
|
||||
agent._compress_context(messages, "sys", approx_tokens=120_000)
|
||||
|
||||
compacting = [event for event, _message in status_events if event == "compacting"]
|
||||
# One emit at heartbeat start plus at least one periodic tick while the
|
||||
# summary call blocks.
|
||||
assert len(compacting) >= 2
|
||||
|
||||
|
||||
def test_lock_contender_preserves_terminal_compaction_lifecycle(tmp_path: Path) -> None:
|
||||
"""A lock loser still closes the structured compaction lifecycle.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user