diff --git a/agent/conversation_compression.py b/agent/conversation_compression.py index c2e687413d..52d486866f 100644 --- a/agent/conversation_compression.py +++ b/agent/conversation_compression.py @@ -667,6 +667,46 @@ _compress_timeout_executor_lock = threading.Lock() # ceilings so overrun reporting stays observable at test timescales. _COMMIT_OVERRUN_WAIT_SLICE_SECONDS = 30.0 +# Bounded admission for the shared compress-timeout pool (#76354 review F6). +# The stdlib executor queue is unbounded: with all four workers wedged in hung +# summaries, a fifth compression would queue silently, wait out its whole +# timeout without ever starting, and remain eligible to run as a stale job +# whenever a worker recovered. Admission is therefore capped at the worker +# count — when every worker slot is occupied (running OR admitted-not-started) +# submission FAILS FAST and the caller continues without compression. +# +# Recovery contract when all workers are wedged: new compressions fail fast +# (no queue growth, conversation continues uncompressed, a warning is logged +# each attempt); wedged workers are fence-cancelled so they cannot publish +# anything when they eventually return, and each recovery frees its admission +# slot via the future done-callback, restoring normal service. If a worker +# NEVER returns, its slot is lost for the process lifetime — bounded, +# observable degradation instead of an unbounded stale-job queue. +_COMPRESS_EXECUTOR_MAX_WORKERS = 4 +_compress_admission_lock = threading.Lock() +_compress_admitted_count = 0 + + +class CompressionExecutorSaturatedError(RuntimeError): + """All compression pool slots are occupied; submission was refused.""" + + +def _try_admit_compression_job() -> bool: + """Reserve one bounded compression-pool admission slot (F6).""" + global _compress_admitted_count + with _compress_admission_lock: + if _compress_admitted_count >= _COMPRESS_EXECUTOR_MAX_WORKERS: + return False + _compress_admitted_count += 1 + return True + + +def _release_compression_admission(_future=None) -> None: + """Free an admission slot (future done-callback or failed submit).""" + global _compress_admitted_count + with _compress_admission_lock: + if _compress_admitted_count > 0: + _compress_admitted_count -= 1 def _get_compress_timeout_executor(): @@ -683,7 +723,7 @@ def _get_compress_timeout_executor(): # overlapping calls (live compress + fence-cancelled workers # still winding down), not asyncio's min(32, cpu+4) fan-out. _compress_timeout_executor = DaemonThreadPoolExecutor( - max_workers=4, + max_workers=_COMPRESS_EXECUTOR_MAX_WORKERS, thread_name_prefix="compress-ctx-timeout", ) return _compress_timeout_executor @@ -794,9 +834,43 @@ def run_compress_context_with_progress_timeout( from tools.thread_context import propagate_context_to_thread executor = _get_compress_timeout_executor() + # Bounded admission (#76354 F6): refuse rather than queue when every pool + # slot is occupied. A queued job would silently wait out its whole budget + # without starting and stay eligible to run as a stale cancelled job when + # a worker recovers. Fail fast: continue without compression this cycle. + if not _try_admit_compression_job(): + logger.warning( + "Context compression pool saturated (%d workers busy) — " + "refusing new compression this cycle and continuing without " + "compression. Wedged workers are fence-cancelled and free their " + "slot when they return; if this persists, check the summary " + "provider health.", + _COMPRESS_EXECUTOR_MAX_WORKERS, + ) + return messages, _resolve_fallback_prompt() + + def _fence_gated_worker(worker_fence: CompressionCommitFence): + # F6: an admitted job can still start after the host stopped waiting + # (worker slot freed late). Check the fence BEFORE any expensive + # summary work so a stale job never burns an LLM call; its return + # value is discarded by the already-departed host. + if worker_fence.is_cancelled: + logger.info( + "Skipping stale compression job: fence cancelled before start" + ) + return messages, "" + return worker(worker_fence) + # Bare pool workers start with an empty ContextVar map; propagate the # parent conversation/approval context into the worker. - future = executor.submit(propagate_context_to_thread(worker), fence) + try: + future = executor.submit( + propagate_context_to_thread(_fence_gated_worker), fence + ) + except BaseException: + _release_compression_admission() + raise + future.add_done_callback(_release_compression_admission) wait_started = time.monotonic() # F2: EVERY host unwind (KeyboardInterrupt, task cancellation, unexpected # exception while waiting) must revoke future commit admission before the @@ -831,6 +905,9 @@ def run_compress_context_with_progress_timeout( continue break + # F6: a not-yet-started future must not linger as a stale queued job. + # cancel() is a no-op for a running worker (fence handles that path). + future.cancel() cancelled: Optional[bool] = None while cancelled is None: @@ -2678,8 +2755,18 @@ def compress_context( except Exception: pass try: - with aux_progress_hook(_progress_hook): - compressed = compress_fn(messages, **compress_kwargs) + # F6: never start expensive summary work for an already-cancelled + # fence (a stale queued job admitted after host departure). + if commit_fence is not None and commit_fence.is_cancelled: + logger.info( + "Compression cancelled before summary dispatch " + "(session=%s) — skipping summary work.", + agent.session_id or "none", + ) + compressed = messages + else: + with aux_progress_hook(_progress_hook): + compressed = compress_fn(messages, **compress_kwargs) finally: if commit_fence is not None: try: diff --git a/tests/agent/test_compression_review_76354.py b/tests/agent/test_compression_review_76354.py index 748e4f8c22..6b63e5e9ac 100644 --- a/tests/agent/test_compression_review_76354.py +++ b/tests/agent/test_compression_review_76354.py @@ -34,8 +34,13 @@ from agent.conversation_compression import ( def _drain_admission_slots(): - """Placeholder until bounded admission (F6) lands in a later commit.""" - return + """Best-effort wait for pool admission slots to free between tests.""" + deadline = time.time() + 5 + while time.time() < deadline: + with cc._compress_admission_lock: + if cc._compress_admitted_count == 0: + return + time.sleep(0.02) class TestF1CommitOverrunWhileHung: @@ -265,3 +270,142 @@ class TestF4CooldownClearOrdering: fake2._compression_cancelled_check = staticmethod(lambda: False) ContextCompressor._clear_compression_failure_cooldown(fake2) assert fake2._summary_failure_cooldown_until == 0.0 + + +class TestF6ExecutorSaturation: + def test_saturated_pool_fails_fast_and_never_runs_stale_job(self): + """4 blocked summaries + 5th submission fails fast; recovery does not + run the refused job.""" + _drain_admission_slots() + release = threading.Event() + started = threading.Barrier(5, timeout=10) # 4 workers + main + + def blocked_worker(fence: CompressionCommitFence): + started.wait() + assert release.wait(timeout=30) + return ([], "done") + + hosts = [] + results = {} + + def host(i): + results[i] = run_compress_context_with_progress_timeout( + worker=blocked_worker, + messages=[{"role": "user", "content": f"m{i}"}], + system_prompt_fallback=f"fb{i}", + idle_timeout_seconds=0.05, + total_ceiling_seconds=0.1, + ) + + try: + for i in range(4): + t = threading.Thread(target=host, args=(i,), name=f"sat-{i}") + t.start() + hosts.append(t) + started.wait() # all 4 workers occupy the pool + for t in hosts: + t.join(timeout=5) # hosts time out; workers stay wedged + assert not t.is_alive() + + # All 4 slots still admitted (workers blocked). + with cc._compress_admission_lock: + assert cc._compress_admitted_count == 4 + + fifth_ran = threading.Event() + + def fifth_worker(fence): + fifth_ran.set() + return ([], "5th") + + fifth_msgs = [{"role": "user", "content": "fifth"}] + t0 = time.monotonic() + msgs, prompt = run_compress_context_with_progress_timeout( + worker=fifth_worker, + messages=fifth_msgs, + system_prompt_fallback="fifth-fallback", + idle_timeout_seconds=5.0, + total_ceiling_seconds=5.0, + ) + elapsed = time.monotonic() - t0 + # ── Assert while the 4 workers are STILL wedged ─────────────── + assert not release.is_set() + assert elapsed < 1.0, ( + f"saturated submission must fail fast, took {elapsed:.2f}s" + ) + assert msgs is fifth_msgs + assert prompt == "fifth-fallback" + assert not fifth_ran.is_set() + finally: + release.set() + + # Worker recovery: slots free, and the refused fifth job never runs. + _drain_admission_slots() + time.sleep(0.1) + assert not fifth_ran.is_set(), ( + "recovered workers must not run the refused stale job" + ) + # Recovery restores service: a new submission is admitted and runs. + msgs, prompt = run_compress_context_with_progress_timeout( + worker=lambda fence: ([{"role": "user", "content": "ok"}], "ok"), + messages=[{"role": "user", "content": "after"}], + system_prompt_fallback="fb", + idle_timeout_seconds=1.0, + total_ceiling_seconds=2.0, + ) + assert prompt == "ok" + _drain_admission_slots() + + def test_cancelled_fence_skips_summary_work_before_start(self): + """A stale job whose fence was already cancelled never runs summary. + + Drives compress_context's pre-summary fence gate directly: the fence + is cancelled BEFORE dispatch, so the expensive compress() call must + not run and the transcript must come back unchanged. + """ + import os + from pathlib import Path + from unittest.mock import MagicMock, patch + import tempfile + + from hermes_state import SessionDB + + with tempfile.TemporaryDirectory() as td: + db = SessionDB(db_path=Path(td) / "state.db") + session_id = "F6_PRESTART_FENCE" + 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, + ) + compressor = MagicMock() + compressor.compress.return_value = [ + {"role": "user", "content": "should-not-run"} + ] + compressor._last_summary_error = None + compressor._last_compress_aborted = False + compressor._last_aux_model_failure_model = None + compressor._last_aux_model_failure_error = None + agent.context_compressor = compressor + agent._cached_system_prompt = "sys" + + fence = CompressionCommitFence() + assert fence.cancel_before_commit() is True + + messages = [{"role": "user", "content": f"m{i}"} for i in range(20)] + returned, _sp = agent._compress_context( + messages, "sys", approx_tokens=120_000, commit_fence=fence + ) + + compressor.compress.assert_not_called() + assert returned is messages + # The cancelled attempt must not leave the durable lock held. + assert db.get_compression_lock_holder(session_id) is None