diff --git a/README.md b/README.md index fa97463..82c2800 100644 --- a/README.md +++ b/README.md @@ -105,7 +105,9 @@ out = src( ) ``` -See [Large problems](https://algorithmiq.github.io/src-method/features/large-problems) for the budgets, the scratch directory and how to read the logged plan. +On a node with several GPUs, `Resources(devices=4)` splits the sweep among four of them, in the same process. + +See [Large problems](https://algorithmiq.github.io/src-method/features/large-problems) for the budgets, the scratch directory, several GPUs and how to read the logged plan. ### Logging diff --git a/benches/large/README.md b/benches/large/README.md index 24e15f4..916d8a3 100644 --- a/benches/large/README.md +++ b/benches/large/README.md @@ -1,7 +1,8 @@ # Out-of-core benchmark `bench_large.py` compresses the stack `N . V . M . U` of Pauli-transfer-matrix MPOs -(physical legs of 4) on one GPU, where `M` has a large bond and is read from disk. +(physical legs of 4) on one GPU or several, where `M` has a large bond and is read +from disk. ```bash # Write the stack to node-local NVMe (205 GB for M at the defaults, complex128). @@ -15,9 +16,13 @@ uv run python benches/large/bench_large.py run /local/stack --chi-out 2000 \ uv run python benches/large/bench_large.py generate /local/medium --bond-m 1000 uv run python benches/large/bench_large.py compare /local/medium --chi-out 500 \ --small 40GB --large 80GB --scratch-dir /local/scratch + +# Split the sweep among the GPUs of the node: 1, 2, 4 and 8 of them, in turn. +uv run python benches/large/bench_large.py scaling /local/stack --chi-out 2000 \ + --devices 1 2 4 8 --scratch-dir /local/scratch --results scaling.jsonl ``` -`run` logs the wall time, the size of CuPy's pool at the end (its high-water mark, +`run` logs the wall time, the size of CuPy's pool on each GPU at the end (its high-water mark, since the pool keeps its blocks), the host peak (`ru_maxrss`), the planned peaks, the bytes spilled and the tiers. With `--debug` before the command, `src_method` also logs its plan, the time of each pass and the stall time (`SRC stalls`). @@ -31,3 +36,12 @@ memory and node-local NVMe: budget. 3. The stall time is below 10% of the wall time. 4. `compare` at `D_M = 1000`, `chi_out = 500` reports a distance below `1e-10`. + +`scaling` logs, for each GPU count, the wall time, the speed-up and the parallel +efficiency `T_1 / (G T_G)` against the first count, and the relative distance to +its output; `--results` appends the same as JSON lines. `--devices` with +`run` splits a single run. Acceptance of the multi-GPU sweep, at the reference +size on 8xH100 or 8xH200 (and 4xA100 on Leonardo): + +1. Every distance is below `1e-10`. +2. The parallel efficiency is at least 80% on all the GPUs of the node. diff --git a/benches/large/bench_large.py b/benches/large/bench_large.py index 8ad4143..59b8a2e 100644 --- a/benches/large/bench_large.py +++ b/benches/large/bench_large.py @@ -1,13 +1,16 @@ """Out-of-core SRC of ``N . V . M . U`` with a large ``M``, read from disk. -Three commands: +Four commands: - ``generate``: write the four MPOs site by site as ``.npy`` files, so that ``M`` is never whole in memory. -- ``run``: compress the stack with automatic (or given) budgets and report the - plan, the wall time, the device pool size, the host peak and the stall time. +- ``run``: compress the stack with automatic (or given) budgets, on one or several + GPUs, and report the plan, the wall time, the device pool sizes, the host peak + and the stall time. - ``compare``: run twice with two GPU budgets and report the relative distance between the two outputs, which checks that the plan does not change the result. +- ``scaling``: run on 1, 2, 4, ... GPUs of the node and report the wall time, the + parallel efficiency and the distance to the single-GPU output. The plan, the per-pass times and the stall time are also logged by `src_method` itself at ``DEBUG``; pass ``--debug`` to see them. @@ -15,6 +18,7 @@ from __future__ import annotations +import json import logging import resource from pathlib import Path # noqa: TC003 (cyclopts reads the annotations at runtime) @@ -94,8 +98,12 @@ def _load(directory: Path) -> list[list[np.ndarray]]: def _run( - directory: Path, chi_out: int, resources: Resources, seed: int -) -> list[np.ndarray]: + directory: Path, + chi_out: int, + resources: Resources, + seed: int, + device: str = "gpu", +) -> tuple[list[np.ndarray], float]: layers = _load(directory) plans = [] @@ -111,21 +119,27 @@ def spy(*args: object, **kwargs: object) -> object: chi_out=chi_out, dtype=layers[0][0].dtype, seed=seed, - device="gpu", + device=device, resources=resources, ) seconds = perf_counter() - start finally: sweep_module.make_plan = make_plan - import cupy # noqa: PLC0415 (GPU-only benchmark) - (plan,) = plans tiers = [site.tier for site in plan.sites] + pools = [] + if device == "gpu": + import cupy # noqa: PLC0415 (optional dependency) + + devices = resources.devices or 1 + for gpu in range(devices) if isinstance(devices, int) else devices: + with cupy.cuda.Device(gpu): + pools.append(cupy.get_default_memory_pool().total_bytes()) logger.info( - "Run complete: %.1f s, pool %d B, host peak %d B, planned device peak %d B, " + "Run complete: %.1f s, pools %s B, host peak %d B, planned device peak %d B, " "planned host peak %d B, disk %d B, tiers %s, sketch batches %s", seconds, - cupy.get_default_memory_pool().total_bytes(), + pools, resource.getrusage(resource.RUSAGE_SELF).ru_maxrss * 1024, plan.device_peak, plan.host_peak, @@ -133,7 +147,7 @@ def spy(*args: object, **kwargs: object) -> object: {tier: tiers.count(tier) for tier in ("device", "host", "disk")}, sorted({site.sketch_batch for site in plan.sites[1:]}), ) - return out + return out, seconds @app.command @@ -144,6 +158,7 @@ def run( gpu_memory: str | None = None, host_memory: str | None = None, scratch_dir: Path | None = None, + devices: int = 1, seed: int = 0, ) -> None: """Compress the stack written by ``generate`` on the GPU. @@ -151,12 +166,14 @@ def run( Args: directory: The directory given to ``generate``. chi_out: The output bond dimension. - gpu_memory: GPU budget, e.g. ``36GB``; detected when omitted. + gpu_memory: GPU budget of each device, e.g. ``36GB``; detected when omitted. host_memory: Host budget; detected when omitted. scratch_dir: Where to spill environments; ``$TMPDIR`` when omitted. + devices: The number of GPUs to split the sweep among. seed: Seed of the sketch. """ - _run(directory, chi_out, Resources(gpu_memory, host_memory, scratch_dir), seed) + resources = Resources(gpu_memory, host_memory, scratch_dir, devices=devices) + _run(directory, chi_out, resources, seed) def _inner(a: list[np.ndarray], b: list[np.ndarray]) -> complex: @@ -187,11 +204,82 @@ def compare( scratch_dir: Where to spill environments; ``$TMPDIR`` when omitted. seed: Seed of the sketch, the same for both runs. """ - first = _run(directory, chi_out, Resources(small, None, scratch_dir), seed) - second = _run(directory, chi_out, Resources(large, None, scratch_dir), seed) - aa, bb, ab = _inner(first, first), _inner(second, second), _inner(first, second) - distance = np.sqrt(max((aa + bb - 2 * ab).real, 0.0) / aa.real) - logger.info("Relative distance between the runs: %.3e", distance) + first, _ = _run(directory, chi_out, Resources(small, None, scratch_dir), seed) + second, _ = _run(directory, chi_out, Resources(large, None, scratch_dir), seed) + logger.info("Relative distance between the runs: %.3e", _distance(first, second)) + + +def _distance(a: list[np.ndarray], b: list[np.ndarray]) -> float: + """Relative Frobenius distance ``|a - b| / |a|`` of two MPOs.""" + aa, bb, ab = _inner(a, a), _inner(b, b), _inner(a, b) + return float(np.sqrt(max((aa + bb - 2 * ab).real, 0.0) / aa.real)) + + +@app.command +def scaling( + directory: Path, + *, + chi_out: int = 2000, + devices: Annotated[tuple[int, ...], cyclopts.Parameter(consume_multiple=True)] = ( + 1, + 2, + 4, + 8, + ), + gpu_memory: str | None = None, + host_memory: str | None = None, + scratch_dir: Path | None = None, + results: Path | None = None, + seed: int = 0, + device: str = "gpu", +) -> None: + """Run on increasing numbers of GPUs and report the strong scaling. + + The first count is the reference: the efficiency of ``G`` GPUs is + ``T_ref * G_ref / (T_G * G)``, and every output is compared with the + reference output. Host memory holds two outputs at a time. + + Args: + directory: The directory given to ``generate``. + chi_out: The output bond dimension. + devices: The GPU counts, the reference first. + gpu_memory: GPU budget of each device; detected when omitted. + host_memory: Host budget of the node; detected when omitted. + scratch_dir: Where to spill environments; ``$TMPDIR`` when omitted. + results: A JSON-lines file to append one record per run to. + seed: Seed of the sketch, the same for every run. + device: ``gpu``, or ``cpu`` to smoke-test with simulated devices. + """ + reference: list[np.ndarray] | None = None + ref_seconds = ref_devices = 0.0 + for count in devices: + resources = Resources(gpu_memory, host_memory, scratch_dir, devices=count) + out, seconds = _run(directory, chi_out, resources, seed, device) + if reference is None: + reference, ref_seconds, ref_devices = out, seconds, count + distance = 0.0 + else: + distance = _distance(reference, out) + efficiency = ref_seconds * ref_devices / (seconds * count) + logger.info( + "GPUs %d: %.1f s, speed-up %.2f, efficiency %.0f %%, distance %.3e", + count, + seconds, + ref_seconds / seconds, + 100 * efficiency, + distance, + ) + if results is not None: + record = { + "gpus": count, + "seconds": seconds, + "efficiency": efficiency, + "distance": distance, + "chi_out": chi_out, + } + with results.open("a") as f: + f.write(json.dumps(record) + "\n") + del out @app.meta.default diff --git a/docs/content/docs/contributing/testing.mdx b/docs/content/docs/contributing/testing.mdx index 02e5b48..ce1aa0c 100644 --- a/docs/content/docs/contributing/testing.mdx +++ b/docs/content/docs/contributing/testing.mdx @@ -23,6 +23,7 @@ Always run tests through `pytest`, never as `python tests/test_.py`. | `test_plan.py` | the planner: budgets, batch sizes and environment tiers | | `test_store.py` | environment storage on the device, in host memory and on disk | | `test_sites.py` | reading the cores one site at a time, with prefetching | +| `test_group.py` | the device group: work splitting, collectives, failures | | `test_logging.py` | the package leaves the host's logging configuration alone | | `test_gpu_backend.py` | the CuPy path; skipped without CuPy or a GPU | @@ -88,6 +89,16 @@ pure function, so `tests/test_plan.py` checks batch sizes and tiers from shapes alone, and `tests/test_gpu_backend.py` derives GPU budgets from a plan made with `make_plan` to reach every tier. +## Exercising several devices + +`Resources(devices=n)` on the CPU simulates `n` devices with threads that share +the host, through the same code as on GPUs: the split of the sketch columns, the +collectives and the shared reading of the cores. `tests/test_stack.py` compares +such runs with a single device, also with uneven blocks (`chi_out` not a +multiple of `n`, or smaller than it), a cutoff and spilling. Only the peer copies +and device selection are GPU-specific; their tests in `tests/test_gpu_backend.py` +skip unless at least two GPUs are visible. + ## Logging in tests `log_cli` is on, so log records show up live while the tests run. To assert on diff --git a/docs/content/docs/features/gpu.mdx b/docs/content/docs/features/gpu.mdx index 47c7526..dde982a 100644 --- a/docs/content/docs/features/gpu.mdx +++ b/docs/content/docs/features/gpu.mdx @@ -27,3 +27,7 @@ result = apply(H, psi, chi_out=256, device="gpu") The inputs can be NumPy or CuPy arrays; they are moved to the device for the sweep, and the result always comes back as NumPy arrays on the host. Without CuPy installed, `device="gpu"` raises an `ImportError`. + +To split one sweep among several GPUs of a node, pass +`resources=Resources(devices=n)`; see +[Large problems](/features/large-problems#several-gpus). diff --git a/docs/content/docs/features/large-problems.mdx b/docs/content/docs/features/large-problems.mdx index 3cc4bd7..165bb01 100644 --- a/docs/content/docs/features/large-problems.mdx +++ b/docs/content/docs/features/large-problems.mdx @@ -1,6 +1,6 @@ --- title: Large problems -description: Run stacks that fit neither on the GPU nor in host memory. +description: Run stacks that fit neither on the GPU nor in host memory, on one GPU or several. --- A single SRC sweep keeps, besides the input and output trains, one sketched @@ -75,13 +75,68 @@ reference size of 50 sites, `D_M = 4000` and `chi_out = 2000` in complex128, up returns, also after an error; a process killed with `SIGKILL` leaves it behind, and its name identifies the process. +## Several GPUs + +On a node with several GPUs, one sweep can be split among them: + +```python notest +from src_method import Resources, src + +out = src( + N, V, M, U, + chi_out=2000, + dtype=np.complex128, + device="gpu", + resources=Resources(devices=4, scratch_dir="/local/scratch"), +) +``` + +`devices` is a count, for the first `n` visible GPUs, or a list of device ids; +the first one gathers the sketches and runs the QR. Everything runs in the calling +process, one thread per GPU, so nothing changes for the caller: the call returns +the output cores as with one GPU, errors are raised as usual, and the inputs are +read once and shared. + +The sketch index is split: each GPU owns a contiguous block of the `chi_out` +sketch columns and computes only those columns of every environment and of every +sketch, so the environments, the dominant memory, are split among the GPUs too. +At every site of the right-to-left pass, the first GPU gathers the sketch, runs +the QR and gives each GPU a block of rows of the output core; each GPU projects +its rows of the next projected environment, and those rows are gathered on every +GPU. The random draws are those of one GPU, so the result is that of one GPU up to +rounding. + +What is not split: + +- every GPU holds a whole site: the cores, the projected environment and the + output core of the site, so the largest site that fits is that of one GPU; +- every core is read from the source once, but every GPU needs a copy. When the + GPUs reach each other directly (NVLink, or PCIe peer-to-peer), each uploads a + slice of the core and gathers the rest from the others; otherwise each uploads + the whole core; +- the QR runs on the first GPU while the others wait. + +These costs grow more slowly with the problem than the contractions do, so large +problems, the ones this is for, scale close to linearly, while small ones may run +slower on several GPUs than on one. + +The budgets apply as follows: `gpu_memory` is the budget of each GPU (detected: +the smallest free memory of the GPUs), while `host_memory` and `scratch_dir` are +shared by all of them. Every GPU writes its environments to its own scratch +directory. + +On the CPU, `devices` simulates the GPUs with threads that share the host; the +result is the same, and the mode exists to test the multi-GPU sweep without GPUs. + ## Reading the plan With `DEBUG` logging on for the `src_method` logger (see [Logging](/features/logging)), every sweep logs its plan as `SRC plan: ...`: - `prefetch`: whether the next site is read ahead; -- `device peak`, `host peak`, `disk`: the planned peaks, in bytes; +- `devices`: the device ids the sweep is split among; +- `device peak`, `host peak`, `disk`: the planned peaks, in bytes; the device + peak is that of each device, the others are totals; - `scratch`: the scratch directory, when anything spills to disk; - `tiers`: where the environment of each site lives (`device`, `host` or `disk`); - `batches`: for each site, the batch of the environment, sketch and projection @@ -89,4 +144,5 @@ With `DEBUG` logging on for the `src_method` logger (see It then logs the time of each pass and, at the end, `SRC stalls`: the seconds the sweep waited for input cores and for environments. If they are a large -fraction of the run, the disk is too slow for the compute. +fraction of the run, the disk is too slow for the compute. With several devices, +these lines start with the device they refer to. diff --git a/docs/superpowers/plans/2026-09-25-src-multi-gpu.md b/docs/superpowers/plans/2026-09-25-src-multi-gpu.md new file mode 100644 index 0000000..c97ab9b --- /dev/null +++ b/docs/superpowers/plans/2026-09-25-src-multi-gpu.md @@ -0,0 +1,503 @@ +# SRC Sweep over the GPUs of a Node (Phase 2): Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Run one `src` sweep on the GPUs of a node, one process per GPU, with the sketch index split among them, so that the Phase 1 reference problem finishes close to `G` times faster on `G` GPUs. + +**Architecture:** A small `Communicator` interface in `_comm.py` (single process, `mpi4py` on host arrays, NCCL on CuPy arrays) carries the few collectives the sweep needs. `src` builds the communicator from `Resources(comm=...)`, agrees on the call and the seed, and hands the output to a sink (rank 0 in memory, or `.npy` files and memmaps on every rank). The Phase 1 driver keeps one copy of the maths: each rank owns a cyclic subset of the sketch columns, and the right-to-left pass all-gathers the sketch columns and the rows of the projected environment at every site. The planner plans each rank's share. + +**Tech Stack:** Python 3.11-3.14, NumPy and CuPy (through `src_method.utils._backend`), `opt_einsum`, standard-library `logging`, `mpi4py>=4` (new `mpi` extra), NCCL through `cupy.cuda.nccl` (`nvidia-nccl-cu12` in the `gpu-nvidia` extra), `pytest`, `cyclopts` (benchmark), `uv`, `ruff` through `prek`, GitButler (`but`). + +**Spec:** `docs/superpowers/specs/2026-09-25-src-multi-gpu-design.md` (read it first; this plan argues from it). Phase 1 is the out-of-core sweep merged in #42 (`_plan.py`, `_kernels.py`, `_sites.py`, `_store.py`, `_sweep.py`). + +## How to read this plan + +This plan was written under a rule of the design phase: no code in full. Every task +gives the files, the exact interfaces (signatures, types, error messages), the +behaviour in prose or pseudocode, and the tests to write, each with its setup and +its assertions. The executor writes the code, test first, and runs the commands +given. Where a detail is left to the executor, the plan says so. + +## Global Constraints + +- Work on branch `feat/src-multi-gpu`, on top of `main` (Phase 1 is merged). Commit only the files of the task, never the untracked `.codegraph/`. +- `from __future__ import annotations` at the top of every module under `src/`; type hints everywhere; Google-style docstrings without types on every public function, class and module; `ruff` is the source of truth for style. +- Backend-specific calls (CuPy, NCCL, device selection) live in `src_method.utils._backend` or in `_comm.py`; the algorithms get `xp` and a `Communicator`. `mpi4py` and `cupy.cuda.nccl` are imported lazily, only when a communicator is given: a single-process run imports neither. +- Log with the standard library (`logger = logging.getLogger(__name__)`, lazy `%`-style arguments, no handlers in the package), never `print`, and never `warnings.warn` for run-time advice (`filterwarnings = ["error"]` turns warnings into test failures). +- `src`, `apply`, `compress` stay pure: never mutate inputs. +- Commits: `(): `, ending with `Assisted-by: Pi:claude-opus-5-5`, no `Co-authored-by`. +- Before pushing: `uv run prek run --all-files` and `uv run pytest` pass. +- Decisions of the spec, verbatim: one process per GPU (SPMD) launched with `srun` or `mpirun`; NCCL through `cupy.cuda.nccl` with `mpi4py` for bootstrap; `Resources(comm=...)`, nothing collective without it; output on rank 0 in memory, or `.npy` files and memmaps on every rank with `Resources(output_dir=...)`; large inputs file-backed; success is strong scaling (at least 80% parallel efficiency at 8 GPUs on 8xA100 and 8xH100 at the reference size, at least 90% on 2xA100). +- Without a communicator, or with one of size one, the sweep is exactly Phase 1: the existing tests must pass unchanged. + +## Refinements of the spec + +Settled while planning; the executor should not undo them. + +1. **`Communicator.abort(code)`** joins the interface: the driver calls it when an exception escapes during the sweep on more than one rank. `SingleComm.abort` is never called. +2. **The truncation rank is agreed through a callback**: `truncated_qr(matrix, cutoff, xp, agree_rank=None)` computes its local rank from the singular values of `R` and passes it through `agree_rank` (rank 0's value, broadcast) before truncating. The spec's "optional `rank` argument" cannot work, because the rank is only known inside the function. +3. **Output sinks and `unpad_site`**: the output goes through a small sink object, and each core is unpadded as it is produced, with a new `unpad_site` that mirrors `pad_site`. +4. **Large in-memory inputs are logged, not warned about**: when a distributed call on more than one rank per node gets an in-memory `np.ndarray` layer (not a memmap) above 1 GiB, `src` logs a warning (`logger.warning`) that every rank holds a copy. +5. **CI uses the system MPICH** on the Ubuntu runner (`apt-get install mpich`) with the `mpi4py` wheel from PyPI; the MPI tests skip wherever `mpiexec` or a working `mpi4py` is missing. + +## File map + +| File | Status | Responsibility | +|---|---|---| +| `src/src_method/_comm.py` | create | `Communicator`, `SingleComm`, `MpiHostComm`, `NcclComm`, `make_communicator`; block and share helpers; `agree`, `collectively` | +| `src/src_method/_output.py` | create | output sinks: memory, directory, none | +| `src/src_method/_tensor_train.py` | modify | `unpad_site` | +| `src/src_method/utils/linalg.py` | modify | `truncated_qr(..., agree_rank=None)` | +| `src/src_method/utils/_backend.py` | modify | `select_device(local_rank, xp)`, `nccl_module(xp)` | +| `src/src_method/_plan.py` | modify | `Resources.comm`, `Resources.output_dir`; per-rank shares in `resolve_budgets`; `make_plan(..., columns=, rows=, gather=, holds_output=)` | +| `src/src_method/_sweep.py` | modify | column ownership, the two all-gathers, rank-0 output, collective error handling | +| `src/src_method/stack.py` | modify | communicator, fingerprint, seed, device, sinks | +| `pyproject.toml`, `uv.lock`, `flake.nix`, `.github/workflows/test.yml` | modify | `mpi` extra, NCCL wheel, MPICH in the dev shell and in CI | +| `tests/conftest.py` | create | `mpirun` fixture | +| `tests/mpi_scripts/*.py` | create | scripts run under `mpiexec` | +| `tests/test_comm.py`, `tests/test_output.py` | create | unit and MPI tests | +| `tests/test_plan.py`, `tests/test_stack.py`, `tests/test_tensor_train.py`, `tests/test_backend.py`, `tests/test_gpu_backend.py`, `tests/test_linalg.py` | modify / create | planner, integration, helpers, GPU | +| `benches/large/bench_large.py`, `benches/large/README.md` | modify | MPI mode, `scaling` and `efficiency` commands | +| `docs/content/docs/features/large-problems.mdx`, `README.md`, `docs/content/docs/contributing/dependencies.mdx`, `docs/content/docs/contributing/testing.mdx` | modify | documentation | + +--- + +### Task 1: Work splitting and block helpers, and the single-process communicator + +**Files:** +- Create: `src/src_method/_comm.py` +- Test: `tests/test_comm.py` (create) + +**Interfaces (produced):** + +```python +class Communicator(Protocol): + rank: int + size: int + local_rank: int + local_size: int + def allgather(self, local: NDArray, axis: int) -> NDArray: ... + def allgather_objects(self, value: object) -> list[object]: ... + def bcast_int(self, value: int, root: int = 0) -> int: ... + def barrier(self) -> None: ... + def abort(self, code: int) -> None: ... + +class SingleComm: # rank 0 of 1, local rank 0 of 1 + ... # allgather returns `local` itself; allgather_objects -> [value]; + # bcast_int -> value; barrier no-op; abort raises SystemExit(code) + +def owned_columns(n: int, rank: int, size: int) -> slice # slice(rank, n, size) +def row_block(n: int, rank: int, size: int) -> slice # contiguous, sizes differ by <= 1, + # lower ranks take the larger blocks +def pad_blocks(block: NDArray, axis: int, rows: int) -> NDArray + # move `axis` to the front, make contiguous, pad with zeros to `rows` along it +def join_blocks(gathered: NDArray, counts: Sequence[int], axis: int) -> NDArray + # gathered: (size, rows_max, *rest), padded blocks in rank order; + # take counts[r] rows of block r, concatenate, move the axis back +def agree(comm: Communicator, fn: Callable[[], T]) -> T +def collectively(comm: Communicator, fn: Callable[[], T]) -> T +``` + +Behaviour: + +- `agree` runs `fn` on every rank, then all-gathers `None` or `(type name, message)` of the exception it raised. If any rank failed, every rank raises an exception of the same type (one of `ValueError`, `MemoryError`, `TypeError`, otherwise `RuntimeError`) with the message of the lowest failing rank, prefixed `"Rank {r}: "`. Without failures it returns `fn()`'s result. With `SingleComm` it is `fn()`. +- `collectively` runs `fn`; if an exception escapes and `comm.size > 1`, it logs it with `logger.exception("SRC failed on rank %d; aborting", comm.rank)` and calls `comm.abort(1)`; with one rank it re-raises. + +Tests to write (`tests/test_comm.py`): + +- [ ] **Step 1: Write the failing tests** + - `owned_columns`: for `n in (1, 7, 8, 2000)` and `size in (1, 2, 3, 8)`, the ranges of all ranks partition `range(n)` exactly once; rank `r` gets `ceil((n - r) / size)` columns. + - `row_block`: same partition property; blocks are contiguous and in rank order; block sizes differ by at most one. + - `pad_blocks` then `join_blocks` round-trips random blocks of uneven sizes along axis 0 and axis 3 of a rank-4 array, float64 and complex128: joining the padded blocks of three ranks equals `np.concatenate(blocks, axis)`. + - `SingleComm`: `allgather(x, axis)` is `x`; `allgather_objects(v) == [v]`; `bcast_int(5) == 5`; `rank, size, local_rank, local_size == 0, 1, 0, 1`. + - `agree` with `SingleComm`: returns the value; re-raises a `MemoryError` from `fn` unchanged in type. + - `collectively` with `SingleComm`: re-raises. +- [ ] **Step 2: Run** `uv run pytest tests/test_comm.py -o log_cli=false`; expected: `ModuleNotFoundError: No module named 'src_method._comm'`. +- [ ] **Step 3: Implement** `_comm.py` with the interfaces above. Keep `mpi4py` and CuPy imports out of this task. +- [ ] **Step 4: Run** the tests; expected: all pass. +- [ ] **Step 5: Lint** `uv run ruff format src/src_method/_comm.py tests/test_comm.py && uv run ruff check src tests`. +- [ ] **Step 6: Commit** `feat(comm): ✨ work splitting, block helpers and the single-process communicator`. + +--- + +### Task 2: MPI on host arrays, the `mpi` extra and the MPI test harness + +**Files:** +- Modify: `src/src_method/_comm.py` (add `MpiHostComm`, `make_communicator`) +- Modify: `pyproject.toml` (extra `mpi = ["mpi4py>=4"]`), `uv.lock` (`uv lock`), `flake.nix` (add `pkgs.mpich` to the dev shell packages and its `lib` to `LD_LIBRARY_PATH`), `.github/workflows/test.yml` (Ubuntu runner: `sudo apt-get install -y mpich` before syncing, and `--extra mpi` in `uv sync`; the GPU runner is unchanged) +- Create: `tests/conftest.py`, `tests/mpi_scripts/collectives.py` +- Test: `tests/test_comm.py` (append) + +**Interfaces (produced):** + +```python +class MpiHostComm: + def __init__(self, comm: Any) -> None: ... # an mpi4py Comm + # rank/size from comm; local_rank/local_size from comm.Split_type(MPI.COMM_TYPE_SHARED) + # allgather: moveaxis + contiguous uint8 view, Allgatherv with byte counts and + # displacements from an allgather of the block lengths, then join_blocks + # allgather_objects: comm.allgather; bcast_int: comm.bcast; barrier: comm.Barrier + # abort: comm.Abort(code) + +def make_communicator(comm: Any | None, xp: ModuleType) -> Communicator: + # None or comm.Get_size() == 1 -> SingleComm() + # host backend -> MpiHostComm(comm); GPU backend -> NcclComm(comm, xp) (Task 3) +``` + +Test harness (`tests/conftest.py`): + +```python +@pytest.fixture +def mpirun(tmp_path) -> Callable[..., MpiResult]: + # skip if shutil.which("mpiexec") is None, or if + # `python -c "from mpi4py import MPI"` fails in a subprocess + # run(script: str, n: int, *args: str, timeout: float = 120) -> MpiResult + # command: mpiexec [--oversubscribe if Open MPI] -n n sys.executable + # tests/mpi_scripts/