mpi-rs 0.1.0

A pure-Rust implementation of the Message Passing Interface (MPI), API-compatible with rsmpi. No C library required.
Documentation
//! Integration tests: launch the example programs under the bundled `mpiexec`
//! across several process counts and assert they pass.

use std::path::PathBuf;
use std::process::{Command, Stdio};
use std::sync::Once;
use std::time::{Duration, Instant};

/// Directory holding the built example binaries (sibling `examples/` dir of the
/// launcher binary; works for both debug and release profiles).
fn examples_dir() -> PathBuf {
    PathBuf::from(env!("CARGO_BIN_EXE_mpiexec"))
        .parent()
        .unwrap()
        .join("examples")
}

static BUILD: Once = Once::new();

fn ensure_examples_built() {
    // Build all examples (including the `derive`-gated one) exactly once, even
    // when tests run in parallel.
    BUILD.call_once(|| {
        // Build the examples with the same features the test was built with, so
        // e.g. `cargo test --features "derive shm"` exercises the shm transport.
        let features = if cfg!(feature = "shm") {
            "derive,shm"
        } else {
            "derive"
        };
        let status = Command::new(env!("CARGO"))
            .args(["build", "--examples", "--features", features, "--quiet"])
            .status()
            .expect("failed to invoke cargo build");
        assert!(status.success(), "cargo build --examples failed");
    });
}

/// Run a command, killing it if it exceeds `timeout`. Returns (success, output).
fn run_with_timeout(cmd: &mut Command, timeout: Duration) -> (bool, String) {
    let mut child = cmd
        .stdout(Stdio::piped())
        .stderr(Stdio::piped())
        .spawn()
        .expect("failed to spawn");
    let start = Instant::now();
    loop {
        if let Some(status) = child.try_wait().expect("try_wait failed") {
            let out = child.wait_with_output().expect("wait_with_output");
            let mut s = String::from_utf8_lossy(&out.stdout).into_owned();
            s.push_str(&String::from_utf8_lossy(&out.stderr));
            return (status.success(), s);
        }
        if start.elapsed() > timeout {
            let _ = child.kill();
            let _ = child.wait();
            return (false, "TIMEOUT".to_string());
        }
        std::thread::sleep(Duration::from_millis(20));
    }
}

/// Launch `example` under `mpiexec -n nprocs` and assert it prints `expect`.
fn run_example(example: &str, nprocs: usize, expect: &str) {
    ensure_examples_built();
    let bin = examples_dir().join(example);
    assert!(bin.exists(), "example {example} not found at {bin:?}");

    let mut cmd = Command::new(env!("CARGO_BIN_EXE_mpiexec"));
    cmd.arg("-n").arg(nprocs.to_string()).arg(&bin);

    let (ok, output) = run_with_timeout(&mut cmd, Duration::from_secs(30));
    assert!(ok, "{example} n={nprocs} failed. Output:\n{output}");
    assert!(
        output.contains(expect),
        "{example} n={nprocs} missing {expect:?}. Output:\n{output}"
    );
}

#[test]
fn selftest_1() {
    run_example("selftest", 1, "SELFTEST PASS");
}
#[test]
fn selftest_2() {
    run_example("selftest", 2, "SELFTEST PASS");
}
#[test]
fn selftest_4() {
    run_example("selftest", 4, "SELFTEST PASS");
}
#[test]
fn selftest_7() {
    run_example("selftest", 7, "SELFTEST PASS");
}

#[test]
fn rma_2() {
    run_example("rma", 2, "RMA PASS");
}
#[test]
fn rma_4() {
    run_example("rma", 4, "RMA PASS");
}

#[test]
fn parallel_io_2() {
    run_example("parallel_io", 2, "IO PASS");
}
#[test]
fn parallel_io_5() {
    run_example("parallel_io", 5, "IO PASS");
}

#[test]
fn topo_inter_3() {
    run_example("topo_inter", 3, "TOPO/INTER PASS");
}
#[test]
fn topo_inter_6() {
    run_example("topo_inter", 6, "TOPO/INTER PASS");
}

#[test]
fn derive_2() {
    run_example("derive", 2, "derive(Equivalence) round-trip OK");
}

/// A blocked peer must not hang when another rank aborts: the job should exit
/// (non-zero) well within the timeout rather than being killed by it.
#[test]
fn abort_does_not_hang() {
    ensure_examples_built();
    let bin = examples_dir().join("abort");
    let mut cmd = Command::new(env!("CARGO_BIN_EXE_mpiexec"));
    cmd.arg("-n").arg("2").arg(&bin);
    let (ok, output) = run_with_timeout(&mut cmd, Duration::from_secs(15));
    assert!(
        !ok,
        "abort example should exit non-zero, got success:\n{output}"
    );
    assert!(
        !output.contains("TIMEOUT"),
        "abort example hung (a blocked peer was not torn down)"
    );
}

#[test]
fn spawn_singleton() {
    run_example("spawn", 1, "SPAWN PASS");
}
#[test]
fn spawn_2() {
    run_example("spawn", 2, "SPAWN PASS");
}

#[test]
fn thread_multiple_1() {
    run_example("thread_multiple", 1, "THREAD PASS");
}
#[test]
fn thread_multiple_4() {
    run_example("thread_multiple", 4, "THREAD PASS");
}

#[test]
fn bigmsg_2() {
    run_example("bigmsg", 2, "BIGMSG PASS");
}
#[test]
fn bigmsg_4() {
    run_example("bigmsg", 4, "BIGMSG PASS");
}