From 9bbd956b739d9997751d0b33bcac4274f35d00a6 Mon Sep 17 00:00:00 2001 From: kshitij <82637225+kshitijk4poor@users.noreply.github.com> Date: Sun, 2 Aug 2026 21:54:40 +0530 Subject: [PATCH] fix(memory/hindsight): evict timed-out retain ops + coarser status polls MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- plugins/memory/hindsight/__init__.py | 55 ++++++++++++++++--- .../plugins/memory/test_hindsight_provider.py | 32 +++++++++++ 2 files changed, 78 insertions(+), 9 deletions(-) diff --git a/plugins/memory/hindsight/__init__.py b/plugins/memory/hindsight/__init__.py index c12238518c..8c93df87f6 100644 --- a/plugins/memory/hindsight/__init__.py +++ b/plugins/memory/hindsight/__init__.py @@ -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. diff --git a/tests/plugins/memory/test_hindsight_provider.py b/tests/plugins/memory/test_hindsight_provider.py index fef2bf3e7f..e3ba1b8a55 100644 --- a/tests/plugins/memory/test_hindsight_provider.py +++ b/tests/plugins/memory/test_hindsight_provider.py @@ -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