Skip to content

Instantly share code, notes, and snippets.

@Vizonex
Created October 31, 2025 20:01
Show Gist options
  • Select an option

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

Select an option

Save Vizonex/f3b1cf2bccff392a26b64c9b62373e19 to your computer and use it in GitHub Desktop.
Mainly used so I can demonstrate something I wanted implemented in apsheduler before learning it existed. This code runs correctly and I have a quick and dirty example of the way it is used.
"""A Small and compact scheduler built off anyio and reedited around python-std shed"""
from __future__ import annotations
import inspect
import sys
from collections import deque
from datetime import datetime, timedelta
from itertools import count
from typing import (
Any,
Awaitable,
Callable,
Sequence,
TypeVar,
Generic,
ParamSpec
)
from attrs import define, field
import anyio
from aiosignal import Signal
from anyio.lowlevel import checkpoint
from enum import Enum
from apscheduler import schedulers
# typing_extensions is included with anyio by default so we can use it...
if sys.version_info >= (3, 11):
from typing import TypeVarTuple, Unpack
else:
from typing_extensions import TypeVarTuple, Unpack
P = ParamSpec("P")
R = TypeVar("R")
R_cov = TypeVar("R_cov", covariant=True)
T_contra = TypeVar("T_contra", contravariant=True)
PosArgsT = TypeVarTuple("PosArgsT", default=Unpack[tuple[()]])
"""Allows for typehinting out positional arguments such as context variables"""
AsyncFunc = Callable[P, Awaitable[R]]
AsyncWrapper = Callable[[AsyncFunc[P, R]], AsyncFunc[P, R]]
# State is used for controlling Scheduler so that it can properly
# execute on running it thus preventing newer events from being
# taken in while in the middle of running
class _State(Enum):
OFF = "OFF"
RUNNING = "RUNNING"
@define(slots=True)
class Event(Generic[Unpack[PosArgsT], R]):
# instead of using heaps we can be a little more accurate or persise
# by treating all events as handles...
# Making up for the fact that users might want something more accurate than
# a float as rounding that float, this can be a problem if the task is set
# to weekly monthly or yearly...
action: Callable[..., Awaitable[R]]
"""Executing the event means executing await action(*args, **kwargs)"""
args: Sequence[Any] = field(factory=list)
"""args is a sequence holding the positional argss for the action."""
kwargs: dict[str, Any] = field(factory=dict)
"""kwargs is a dictionary holding the keyword argss for the action."""
time: timedelta | None = None
"""Numeric type compatible with the return value
of the timefunc function passed to the constructor."""
name: str | None = None
"""A Name or tag to give to this event, goes well with callback functions"""
instant: bool = True
"""When ran the first time do not wait, if repeated run,
again with the timeout in mind afterwards"""
repeat: bool = True
"""flags the event as being a repeatable event, this will make it so that the listener
runs forever..."""
is_wrapper:bool = False
"""flags that our event was setup via wrapper, hinting that we will need to pass
in our own context via run(...) in order to properly use it"""
traces:list["TraceEvent[Unpack[PosArgsT], R]"] = field(factory=list)
"""A list of listeners to listen for events starting or ending"""
eta: datetime | None = field(default=None, init=False) # it hides...
def new_eta(self):
"""Schedules a new deadline for when this event can be ran"""
if self.time is None:
raise RuntimeError(
"A timeout is required for this Event to be ran before setting a new eta"
)
self.eta = datetime.now() + self.time
def __hash__(self):
return hash(self.sequence)
# TODO: Signal events for startup and finishing...
async def run(self) -> R:
"""Calls the function and it's arguments once"""
return await self.action(*self.args, **self.kwargs)
async def run_with_context(self, *args:Any) -> R:
"""Runs an event but with context varaibles as positional arguments"""
return await self.action(*args)
class RunningScheduler(Generic[Unpack[PosArgsT], R]):
"""Used to simplify tasks and grouping by treating events as handles
and going through them consecutatively all until all events
have been fully exhausted."""
def __init__(
self,
events: Sequence[Event[Unpack[PosArgsT], R]],
*contexts:Unpack[PosArgsT]
):
# Loop through the events and treat them as handles simillar to how uvloop and asyncio both work...
self._events = deque(events)
# will use this for all of our tasks we run since it will join them all before closing out...
self._group = anyio.create_task_group()
# flag for checking still needs to be in loop for later or not...
self._empty = False
self._context_args = contexts
async def run_event(self, event: Event[Unpack[PosArgsT], R], task_status: anyio.TaskStatus[None] = anyio.TASK_STATUS_IGNORED) -> None:
# if we set this event as a wrapper system then call it with any positional
# arguments passed in from run(...)
# (likely was setup from click | anyio-typer | or another commandline tool)
# We made wrappers to stop a deep problem with most tools which is retaining
# a clean frontend.
task_status.started()
if traces := event.traces:
# Queue up all trace objets to be ran...
for t in traces:
await t.on_start.send(event, *self._context_args)
if event.is_wrapper:
ret = await event.run_with_context(*self._context_args)
else:
# call with events already sheduled and queued items instead of context arguments
ret = await event.run()
if traces := event.traces:
for t in traces:
await t.on_finish.send(event, ret, *self._context_args)
# See if the handle needs to be added back into the loop
if event.repeat:
# send back for calling later and assign a new eta if possible...
if event.time:
event.eta = datetime.now() + event.time
# Run later with or without eta or timeout...
self._events.append(event)
if self._empty:
# turn back on as we still have stuff left to complete...
self._empty = False
self._group.start_soon(self.check_event)
async def shedule_event(self, event: Event, task_status: anyio.TaskStatus[None] = anyio.TASK_STATUS_IGNORED) -> None:
"""queues an event to be ran later..."""
task_status.started()
self._group.start_soon(self.run_event, event)
await checkpoint()
async def check_event(self, task_status: anyio.TaskStatus[None] = anyio.TASK_STATUS_IGNORED) -> None:
"""Main loop for checking and finishing off handles one by one
it does just one event at a time so that background tasks are being
visited at the same time..."""
task_status.started()
if not self._events:
# it's done just need to wait on the queued events to finish cycling...
# There is no call-soon-like handle
# (everything is async so we should just noop)
self._empty = True
return await checkpoint()
# obtain from the front rather than the
# back to keep proper ordering in check...
event = self._events.popleft()
if event.instant:
event.instant = False
self.shedule_event(event)
return await checkpoint()
elif event.time is not None:
if event.eta is not None:
# is the event ready to be initiated?
if event.eta < datetime.now():
self.shedule_event(event)
return await checkpoint()
else:
# first time?
event.new_eta()
# recycle this event again for later...
self._events.append(event)
self._group.start_soon(self.check_event)
return await checkpoint()
else:
# Event must be sheduled since it has no timeout and isn't treated as running immediately...
await self.shedule_event(event)
self._group.start_soon(self.check_event)
return await checkpoint()
async def run(self) -> None:
"""starts listening for events to run..."""
await self._group.start(self.check_event)
async def __aenter__(self):
await self._group.__aenter__()
return self
async def __aexit__(self, *args):
return await self._group.__aexit__(*args)
class TraceEvent(Generic[Unpack[PosArgsT], R_cov]):
"""A Callback system for calling back to trace or debug events
that have started or have been closed making it extemely easy to
write hooks to either specific groups of events or all of them,
the context arguments passed to Scheduler.start(...) will be
passed to these trace signals"""
def __init__(self):
# TODO: Reorganize argument order if needed...
self._on_start: Signal[Event, Unpack[PosArgsT]] = Signal(self)
self._on_finish: Signal[Event, R_cov, Unpack[PosArgsT]] = Signal(self)
self._frozen = False
@property
def on_start(self) -> Signal[Event, Unpack[PosArgsT]]:
"""wraps a startup function for an event that is about to get queued"""
return self._on_start
@property
def on_finish(self) -> Signal[Event, R_cov, Unpack[PosArgsT]]:
"""Called when an event is finished it will return in the order of
`passed context variables, event, returned item`"""
return self._on_finish
def freeze(self):
"""Freezes TraceEvent's callbacks to prevent misuse"""
if not self._frozen:
self._on_start.freeze()
self._on_finish.freeze()
self._frozen = True
class Scheduler(Generic[Unpack[PosArgsT]]):
"""Anyio version of the python standary library's sched module
https://docs.python.org/3/library/sched.html#module-sched
"""
def __init__(self) -> None:
self._queue: deque[Event[Unpack[PosArgsT], Any]] = deque()
self._state = _State.OFF
def __update_state(self, new_state:_State) -> None:
"""Updates scheduler's state
:raises RuntimeError: if state is invalid or impossible to move to this helps capture invalid movements in the Sheduler's
system..."""
if self._state == new_state:
# Show exactly why __update_state screwed up.
# Redundant updates to the same value are unacceptable and this w
# it will let us catch it immediately...
raise RuntimeError(f"invalid state {self._state.value}")
self._state = new_state
# I am often ridiculed for having redundant ideas but this time it's pretty rational
# as this library has an option for a wrapper-like frontend for funtions
# we want to spawn in upon starting up (Nobody likes writing redundant things over and over again believe me -_-)
def event(
self,
delay: timedelta | float | None = None,
instant: bool = False,
repeat: bool = False,
name: str | None = None,
traces:list[TraceEvent] = []
) -> AsyncWrapper[P, R_cov]:
"""Creates a wrapper for calling a group of context arguments from later"""
# filter delay down to only allowing timedelta or None
if isinstance(delay, (float, int)):
delay = timedelta(seconds=delay)
elif delay and not isinstance(delay, timedelta):
raise TypeError(
f"delay requires a float in seconds or a datetime.timedelta type not {type(delay)}"
)
if instant:
if not repeat:
raise ValueError("setting instant without repeat is invalid")
if not delay:
raise ValueError("repeating and being instant without a delay is invalid")
def wrapper(func: AsyncFunc[P, R_cov]) -> AsyncFunc[P, R_cov]:
if not inspect.iscoroutinefunction(func):
raise TypeError(
f"only asynchronous actions can be added not {type(func)}"
)
for t in traces:
t.freeze()
self._queue.append(
Event(func, time=delay, name=name, instant=instant, repeat=repeat, is_wrapper=True, traces=traces)
)
return func
return wrapper
def enter(
self,
action: Callable[..., Awaitable[Any]],
args: Sequence[Any],
kwargs: dict[str, Any],
delay: timedelta | int | float | None = None,
instant: bool = False,
repeat: bool = False,
name: str | None = None,
traces:list[TraceEvent[Unpack[PosArgsT], R_cov]] = []
) -> 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, this number could also be a
timedelta or float time (timedelta is preferred for extra persision)
"""
if not inspect.iscoroutinefunction(action):
raise TypeError(
f"only asynchronous actions can be added not {type(action)}"
)
if isinstance(delay, (float, int)):
delay = timedelta(seconds=delay)
elif delay and not isinstance(delay, timedelta):
raise TypeError(
f"delay requires a float in seconds or a datetime.timedelta type not {type(delay)}"
)
if kwargs is None:
kwargs = {}
if instant:
if not repeat:
raise ValueError("setting instant without repeat is invalid")
if not delay:
raise ValueError("repeating and being instant without a delay is invalid")
# Freeze all callbacks we wish to have used
for t in traces:
t.freeze()
return Event(action, args, kwargs, delay, name, instant, repeat, is_wrapper=False)
def cancel(self, event: Event):
"""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.
"""
self._queue.remove(event)
def empty(self) -> bool:
"""Check whether the queue is empty."""
return not self._queue
async def run(self, *args:Unpack[PosArgsT]) -> None:
"""Execute events until the queue is empty otherwise this function shall loop forever...
:param args: A list of arguments to pass to all TraceEvent listeners and sheduled tasks
that were added from using a wrapper structure.
Good case scenario for this is arguments that were parsed from a commandline-parser such as
argparse, click, typer, asyncclick or anyio-typer
"""
self.__update_state(_State.RUNNING)
async with RunningScheduler(self._queue, *args) as running:
# run everything and let closure take care of the joining portions...
await running.run()
self.__update_state(_State.OFF)
def start(self, *args:Unpack[PosArgsT], backend: str = "asyncio", backend_options: dict[str, Any] | None = None) -> None:
"""uses `anyio.run` to call `start(*args)` to begin execution of the application's task setups
:param args: A list of arguments to pass to all TraceEvent listeners and sheduled tasks
that were added from using a wrapper structure.
Good case scenario for this is arguments that were parsed from a commandline-parser such as
argparse, click, typer, asyncclick or anyio-typer
NOTE: it's smarter to use the non-asynchronous versions if this approch
"""
return anyio.run(self.run, *args, backend=backend, backend_options=backend_options)
@property
def queue(self) -> list[Event[Unpack[PosArgsT], R]]:
"""An ordered list of upcoming events.
Events are dataclasses with fields for:
action, priority, action, argss, kwargs
"""
return list(self._queue[:])
# This was mainly used for illistrating what it is that I wanted to develop for a user to end up working with.
# it is clean and very easy to navigate without too many problems, you can typehint important arguments
# to send accross multiple different callbacks without being overly complex and having
# too much logic exposed.
from anyio_sched import Scheduler, TraceEvent
from typer import run
from httpx import AsyncClient
sched: Sheduler[str] = Scheduler()
trace_event = TraceEvent()
@trace_event.on_start
async def on_start(event, host:str):
print(f"{event.name} has started get request to {host}")
@trace_event.on_finish
async def on_end(event, data:bytes, host:str):
print(f"Got: {data} from {host} after calling {event.name}")
@sched.event(traces=[trace_event])
async def simple_get_request(host:str):
async with AsyncClient() as client:
resp = await client.get(host)
data = await resp.aread()
return data
def main(host: str):
return sched.start(host, backend='asyncio')
if __name__ == "__main__":
run(main)
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment