From aa0608bca3a5b50d8d78b812e837b366744eed3a Mon Sep 17 00:00:00 2001 From: Matt Wozniski Date: Wed, 9 Sep 2026 00:22:44 +0000 Subject: [PATCH] tests: Harden subprocess management We've seen the test suite hang occasionally in CI, and it's not clear what's going wrong. Add timeouts to some things that would wait indefinitely, and add extra logging to try to track down where things are going wrong. Signed-off-by: Matt Wozniski --- tests/integration/subinterpreters_program.py | 5 +- tests/integration/subinterpreters_shim.py | 55 +++++++++++-- .../subinterpreters_two_chains_program.py | 5 +- .../subinterpreters_with_threads_program.py | 5 +- tests/utils.py | 77 ++++++++++++++++--- 5 files changed, 124 insertions(+), 23 deletions(-) diff --git a/tests/integration/subinterpreters_program.py b/tests/integration/subinterpreters_program.py index 086fe9fe..b2c0edd4 100644 --- a/tests/integration/subinterpreters_program.py +++ b/tests/integration/subinterpreters_program.py @@ -3,6 +3,7 @@ import threading import time +from subinterpreters_shim import read_n from subinterpreters_shim import run_in_new_interpreter NUM_INTERPRETERS = 3 @@ -28,9 +29,7 @@ def start_interpreter_async(code): for _ in range(NUM_INTERPRETERS): start_interpreter_async(CODE) -data = b"" -while len(data) < NUM_INTERPRETERS: - data += os.read(r_fd, NUM_INTERPRETERS - len(data)) +read_n(r_fd, NUM_INTERPRETERS, timeout=30.0) os.close(r_fd) os.close(w_fd) diff --git a/tests/integration/subinterpreters_shim.py b/tests/integration/subinterpreters_shim.py index 3a009098..840c684b 100644 --- a/tests/integration/subinterpreters_shim.py +++ b/tests/integration/subinterpreters_shim.py @@ -1,20 +1,65 @@ +import os +import select +import sys +import time +import traceback + + +def abort(reason): + """Terminate this process immediately, with a diagnostic on stderr.""" + sys.stderr.write(f"subinterpreter test program failed: {reason}\n") + sys.stderr.flush() + os._exit(1) + + +def read_n(read_fd, count, *, timeout): + """Read exactly *count* bytes from *read_fd*.""" + deadline = time.monotonic() + timeout + data = b"" + while len(data) < count: + remaining = deadline - time.monotonic() + if remaining <= 0: + abort(f"only {len(data)} of {count} bytes arrived within {timeout} seconds") + if not select.select([read_fd], [], [], remaining)[0]: + continue + chunk = os.read(read_fd, count - len(data)) + if not chunk: + abort(f"EOF with {count - len(data)} bytes missing") + data += chunk + return data + + try: from concurrent import interpreters # type: ignore def run_in_new_interpreter(code): - interpreters.create().exec(code) + try: + interpreters.create().exec(code) + except BaseException: + traceback.print_exc(file=sys.stderr) + abort("Terminating due to the above exception") except ImportError: try: import _interpreters # type: ignore def run_in_new_interpreter(code): - _interpreters.exec(_interpreters.create(), code) + excinfo = _interpreters.exec(_interpreters.create(), code) + if excinfo is not None: + abort( + getattr(excinfo, "errdisplay", None) + or getattr(excinfo, "formatted", None) + or repr(excinfo) + ) except ImportError: import _xxsubinterpreters # type: ignore def run_in_new_interpreter(code): - _xxsubinterpreters.run_string( - _xxsubinterpreters.create(isolated=False), code - ) + try: + _xxsubinterpreters.run_string( + _xxsubinterpreters.create(isolated=False), code + ) + except BaseException: + traceback.print_exc(file=sys.stderr) + abort("Terminating due to the above exception") diff --git a/tests/integration/subinterpreters_two_chains_program.py b/tests/integration/subinterpreters_two_chains_program.py index ca984b69..aea3ce80 100644 --- a/tests/integration/subinterpreters_two_chains_program.py +++ b/tests/integration/subinterpreters_two_chains_program.py @@ -3,6 +3,7 @@ import threading import time +from subinterpreters_shim import read_n from subinterpreters_shim import run_in_new_interpreter r_fd, w_fd = os.pipe() @@ -28,9 +29,7 @@ def launch_chain(): t1.start() t2.start() -data = b"" -while len(data) < 2: - data += os.read(r_fd, 2 - len(data)) +read_n(r_fd, 2, timeout=30.0) os.close(r_fd) os.close(w_fd) diff --git a/tests/integration/subinterpreters_with_threads_program.py b/tests/integration/subinterpreters_with_threads_program.py index 22991993..57244c44 100644 --- a/tests/integration/subinterpreters_with_threads_program.py +++ b/tests/integration/subinterpreters_with_threads_program.py @@ -3,6 +3,7 @@ import threading import time +from subinterpreters_shim import read_n from subinterpreters_shim import run_in_new_interpreter NUM_INTERPRETERS = 2 @@ -48,9 +49,7 @@ def worker(): TOTAL_EXPECTED = NUM_INTERPRETERS * (NUM_THREADS_PER_SUBINTERPRETER + 1) -data = b"" -while len(data) < TOTAL_EXPECTED: - data += os.read(r_fd, TOTAL_EXPECTED - len(data)) +read_n(r_fd, TOTAL_EXPECTED, timeout=30.0) os.close(r_fd) os.close(w_fd) diff --git a/tests/utils.py b/tests/utils.py index 66a3d513..fc5eaebc 100644 --- a/tests/utils.py +++ b/tests/utils.py @@ -3,6 +3,7 @@ import itertools import os import pathlib +import select import shutil import subprocess import sys @@ -42,6 +43,70 @@ Interpreter = collections.namedtuple("Interpreter", "version path has_symbols") +def _format_child_output(output) -> str: + if isinstance(output, bytes): + return output.decode(errors="replace") + return output or "" + + +def _raise_child_startup_error(process: "subprocess.Popen", error: str) -> None: + exit_code = process.poll() + if exit_code is None: + process.terminate() + try: + stdout, stderr = process.communicate(timeout=1) + except subprocess.TimeoutExpired: + process.kill() + stdout, stderr = process.communicate(timeout=TIMEOUT) + status = "still running and was terminated by the test" + else: + stdout, stderr = process.communicate(timeout=TIMEOUT) + status = f"exited with status {exit_code}" + + raise AssertionError( + f"Child process {process.pid} {error}; it {status}.\n" + f"Command: {process.args!r}\n" + f"stdout:\n{_format_child_output(stdout)}\n" + f"stderr:\n{_format_child_output(stderr)}" + ) + + +def _wait_for_child_ready( + process: "subprocess.Popen", fifo: pathlib.Path, timeout: float = TIMEOUT +) -> None: + deadline = time.monotonic() + timeout + response = bytearray() + + read_fd = os.open(fifo, os.O_RDONLY | os.O_NONBLOCK) + # Ensure the FIFO has been opened for writing before we try to read it. + write_fd = os.open(fifo, os.O_WRONLY | os.O_NONBLOCK) + try: + while True: + if process.poll() is not None: + _raise_child_startup_error(process, "exited before reporting readiness") + + remaining = deadline - time.monotonic() + if remaining <= 0: + _raise_child_startup_error( + process, f"did not report readiness within {timeout:g} seconds" + ) + + readable, _, _ = select.select([read_fd], [], [], min(remaining, 0.1)) + if not readable: + continue + + response.extend(os.read(read_fd, 4096)) + if response == b"ready": + return + if response and not b"ready".startswith(response): + _raise_child_startup_error( + process, f"reported unexpected readiness value {bytes(response)!r}" + ) + finally: + os.close(write_fd) + os.close(read_fd) + + def find_all_available_pythons() -> Iterable[Interpreter]: # pragma: no cover versions: List[Tuple[Tuple[int, int], str]] test_version = os.getenv("PYTHON_TEST_VERSION") @@ -98,12 +163,9 @@ def spawn_child_process( stderr=subprocess.PIPE, text=True, ) as process: - with open(fifo, "r") as fifo_file: - response = fifo_file.read() - - assert response == "ready" - time.sleep(0.1) try: + _wait_for_child_ready(process, fifo) + time.sleep(0.1) yield process finally: os.remove(fifo) @@ -126,10 +188,7 @@ def generate_core_file( stdout=subprocess.PIPE, stderr=subprocess.PIPE, ) as process: - with open(fifo, "r") as fifo_file: - response = fifo_file.read() - - assert response == "ready" + _wait_for_child_ready(process, fifo) subprocess.run( ["gcore", str(process.pid)], check=True,