fix(agent): bounded admission + stale-job cancellation for the compress pool (review F6)
The process-wide 4-worker pool retained the stdlib executor's unbounded queue: four hung summaries wedged every slot, a fifth compression queued silently, waited out its whole budget without starting, and remained eligible to run later as an expensive stale job whose fence was already cancelled (the first fence check used to sit AFTER the summary call). - Bounded admission: submission fails fast (messages returned unchanged, loud warning) when all pool slots are occupied; slots are freed by a future done-callback. Recovery contract documented at the constant: new work fails fast while wedged, wedged workers are fence-cancelled and restore service when they return; a worker that never returns costs its slot — bounded, observable degradation instead of unbounded queueing. - Not-yet-started futures are cancel()ed on timeout. - The cancelled fence is checked BEFORE any expensive summary work, both in the pooled wrapper (stale queued job) and inside compress_context (pre-summary gate), so a stale job never burns an LLM call or acquires session state. Saturation regression: 4 event-blocked summaries wedge the pool, a 5th submission fails fast (asserted while the four are provably still blocked), the refused job never runs after worker recovery, and a fresh submission after recovery succeeds. PR #76354 review, blocking finding 6 / merge gate 7.
This commit is contained in:
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user