66 lines
1.7 KiB
Python
66 lines
1.7 KiB
Python
# retoor <retoor@molodetz.nl>
|
|||
|
|
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()
|