Last active
October 29, 2025 15:59
-
-
Save Vizonex/3b94c59f01ba441f6582a8324d51ccb7 to your computer and use it in GitHub Desktop.
python's standard scheduler made for asyncio with a couple of upgrades. IF anyone wants this I'll make it into a pypi package and I'll go and add some pytests.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| """A Small and compact scheduler built off asyncio and reedited around python-std sched""" | |
| from __future__ import annotations | |
| import asyncio | |
| import heapq | |
| import inspect | |
| import types | |
| from itertools import count | |
| from typing import ( | |
| Any, | |
| Awaitable, | |
| Callable, | |
| Generic, | |
| NamedTuple, | |
| ParamSpec, # 3.9 died recently. | |
| Sequence, | |
| TypeVar, | |
| final, | |
| overload, | |
| ) | |
| # only depedency (and I maintain this one as well) | |
| from aiocallback import EventList, contextevent | |
| _T = TypeVar("_T") | |
| _P = ParamSpec("_P") | |
| @types.coroutine | |
| def __heartbeat(): | |
| """implements a small checkpoint system | |
| for the eventloop simillar to what asyncio.sleep() | |
| does but at a lower level aka `asyncio.tasks.__sleep0()`""" | |
| yield | |
| class Event(NamedTuple): | |
| time: float | |
| """Numeric type compatible with the return value | |
| of the timefunc function passed to the constructor.""" | |
| priority: int | |
| """Events scheduled for the same time will be executed in the order of their priority.""" | |
| sequence: int | |
| """A continually increasing sequence number that | |
| separates events if time and priority are equal. | |
| (you can consider it as being a special id)""" | |
| action: Callable[..., Awaitable[Any]] | |
| """Executing the event means executing await action(*args, **kwargs)""" | |
| args: Sequence[Any] | |
| """args is a sequence holding the positional arguments for the action.""" | |
| kwargs: dict[str, Any] | |
| """kwargs is a dictionary holding the keyword args for the action.""" | |
| name: str | None | |
| """An event's given name""" | |
| # inspired by aiohttp's TraceConfig object but with an easier to understand interface thanks to aiocallback. | |
| @final | |
| class Trace(EventList, Generic[_T]): | |
| """Can be utilized for reporting on important events taking place in or around another context object | |
| this can be used for setting up debugging or other functionallities involved without needing to subclass | |
| the main Scheduler object for enhanced performance.""" | |
| __class_getitem__ = classmethod(types.GenericAlias) | |
| @overload | |
| def __init__(self: "Trace[None]") -> None: ... | |
| @overload | |
| def __init__(self, context: _T) -> None: ... | |
| def __init__(self, context: _T | None = None) -> None: | |
| self._context = context | |
| @property | |
| def context(self) -> _T: | |
| """The context passed in the Trace's events""" | |
| return self._context | |
| @contextevent | |
| async def on_started(self, sched: Scheduler, event: Event, ctx: _T) -> None: | |
| """A Callback for executing and reporting a given event | |
| being executed:: | |
| from aiosched import Trace, Event | |
| trace = Trace() # None pass by default | |
| @trace.on_started | |
| async def started_event(sched, event:Event, ctx): | |
| # do something... | |
| """ | |
| async def send_on_started(self, sched: Scheduler, event: Event) -> None: | |
| await self.on_started.send(sched, event, self.context) | |
| @contextevent | |
| async def on_finished(self, sched: Scheduler, event: Event, ctx: _T) -> None: | |
| """A callback for reporting a given event as being finished...""" | |
| async def send_on_finished(self, sched: Scheduler, event: Event) -> None: | |
| await self.on_finished.send(sched, event, self.context) | |
| @final | |
| class Scheduler: | |
| """asyncio version of the python standard library's sched module | |
| with a few edits on how some things are done such as using default | |
| timeouts instead of custom ones... | |
| https://docs.python.org/3/library/sched.html#module-sched | |
| """ | |
| def __init__(self, traces: list[Trace[_T]] = []) -> None: | |
| """ | |
| :param traces: hook objects to initiate when an event has executed or has finished. | |
| """ | |
| self._queue: list[Event] = [] | |
| self._loop = asyncio.get_event_loop() | |
| self._lock = asyncio.Lock() | |
| self._sequence_generator = count() | |
| self._traces = traces | |
| # Freeze immediately... | |
| for t in self._traces: | |
| t.freeze() | |
| async def enterabs( | |
| self, | |
| time: float, | |
| priority: int, | |
| action: Callable[..., Awaitable[_T]], | |
| args: Sequence[Any] = (), | |
| kwargs: dict[str, Any] | None = None, | |
| name: str | None = None, | |
| ): | |
| """Enter a new event in the queue at an absolute time. | |
| Returns an ID for the event which can be used to remove it, | |
| if necessary. | |
| :param time: interval to wait on this event for... | |
| :param priority: determines when should it be queued? | |
| """ | |
| if not inspect.iscoroutinefunction(action): | |
| raise TypeError( | |
| f"only asynchronous actions can be added not {type(action)}" | |
| ) | |
| if kwargs is None: | |
| kwargs = {} | |
| async with self._lock: | |
| event = Event( | |
| time, | |
| priority, | |
| next(self._sequence_generator), | |
| action, | |
| args, | |
| kwargs, | |
| name, | |
| ) | |
| heapq.heappush(self._queue, event) | |
| return event # The ID | |
| async def enter( | |
| self, | |
| delay: float, | |
| priority: int, | |
| action: Callable[..., Awaitable[_T]], | |
| args: Sequence[Any] = (), | |
| kwargs: dict[str, Any] | None = None, | |
| name: str | None = None, | |
| ): | |
| """A variant that specifies the time as a relative time. | |
| This is actually the more commonly used interface. | |
| """ | |
| return await self.enterabs( | |
| self._loop.time() + delay, priority, action, args, kwargs, name | |
| ) | |
| async def cancel(self, event: Event) -> None: | |
| """Remove an event from the queue. | |
| This must be presented the ID as returned by enter(). | |
| If the event is not in the queue, this raises ValueError. | |
| """ | |
| async with self._lock: | |
| self._queue.remove(event) | |
| heapq.heapify(self._queue) | |
| async def empty(self) -> bool: | |
| """Check whether the queue is empty.""" | |
| async with self._lock: | |
| return not self._queue | |
| # NOTE: This function is not derrived from the standard library but important for tracing Schedulers and adding customization | |
| async def __run_event(self, event: Event) -> None: | |
| if self._traces: | |
| for t in self._traces: | |
| await t.send_on_started(self, event) | |
| await event.action(*event.args, **event.kwargs) | |
| for t in self._traces: | |
| await t.send_on_finished(self, event) | |
| else: | |
| await event.action(*event.args, **event.kwargs) | |
| @overload | |
| async def run(self) -> None: ... | |
| @overload | |
| async def run(self, blocking: bool = False) -> float: ... | |
| async def run(self, blocking: bool = True) -> float | None: | |
| """Execute events until the queue is empty. | |
| If blocking is False executes the scheduled events due to | |
| expire soonest (if any) and then return the deadline of the | |
| next scheduled call in the scheduler. | |
| :param blocking: run all actions at one time otherwise return back a float | |
| with how much time is left before the next one can be reasonably queued | |
| :returns: a float if blocking is false and tasks have | |
| not all completely queued yet... | |
| """ | |
| loop = self._loop | |
| lock = self._lock | |
| q = self._queue | |
| pop = heapq.heappop | |
| delayfunc = asyncio.sleep | |
| running_tasks: list[asyncio.Task[Any]] = [] | |
| while True: | |
| async with lock: | |
| if not q: | |
| break | |
| (time, p, s, action, args, kwargs) = q[0] | |
| now = loop.time() | |
| if not (delay := time > now): | |
| pop(q) | |
| await __heartbeat() | |
| if delay: | |
| if not blocking: | |
| if running_tasks: | |
| # join them all including finished or pending... | |
| await asyncio.gather(*running_tasks) | |
| return time - now | |
| else: | |
| _task = loop.create_task(action(*args, **kwargs)) | |
| running_tasks.append(_task) | |
| # rejoin all queued tasks even if they run forever... | |
| if running_tasks: | |
| await asyncio.gather(*running_tasks) | |
| @property | |
| def queue(self) -> list[Event]: | |
| """An ordered list of upcoming events. | |
| Events are attr classes with fields for: | |
| time, priority, action, argss, kwargs | |
| """ | |
| # Use heapq to sort the queue rather than using 'sorted(self._queue)'. | |
| # With heapq, two events scheduled at the same time will show in | |
| # the actual order they would be retrieved. | |
| with self._lock: | |
| events = self._queue[:] | |
| return list(map(heapq.heappop, [events] * len(events))) | |
| class PreparedScheduler(Generic[_T]): | |
| """A version of a scheduler used for preparing events via wrappers or through other means of access | |
| before asynchronous execution takes place...""" | |
| __class_getitem__ = classmethod(types.GenericAlias) | |
| @overload | |
| def __init__(self: "PreparedScheduler[None]") -> None: ... | |
| @overload | |
| def __init__(self, context: _T) -> None: ... | |
| def __init__(self, context: _T | None = None) -> None: | |
| self._context = context | |
| self._events: list[Event] = [] | |
| @property | |
| def context(self) -> _T: | |
| """The context passed in the PreparedScheduler for handling callbacks during pending so that other tasks can take place...""" | |
| return self._context | |
| def add_event( | |
| self, | |
| time: float, | |
| priority: int, | |
| action: Callable[_P, Awaitable[_T]], | |
| args: Sequence[Any], | |
| kwargs: dict[str, Any], | |
| name: str | None, | |
| ): | |
| if not inspect.iscoroutinefunction(action): | |
| raise TypeError( | |
| f"only asynchronous actions can be added not {type(action)}" | |
| ) | |
| heapq.heappush( | |
| self._events, | |
| Event( | |
| time, | |
| priority, | |
| next(self._sequence_generator), | |
| action, | |
| args, | |
| kwargs, | |
| name, | |
| ), | |
| ) | |
| def enter( | |
| self, | |
| delay: float, | |
| priority: int, | |
| args: Sequence[Any] = (), | |
| kwargs: dict[str, Any] | None = None, | |
| name: str | None = None, | |
| ) -> Callable[[Callable[_P, Awaitable[_T]]], Callable[_P, Awaitable[_T]]]: | |
| """Enter a new event in the queue at an absolute time. | |
| Returns an ID for the event which can be used to remove it, | |
| if necessary.:: | |
| sched = PreparedScheduler() | |
| @sched.enter(10, 1, args=('data'), name="run") | |
| async def run(data: str): | |
| pass | |
| :param time: interval to wait on this event for... | |
| :param priority: determines when should it be queued? | |
| :param args: the arguments that should get passed for this event | |
| :param kwargs: keywords arguments that should get passed for this event | |
| :param name: an optional name to give the event being assigned | |
| """ | |
| if kwargs is None: | |
| kwargs = {} | |
| def decorator(func: Callable[_P, Awaitable[_T]]) -> Callable[_P, Awaitable[_T]]: | |
| self.add_event(delay, priority, func, args, kwargs, name) | |
| return func | |
| return decorator | |
| @contextevent | |
| async def on_pending(self, ctx: _T, waiting: float) -> None: | |
| """Called while tasks are in pending and not blocking, | |
| using this wrapper is completely optional""" | |
| def freeze(self): | |
| """Freezes on_pending callback""" | |
| if not self.on_pending.frozen: | |
| self.on_pending.freeze() | |
| async def start(self, traces: list[Trace[_T]]): | |
| scheduler = Scheduler(traces) | |
| scheduler._queue = self.events.copy() | |
| self.freeze() | |
| while ret := await scheduler.run(self._blocking): | |
| await self.on_pending.send(self._context, ret) | |
| def run( | |
| self, | |
| runner: Callable[ | |
| [Callable[[list[Trace[_T]]], Awaitable[None]]], _T | |
| ] = asyncio.run, | |
| traces: list[Trace[_T]] = [], | |
| ) -> None: | |
| """Runs the event through this eventloop you can replace runner with another runner like winloop or uvloop if needed""" | |
| return runner(self.start(traces)) |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment