fix: return a raced-successful compression from _await_worker_within_budget
The `future.done()` guard added for dead workers (#117261 / #63892) unconditionally logged `future.exception()` and returned `(False, None)`. If the worker finished SUCCESSFULLY in the window between `result(timeout=)` expiring and the `done()` check, that discarded a completed compression and sent the caller down the stall/fallback path — `future.cancel()` becomes a no-op on a settled future and the fallback route costs a second LLM call. Siblings `_await_in_flight_commit` and tool_executor's `_poll_sequential_future` already re-read the result on a settled future; do the same here: `exception() is None` → return `(True, result)`, otherwise take the stall path as before. The compression-seam test grows a settled-successful case through the same helper (a Future subclass whose first timed `result()` still raises TimeoutError, modelling the race); it fails on the previous guard and passes with this one.
This commit is contained in:
@@ -1064,7 +1064,10 @@ def _await_worker_within_budget(
|
||||
# Aliases builtin TimeoutError (3.11+): also fires when the WORKER died with one (#63892).
|
||||
# A settled future never unsettles — re-waiting spun ~2k iter/s; take the stall path now.
|
||||
if future.done():
|
||||
logger.info("Context compression worker exited with %r — taking the stall path", future.exception())
|
||||
exc = future.exception()
|
||||
if exc is None:
|
||||
return True, future.result()
|
||||
logger.info("Context compression worker exited with %r — taking the stall path", exc)
|
||||
return False, None
|
||||
waited = time.monotonic() - wait_started
|
||||
since_progress = fence.seconds_since_progress()
|
||||
|
||||
@@ -53,6 +53,34 @@ def test_compression_wait_returns_promptly_when_worker_dies():
|
||||
assert (settled, result) == (False, None)
|
||||
assert _join_cancelled_worker(future, 0.5) is True
|
||||
|
||||
# Success race: the worker finishes in the window between ``result(timeout=)`` expiring and
|
||||
# the ``done()`` check. The guard must hand back the completed compression, not discard it
|
||||
# onto the stall/fallback path (which would cost a second LLM call).
|
||||
raced = _RacedSuccessFuture()
|
||||
raced.set_result("summary")
|
||||
fence = CompressionCommitFence()
|
||||
fence.touch_progress()
|
||||
started = time.monotonic()
|
||||
assert _await_worker_within_budget(raced, fence, idle=1.0, ceiling=2.0, wait_started=started) == (
|
||||
True,
|
||||
"summary",
|
||||
)
|
||||
|
||||
|
||||
class _RacedSuccessFuture(concurrent.futures.Future):
|
||||
"""Already-settled future whose FIRST timed ``result()`` still raises TimeoutError, modelling
|
||||
a worker that completed just after the wait slice expired."""
|
||||
|
||||
def __init__(self):
|
||||
super().__init__()
|
||||
self._timed_out_once = False
|
||||
|
||||
def result(self, timeout=None):
|
||||
if timeout is not None and not self._timed_out_once:
|
||||
self._timed_out_once = True
|
||||
raise concurrent.futures.TimeoutError()
|
||||
return super().result(timeout)
|
||||
|
||||
|
||||
def test_in_flight_commit_surfaces_worker_timeout_instead_of_looping():
|
||||
"""_await_in_flight_commit has no ceiling on this path — pre-fix a dead worker spun forever."""
|
||||
|
||||
Reference in New Issue
Block a user