fix(compression): stop timeout paths from blocking retries

This commit is contained in:
theo
2026-08-30 23:53:43 +02:00
committed by Teknium
parent 0a8b25e0ab
commit 8207862212
4 changed files with 142 additions and 14 deletions
+16 -11
View File
@@ -151,19 +151,24 @@ _SUMMARY_MISSING_CREDENTIAL_MARKERS: tuple[str, ...] = (
"no api key found",
)
_HYGIENE_IDLE_TIMEOUT_MARKERS: tuple[str, ...] = (
_HYGIENE_PREAGENT_ONLY_COOLDOWN_MARKERS: tuple[str, ...] = (
"session hygiene compression timed out",
"hygiene compression deferred: turn-hold budget expired",
)
def _is_hygiene_idle_timeout_error(error: object) -> bool:
"""Return True when the durable cooldown came from a hygiene watchdog timeout.
def _is_hygiene_preagent_only_cooldown(error: object) -> bool:
"""Return True for a cooldown that belongs only to pre-agent hygiene.
That persist is intentional for the pre-agent hygiene pass (#74136) but
must not block the in-conversation compressor (#86972).
Hygiene watchdog timeouts and turn-hold deferrals intentionally persist
retry spacing for the pre-agent pass (#74136), but neither is evidence of
an auxiliary-model failure and neither may block the in-agent compressor
(#86972).
"""
text = str(error or "").strip().casefold()
return any(marker in text for marker in _HYGIENE_IDLE_TIMEOUT_MARKERS)
return any(
marker in text for marker in _HYGIENE_PREAGENT_ONLY_COOLDOWN_MARKERS
)
def _response_finish_reason(response: Any) -> str:
@@ -2974,11 +2979,11 @@ class ContextCompressor(ContextEngine):
self._last_summary_error = None
return None
# Hygiene idle-watchdog timeouts persist the same column so the
# pre-agent pass can skip (#74136), but they are not evidence of a
# 429/aux-model fault. The in-conversation compressor has its own
# budget and must still be allowed to run (#86972).
if _is_hygiene_idle_timeout_error(state.get("error")):
# Hygiene watchdog timeouts and turn-hold deferrals persist the same
# column so the pre-agent pass can skip (#74136), but they are not
# evidence of a 429/aux-model fault. The in-conversation compressor has
# its own budget and must still be allowed to run (#86972).
if _is_hygiene_preagent_only_cooldown(state.get("error")):
# A later hygiene write can overwrite a previous aux-model row
# on the shared column. Drop any in-memory cooldown so the
# in-agent compressor is not still blocked after this refresh.
+21 -2
View File
@@ -4010,8 +4010,27 @@ def compress_context(
)
# Incoming-message interrupts and active-turn redirects must not tear an
# atomic summary in half (#23975). Explicit stop surfaces set a separate
# Event atomically; never infer cause from the racy message fields.
# Event atomically; never infer cause from the racy message fields. A
# host timeout also cancels the attempt's commit fence. Feed BOTH into
# the protected auxiliary-call seam so the compression owner unwinds
# promptly while an isolated provider stream finishes or closes in its
# daemon worker. Otherwise four timed-out streams retain all four shared
# compression-pool slots until the auxiliary stream's longer absolute
# ceiling expires.
_hard_cancel_event = getattr(agent, "_hard_interrupt_requested", None)
def _compression_cancel_requested() -> bool:
return bool(
(
_hard_cancel_event is not None
and _hard_cancel_event.is_set()
)
or (
commit_fence is not None
and commit_fence.is_cancelled
)
)
try:
# F6: never start expensive summary work for an already-cancelled
# fence (a stale queued job admitted after host departure).
@@ -4024,7 +4043,7 @@ def compress_context(
compressed = messages
else:
with aux_progress_hook(_progress_hook), aux_interrupt_protection(
cancel_event=_hard_cancel_event
cancel_check=_compression_cancel_requested
):
compressed = compress_fn(messages, **compress_kwargs)
# Freeze a hard stop that arrived after the final provider
@@ -130,6 +130,82 @@ def test_f3_mutating_engine_cannot_touch_live_transcript_after_timeout(
assert live == baseline
def test_host_timeout_releases_pool_slot_while_protected_provider_is_still_blocked(
tmp_path: Path, monkeypatch
) -> None:
"""Fence timeout must unwind the compression owner, not occupy the pool.
Protected auxiliary calls isolate their provider stream on a daemon thread.
The compression owner must observe its commit-fence cancellation and unwind
immediately; otherwise four slow streams consume all four shared compression
workers until the auxiliary stream's much longer absolute ceiling expires.
"""
from agent import auxiliary_client as aux
from agent import conversation_compression as cc
deadline = time.time() + 5
while time.time() < deadline:
with cc._compress_admission_lock:
if cc._compress_admitted_count == 0:
break
time.sleep(0.02)
with cc._compress_admission_lock:
assert cc._compress_admitted_count == 0
db = SessionDB(db_path=tmp_path / "state.db")
session_id = "F3_PROVIDER_OWNER_RELEASE"
db.create_session(session_id, source="cli")
agent = _build_agent_with_db(db, session_id)
agent._cached_system_prompt = "sys"
monkeypatch.setattr(
"agent.conversation_compression.resolve_context_compression_timeouts",
lambda cfg=None: (0.05, 0.1),
)
provider_started = threading.Event()
release_provider = threading.Event()
def _blocked_provider(_kwargs):
provider_started.set()
assert release_provider.wait(timeout=10)
return "late-provider-result"
def _compress_with_protected_provider(msgs, **_kwargs):
aux._run_protected_sync_provider_call(_blocked_provider, {})
return msgs
agent.context_compressor.compress.side_effect = _compress_with_protected_provider
live = [{"role": "user", "content": f"m{i}"} for i in range(20)]
try:
returned, _sp = agent._compress_context(
live, "sys", approx_tokens=120_000
)
assert returned is live
assert provider_started.wait(timeout=1)
assert not release_provider.is_set()
deadline = time.time() + 1
while time.time() < deadline:
with cc._compress_admission_lock:
if cc._compress_admitted_count == 0:
break
time.sleep(0.01)
with cc._compress_admission_lock:
assert cc._compress_admitted_count == 0, (
"timed-out compression owner retained its shared pool slot "
"while the isolated provider stream was still blocked"
)
finally:
release_provider.set()
deadline = time.time() + 5
while time.time() < deadline:
with cc._compress_admission_lock:
if cc._compress_admitted_count == 0:
break
time.sleep(0.02)
def test_f4_five_step_stale_holder_regression(tmp_path: Path) -> None:
"""Reviewer's exact 5-step durable-lease regression (#76354 F4).
@@ -1,4 +1,4 @@
"""Hygiene idle-timeout cooldowns must not block the in-agent compressor."""
"""Pre-agent hygiene retry spacing must not block in-agent compression."""
from __future__ import annotations
@@ -13,6 +13,11 @@ _HYGIENE_TIMEOUT_ERROR = (
"session hygiene compression timed out with no output from the summary model"
)
_HYGIENE_TURNHOLD_ERROR = (
"hygiene compression deferred: turn-hold budget expired while the "
"summary was still streaming"
)
def _bound_compressor(db: SessionDB, session_id: str) -> ContextCompressor:
with patch(
@@ -49,6 +54,29 @@ def test_hygiene_idle_timeout_does_not_block_in_agent_compressor(tmp_path: Path)
db.close()
def test_hygiene_turnhold_deferral_does_not_block_in_agent_compressor(
tmp_path: Path,
):
db = SessionDB(db_path=tmp_path / "state.db")
session_id = "hyg-turnhold-sid"
db.create_session(session_id, source="telegram")
db.record_compression_failure_cooldown(
session_id,
time.time() + 300,
_HYGIENE_TURNHOLD_ERROR,
)
compressor = _bound_compressor(db, session_id)
# Hygiene still sees its flat retry-spacing row and avoids respawning a
# doomed 15-second turn-hold attempt on every incoming message.
assert db.get_compression_failure_cooldown(session_id) is not None
# The in-conversation compressor has a separate budget. A healthy hygiene
# stream being deferred must not suppress that path for another minute.
assert compressor.get_active_compression_failure_cooldown() is None
db.close()
def test_aux_model_fault_cooldown_still_blocks_in_agent_compressor(tmp_path: Path):
db = SessionDB(db_path=tmp_path / "state.db")
session_id = "aux-fault-sid"