"""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