Created
October 31, 2025 20:01
-
-
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.
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 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 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
| # 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