Skip to content

Instantly share code, notes, and snippets.

@Vizonex
Last active October 29, 2025 15:59
Show Gist options
  • Select an option

  • Save Vizonex/3b94c59f01ba441f6582a8324d51ccb7 to your computer and use it in GitHub Desktop.

Select an option

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.
"""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