fix(memory/hindsight): evict timed-out retain ops + coarser status polls
Review follow-up on the #62871 salvage (simplify pass, HIGH): 1. Ops unresolved at the wait deadline were RETAINED in the pending set. A permanently failing status endpoint (auth error, endless 500s, or a server that loses ops without 404) would grow the set forever and make EVERY later prefetch burn the full 10s budget re-polling it — and prefetch()'s bounded 3s join sits on the reply path, so that money-quote 'adds no response latency' claim breaks. Timed-out ops are now dropped (identical degradation to prefetch_waits_for_retain=False: possibly stale recall) with a WARNING so persistent server trouble is visible. Guard test mutation-checked (fails with eviction disabled). 2. Status polls now spaced 0.5s (was 0.05s shared with the local drain poll): a wedged op cost up to ~200 get_operation_status round trips per prefetch; now ~20 max over the default 10s budget.
This commit is contained in:
@@ -737,6 +737,11 @@ class HindsightMemoryProvider(MemoryProvider):
|
||||
self._pending_retain_ops: set[str] = set()
|
||||
self._pending_retain_ops_lock = threading.Lock()
|
||||
self._retain_ops_bank_id = ""
|
||||
# Seconds between get_operation_status polls while waiting for server-
|
||||
# side retain completion. Each poll is a server round trip, so this is
|
||||
# deliberately coarser than the 0.05s local queue-drain poll: ~20 calls
|
||||
# max over the default 10s budget instead of ~200.
|
||||
self._RETAIN_OP_POLL_INTERVAL_S = 0.5
|
||||
# Legacy alias — older tests/callers reference _sync_thread directly.
|
||||
# Points at _writer_thread once the writer is running.
|
||||
self._sync_thread = None
|
||||
@@ -1282,6 +1287,21 @@ class HindsightMemoryProvider(MemoryProvider):
|
||||
*deadline* is a ``time.monotonic()`` value (None = no bound). Completed
|
||||
ops are removed from the pending set as they finish so a later prefetch
|
||||
doesn't re-poll them.
|
||||
|
||||
Ops still pending when the deadline expires are DROPPED, not retained:
|
||||
keeping them would make a permanently failing status endpoint (auth
|
||||
error, endless 500s, server that loses ops without a 404) grow the
|
||||
pending set forever and burn the full timeout on EVERY subsequent
|
||||
prefetch — turning "bounded wait per prefetch" into unbounded
|
||||
session-wide degradation (and, via prefetch()'s bounded join on the
|
||||
reply path, a per-turn reply-latency penalty). Dropping trades a
|
||||
possibly-stale recall NOW (identical to prefetch_waits_for_retain=False
|
||||
behavior) for guaranteed liveness; the drop is logged at WARNING once
|
||||
per prefetch so persistent server trouble is visible.
|
||||
|
||||
Status polls are spaced by _RETAIN_OP_POLL_INTERVAL_S (0.5s) — server
|
||||
round trips per op are bounded (~20 over a 10s budget), unlike the
|
||||
cheap 0.05s local queue-drain poll in _wait_for_retains_drained.
|
||||
"""
|
||||
while True:
|
||||
with self._pending_retain_ops_lock:
|
||||
@@ -1293,29 +1313,46 @@ class HindsightMemoryProvider(MemoryProvider):
|
||||
return False
|
||||
|
||||
done: set[str] = set()
|
||||
expired = False
|
||||
for op_id in pending:
|
||||
if self._shutting_down.is_set():
|
||||
return False
|
||||
if deadline is not None and time.monotonic() >= deadline:
|
||||
logger.debug(
|
||||
"Prefetch: server retain visibility timed out after %.1fs "
|
||||
"(%d op(s) still pending)",
|
||||
timeout, len(pending) - len(done),
|
||||
)
|
||||
with self._pending_retain_ops_lock:
|
||||
self._pending_retain_ops.difference_update(done)
|
||||
return False
|
||||
expired = True
|
||||
break
|
||||
if self._is_retain_op_complete(bank_id, op_id):
|
||||
done.add(op_id)
|
||||
|
||||
if expired:
|
||||
with self._pending_retain_ops_lock:
|
||||
self._pending_retain_ops.difference_update(done)
|
||||
dropped = len(self._pending_retain_ops)
|
||||
self._pending_retain_ops.clear()
|
||||
logger.warning(
|
||||
"Prefetch: server retain visibility timed out after %.1fs; "
|
||||
"dropping %d unresolved op(s) so later prefetches stay "
|
||||
"bounded (recall may miss the just-completed turn)",
|
||||
timeout, dropped,
|
||||
)
|
||||
return False
|
||||
|
||||
with self._pending_retain_ops_lock:
|
||||
self._pending_retain_ops.difference_update(done)
|
||||
still_pending = bool(self._pending_retain_ops)
|
||||
if not still_pending:
|
||||
return True
|
||||
if deadline is not None and time.monotonic() >= deadline:
|
||||
with self._pending_retain_ops_lock:
|
||||
dropped = len(self._pending_retain_ops)
|
||||
self._pending_retain_ops.clear()
|
||||
logger.warning(
|
||||
"Prefetch: server retain visibility timed out after %.1fs; "
|
||||
"dropping %d unresolved op(s) so later prefetches stay "
|
||||
"bounded (recall may miss the just-completed turn)",
|
||||
timeout, dropped,
|
||||
)
|
||||
return False
|
||||
time.sleep(0.05)
|
||||
time.sleep(self._RETAIN_OP_POLL_INTERVAL_S)
|
||||
|
||||
def _writer_loop(self) -> None:
|
||||
"""Drain the retain queue serially. Exits on sentinel.
|
||||
|
||||
@@ -649,6 +649,38 @@ class TestPrefetchServerRetainVisibility:
|
||||
assert order == ["recall"], "prefetch should recall after the timeout"
|
||||
assert elapsed < 3.0, "prefetch must not block well past the drain budget"
|
||||
|
||||
def test_timed_out_ops_are_dropped_not_repolled(self, provider_with_config):
|
||||
"""Ops unresolved at deadline must be EVICTED so a permanently failing
|
||||
status endpoint can't make every later prefetch re-burn the full
|
||||
timeout on a growing pending set (unbounded session-wide degradation
|
||||
+ reply-path join penalty)."""
|
||||
p = provider_with_config(prefetch_retain_drain_timeout=0.3)
|
||||
p._client = self._client_with_ops(["pending"]) # never completes
|
||||
p._client.arecall = AsyncMock(
|
||||
return_value=SimpleNamespace(results=[SimpleNamespace(text="m")])
|
||||
)
|
||||
|
||||
p.sync_turn("hello", "world")
|
||||
p._retain_queue.join()
|
||||
assert p._pending_retain_ops, "op should be tracked before the wait"
|
||||
|
||||
# First prefetch burns the budget and must DROP the wedged op.
|
||||
p.queue_prefetch("q1")
|
||||
if p._prefetch_thread:
|
||||
p._prefetch_thread.join(timeout=5.0)
|
||||
assert p._pending_retain_ops == set(), (
|
||||
"unresolved ops must be evicted at deadline, not retained"
|
||||
)
|
||||
|
||||
# A later prefetch with nothing pending must be near-instant.
|
||||
start = time.monotonic()
|
||||
p.queue_prefetch("q2")
|
||||
if p._prefetch_thread:
|
||||
p._prefetch_thread.join(timeout=5.0)
|
||||
assert time.monotonic() - start < 0.25, (
|
||||
"second prefetch re-polled dropped ops — eviction regressed"
|
||||
)
|
||||
|
||||
def test_operation_notfound_treated_as_complete(self, provider):
|
||||
"""A NotFound (completed+evicted) op is treated as done, not pending."""
|
||||
from hindsight_client_api.exceptions import NotFoundException
|
||||
|
||||
Reference in New Issue
Block a user