131 lines
4.9 KiB
Python
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
|