Files
hermes-agent/agent/memory_work_queue.py
T

131 lines
4.9 KiB
Python

"""One acceptance order and owner for MemoryManager's provider work.
Executor submissions only wake drain(); they never carry individual work.
The condition is never held while a provider or a receipt callback runs.
"""
from collections import deque
from concurrent.futures import Future
import logging
import threading
logger = logging.getLogger(__name__)
class MemoryWorkQueue:
def __init__(self):
self.condition = threading.Condition()
self.pending = {}
self._queue = deque()
self._sequence = 0
self._owner = None
self._accepting = True
self._finalizer = None
self._finalizer_set = False
self._finalizer_receipt = None
def submit(self, fn, *, kind="write", inline_if_idle=False):
with self.condition:
if not self._accepting:
return None, False
inline = inline_if_idle and self._owner is None and not self._queue
self._sequence += 1
receipt = Future()
self._queue.append((self._sequence, fn, receipt))
self.pending[receipt] = kind
if inline:
self._owner = threading.get_ident()
return receipt, inline
def needs_wakeup(self):
with self.condition:
return (bool(self._queue) or self._finalizer is not None) and self._owner is None
def owned_by_current_thread(self):
with self.condition:
return self._owner == threading.get_ident()
def drain(self, *, claimed=False):
with self.condition:
if not claimed:
if self._owner is not None:
return
self._owner = threading.get_ident()
while True:
with self.condition:
if self._queue:
_, fn, receipt = self._queue.popleft()
running = receipt.set_running_or_notify_cancel()
elif self._finalizer is not None:
fn, self._finalizer = self._finalizer, None
receipt = self._finalizer_receipt
assert receipt is not None
running = True
else:
self._owner = None
self.condition.notify_all()
return
if running:
try:
result = fn()
except BaseException as exc:
# Keep the owner alive for accepted work and deferred cleanup,
# even when a provider raises SystemExit on the daemon worker.
logger.warning("Memory work failed: %s", exc, exc_info=True)
receipt.set_exception(exc)
else:
receipt.set_result(result)
with self.condition:
self.pending.pop(receipt, None)
self.condition.notify_all()
def flush(self, timeout=None, *, snapshot=None):
with self.condition:
receipts = tuple(self.pending) if snapshot is None else snapshot
if self._owner == threading.get_ident() and any(not f.done() for f in receipts):
return False
return self.condition.wait_for(lambda: all(f.done() for f in receipts), timeout)
def stop(self):
with self.condition:
self._accepting = False
return tuple(self.pending)
def abandon_queued(self):
# Remove before cancel: no worker can race cancellation notification.
with self.condition:
queued = tuple(self._queue)
self._queue.clear()
kinds = [self.pending[receipt] for _, _, receipt in queued]
for _, _, receipt in queued:
receipt.cancel()
receipt.set_running_or_notify_cancel()
with self.condition:
for _, _, receipt in queued:
self.pending.pop(receipt, None)
active = sum(not f.done() for f in self.pending)
self.condition.notify_all()
return kinds, active
def finish(self, finalizer, *, schedule_only=False):
"""Register one non-cancellable close receipt; defer to any active owner.
Registration is not completion. External callers can flush this receipt;
an owner must return to drain before its finalizer can run.
schedule_only lets bounded shutdown arrange a wake without running close
on its caller if ownership was released just before registration.
"""
with self.condition:
self._accepting = False
if self._finalizer_set:
return self._finalizer_receipt
receipt = self._finalizer_receipt = Future()
receipt.set_running_or_notify_cancel()
self.pending[receipt] = "finalizer"
self._finalizer_set = True
self._finalizer = finalizer
self.condition.notify_all()
if not schedule_only:
self.drain()
return receipt