Skip to main content

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 returns True for 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. 1 or 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; see cadence for 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: Whether interval_seconds is 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.

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.

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.