From 64dff37edef89bf351a0beb36f12b2e03a03c284 Mon Sep 17 00:00:00 2001 From: lmoresi Date: Wed, 2 Sep 2026 22:49:18 -0700 Subject: [PATCH] Supervise every parallel batch: bound a collective hang and name the rank that caused it A rank blocked inside a collective produces no output and never returns, so an unsupervised batch spends the whole job budget in silence and is cancelled from outside with no diagnosis. Measured at 76 minutes on one run (#675), which is also how test_0855 and test_0873 managed to fail without leaving a traceback. Every mechanism we already had asks the stuck process to cooperate, which is the one thing it has stopped being able to do. A threading.Timer needs the GIL and fired zero times on blocked ranks. faulthandler.dump_traceback_later is a C thread walking other threads' live frames and can wedge or crash the process it is meant to diagnose (#661). MPI_Abort from a watchdog thread is not safe from inside a collective. mpirun --timeout is spelled differently by Open MPI and MPICH and tells you nothing when it fires. scripts/mpi_supervisor.py asks nothing of the job. It is a parent process holding a clock and a signal, so it behaves the same under either MPI and on either platform: it is POSIX process management and knows nothing about MPI. It triggers on SILENCE rather than elapsed time. A batch that is still printing is alive, and judging on total runtime would mean inventing a wall-clock budget for a matrix nobody has measured -- wrong the first time someone adds a slow test. The #675 failure was 76 minutes of nothing, not of slow progress. On silence it arms each rank through a usercustomize module that registers a SIGUSR1 handler (a signal handler runs on the blocked thread itself and never walks live frames concurrently, so it cannot reproduce #661), samples twice, and reports three independent pieces of evidence: whether anything moved between samples (stuck vs slow), which rank is somewhere the others are not, and CPU per rank. It then kills descendants individually before the group -- mpirun puts its children in a group of their own and killpg has been seen to refuse with EPERM on macOS -- and re-checks for survivors, because an unsupervised kill orphans mpirun and pytest children at 100% CPU (#639). Validated by a planted hang whose answer is known: rank 0 sits in a function no other rank can be inside while the rest wait in a barrier. The supervisor must end the job, call it stuck, and name rank 0. The control also asserts the other direction, that a healthy job passes through with its own exit status -- without it, every other assertion could be met by a supervisor that kills everything. At np=2 there is no majority, and the report says so rather than nominating a rank on a one-all split. Two findings worth recording from writing it: the spinning ranks are the INNOCENT ones (an MPI barrier busy-waits) while the guilty rank sleeps, and an earlier version that guessed the rank-to-CPU mapping reported that backwards -- the handler now records its own pid so the attribution is read, not inferred. tests/test_0063_mpi_hang_supervisor.py: 3 passed in 109s. Underworld development team with AI support from Claude Code Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01E87Q7KrpapxeQiLD1RiNXv --- docs/developer/guides/mpi-hang-supervision.md | 117 +++++ docs/developer/index.md | 1 + scripts/mpi_supervisor.py | 400 ++++++++++++++++++ scripts/test.sh | 12 +- .../rank_zero_misses_the_collective.py | 50 +++ tests/test_0063_mpi_hang_supervisor.py | 113 +++++ 6 files changed, 691 insertions(+), 2 deletions(-) create mode 100644 docs/developer/guides/mpi-hang-supervision.md create mode 100755 scripts/mpi_supervisor.py create mode 100644 tests/parallel/hang_controls/rank_zero_misses_the_collective.py create mode 100644 tests/test_0063_mpi_hang_supervisor.py diff --git a/docs/developer/guides/mpi-hang-supervision.md b/docs/developer/guides/mpi-hang-supervision.md new file mode 100644 index 000000000..50ae9c0e5 --- /dev/null +++ b/docs/developer/guides/mpi-hang-supervision.md @@ -0,0 +1,117 @@ +# Bounding and diagnosing an MPI collective hang + +A rank blocked inside a collective is the worst kind of test failure: it produces +no output, returns no status, and takes the whole job with it. In CI it costs the +entire job budget and then gets cancelled from outside, so you learn nothing at +all. That happened for 76 minutes on a single run ([#675]). + +`scripts/mpi_supervisor.py` bounds that, and says what went wrong. + +## Why not the mechanisms we already had + +Everything we tried first asks the stuck process to cooperate, which is precisely +what it has stopped being able to do. + +| mechanism | why it fails | +|---|---| +| `threading.Timer` and a report | needs the GIL. Measured at np=4 against a 4 s block in `comm.allreduce`: a re-arming timer at 0.5 s fired **zero** times on the blocked ranks. | +| `faulthandler.dump_traceback_later` | a C thread walking other threads' live frames. It can wedge or crash the process it is diagnosing ([#661]). | +| `MPI_Abort` from a watchdog thread | needs `THREAD_MULTIPLE` and a schedulable thread; not safe from inside a collective. | +| `mpirun --timeout` | Open MPI spells it `--timeout`, MPICH spells it `-timeout`, and the semantics differ. Kills the job without telling you anything. | + +The supervisor asks nothing of the job. It is a parent process holding a clock +and a signal, so it behaves identically under Open MPI and MPICH, on macOS and +Linux, because it is POSIX process management and knows nothing about MPI. + +## Usage + +```bash +scripts/mpi_supervisor.py --silence 300 -- mpirun -n 4 python -m pytest --with-mpi tests/parallel/ +``` + +The exit status is the job's own, or `124` if the supervisor had to kill it — +the same number `timeout(1)` uses, so it reads the way people expect. + +`scripts/test.sh` already routes its parallel batches through it. Override the +budget with `PARALLEL_SILENCE` if a batch legitimately goes quiet for longer. + +## It watches for silence, not for elapsed time + +A batch that is still printing is alive. Judging on total runtime would mean +inventing a wall-clock budget for a test matrix nobody has measured, and that +number is wrong the first time somebody adds a slow test. Silence is both the +stronger signal and the one that needs no calibration — the failure in [#675] +was 76 minutes of *nothing*, not 76 minutes of slow progress. + +## What it reports + +On silence it signals every rank to dump its Python stack, samples twice, and +compares: + +``` +=== MPI SUPERVISOR: no output for 300 s === + no rank moved in 20 s: this is stuck, not slow + rank(s) 0 are somewhere the other 3 are not — that is where to look first + rank 0 is in _rank_zero_misses_the_collective() + the other 3 are in main() + rank 0: 0% CPU, idle — blocked rather than polling + rank 1: 100% CPU, spinning — a busy-wait, which is what an MPI progress engine does +``` + +Three separate pieces of evidence, and each answers a different question: + +- **Two samples, not one.** Identical stacks 20 s apart mean stuck; one stack + only ever means slow. +- **The odd rank out.** The ranks that agree are waiting for the one that does + not. At np=2 there is no majority to appeal to, and the report says so rather + than nominating a rank on a one-all split. +- **CPU per rank.** A rank inside a busy-waiting MPI progress engine burns a + whole core; one blocked on a lock sleeps. Note the direction of that signal: + in the example above the *spinning* ranks are the innocent ones, waiting in + the barrier, and the idle rank is the guilty one. + +Stacks are written to `--artifacts` so CI can keep them; a diagnosis that dies +with the job is no better than the silence it replaced. + +## How the dump is armed + +The supervisor writes a `usercustomize` module onto the job's `PYTHONPATH`, and +that registers a `SIGUSR1` handler through `faulthandler.register`. The job's own +code knows nothing about it. + +Registering a *signal handler* is the important difference from +`dump_traceback_later`. The handler runs on the thread that receives the signal, +only when the supervisor asks, and never walks live frames from a concurrent +thread — so arming it costs a healthy run nothing and it cannot reproduce +[#661]. + +## Killing the tree, and checking + +Killing the launcher is not killing the job. Descendants are enumerated and +killed individually before the process group is signalled, because the group +call is the one that fails: `mpirun` puts its children in a group of their own, +and `killpg` on macOS has been seen to refuse with `EPERM`. + +The supervisor then re-checks and reports anything that survived. This is not +ceremony — an unsupervised kill leaves orphaned `mpirun`, `hydra_pmi_proxy` and +pytest children reparented to init and spinning at 100% CPU, which is how a +"finished" run keeps eating cores for hours ([#639]). + +## The control + +`tests/test_0063_mpi_hang_supervisor.py` plants a hang whose answer is known: +rank 0 sits in a function no other rank can be inside while the others wait in a +barrier. The supervisor must end the job, report it as stuck, and name rank 0. + +It also asserts the other direction — a healthy job passes through untouched +with its own exit status. Without that, every other assertion could be satisfied +by a supervisor that simply kills everything it is given. + +A harness that catches hangs is worth nothing until it has caught one, so the +planted hang runs whenever the supervisor changes. It lives outside the +`test_*.py` pattern so pytest can never collect it by accident: running it under +a harness that does not kill it is the exact failure this all exists to prevent. + +[#639]: https://github.com/underworldcode/underworld3/issues/639 +[#661]: https://github.com/underworldcode/underworld3/issues/661 +[#675]: https://github.com/underworldcode/underworld3/issues/675 diff --git a/docs/developer/index.md b/docs/developer/index.md index c32976e1f..b88e30519 100644 --- a/docs/developer/index.md +++ b/docs/developer/index.md @@ -137,6 +137,7 @@ guides/branching-strategy guides/state-as-dataclass guides/BINDER_CONTAINER_SETUP guides/hpc-cluster-setup +guides/mpi-hang-supervision ``` ```{toctree} diff --git a/scripts/mpi_supervisor.py b/scripts/mpi_supervisor.py new file mode 100755 index 000000000..0e1f5c110 --- /dev/null +++ b/scripts/mpi_supervisor.py @@ -0,0 +1,400 @@ +#!/usr/bin/env python3 +"""Run an MPI job under a supervisor that bounds a collective hang and diagnoses it. + +A rank blocked inside a collective cannot help you. It does not reach a +``threading.Timer`` (measured: a re-arming timer at 0.5 s fired zero times on +blocked ranks), it cannot be reasoned with by ``MPI_Abort`` from another +thread, and ``faulthandler.dump_traceback_later`` — a C thread walking other +threads' live frames — can wedge or crash the process it is meant to diagnose +(underworld3#661). Every one of those mechanisms asks the stuck process to +cooperate, which is the one thing it has stopped being able to do. + +This supervisor asks nothing of it. It is a parent process holding a clock and +a signal, so it works the same under Open MPI and MPICH, on macOS and Linux, +because it is POSIX process management and knows nothing about MPI. + +What it does +------------ +1. Launches the job in a new session, so the whole tree can be signalled as one + process group. +2. Watches the job's OUTPUT, not its total runtime. A batch still printing is + alive; silence is the signal. This avoids inventing a wall-clock budget for a + test matrix nobody has timed, and it does not punish a legitimately slow run. +3. On silence: signals every rank to dump its Python stack (a signal handler + runs on the blocked thread itself and never runs while things are healthy), + samples twice to separate "stuck" from "slow", then compares the ranks and + names the one that is somewhere the others are not. +4. Kills the process group and verifies nothing survived. Killing the launcher + is not killing the job: an unsupervised kill leaves orphaned ``mpirun``, + ``hydra_pmi_proxy`` and pytest children spinning at 100% CPU. + +Usage +----- +:: + + scripts/mpi_supervisor.py --silence 300 -- mpirun -n 2 python -m pytest ... + +Exit status is the job's own, or 124 if the supervisor had to kill it. +""" + +import argparse +import os +import re +import signal +import subprocess +import sys +import tempfile +import threading +import time +from collections import Counter +from pathlib import Path + +# The job outlived its silence budget and was killed. Matches coreutils +# `timeout`, which is what a reader will assume this number means. +EXIT_KILLED = 124 + +# Between the two stack samples. Two identical stacks say "stuck"; one says +# only "slow", and that distinction is the whole point of sampling twice. +SECONDS_BETWEEN_SAMPLES = 20.0 + +# After the dump request, before escalating. Enough for a signal handler to +# write a few frames to a file, not enough to matter to a CI budget. +SECONDS_TO_DUMP = 5.0 + +# Between SIGTERM and SIGKILL. +SECONDS_TO_TERMINATE = 5.0 + + +def _dump_handler_module(dump_dir): + """Source for the module that arms each rank to dump on request. + + Installed as ``usercustomize``, which the ``site`` machinery imports at + interpreter start-up, so every rank is armed without the job's own code + knowing anything about it. + + This registers a SIGNAL handler, which is the important difference from + ``faulthandler.dump_traceback_later``: it runs on the thread that receives + the signal, only when we ask for it, and never walks live frames from a + concurrent thread. That is why arming it costs a healthy run nothing and + cannot reproduce underworld3#661. + """ + return f'''"""Arm this rank to dump its stack on SIGUSR1. Written by mpi_supervisor.""" +import faulthandler +import os +import signal + +# Rank as the launchers advertise it, so a dump can be attributed without +# asking mpi4py (importing it here would run before the job's own import). +_rank = (os.environ.get("OMPI_COMM_WORLD_RANK") + or os.environ.get("PMI_RANK") + or os.environ.get("MV2_COMM_WORLD_RANK") + or "unknown") + +# Held open for the life of the process on purpose: faulthandler keeps the +# descriptor, not this object, so letting it close would point the handler at +# whatever reuses the fd. +_path = os.path.join({str(dump_dir)!r}, "rank_%s_pid_%d.stack" % (_rank, os.getpid())) +_stream = open(_path, "w") +faulthandler.register(signal.SIGUSR1, file=_stream, all_threads=True, chain=False) +''' + + +def _install_dump_handler(dump_dir, env): + """Put the handler module on the job's import path. + + ``usercustomize`` rather than ``sitecustomize`` because it is imported + later and is far less likely to shadow one the environment already + provides. + """ + handler_dir = Path(tempfile.mkdtemp(prefix="uw-supervisor-")) + (handler_dir / "usercustomize.py").write_text(_dump_handler_module(dump_dir)) + + existing = env.get("PYTHONPATH") + env["PYTHONPATH"] = f"{handler_dir}{os.pathsep}{existing}" if existing else str(handler_dir) + return handler_dir + + +def _descendants(pid): + """Every live descendant of *pid*, deepest last. + + ``ps`` rather than psutil: this must run on a bare CI image, and the + supervisor has to work when the thing it supervises is the broken part. + """ + listing = subprocess.run( + ["ps", "-Ao", "pid=,ppid="], capture_output=True, text=True + ).stdout + + children = {} + for line in listing.splitlines(): + fields = line.split() + if len(fields) == 2: + child, parent = int(fields[0]), int(fields[1]) + children.setdefault(parent, []).append(child) + + found, frontier = [], [pid] + while frontier: + current = frontier.pop() + for child in children.get(current, []): + found.append(child) + frontier.append(child) + return found + + +def _cpu_percent(pids): + """Per-PID CPU, which separates a spinning collective from a still one. + + A rank inside a busy-waiting MPI progress engine burns a whole core; one + blocked on a lock sleeps. That is a cheap and surprisingly sharp signal + about which kind of stuck we are looking at. + """ + if not pids: + return {} + listing = subprocess.run( + ["ps", "-Ao", "pid=,pcpu="], capture_output=True, text=True + ).stdout + + wanted = set(pids) + usage = {} + for line in listing.splitlines(): + fields = line.split() + if len(fields) == 2 and int(fields[0]) in wanted: + usage[int(fields[0])] = float(fields[1]) + return usage + + +def _request_dumps(pids): + """Ask every rank for a stack. Ranks that are gone are simply not asked.""" + for pid in pids: + try: + os.kill(pid, signal.SIGUSR1) + except ProcessLookupError: + # The rank exited between listing and signalling. Nothing to dump. + pass + + +def _dump_blocks(text): + """The individual dumps in one rank's file. + + ``faulthandler`` appends, so a file accumulates every sample we have asked + for. Splitting them apart is what lets one read answer both "where is this + rank" and "has it moved since last time". + """ + blocks = re.split(r"^Current thread ", text, flags=re.MULTILINE) + return [block for block in blocks if block.strip()] + + +def _stack_signature(block): + """The innermost frames of one dump, as the identity of where a rank is. + + Function names only: paths and line numbers differ between ranks running + from different working directories and would make two ranks in the same + place look different. ``faulthandler`` prints most-recent-call-first, so + the innermost frame is the FIRST one. + """ + frames = re.findall(r'File "[^"]*", line \d+ in (\S+)', block) + return tuple(frames[:4]) + + +def _read_dumps(dump_dir): + """Every rank's samples, as ``{rank: (pid, [block, ...])}``.""" + dumps = {} + for path in sorted(Path(dump_dir).glob("rank_*_pid_*.stack")): + rank, _, pid = path.stem.removeprefix("rank_").partition("_pid_") + blocks = _dump_blocks(path.read_text()) + if blocks: + dumps[rank] = (int(pid), blocks) + return dumps + + +def _diagnose(dumps, cpu_by_rank): + """Say which rank missed the collective, and whether it is stuck or slow. + + This is the reading a human would make from a pile of stack dumps, made + here instead: the ranks that agree are waiting for the one that does not. + Naming that rank is the difference between a diagnosis and more log volume. + + *dumps* is ``{rank: (pid, [block, ...])}`` with the samples in the order + they were taken. + """ + lines = [] + + moved = {rank: _stack_signature(blocks[-2]) != _stack_signature(blocks[-1]) + for rank, (_pid, blocks) in dumps.items() if len(blocks) >= 2} + if moved and not any(moved.values()): + lines.append(f" no rank moved in {SECONDS_BETWEEN_SAMPLES:.0f} s: " + "this is stuck, not slow") + elif any(moved.values()): + still = sorted(rank for rank, m in moved.items() if not m) + busy = sorted(rank for rank, m in moved.items() if m) + lines.append(f" ranks {', '.join(busy)} are still moving; " + f"ranks {', '.join(still)} are not — this may be slow, not stuck") + + signatures = {rank: _stack_signature(blocks[-1]) for rank, (_pid, blocks) in dumps.items()} + tally = Counter(signatures.values()) + majority, majority_count = tally.most_common(1)[0] + + def _where(rank): + return signatures[rank][0] if signatures[rank] else "an unknown frame" + + if len(tally) == 1: + lines.append(f" every rank is in {majority[0] if majority else 'the same place'}()" + " — no rank is the odd one out, so suspect a collective they " + "have all entered") + elif any(count < majority_count for count in tally.values()): + odd = sorted(rank for rank, sig in signatures.items() if sig != majority) + lines.append(f" rank(s) {', '.join(odd)} are somewhere the other " + f"{majority_count} are not — that is where to look first") + for rank in odd: + lines.append(f" rank {rank} is in {_where(rank)}()") + lines.append(f" the other {majority_count} are in " + f"{majority[0] if majority else 'an unknown frame'}()") + else: + # No majority to appeal to — the usual case at np=2. Still worth + # reporting where each rank is; a reader can see the mismatch even + # when the supervisor cannot vote on it. + lines.append(" the ranks are in different places, and with no majority " + "none can be called the odd one out:") + for rank in sorted(signatures): + lines.append(f" rank {rank} is in {_where(rank)}()") + + for rank in sorted(cpu_by_rank): + percent = cpu_by_rank[rank] + state = ("spinning — a busy-wait, which is what an MPI progress engine does" + if percent > 50 else "idle — blocked rather than polling") + lines.append(f" rank {rank}: {percent:.0f}% CPU, {state}") + + return lines + + +def _kill_the_job(pid): + """End the whole tree and confirm it is gone. Returns any survivors. + + Descendants are killed individually before the group, because the group + call is the one that fails: ``mpirun`` puts its children in a group of + their own, and ``killpg`` on macOS has been seen to refuse with EPERM. + Enumerating the tree does not depend on either. + + The confirmation is not ceremony. Killing only the launcher leaves the + children orphaned to init and spinning at 100% CPU, which is how a + "finished" run keeps eating cores. + """ + for sig in (signal.SIGTERM, signal.SIGKILL): + for victim in reversed(_descendants(pid)): + try: + os.kill(victim, sig) + except (ProcessLookupError, PermissionError): + # Already gone, or not ours to kill: either way the group + # sweep below is the remaining chance, and survivors are + # reported rather than hidden. + pass + try: + os.kill(pid, sig) + except ProcessLookupError: + pass + try: + os.killpg(os.getpgid(pid), sig) + except (ProcessLookupError, PermissionError): + pass + + deadline = time.monotonic() + SECONDS_TO_TERMINATE + while time.monotonic() < deadline and _descendants(pid): + time.sleep(0.5) + if not _descendants(pid): + break + + return _descendants(pid) + + +def main(): + parser = argparse.ArgumentParser( + description="Run an MPI job, bound any collective hang, and diagnose it.") + parser.add_argument( + "--silence", type=float, default=300.0, metavar="SECONDS", + help="kill the job after this long with no output (default: 300)") + parser.add_argument( + "--hard-cap", type=float, default=0.0, metavar="SECONDS", + help="also kill it after this much total runtime (default: no cap)") + parser.add_argument( + "--artifacts", type=Path, default=None, metavar="DIR", + help="where to keep the stack dumps (default: a temporary directory)") + parser.add_argument("command", nargs=argparse.REMAINDER, + help="the job, after a bare --") + args = parser.parse_args() + + command = args.command[1:] if args.command[:1] == ["--"] else args.command + if not command: + parser.error("no command given; put it after a bare --") + + dump_dir = args.artifacts or Path(tempfile.mkdtemp(prefix="uw-stacks-")) + dump_dir.mkdir(parents=True, exist_ok=True) + + env = os.environ.copy() + _install_dump_handler(dump_dir, env) + + job = subprocess.Popen( + command, env=env, start_new_session=True, + stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, bufsize=1) + + # The job's output is both the thing we forward and the heartbeat we judge + # it by, so it is read on a thread rather than polled. + last_output = [time.monotonic()] + + def forward_output(): + for line in job.stdout: + last_output[0] = time.monotonic() + sys.stdout.write(line) + sys.stdout.flush() + + reader = threading.Thread(target=forward_output, daemon=True) + reader.start() + + started = time.monotonic() + while job.poll() is None: + time.sleep(1.0) + silent_for = time.monotonic() - last_output[0] + overran = args.hard_cap and (time.monotonic() - started) > args.hard_cap + if silent_for <= args.silence and not overran: + continue + + reason = (f"no output for {silent_for:.0f} s" + if not overran else f"exceeded the {args.hard_cap:.0f} s cap") + print(f"\n=== MPI SUPERVISOR: {reason} ===", flush=True) + + ranks = _descendants(job.pid) + _request_dumps(ranks) + time.sleep(SECONDS_TO_DUMP) + + time.sleep(SECONDS_BETWEEN_SAMPLES) + _request_dumps(ranks) + time.sleep(SECONDS_TO_DUMP) + + # One read: faulthandler appends, so each rank's file already holds + # both samples in order. + dumps = _read_dumps(dump_dir) + + # Attribute CPU through the pid the handler recorded. Guessing which + # rank owned which figure would produce a confident, wrong label. + usage = _cpu_percent([pid for pid, _blocks in dumps.values()]) + cpu_by_rank = {rank: usage[pid] for rank, (pid, _blocks) in dumps.items() + if pid in usage} + + if dumps: + for line in _diagnose(dumps, cpu_by_rank): + print(line, flush=True) + print(f" stacks written to {dump_dir}", flush=True) + else: + print(" no rank answered the dump request — no stacks available; " + "killing the job regardless", flush=True) + + survivors = _kill_the_job(job.pid) + if survivors: + print(f" WARNING: {len(survivors)} process(es) survived the kill: " + f"{survivors}", flush=True) + return EXIT_KILLED + + reader.join(timeout=5.0) + return job.returncode + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/scripts/test.sh b/scripts/test.sh index 30c379132..22ff54768 100755 --- a/scripts/test.sh +++ b/scripts/test.sh @@ -156,8 +156,16 @@ if [ $PARALLEL_RANKS -gt 0 ]; then # - Solver operations # - Global evaluations + # Every mpirun goes through the supervisor. A rank blocked in a collective + # produces nothing and never returns, so an unsupervised batch spends the + # whole job budget in silence and is cancelled from outside with no + # diagnosis -- measured at 76 minutes (#675). The supervisor bounds that on + # silence rather than total runtime, so a legitimately slow batch is not + # punished, and it dumps and compares the ranks before killing them. + SUPERVISE="python $(dirname "$0")/mpi_supervisor.py --silence ${PARALLEL_SILENCE:-300} --" + echo "Testing global statistics and parallel operations..." - mpirun -n $PARALLEL_RANKS python -m pytest --with-mpi tests/parallel/test_075*py || status=1 + $SUPERVISE mpirun -n $PARALLEL_RANKS python -m pytest --with-mpi tests/parallel/test_075*py || status=1 # Parallel SOLVER tests. This line was commented out, so test_1017 and # test_1062..test_1069 — the whole rotated / constrained / MG parallel set, @@ -165,7 +173,7 @@ if [ $PARALLEL_RANKS -gt 0 ]; then # (#560) and the mesh boundary normal (#564) — executed at NO rank count # in CI. echo "Testing parallel solvers..." - mpirun -n $PARALLEL_RANKS python -m pytest --with-mpi tests/parallel/test_10*py || status=1 + $SUPERVISE mpirun -n $PARALLEL_RANKS python -m pytest --with-mpi tests/parallel/test_10*py || status=1 # echo "Testing parallel I/O..." # mpirun -n $PARALLEL_RANKS python -m pytest --with-mpi tests/parallel/test_io*py || status=1 diff --git a/tests/parallel/hang_controls/rank_zero_misses_the_collective.py b/tests/parallel/hang_controls/rank_zero_misses_the_collective.py new file mode 100644 index 000000000..7441537f4 --- /dev/null +++ b/tests/parallel/hang_controls/rank_zero_misses_the_collective.py @@ -0,0 +1,50 @@ +"""A planted collective hang, for validating the MPI supervisor. + +Every rank but 0 enters a barrier. Rank 0 does not, and instead sits in a +function with a name no other rank can be inside. The job therefore hangs +forever, and the correct diagnosis is unambiguous: rank 0 is the one that +missed the collective. + +This is the negative control for ``scripts/mpi_supervisor.py``. A supervisor +that cannot bound this run, and cannot name rank 0 from the stacks it collects, +has no business gating a test suite. + +It lives outside the ``test_*.py`` pattern deliberately: pytest must never +collect it, because running it under any harness that does not kill it is the +exact failure this whole exercise exists to prevent. It is launched by name, +only by ``tests/test_0056_mpi_hang_supervisor.py``. +""" + +import sys +import time + +from mpi4py import MPI + + +def _rank_zero_misses_the_collective(): + """Rank 0 spends forever here while the others wait in the barrier. + + Named so that it is visible in a stack dump and cannot be confused with + the MPI frames the other ranks are in. + """ + while True: + time.sleep(3600) + + +def main(): + comm = MPI.COMM_WORLD + # Unbuffered, and flushed, so the supervisor sees output stop at a known + # point: the silence after this line is the hang it must detect. + print(f"rank {comm.rank} of {comm.size} ready", flush=True) + + if comm.rank == 0: + _rank_zero_misses_the_collective() + else: + comm.barrier() + + print("this line is unreachable while rank 0 is missing", flush=True) + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/tests/test_0063_mpi_hang_supervisor.py b/tests/test_0063_mpi_hang_supervisor.py new file mode 100644 index 000000000..1ce163d79 --- /dev/null +++ b/tests/test_0063_mpi_hang_supervisor.py @@ -0,0 +1,113 @@ +"""The MPI supervisor must bound a planted collective hang and diagnose it. + +A harness that catches hangs is worth nothing until it has caught one. These +tests plant a hang whose answer is known — rank 0 sits in a function no other +rank can be inside, while the rest wait in a barrier — and require the +supervisor to end the job and name rank 0 from the stacks it collects. + +The second requirement matters as much as the first: a supervisor that kills +healthy jobs, or that reports a diagnosis it cannot support, is worse than no +supervisor. So a clean run must pass through with its own exit status, and a +two-rank hang, where there is no majority to appeal to, must say so rather than +pick a rank arbitrarily. +""" + +import shutil +import subprocess +import sys +import time +from pathlib import Path + +import pytest + +pytestmark = [pytest.mark.level_2, pytest.mark.tier_b] + +REPO = Path(__file__).resolve().parent.parent +SUPERVISOR = REPO / "scripts" / "mpi_supervisor.py" +PLANTED_HANG = REPO / "tests" / "parallel" / "hang_controls" / "rank_zero_misses_the_collective.py" + +# Short enough to keep the test quick, long enough that a slow import cannot +# be mistaken for a hang. The supervisor adds its two sampling pauses on top. +SILENCE = 15.0 + +# The supervisor's own budget: silence, two dumps, the gap between them, and +# the kill ladder. Generous, because this asserts "it ends", not "it is fast". +BUDGET = 180.0 + +needs_mpirun = pytest.mark.skipif( + shutil.which("mpirun") is None, reason="no mpirun on this machine") + + +def _run_supervised(job, artifacts, ranks=2): + started = time.monotonic() + completed = subprocess.run( + [sys.executable, str(SUPERVISOR), "--silence", str(SILENCE), + "--artifacts", str(artifacts), + "--", "mpirun", "-n", str(ranks), sys.executable, str(job)], + capture_output=True, text=True, timeout=BUDGET) + return completed, time.monotonic() - started + + +@needs_mpirun +def test_a_planted_hang_is_bounded_and_the_guilty_rank_named(tmp_path): + """Four ranks, one of which skips the barrier. Rank 0 must be named.""" + completed, elapsed = _run_supervised(PLANTED_HANG, tmp_path, ranks=4) + + assert completed.returncode == 124, ( + f"the supervisor did not report a kill (got {completed.returncode}):\n" + f"{completed.stdout}\n{completed.stderr}") + assert elapsed < BUDGET, "the supervisor did not end the job within its budget" + + report = completed.stdout + assert "stuck, not slow" in report, ( + f"two identical samples should have read as stuck:\n{report}") + assert "rank(s) 0 are somewhere the other 3 are not" in report, ( + f"the rank that missed the collective was not named:\n{report}") + assert "_rank_zero_misses_the_collective" in report, ( + f"the diagnosis did not reach the frame that proves it:\n{report}") + + stacks = sorted(tmp_path.glob("rank_*_pid_*.stack")) + assert len(stacks) == 4, f"expected a dump per rank, got {[p.name for p in stacks]}" + + +@needs_mpirun +def test_two_ranks_have_no_majority_and_the_report_says_so(tmp_path): + """At np=2 neither rank outvotes the other. + + The supervisor must report both positions and decline to nominate one. + Claiming a culprit from a one-all split would be a confident guess, which + is the failure mode this whole harness exists to avoid. + """ + completed, _elapsed = _run_supervised(PLANTED_HANG, tmp_path, ranks=2) + + assert completed.returncode == 124 + report = completed.stdout + assert "no majority" in report, f"the tie was not reported as a tie:\n{report}" + assert "rank 0 is in _rank_zero_misses_the_collective()" in report, report + assert "rank 1 is in main()" in report, report + + +@needs_mpirun +def test_a_healthy_job_is_left_alone(tmp_path): + """The negative control: no false positives, and the exit status passes through. + + Without this, every assertion above could be satisfied by a supervisor that + simply kills everything it is given. + """ + healthy = tmp_path / "healthy_job.py" + healthy.write_text( + "from mpi4py import MPI\n" + "comm = MPI.COMM_WORLD\n" + "print(f'rank {comm.rank} reporting', flush=True)\n" + "comm.barrier()\n" + "print('all ranks through the barrier', flush=True)\n") + + completed, elapsed = _run_supervised(healthy, tmp_path / "stacks", ranks=2) + + assert completed.returncode == 0, ( + f"a healthy job was not left alone:\n{completed.stdout}\n{completed.stderr}") + assert "MPI SUPERVISOR" not in completed.stdout, ( + f"the supervisor intervened in a healthy run:\n{completed.stdout}") + assert elapsed < SILENCE + 30.0, "a healthy job should finish well inside the silence budget" + assert "all ranks through the barrier" in completed.stdout, ( + "the job's own output was not forwarded")