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
11 changes: 11 additions & 0 deletions Doc/library/asyncio-task.rst
Original file line number Diff line number Diff line change
Expand Up @@ -288,6 +288,17 @@ Creating tasks
# completion:
task.add_done_callback(background_tasks.discard)

Note that this approach never awaits the tasks, so if a task
fails, its exception is never retrieved and asyncio logs a
"Task exception was never retrieved" message when the task is
garbage collected. To avoid this, use :class:`asyncio.TaskGroup`
which keeps a strong reference to each task, awaits them and
propagates their exceptions::

async with asyncio.TaskGroup() as tg:
for i in range(10):
tg.create_task(some_coro(param=i))

.. versionadded:: 3.7

.. versionchanged:: 3.8
Expand Down
75 changes: 42 additions & 33 deletions Lib/asyncio/windows_events.py
Original file line number Diff line number Diff line change
Expand Up @@ -760,6 +760,46 @@ def _get_accept_socket(self, family):
s.settimeout(0)
return s

def _process_completion_status(self, status):
"""Process a single status from the completion port.

A caller that waits on the completion port itself can pass each
status it receives here.
"""
err, transferred, key, address = status
try:
f, ov, obj, callback = self._cache.pop(address)
except KeyError:
if self._loop.get_debug():
self._loop.call_exception_handler({
'message': ('GetQueuedCompletionStatus() returned an '
'unexpected event'),
'status': ('err=%s transferred=%s key=%#x address=%#x'
% (err, transferred, key, address)),
})

# key is either zero, or it is used to return a pipe
# handle which should be closed to avoid a leak.
if key not in (0, _overlapped.INVALID_HANDLE_VALUE):
_winapi.CloseHandle(key)
return

if obj in self._stopped_serving:
f.cancel()
# Don't call the callback if _register() already read the result or
# if the overlapped has been cancelled
elif not f.done():
try:
value = callback(transferred, key, ov)
except OSError as e:
f.set_exception(e)
self._results.append(f)
else:
f.set_result(value)
self._results.append(f)
finally:
f = None

def _poll(self, timeout=None):
if timeout is None:
ms = INFINITE
Expand All @@ -778,39 +818,8 @@ def _poll(self, timeout=None):
break
ms = 0

err, transferred, key, address = status
try:
f, ov, obj, callback = self._cache.pop(address)
except KeyError:
if self._loop.get_debug():
self._loop.call_exception_handler({
'message': ('GetQueuedCompletionStatus() returned an '
'unexpected event'),
'status': ('err=%s transferred=%s key=%#x address=%#x'
% (err, transferred, key, address)),
})

# key is either zero, or it is used to return a pipe
# handle which should be closed to avoid a leak.
if key not in (0, _overlapped.INVALID_HANDLE_VALUE):
_winapi.CloseHandle(key)
continue

if obj in self._stopped_serving:
f.cancel()
# Don't call the callback if _register() already read the result or
# if the overlapped has been cancelled
elif not f.done():
try:
value = callback(transferred, key, ov)
except OSError as e:
f.set_exception(e)
self._results.append(f)
else:
f.set_result(value)
self._results.append(f)
finally:
f = None
# gh-154971: split out so custom event loops can call it directly
self._process_completion_status(status)

# Remove unregistered futures
for ov in self._unregistered:
Expand Down
32 changes: 32 additions & 0 deletions Lib/test/_isolated_sample.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@
a subprocess. Several of these tests fail, error or are skipped on purpose.
"""

import atexit
import os
import time
import unittest
from test.support import isolation
Expand Down Expand Up @@ -109,3 +111,33 @@ class BrokenSubclassSample(SubclassingSample):
@classmethod
def setUpClass(cls):
pass


# The exit code the samples below die with, after their tests have run.
EXIT_CODE = 3


def _die_at_exit():
atexit.register(os._exit, EXIT_CODE)


class MethodExitSample(unittest.TestCase):

@isolation.runInSubprocess()
def test_passes_then_dies(self):
_die_at_exit()

@isolation.runInSubprocess()
def test_fails_and_dies(self):
_die_at_exit()
self.fail('the test itself failed')


@isolation.runInSubprocess()
class ClassExitSample(unittest.TestCase):

def test_pass(self):
pass

def test_dies(self):
_die_at_exit()
4 changes: 2 additions & 2 deletions Lib/test/list_tests.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
from functools import cmp_to_key

from test import seq_tests
from test.support import ALWAYS_EQ, NEVER_EQ, skip_if_huge_c_stack
from test.support import ALWAYS_EQ, NEVER_EQ, run_with_limited_c_stack
from test.support import skip_emscripten_stack_overflow, skip_wasi_stack_overflow


Expand Down Expand Up @@ -60,7 +60,7 @@ def test_repr(self):
self.assertEqual(str(a2), "[0, 1, 2, [...], 3]")
self.assertEqual(repr(a2), "[0, 1, 2, [...], 3]")

@skip_if_huge_c_stack(200_000)
@run_with_limited_c_stack(200_000)
@skip_wasi_stack_overflow()
@skip_emscripten_stack_overflow()
def test_repr_deep(self):
Expand Down
2 changes: 1 addition & 1 deletion Lib/test/mapping_tests.py
Original file line number Diff line number Diff line change
Expand Up @@ -629,7 +629,7 @@ def __repr__(self):
d = self._full_mapping({1: BadRepr()})
self.assertRaises(Exc, repr, d)

@support.skip_if_huge_c_stack()
@support.run_with_limited_c_stack()
@support.skip_wasi_stack_overflow()
@support.skip_emscripten_stack_overflow()
@support.skip_if_sanitizer("requires deep stack", ub=True)
Expand Down
95 changes: 80 additions & 15 deletions Lib/test/support/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@
"check_disallow_instantiation", "check_sanitizer", "skip_if_sanitizer",
"requires_limited_api", "requires_specialization", "thread_unsafe",
"skip_if_unlimited_stack_size", "skip_if_huge_c_stack",
"run_with_limited_c_stack",
# sys
"MS_WINDOWS", "is_jython", "is_android", "is_emscripten", "is_wasi",
"is_apple_mobile", "check_impl_detail", "unix_shell", "setswitchinterval",
Expand Down Expand Up @@ -2839,30 +2840,94 @@ def exceeds_recursion_limit():
return 150_000


def _has_huge_c_stack(depth):
"""Check that *depth* recursive calls cannot exhaust the C stack."""
try:
from _testinternalcapi import get_c_recursion_remaining
except ImportError:
# Fall back to checking for an unlimited stack size.
if is_emscripten or is_wasi or os.name == "nt":
return False
import resource
soft, hard = resource.getrlimit(resource.RLIMIT_STACK)
return soft == hard and soft in (-1, 0xFFFF_FFFF_FFFF_FFFF)
else:
remaining = get_c_recursion_remaining()
# A negative value means integer overflow in the estimate
# (e.g. with an unlimited RLIMIT_STACK). The estimate is based on
# the size of the interpreter loop frame, so it is only a lower
# bound for recursion with smaller C frames.
return remaining >= depth or remaining < 0


def skip_if_huge_c_stack(depth=150_000):
"""Skip decorator for tests which cannot overflow the C stack.

Tests exhausting the C stack with *depth* recursive calls cannot
trigger the recursion protection if the C stack is too large (e.g.
with a large or unlimited RLIMIT_STACK), and either fail, or run
for a very long time, or crash, or consume all memory.

Prefer run_with_limited_c_stack() for tests recursing to a fixed depth.
"""
try:
from _testinternalcapi import get_c_recursion_remaining
except ImportError:
# Fall back to checking for an unlimited stack size.
huge = False
if not (is_emscripten or is_wasi) and os.name != "nt":
import resource
soft, hard = resource.getrlimit(resource.RLIMIT_STACK)
huge = soft == hard and soft in (-1, 0xFFFF_FFFF_FFFF_FFFF)
else:
remaining = get_c_recursion_remaining()
# A negative value means integer overflow in the estimate
# (e.g. with an unlimited RLIMIT_STACK).
huge = remaining >= depth or remaining < 0
return unittest.skipIf(
huge, f"the C stack is large enough for {depth} recursive calls")
_has_huge_c_stack(depth),
f"the C stack is large enough for {depth} recursive calls")


# Small enough to be exhausted by tens of thousands of recursive calls,
# but not smaller than Py_C_STACK_SIZE (4 MiB) which the interpreter
# assumes if it cannot query the thread stack size.
C_STACK_SIZE = 8 * 1024 * 1024


def run_with_limited_c_stack(depth=150_000, size=C_STACK_SIZE):
"""Decorator for tests exhausting the C stack with *depth* recursive calls.

Run the test in a separate thread with the C stack of *size* bytes, so
that the outcome does not depend on the C stack size of the main thread
(which can be large or unlimited, see RLIMIT_STACK).

If a thread with the limited C stack cannot be created, run the test in
the current thread, but skip it if the C stack is too large.
"""
reason = f"the C stack is large enough for {depth} recursive calls"
def decorator(test):
@functools.wraps(test)
def wrapper(*args, **kwargs):
def run_test():
# The C stack can still be too large if limiting it failed.
if _has_huge_c_stack(depth):
raise unittest.SkipTest(reason)
test(*args, **kwargs)

try:
import threading
old_size = threading.stack_size(size)
except (ImportError, ValueError, RuntimeError):
# Setting the thread stack size is not supported.
return run_test()

exceptions = []
def run():
try:
run_test()
except BaseException as exc:
exceptions.append(exc)

thread = threading.Thread(target=run)
try:
thread.start()
except RuntimeError:
# Threads are not supported.
return run_test()
finally:
threading.stack_size(old_size)
thread.join()
if exceptions:
raise exceptions[0]
return wrapper
return decorator


# Windows doesn't have os.uname() but it doesn't support s390x.
Expand Down
25 changes: 22 additions & 3 deletions Lib/test/support/isolation.py
Original file line number Diff line number Diff line change
Expand Up @@ -163,6 +163,16 @@ def _raise_fixture_outcome(outcome):
raise exc from _remote(outcome['detail'])


def _check_returncode(returncode, output, what):
# The subprocess writes its result before exiting, so a non-zero exit code
# means it died afterwards, during finalization, unnoticed by the result.
if returncode:
exc = _SubprocessTestError(
f'the subprocess exited with code {returncode} '
f'after running the {what}')
raise exc from _remote(output)


def _isolate_method(func):
@functools.wraps(func)
def wrapper(self, /, *args, **kwargs):
Expand All @@ -180,7 +190,9 @@ def wrapper(self, /, *args, **kwargs):
raise exc from _remote(output)
# The parent measures this method's own duration (the real cost of the
# isolated run, subprocess startup included), so nothing to forward here.
# Replay the outcomes first: a failure of the test itself is more useful.
_replay_outcomes(self, payload['outcomes'])
_check_returncode(returncode, output, 'test')
return wrapper


Expand Down Expand Up @@ -219,13 +231,20 @@ def setUpClass(cls):
by_id.setdefault(outcome['id'], []).append(outcome)
cls._isolated_outcomes = by_id
cls._isolated_durations = dict(payload.get('durations', ()))
# Report the crash from tearDownClass(), after replaying the outcomes.
cls._isolated_exit = (returncode, output)

def tearDownClass(cls):
if runningInSubprocess:
orig_tearDownClass(cls)
else:
cls._isolated_outcomes = None
cls._isolated_durations = None
return
cls._isolated_outcomes = None
cls._isolated_durations = None
# Missing if an overriding setUpClass() bypassed the subprocess.
exited = getattr(cls, '_isolated_exit', None)
cls._isolated_exit = None
if exited is not None:
_check_returncode(*exited, 'class')

def _callSetUp(self):
# In the parent the real test does not run, so neither should setUp().
Expand Down
3 changes: 2 additions & 1 deletion Lib/test/test_ast/test_ast.py
Original file line number Diff line number Diff line change
Expand Up @@ -1025,7 +1025,8 @@ def next(self):
enum._test_simple_enum(_Precedence, _ast_unparse._Precedence)

@support.cpython_only
@support.skip_if_huge_c_stack(100_000 if sys.platform == "android" else 500_000)
@support.run_with_limited_c_stack(
100_000 if sys.platform == "android" else 500_000)
@skip_wasi_stack_overflow()
@skip_emscripten_stack_overflow()
def test_ast_recursion_limit(self):
Expand Down
22 changes: 22 additions & 0 deletions Lib/test/test_asyncio/test_windows_events.py
Original file line number Diff line number Diff line change
Expand Up @@ -327,6 +327,28 @@ def threadMain():
stop.set()
thr.join()

def test_custom_poll_integration(self):
# gh-154971: a caller can wait on the completion port and process statuses itself
proactor = self.loop._proactor

a, b = socket.socketpair()
self.addCleanup(a.close)
self.addCleanup(b.close)

fut = proactor.recv(a, 100)
self.assertFalse(fut.done())

b.send(b'data')

deadline = time.monotonic() + support.SHORT_TIMEOUT
while not fut.done() and time.monotonic() < deadline:
status = _overlapped.GetQueuedCompletionStatus(proactor._iocp, 100)
if status is not None:
proactor._process_completion_status(status)

self.assertTrue(fut.done())
self.assertEqual(fut.result(), b'data')


class ProactorPipeObjectSupportTests(unittest.TestCase):

Expand Down
Loading
Loading