Visitar URL original
Update concurrent.futures for Python 3.14 by 1ndahous3 · Pull Request #8904 · RustPython/RustPython · GitHub
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 21 additions & 10 deletions Lib/concurrent/futures/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@
wait,
as_completed)

__all__ = (
__all__ = [
'FIRST_COMPLETED',
'FIRST_EXCEPTION',
'ALL_COMPLETED',
Expand All @@ -31,24 +31,35 @@
'as_completed',
'ProcessPoolExecutor',
'ThreadPoolExecutor',
)
]


try:
import _interpreters
except ImportError:
_interpreters = None

if _interpreters:
__all__.append('InterpreterPoolExecutor')


def __dir__():
return __all__ + ('__author__', '__doc__')
return __all__ + ['__author__', '__doc__']


def __getattr__(name):
global ProcessPoolExecutor, ThreadPoolExecutor
global ProcessPoolExecutor, ThreadPoolExecutor, InterpreterPoolExecutor

if name == 'ProcessPoolExecutor':
from .process import ProcessPoolExecutor as pe
ProcessPoolExecutor = pe
return pe
from .process import ProcessPoolExecutor
return ProcessPoolExecutor

if name == 'ThreadPoolExecutor':
from .thread import ThreadPoolExecutor as te
ThreadPoolExecutor = te
return te
from .thread import ThreadPoolExecutor
return ThreadPoolExecutor

if _interpreters and name == 'InterpreterPoolExecutor':
from .interpreter import InterpreterPoolExecutor
return InterpreterPoolExecutor

raise AttributeError(f"module {__name__!r} has no attribute {name!r}")
146 changes: 86 additions & 60 deletions Lib/concurrent/futures/_base.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@
import threading
import time
import types
import weakref
from itertools import islice

FIRST_COMPLETED = 'FIRST_COMPLETED'
FIRST_EXCEPTION = 'FIRST_EXCEPTION'
Expand All @@ -23,14 +25,6 @@
CANCELLED_AND_NOTIFIED = 'CANCELLED_AND_NOTIFIED'
FINISHED = 'FINISHED'

_FUTURE_STATES = [
PENDING,
RUNNING,
CANCELLED,
CANCELLED_AND_NOTIFIED,
FINISHED
]

_STATE_TO_DESCRIPTION_MAP = {
PENDING: "pending",
RUNNING: "running",
Expand Down Expand Up @@ -200,15 +194,15 @@ def as_completed(fs, timeout=None):
"""An iterator over the given futures that yields each as it completes.

Args:
fs: The sequence of Futures (possibly created by different Executors) to
iterate over.
timeout: The maximum number of seconds to wait. If None, then there
is no limit on the wait time.
fs: The sequence of Futures (possibly created by different
Executors) to iterate over.
timeout: The maximum number of seconds to wait. If None, then
there is no limit on the wait time.

Returns:
An iterator that yields the given Futures as they complete (finished or
cancelled). If any given Futures are duplicated, they will be returned
once.
An iterator that yields the given Futures as they complete
(finished or cancelled). If any given Futures are duplicated,
they will be returned once.

Raises:
TimeoutError: If the entire result iterator could not be generated
Expand Down Expand Up @@ -264,19 +258,20 @@ def wait(fs, timeout=None, return_when=ALL_COMPLETED):
"""Wait for the futures in the given sequence to complete.

Args:
fs: The sequence of Futures (possibly created by different Executors) to
wait upon.
timeout: The maximum number of seconds to wait. If None, then there
is no limit on the wait time.
return_when: Indicates when this function should return. The options
are:
fs: The sequence of Futures (possibly created by different
Executors) to wait upon.
timeout: The maximum number of seconds to wait. If None, then
there is no limit on the wait time.
return_when: Indicates when this function should return.
The options are:

FIRST_COMPLETED - Return when any future finishes or is
cancelled.
FIRST_EXCEPTION - Return when any future finishes by raising an
exception. If no future raises an exception
exception. If no future raises an exception
then it is equivalent to ALL_COMPLETED.
ALL_COMPLETED - Return when all futures finish or are cancelled.
ALL_COMPLETED - Return when all futures finish or are
cancelled.

Returns:
A named 2-tuple of sets. The first set, named 'done', contains the
Expand Down Expand Up @@ -410,11 +405,12 @@ def add_done_callback(self, fn):

Args:
fn: A callable that will be called with this future as its only
argument when the future completes or is cancelled. The callable
will always be called by a thread in the same process in which
it was added. If the future has already completed or been
cancelled then the callable will be called immediately. These
callables are called in the order that they were added.
argument when the future completes or is cancelled. The
callable will always be called by a thread in the same
process in which it was added. If the future has already
completed or been cancelled then the callable will be
called immediately. These callables are called in the
order that they were added.
"""
with self._condition:
if self._state not in [CANCELLED, CANCELLED_AND_NOTIFIED, FINISHED]:
Expand All @@ -429,17 +425,19 @@ def result(self, timeout=None):
"""Return the result of the call that the future represents.

Args:
timeout: The number of seconds to wait for the result if the future
isn't done. If None, then there is no limit on the wait time.
timeout: The number of seconds to wait for the result if the
future isn't done. If None, then there is no limit on the
wait time.

Returns:
The result of the call that the future represents.

Raises:
CancelledError: If the future was cancelled.
TimeoutError: If the future didn't finish executing before the given
timeout.
Exception: If the call raised then that exception will be raised.
TimeoutError: If the future didn't finish executing before the
given timeout.
Exception: If the call raised then that exception will be
raised.
"""
try:
with self._condition:
Expand All @@ -465,17 +463,17 @@ def exception(self, timeout=None):

Args:
timeout: The number of seconds to wait for the exception if the
future isn't done. If None, then there is no limit on the wait
time.
future isn't done. If None, then there is no limit on the
wait time.

Returns:
The exception raised by the call that the future represents or None
if the call completed without raising.
The exception raised by the call that the future represents or
None if the call completed without raising.

Raises:
CancelledError: If the future was cancelled.
TimeoutError: If the future didn't finish executing before the given
timeout.
TimeoutError: If the future didn't finish executing before the
given timeout.
"""

with self._condition:
Expand All @@ -500,22 +498,23 @@ def set_running_or_notify_cancel(self):
Should only be used by Executor implementations and unit tests.

If the future has been cancelled (cancel() was called and returned
True) then any threads waiting on the future completing (though calls
to as_completed() or wait()) are notified and False is returned.
True) then any threads waiting on the future completing (though
calls to as_completed() or wait()) are notified and False is
returned.

If the future was not cancelled then it is put in the running state
(future calls to running() will return True) and True is returned.

This method should be called by Executor implementations before
executing the work associated with this future. If this method returns
False then the work should not be executed.
executing the work associated with this future. If this method
returns False then the work should not be executed.

Returns:
False if the Future was cancelled, True otherwise.

Raises:
RuntimeError: if this method was already called or if set_result()
or set_exception() was called.
RuntimeError: if this method was already called or if
set_result() or set_exception() was called.
"""
with self._condition:
if self._state == CANCELLED:
Expand Down Expand Up @@ -572,40 +571,61 @@ class Executor(object):
def submit(self, fn, /, *args, **kwargs):
"""Submits a callable to be executed with the given arguments.

Schedules the callable to be executed as fn(*args, **kwargs) and returns
a Future instance representing the execution of the callable.
Schedules the callable to be executed as fn(*args, **kwargs) and
returns a Future instance representing the execution of the
callable.

Returns:
A Future representing the given call.
"""
raise NotImplementedError()

def map(self, fn, *iterables, timeout=None, chunksize=1):
def map(self, fn, *iterables, timeout=None, chunksize=1, buffersize=None):
"""Returns an iterator equivalent to map(fn, iter).

Args:
fn: A callable that will take as many arguments as there are
passed iterables.
timeout: The maximum number of seconds to wait. If None, then there
is no limit on the wait time.
chunksize: The size of the chunks the iterable will be broken into
before being passed to a child process. This argument is only
used by ProcessPoolExecutor; it is ignored by
timeout: The maximum number of seconds to wait. If None, then
there is no limit on the wait time.
chunksize: The size of the chunks the iterable will be broken
into before being passed to a child process. This argument
is only used by ProcessPoolExecutor; it is ignored by
ThreadPoolExecutor.
buffersize: The number of submitted tasks whose results have not
yet been yielded. If the buffer is full, iteration over the
iterables pauses until a result is yielded from the buffer.
If None, all input elements are eagerly collected, and
a task is submitted for each.

Returns:
An iterator equivalent to: map(func, *iterables) but the calls may
be evaluated out-of-order.
An iterator equivalent to: map(func, *iterables) but the calls
may be evaluated out-of-order.

Raises:
TimeoutError: If the entire result iterator could not be generated
before the given timeout.
TimeoutError: If the entire result iterator could not be
generated before the given timeout.
Exception: If fn(*args) raises for any values.
"""
if buffersize is not None and not isinstance(buffersize, int):
raise TypeError("buffersize must be an integer or None")
if buffersize is not None and buffersize < 1:
raise ValueError("buffersize must be None or > 0")

if timeout is not None:
end_time = timeout + time.monotonic()

fs = [self.submit(fn, *args) for args in zip(*iterables)]
zipped_iterables = zip(*iterables)
if buffersize:
fs = collections.deque(
self.submit(fn, *args) for args in islice(zipped_iterables, buffersize)
)
else:
fs = [self.submit(fn, *args) for args in zipped_iterables]

# Use a weak reference to ensure that the executor can be garbage
# collected independently of the result_iterator closure.
executor_weakref = weakref.ref(self)

# Yield must be hidden in closure so that the futures are submitted
# before the first iterator value is required.
Expand All @@ -614,6 +634,12 @@ def result_iterator():
# reverse to keep finishing order
fs.reverse()
while fs:
if (
buffersize
and (executor := executor_weakref())
and (args := next(zipped_iterables, None))
):
fs.appendleft(executor.submit(fn, *args))
# Careful not to keep a reference to the popped future
if timeout is None:
yield _result_or_cancel(fs.pop())
Expand All @@ -632,8 +658,8 @@ def shutdown(self, wait=True, *, cancel_futures=False):

Args:
wait: If True then shutdown will not return until all running
futures have finished executing and the resources used by the
executor have been reclaimed.
futures have finished executing and the resources used by
the executor have been reclaimed.
cancel_futures: If True then shutdown will cancel all pending
futures. Futures that are completed or running will not be
cancelled.
Expand Down
Loading
Loading