diff --git a/agent/context_compressor.py b/agent/context_compressor.py index b20199620e..3f4b1f3530 100644 --- a/agent/context_compressor.py +++ b/agent/context_compressor.py @@ -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. diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index 322d874411..12c8304892 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -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 diff --git a/tests/agent/test_compression_worker_isolation_76354.py b/tests/agent/test_compression_worker_isolation_76354.py index df0c483d70..fe976e5e5b 100644 --- a/tests/agent/test_compression_worker_isolation_76354.py +++ b/tests/agent/test_compression_worker_isolation_76354.py @@ -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). diff --git a/tests/agent/test_hygiene_timeout_cooldown_isolation.py b/tests/agent/test_hygiene_timeout_cooldown_isolation.py index e1493bbab5..25c4ffcc74 100644 --- a/tests/agent/test_hygiene_timeout_cooldown_isolation.py +++ b/tests/agent/test_hygiene_timeout_cooldown_isolation.py @@ -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"