forked from retoor/devplacepy
fix: correct "bugs" to "issues" in routing table and README references across multiple documentation files
This commit is contained in:
@@ -0,0 +1,103 @@
|
||||
# retoor <retoor@molodetz.nl>
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
from typing import Any, Callable, Optional
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
QUEUE_MAXSIZE = 10000
|
||||
DRAIN_BATCH = 64
|
||||
|
||||
|
||||
class BackgroundQueue:
|
||||
def __init__(self, maxsize: int = QUEUE_MAXSIZE) -> None:
|
||||
self._queue: asyncio.Queue = asyncio.Queue(maxsize=maxsize)
|
||||
self._task: Optional[asyncio.Task] = None
|
||||
self._submitted = 0
|
||||
self._processed = 0
|
||||
self._failed = 0
|
||||
self._inline = 0
|
||||
|
||||
@property
|
||||
def running(self) -> bool:
|
||||
return self._task is not None and not self._task.done()
|
||||
|
||||
def submit(self, fn: Callable[..., Any], *args: Any, **kwargs: Any) -> None:
|
||||
self._submitted += 1
|
||||
if not self.running:
|
||||
self._execute(fn, args, kwargs, inline=True)
|
||||
return
|
||||
try:
|
||||
self._queue.put_nowait((fn, args, kwargs))
|
||||
except asyncio.QueueFull:
|
||||
logger.warning("background queue full; running task inline")
|
||||
self._execute(fn, args, kwargs, inline=True)
|
||||
|
||||
def _execute(self, fn: Callable[..., Any], args: tuple, kwargs: dict, inline: bool = False) -> None:
|
||||
if inline:
|
||||
self._inline += 1
|
||||
try:
|
||||
fn(*args, **kwargs)
|
||||
self._processed += 1
|
||||
except Exception as exc:
|
||||
self._failed += 1
|
||||
logger.warning("background task %s failed: %s", getattr(fn, "__name__", fn), exc)
|
||||
|
||||
async def start(self) -> None:
|
||||
if self.running:
|
||||
return
|
||||
self._task = asyncio.create_task(self._drain())
|
||||
logger.info("background task queue started")
|
||||
|
||||
async def stop(self) -> None:
|
||||
task = self._task
|
||||
self._task = None
|
||||
if task is not None:
|
||||
task.cancel()
|
||||
try:
|
||||
await task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
self._flush_remaining()
|
||||
logger.info(
|
||||
"background task queue stopped processed=%s failed=%s inline=%s",
|
||||
self._processed,
|
||||
self._failed,
|
||||
self._inline,
|
||||
)
|
||||
|
||||
def _flush_remaining(self) -> None:
|
||||
while True:
|
||||
try:
|
||||
fn, args, kwargs = self._queue.get_nowait()
|
||||
except asyncio.QueueEmpty:
|
||||
break
|
||||
self._execute(fn, args, kwargs)
|
||||
self._queue.task_done()
|
||||
|
||||
async def _drain(self) -> None:
|
||||
while True:
|
||||
fn, args, kwargs = await self._queue.get()
|
||||
self._execute(fn, args, kwargs)
|
||||
self._queue.task_done()
|
||||
for _ in range(DRAIN_BATCH - 1):
|
||||
try:
|
||||
fn, args, kwargs = self._queue.get_nowait()
|
||||
except asyncio.QueueEmpty:
|
||||
break
|
||||
self._execute(fn, args, kwargs)
|
||||
self._queue.task_done()
|
||||
|
||||
def stats(self) -> dict:
|
||||
return {
|
||||
"running": self.running,
|
||||
"depth": self._queue.qsize(),
|
||||
"submitted": self._submitted,
|
||||
"processed": self._processed,
|
||||
"failed": self._failed,
|
||||
"inline": self._inline,
|
||||
}
|
||||
|
||||
|
||||
background = BackgroundQueue()
|
||||
Reference in New Issue
Block a user