# retoor import asyncio import logging logger = logging.getLogger(__name__) class BackgroundQueue: def __init__(self, max_size=10000): self.max_size = max_size self._queue = None self._task = None @property def running(self): return self._task is not None and not self._task.done() def start(self): if self.running: return self._queue = asyncio.Queue(maxsize=self.max_size) self._task = asyncio.get_running_loop().create_task(self._consume()) def submit(self, fn, *args, **kwargs): if not self.running: self._run(fn, args, kwargs) return try: self._queue.put_nowait((fn, args, kwargs)) except asyncio.QueueFull: self._run(fn, args, kwargs) async def stop(self): if self._task is None: return self._task.cancel() try: await self._task except asyncio.CancelledError: pass self._task = None self.drain() def drain(self): if self._queue is None: return while not self._queue.empty(): fn, args, kwargs = self._queue.get_nowait() self._run(fn, args, kwargs) async def _consume(self): while True: fn, args, kwargs = await self._queue.get() self._run(fn, args, kwargs) self._queue.task_done() @staticmethod def _run(fn, args, kwargs): try: fn(*args, **kwargs) except Exception: logger.exception("background task %s failed", getattr(fn, "__name__", fn)) background = BackgroundQueue()