concurrency_utils
Useful components related to asyncio, multithreading or multiprocessing.
Module
Functions
asyncnullcontext
async def asyncnullcontext() ‑> collections.abc.AsyncGenerator[None, None]:Async version of contextlib.nullcontext().
await_event_with_stop
async def await_event_with_stop( wait_event: threading.Event, stop_event: threading.Event, wait_event_name: str, polling_timeout: float = 5,) ‑> bool:Helper function that waits on a Threading Event but allows exiting early.
Avoids blocking the async event loop.
Monitors the wait_event but with a polling timeout, so it can periodically check if it should stop early.
Arguments
wait_event: The event to wait to be set.stop_event: An event indicating we should stop early.wait_event_name: The name of the wait event, used for debugging.polling_timeout: The amount of time in seconds between checks for whether the stop_event has been set.
Returns
True if the wait_event has been set, False if the stop_event is
used to cancel waiting before the wait_event is set.
await_threading_event
async def await_threading_event( event: threading.Event, event_name: str, timeout: float | None = None, polling_timeout: float = 5,) ‑> bool:Event.wait() that doesn't block the event loop.
Avoids saturation of the ThreadPoolExecutor by either limiting run time directly
(via timeout) or by occassionally rescheduling the wait so other thread-based
tasks get a chance to run (via polling_timeout).
If both timeout and polling_timeout are supplied, timeout takes precedence
and no polling/rescheduling will be performed.
Arguments
event: The event to wait on.event_name: The name of the event, used in debug logging.timeout: The timeout to wait for the event to be set.polling_timeout: How frequently to return control to this function from the thread execution. If the event is still not set, the thread execution is rescheduled. Useful for handling cancellation events as it places a bound on the time a thread can run.
Returns
True if the event is set within any timeout constraints specified,
False otherwise.
pooled_breadth_first
def pooled_breadth_first( root: _C, visit: Callable[[_C], tuple[_T, Iterable[_C]]], *, workers: int, admit: Callable[[_C], bool] | None = None,) ‑> collections.abc.Generator[~_T, None, None]:Walk a tree breadth-first, visiting up to workers nodes at once.
For trees whose nodes cost a round trip to expand, such as directories on a network share: each node is expanded by visit on a pool thread as soon as its parent has named it, rather than one node after another. Results are yielded as nodes finish, so their order is not defined.
admit runs on the consuming thread, so it may keep unsynchronised state, such as the set of symlink targets already followed. Closing the generator early cancels the queued visits and waits for the running ones.
Arguments
root: The node to start from.visit: Expands a node into(result, children). Runs on the pool, so it must be safe to run concurrently and should not raise: an exception it raises propagates to the consumer and ends the walk.workers: How many nodes are expanded at once.admit: When set, a child is expanded only if this returnsTruefor it.
prefetch_map
def prefetch_map( fn: Callable[[_T], _R], items: Iterable[_T], *, workers: int, window: int | None = None,) ‑> collections.abc.Generator[tuple[~_T, typing.Union[~_R, Exception]], None, None]:Apply fn to items on a thread pool, yielding results in input order.
Built for I/O that waits on round trips, such as file reads or stats on a
network share: up to workers calls wait at once instead of one after
another. At most window items are in flight, so a long items is never
submitted all at once and results do not pile up unconsumed.
An exception raised by fn is yielded as that item's result rather than raised, so one bad item neither stops the others nor loses its place.
fn runs on the pool's threads; whatever the consumer does with a result runs on the consuming thread. Closing the generator early cancels the queued calls and waits for the running ones.
items is owned once passed: when this generator ends or is closed, so is
items if it has a close(). A lazy source with work of its own, such as a
pooled directory walk, therefore stops as soon as its consumer does rather
than whenever it is garbage-collected.
Arguments
fn: The call to make per item. Must be safe to run concurrently.items: The inputs, consumed lazily and closed when this generator is.workers: How many calls run at once.1or less calls fn serially on the consuming thread, with no pool.window: The most items submitted but not yet yielded. Defaults to twice workers; never less than workers.
run_on_daemon_thread
async def run_on_daemon_thread( func: Callable[..., Any], kwargs: dict[str, Any], thread_name: str,) ‑> Any:Run func on a daemon thread and await its result.
asyncio.to_thread submits to the event loop's default executor, whose
worker threads are not daemons. Cancelling the awaiting coroutine does not
stop a work item already running, and the thread is then joined twice over:
asyncio.run calls loop.shutdown_default_executor() on teardown, and the
interpreter joins surviving pool threads at exit. Long-running work
therefore held the process open for the rest of its run even once nothing
was waiting on it. A daemon thread is abandoned at interpreter exit instead.
Use this for work that a caller may need to walk away from. Work that must
finish belongs on asyncio.to_thread, whose join is the point.
The calling context is copied in, as asyncio.to_thread does.
Arguments
func: Callable to run off the loop.kwargs: Arguments to call it with.thread_name: Worker thread name, for debugging.
Returns Whatever func returned.
Raises
BaseException: Whatever func raised.
run_periodically
async def run_periodically( tick: Callable[[], Coroutine[Any, Any, Any]], interval_seconds: float, *, name: str, startup_delay: float = 0.0, stop_on: tuple[type[BaseException], ...] = (), cadence: Cadence = fixed,) ‑> None:Run tick every interval_seconds until cancelled.
The asyncio-native replacement for a "stoppable thread running its own event
loop" ticker. Shutdown is task.cancel() rather than setting a
threading.Event and hoping a join lands inside its timeout, and the caller
needs no thread, no second event loop, and no bridge between the two.
A tick that raises is logged and the loop continues. A ticker that died on
one bad tick would take its job down silently for the pod's whole lifetime,
which is the failure mode this guards. Pass stop_on for the exceptions
that should end the loop.
cadence decides what interval_seconds is measured from; see Cadence.
Arguments
tick: One pass of the periodic work.interval_seconds: Delay between ticks; seecadencefor what it is measured from.name: Used in this loop's log lines to identify it.startup_delay: Delay before the first tick. Defaults to 0, i.e. the first tick runs immediately.stop_on: Exceptions that end the loop instead of being logged and swallowed. Re-raised to the caller.cadence: Whetherinterval_secondsis measured from the start of a tick (Cadence.FIXED, the default) or the end of one (Cadence.SPACED).
Raises
asyncio.CancelledError: When the task running this loop is cancelled.
Classes
Cadence
class Cadence(*args, **kwds):What a periodic loop's interval is measured from.
Only distinguishable when a tick takes a comparable time to the interval; below that the two behave identically. The distinction exists because the pod's tickers and its DAG pollers want opposite things from a slow tick.
Ancestors
Variables
- static
FIXED- Measure from the start of a tick, making the interval a floor. A slow tick does not stretch the cadence. What a heartbeat wants: a 10s heartbeat should stay 10s whether or not the hub is slow.
- static
SPACED- Measure from the end of a tick, making the interval a gap. A slow tick is always followed by a full interval of quiet. What a poller against a shared service wants: an overrunning tick backs off rather than querying it back to back.
HoldCounter
class HoldCounter( *, on_first_acquire: Callable[[], None], on_last_release: Callable[[], None], name: str = 'resource',):Counts holders of a resource that must stay idle whilst anything holds it.
Acquiring never blocks. The count drives two side effects instead:
on_first_acquire as it rises off zero, on_last_release as it returns
there. A hold taken in between joins the suspension already in place. This
is a reference count with a callback at each end of its range, not a
semaphore -- nothing here waits on anything.
A count rather than a flag, because holders overlap: one can take a hold whilst another is part way through its own, and whichever finishes first must not wake the resource under the other.
The callbacks run outside the internal lock, so either is free to read
held or count, or to take a hold of its own.
Arguments
on_first_acquire: Called as the count rises from zero to one.on_last_release: Called as the count returns to zero.name: Used in the warning logged for an unbalanced release.
Variables
count : int- How many holders there are right now.
held : bool- Whether anything is holding at the moment.
Methods
acquire
def acquire(self) ‑> None:Take a hold, suspending the resource if this is the first.
hold
def hold(self) ‑> collections.abc.Generator[None, None, None]:Hold for the duration of the block.
For a holder whose lifetime is a scope. A holder that outlives the call
that starts it -- a task, a spawned run -- pairs acquire with a
release on whatever signals its end instead.
release
def release(self) ‑> None:Drop a hold, waking the resource if this is the last.
A release with no hold behind it is logged and ignored rather than taking the count negative, which would leave every later hold unable to reach zero again and the resource suspended for good.
ThreadWithException
class ThreadWithException( group=None, target=None, name=None, args=(), kwargs=None, *, daemon=None,):A thread subclass that captures exceptions to be reraised in the main thread.
This constructor should always be called with keyword arguments. Arguments are:
group should be None; reserved for future extension when a ThreadGroup class is implemented.
target is the callable object to be invoked by the run() method. Defaults to None, meaning nothing is called.
name is the thread name. By default, a unique name is constructed of the form "Thread-N" where N is a small decimal number.
args is a list or tuple of arguments for the target invocation. Defaults to ().
kwargs is a dictionary of keyword arguments for the target invocation. Defaults to .
If a subclass overrides the constructor, it must make sure to invoke the base class constructor (Thread.init()) before doing anything else to the thread.
Ancestors
Methods
join
def join(self, timeout: float | None = None) ‑> None:See parent method for documentation.
If an exception occurs in the joined thread it will be reraised in the calling thread.
run
def run(self) ‑> None:See parent method for documentation.
Captures exceptions raised during the run call and stores them as an attribute.