Last active
June 12, 2026 14:49
-
-
Save x42005e1f/4f18c3c62da9135020bdea8c44c248a2 to your computer and use it in GitHub Desktop.
A fast future implementation (thread-safe, signal-safe, non-blocking)
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
| #!/usr/bin/env python3 | |
| # SPDX-FileCopyrightText: 2026 Ilya Egorov <0x42005e1f@gmail.com> | |
| # SPDX-License-Identifier: ISC | |
| # requires-python = ">=3.8" | |
| # dependencies = [ | |
| # "exceptiongroup>=1.3.0; python_version<'3.11'", | |
| # "typing-extensions>=4.5.0; python_version<'3.12'", | |
| # ] | |
| from __future__ import annotations | |
| import sys | |
| import weakref | |
| from abc import ABC, abstractmethod | |
| from logging import getLogger | |
| from typing import TYPE_CHECKING, Generic, TypeVar | |
| if TYPE_CHECKING: | |
| from typing import Any | |
| if sys.version_info >= (3, 11): # PEP 654 | |
| from builtins import BaseExceptionGroup | |
| else: # exceptiongroup>=1.3.0 | |
| from exceptiongroup import BaseExceptionGroup | |
| if TYPE_CHECKING: | |
| if sys.version_info >= (3, 9): # PEP 585 | |
| from collections.abc import Callable | |
| else: | |
| from typing import Callable | |
| if sys.version_info >= (3, 11): # python/cpython#30842 | |
| from typing import Never | |
| else: # typing-extensions>=4.1.0 | |
| from typing_extensions import Never | |
| if sys.version_info >= (3, 11): # PEP 673 | |
| from typing import Self | |
| else: # typing-extensions>=4.0.0 | |
| from typing_extensions import Self | |
| if sys.version_info >= (3, 11): # PEP 646 | |
| from typing import TypeVarTuple, Unpack | |
| else: # typing-extensions>=4.1.0 | |
| from typing_extensions import TypeVarTuple, Unpack | |
| if sys.version_info >= (3, 11): # python/cpython#30530: introspectable | |
| from typing import final | |
| else: # typing-extensions>=4.1.0 | |
| from typing_extensions import final | |
| if sys.version_info >= (3, 12): # PEP 698 | |
| from typing import override | |
| else: # typing-extensions>=4.5.0 | |
| from typing_extensions import override | |
| _T = TypeVar("_T") | |
| _T_co = TypeVar("_T_co", covariant=True) | |
| if TYPE_CHECKING: | |
| _Ts = TypeVarTuple("_Ts") | |
| _LOGGER = getLogger(__name__) | |
| _sentinel = object() | |
| class CancelledError(Exception): | |
| """ | |
| Exception raised when the :meth:`Future.get` method is called after the | |
| call was successfully cancelled. | |
| """ | |
| class InvalidStateError(Exception): | |
| """ | |
| Exception raised when an operation is performed on a :class:`Future` | |
| instance that is not allowed in the current state. | |
| """ | |
| class Handle(ABC, Generic[_T_co]): | |
| """ | |
| A base class for :class:`UnboundHandle` and :class:`WeakBoundHandle`. | |
| Allows to track and control the state of the call; it is thread-safe, | |
| signal-safe, and non-blocking. | |
| It is typically the result of calling :meth:`future.add_done_callback() | |
| <Future.add_done_callback>`. | |
| .. note:: | |
| If the handle was not cancelled or finished before it was destroyed, | |
| the event is logged as an error. | |
| """ | |
| __slots__ = ( | |
| "__weakref__", | |
| "_args", | |
| "_callback", | |
| "_cancelled", | |
| "_not_done", | |
| "_not_removed", | |
| "_not_running", | |
| "_owner_ref", | |
| ) | |
| def __reduce__(self, /) -> Never: | |
| cls = type(self) | |
| cls_repr = f"{cls.__module__}.{cls.__qualname__}" | |
| msg = f"cannot reduce {cls_repr} objects" | |
| raise TypeError(msg) | |
| def __repr__(self, /) -> str: | |
| cls = type(self) | |
| cls_repr = f"{cls.__module__}.{cls.__qualname__}" | |
| owner = self._owner_ref() | |
| if owner is None: | |
| owner_repr = "<? object (dead)>" | |
| else: | |
| owner_cls = type(owner) | |
| owner_cls_repr = f"{owner_cls.__module__}.{owner_cls.__qualname__}" | |
| owner_repr = f"<{owner_cls_repr} object at {id(owner):#x}>" | |
| args_repr = ", ".join([ | |
| repr(self._callback), | |
| owner_repr, | |
| *map(repr, self._args), | |
| ]) | |
| if getattr(self, "_not_done", False): | |
| cancelled = self._cancelled | |
| if cancelled is True: | |
| state = "cancelled" | |
| elif cancelled is False: | |
| state = "running" | |
| else: | |
| state = "pending" | |
| else: | |
| state = "finished" | |
| return f"<{cls_repr}({args_repr}) at {id(self):#x}: {state}>" | |
| def __del__(self, /) -> None: | |
| obj_repr = repr(self).rpartition(":")[0] + ">" | |
| if not self.done(): | |
| if self.running(): | |
| msg = "%s was destroyed but it is running!" | |
| else: | |
| msg = "%s was destroyed but it is pending!" | |
| _LOGGER.error(msg, obj_repr, stacklevel=2) | |
| @abstractmethod | |
| def _run(self, /) -> None: | |
| raise NotImplementedError | |
| def run(self, /) -> bool: | |
| """ | |
| Execute the associated callback once. | |
| If the call was successfully cancelled, is currently being executed, or | |
| has finished running, then the method will return :data:`False`; | |
| otherwise, the method will put the handle into the running state, | |
| execute the call, and return :data:`True`. | |
| Any return value will be suppressed. And any raised exception that is | |
| an instance of the :exc:`Exception` class or any of its subclasses will | |
| be suppressed and logged. Everything else, such as :exc:`SystemExit` | |
| and :exc:`KeyboardInterrupt`, will be propagated as is. | |
| """ | |
| cancelled = self._cancelled | |
| if isinstance(cancelled, list): | |
| if not cancelled: | |
| cancelled.append(False) | |
| self._cancelled = cancelled = cancelled[0] | |
| if cancelled is True: | |
| return False | |
| # One point that is not immediately obvious is where this piece of code | |
| # should be placed. If the method returns `False`, then either | |
| # `handle.running()` or `handle.done()` must return `True`. If the code | |
| # were placed at the beginning (which would result in some | |
| # optimization), the behavior would not be observed in a race | |
| # condition. | |
| try: | |
| del self._not_running | |
| except AttributeError: # not the first call | |
| return False | |
| try: | |
| self._run() | |
| except Exception: | |
| msg = "exception from %r" | |
| _LOGGER.exception(msg, self) | |
| finally: | |
| del self._not_done | |
| return True | |
| def cancel(self, /) -> bool: | |
| """ | |
| Attempt to cancel the call. | |
| If the call is currently being executed or finished running and cannot | |
| be cancelled then the method will return :data:`False`, otherwise the | |
| call will be cancelled and the method will return :data:`True`. Note, | |
| if the call has already been cancelled, the method will also return | |
| :data:`True`. | |
| """ | |
| cancelled = self._cancelled | |
| if isinstance(cancelled, list): | |
| if not cancelled: | |
| cancelled.append(True) | |
| self._cancelled = cancelled = cancelled[0] | |
| return cancelled is True | |
| def cancelled(self, /) -> bool: | |
| """ | |
| Return :data:`True` if the call was successfully cancelled, | |
| :data:`False` otherwise. | |
| """ | |
| return self._cancelled is True | |
| def running(self, /) -> bool: | |
| """ | |
| Return :data:`True` if the call is currently being executed and cannot | |
| be cancelled, :data:`False` otherwise. | |
| """ | |
| return self._cancelled is False and getattr(self, "_not_done", False) | |
| def done(self, /) -> bool: | |
| """ | |
| Return :data:`True` if the call was successfully cancelled or finished | |
| running, :data:`False` otherwise. | |
| """ | |
| return self._cancelled is True or not getattr(self, "_not_done", False) | |
| @property | |
| def owner(self, /) -> _T_co: | |
| """ | |
| The object that created the handle (typically, an instance of the | |
| :class:`Future` or one of its subclasses). | |
| Raises: | |
| ReferenceError: | |
| if the object is garbage collected. | |
| """ | |
| owner = self._owner_ref() | |
| if owner is None: | |
| msg = "weakly-referenced object no longer exists" | |
| raise ReferenceError(msg) | |
| return owner | |
| @final | |
| class UnboundHandle(Handle[_T_co]): | |
| """ | |
| Represents the execution of a *callback* unbound from a weakly-referenced | |
| *owner*: :meth:`run` will call ``callback(*args)``. | |
| """ | |
| __slots__ = () | |
| def __init__( | |
| self, | |
| callback: Callable[[Unpack[_Ts]], Any], | |
| owner: _T_co, | |
| /, | |
| *args: Unpack[_Ts], | |
| ) -> None: | |
| # A list that will be replaced with a specific literal (`False`, or | |
| # `True`) when the handle is put into a different state ("running", or | |
| # "cancelled") . The result is always the first element of the list, | |
| # which ensures thread-safety. | |
| self._cancelled = [] | |
| # One-time attributes that are deleted when the condition is no longer | |
| # true within some code block. Since an attribute can only be deleted | |
| # once, this also ensures thread-safety and is used for entry control. | |
| self._not_removed = True # used by `Future.remove_done_callback()` | |
| self._not_running = True | |
| self._not_done = True | |
| self._callback = callback | |
| self._args = args | |
| self._owner_ref = weakref.ref(owner) | |
| def __init_subclass__(cls, /, **kwargs: Any) -> Never: | |
| bcs = __class__ # an implicit closure reference | |
| bcs_name = bcs.__name__ | |
| msg = f"type {bcs_name!r} is not an acceptable base type" | |
| raise TypeError(msg) | |
| @override | |
| def _run(self, /) -> None: | |
| self._callback(*self._args) | |
| @final | |
| class WeakBoundHandle(Handle[_T_co]): | |
| """ | |
| Represents the execution of a *callback* bound to a weakly-referenced | |
| *owner*: :meth:`run` will call ``callback(owner, *args)``. | |
| """ | |
| __slots__ = () | |
| def __init__( | |
| self, | |
| callback: Callable[[_T_co, Unpack[_Ts]], Any], | |
| owner: _T_co, | |
| /, | |
| *args: Unpack[_Ts], | |
| ) -> None: | |
| # A list that will be replaced with a specific literal (`False`, or | |
| # `True`) when the handle is put into a different state ("running", or | |
| # "cancelled") . The result is always the first element of the list, | |
| # which ensures thread-safety. | |
| self._cancelled = [] | |
| # One-time attributes that are deleted when the condition is no longer | |
| # true within some code block. Since an attribute can only be deleted | |
| # once, this also ensures thread-safety and is used for entry control. | |
| self._not_removed = True # used by `Future.remove_done_callback()` | |
| self._not_running = True | |
| self._not_done = True | |
| self._callback = callback | |
| self._args = args | |
| self._owner_ref = weakref.ref(owner) | |
| def __init_subclass__(cls, /, **kwargs: Any) -> Never: | |
| bcs = __class__ # an implicit closure reference | |
| bcs_name = bcs.__name__ | |
| msg = f"type {bcs_name!r} is not an acceptable base type" | |
| raise TypeError(msg) | |
| @override | |
| def _run(self, /) -> None: | |
| owner = self._owner_ref() | |
| if owner is None: | |
| msg = "weakly-referenced object no longer exists" | |
| raise ReferenceError(msg) | |
| self._callback(owner, *self._args) | |
| class Future(Generic[_T]): | |
| """ | |
| Represents the execution (and eventual result) of an asynchronous | |
| operation (a call). | |
| Designed for use in :abbr:`SPSC (single-producer, single-consumer)` | |
| scenarios; it is thread-safe, signal-safe, and non-blocking. | |
| It is similar to :class:`concurrent.futures.Future`, but is faster, safer, | |
| and more flexible. | |
| .. note:: | |
| If the future was not cancelled or finished (and used) before it was | |
| destroyed, the event is logged as an error, including information about | |
| the result/exception (if available). | |
| """ | |
| __slots__ = ( | |
| "__weakref__", | |
| "_callbacks", | |
| "_cancelled", | |
| "_not_cancelled", | |
| "_not_running", | |
| "_not_set", | |
| "_outcome", | |
| ) | |
| def __init__(self, /) -> None: | |
| super().__init__() # for cooperative multiple inheritance | |
| # A list that will be replaced with a specific literal (`None`, | |
| # `False`, or `True`) when the future is put into a different state | |
| # ("finished", "running", or "cancelled") . The result is always the | |
| # first element of the list, which ensures thread-safety. | |
| self._cancelled = [] | |
| # One-time attributes that are deleted when the condition is no longer | |
| # true within some code block. Since an attribute can only be deleted | |
| # once, this also ensures thread-safety and is used for entry control. | |
| self._not_running = True | |
| self._not_cancelled = True | |
| self._not_set = True | |
| self._outcome = None # deleted when used by the consumer | |
| self._callbacks = [] # deleted when finished/cancelled by the producer | |
| def __reduce__(self, /) -> Never: | |
| cls = type(self) | |
| cls_repr = f"{cls.__module__}.{cls.__qualname__}" | |
| msg = f"cannot reduce {cls_repr} objects" | |
| raise TypeError(msg) | |
| def __repr__(self, /) -> str: | |
| cls = type(self) | |
| cls_repr = f"{cls.__module__}.{cls.__qualname__}" | |
| outcome = getattr(self, "_outcome", _sentinel) | |
| if outcome is None: | |
| cancelled = self._cancelled | |
| if cancelled is None: | |
| state = "finished" | |
| elif cancelled is True: | |
| state = "cancelled" | |
| elif cancelled is False: | |
| state = "running" | |
| else: | |
| state = "pending" | |
| elif outcome is _sentinel: | |
| state = "finished (used)" | |
| else: | |
| result, exception = outcome | |
| if exception is None: | |
| state = f"finished (returned {type(result).__name__})" | |
| else: | |
| state = f"finished (raised {type(exception).__name__})" | |
| return f"<{cls_repr} object at {id(self):#x}: {state}>" | |
| def __del__(self, /) -> None: | |
| cls = type(self) | |
| cls_repr = f"{cls.__module__}.{cls.__qualname__}" | |
| obj_repr = f"<{cls_repr} object at {id(self):#x}>" | |
| if not self.done(): | |
| if self.running(): | |
| msg = "%s was destroyed but it is running!" | |
| else: | |
| msg = "%s was destroyed but it is pending!" | |
| _LOGGER.error(msg, obj_repr, stacklevel=2) | |
| elif not self.cancelled() and not self.used(): | |
| result, exception = self._outcome | |
| if exception is None: | |
| msg = "%s's result was never retrieved:\n%r" | |
| _LOGGER.error(msg, obj_repr, result, stacklevel=2) | |
| else: | |
| msg = "%s's exception was never retrieved:" | |
| _LOGGER.error(msg, obj_repr, exc_info=exception, stacklevel=2) | |
| def get(self, /) -> _T: | |
| """ | |
| Return the result (or raise the exception) set by the :meth:`set` | |
| method. | |
| Raises: | |
| InvalidStateError: | |
| if the call has not yet finished running or the method was already | |
| called. | |
| CancelledError: | |
| if the call was successfully cancelled. | |
| """ | |
| try: | |
| outcome = self._outcome | |
| if outcome is None and self._cancelled is not True: | |
| msg = "future is not done" | |
| raise InvalidStateError(msg) | |
| del self._outcome | |
| except AttributeError: # not the first call | |
| msg = "future is already used" | |
| raise InvalidStateError(msg) from None | |
| if self._cancelled is True: | |
| msg = "future is cancelled" | |
| raise CancelledError(msg) | |
| try: | |
| result, exception = outcome | |
| if exception is None: | |
| return result | |
| try: | |
| raise exception | |
| finally: | |
| del exception # break reference cycles | |
| finally: | |
| del outcome # break reference cycles | |
| def set( | |
| self, | |
| result: _T | None, | |
| exception: BaseException | None = None, | |
| /, | |
| ) -> None: | |
| """ | |
| Set the *result* (or the *exception*) of the work associated with the | |
| future (by the producer). | |
| Regardless of whether :meth:`set_running_or_notify_cancel` was called, | |
| the method will put the future into the finished state and invoke all | |
| attached callbacks. | |
| Raises: | |
| InvalidStateError: | |
| if the call was successfully cancelled or finished running. | |
| BaseExceptionGroup: | |
| if any callback raised a :exc:`BaseException` but not an | |
| :exc:`Exception`. | |
| """ | |
| cancelled = self._cancelled | |
| if isinstance(cancelled, list): | |
| if not cancelled: | |
| cancelled.append(None) | |
| self._cancelled = cancelled = cancelled[0] | |
| if cancelled is True: | |
| msg = "future is cancelled" | |
| raise InvalidStateError(msg) | |
| try: | |
| del self._not_set | |
| except AttributeError: # not the first call | |
| msg = "future is already finished" | |
| raise InvalidStateError(msg) from None | |
| self._outcome = (result, exception) | |
| self._invoke_callbacks() | |
| def set_running_or_notify_cancel(self, /) -> bool: | |
| """ | |
| Attempt to indicate that the call is currently being executed (by the | |
| producer). | |
| If the call was successfully cancelled then the method will invoke all | |
| attached callbacks and return :data:`False`, otherwise the method will | |
| put the future into the running state and return :data:`True`. | |
| Raises: | |
| InvalidStateError: | |
| if the call has finished running or the method was already called. | |
| BaseExceptionGroup: | |
| if any callback raised a :exc:`BaseException` but not an | |
| :exc:`Exception`. | |
| """ | |
| cancelled = self._cancelled | |
| if isinstance(cancelled, list): | |
| if not cancelled: | |
| cancelled.append(False) | |
| self._cancelled = cancelled = cancelled[0] | |
| try: | |
| del self._not_running | |
| except AttributeError: # not the first call | |
| msg = "method is already called" | |
| raise InvalidStateError(msg) from None | |
| if cancelled is None: | |
| msg = "future is already finished" | |
| raise InvalidStateError(msg) | |
| if cancelled is True: | |
| try: | |
| del self._not_cancelled | |
| except AttributeError: # a race condition | |
| pass | |
| else: | |
| self._invoke_callbacks() | |
| return False | |
| return True | |
| def abort(self, /) -> bool: | |
| """ | |
| Attempt to cancel the call (by the producer). | |
| If the call is currently being executed or finished running and cannot | |
| be cancelled then the method will return :data:`False`, otherwise the | |
| call will be cancelled and the method will return :data:`True`. Note, | |
| if the call has already been cancelled, the method will also return | |
| :data:`True`. | |
| Unlike the :meth:`cancel` method, it invokes callbacks. | |
| Raises: | |
| BaseExceptionGroup: | |
| if any callback raised a :exc:`BaseException` but not an | |
| :exc:`Exception`. | |
| """ | |
| cancelled = self._cancelled | |
| if isinstance(cancelled, list): | |
| if not cancelled: | |
| cancelled.append(True) | |
| self._cancelled = cancelled = cancelled[0] | |
| if cancelled is True: | |
| try: | |
| del self._not_cancelled | |
| except AttributeError: # a race condition | |
| pass | |
| else: | |
| self._invoke_callbacks() | |
| return True | |
| return False | |
| def cancel(self, /) -> bool: | |
| """ | |
| Attempt to cancel the call. | |
| If the call is currently being executed or finished running and cannot | |
| be cancelled then the method will return :data:`False`, otherwise the | |
| call will be cancelled and the method will return :data:`True`. Note, | |
| if the call has already been cancelled, the method will also return | |
| :data:`True`. | |
| Unlike the standard library, it does not invoke callbacks. | |
| """ | |
| cancelled = self._cancelled | |
| if isinstance(cancelled, list): | |
| if not cancelled: | |
| cancelled.append(True) | |
| self._cancelled = cancelled = cancelled[0] | |
| return cancelled is True | |
| def cancelled(self, /) -> bool: | |
| """ | |
| Return :data:`True` if the call was successfully cancelled, | |
| :data:`False` otherwise. | |
| """ | |
| return self._cancelled is True | |
| def running(self, /) -> bool: | |
| """ | |
| Return :data:`True` if the call is currently being executed and cannot | |
| be cancelled, :data:`False` otherwise. | |
| """ | |
| return ( | |
| self._cancelled is False | |
| and getattr(self, "_outcome", _sentinel) is None | |
| ) | |
| def done(self, /) -> bool: | |
| """ | |
| Return :data:`True` if the call was successfully cancelled or finished | |
| running, :data:`False` otherwise. | |
| """ | |
| return ( | |
| self._cancelled is True | |
| or getattr(self, "_outcome", _sentinel) is not None | |
| ) | |
| def used(self, /) -> bool: | |
| """ | |
| Return :data:`True` if the :meth:`get` was successfully called, | |
| :data:`False` otherwise. | |
| """ | |
| return not hasattr(self, "_outcome") | |
| def add_done_callback( | |
| self, | |
| callback: Callable[[Unpack[_Ts]], Any], | |
| /, | |
| *args: Unpack[_Ts], | |
| ) -> Handle[Self]: | |
| """ | |
| Attach the *callback* to the future with the specified *args*. | |
| When the call was successfully cancelled or finished running, callbacks | |
| are executed by the methods for the producer in the order that they | |
| were added. The returned :class:`Handle` instance allows to track the | |
| state of the callback and, if necessary, either cancel it or execute it | |
| early. | |
| If the condition has already been met, the callback will be executed | |
| immediately (regardless of whether the previous attached callbacks have | |
| already been executed). | |
| """ | |
| if args and args[0] is self: | |
| handle = WeakBoundHandle(callback, *args) | |
| else: | |
| handle = UnboundHandle(callback, self, *args) | |
| try: | |
| callbacks = self._callbacks | |
| except AttributeError: # future is done | |
| pass | |
| else: | |
| callbacks.append(handle) | |
| if hasattr(self, "_callbacks"): | |
| return handle | |
| # future is done (a race condition) | |
| handle.run() | |
| return handle | |
| def remove_done_callback(self, handle: Handle[Any], /) -> bool: | |
| """ | |
| Detach the callback (represented by the *handle*) from the future. | |
| The method has the same effect as calling :meth:`handle.cancel() | |
| <Handle.cancel>` and returns the same value. | |
| Unlike :meth:`handle.cancel() <Handle.cancel>`, it actually removes the | |
| handle from the callback list when possible, which can be used for | |
| memory optimization. | |
| Raises: | |
| ValueError: | |
| if the handle belongs to another object. | |
| """ | |
| if handle.owner is self: | |
| # Removing a callback from the list is a relatively expensive | |
| # operation, especially due to the friction it creates with the | |
| # worker thread, but without it, any temporary registrations would | |
| # lead to memory leaks. This method is implemented as a complement | |
| # to `handle.cancel()`, removing the callback from the list when | |
| # possible and returning the same result. Of course, other | |
| # behaviors could be devised, but the current one seems the most | |
| # justified. | |
| cancelled = handle.cancel() | |
| try: | |
| # O(nm) -> O(n), where | |
| # n = len(self._callbacks), | |
| # m = number of calls (or concurrent threads). | |
| del handle._not_removed | |
| except AttributeError: # not the first call | |
| pass | |
| else: | |
| try: | |
| callbacks = self._callbacks | |
| except AttributeError: # future is done | |
| pass | |
| else: | |
| try: | |
| callbacks.remove(handle) | |
| except ValueError: # future is done (a race condition) | |
| pass | |
| return cancelled | |
| msg = "handle belongs to another object" | |
| raise ValueError(msg) | |
| def _invoke_callbacks(self, /) -> None: | |
| # The current implementation does not guarantee any deterministic order | |
| # for running callbacks. In fact, there is an alternative strategy in | |
| # which, after the worker thread exits, each handle is processed while | |
| # holding its own once lock, and is removed from the `callbacks` by the | |
| # first thread that acquires the lock. However, this would lead to | |
| # increased memory overhead (in particular, the `callbacks` would have | |
| # to be implemented as a deque), synchronization overhead, and re-entry | |
| # issues. Therefore, we use a mixed strategy in which the FIFO order is | |
| # maintained for old callbacks, while for concurrent ones (when the | |
| # future is done) it remains undefined (subject to race conditions). | |
| callbacks = self._callbacks | |
| del self._callbacks | |
| callbacks.reverse() # LIFO -> FIFO | |
| # Handling errors from callbacks is a more serious problem than it | |
| # might seem at first glance. Callbacks can be used to register wakeups | |
| # and other events that are important for proper finalization. | |
| # Therefore, we must execute all callbacks regardless of what | |
| # exceptions they raise, but still raise them at the end to avoid | |
| # unnecessary (and incorrect) suppression. And by using exception | |
| # groups, we provide the user with context information, which can be | |
| # useful for debugging. | |
| exceptions = [] | |
| while callbacks: | |
| try: | |
| handle = callbacks.pop() | |
| except IndexError: # a race condition | |
| break | |
| try: | |
| handle.run() | |
| except BaseException as exc: # noqa: BLE001 | |
| exceptions.append(exc) | |
| if exceptions: | |
| # If the user interrupts callbacks, and each interrupted callback | |
| # successfully propagates a `KeyboardInterrupt`, then propagating | |
| # any `KeyboardInterrupt` further would be more intuitive than | |
| # using an exception group instead. However, context also matters, | |
| # so we set the exception group as the cause of our | |
| # `KeyboardInterrupt` when there is more than one exception. | |
| if all(isinstance(exc, KeyboardInterrupt) for exc in exceptions): | |
| try: | |
| if len(exceptions) == 1: | |
| raise exceptions[0] | |
| else: | |
| msg = "unhandled errors from callbacks" | |
| exceptions = BaseExceptionGroup(msg, exceptions) | |
| raise KeyboardInterrupt from exceptions | |
| finally: | |
| del exceptions # break reference cycles | |
| else: | |
| try: | |
| msg = "unhandled errors from callbacks" | |
| exceptions = BaseExceptionGroup(msg, exceptions) | |
| raise exceptions | |
| finally: | |
| del exceptions # break reference cycles |
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
| #!/usr/bin/env python3 | |
| # SPDX-FileCopyrightText: 2026 Ilya Egorov <0x42005e1f@gmail.com> | |
| # SPDX-License-Identifier: ISC | |
| import sys | |
| from abc import ABC, abstractmethod | |
| from typing import Any, Generic, TypeVar | |
| if sys.version_info >= (3, 9): # PEP 585 | |
| from collections.abc import Callable | |
| else: | |
| from typing import Callable | |
| if sys.version_info >= (3, 11): # python/cpython#30842 | |
| from typing import Never | |
| else: # typing-extensions>=4.1.0 | |
| from typing_extensions import Never | |
| if sys.version_info >= (3, 11): # PEP 673 | |
| from typing import Self | |
| else: # typing-extensions>=4.0.0 | |
| from typing_extensions import Self | |
| if sys.version_info >= (3, 11): # PEP 646 | |
| from typing import TypeVarTuple, Unpack | |
| else: # typing-extensions>=4.1.0 | |
| from typing_extensions import TypeVarTuple, Unpack | |
| if sys.version_info >= (3, 11): # python/cpython#30530: introspectable | |
| from typing import final | |
| else: # typing-extensions>=4.1.0 | |
| from typing_extensions import final | |
| if sys.version_info >= (3, 12): # PEP 698 | |
| from typing import override | |
| else: # typing-extensions>=4.5.0 | |
| from typing_extensions import override | |
| _T = TypeVar("_T") | |
| _T_co = TypeVar("_T_co", covariant=True) | |
| _Ts = TypeVarTuple("_Ts") | |
| class CancelledError(Exception): ... | |
| class InvalidStateError(Exception): ... | |
| class Handle(ABC, Generic[_T_co]): | |
| __slots__ = ( | |
| "__weakref__", | |
| "_args", | |
| "_callback", | |
| "_cancelled", | |
| "_not_done", | |
| "_not_removed", | |
| "_not_running", | |
| "_owner_ref", | |
| ) | |
| def __reduce__(self, /) -> Never: ... | |
| def __repr__(self, /) -> str: ... | |
| def __del__(self, /) -> None: ... | |
| @abstractmethod | |
| def _run(self, /) -> None: ... | |
| def run(self, /) -> bool: ... | |
| def cancel(self, /) -> bool: ... | |
| def cancelled(self, /) -> bool: ... | |
| def running(self, /) -> bool: ... | |
| def done(self, /) -> bool: ... | |
| @property | |
| def owner(self, /) -> _T_co: ... | |
| @final | |
| class UnboundHandle(Handle[_T_co]): | |
| __slots__ = () | |
| def __init__( | |
| self, | |
| callback: Callable[[Unpack[_Ts]], Any], | |
| owner: _T_co, | |
| /, | |
| *args: Unpack[_Ts], | |
| ) -> None: ... | |
| def __init_subclass__(cls, /, **kwargs: Any) -> Never: ... | |
| @override | |
| def _run(self, /) -> None: ... | |
| @final | |
| class WeakBoundHandle(Handle[_T_co]): | |
| __slots__ = () | |
| def __init__( | |
| self, | |
| callback: Callable[[_T_co, Unpack[_Ts]], Any], | |
| owner: _T_co, | |
| /, | |
| *args: Unpack[_Ts], | |
| ) -> None: ... | |
| def __init_subclass__(cls, /, **kwargs: Any) -> Never: ... | |
| @override | |
| def _run(self, /) -> None: ... | |
| class Future(Generic[_T]): | |
| __slots__ = ( | |
| "__weakref__", | |
| "_callbacks", | |
| "_cancelled", | |
| "_not_cancelled", | |
| "_not_running", | |
| "_not_set", | |
| "_outcome", | |
| ) | |
| def __init__(self, /) -> None: ... | |
| def __reduce__(self, /) -> Never: ... | |
| def __repr__(self, /) -> str: ... | |
| def __del__(self, /) -> None: ... | |
| def get(self, /) -> _T: ... | |
| def set( | |
| self, | |
| result: _T | None, | |
| exception: BaseException | None = None, | |
| /, | |
| ) -> None: ... | |
| def set_running_or_notify_cancel(self, /) -> bool: ... | |
| def abort(self, /) -> bool: ... | |
| def cancel(self, /) -> bool: ... | |
| def cancelled(self, /) -> bool: ... | |
| def running(self, /) -> bool: ... | |
| def done(self, /) -> bool: ... | |
| def used(self, /) -> bool: ... | |
| def add_done_callback( | |
| self, | |
| callback: Callable[[Unpack[_Ts]], Any], | |
| /, | |
| *args: Unpack[_Ts], | |
| ) -> Handle[Self]: ... | |
| def remove_done_callback(self, handle: Handle[Any], /) -> bool: ... | |
| def _invoke_callbacks(self, /) -> None: ... |
Author
Author
As I now see it, future objects actually combine two distinct entities. The first entity is responsible for setting the result by the worker. The second entity is responsible for presenting the result to the waiter. This violates the single-responsibility principle. If we want to have safe future objects, these two responsibilities should be separated. Similar to how pipes are represented.
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Here is one possible implementation of future objects for Python. I could go into detail about how it came about and all that, but... I am not very good at writing (writing comments & docstrings has been particularly labor-intensive for me). So let us skip straight to the notable features.
Unlike the standard library,
fasture.Futureis completely non-blocking. It uses no locks — neither under the hood (lockless, even effectively wait-free design) nor for implementing blocking methods (there simply are not any!). Therefore, if you want to wait for completion, you will have to manually use callbacks, for example with aiologic:But that is not cheap. It would be more efficient to override the
_invoke_callbacks()method to handle low-level aiologic events specifically, and add the correspondingwait()and__await__()methods. You can try doing this in your subclass.And now for something completely different:
future.cancel(), callbacks will only be invoked after the producer calls eitherfuture.set_running_or_notify_cancel()orfuture.abort()(the latter is relevant if the producer no longer wants to invoke callbacks; however, this is somewhat dangerous behavior, so do not abuse it).future.cancel(), so callbacks of a newer object may be executed before callbacks of an older one.fasture.Futureis not like that; itsfuture._invoke_callbacks()is always called only on the producer side.future._invoke_callback()executes all callbacks regardless of whether they completed successfully or not, even in the case of keyboard interrupts. In particular, if you add tentime.sleep(10)calls as callbacks and want to interruptfuture._invoke_callback(), you will have to press Control-C all ten times. This might be a bit surprising, but that is the price of safe cancellation.future.add_done_callback()returns a special object similar to asyncio loop's handles. Hmm. I am not sure if it is worth describing it here. You are better off reading the docstrings.As for speed, this means that attempts to interact with an object in any way will not cause friction between threads (your
time.sleep(10)callback will not block parallel state checks, use of the object, or wakeups in general, unlikeconcurrent.futures.Future). Of course, single-threaded use also offers performance gains (a 17x speedup was observed on PyPy for a simple produce-consume pattern), but handles create additional overhead, and therefore in normal scenarios,fasture.Futuremay perform slower thanconcurrent.futures.Future, and it is not even comparable toasyncio.Future. Usefasture.Futurenot for speed, but for the features it offers (in particular, it can be safely used from within signal handlers and destructors).Memory usage by one instance (on CPython 3.12)
concurrent.futures.Futureasyncio.Futurefasture.Handlefasture.Future