processkit

Async child-process management for Rust + tokio with kernel-backed whole-tree containment: every process you start — and everything it spawns — lives in a kill-on-drop container (a Linux cgroup v2, a Windows Job Object, or a POSIX process group). The active mechanism is observable, including the weaker POSIX fallback and its documented caveats.

Beyond spawning a subprocess: run-and-capture, line streaming, interactive stdin, shell-free pipelines, readiness probes, timeouts & cancellation, supervision with restart/backoff, and a mockable runner seam for subprocess-free tests.

Crates.io Docs.rs CI License: MIT MSRV

use processkit::Command;

#[tokio::main]
async fn main() -> processkit::Result<()> {
    let version = Command::new("cargo").arg("--version").run().await?;
    println!("{version}");
    Ok(())
}

Cover

Why processkit?

std::process and tokio::process reach (at most) the direct child. The processes it spawned — a build tool's compiler children, the real payload behind a wrapper (cmd /c …, sh -c …), a test's helper servers — survive a timeout, a panic, or a dropped future, and keep running as orphans.

processkit spawns every child into the operating system's own containment primitive — a Job Object on Windows, a cgroup v2 on Linux (with a process-group fallback), a POSIX process group on macOS/BSD — so teardown is a kernel operation over the whole tree, not a best-effort signal to one pid:

  • Nothing escapes silently. Dropping the handle or group reaps every descendant, grandchildren included. Where a mechanism has a genuine weakness (a setsid child escapes a POSIX process group), the active Mechanism is reported instead of pretending — never a silent downgrade.
  • Async-first. Run-and-capture, line streaming, interactive stdin, readiness probes, shell-free pipelines, supervision — all tokio futures.
  • Honest results. A non-zero exit is data (ProcessResult) until you ask for success; a timeout is captured in the result; a cancellation is always an error; every platform divergence is typed or documented.
  • Testable. One trait seam (ProcessRunner) swaps the real spawner for scripted doubles or record/replay cassettes — no subprocess in your tests.

How it compares

whole-tree kill-on-dropasynclimits / statsstreaming · pipelines · supervision
std::process
tokio::process
command-group
async-process✓ (smol)
ductpipelines
processkit (tokio)

The first column is the differentiator: a child's descendants are contained and reaped as a unit (Job Object / cgroup v2 / process group), not just the direct child.

Status: stable — current release 2.2.1. The public API follows Semantic Versioning: breaking changes land only in a new major version, so 2.x upgrades are backward-compatible. See CHANGELOG.md.

Install

cargo add processkit

This crate requires a tokio runtime. Minimum Supported Rust Version: 1.88.

Picking a verb

Every run starts with the same builder; the verb you finish with decides what you get back:

You wantCallYou get
stdout, success required.run()trimmed String; non-zero exit / timeout / kill → typed Error
the full outcome, exit code as data.output_string() / .output_bytes()ProcessResult — code, stdout, stderr, timed_out; never errors on non-zero
just the exit code.exit_code()i32 (a timed-out / killed run errors instead of inventing -1)
a yes/no answer.probe()bool — exit 0 → true, 1 → false, anything else errors
a typed value from stdout.parse(|s| …) / .try_parse(|s| …)T — success required; fails loud on a truncated capture
the first matching output line.first_line(|l| …)Option<String>None when stdout closes without a match
a live handle — streaming, stdin, probes.start()RunningProcess

The same vocabulary repeats on every layer (ProcessRunner, CliClient), and processkit::run("git", ["status"]) / processkit::output_string(…) skip the builder for one-liners.

Quick start

use processkit::{Command, ProcessGroup, Stdin};

#[tokio::main]
async fn main() -> processkit::Result<()> {
    // Capture output; a non-zero exit does not error on its own.
    let result = Command::new("git").args(["rev-parse", "HEAD"]).output_string().await?;
    println!("HEAD is {}", result.stdout().trim());

    // Require success and get trimmed stdout directly.
    let version = Command::new("cargo").arg("--version").run().await?;
    println!("{version}");

    // Feed stdin.
    let sorted = Command::new("sort")
        .stdin(Stdin::from_string("banana\napple\n"))
        .output_string()
        .await?;
    println!("{}", sorted.stdout());

    // Share one kill-on-drop group across several children; dropping the group
    // reaps the whole tree.
    let group = ProcessGroup::new()?;
    let _server = group.start(&Command::new("some-server")).await?;
    // ... work ...
    group.shutdown().await?; // graceful SIGTERM → wait → SIGKILL (Unix); atomic on Windows

    Ok(())
}

Documentation

This page is the quick tour. The guide set in the sidebar goes deeper on every capability, with more examples and the platform fine print collected in one place. New here? Skim the Cookbook first — it maps "I want to …" tasks to working snippets — then read Running commands end to end:

GuideCovers
CookbookTask → snippet recipes for everything below; the fastest way in
Running commandsThe full Command builder and every consuming verb, with error semantics
Process groupsContainment, teardown, signals, suspend/resume, members, limits, stats
Streaming & interactive I/OLine streaming, conversational stdin, readiness probes, wait_any, profiling
PipelinesShell-free a | b | c, pipefail attribution, chain timeouts
Timeouts, retries & cancellationCaptured vs raised deadlines, retry classifiers, CancellationToken
SupervisionRestart policies, backoff & jitter, stop conditions, outcomes
Testing your codeThe ProcessRunner seam, scripted/recording/mock doubles, cassettes, CliClient
Platform supportMechanisms, all capability matrices, every caveat
UpgradingPer-version consumer upgrade notes — what changed on each release and how to migrate across a major bump

API reference: docs.rs/processkit.

Using Python instead of Rust? The processkit-py wrapper provides matching synchronous and asyncio APIs over this same native core, including context-manager teardown, streaming, pipelines, supervision, resource limits, and the runner testing seam.

Source, changelog, and development history live in the Rust repository.

Feature flags

Each flag is additive and only gates visibility — the kill-on-drop tree guarantee is unconditional in every configuration.

FeatureDefaultAdds
process-controlSignal, ProcessGroup::{signal, suspend, resume, members, adopt}
statsgroup/per-run resource measurement, sample_stats, profile (enables an additional windows-sys API feature on Windows)
limitswhole-tree resource caps (implies stats)
recordrecord/replay cassettes (pulls serde)
mockmockall-generated MockRunner (test-only; its surface is semver-exempt — prefer ScriptedRunner/RecordingRunner)
tracinglifecycle events: spawn/exit, timeout/cancel, teardown, retries, storms (never argv/env)

Capping a group's resources

Requires the limits feature (off by default) — add it to the dependency:

processkit = { version = "…", features = ["limits"] }

ProcessGroupOptions can then bound the whole tree's memory, process count, and CPU at creation, so a runaway or untrusted child tree can't exhaust the host:

use processkit::{Command, ProcessGroup, ProcessGroupOptions};

#[tokio::main]
async fn main() -> processkit::Result<()> {
    let group = ProcessGroup::with_options(
        ProcessGroupOptions::default()
            .max_memory(512 * 1024 * 1024) // 512 MiB across the tree
            .max_processes(64)
            .cpu_quota(0.5),               // half of one core
    )?;
    let _job = group.start(&Command::new("untrusted-tool")).await?;
    // ... work ...
    Ok(())
}

cpu_quota is a fraction of a single core (0.5 = half a core, 2.0 = two cores); on Windows it is converted against the host's CPU count and is approximate.

Limits need a real container — a Windows Job Object or a Linux cgroup v2. There is no whole-tree limit on macOS/the BSDs or the Linux process-group fallback. A Linux cgroup limit additionally needs this process to run at the real cgroup-v2 root: the crate creates the limit cgroup under this process's own cgroup and enables the controllers there, which cgroup v2's "no internal processes" rule allows only for the real hierarchy root — not a cgroup-namespace root (so an ordinary container EBUSYs too), and not under a systemd session/scope/service. The crate does not migrate your process to work around it, so in practice limits apply only at a minimal non-systemd init. When a requested limit can't be enforced, with_options returns Error::ResourceLimit instead of a silently-unbounded group — an unapplied cap is no protection.

Deeper: Process groups → resource limits.

Signalling and pausing the whole tree

Beyond the kill/shutdown teardown verbs, a group can broadcast a signal to every member or freeze and thaw the whole tree:

use processkit::{Command, ProcessGroup, Signal};

#[tokio::main]
async fn main() -> processkit::Result<()> {
    let group = ProcessGroup::new()?;
    let _server = group.start(&Command::new("my-server")).await?;

    group.signal(Signal::Hup)?; // e.g. "reload configuration"
    group.suspend()?;           // freeze the whole tree…
    group.resume()?;            // …and let it run again
    Ok(())
}

Signals are POSIX-only: on Windows just Signal::Kill is deliverable (it maps to the Job Object terminate) and anything else returns Error::Unsupported. Signal::Kill always takes the same whole-tree hard-kill path as kill_all(). Suspend/resume work everywhere a container exists — one cgroup.freeze write covering the subtree on Linux, SIGSTOP/SIGCONT on macOS/BSD and the Linux process-group fallback (both idempotent), and per-thread suspension on Windows (best-effort; only there nested suspends stack and need matching resumes).

Deeper: Process groups → signals, suspend/resume.

Inspecting the tree

members() snapshots the live member pids, and wait_any races several running processes, reporting whichever exits first — the natural primitive for supervising a few long-lived children:

use processkit::{Command, ProcessGroup, wait_any};

#[tokio::main]
async fn main() -> processkit::Result<()> {
    let group = ProcessGroup::new()?;
    let mut a = group.start(&Command::new("server-a")).await?;
    let mut b = group.start(&Command::new("server-b")).await?;

    println!("live pids: {:?}", group.members()?);

    // Borrows only: the loser stays usable after the race.
    let (idx, outcome) = wait_any(&mut [&mut a, &mut b]).await?;
    println!("contender #{idx} exited first with {outcome:?}");
    Ok(())
}

members() lists the whole tree on Windows (Job Object) and Linux (cgroup); the POSIX process-group backends list the tracked group leaders only. (members is part of the default-on process-control feature; wait_any is always available.) wait_any applies no per-process timeout (bound the race with tokio::time::timeout) and does no output pumping — drain chatty children first.

Deeper: Process groups → members · Streaming → racing children.

Running many at once

wait_any's siblings cover the join and fan-out cases. wait_all joins a fixed set of handles you already hold, returning every outcome in order; output_all runs a whole batch of commands with a concurrency cap, so fanning out hundreds of commands can't exhaust file descriptors or the process table:

use processkit::{Command, JobRunner, wait_all, output_all};

#[tokio::main]
async fn main() -> processkit::Result<()> {
    let group = processkit::ProcessGroup::new()?;
    let mut a = group.start(&Command::new("worker-a")).await?;
    let mut b = group.start(&Command::new("worker-b")).await?;
    let outcomes = wait_all(&mut [&mut a, &mut b]).await?; // both, in input order

    // 200 conversions, but never more than 8 processes alive at once.
    let cmds = (0..200).map(|i| Command::new("convert").arg(format!("{i}.png")));
    let results = output_all(cmds, 8, &JobRunner).await;
    let failed = results.iter().filter(|r| !matches!(r, Ok(o) if o.is_success())).count();
    println!("{:?}; {failed} conversions failed", outcomes);
    Ok(())
}

output_all is collect-all: each element is one command's independent Result, so a non-zero exit (an Ok with a non-zero code) never short-circuits the batch — the caller folds the outcomes. Pass &group instead of &JobRunner to keep every child in one shared kill-on-drop group. It is deliberately not a pool, scheduler, or retrier. For batches whose stdout is binary, output_all_bytes is the identical fan-out with each result captured as Vec<u8>. wait_all shares wait_any's two non-features (no per-process timeout, no output pumping).

Sampling stats over time

A point-in-time stats() becomes a series with sample_stats, and a single run can be profiled end-to-end (requires the opt-in stats feature):

use processkit::prelude::StreamExt;
use processkit::{Command, ProcessGroup};
use std::time::Duration;

#[tokio::main]
async fn main() -> processkit::Result<()> {
    // A CPU/RSS/process-count series for a whole group:
    let group = ProcessGroup::new()?;
    let _worker = group.start(&Command::new("worker")).await?;
    let mut samples = group.sample_stats(Duration::from_millis(250));
    if let Some(s) = samples.next().await {
        println!("procs={} rss={:?}", s.active_process_count, s.peak_memory_bytes);
    }
    drop(samples);

    // …or a one-shot summary of a single run:
    let profile = Command::new("crunch")
        .start().await?
        .profile(Duration::from_millis(100)).await?;
    println!(
        "outcome={:?} took={:?} peak_rss={:?} avg_cpu_cores={:?}",
        profile.outcome, profile.duration, profile.peak_memory_bytes, profile.avg_cpu_cores(),
    );
    Ok(())
}

The series inherits stats()'s platform matrix (full CPU/memory on Windows and Linux cgroup; counts only on the POSIX process-group backends); profile samples the started child process itself and applies the run's normal timeout/output handling.

Deeper: Process groups → stats · Streaming → profiling a run.

Supervising a long-lived child

Where Command::retry replays one run until it succeeds, a Supervisor keeps a child alive: it restarts the command per policy whenever it exits, with bounded restarts and exponential backoff (jittered by default so a restarted fleet doesn't stampede):

use processkit::{Command, RestartPolicy, Supervisor};
use std::time::Duration;

#[tokio::main]
async fn main() -> processkit::Result<()> {
    let outcome = Supervisor::new(Command::new("my-server").args(["--port", "8080"]))
        .restart(RestartPolicy::OnCrash)          // Always | OnCrash | Never
        .max_restarts(5)
        .backoff(Duration::from_millis(200), 2.0) // base, multiplier (cap: .max_backoff)
        .storm_pause(Duration::from_secs(15))     // crash-loop guard (off by default)
        .stop_when(|res| res.code() == Some(0))   // a clean exit ends supervision
        .run()
        .await?;
    println!("ended after {} restarts: {:?}", outcome.restarts, outcome.stopped);
    Ok(())
}

run() reports a SupervisionOutcome — the final run's result, the restart count, and why supervision stopped. The opt-in failure-storm guard distinguishes "fails rarely" from "crash-looping": each failure feeds a score that halves every failure_decay; past failure_threshold the supervisor takes one collective storm_pause instead of hammering restarts at backoff speed. Supervision is platform-agnostic and runs through the ProcessRunner seam: pass .with_runner(&group) to keep every incarnation in one shared kill-on-drop group, or a ScriptedRunner to test supervision logic hermetically.

Deeper: Supervision.

Waiting for a child to be ready

"Start a server, then use it" needs the server to be ready, not merely started. Three probes replace the arbitrary sleep:

use processkit::Command;
use std::time::Duration;

#[tokio::main]
async fn main() -> processkit::Result<()> {
    let mut run = Command::new("my-server").start().await?;

    // Wait for the startup banner (returns the matching line)…
    let banner = run
        .wait_for_line(|l| l.contains("listening on"), Duration::from_secs(10))
        .await?;
    println!("server says: {banner}");

    // …or for the port to accept connections…
    let addr = "127.0.0.1:8080".parse().expect("valid socket address");
    run.wait_for_port(addr, Duration::from_secs(10)).await?;

    // …or for any async health check to pass.
    run.wait_for(|| async { health_check().await }, Duration::from_secs(10)).await?;

    // ready — use the server …
    Ok(())
}

async fn health_check() -> bool {
    // e.g. probe an HTTP /health endpoint
    true
}

A probe that doesn't pass within its deadline — or that can no longer pass (the child exits; for wait_for_line, its stdout closes) — fails with Error::NotReady (distinct from Error::Timeout, which is the run's own Command::timeout) and does not kill the child: the caller decides what happens next. wait_for_line consumes stdout up to the match (continue with finish); wait_for_port / wait_for don't touch the pipes at all.

Deeper: Streaming → readiness probes.

Pipelines without a shell

a | b | c without a shell string — stages connected in-process (a relay, not a shell), so no quoting or injection surface, and every stage lives in one shared kill-on-drop group:

use processkit::Command;

#[tokio::main]
async fn main() -> processkit::Result<()> {
    let out = Command::new("git").args(["log", "--format=%an"])
        .pipe(Command::new("sort"))
        .pipe(Command::new("uniq").arg("-c"))
        .output_string()
        .await?;
    println!("{}", out.stdout());
    Ok(())
}

The | operator is equivalent sugar: (a | b | c).run().

The outcome is pipefail: stdout is the last stage's output, while the exit code, stderr, and reported program come from the first stage that didn't exit cleanly (or the last stage when all succeed). For a consumer that legitimately stops reading early — the producer | head -1 shape, where the producer's broken-pipe death (its next write fails once the downstream closes) is expected — mark that stage .unchecked_in_pipe() and pipefail skips it (a checked failure still always wins). The first stage's stdin source is honored; inner stages read from the pipe. .timeout(d) bounds the whole chain (killing every stage at the deadline), and run() requires every stage to succeed, returning the trimmed final stdout.

Deeper: Pipelines.

Environment and privileges

Spawn-time controls for sandboxing and service launch:

use processkit::Command;

#[tokio::main]
async fn main() -> processkit::Result<()> {
    Command::new("worker")
        .inherit_env(["PATH", "HOME", "LANG"]) // allow-list on a cleared env
        .uid(1000).gid(1000)                   // Unix: drop privileges
        .setsid()                              // Unix: new session
        .run()
        .await?;

    Command::new("helper")
        .create_no_window()                    // Windows: no console window
        .run()
        .await?;

    Command::new("daemonish")
        .kill_on_parent_death()                // die with a SIGKILLed parent
        .start()
        .await?;
    Ok(())
}

inherit_env clears the environment and copies only the named parent vars (explicit env/env_remove still apply on top); it works everywhere. uid / gid (group id is set before user id) and setsid are POSIX-only — on other targets the run fails with Error::Unsupported rather than silently skipping a privilege drop. One Linux caveat: under the cgroup mechanism the child joins its cgroup after the uid has already dropped, and the auto-created cgroup isn't writable by the target user — the spawn fails with a permission error (never an uncontained child); privilege drop currently composes cleanly with the process-group mechanism. setsid keeps containment: the group tracks the new session's process group. create_no_window is a harmless no-op outside Windows and, unlike the raw ProcessGroup::spawn escape hatch, survives the group's CREATE_SUSPENDED containment flag (they are OR'd together). kill_on_parent_death hardens the one case Drop can't cover — the parent dying abruptly (SIGKILL): Windows already guarantees it (the job handle closes with the process), Linux arms PDEATHSIG on the direct child, macOS/BSD have no equivalent (documented no-op).

Deeper: Running commands → privileges and spawn flags.

Cancelling a run

Hand a command a CancellationToken (re-exported from tokio-util); cancelling the token kills the process tree, and every consuming path reports Error::Cancelled:

use processkit::{CancellationToken, Command};

#[tokio::main]
async fn main() -> processkit::Result<()> {
    let token = CancellationToken::new();

    let job = tokio::spawn({
        let token = token.child_token();
        async move {
            Command::new("long-job").cancel_on(token).run().await
        }
    });

    // elsewhere — a shutdown signal, a sibling failure, a UI button:
    token.cancel();

    assert!(matches!(job.await.unwrap(), Err(processkit::Error::Cancelled { .. })));
    Ok(())
}

Unlike a timeout — whose expiry is captured in the result as timed_out — cancellation is always an error: the run was abandoned, so there is no result to inspect. When a cancel and a timeout land together, cancellation wins. A token cancelled before the run starts short-circuits without spawning anything. On a shared ProcessGroup handle, cancelling kills the child itself but leaves the group's siblings alone (same scope as a timeout), and a supervised command that gets cancelled stops its Supervisor for good — restarting into a still-cancelled token would loop futilely.

For a typed wrapper whose commands never cross your code, set the token once on the client: CliClient::new("gh").default_cancel_on(token.child_token()) — cancelling it kills every in-flight command of that client.

Deeper: Timeouts, retries & cancellation.

Async streaming and interactive I/O

The one-shot helpers above buffer the whole output. For long-running or conversational children, start() returns a live RunningProcess you can drive asynchronously.

Stream stdout line by line

Process each line as it arrives — no waiting for the child to exit, no buffering the full output. StreamExt (from processkit::prelude, re-exported from tokio-stream) provides .next():

use processkit::prelude::StreamExt;
use processkit::{Command, Outcome, Finished};

#[tokio::main]
async fn main() -> processkit::Result<()> {
    let mut run = Command::new("git")
        .args(["log", "--oneline", "-n", "50"])
        .start()
        .await?;

    let mut lines = run.stdout_lines()?;
    while let Some(line) = lines.next().await {
        println!("commit: {line}");
    }

    // After the stream ends, collect the outcome and whatever went to stderr
    // (drained in the background while you streamed stdout). The `Outcome`
    // distinguishes a clean exit, a signal kill, and a timeout.
    let Finished { outcome, stderr, .. } = run.finish().await?;
    if outcome != Outcome::Exited(0) {
        eprintln!("git ended {outcome:?}: {stderr}");
    }
    Ok(())
}

The command's timeout bounds the stream: at the deadline the tree is killed, the pipes close, and the stream ends (on a handle that owns its group — the start() path). A cancel_on token ends the stream the same way, and the following finish reports Error::Cancelled. For an ad-hoc bound, wrapping the loop in tokio::time::timeout and dropping the handle (which kills the tree) still works.

Interactive stdin — write requests, read responses

Keep stdin open with keep_stdin_open(), take the writer with take_stdin(), then interleave async writes and reads:

use processkit::prelude::StreamExt;
use processkit::Command;

// `ProcessStdin`'s writer methods return `std::io::Result` (idiomatic for a
// writer), so this example uses `Box<dyn std::error::Error>` to mix them with the
// crate's `Result`; in a `processkit::Result` function, `.map_err(processkit::Error::Io)?`.
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // `bc` evaluates each stdin line and prints the result on stdout.
    let mut run = Command::new("bc").keep_stdin_open().start().await?;

    let mut stdin = run.take_stdin().expect("stdin was kept open");
    stdin.write_line("2 + 2").await?;
    stdin.write_line("6 * 7").await?;
    stdin.finish().await?; // send EOF so bc finishes

    let mut answers = run.stdout_lines()?;
    while let Some(answer) = answers.next().await {
        println!("bc says: {answer}");
    }
    Ok(())
}

Writing the whole input before reading is fine for a small payload like this. For a large interactive stdin, drain stdout concurrently (write from one task, read stdout_lines from another) — otherwise the child can block writing stdout while you block writing stdin, a full-duplex deadlock. See the note in Streaming & interactive I/O. (The non-interactive Stdin::from_* sources are written on a background task and never deadlock.)

Feed stdin from an async stream, react to stdout as it's read

Stdin::from_lines writes each item of any Stream<Item = String> as a line — back it with a channel, a file tail, or a network source. Pair it with on_stdout_line / on_stderr_line to handle output inline (the handler runs on the read pump, in addition to capture):

use processkit::{Command, Stdin};
use tokio_stream::iter; // any `Stream<Item = String>` works

#[tokio::main]
async fn main() -> processkit::Result<()> {
    let input = iter(vec!["banana".to_owned(), "apple".to_owned(), "cherry".to_owned()]);

    let result = Command::new("sort")
        .stdin(Stdin::from_lines(input))
        .on_stdout_line(|line| println!("sorted: {line}"))
        .output_string()
        .await?;
    let _ = result; // already printed line by line above
    Ok(())
}

Deeper: Streaming & interactive I/O.

Wrapping a CLI tool

CliClient + the cli_client! macro turn a typed wrapper around an external tool (git, jj, gh, …) into just its parsers — the runner is injectable, so the wrapper is hermetically testable with a ScriptedRunner (no subprocess). The seam covers streaming too: a scripted start() feeds canned lines through the same pump machinery, so stdout_lines/wait_for_line-based orchestration tests hermetically as well:

#![allow(unused)]
fn main() {
use processkit::{cli_client, ProcessRunner, Result};
use std::path::Path;

cli_client!(pub struct Git => "git");

impl<R: ProcessRunner> Git<R> {
    async fn head(&self, dir: &Path) -> Result<String> {
        self.core.run(self.core.command_in(dir, ["rev-parse", "HEAD"])).await
    }
}
}

Deeper: Testing your code → CliClient.

Recording and replaying runs

Requires the record feature (off by default). RecordReplayRunner turns real runs into a JSON cassette once, then replays them deterministically — fast, hermetic, no subprocess in CI:

use processkit::{Command, JobRunner, ProcessRunnerExt};
use processkit::testing::RecordReplayRunner;

#[tokio::main]
async fn main() -> processkit::Result<()> {
    // Record once against the real tool (e.g. an opt-in `--record` test run):
    let runner = RecordReplayRunner::record("fixtures/git.json", JobRunner::new());
    let version = runner.run(&Command::new("git").arg("--version")).await?;
    runner.save()?; // or best-effort on drop

    // Replay everywhere else — no subprocess, identical results:
    let runner = RecordReplayRunner::replay("fixtures/git.json")?;
    assert_eq!(runner.run(&Command::new("git").arg("--version")).await?, version);
    Ok(())
}

Entries are matched by program + args + cwd + a stdin source digest (hashed, never persisted: in-memory bytes hash their content, a from_file source hashes its path). Environment override values never reach the file — only the sorted variable names, so env values can't leak through a committed fixture (and env differences can't cause spurious misses). Note: argv, cwd, stdout, and stderr are stored verbatim and can carry secrets (a --password=… flag, a token echoed to output) — review a fixture before committing it; on Unix the file is written 0600. When one invocation was recorded several times, replay serves the entries in capture order and then repeats the last one — a recorded sequence of changing outputs replays faithfully, while retry/probe loops keep getting a stable final answer. An invocation absent from the cassette is a strict Error::CassetteMiss (distinct from a missing program, so is_not_found() is false; replay never spawns a surprise subprocess), and the file carries a format version so future readers fail loudly instead of misreading old fixtures.

Deeper: Testing your code → record/replay.

Contributing

Running the tests and the maintainer-only release process are documented in the contribution guide.

License

Licensed under the MIT License.

Python wrapper

‹ docs index

processkit-py brings ProcessKit's kernel-backed process-tree containment to Python. It is a thin, typed PyO3 binding over the Rust crate rather than a separate implementation: Job Objects, cgroup v2, POSIX process groups, process I/O, and the teardown state machine remain in Rust.

The Python layer adds the conventions Python applications need:

  • matching synchronous and asyncio APIs;
  • deterministic with / async with teardown;
  • Python exception types with structured fields;
  • complete type stubs and a ProcessRunner protocol;
  • prebuilt wheels with the optional Rust capabilities already enabled.

processkit-py on PyPI

Published package: processkit-py on PyPI
Source: ZelAnton/processkit-py
Full guides: processkit-py/docs.

Install

The distribution name is processkit-py; the import name is processkit:

pip install processkit-py

Version 1.0 supports CPython 3.10 and newer. The project publishes abi3 wheels for regular CPython and a version-specific wheel for free-threaded CPython 3.14t. A source distribution is available for platforms without a matching wheel.

The wrapper currently pins Rust processkit 1.2.0 and enables process-control, limits, record, and tracing. Unlike Cargo's optional feature model, every Python wheel exposes one complete API; there are no extras to select at installation time.

Sync and asyncio

Run-to-completion operations have a synchronous name and an a-prefixed asyncio twin. Both return the same result types and use the same Rust runtime semantics:

OperationSynchronousAsyncio
capture a full resultoutput()await aoutput()
require successrun()await arun()
read the exit codeexit_code()await aexit_code()
run a 0/1 predicateprobe()await aprobe()
start a live processstart()await astart()
from processkit import Command

# Non-zero exit is data: inspect code, stdout, stderr, and timed_out.
result = Command("git", ["status", "--short"]).output()

# Require an accepted exit and return trimmed stdout.
version = Command("python", ["--version"]).run()

The same operations compose naturally under asyncio:

import asyncio
from processkit import Command

async def main() -> None:
    result = await Command("git", ["rev-parse", "HEAD"]).aoutput()
    print(result.stdout.strip())

asyncio.run(main())

Cancelling an awaited run drops the underlying Rust future and tears down its owned process tree. Blocking calls release the GIL while Rust waits and check for Python signals, so Ctrl+C can interrupt a synchronous run instead of waiting for the child to finish.

Whole-tree containment

ProcessGroup is a synchronous and asynchronous context manager. Everything started into it shares one containment object and is torn down when the block exits:

from processkit import Command, ProcessGroup

with ProcessGroup() as group:
    group.start(Command("dev-server"))
    group.start(Command("worker"))
    print(group.mechanism)  # job_object | cgroup_v2 | process_group
    print(group.members())

Use the context managers for deterministic cleanup. Windows Job Objects also close when the process loses its last kernel handle. On Unix, cleanup requires Python to execute the teardown path; an uncatchable SIGKILL cannot run it. The POSIX process-group backend also has the same documented setsid escape caveat as the Rust crate. Inspect group.mechanism when the stronger cgroup or Job Object guarantee matters.

Group resource limits, statistics, signals, suspend/resume, kill_all(), and graceful shutdown are exposed directly. Limits require a Windows Job Object or a writable Linux cgroup-v2 root; an unenforceable limit raises ResourceLimit instead of being silently ignored.

Streaming and interactive processes

RunningProcess exposes asyncio-native streaming and stdin. Its result methods consume the handle, matching Rust ownership explicitly: after wait(), finish(), output(), output_bytes(), profile(), or shutdown(), the Python handle is spent.

from processkit import Command

async with await Command("my-build", ["--watch"]).astart() as proc:
    async for line in proc.stdout_lines():
        print(line)

Interactive commands use keep_stdin_open() and take_stdin(). Readiness is provided by the composable Python helpers wait_for, wait_for_line, and wait_for_port.

Pipelines, supervision, and batch runs

The higher-level Rust features have Python-native surfaces:

  • Command(...) | Command(...) creates a shell-free Pipeline with pipefail attribution;
  • Supervisor keeps a command alive using restart, backoff, and storm-guard policies;
  • output_all / aoutput_all run a bounded batch and preserve input order;
  • CliClient binds one executable to shared timeout and environment defaults.

Builders are immutable: configuration methods return a new Command, so a base command can be safely reused for several variants.

Errors and results

Like the Rust API, output() treats a non-zero exit and timeout as inspectable result data. Success-checking verbs raise typed exceptions when the outcome is not accepted.

Python exceptionAdditional compatibility
NonZeroExitcarries program, code, stdout, and stderr
Timeoutalso a built-in TimeoutError
ProcessNotFoundalso a built-in FileNotFoundError
PermissionDeniedalso a built-in PermissionError
Signalledcarries the signal and captured output
ResourceLimit / Unsupported / OutputTooLargepreserves the corresponding Rust failure category

Every package-specific exception derives from ProcessError.

Testing without subprocesses

Application code can depend on the typed ProcessRunner protocol and receive a real Runner in production. Tests can inject the doubles exported from processkit.testing:

  • ScriptedRunner returns canned replies and supports scripted streaming;
  • RecordingRunner records invocations for assertions;
  • RecordReplayRunner records real runs to a cassette and replays them without spawning;
  • Reply describes successful, failed, timed-out, signalled, pending, or line-streamed outcomes.
from processkit import Command
from processkit.testing import Reply, ScriptedRunner

runner = ScriptedRunner()
runner.on(["git", "rev-parse"], Reply.ok("abc123\n"))

assert runner.run(Command("git", ["rev-parse", "HEAD"])) == "abc123"

The package ships py.typed and a complete hand-written .pyi surface, checked against the compiled extension with mypy.stubtest.

Where to continue

The Python repository contains the detailed Cookbook, migration guide from subprocess, command reference, streaming guide, and platform matrix.

Cookbook

‹ docs index

Task-oriented recipes: find the thing you're trying to do, copy the snippet, follow the link when you need the fine print. Every snippet assumes a tokio runtime and use processkit::Command; unless shown otherwise.

Run a command and get its output

#![allow(unused)]
fn main() {
let head = Command::new("git").args(["rev-parse", "HEAD"]).run().await?;
}

run() requires a zero exit and returns stdout with trailing whitespace trimmed; a non-zero exit, spawn failure, or timeout is a typed Error. For a one-liner without the builder: processkit::run("git", ["rev-parse", "HEAD"]).

Fine print: Running commands → consuming verbs.

Inspect a failure instead of erroring

#![allow(unused)]
fn main() {
let result = Command::new("git").args(["merge", "topic"]).output_string().await?;
if !result.is_success() {
    eprintln!("merge exited {:?}: {}", result.code(), result.stderr());
}
}

output_string() (and output_bytes() for raw bytes) treats the exit code as data — Err means the run couldn't happen at all. Call result.ensure_success()? later to convert a stored failure into the same typed error run() would have produced.

Fine print: Running commands → results and errors.

Ask a yes/no question

#![allow(unused)]
fn main() {
let dirty = !Command::new("git").args(["diff", "--quiet"]).probe().await?;
}

probe() maps exit 0 → true, exit 1 → false, and anything else to an error — the git diff --quiet / grep -q convention without manual code matching.

Accept non-zero exit codes as success

#![allow(unused)]
fn main() {
// `grep` exits 1 when it finds no match — not a failure for this call.
let found = Command::new("grep")
    .args(["needle", "haystack.txt"])
    .ok_codes([0, 1])
    .output_string()
    .await?;
let matched = found.code() == Some(0); // 0 = matched, 1 = no match (both "success")
}

ok_codes widens what the checking verbs (run/run_unit) and is_success/ensure_success treat as success — for tools whose non-zero exit is a normal result (grep 1 = no match, diff 1 = differs, rsync's code families). It does not change exit_code (always the raw code) or probe (always the 0/1 convention). An empty set is ignored, so the default stays exit 0.

Bound a run with a timeout

#![allow(unused)]
fn main() {
use std::time::Duration;

let result = Command::new("slow-tool")
    .timeout(Duration::from_secs(30))
    .output_string()
    .await?;
if result.timed_out() {
    eprintln!("gave up after 30s; partial output: {}", result.stdout());
}
}

At the deadline the whole tree is killed. On the capture verbs the timeout is captured (timed_out(), partial output kept); on the success-checking verbs (run, exit_code) it surfaces as Error::Timeout.

Let a tool clean up on timeout

#![allow(unused)]
fn main() {
use std::time::Duration;

let result = Command::new("dev-server")
    .timeout(Duration::from_secs(30))
    .timeout_grace(Duration::from_secs(5)) // SIGTERM, wait up to 5s, then SIGKILL
    .output_string()
    .await?;
}

timeout_grace turns the hard deadline kill into a graceful one: SIGTERM (or the signal from timeout_signal, with the process-control feature), up to the grace window to exit, then SIGKILL. A signal-handling child exits early; timed_out() stays true. Windows has no signal tier — the deadline kills atomically.

Fine print: Timeouts → graceful timeout.

Fine print: Timeouts, retries & cancellation.

Show a useful error message

#![allow(unused)]
fn main() {
if let Err(e) = Command::new("git").args(["merge", "topic"]).run().await {
    eprintln!("merge failed: {}", e.diagnostic().unwrap_or("(no output)"));
}
}

Error::diagnostic() picks the most explanatory captured text — stderr, falling back to stdout (git writes CONFLICT … there) — so callers don't re-implement the same heuristic.

Feed the child's stdin

#![allow(unused)]
fn main() {
use processkit::Stdin;

// A string you already have:
let sorted = Command::new("sort")
    .stdin(Stdin::from_string("banana\napple\n"))
    .run()
    .await?;

// …or any async source: a reader (file, socket) or a stream of lines.
let from_file = Stdin::from_reader(tokio::fs::File::open("input.txt").await?);
let from_chan = Stdin::from_lines(tokio_stream::iter(vec!["one".to_owned()]));
}

One-shot sources (from_reader/from_lines) feed a single run; re-running the same Command afterwards fails loud (an Error::Io at launch, D10) instead of silently seeing empty stdin. For a conversation, see the next recipe but one.

Fine print: Running commands → standard input.

Stream output as it arrives

#![allow(unused)]
fn main() {
use processkit::Finished;
use processkit::prelude::StreamExt; // provides `.next()`

let mut run = Command::new("cargo").args(["build", "--verbose"]).start().await?;
let mut lines = run.stdout_lines()?;
while let Some(line) = lines.next().await {
    println!("build: {line}");
}
let Finished { outcome, stderr, .. } = run.finish().await?; // outcome + buffered stderr
}

No waiting for exit, no full-output buffering; stderr is drained in the background so the child can't block. A timeout on the command bounds the stream itself. Prefer a callback? .on_stdout_line(|l| …) runs one per line while any capture verb drives the run.

Fine print: Streaming & interactive I/O.

Talk to an interactive child

#![allow(unused)]
fn main() {
use processkit::prelude::StreamExt;

let mut run = Command::new("bc").keep_stdin_open().start().await?;
let mut stdin = run.take_stdin().expect("stdin was kept open");
stdin.write_line("2 + 2").await?;
stdin.finish().await?; // EOF — bc exits

let mut answers = run.stdout_lines()?;
while let Some(answer) = answers.next().await {
    println!("{answer}");
}
}

keep_stdin_open() hands you an async writer instead of closing stdin at spawn; interleave writes with reads for request/response tools. Its writer methods return std::io::Result (idiomatic for a writer) — convert with .map_err(processkit::Error::Io)? in a processkit::Result function, or use Box<dyn std::error::Error>.

Fine print: Streaming & interactive I/O → interactive stdin.

Pipe commands without a shell

#![allow(unused)]
fn main() {
let authors = Command::new("git").args(["log", "--format=%an"])
    .pipe(Command::new("sort"))
    .pipe(Command::new("uniq").arg("-c"))
    .output_string()
    .await?;
}

Native pipes — no shell string, no quoting, no injection surface. The outcome is pipefail: stdout comes from the last stage, the reported failure from the first stage that didn't exit cleanly. All stages share one kill-on-drop group. The | operator is equivalent sugar: (a | b | c).output_string().

For a consumer that legitimately stops reading early — the | head -1 shape, where the producer's broken-pipe death (its next write fails once the downstream closes, or SIGPIPE where the OS delivers it) is expected — mark the producer unchecked_in_pipe() so that death doesn't fail the chain:

#![allow(unused)]
fn main() {
let first = (Command::new("seq").args(["1", "1000000"]).unchecked_in_pipe()
    | Command::new("head").args(["-n", "1"]))
    .run()
    .await?;
}

Fine print: Pipelines → unchecked stages.

Start a server and wait until it's ready

#![allow(unused)]
fn main() {
use std::time::Duration;

let mut server = Command::new("my-server").args(["--port", "8080"]).start().await?;

// Pick the probe that matches how the server announces readiness:
server.wait_for_line(|l| l.contains("listening"), Duration::from_secs(10)).await?;
// server.wait_for_port("127.0.0.1:8080".parse().unwrap(), Duration::from_secs(10)).await?;
// server.wait_for(|| async { http_health().await }, Duration::from_secs(10)).await?;

// …use the server; dropping `server` kills its whole tree.
}

A probe that can't succeed fails fast with Error::NotReady and never kills the child — you decide what happens next. No more sleep(2) and hoping.

Fine print: Streaming & interactive I/O → readiness probes.

Tear down several children as a unit

#![allow(unused)]
fn main() {
use processkit::ProcessGroup;

let group = ProcessGroup::new()?;
let _db = group.start(&Command::new("dev-db")).await?;
let _api = group.start(&Command::new("dev-api")).await?;

// Either: graceful — SIGTERM, bounded wait, optional SIGKILL escalation…
group.shutdown().await?;
// …or just drop(group): hard kill-on-drop of everything, grandchildren included.
}

The group is the unit of fate: a panic or early return anywhere reaps every member. Configure the grace window via ProcessGroupOptions.

Fine print: Process groups.

React to whichever child exits first

#![allow(unused)]
fn main() {
use processkit::{ProcessGroup, wait_any};

let group = ProcessGroup::new()?;
let mut a = group.start(&Command::new("worker-a")).await?;
let mut b = group.start(&Command::new("worker-b")).await?;

let (idx, outcome) = wait_any(&mut [&mut a, &mut b]).await?;
println!("worker #{idx} exited first with {outcome:?}");
// `a` and `b` are only borrowed — the loser is still usable here.
}

Fine print: Streaming & interactive I/O → racing children.

Sandbox an untrusted tool

#![allow(unused)]
fn main() {
use processkit::{ProcessGroup, ProcessGroupOptions};

// Cap the whole tree (requires the `limits` feature; Windows Job / Linux cgroup):
let group = ProcessGroup::with_options(
    ProcessGroupOptions::default()
        .max_memory(512 * 1024 * 1024)
        .max_processes(64)
        .cpu_quota(0.5),
)?;

let result = group
    .start(
        &Command::new("untrusted-tool")
            .inherit_env(["PATH"]) // allow-list: everything else is cleared
            .timeout(std::time::Duration::from_secs(60)),
    )
    .await?
    .output_string()
    .await?;
}

Unenforceable limits are a hard Error::ResourceLimit, never a silently unbounded group. On Unix, add .uid(…)/.gid(…) to drop privileges (note the cgroup-mechanism caveat in the guide).

Fine print: Process groups → resource limits · Running commands → privileges.

Keep a crash-prone service running

#![allow(unused)]
fn main() {
use processkit::{RestartPolicy, Supervisor};
use std::time::Duration;

let outcome = Supervisor::new(Command::new("my-service"))
    .restart(RestartPolicy::OnCrash)
    .max_restarts(5)
    .backoff(Duration::from_millis(200), 2.0)
    .storm_pause(Duration::from_secs(15)) // crash-loop guard (off by default)
    .run()
    .await?;
println!(
    "stopped after {} restarts ({} storm pauses): {:?}",
    outcome.restarts, outcome.storm_pauses, outcome.stopped
);
}

Exponential backoff with jitter by default; stop_when(…) ends supervision on a condition; .with_runner(&group) keeps every incarnation inside one shared kill-on-drop group. storm_pause arms the failure-storm guard: failures feed a decaying score, and past the threshold the supervisor takes one collective pause instead of hammering restarts — "fails rarely" and "crash-looping" stop being the same case.

Fine print: Supervision, failure storms.

Retry a flaky command

#![allow(unused)]
fn main() {
use processkit::Error;
use std::time::Duration;

let fetched = Command::new("git")
    .args(["fetch", "--quiet"])
    .timeout(Duration::from_secs(10))
    .retry(3, Duration::from_millis(200), |e| {
        matches!(e, Error::Timeout { .. })
            || e.diagnostic().is_some_and(|m| m.contains("Could not resolve host"))
    })
    .run()
    .await?;
}

The classifier sees the typed error and decides whether this failure is worth another attempt; each attempt is a fresh process. retry replays a run to success — for keeping a process alive, use a Supervisor (previous recipe).

Fine print: Timeouts, retries & cancellation → retry.

Cancel runs on shutdown

#![allow(unused)]
fn main() {
use processkit::CancellationToken;

let token = CancellationToken::new();

let job = tokio::spawn({
    let token = token.child_token();
    async move { Command::new("long-job").cancel_on(token).run().await }
});

// On Ctrl-C / shutdown signal / sibling failure:
token.cancel(); // kills the tree; the run resolves to Error::Cancelled
let outcome = job.await; // Err(Error::Cancelled { .. }) inside
}

Cancellation is always an error (the run was abandoned, there is no result), beats a simultaneous timeout, and is terminal for retry and Supervisor alike.

For a typed wrapper whose commands never cross your code, set the token once on the client — every command it builds carries it:

#![allow(unused)]
fn main() {
use processkit::{CancellationToken, CliClient};

let token = CancellationToken::new();
let gh = CliClient::new("gh").default_cancel_on(token.child_token());
// token.cancel() → every in-flight command of THIS client dies.
}

Fine print: Timeouts, retries & cancellation → cancellation, client-level default.

Measure what a run cost

#![allow(unused)]
fn main() {
use std::time::Duration;

// One run, summarized (requires the opt-in `stats` feature):
let profile = Command::new("crunch").start().await?.profile(Duration::from_millis(100)).await?;
println!("outcome={:?} took={:?} peak_rss={:?} avg_cpu_cores={:?}",
    profile.outcome, profile.duration, profile.peak_memory_bytes, profile.avg_cpu_cores());
}

For a live series over a whole group, group.sample_stats(every) yields a Stream of snapshots. CPU/memory need a real container (Windows Job / Linux cgroup); elsewhere you still get process counts.

Fine print: Process groups → stats · Streaming → profiling.

Contain a process you didn't spawn

#![allow(unused)]
fn main() {
use processkit::{Error, ProcessGroup};

// `tokio`/`std` calls return `io::Error`, which the crate does NOT auto-convert
// into `processkit::Error` (there is no blanket `From<io::Error>`, by design) —
// map it explicitly, or use a `Box<dyn std::error::Error>` / `anyhow` return in
// your own code so both error types `?` freely.
let child = tokio::process::Command::new("legacy-launcher")
    .spawn()
    .map_err(Error::Io)?;

let group = ProcessGroup::new()?; // `adopt` is part of `process-control` (default-on)
group.adopt(&child)?;            // from now on the group's teardown covers it
}

Adoption is best-effort by mechanism — on Windows/cgroup the whole running tree joins; on the POSIX process-group backends an exec'd child is contained individually (its future forks too, where it could be re-grouped). The guide spells out exactly what each mechanism can promise.

Fine print: Process groups → adopt · Platform support.

Test code that runs processes — without processes

#![allow(unused)]
fn main() {
use processkit::testing::{Reply, ScriptedRunner};

// Your code takes any `R: ProcessRunner`; in tests, hand it a script.
// Rules match on a prefix of the *program name followed by its arguments*
// (the first element is the program):
let runner = ScriptedRunner::new()
    .on(["git", "rev-parse"], Reply::ok("abc123\n"))
    .on(["git", "push"], Reply::fail(128, "remote: permission denied"))
    .fallback(Reply::ok(""));

// my_deploy(&runner).await? — no subprocess, fully deterministic.
}

RecordingRunner wraps any runner and captures every Invocation for assertions; MockRunner (feature mock) gives mockall expectations; DryRunRunner records the commands it would have run without executing anything, for the simplest "did we call the right thing" checks; and the record feature's RecordReplayRunner records real runs into a JSON cassette once and replays them hermetically in CI.

Fine print: Testing your code.

Test streaming code — without processes

#![allow(unused)]
fn main() {
use processkit::{Command, Outcome, ProcessRunner, Finished};
use processkit::testing::{Reply, ScriptedRunner};
use std::time::Duration;

let runner = ScriptedRunner::new()
    .on(["gh", "run", "watch"], Reply::lines(["queued", "in_progress", "completed"])
        .with_line_delay(Duration::from_millis(50))); // paced delivery

let mut run = runner.start(&Command::new("gh").args(["run", "watch", "123"])).await?;
run.wait_for_line(|l| l.contains("completed"), Duration::from_secs(5)).await?;
let Finished { outcome, .. } = run.finish().await?;
assert_eq!(outcome, Outcome::Exited(0));
}

A scripted start() feeds the canned lines through the same pump machinery a real child uses, so stdout_lines, the readiness probes, and finish behave identically — and with_line_delay is deterministic under #[tokio::test(start_paused = true)]. Canned output also replays through on_stdout_line/on_stderr_line handlers on the bulk verbs, so progress-reporting paths test hermetically too.

Fine print: Testing → scripted streaming.

Wrap a CLI tool behind a typed API

#![allow(unused)]
fn main() {
use processkit::{cli_client, ProcessRunner, Result};

cli_client!(pub struct Git => "git");

impl<R: ProcessRunner> Git<R> {
    pub async fn current_branch(&self) -> Result<String> {
        // A verb takes the args directly (D7); pass a built `command(..)` only
        // when you need to customize it (per-call timeout, stdin, …).
        self.core.run(["branch", "--show-current"]).await
    }
    pub async fn is_clean(&self) -> Result<bool> {
        self.core.probe(["diff", "--quiet"]).await
    }
}
}

The generated struct carries a runner and per-client defaults (default_timeout, default_env); your methods are just argument lists and parsers — and because the runner is injectable, the whole wrapper is testable with the previous recipe's ScriptedRunner.

Fine print: Testing your code → CliClient.

Running commands

‹ docs index

Command is the entry point of the runner layer: a builder describing what to run and how, plus a family of consuming verbs that decide what you get back. Every one-shot verb spawns the child into a fresh, private kill-on-drop process group, so an early return, panic, or dropped future can never leak a process tree.

Program, arguments, working directory

#![allow(unused)]
fn main() {
use processkit::Command;

let out = Command::new("git")
    .arg("log")                          // one at a time…
    .args(["--oneline", "-n", "10"])     // …or in bulk
    .current_dir("/path/to/repo")        // run there
    .run()
    .await?;
}

Arguments are passed as an array — there is no shell between you and the child, so there is no quoting, no word-splitting, and no injection surface. (When you actually want a | b | c, use a pipeline, which connects the stages in-process instead of invoking a shell.)

The program name reaches the OS verbatim — two deliberate non-goals (conveniences some libraries layer on, e.g. duct): a bare name is resolved on PATH by the OS, never rewritten to ./name; and current_dir does not re-anchor a relative program path against the new directory — whether Command::new("./tool").current_dir(dir) resolves tool relative to dir is the platform's behavior (Unix: yes; Windows: the parent's directory may win). Pass absolute program paths when combining the two.

prefer_local(dir) steers bare-name resolution without touching any of that verbatim contract: it probes a directory before the system PATH for this one run — a project's node_modules/.bin, a target/debug, a vendored toolchain — so Command::new("eslint").prefer_local("node_modules/.bin") picks up the local binary ahead of any global one. Repeated calls accumulate in priority order, all ahead of the PATH fallback, and the lookup reuses the same PATHEXT-aware machinery, so a .exe/.cmd/.bat on Windows resolves exactly as it would on PATH. It affects only a bare name (a path-form program is untouched) and never rewrites the child's own PATH; when resolution still fails, these directories appear in Error::NotFound's searched alongside the PATH entries.

For quick one-liners the free functions skip the builder:

#![allow(unused)]
fn main() {
let version = processkit::run("cargo", ["--version"]).await?;       // trimmed stdout, success required
let result  = processkit::output_string("git", ["status", "-s"]).await?;   // full ProcessResult
}

Environment

Four builders compose, applied in a fixed order at spawn:

#![allow(unused)]
fn main() {
use processkit::Command;

Command::new("worker")
    .env("RUST_LOG", "debug")        // set one variable
    .env_remove("GIT_DIR")           // unset one inherited variable
    .run().await?;

// Allow-list mode: clear everything, copy only the named parent variables.
Command::new("sandboxed-tool")
    .inherit_env(["PATH", "HOME", "LANG"])
    .env("MODE", "ci")               // explicit env/env_remove still apply on top
    .run().await?;

// Scorched earth: the child starts with an empty environment.
Command::new("hermetic-tool").env_clear().run().await?;
}

inherit_env is the sandboxing middle ground: it implies env_clear, then copies the listed variables from the parent at each spawn (so a retry sees fresh values), and repeated calls accumulate names. A name the parent doesn't have is skipped, not set to empty.

Standard input

By default stdin is closed at spawn — the child reads EOF immediately and can never hang waiting for input. Everything else is opt-in via stdin(Stdin::…):

SourceReusable on re-run?Use for
Stdin::empty()The default, explicit
Stdin::from_string("…")Text payloads
Stdin::from_bytes(vec![…])Binary payloads
Stdin::from_iter_lines(["a", "b"])Anything iterable; each item is written \n-terminated
Stdin::from_file(path)✅ (re-opened per run)Large inputs streamed from disk
Stdin::from_reader(reader)❌ one-shotAny AsyncRead — a socket, a decompressor, …
Stdin::from_lines(stream)❌ one-shotAny Stream<Item = String> — a channel, a tail, …
#![allow(unused)]
fn main() {
use processkit::{Command, Stdin};

let sorted = Command::new("sort")
    .stdin(Stdin::from_iter_lines(["banana", "apple", "cherry"]))
    .run()
    .await?;
assert_eq!(sorted, "apple\nbanana\ncherry");
}

The payload is written on a background task (so a large input can't deadlock against the child's output) and the pipe is dropped afterwards to signal EOF. The two one-shot sources are consumed by their first run: a retried or cloned command reusing them fails loud the second time — re-running a consumed from_reader/from_lines source is an Error::Io (InvalidInput) at launch (D10), not a silent empty stdin. Prefer the reusable sources when a command may run more than once.

For conversational, request/response stdin — write a line, read the answer, repeat — use keep_stdin_open() and the streaming API instead: see Streaming & interactive I/O.

Output handling

Encodings

Output is decoded line by line, UTF-8 by default (invalid bytes become U+FFFD, never an error). Legacy-encoding tools can override per stream:

#![allow(unused)]
fn main() {
use processkit::Command;

let out = Command::new("legacy-tool")
    .encoding(encoding_rs::SHIFT_JIS)          // both streams…
    // .stdout_encoding(…) / .stderr_encoding(…) // …or each its own
    .output_string()
    .await?;
}

(processkit::prelude::Encoding re-exports encoding_rs::Encoding, so any of its encodings works — the single-byte and ASCII-compatible multibyte ones (WINDOWS_1252, GBK, SHIFT_JIS, …) and the non-ASCII-compatible ones (UTF_16LE/UTF_16BE): output is fed through one persistent decoder and split on decoded newlines, so a 0x0A byte inside a UTF-16 code unit is not mistaken for a line break. A leading byte-order mark of the chosen encoding is stripped once at the stream start.)

Line terminators

Output is split into lines on \n by default. Carriage-return progress output — curl/pip/apt redrawing \rProgress: 50%\rProgress: 100% over one visual line — otherwise accumulates as a single ever-growing unstreamed line (and can be dropped whole under a byte cap). Switch the boundary to \r per stream:

#![allow(unused)]
fn main() {
use processkit::{Command, LineTerminator};

let out = Command::new("pip")
    .args(["install", "big-package"])
    .line_terminator(LineTerminator::CarriageReturn)   // both streams…
    // .stdout_line_terminator(…) / .stderr_line_terminator(…) // …or each its own
    .on_stderr_line(|line| eprintln!("[pip] {line}"))
    .output_string()
    .await?;
}

CarriageReturn treats a bare \r as the boundary, while a \r\n pair still counts as one terminator (whole or split across reads). The choice threads through every sink that consumes "a line" — the streaming verbs, the stdout_tee/stderr_tee sinks, the on_stdout_line/on_stderr_line handlers, and output_string. The default (Newline, \n-only) is unchanged.

Buffer policies — bounding memory on chatty children

Captured lines are held in memory; a multi-gigabyte log would normally grow the buffer to match. output_buffer bounds retention (the pipe is always fully drained, so the child never blocks):

#![allow(unused)]
fn main() {
use processkit::{Command, OutputBufferPolicy, OverflowMode};

let tail = Command::new("verbose-build")
    .output_buffer(OutputBufferPolicy::bounded(1_000)) // keep the newest 1000 lines
    .output_string()
    .await?;

// …or keep the head instead of the tail:
let head_policy = OutputBufferPolicy::bounded(1_000).with_overflow(OverflowMode::DropNewest);
}

DropOldest (the default) keeps a rolling tail; DropNewest freezes the head. bounded(0) retains nothing — useful when a line handler (below) is the real consumer. Under a line cap, dropped or not, every line still feeds the handlers and the line counters.

The line cap alone does not bound memory — one enormous newline-free "line" (base64 -w0) is held whole. Add with_max_bytes to cap the retained bytes too (either ceiling, or both); the byte cap also bounds the pump's in-flight assembly buffer, so a never-terminated flood can't exhaust memory. One consequence: a line whose own length exceeds the byte cap can't be assembled, so it is dropped whole — counted, but not delivered to a per-line handler or stdout_tee (don't set a byte cap if a tee must see arbitrarily long lines):

#![allow(unused)]
fn main() {
use processkit::{Command, OutputBufferPolicy};
let policy = OutputBufferPolicy::unbounded().with_max_bytes(8 << 20); // 8 MiB ring
let strict = OutputBufferPolicy::fail_loud(10_000).with_max_bytes(8 << 20); // error on either
}

fail_loud makes the ceiling error instead of dropping: the run fails with Error::OutputTooLarge once the cumulative output (lines or bytes) crosses the cap — even when a streaming consumer is draining lines as they arrive. It bounds memory, not wall-time, so pair it with timeout against a flooding child.

The byte ceiling now reaches output_bytes' raw stdout capture too, not just the line-pumped streams: under OverflowMode::Error a raw capture past max_bytes fails with Error::OutputTooLarge (reported with max_lines: None — raw bytes carry no line count), and the drop modes bound the retained bytes to a head/tail and set ProcessResult::truncated. Without a byte cap output_bytes stays unbounded, exactly as before.

Even under a drop policy (DropOldest/DropNewest), the checking verbs that hand back stdout as if complete — run, parse, try_parserefuse silently-truncated output (B12): if the policy dropped lines they fail with Error::OutputTooLarge rather than feed a parser a truncated tail. The lenient capture verbs (output_string / output_bytes) are unaffected — they return the partial result with truncated() set for you to inspect.

Line handlers — tee output as it arrives

on_stdout_line / on_stderr_line run a callback on each decoded line in addition to capture or streaming — logging, progress bars, metrics:

#![allow(unused)]
fn main() {
use processkit::Command;

let result = Command::new("cargo")
    .args(["build", "--release"])
    .on_stderr_line(|line| eprintln!("[build] {line}"))
    .output_string()
    .await?;
}

The handler runs on the read pump — keep it cheap. The contract is forgiving and precisely specified:

  • A panicking handler does not poison the run. The panic is caught, the handler is disabled for the rest of the run (surfaced as a tracing warn when that feature is on), and pumping continues — the final result still carries every line. You can safely re-export this callback seam to your own users without auditing their closures.
  • Ordering: invocations are FIFO within a stream; there is no ordering between stdout and stderr handlers (two independent pumps). On the consuming verbs, all handler calls happen-before the awaited future resolves — finalize a progress bar the moment the call returns. (One documented exception: a leaked pipe held open past the child's death is cut off after a bounded teardown grace.)
  • Handlers are hermetically testable: ScriptedRunner replays canned output through them — see Testing → scripting replies.

For a ready-made tee to an async sink — a file, socket, or any [tokio::io::AsyncWrite] — reach for stdout_tee / stderr_tee instead of hand-writing a handler. Each decoded line is written to the sink (plus a \n) as it is produced, awaited on the pump so a slow sink applies backpressure (the pump slows, the pipe fills, the child blocks) rather than blocking the runtime; a write error disables the tee with a tracing warn instead of being swallowed. It runs independently of on_stdout_line — set both and both fire per line.

Timeouts and retries

#![allow(unused)]
fn main() {
use processkit::{Command, Error};
use std::time::Duration;

let out = Command::new("flaky-network-tool")
    .timeout(Duration::from_secs(30))                 // kill the tree at the deadline
    .retry(3, Duration::from_millis(200), |e| {       // up to 3 attempts total
        matches!(e, Error::Timeout { .. })            // …but only retry timeouts
    })
    .run()
    .await?;
}
  • timeout kills the whole process tree at the deadline. On the capturing verbs the expiry is captured (ProcessResult::timed_out), on the success-checking verbs it raises Error::Timeout — the full decision table lives in Timeouts, retries & cancellation.
  • retry applies to the success-checking verbs only — run, run_unit, exit_code, probe, checked, parse, and try_parse (seven in all; each runs through the retry loop). The classifier sees the typed error and decides. The non-erroring output_string/output_bytes paths never retry, and neither does first_line (its stream search is single-attempt).
  • Config-driven variants. timeout_opt(Option<Duration>) folds the Some(d) => timeout(d) / None => no_timeout() match into one call for a config-sourced deadline, and retry_never() is an explicit per-command opt-out of a client default_retry (symmetric with no_timeout()). Both are covered in full on Timeouts, retries & cancellation.

Privileges and spawn flags

Spawn-time controls for sandboxing and service launch:

#![allow(unused)]
fn main() {
use processkit::Command;

// Unix: drop privileges (uid + gid + supplementary groups) and detach.
Command::new("worker")
    .gid(1000)            // applied before uid (a gid change needs privilege)
    .groups([1000])       // replace the inherited (often root's) supplementary groups
    .uid(1000)            // dropped last
    .setsid()             // new session: survives the controlling terminal
    .run().await?;

// Windows: no console window flashing up from a GUI app.
Command::new("helper").create_no_window().run().await?;

// Hardening: take the direct child down even if THIS process is SIGKILLed
// (Drop never runs). Windows has this for free; Linux arms PDEATHSIG.
Command::new("worker").kill_on_parent_death().start().await?;
}

uid / gid / groups / setsid are POSIX-only — on Windows the run fails with Error::Unsupported rather than silently skipping a privilege drop. A correct drop sets all three of uid/gid/groups: dropping the uid alone leaves the child holding the parent's (often root's) supplementary groups. create_no_window is a harmless no-op outside Windows. kill_on_parent_death is best-effort by design: guaranteed on Windows (regardless of the knob), direct-child-only on Linux, unavailable on macOS/BSD — the graceful-exit guarantee via Drop holds everywhere either way. Containment is preserved in every combination; the platform fine print (the Linux cgroup × uid interaction, setsid × process-group coordination, the pdeathsig thread caveat) is collected in Platform support.

Scheduling priority and umask. priority(Priority) launches the child at a lower — or higher — CPU-scheduling priority, for background/batch work that shouldn't starve the foreground (or a task that should win over it); umask sets its file-mode creation mask:

#![allow(unused)]
fn main() {
use processkit::{Command, Priority};

Command::new("ffmpeg")
    .args(["-i", "in.mov", "out.mp4"])
    .priority(Priority::BelowNormal)   // Idle · BelowNormal · Normal · AboveNormal · High
    .umask(0o077)                      // Unix: files this child creates are owner-only
    .run().await?;
}

Idle/BelowNormal/Normal/AboveNormal/High map to nice/setpriority on Unix and a priority-class flag on Windows. Unlike the privilege builders, priority never yields Error::Unsupported — every variant is covered on both platforms. The Unix catch is privilege, not support: lowering nice below the value the child inherits needs CAP_SYS_NICE/root, or the spawn fails as Error::Spawn. That bites both High and AboveNormal (nice(-5) is already a decrease), and — under a positively-niced parent (a niced CI/batch launcher) — even Priority::Normal: nice(0) there is itself a privileged decrease from the inherited value, not the no-op it is from an un-niced parent. Raising nice (Idle/BelowNormal) is always allowed.

umask(u32) sets the child's file-mode creation mask (umask(2)), governing the default permissions of the files it creates. Like uid/gid/setsid it is Unix-only (applied via pre_exec); on other targets the run fails with Error::Unsupported rather than silently ignoring the requested mask.

Interactive auth / TTY. processkit wires pipes, not a pseudo-terminal, so a tool that demands a tty — an ssh/sudo password prompt, some credential helpers — won't get one (PTY support is not implemented; the trade-off is recorded in decisions/permissions-privileges-pty-network.md). Drive such tools non-interactively instead: key-based auth, ssh -o BatchMode=yes, GIT_SSH_COMMAND / GIT_TERMINAL_PROMPT=0, or feed a known answer over interactive stdin. Conversational tools that read stdin without needing a tty already work today via keep_stdin_open + stdout_lines.

Consuming verbs

VerbReturnsNon-zero exitTimeoutUse when
output_string()ProcessResult<String>capturedcaptured (timed_out)You want to inspect the outcome yourself
output_bytes()ProcessResult<Vec<u8>>capturedcapturedBinary stdout (images, archives, …)
run()trimmed stdout StringError::ExitError::Timeout"Give me the answer or fail"
exit_code()i32the code, OkError::TimeoutThe code is the answer
probe()bool0true, 1false, else Error::ExitError::TimeoutPredicate commands: git diff --quiet, grep -q
first_line(pred)Option<String>— (stream-based)Error::TimeoutGrab one matching line, kill the rest
start()live RunningProcessbounds the streamStreaming, interactive I/O, probes
#![allow(unused)]
fn main() {
use processkit::Command;

// probe(): the exit code as a boolean.
let clean = Command::new("git").args(["diff", "--quiet"]).probe().await?;

// first_line(): stop as soon as the interesting line appears.
let first_match = Command::new("git")
    .args(["log", "--oneline"])
    .first_line(|l| l.contains("fix:"))
    .await?;
}

first_line returns Ok(None) when stdout closes without a match, and kills the (private-group) child once it has its answer — you never wait out a long log for one line. A cancel_on token that fires while the search is still running surfaces as Error::Cancelled, so a readiness probe with a shutdown token can't misread token-driven teardown as "the line never appeared" — while a run that genuinely ends with no match still reports Ok(None), even if the token happens to fire an instant later.

Results and errors

The capturing verbs hand back a ProcessResult:

#![allow(unused)]
fn main() {
use processkit::Command;

let result = Command::new("git").args(["merge", "feature"]).output_string().await?;

result.code();         // Option<i32> — None = killed (timeout/signal), no code
result.signal();       // Option<i32> — the signal number (Unix), else None
result.is_success();   // code in ok_codes (default {0})
result.timed_out();    // the run's own deadline expired
result.outcome();      // the explicit three-way enum behind the accessors above
result.stdout();       // &str (or &[u8] from output_bytes)
result.stderr();       // &str
result.combined();     // stdout + stderr concatenated
result.diagnostic();   // stderr if non-empty, else stdout — the human-facing line
                       // (git/jj put "CONFLICT …" on stdout!)

// Opt into erroring whenever you're ready:
let ok = result.ensure_success()?; // Exit / Timeout / Signalled (signal-kill) as typed errors
}

When the three-way distinction matters, match on Outcome instead of mentally decoding the code()/timed_out() pair:

#![allow(unused)]
fn main() {
use processkit::Outcome;

match result.outcome() {
    Outcome::Exited(0) => println!("clean"),
    Outcome::Exited(code) => println!("failed with {code}"),
    Outcome::Signalled(signal) => println!("killed by signal {signal:?}"),
    Outcome::TimedOut => println!("hit its deadline"),
    _ => {} // non_exhaustive: future dispositions
}
}

For a single query you usually don't need the match (and its #[non_exhaustive] wildcard): Outcome carries the same code() / signal() / timed_out() accessors as ProcessResult, so a bare Outcome (from RunningProcess::wait or Finished::outcome) answers directly — outcome.code(), outcome.signal(), outcome.timed_out(). There is no Outcome::is_success (success is ok_codes-aware — use ProcessResult::is_success).

The error enum is structured and #[non_exhaustive]:

VariantMeaning
Error::Spawn { program, source }The program was located but the OS couldn't start it (permissions, a bad working directory, a Windows .cmd/.bat needing cmd.exe, …) — not is_not_found()
Error::NotFound { program, searched }The program couldn't be located (the single "not found" representation — is_not_found() is true); searched is Some(dirs) for a bare-name PATH lookup, None otherwise
Error::Exit { program, code, stdout, stderr, stdout_bytes }Non-zero exit, both streams attached in full (the Display message is bounded, but the fields carry the complete captured text for classification); stdout_bytes: Option<Vec<u8>> holds the raw stdout for an error built over output_bytes (None on the text path) — reach it via stdout_bytes()
Error::Signalled { program, signal, stdout, stderr, stdout_bytes }The process was killed by a signal (no exit code); signal carries the number on Unix, None elsewhere; the partial streams captured before the kill are attached (reach them via diagnostic()), with the raw stdout_bytes when built over output_bytes
Error::OutputTooLarge { program, max_lines, max_bytes, total_lines, total_bytes }A fail_loud buffer's line or byte ceiling was exceeded; max_lines/max_bytes are each Option<usize> — the configured cap that tripped, None on an axis with no cap (e.g. max_lines: None from a raw output_bytes capture) — mirroring the OutputBufferPolicy knobs, while total_lines/total_bytes report what was seen
Error::Timeout { program, timeout, stdout, stderr, stdout_bytes }The run's own deadline killed it; whatever the run captured before the kill is attached — a hung tool's last stderr line tails the Display and is reachable via diagnostic(), with the raw stdout_bytes when built over output_bytes
Error::NotReady { program, timeout }A readiness probe gave up
Error::Parse { program, message }A try_parse parser (on Command, ProcessRunnerExt, CliClient, or Pipeline) rejected the output (the Display/Debug of message is bounded to a 200-byte preview; the field carries the full text)
Error::Stdin { program, source }Feeding the child's stdin failed for a non-broken-pipe reason on an otherwise-successful run (a louder failure — exit/signal/timeout — wins instead); a routine broken pipe never surfaces
Error::CassetteMiss { program }(record feature) a cassette replay found no matching recording (stale/incomplete cassette) — kept distinct from a missing program, so is_not_found() is false
Error::Unsupported { operation }The platform can't do what was asked (and silently skipping would be wrong)
Error::Cancelled { program }the run's token was cancelled
Error::ResourceLimit { kind, reason, detail }(limits feature) a requested resource cap couldn't be enforced — kind: LimitKind (Memory/Processes/Cpu) is which limit, reason: LimitReason (Invalid/Unsupported/Unenforceable) is why, detail: String is the human-facing explanation; classify without parsing text via limit_kind()/limit_reason()
Error::Io(source)A low-level IO error from the crate's own machinery (driving a child, group control, cassette files) — never an arbitrary foreign io::Error (no blanket From, D13)

Beyond the enum itself, each data-carrying variant — Exit, Timeout, Signalled, Spawn, NotFound, Parse, OutputTooLarge, Stdin, and ResourceLimit — is individually #[non_exhaustive]: outside the crate you can neither build one with a struct literal nor destructure it field-exhaustively (a future release may add fields). Match with a trailing .. and read through the accessors instead — program(), stdout()/stderr()/combined(), code(), signal(), the is_*() predicates, stdout_bytes() for the raw stdout of a checking verb over output_bytes, and limit_kind()/limit_reason() for a ResourceLimit. Going the other way, building an Error::Parse from your own output parser (outside the crate's try_parse helpers) is a supported path: the constructor Error::parse(program, message) is on the public surface.

Error::diagnostic() returns the most useful human-facing line out of a failure that captured output — Exit, and (D12) Timeout / Signalled (the partial streams of a hung-then-killed or crashed tool). Each of those variants' one-line Display also appends a bounded excerpt of that diagnostic (the last non-empty line, capped at 200 bytes), so a bare eprintln!("{e}") reads `git` exited with code 2: fatal: boom — actionable in a log line without dumping multi-KiB streams into it.


Next: Streaming & interactive I/O · Timeouts, retries & cancellation · Process groups

Process groups

‹ docs index

A ProcessGroup ties the lifetime of a whole child-process tree to a Rust value: every process spawned into the group — and everything those processes spawn — is killed when the group is dropped. An exiting, panicking, or ?-returning owner never leaks subprocesses; the kernel object enforcing this (Job Object / cgroup / POSIX process group) catches even grandchildren you never knew about. (Killing grandchildren is the problem duct.py's gotchas list files under "currently unsolved" for pipe-based designs — kernel containment is the solution, and the reason this crate exists.)

Creating a group

#![allow(unused)]
fn main() {
use processkit::{ProcessGroup, ProcessGroupOptions};
use std::time::Duration;

// Defaults: 2s graceful-shutdown grace, escalate to SIGKILL.
let group = ProcessGroup::new()?;

// Tuned:
let group = ProcessGroup::with_options(
    ProcessGroupOptions::default()
        .shutdown_timeout(Duration::from_secs(10))
        .escalate_to_kill(true),
)?;

// Which kernel mechanism is actually containing the tree?
println!("{:?}", group.mechanism()); // JobObject | CgroupV2 | ProcessGroup
}

mechanism() reports what you actually got: CgroupV2 quietly falls back to ProcessGroup on Linux hosts without cgroup delegation (see Platform support).

You rarely create a group explicitly for one-shot runs: every Command::run()-style call makes a private group automatically. Reach for an explicit group when several children should share one fate, or when you need the group verbs below.

Putting processes in

Three doors, in order of preference:

#![allow(unused)]
fn main() {
use processkit::{Command, ProcessGroup};

let group = ProcessGroup::new()?;

// 1. start(): the full Command experience (capture, streaming, timeouts) in a
//    SHARED group. The handle does not own the group — dropping the handle
//    kills that child, dropping the group kills everyone.
let server = group.start(&Command::new("dev-server")).await?;

// 2. spawn(): the raw escape hatch for a tokio::process::Command you already
//    have. You get the bare Child back; pipes and reaping are your problem.
//    spawn() takes the command BY VALUE (reuse would stack pre-exec hooks).
let raw = tokio::process::Command::new("background-helper");
let child = group.spawn(raw)?;

// 3. adopt(): contain a child that was spawned OUTSIDE the group.
let external = tokio::process::Command::new("legacy-launcher").spawn()?;
group.adopt(&external)?;
}

adopt moves only the named process: descendants it already has keep their old containment (future forks are captured — on Windows/cgroup). A few sharp edges worth knowing:

  • A child that already exited but has not been reaped (no wait() yet — a zombie whose pid/handle is still valid) is a successful no-op: there is nothing left to contain, so adopt returns Ok on the containment backends.
  • A child that already exited and was reaped (wait()ed) has no pid left — adopt returns an error rather than silently tracking nothing.
  • On the POSIX process-group mechanism, a child that has already exec'd can't be re-grouped (POSIX forbids it), so it is tracked individually: the child itself is signalled/killed with the group, but its future forks are not. The caller keeps the Child handle and is responsible for reaping.

Tearing down: drop, terminate, shutdown

VerbWhat happensWhen
drop(group)Immediate hard kill of the whole tree (kill-on-close)The safety net — always on
group.kill_all()The same hard kill, group stays usable (cgroup-kill / Job Object / process-group backends). On a pre-5.14 Linux kernel lacking cgroup.kill, the per-pid SIGKILL fallback returns Err if the tree doesn't drain (a fork bomb still out-spawning, or D-state zombies)Explicit teardown mid-flight; idempotent
group.shutdown().awaitUnix: SIGTERM → wait shutdown_timeoutSIGKILL survivors (if escalate_to_kill); Windows: atomic job kill when escalate_to_kill, else the survivors are spared (handle closed without kill-on-close). Consumes the group (shutdown_ref(&self) is the same teardown, borrowing — for a group held behind an Arc/supervisor)Graceful service stop
#![allow(unused)]
fn main() {
use processkit::{Command, ProcessGroup, ProcessGroupOptions};
use std::time::Duration;

let group = ProcessGroup::with_options(
    ProcessGroupOptions::default()
        .shutdown_timeout(Duration::from_secs(5))
        .escalate_to_kill(true),
)?;
let _service = group.start(&Command::new("my-service")).await?;

// SIGTERM, give it 5s to flush and exit, SIGKILL stragglers:
group.shutdown().await?;
}

A child that handles SIGTERM ends the grace earlyshutdown returns as soon as the tree is empty, not after the full timeout. One subtlety: the liveness probe sees an exited-but-unreaped child (a zombie) as alive on the process-group backends, so keep wait()ing your handles concurrently if you want the early return. Drop can't await, which is why the graceful tier lives in this async method — dropping without calling it performs only the hard kill.

Signalling the whole tree

signal/suspend/resume/members/adopt — this section and the two below — require the default-on process-control feature. The teardown verbs above are core and always present.

#![allow(unused)]
fn main() {
use processkit::{Command, ProcessGroup, Signal};

let group = ProcessGroup::new()?;
let _server = group.start(&Command::new("my-server")).await?;

group.signal(Signal::Hup)?;        // "reload your configuration"
group.signal(Signal::Usr1)?;       // whatever the tool defines
group.signal(Signal::Other(34))?;  // raw signal number escape hatch
}
PlatformDeliverable signals
Linux (cgroup or pgroup), macOS/BSDAny — Term, Kill, Int, Hup, Quit, Usr1, Usr2, Other(n)
WindowsKill only (maps to the Job Object terminate); anything else → Error::Unsupported

Signal::Kill always takes the same atomic whole-tree kill path as kill_all (cgroup.kill / killpg / job terminate), so it cannot miss a process forked mid-broadcast. Other signals are a per-member broadcast — best-effort against a tree that is forking at that exact moment. An empty group accepts any deliverable signal trivially. On the cgroup mechanism a real per-member delivery failure (e.g. EPERM from a member that changed uid, or a seccomp/container restriction) is surfaced as an Err rather than swallowed — an ESRCH race (the member already exited) is still success; the pgroup (macOS/BSD, Linux-without-cgroup) backend remains purely best-effort.

Suspending and resuming

Freeze a tree (to snapshot it, to starve a runaway while you investigate, to pause background work), then thaw it:

#![allow(unused)]
fn main() {
use processkit::{Command, ProcessGroup};

let group = ProcessGroup::new()?;
let _cruncher = group.start(&Command::new("cpu-hog")).await?;

group.suspend()?;   // the whole tree stops consuming CPU
// … inspect, snapshot, wait for the user …
group.resume()?;
}

Per-platform machinery — and its visible differences:

PlatformMechanismNotes
Linux cgroupone cgroup.freeze writeAtomic over the subtree; freeze is group state
Linux pgroup, macOS/BSDSIGSTOP / SIGCONT broadcastIdempotent (level-triggered)
Windowsper-thread SuspendThread walkCounted: N suspends need N resumes; best-effort against mid-walk thread churn

Two caveats that bite in practice:

  • Spawning into a suspended group diverges. Under the cgroup mechanism a child spawned or adopted while the group is frozen starts frozen — and start() may never return until resume (the forked child joins the cgroup before exec, so it can freeze before completing the spawn handshake). Windows and the pgroup backends freeze only members present at the call. Rule of thumb: resume before starting new work.
  • A suspended tree can still be hard-killed (drop / kill_all / Signal::Kill all act on frozen processes), but a graceful shutdown starts with a SIGTERM the frozen tree can't act on — it would wait out the whole grace. Resume first for a clean shutdown.

Listing members

#![allow(unused)]
fn main() {
use processkit::{Command, ProcessGroup};

let group = ProcessGroup::new()?;
let _a = group.start(&Command::new("worker-a")).await?;
let _b = group.start(&Command::new("worker-b")).await?;

let pids: Vec<u32> = group.members()?;
println!("live members: {pids:?}");
}

What "members" means depends on the mechanism: Windows and Linux-cgroup list the whole tree (every descendant pid); the POSIX process-group backends list the tracked group leaders (one pid per started/adopted child) — their descendants are contained but not enumerated. An exited child still counts until it is reaped. The snapshot is point-in-time: a tree that is forking races it.

To wait on members rather than list them, race the handles with wait_any.

Resource limits

Requires the limits feature. Caps are a property of the group, set once at creation and enforced by the same kernel object that contains the tree:

#![allow(unused)]
fn main() {
use processkit::{Command, ProcessGroup, ProcessGroupOptions};

let group = ProcessGroup::with_options(
    ProcessGroupOptions::default()
        .max_memory(512 * 1024 * 1024) // bytes, whole tree
        .max_processes(64)             // fork-bomb ceiling
        .cpu_quota(0.5),               // half of one core
)?;
let _sandboxed = group.start(&Command::new("untrusted-tool")).await?;
}
CapabilityWindows Job ObjectLinux cgroup v2pgroup / macOS / BSD
Memory cap✅ whole-tree✅ whole-tree (memory.max)
Process-count cap✅ (pids.max)
CPU quota🟡 approximate (rate vs. total CPU)✅ (cpu.max)

cpu_quota is a fraction of a single core (2.0 = two cores). Limits need a real container; when a requested cap can't be enforced — no Job Object/cgroup, or a Linux cgroup whose controllers can't be enabled — with_options returns the typed Error::ResourceLimit { kind, reason, detail } instead of handing back a silently-unbounded group: kind says which limit (LimitKind::Memory/Processes/Cpu), and reason classifies why it was rejected via LimitReasonUnsupported (no containment mechanism at all on this platform), Unenforceable (the mechanism exists but the request was refused, e.g. missing cgroup delegation), or Invalid (a nonsensical value, like a zero quota). Match on the error or read Error::limit_kind() / Error::limit_reason() without destructuring.

On Linux this needs the process to run at the real cgroup-v2 root: the crate enables the controllers in this process's own cgroup, which cgroup v2's "no internal processes" rule allows only for the real hierarchy root — not a cgroup-namespace root (so an ordinary container fails too), not under systemd — and the crate doesn't migrate your process. See the limits prerequisites in Platform support. The uid()-drop interaction lives under its Caveats.

Stats and sampling

Requires the opt-in stats feature (features = ["stats"], or limits).

#![allow(unused)]
fn main() {
use processkit::{Command, ProcessGroup};
use processkit::prelude::StreamExt;
use std::time::Duration;

let group = ProcessGroup::new()?;
let _worker = group.start(&Command::new("worker")).await?;

// Point-in-time:
let snap = group.stats()?;
println!(
    "procs={} cpu={:?} peak_rss={:?}",
    snap.active_process_count, snap.total_cpu_time, snap.peak_memory_bytes,
);

// …or a series: first sample immediate, then every 250ms; missed ticks are
// skipped; the stream ends when the group can no longer report.
let mut samples = group.sample_stats(Duration::from_millis(250));
while let Some(s) = samples.next().await {
    println!("rss now: {:?}", s.peak_memory_bytes);
}
}

CPU time and peak memory are available where the kernel accounts for the whole tree (Windows, Linux cgroup); the process-group backends report the member count only — the Option fields stay None. The sampler borrows the group, so it can neither outlive it nor keep it (and the kill-on-drop guarantee) alive. For a single run's end-to-end summary, see profile.


Next: Streaming & interactive I/O · Platform support · Supervision

Streaming & interactive I/O

‹ docs index

The one-shot verbs in Running commands buffer the whole output. For long-running or conversational children, Command::start() returns a live RunningProcess you drive yourself: stream stdout as it arrives, write stdin incrementally, probe for readiness, race several children, or profile a run.

Lifecycle

#![allow(unused)]
fn main() {
use processkit::Command;

let mut run = Command::new("dev-server").start().await?;

run.pid();        // Option<u32> — None once the child is reaped
run.elapsed();    // time since spawn

// Consume the handle exactly one way:
//   output_string() / output_bytes()  → capture everything (same as the one-shot verbs)
//   wait()                            → just the Outcome; output is discarded
//   finish()                 → after streaming stdout (below)
//   profile(every)                    → resource samples; output discarded, like wait() (stats feature)
let outcome = run.wait().await?;   // Outcome: Exited(code) / Signalled(sig) / TimedOut
}

start() puts the child in a private group the handle owns: dropping the RunningProcess kills the whole tree, exactly like dropping a one-shot run's future. The shared-group variant — group.start(&cmd) — gives the same handle but the group controls the tree's fate (see Process groups).

There is also an explicit run.start_kill() for "stop it now, I'll wait() for the code myself".

Streaming stdout

stdout_lines() yields decoded lines as the child produces them — no waiting for exit, no full-output buffering. StreamExt (from processkit::prelude, re-exported from tokio-stream) provides .next():

use processkit::prelude::StreamExt;
use processkit::{Command, Outcome, Finished};

#[tokio::main]
async fn main() -> processkit::Result<()> {
    let mut run = Command::new("cargo")
        .args(["build", "--release"])
        .start()
        .await?;

    let mut lines = run.stdout_lines()?;
    while let Some(line) = lines.next().await {
        println!("build: {line}");
    }

    // The stream ended (stdout closed). Collect the outcome and stderr —
    // stderr was drained in the background the whole time, so a noisy child
    // could never block on a full pipe.
    let Finished { outcome, stderr, .. } = run.finish().await?;
    if outcome != Outcome::Exited(0) {
        eprintln!("build failed ({outcome:?}):\n{stderr}");
    }
    Ok(())
}

Things to know:

  • Call stdout_lines() once. It is fallible: a second stdout_lines / output_events call (stdout is consumed once), or a non-piped stdout (StdioMode::Inherit/Null), returns Err rather than a silently-empty stream.
  • The command's timeout bounds the stream on an own-group handle: at the deadline the tree is killed, the pipes close, and the stream ends — a streamed run can't hang past its deadline. A cancel_on token ends it the same way; the following finish then reports Error::Cancelled. Details in Timeouts & cancellation.
  • Line counters tick live: run.stdout_line_count() / stderr_line_count() are cheap progress gauges even while you stream.
  • The buffer policy and line handlers apply to streamed runs too — a handler sees each line on the pump, in addition to your loop.
  • The whole streaming surface is hermetically testable: a ScriptedRunner's start() returns a handle whose canned lines flow through the same pump machinery — stdout_lines, the readiness probes, and finish behave identically with no subprocess. See Testing → scripted streaming.

Interactive stdin

Conversational tools — write a request, read the response, repeat. Keep stdin open with keep_stdin_open(), take the writer with take_stdin():

use processkit::prelude::StreamExt;
use processkit::{Command, Outcome, Finished};

// `ProcessStdin`'s writer methods return `std::io::Result`; `Box<dyn Error>`
// mixes them with the crate's `Result` (or `.map_err(processkit::Error::Io)?`).
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // `bc` evaluates each stdin line and prints the result.
    let mut run = Command::new("bc").keep_stdin_open().start().await?;
    let mut stdin = run.take_stdin().expect("stdin was kept open");
    let mut answers = run.stdout_lines()?;

    stdin.write_line("2 + 2").await?;             // writes "2 + 2\n", flushed
    println!("= {}", answers.next().await.unwrap());

    stdin.write_line("6 * 7").await?;
    println!("= {}", answers.next().await.unwrap());

    stdin.finish().await?;                        // send EOF — bc exits
    let Finished { outcome, .. } = run.finish().await?;
    assert_eq!(outcome, Outcome::Exited(0));
    Ok(())
}

ProcessStdin offers write(&[u8]), write_line(&str) (newline + flush), flush(), and finish() (EOF). Dropping the writer — or the whole RunningProcess — closes stdin too; finish() just makes the EOF explicit and awaitable.

To drive a REPL with a control character, send_control(char) writes the single mapped byte for you — send_control('c') → Ctrl-C (\x03), send_control('d') → Ctrl-D (\x04) — instead of spelling out the raw escape with write. This writes one byte into a plain pipe, not a real terminal signal. A genuine SIGINT/SIGTSTP is something the OS delivers to a foreground process group attached to a TTY; there is no TTY here, so the child only "reacts" if it inspects its own stdin for the byte (as many REPLs and readline-based tools do). Actual signal delivery needs a pseudo-terminal, which this crate does not yet provide.

Avoid the full-duplex deadlock. A child's stdout pipe has a finite OS buffer; once it fills, the child blocks writing stdout until something reads it. If you push a large interactive stdin while nothing drains the child's stdout, the child stops reading stdin (blocked on stdout), your write parks waiting for stdin buffer space, and neither side progresses. The bc example above is safe because it interleaves one write with one read; when you both feed a sizable stdin and the child produces output, drain stdout_lines from one task while writing stdin from another. (The non-interactive Stdin::from_* sources are safe — the crate writes them on a background task that runs concurrently with the output pumps.)

For one-directional streamed input (a channel, a file tail) you don't need interactivity — give the command Stdin::from_lines(stream) / Stdin::from_reader(reader) and let the background writer feed it; see the stdin source table.

Readiness probes

"Start a server, then use it" needs ready, not merely started. Three probes replace the arbitrary sleep, each bounded by its own deadline:

#![allow(unused)]
fn main() {
use processkit::Command;
use std::time::Duration;

let mut run = Command::new("my-server").start().await?;

// 1. A line on stdout (returns the matching line):
let banner = run
    .wait_for_line(|l| l.contains("listening on"), Duration::from_secs(10))
    .await?;

// 2. A TCP port accepting connections:
run.wait_for_port("127.0.0.1:8080".parse().unwrap(), Duration::from_secs(10))
    .await?;

// 3. Any async predicate (an HTTP /health endpoint, a file appearing, …):
run.wait_for(|| async { health_check().await }, Duration::from_secs(10))
    .await?;

// ready — use the server…
}

Probe semantics, deliberately uniform:

  • A probe that can't pass within its deadline fails with Error::NotReady — distinct from Error::Timeout, which is the run's own deadline.
  • A probe also fails fast once readiness can no longer happen: the child exits, or (for wait_for_line) its stdout closes — no waiting out a 30s deadline on a dead server.
  • A failed probe never kills the child. You decide: retry, log and continue, or tear down.
  • All three probes background-drain the child's stdout/stderr while they poll — a child that writes a large startup burst before becoming ready can no longer stall in write() on a full OS pipe and spuriously fail the probe with Error::NotReady. wait_for_line consumes stdout up to (and including) the match — continue with finish or further streaming. wait_for_port / wait_for drain the same way but hand nothing back mid-probe; wait / output_string afterward still see the full output, but a fresh stdout_lines / output_events call doesn't compose with any of the three probes (same restriction as calling wait_for_line first).

Racing children with wait_any

The free function wait_any races several running processes and reports whichever exits first — the natural primitive for "restart whatever died" or "first answer wins":

#![allow(unused)]
fn main() {
use processkit::{Command, ProcessGroup, wait_any};

let group = ProcessGroup::new()?;
let mut a = group.start(&Command::new("replica-a")).await?;
let mut b = group.start(&Command::new("replica-b")).await?;

let (index, outcome) = wait_any(&mut [&mut a, &mut b]).await?;
println!("contender #{index} exited first with {outcome:?}");

// Only borrows: the loser is still usable.
let survivor = if index == 0 { &mut b } else { &mut a };
}

wait_any takes &mut borrows, applies no timeout of its own (wrap it in tokio::time::timeout to bound the race), and does no output pumping — drain chatty children first or give them bounded buffer policies.

Per-run telemetry

With the opt-in stats feature, a running child reports its own resource usage, and profile() turns a whole run into a summary:

#![allow(unused)]
fn main() {
use processkit::Command;
use std::time::Duration;

let run = Command::new("crunch").start().await?;
run.cpu_time();          // Option<Duration> — user+kernel so far
run.peak_memory_bytes(); // Option<u64>

// …or capture + sample on an interval until exit:
let profile = Command::new("crunch")
    .start().await?
    .profile(Duration::from_millis(100))
    .await?;

println!(
    "outcome={:?} wall={:?} cpu={:?} peak_rss={:?} avg_cpu_cores={:?} ({} samples)",
    profile.outcome,            // Exited(code) / Signalled(sig) / TimedOut
    profile.duration,
    profile.cpu_time,
    profile.peak_memory_bytes,
    profile.avg_cpu_cores(),    // cpu / wall — e.g. Some(1.7) ≈ 1.7 cores busy
    profile.samples,
);
}

These read the child process itself (not a whole tree — that's ProcessGroup::stats), and availability follows the platform: full CPU/memory on Windows and Linux, None where the kernel doesn't account per-process cheaply — see Platform support.


Next: Pipelines · Timeouts, retries & cancellation · Supervision

Pipelines

‹ docs index

a | b | c without a shell. Each stage's stdout feeds the next stage's stdin through an in-process relay (a tokio::io::copy task per boundary) — there is no shell string anywhere, so no quoting rules, no word splitting, no injection surface. Each stage spawns into its own kill-on-drop process group sub-group rather than one shared across the chain, so a per-stage Command::timeout tears down that stage's whole subtree — including a forking stage's grandchildren (sh -c …), not just its direct child — while a chain-wide Pipeline::timeout still tears every stage down together, so the chain lives and dies as a unit. (The relay is an implementation detail, not a kernel splice: a producer whose consumer exits early stops on a broken pipe when the relay's next write fails, rather than instantly via SIGPIPE.)

Building and running

Command::pipe(next) starts a Pipeline; chain more stages with Pipeline::pipe; drive it with output_string() or run():

use processkit::Command;

#[tokio::main]
async fn main() -> processkit::Result<()> {
    // git log --format=%an | sort | uniq -c
    let authors = Command::new("git").args(["log", "--format=%an"])
        .pipe(Command::new("sort"))
        .pipe(Command::new("uniq").arg("-c"))
        .run()                         // require every stage to succeed
        .await?;
    println!("{authors}");
    Ok(())
}

The verbs mirror Command's, each operating on the pipefail outcome:

VerbReturnsA failing stage is…
output_string()ProcessResult<String>…reported in the result (code/stderr/program of the first unclean stage)
output_bytes()ProcessResult<Vec<u8>>…same, with the last stage's stdout captured raw (binary pipes)
run()trimmed final stdout…raised as that stage's Error::Exit; fails loud on a truncated capture
checked()full ProcessResult<String>…raised as Error::Exit (untrimmed stdout)
run_unit()()…raised as Error::Exit (output discarded)
exit_code()i32…its attributed code (no code → Error::Timeout/Signalled)
probe()bool0true, 1false, else Err
parse(|s| …) / try_parse(|s| …)T…raised as Error::Exit; fails loud on a truncated capture

Err from output_string itself means a stage couldn't be started or driven at all (spawn failure, broken plumbing) — never a mere non-zero exit.

The streaming first_line probe is deliberately not a pipeline verb: a chain consumes its last stage in full to fold the pipefail outcome. To capture the first matching line of a finished chain, add a | head -n1 (Unix) / grep -m1 / findstr stage and capture. This does not cover a streaming readiness probe over a chain that must keep running (e.g. wait for a banner line, then leave the chain alive) — | head would tear it down; use a single Command with first_line for that.

The | operator is sugar for the same thing — a | b | ca.pipe(b).pipe(c). Parenthesize the chain before a terminal verb, since method calls bind tighter than |:

#![allow(unused)]
fn main() {
use processkit::Command;

let authors = (Command::new("git").args(["log", "--format=%an"])
    | Command::new("sort")
    | Command::new("uniq").arg("-c"))
    .run()
    .await?;
}

Semantics: pipefail and the ends

The outcome is pipefail, like set -o pipefail in a shell:

  • stdout is always the last stage's output — that's what the chain produced.
  • code, stderr, and the reported program come from the culprit stage: the leftmost stage that didn't exit cleanly (non-zero, signal-killed, or timed out), but preferring a real failure over a downstream SIGPIPE victim — a stage killed only because a later stage closed the pipe early. A checked stage's failure also fires an internal teardown that kills every other stage's subtree proactively, concurrently with the ordered gather — not only passively through pipe EOF — so a quiet upstream producer that never writes (and so never dies of a broken pipe) can't hold the run open indefinitely after a downstream failure. Stages ended by that teardown are flagged torn_down and deprioritized in the fold the same way a SIGPIPE victim is, so the blame stays on the stage that actually failed. If every failure is such a victim (broken-pipe or torn-down), the leftmost one wins; when every stage succeeded, the last stage speaks.
#![allow(unused)]
fn main() {
use processkit::Command;

let result = Command::new("cat").arg("data.txt")
    .pipe(Command::new("grep").arg("ERROR"))      // suppose grep exits 2 (bad pattern)
    .pipe(Command::new("wc").arg("-l"))
    .output_string()
    .await?;

// Diagnostics point at grep — the first unclean stage — while stdout is
// whatever wc managed to print:
assert_eq!(result.code(), Some(2));
println!("blamed: {}", result.ensure_success().unwrap_err()); // names `grep`
}

The ends of the chain behave like a single Command:

  • The first stage's configured stdin source is honored — feed the whole pipeline from a string, file, or stream.
  • Inner stages read from the pipe, full stop: any stdin source or keep_stdin_open configured on them is overridden.
  • Inner stages' stderr is captured per-stage for pipefail diagnostics; only the last stage's stdout reaches you.
#![allow(unused)]
fn main() {
use processkit::{Command, Stdin};

let unique_count = Command::new("sort")
    .stdin(Stdin::from_iter_lines(["b", "a", "b", "c"]))
    .pipe(Command::new("uniq"))
    .pipe(Command::new("wc").arg("-l"))
    .run()
    .await?;
assert_eq!(unique_count.trim(), "3");
}

Unchecked stages

Strict pipefail has one classic false positive: a consumer that legitimately stops reading early. In producer | head -1 the consumer exits 0 after one line and closes the pipe; the producer then stops on a broken pipe — its next write fails once the relay's downstream is gone (a broken-pipe write error, or SIGPIPE where the OS delivers it) — a perfectly normal death that strict pipefail would blame the chain for. Mark that stage unchecked_in_pipe():

#![allow(unused)]
fn main() {
use processkit::Command;

// seq 1 1000000 | head -1 — the producer's broken-pipe death is expected.
let first = (Command::new("seq").args(["1", "1000000"]).unchecked_in_pipe()
    | Command::new("head").args(["-n", "1"]))
    .run()
    .await?;
assert_eq!(first.trim(), "1");
}

The rules (a design borrowed from duct's unchecked() — the idea, not the code):

  • An unchecked stage's unclean exit — a non-zero code, a broken-pipe write failure (or SIGPIPE where the OS delivers it) from a consumer that closed early, or its own per-stage timeout kill — is skipped when the chain decides what to report.
  • A checked failure always trumps an unchecked one, regardless of position: unchecked never shields another stage's real failure.
  • A chain whose only failures are unchecked reports success with the last stage's stdout and its real exit code preserved (not a fabricated 0 — the accepted-code set is widened to include it). The carve-out is for an exit status only: a last stage killed by a signal or its own timeout is not a status to forgive and still surfaces as the failure.
  • unchecked forgives exit status only — never a whole-chain Pipeline::timeout, and it has no effect on a Command run outside a pipeline (a single run's status is already plain data in its ProcessResult).

Timeouts

Two scopes, deliberately distinct:

#![allow(unused)]
fn main() {
use processkit::Command;
use std::time::Duration;

let out = Command::new("producer")
    .timeout(Duration::from_secs(10))      // per-STAGE: kills just `producer`
    .pipe(Command::new("consumer"))
    .timeout(Duration::from_secs(30))      // whole-CHAIN: Pipeline::timeout
    .output_string()
    .await?;
}
  • Pipeline::timeout bounds the whole chain: at the deadline every stage's sub-group is torn down together and the result reports timed_out (no partial stdout — unlike a single command's captured timeout).
  • A per-stage Command::timeout tears down just that stage's own sub-group — its whole subtree, including a forking stage's grandchildren (sh -c …), not just its direct child. Every stage is evaluated by the same pipefail rule (D14): a stage that hit its own deadline — inner or last — surfaces on run() as that stage's Error::Timeout, reporting that stage's own deadline (not the chain's, and never 0ns).

Cancellation has two forms. Pipeline::cancel_on(token) is the chain-level control: the token gap-fills into every stage that doesn't already carry its own [Command::cancel_on] (an explicit per-stage token is left intact), so firing it tears the whole chain down and the run resolves to Error::Cancelled. (A cancel_on token on an individual stage Command also cancels that stage and errors the pipeline, but the pipeline-level builder is the clearer authority.) See Timeouts & cancellation.

Re-running a pipeline

A Pipeline is Clone and re-runnable — stages are re-cloned per run. The one caveat is inherited from Command: a one-shot stdin source on the first stage (Stdin::from_reader / from_lines) is consumed by the first run; re-running then fails loud (an Error::Io at launch, D10) rather than silently feeding empty stdin. Use the reusable sources (from_string / from_bytes / from_iter_lines / from_file) when a chain runs more than once.


Next: Timeouts, retries & cancellation · Running commands · Process groups

Timeouts, retries & cancellation

‹ docs index

Three ways a run ends early, with three different philosophies:

  • a timeout is data — the deadline was part of the run's contract, so its expiry is captured in the result (and only the success-checking verbs turn it into an error);

  • a retry is a policy — the success-checking verbs replay the run while your classifier says the failure is transient;

  • a cancellation is an abandonment — the caller changed its mind, so every path reports an error; there is no result worth inspecting.

  • Timeouts

  • Retries

  • Cancellation

  • Precedence and interactions

Timeouts

Command::timeout(d) kills the whole process tree at the deadline — not just the direct child, so a wrapper script's grandchildren die too.

#![allow(unused)]
fn main() {
use processkit::Command;
use std::time::Duration;

// Captured: inspect the flag yourself.
let result = Command::new("slow-tool")
    .timeout(Duration::from_secs(5))
    .output_string()
    .await?;
if result.timed_out() {
    println!("partial output before the kill: {}", result.stdout());
}

// Raised: the checking verbs convert the flag into a typed error.
let err = Command::new("slow-tool")
    .timeout(Duration::from_secs(5))
    .run()
    .await
    .unwrap_err();
assert!(matches!(err, processkit::Error::Timeout { .. }));
}

Where each verb lands:

VerbDeadline expiry becomes
output_string() / output_bytes()Ok result with timed_out() == true, code() == None, partial output kept
run() / exit_code() / probe() / checked()Error::Timeout { program, timeout, stdout, stderr } — the partial output captured before the kill is attached (err.diagnostic() surfaces a hung tool's last words)
first_line(pred)Error::Timeout (the line never arrived in time)
start() + streamingthe stream ends at the deadline (tree killed, pipes closed); finish then reports the kill (outcome == Outcome::TimedOut)
ensure_success() on a captured resultError::Timeout, checked before the exit code
Pipelinechain deadline → timed_out result; per-stage deadlines fold into pipefail

Error::Timeout is #[non_exhaustive] (like Exit and Signalled), so read it through accessors instead of destructuring { program, timeout, stdout, stderr } — a full literal no longer compiles, and a later field would break one. Alongside the lossy-UTF-8 err.stdout() / err.stderr() strings it now carries err.stdout_bytes() -> Option<&[u8]>: the exact captured bytes when the failure rose from a byte path (output_bytes().await?.ensure_success()?), and None on the text verbs (run / output_string / checked), whose decoded stdout is already complete. The full error table lists every variant.

Two distinct deadline families to keep apart:

  • Command::timeout — the run's own contract, this section.
  • The readiness probes' within parameter — gives Error::NotReady and never kills the child.

Optional and explicitly-unbounded deadlines

A CliClient's default_timeout is a gap-fill: it lends its deadline to any command that didn't set one, and a per-command timeout overrides it. Two verbs shape that fill:

  • Command::no_timeout() marks a command explicitly unbounded — a tail -f, a watch loop — so the client's default_timeout leaves it alone. That is not the same as simply never setting a timeout (which the client does fill): you are opting out on purpose. The last of timeout / no_timeout on the chain wins, so a late no_timeout() clears an earlier deadline and a late timeout(d) restores one.
  • Command::timeout_opt(Option<Duration>) folds a config value straight in: Some(d) is exactly timeout(d), None is exactly no_timeout() — the deliberate opt-out, not a silent fall back to the client default. It collapses the match cfg { Some(d) => c.timeout(d), None => c.no_timeout() } dance at a config-driven call site into one call.
#![allow(unused)]
fn main() {
use processkit::Command;
use std::time::Duration;

// Config-driven deadline: `Some(d)` behaves like `timeout(d)`, `None` like
// `no_timeout()` (stay unbounded — do NOT fall back to a client default_timeout).
let deadline: Option<Duration> = None;
Command::new("worker")
    .timeout_opt(deadline)
    .run()
    .await?;

// A follower that must outlive an ambient default_timeout — opt out explicitly:
Command::new("tail").args(["-f", "app.log"]).no_timeout().run().await?;
}

Graceful timeout

By default the deadline hard-kills at once. Add timeout_grace(d) to give the tree a chance to clean up: at the deadline it is sent SIGTERM (or the signal chosen with timeout_signal, which needs the process-control feature), allowed up to the grace window to exit, then SIGKILLed — the same SIGTERM → wait → SIGKILL tier as ProcessGroup::shutdown. A signal-handling child that exits ends the grace early.

#![allow(unused)]
fn main() {
use processkit::Command;
use std::time::Duration;

let result = Command::new("slow-tool")
    .timeout(Duration::from_secs(30))
    .timeout_grace(Duration::from_secs(5)) // SIGTERM, wait up to 5s, then SIGKILL
    .output_string()
    .await?;
}

timed_out() is true regardless of whether the child exited on the signal or was SIGKILLed after the grace — the deadline is what fired. Windows has no signal tier: timeout_grace is accepted but the deadline kills the job atomically.

The explicit RunningProcess::shutdown(grace) verb (stop a started handle on demand) composes with a Command::timeout: its own SIGTERM → grace → SIGKILL is the single teardown (it does not also fire the run's timeout teardown), and if the deadline has already elapsed when you call shutdown, the outcome is reported as Outcome::TimedOut — the grace you pass governs the teardown timing.

Retries

retry(max_attempts, backoff, classifier) replays a failed run — up to max_attempts total attempts, sleeping backoff between tries, retrying only while the classifier accepts the error:

#![allow(unused)]
fn main() {
use processkit::{Command, Error};
use std::time::Duration;

let out = Command::new("curl")
    .args(["-fsS", "https://example.com/api"])
    .timeout(Duration::from_secs(10))
    .retry(3, Duration::from_millis(250), |e| {
        // transient: network timeouts and curl's "couldn't connect" (7)
        matches!(e, Error::Timeout { .. })
            || matches!(e, Error::Exit { code: 7, .. })
    })
    .run()
    .await?;
}

Ground rules:

  • Retries apply to the success-checking paths only (run, exit_code, probe, ProcessRunnerExt::checked — and everything built on them, e.g. CliClient). The non-erroring output_string capture never retries: it didn't fail.
  • The classifier sees the typed error — match on variants, codes, even the captured stderr.
  • Each attempt re-runs the same Command: a one-shot stdin source (table) is consumed by attempt #1, so attempt #2 fails loud with an Error::Io (InvalidInput) at launch (D10) rather than silently feeding empty stdin. Use reusable sources for retried commands.
  • A Cancelled error is never retried, classifier or not — the token stays cancelled forever, so another attempt could only fail the same way.

Exponential backoff with RetryPolicy

retry's backoff is a fixed delay — the same pause between every attempt, unchanged. For network-style load, reach for retry_with(policy, classifier) and a RetryPolicy: exponential growth, a per-delay cap, and optional AWS-style full jitter (each wait spread uniformly over [0, delay], so a herd of retriers doesn't resynchronize). It drives the same retry loop as retry — only the delay schedule differs.

#![allow(unused)]
fn main() {
use processkit::{Command, Error, RetryPolicy};
use std::time::Duration;

let policy = RetryPolicy::new()   // == RetryPolicy::default()
    .max_retries(5)               // 5 retries after the first → 6 attempts total
    .initial_backoff(Duration::from_millis(100))
    .multiplier(2.0)              // 100ms, 200ms, 400ms, … (before cap + jitter)
    .max_backoff(Duration::from_secs(30))
    .jitter(true);

let out = Command::new("curl")
    .args(["-fsS", "https://example.com/api"])
    .timeout(Duration::from_secs(10))
    .retry_with(policy, |e| matches!(e, Error::Timeout { .. }))
    .run()
    .await?;
}

RetryPolicy is #[non_exhaustive]; build it with new() / default() and the setters above. Mind the counting difference: retry(n, …) counts n total attempts, but a RetryPolicy counts max_retries — the runs after the first — so retry(3, …) corresponds to RetryPolicy::new().max_retries(2). Default is a sane cloud client: 3 retries (4 attempts) / 100 ms initial / ×2 / 30 s cap / jitter on.

Client-wide retries and opting out

Set a policy once on a CliClient and every verb it runs retries on it — the client-wide analogue of a per-command retry:

#![allow(unused)]
fn main() {
use processkit::{CliClient, RetryPolicy};

let gh = CliClient::new("gh")
    .default_retry(RetryPolicy::default(), |e| e.is_transient());
// every gh.run(...) / gh.checked(...) now retries a transient failure
}

A per-command retry or retry_with wins over the client default (gap-fill, not override — the same rule as timeout vs. default_timeout). To make one command a single shot despite the default, chain Command::retry_never(): it runs exactly once and suppresses the fill, symmetric with no_timeout() (and tidier than the retry(1, Duration::ZERO, |_| false) idiom it replaces).

For "keep it alive" (restart a service whenever it exits) rather than "replay this one operation", use a Supervisor — same backoff shape, different loop condition.

Every backoff sleep is cancelable: it is raced against the command's cancellation token, so a cancel landing mid-backoff resolves promptly with Error::Cancelled instead of waiting out a (possibly 30 s) delay. The same holds for a Supervisor's restart backoff and its storm-pause windows.

Cancellation

Hand any command a CancellationToken (re-exported at the crate root); cancelling the token kills the run's tree and makes every consuming path report Error::Cancelled:

use processkit::{CancellationToken, Command};

#[tokio::main]
async fn main() -> processkit::Result<()> {
    let shutdown = CancellationToken::new();

    // Wire the same parent token into many jobs via child tokens:
    let job = tokio::spawn({
        let token = shutdown.child_token();
        async move {
            Command::new("long-export").cancel_on(token).run().await
        }
    });

    // Ctrl-C handler, sibling failure, UI button, …
    shutdown.cancel();

    assert!(matches!(
        job.await.unwrap(),
        Err(processkit::Error::Cancelled { .. })
    ));
    Ok(())
}

The contract, path by path:

SituationBehavior
Cancel during run / output_string / output_bytes / wait / profile / exit_code / probetree killed, Error::Cancelled { program }
Cancel during streaming (stdout_lines)the stream ends; the following finish reports Error::Cancelled
Token already cancelled before the runshort-circuits before spawning — no process is ever created
Cancel on a shared-ProcessGroup handlekills the child itself, leaves the group's siblings alone (same scope as a timeout)
A Pipeline stage's token cancelsthat stage dies; the cancellation errors the whole pipeline and the private group reaps the other stages
Under retryterminal — never retried, whatever the classifier says
Under a Supervisorterminal — supervision returns Err(Cancelled) instead of restarting into a still-cancelled token
wait_any mid-runsurfaces Err(Cancelled) — each racer's wait path resolves to Cancelled when its token fires, the same as a bulk verb (a pre-cancelled token still hits the pre-spawn short-circuit)
first_line mid-runsurfaces Error::Cancelled once the token fires — a cancelled stream that closes without a match is reported as cancellation, not Ok(None)

Client-level default

A typed wrapper built on CliClient usually constructs and consumes its Commands internally — there is no place to chain a per-call cancel_on. Set the token once on the client; every command it builds carries it:

#![allow(unused)]
fn main() {
use processkit::{CancellationToken, CliClient};

let token = CancellationToken::new();
let gh = CliClient::new("gh").default_cancel_on(token.child_token());
// ... controller cancels `token` → every in-flight command of THIS client
// dies (whole tree), surfacing Error::Cancelled to the awaiting call.
}

Clients are cheap — scope cancellation by building one client per cancellable scope with its own (child) token, instead of threading tokens through call signatures. cli_client!-generated wrappers re-emit the builder, so Git::new().default_cancel_on(t) works for downstream crates too.

Precedence: a per-command cancel_on chained on a built command replaces the client default (explicit beats default, like a per-command timeout after default_timeout). To honor both sources, wire it explicitly — CancellationToken has no built-in merge: derive a child of the default (let c = default.child_token()), hand the command cancel_on(c.clone()), and have the second source call c.cancel(). Or simpler: build a dedicated client per scope.

Precedence and interactions

Timeout vs. cancellation. A timeout is captured; a cancellation is always an error. When both land on the same run, cancellation wins — you asked the run to stop mattering, so no result is synthesized:

#![allow(unused)]
fn main() {
use processkit::{CancellationToken, Command};
use std::time::Duration;

let token = CancellationToken::new();
token.cancel();

let err = Command::new("tool")
    .timeout(Duration::from_millis(1))   // would have been a Timeout…
    .cancel_on(token)                    // …but cancellation takes priority
    .run()
    .await
    .unwrap_err();
assert!(matches!(err, processkit::Error::Cancelled { .. }));
}

Which knob for which job:

You wantReach for
"This run may not take longer than X"Command::timeout
"This operation is flaky, try a few times"Command::retry
"Stop everything when the app shuts down"cancel_on + one shared token
"Keep this service alive across crashes"Supervisor
"Tell me when it's ready, don't kill it"readiness probes

Next: Supervision · Streaming & interactive I/O · Running commands

Supervision

‹ docs index

Where retry answers "run this once, replaying on failure", a Supervisor answers the different question "keep this alive": restart a child per policy whenever it exits, with bounded restarts, exponential backoff, and jitter — a minimal runit/systemd-style keeper, platform-agnostic because it sits entirely on the ProcessRunner seam.

The shape

use processkit::{Command, RestartPolicy, Supervisor};
use std::time::Duration;

#[tokio::main]
async fn main() -> processkit::Result<()> {
    let outcome = Supervisor::new(Command::new("my-server").args(["--port", "8080"]))
        .restart(RestartPolicy::OnCrash)           // default
        .max_restarts(5)                           // default: unlimited
        .backoff(Duration::from_millis(200), 2.0)  // default: 200ms × 2.0
        .max_backoff(Duration::from_secs(30))      // default: 30s cap
        .jitter(true)                              // default: on
        .stop_when(|res| res.code() == Some(0))    // optional exit condition
        .run()
        .await?;

    println!(
        "ended after {} restarts, reason: {:?}, last exit: {:?}",
        outcome.restarts, outcome.stopped, outcome.final_result.code(),
    );
    Ok(())
}

Each incarnation is one full captured run of the command (so the command's own timeout, stdin, env, … all apply per run — with the usual one-shot-stdin caveat for the second run onward).

Policies: what counts as a crash

A crash is any run that is not a success (ProcessResult::is_success, which honors the command's ok_codes): an exit code outside the accepted set (default {0}), a timeout, a signal-kill, or a spawn failure. A command with ok_codes([0, 2]) that exits 2 is a success, so OnCrash treats it as clean, not a crash.

RestartPolicyRestarts after…
OnCrash (default)crashes only; a clean exit ends supervision (PolicySatisfied)
Alwaysevery completed run, clean or not — pair it with stop_when/max_restarts or it loops forever
Nevernothing: one run, reported as-is

Backoff and jitter

The n-th restart (0-based) sleeps

delay(n) = min(base × factor^n, max_backoff) × jitter

with jitter drawn uniformly from [0.5, 1.5) per restart. Jitter is on by default so a fleet of supervised workers restarted by the same incident doesn't stampede back in lockstep; jitter(false) gives deterministic delays (useful in tests with a paused tokio clock). A non-finite or < 1.0 factor is treated as 1.0 — constant delay, never a shrinking one.

base=200ms, factor=2.0, cap=30s:
restart #0 → ~200ms   #1 → ~400ms   #2 → ~800ms … #7 → ~25.6s   #8+ → 30s (cap)

The exponent n is not the lifetime restart count — it resets whenever a run stays up at least as long as max_backoff. So a long-lived service that crashes now and then restarts near base each time (it demonstrated health between crashes), while only a tight loop — each incarnation shorter than the cap — climbs to the ceiling. The floor is on uptime, not exit kind: under Always a worker that exits (cleanly or not) in under max_backoff is treated as flapping and escalates, which is what stops an exit 0 spin loop from hammering at the base delay. A single jittered delay can reach up to 1.5 × max_backoff.

Failure storms

Backoff spaces individual restarts; max_restarts is a lifetime cap. Neither distinguishes a service that fails once a day from one that is suddenly crash-looping. The opt-in storm guard does (a design borrowed from Go's suture supervisor — the idea, not the code):

#![allow(unused)]
fn main() {
use processkit::{Command, Supervisor};
use std::time::Duration;

let outcome = Supervisor::new(Command::new("worker"))
    .storm_pause(Duration::from_secs(15))     // master switch — off by default
    .failure_decay(Duration::from_secs(30))   // score half-life (default 30s)
    .failure_threshold(5.0)                   // trip point (default 5.0)
    .run()
    .await?;

println!("storm pauses taken: {}", outcome.storm_pauses);
}

Each failed run adds 1 to a score that halves every failure_decay:

score = score × 0.5^(Δt / failure_decay) + 1
  • Fails rarely: the score decays back toward 1 between failures and never reaches the threshold — the guard stays out of the way.
  • Failure storm: failures arrive faster than the half-life drains them, the score climbs past failure_threshold, and the supervisor takes one collective pause of storm_pause (jittered into [0.5, 1.5) like the backoff), resets the score, and resumes.

Only failures feed the score — crashes and spawn errors, not clean exits restarted under RestartPolicy::Always. The pause stacks with (runs before) the per-restart backoff, and the max_restarts budget is checked first, so a storm pause never extends an exhausted budget. Pauses taken are reported in SupervisionOutcome::storm_pauses.

Stopping

Four gates, checked in this order after every completed run:

  1. stop_when(predicate) — sees the run's ProcessResult; returning true ends supervision regardless of policy (→ StopReason::Predicate). "Exit 0 is done, anything else is a crash" is the classic: stop_when(|res| res.code() == Some(0)) under RestartPolicy::Always.
  2. The policyOnCrash stops on a clean exit (→ PolicySatisfied).
  3. give_up_when(classifier) — only consulted for a crash the policy would otherwise restart; recognizing it as permanent ends supervision (→ StopReason::GaveUp). See Giving up on permanent failures.
  4. max_restarts(n) — at most n restarts = n + 1 total runs; an exhausted budget reports the last result (→ RestartsExhausted). max_restarts(0) means exactly one run.

give_up_when is checked before max_restarts, so a recognized-permanent crash reports the more specific GaveUp even when the budget hasn't run out yet — and before the failure-storm guard, so giving up never pays for a storm pause it was going to end anyway.

Giving up on permanent failures

Without more, the supervisor cannot tell a transient crash from a permanent one — a command that can never succeed (a missing binary, a config error that crashes on startup, a permanently-taken port) restarts forever under the default unlimited OnCrash, throttled only by backoff. give_up_when lets you recognize the unrecoverable case and stop instead:

#![allow(unused)]
fn main() {
use processkit::{Command, GiveUpAttempt, Supervisor};

let outcome = Supervisor::new(Command::new("maybe-typo'd-binary"))
    .give_up_when(|attempt| match attempt {
        // The child never even started — e.g. ENOENT for a mistyped name.
        GiveUpAttempt::Failed(err) => err.is_not_found(),
        // A completed crash your own domain knows is permanent, e.g. a
        // documented "config invalid, do not restart" exit code.
        GiveUpAttempt::Crashed(res) => res.code() == Some(78),
        _ => false,
    })
    .run()
    .await?;

println!("stopped: {:?}", outcome.stopped);
}

GiveUpAttempt distinguishes the two shapes a permanent failure can take:

  • Crashed(&ProcessResult<String>) — a completed run that counts as a crash. A match reports StopReason::GaveUp in the SupervisionOutcome, same as any other stop reason.
  • Failed(&Error) — the child never started at all (spawn/IO failure, e.g. ENOENT). There is no ProcessResult to report here, so a match surfaces the classified error directly as run()'s Err — the same contract an exhausted budget already has on this path (see Errors and cancellation).

Unset by default: a permanent failure restarts forever exactly as before, bounded only by max_restarts / the storm guard — adding give_up_when never changes existing behavior until you set it.

Outcomes

run() resolves to a SupervisionOutcome:

#![allow(unused)]
fn main() {
let outcome = Supervisor::new(Command::new("job")).run().await?;

outcome.final_result; // ProcessResult<String> of the LAST run
outcome.restarts;     // how many restarts happened (not counting run #1)
outcome.stopped;      // StopReason::{Predicate, PolicySatisfied, GaveUp, RestartsExhausted}
outcome.storm_pauses; // failure-storm pauses taken (0 unless storm_pause is set)
}

Note run() returning Ok does not mean the child succeeded — it means supervision concluded. Inspect final_result (or ensure_success() it) for the child's own verdict.

Supervising inside a shared group

The supervisor runs through any ProcessRunner. The headline production variant injects a ProcessGroup so every incarnation — and everything it spawns — lives in one kill-on-drop container:

#![allow(unused)]
fn main() {
use processkit::{Command, ProcessGroup, RestartPolicy, Supervisor};

let group = ProcessGroup::new()?;

let outcome = Supervisor::new(Command::new("worker"))
    .with_runner(&group)                 // &group is itself a ProcessRunner
    .restart(RestartPolicy::OnCrash)
    .max_restarts(10)
    .run()
    .await?;

// The group outlives supervision: drop it (or shutdown) to reap any strays.
}

Mind one interaction: don't supervise into a group you've suspended — under the cgroup mechanism the restarted child would start frozen (and the spawn itself can block). Resume first.

The same injection point makes supervision logic hermetically testable — script a sequence of fake results and assert the restart/stop behavior with no real process; see Testing your code.

Errors and cancellation

A run that produces no result at all (spawn/IO failure) can't be judged by stop_when (it needs a ProcessResult) — but it is visible to give_up_when as GiveUpAttempt::Failed. The policy treats it as a crash and restarts (with backoff) unless the policy is Never, give_up_when classifies it as permanent, or the budget is exhausted — any of those surfaces the error itself as run()'s Err.

A cancelled incarnation is terminal: run() returns Err(Error::Cancelled) immediately. The token never un-cancels, so a restart could only produce another instantly-cancelled run — the supervisor refuses the futile loop.

The backoff sleep itself (and a storm pause) is also cancellable: a token firing mid-delay resolves promptly with Error::Cancelled instead of waiting out the full — possibly 30s+ — delay.


Next: Testing your code · Timeouts, retries & cancellation · Process groups

Testing your code

‹ docs index

Code that shells out is miserable to test — unless the subprocess is behind a seam. In processkit that seam is one small trait. Only output_string is required; output_bytes (raw-byte stdout) and start (a live handle for streaming/probes) are defaulted, so a minimal double implements just output_string:

#[async_trait]
pub trait ProcessRunner: Send + Sync {
    async fn output_string(&self, command: &Command) -> Result<ProcessResult<String>>;
    // Defaulted (route through `start`); override for byte/streaming support:
    async fn output_bytes(&self, command: &Command) -> Result<ProcessResult<Vec<u8>>>;
    async fn start(&self, command: &Command) -> Result<RunningProcess>;
}

Production code takes a runner (generically or as &dyn ProcessRunner); tests hand it a double. Five doubles ship with the crate, plus a macro that makes whole CLI wrappers testable for free.

The ProcessRunner seam

JobRunner is the real implementation (each run in a fresh private group); a ProcessGroup is also a runner (runs land in that shared group); and impl ProcessRunner for &R means a borrowed runner works wherever an owned one does — inject &group or &recording without giving ownership away.

Every runner — real or double — gets the convenience helpers of ProcessRunnerExt for free: run (trimmed stdout, success required), run_unit, exit_code, probe (exit code as a boolean), checked (success-checked full result), and parse/try_parse (feed stdout to a closure). These are all callable on a &dyn ProcessRunner; being generic over the closure, parse/try_parse/first_line simply can't be dispatched through a dyn ProcessRunnerExt object (the ext trait isn't object-safe). Retry policies work through the seam too, so a double exercises your retry handling hermetically.

The seam covers streaming as well as bulk runs: ProcessRunner::start returns a live RunningProcess, and a ScriptedRunner's start hands back a scripted handle whose canned lines flow through the same pump machinery a real child uses — stdout_lines, wait_for_line, and finish behave identically, with no subprocess (see Scripted streaming below). An output_string-only custom runner keeps compiling: start is defaulted to Error::Unsupported.

#![allow(unused)]
fn main() {
use processkit::{Command, ProcessRunner, ProcessRunnerExt, Result};

// Production code: generic over the runner.
async fn current_branch(runner: &impl ProcessRunner) -> Result<String> {
    runner
        .run(&Command::new("git").args(["branch", "--show-current"]))
        .await
}
}

Scripting replies

ScriptedRunner returns canned Replys for matched commands — the work-horse double:

#![allow(unused)]
fn main() {
use processkit::{Command, ProcessRunnerExt};
use processkit::testing::{Reply, ScriptedRunner};

#[tokio::test]
async fn detects_the_branch() {
    let runner = ScriptedRunner::new()
        // Match by program + argument PREFIX (element-wise; first element is
        // the program name, in registration order):
        .on(["git", "branch", "--show-current"], Reply::ok("main\n"))
        // …or by any predicate over the full Command:
        .when(
            |cmd| cmd.working_dir().is_some(),
            Reply::fail(128, "fatal: not a git repository"),
        )
        // …with an optional catch-all:
        .fallback(Reply::ok(""));

    assert_eq!(current_branch(&runner).await.unwrap(), "main");
}
}

The pieces:

  • Reply::ok(stdout) — exit 0. Reply::fail(code, stderr) — non-zero with stderr. Reply::lines(["a", "b"]) — exit 0 with the lines joined (and streamed one by one on a scripted start). Reply::timeout() — a timed-out run (the checking helpers raise Error::Timeout from it, carrying the command's own configured deadline). On a scripted start it resolves immediately as timed-out; to exercise a real deadline race, use Reply::pending() + a Command::timeout. .with_stdout(text) — attach stdout to any of them (e.g. the CONFLICT … text git prints on a failing merge). .with_stderr(text) — attach stderr to any reply, including a successful one (Reply::ok("out").with_stderr("warning\n")) to model a CLI that writes warnings to stderr even on exit 0; on a Reply::fail it overrides the fail-supplied stderr. .with_line_delay(d) — pace a scripted stream's lines.
  • Reply::pending() — parks the call until the command's cancellation token (per-command cancel_on or the client-level default_cancel_on) fires, then resolves with Error::Cancelled — so a test can prove an orchestration actually cancels a blocked call, not just that it formats a canned error. With no token it parks forever, like a hung child.
  • Rules are tried in registration order; first match wins. Prefix matching is element-wise over the program name then the arguments (the first element is the program) — on(["git", "foo"]) matches git foo bar but not git foobar (and not rm foo). Use on_sequence to serve an ordered sequence of replies (each once, then the last repeats) for a fail-then-succeed scenario.
  • No match and no fallback is a loud error (Error::Spawn, not-found) — an unexpected invocation can't slip through a test silently.
  • Bulk runs also replay the canned lines through the command's on_stdout_line/on_stderr_line handlers, so a wrapper's progress-reporting path is exercised without a subprocess.

Scripted streaming

ScriptedRunner::start returns a live RunningProcess backed by the canned reply instead of an OS child. The canned stdout/stderr feed the same pump machinery a real child uses, so the whole streaming surface works hermetically — stdout_lines yields the lines, wait_for_line probes them, finish reports the canned outcome and stderr:

#![allow(unused)]
fn main() {
use processkit::{Command, Outcome, ProcessRunner, Finished};
use processkit::prelude::StreamExt;
use processkit::testing::{Reply, ScriptedRunner};
use std::time::Duration;

#[tokio::test]
async fn server_becomes_ready() {
    let runner = ScriptedRunner::new()
        .on(["server", "serve"], Reply::lines(["booting", "listening on 8080"]));

    let mut run = runner.start(&Command::new("server").arg("serve")).await.unwrap();
    run.wait_for_line(|l| l.contains("listening"), Duration::from_secs(5))
        .await
        .unwrap(); // satisfied by the canned banner — no subprocess

    let Finished { outcome, .. } = run.finish().await.unwrap();
    assert_eq!(outcome, Outcome::Exited(0));
}
}

Reply::lines([...]) scripts the stdout lines; .with_line_delay(d) paces them (deterministic under #[tokio::test(start_paused = true)]), and the scripted run "exits" after the last line. The honest boundaries: a scripted handle has no OS identity (pid() is None, profile reports empty samples), does not compose into a real Pipeline, and does not model interactive stdin. Reply::pending() scripts a run that never exits on its own — cancel or time it out through the command's own knobs. A command timeout does bound a scripted stream (it ends at the deadline and reports Outcome::TimedOut, like a real child), but a scripted handle has no signal tier, so — like on Windows — it ignores timeout_grace and ends at once.

Asserting invocations

RecordingRunner wraps another runner and records every Invocation — what was asked — so a test asserts inputs, not just outputs:

#![allow(unused)]
fn main() {
use processkit::{Command, ProcessRunnerExt};
use processkit::testing::{RecordingRunner, Reply, ScriptedRunner};

#[tokio::test]
async fn passes_the_right_flags() {
    let runner = RecordingRunner::new(
        ScriptedRunner::new().fallback(Reply::ok("done")),
    );

    runner
        .run(&Command::new("gh").args(["pr", "create", "--draft"]).current_dir("/repo"))
        .await
        .unwrap();

    let call = runner.only_call(); // panics unless exactly one call
    assert_eq!(call.args_str(), ["pr", "create", "--draft"]);
    assert!(call.has_flag("--draft"));
    assert_eq!(call.cwd.as_deref().map(|c| c.to_str().unwrap()), Some("/repo"));
    assert!(!call.has_stdin);
}
}

An Invocation captures the routing knobs — program, args, cwd, envs (explicit overrides, None = removal), has_stdin — not the I/O-shaping ones (timeout, encodings, buffer policy); assert those through a when predicate over the Command itself. For the captured env overrides, env(name) returns the raw entry (Some(Some(value)) set, Some(None) removed, None untouched), while env_is(name, value) and has_env(name) answer the common yes/no checks directly (a removal is not "having" it, so has_env is false there). calls() returns the full list when more than one run is expected.

Dry-run rendering

DryRunRunner never spawns. It renders each command through the crate's own Command::command_line quoting and returns a synthetic successful result — the seam behind a tool's own --dry-run/--echo mode, so a test can prove what your code would run without running it. The rendered command lines are captured as a collected snapshot (commands() / only_command(), like RecordingRunner) and/or streamed to a live on_invocation callback — usable together or alone.

#![allow(unused)]
fn main() {
use processkit::{Command, ProcessRunnerExt};
use processkit::testing::DryRunRunner;

#[tokio::test]
async fn renders_without_spawning() {
    let runner = DryRunRunner::new();

    runner
        .run(&Command::new("git").args(["push", "--force"]))
        .await
        .unwrap(); // never spawns — returns a synthetic success

    // The rendered line uses the crate's own quoting (an arg with spaces or
    // shell specials comes back safely quoted):
    assert_eq!(runner.only_command(), "git push --force");
}
}

Expectation-style: MockRunner

With the mock feature, mockall generates a MockRunner for expectation-style tests (call counts, argument matchers, ordered expectations) — the right tool when the interaction is the contract.

Note: MockRunner's expect_* surface is generated by mockall and is exempt from this crate's semver guarantees — it tracks the mockall dependency, not a frozen API. For a stable double, prefer ScriptedRunner (canned replies) or RecordingRunner (input assertions) above.

use processkit::testing::MockRunner;

let mut mock = MockRunner::new();
mock.expect_output_string()
    .times(1)
    .returning(|_cmd| /* build a Result<ProcessResult<String>> */ …);

MockRunner does not inherit the defaults. Unlike a hand-written runner (where output_bytes/start are defaulted), mockall::automock replaces every method with an expectation — so a verb that routes through start or output_bytes needs its own expect_start() / expect_output_bytes(), or the unset call panics ("no expectation"). ScriptedRunner provides the defaults and the streaming seam out of the box.

For most tests ScriptedRunner/RecordingRunner read better; reach for the mock when you need mockall's matching machinery.

Record/replay cassettes

With the record feature, RecordReplayRunner closes the loop: record real runs to a JSON cassette once, then replay them deterministically — fast, hermetic, byte-stable, no subprocess in CI:

#![allow(unused)]
fn main() {
use processkit::{Command, JobRunner, ProcessRunnerExt};
use processkit::testing::RecordReplayRunner;

// Record once against the real tool (an opt-in `--record` test run, say):
let runner = RecordReplayRunner::record("fixtures/git.json", JobRunner::new());
let version = runner.run(&Command::new("git").arg("--version")).await?;
runner.save()?;                                  // the error-surfacing flush
                                                 // (best-effort on drop too)

// Replay everywhere else:
let runner = RecordReplayRunner::replay("fixtures/git.json")?;
assert_eq!(runner.run(&Command::new("git").arg("--version")).await?, version);
}

Semantics worth knowing before you commit a cassette:

AspectBehavior
Match keyprogram + args + a stdin source digest (hashed, never persisted: in-memory bytes hash their content, a from_file source hashes its path) — no stdin (absent or Stdin::empty()) keys distinctly; lossy UTF-8 on the text parts. cwd is not matched — it's stored verbatim on each entry for visibility but no longer discriminates two otherwise-identical runs, so a cassette recorded in one working directory replays from another (a dev box → a CI workspace)
Environmentvalues never reach the file — only sorted variable names, so env secrets can't leak through a committed fixture; env is not matched, so env differences can't cause spurious misses
Duplicates of one keyreplay in capture order, then the last entry repeats — a recorded sequence (git rev-parse HEAD before/after a commit) replays faithfully, while retry/probe loops keep getting a stable final answer
Missstrict Error::CassetteMiss (distinct from a missing program — is_not_found() is false) — replay never spawns a surprise subprocess; a stale cassette fails loudly
Timeoutsa recorded timed-out run replays as one, surfacing Error::Timeout with the replaying command's deadline
Formatpretty-printed JSON with a version field; unknown versions / corrupt files / an entry with a contradictory outcome / a file over 64 MiB are Error::Io(InvalidData), a missing file keeps NotFound
Err resultsrecords and replays failed calls — an Err from output_string/start (Error::Spawn/NotFound/Stdin/OutputTooLarge/Unsupported/Io, its ErrorKind preserved by name, plus an Other fallback) is captured and reconstructed on replay, not silently dropped into a spurious CassetteMiss. Error::Cancelled is deliberately never recorded — replay short-circuits on the replaying command's own token first. Non-zero exits and captured timeouts are results and record as before
Verbs (output_string + start)a cassette is verb-agnostic: record through either and replay through either. Replaying start hands back a scripted RunningProcess whose recorded lines flow through the command's real pumps (stdout_lines / wait_for_line / finish), no subprocess. Recording a start captures the run whole (the child runs to completion before the handle returns), so an interactive run fed stdin mid-stream can't be recorded that way — bound it with Command::timeout or script it with ScriptedRunner
output_bytesunsupported (Error::Unsupported) in both modes — a lossy-UTF-8 text fixture can't reproduce exact raw bytes; capture bytes from a real or scripted runner

Only env values are redacted. program, args, cwd, stdout, and stderr are stored verbatim and can carry secrets (a --password=… flag, a token echoed to output), so review a fixture before committing it. On Unix the file is written 0600 and the write refuses to follow a symlink at the cassette path (O_NOFOLLOW, so a planted link can't redirect the secret-bearing write — it fails loud instead). On Windows the file inherits the containing directory's ACL, so restrict that directory (or use a per-user temp dir, not a world-writable shared one) for secret-bearing fixtures.

A neat trick: in tests, record against a ScriptedRunner instead of JobRunner — the whole record→save→replay round trip is then itself hermetic.

Wrapping a CLI tool

CliClient is the foundation for typed wrappers around external tools (git, jj, gh, kubectl, …): it owns the program name, per-client defaults, and the runner; your wrapper contributes only commands and parsers. The cli_client! macro generates the boilerplate:

#![allow(unused)]
fn main() {
use processkit::{cli_client, Error, ProcessRunner, Result};
use std::path::Path;
use std::time::Duration;

cli_client!(
    /// A typed `git` client.
    pub struct Git => "git"
);

impl<R: ProcessRunner> Git<R> {
    /// HEAD's commit id.
    pub async fn head(&self, repo: &Path) -> Result<String> {
        self.core.run(self.core.command_in(repo, ["rev-parse", "HEAD"])).await
    }

    /// Is the work tree clean? (exit code IS the answer)
    pub async fn is_clean(&self, repo: &Path) -> Result<bool> {
        self.core.probe(self.core.command_in(repo, ["diff", "--quiet"])).await
    }

    /// Branch list, parsed — the parser is fallible and returns the crate's
    /// `Result`, typically an `Error::Parse` naming the program.
    pub async fn branches(&self, repo: &Path) -> Result<Vec<String>> {
        self.core
            .try_parse(
                self.core.command_in(repo, ["branch", "--format=%(refname:short)"]),
                |out| {
                    let list: Vec<String> = out.lines().map(str::to_owned).collect();
                    if list.is_empty() {
                        Err(Error::parse("git", "no branches"))
                    } else {
                        Ok(list)
                    }
                },
            )
            .await
    }
}

// Production: the real runner, with per-client defaults.
let git = Git::new().default_timeout(Duration::from_secs(30));
let head = git.head(Path::new(".")).await?;
}

The generated type is Git<R: ProcessRunner = JobRunner> with Git::new(), Git::with_runner(runner), default_timeout / default_env / default_env_remove builders, and a module-private core: CliClient<R> (reach it as self.core from the wrapper's own methods) whose helpers speak the crate-wide verb vocabulary: run (trimmed stdout), output_string (full result), run_unit (success only), exit_code, probe, plus parse (infallible) and try_parse (fallible → Error::Parse).

And the payoff — the wrapper tests hermetically with any double:

#![allow(unused)]
fn main() {
#[tokio::test]
async fn head_is_trimmed() {
    let git = Git::with_runner(
        ScriptedRunner::new().on(["git", "rev-parse", "HEAD"], Reply::ok("abc123\n")),
    );
    assert_eq!(git.head(Path::new("/repo")).await.unwrap(), "abc123");
}
}

…or with a cassette recorded against the real tool once.


Next: Platform support · Supervision · Running commands

Platform support

‹ docs index

processkit supports Unix and Windows only — it requires tokio::process and OS job / process-group primitives that have no equivalent on bare targets like wasm. Building for such a target fails at compile time (a compile_error! guard, or earlier in tokio's own dependencies). Within the supported set, it treats platform support as first-class: every capability is either fully implemented, honestly partial (documented and typed), or refused with Error::Unsupported — never silently skipped. This page collects all the matrices and fine print in one place.

Containment mechanisms

ProcessGroup::mechanism() reports which one you actually got:

MechanismPlatformHow containment works
JobObjectWindowsA Job Object with kill-on-close; children are created suspended, assigned to the job, then resumed — so even a grandchild forked in the first instant is contained
CgroupV2Linux (with delegation)A private cgroup; children join in pre_exec, before exec, so descendants can never escape; teardown is cgroup.kill
ProcessGroupmacOS, BSDs, Linux fallbackPOSIX process groups (setpgid); teardown is killpg; tracked per started/adopted child

On Linux the cgroup backend requires controller delegation, and resource limits specifically need this process to run at the real cgroup-v2 root. The crate creates the limit cgroup under this process's own cgroup and enables the controllers in that cgroup's subtree_control, which cgroup v2's "no internal processes" rule allows only for the real hierarchy root (the one exempt cgroup). A cgroup namespace root does not qualify — it only virtualizes the view — so an ordinary (private-cgroupns) container fails EBUSY just like a systemd session/scope/service. The crate does not migrate your process into a sub-cgroup to work around it, so in practice limits apply only at a minimal non-systemd init sitting at the real root. Without a usable cgroup it quietly falls back to ProcessGroup — unless you requested resource limits, which fail fast instead (Error::ResourceLimit), because an unapplied cap is no protection. That failure is typed rather than left as opaque text: LimitKind (Memory/Processes/Cpu) says which limit, and LimitReason (Unsupported/Unenforceable/Invalid) says whyUnsupported when the platform has no whole-tree containment mechanism at all (macOS/BSD, or Linux with no cgroup v2 mounted), Unenforceable when a capable mechanism exists but this request was rejected (missing cgroup delegation, a failing Windows Job Object call), Invalid for a nonsensical value caught before the OS is ever touched. Read both off the error via Error::limit_kind() / Error::limit_reason() without parsing English text.

Capability matrices

Teardown & containment

CapabilityWindows JobObjectLinux cgroupLinux pgroupmacOS/BSD
Kill-on-drop, whole tree✅ groups-based✅ groups-based
Graceful shutdown (TERM → grace → KILL)🟡 atomic kill only
adopt an external child✅ (future forks contained)✅ (future forks contained)🟡 exec'd child tracked individually🟡 same

Windows has no signal tier, so a graceful shutdown collapses to the atomic Job kill — but it still honors escalate_to_kill: false spares the survivors (closes the Job handle without KILL_ON_JOB_CLOSE) rather than killing them, so the Windows column is "atomic kill when it kills", not an unconditional kill.

Signals & freezing

CapabilityWindowsLinux cgroupLinux pgroupmacOS/BSD
Arbitrary signal (Hup, Usr1, Other(n), …)Kill only
suspend / resume🟡 per-thread countscgroup.freezeSIGSTOP/CONTSIGSTOP/CONT

On the cgroup mechanism, a non-Kill signal (and the SIGSTOP/SIGCONT fallback used for suspend/resume on pre-5.2 kernels without cgroup.freeze) surfaces a real per-member delivery failure (e.g. EPERM) as an Err rather than swallowing it — consistent with the "never silently skipped" philosophy; an ESRCH race (the member already exited) is still success.

Inspection & accounting

CapabilityWindowsLinux cgroupLinux pgroupmacOS/BSD
members()✅ whole tree✅ whole tree🟡 leaders only🟡 leaders only
Group CPU / peak memory❌ count only❌ count only
Per-run cpu_time / peak_memory_bytes / profile✅ (/proc)None

members() is gated on the process-control feature; the CPU / memory / profile rows are gated on the stats feature.

Resource limits (limits feature)

CapabilityWindowsLinux cgroupLinux pgroupmacOS/BSD
max_memory (whole tree)
max_processes
cpu_quota🟡 approximate

Spawn-time controls

CapabilityWindowsUnix (all)
inherit_env allow-list
uid / gid dropUnsupported
setsidUnsupported
create_no_windowno-op
kill_on_parent_death✅ always on (kernel)Linux: direct child; macOS/BSD: no-op
priority
umaskUnsupported

priority never yields Error::Unsupported on either platform — every Priority variant maps to a native primitive (Unix nice/setpriority, Windows a priority-class creation flag). On Unix, raising priority above the inherited value needs CAP_SYS_NICE/root and fails as Error::Spawn without it rather than silently applying a smaller change — see the caveat below. umask is Unix-only (umask(2) applied via pre_exec); non-Unix targets get Error::Unsupported instead of silently ignoring the requested mask.

Everything not listed — capture, streaming, interactive stdin, encodings, buffer policies, timeouts, retry, pipelines, supervision, readiness probes, the test doubles, cassettes, cancellation — is platform-agnostic and behaves identically everywhere.

Caveats

The honest fine print, mostly consequences of OS semantics:

Windows: termination is an exit code, never Signalled (D18). Windows has no signal abstraction, so a killed process reports Outcome::Exited, not Outcome::Signalled. TerminateProcess / TerminateJobObject(_, 1) is Exited(1) — indistinguishable from a voluntary exit(1) — and Ctrl-C surfaces as Exited(-1073741510) (STATUS_CONTROL_C_EXIT as a signed i32). The crate reports the platform truth rather than fabricating a Signalled from an NTSTATUS code (that mapping would be a lossy guess). When you need to know the run was killed, use a ProcessGroup deadline or a cancellation token (which surface as TimedOut / Error::Cancelled on every platform). Outcome::Signalled is therefore Unix-only.

Unix: raising priority needs privilege. Lowering nice below its inherited value — Priority::High (nice(-10)) and Priority::AboveNormal (nice(-5)) — needs CAP_SYS_NICE/root; without it the spawn fails as Error::Spawn rather than silently applying a smaller increase. Priority::Normal (nice(0)) is usually a no-op, but under a positively-niced parent (e.g. a niced CI/batch launcher) it too asks the kernel to raise priority back to the default — the same privileged decrease, and the same failure mode. Windows needs no special privilege for any priority class, so Command::priority never yields Error::Unsupported on either platform — only, on Unix, this conditional Error::Spawn.

Linux cgroup delegation. Creating the per-group cgroup needs write access to the cgroup v2 hierarchy. Dev boxes typically lack it → the pgroup fallback. CI inside containers usually has it. Check mechanism() when behavior must not silently degrade.

uid()/gid() × the cgroup mechanism. The OS applies the uid drop before pre_exec hooks, and the cgroup join runs in pre_exec — as the already-dropped user, who can't write the root-owned cgroup.procs. The spawn fails with a permission error (never an uncontained child). Privilege drop composes cleanly with the process-group mechanism.

setsid() × process groups. A new session implies a new process group; the crate coordinates the two (the containment tracking follows the new session's group), so setsid keeps the kill-on-drop guarantee instead of breaking out of it.

kill_on_parent_death() is thread-scoped on Linux. PR_SET_PDEATHSIG fires when the spawning thread dies, not only the process. On a multi-threaded tokio runtime a retired worker thread could kill the child early; spawn from a current-thread runtime for the strongest guarantee. It covers the direct child only — with the parent SIGKILLed, nothing tears the cgroup/pgroup down, so grandchildren survive. The parent-died-before-arming race is closed by re-checking getppid() in the child against the spawner's pid captured before the fork — which stays correct when the spawner itself is PID 1 (a container entrypoint).

Windows: the suspended-spawn handshake. Children are created CREATE_SUSPENDED, assigned to the job, then resumed — closing the classic race where a fast child forks before it's in the job. A consequence: a raw ProcessGroup::spawn caller passing its own creation flags gets them OR'd with CREATE_SUSPENDED (the Command-driven paths handle this for you, incl. create_no_window).

Windows: nested suspends. SuspendThread keeps per-thread counts — two suspend() calls need two resume()s. The POSIX backends are level-triggered (idempotent). Suspension is also best-effort against a tree that is spawning threads mid-walk. Each thread's owning process is re-verified as still a member of this job (IsProcessInJob) immediately before SuspendThread/ResumeThread — closing a recycle-window race where a member's pid exits and is reused by an unrelated foreign process between the member snapshot and the system-wide thread walk; any failure to open or query the owner is fail-safe (the thread is left alone).

Spawning into a suspended cgroup group. The freeze is group state: a child spawned or adopted while suspended joins frozen — the forked child joins the cgroup before exec, so it can freeze before completing the spawn handshake and start() may never return until resume. Resume before starting new work; details in Process groups.

Frozen trees and graceful shutdown. Hard kills penetrate a frozen tree (SIGKILL / cgroup.kill / job terminate), but a graceful shutdown leads with a SIGTERM the frozen processes can't handle — it waits out the full grace. Resume first.

pgroup backends: leaders, zombies, pid reuse. members() lists tracked group leaders only; an exited-but-unreaped child (zombie) still probes as alive (keep wait()ing handles if you need prompt liveness, e.g. for shutdown's early return); and pid-based signalling is inherently best-effort against pid reuse — the crate prunes dead entries on every probe to keep the window minimal.


Next: Process groups · Running commands · docs index

Upgrading processkit

Per-version notes for consumers moving their dependency forward: what breaks, who it affects, and the exact change to make. The CHANGELOG is the full record; this page is the "I depend on it, what do I do" view.

Versioning. From 1.0.0 onward processkit follows Semantic Versioning: the public API is stable, and any breaking change lands only in a new major version, so 2.x upgrades are backward-compatible. The default Cargo requirement processkit = "2" already does the right thing — it allows 2.* but not 3.0. Skim the relevant section here before each major bump. (The mock feature's mockall-generated expect_* surface stays semver-exempt — it tracks the mockall version.)

2.0.0 (from 1.x)

The 2.0 line lands a batch of breaking changes accumulated over the 2.0-series minors (2.1.0/2.1.1). Most are caught by the compiler, but two are behavioral rather than structural — output_bytes's new byte cap and RecordReplayRunner's dropped cwd match — and are worth a deliberate read even if your build still passes after the bump.

Error variants are now #[non_exhaustive]

Every data-carrying variant of ErrorExit, Timeout, Signalled, Spawn, NotFound, Parse, OutputTooLarge, Stdin, and ResourceLimit (limits feature) — is now individually #[non_exhaustive]. A struct-literal construction or a field-exhaustive match/destructuring of any of these variants outside the crate no longer compiles.

Affected if you build one of these variants directly (e.g. in a custom ProcessRunner double or cassette replay) or match on one without a trailing ... The symptom is a compiler error like "cannot construct ... due to private fields" or "pattern does not mention field ...".

Fix:

  • Add .. to every field pattern: Error::Exit { program, code, .. } => ....
  • Read fields through the existing accessors instead of destructuring — program(), stdout()/stderr()/combined(), code(), signal(), is_*().
  • Build a variant through the new #[doc(hidden)] (but pub, semver-covered) constructors instead of a struct literal: Error::exit(program, code, stdout, stderr), Error::timeout(program, timeout, stdout, stderr), Error::signalled(program, signal, stdout, stderr), Error::spawn(program, io_error), Error::not_found(program, searched), Error::stdin(program, io_error). Error::parse(program, message) -> Error::Parse is the one constructor left on the documented surface — building an Error::Parse from your own output parser is a normal production path, not just a test-doubling convenience.

Error::Exit/Timeout/Signalled carry the raw stdout bytes

Each of the three gains a field, stdout_bytes: Option<Vec<u8>>. Since all three are #[non_exhaustive] (see above), this alone breaks nothing by itself — it only matters if you were already matching one of these variants; follow the fix above. Read the bytes through the new accessor: Error::stdout_bytes(&self) -> Option<&[u8]>. It's Some only for an error built over output_bytes (e.g. output_bytes().await?.ensure_success()?); every text-path error (output_string/run/checked/…) reports None — the decoded stdout string there is already the complete capture.

Error::ResourceLimit restructured again

(limits feature.) The 1.0 shape, Error::ResourceLimit { message: String }, is now Error::ResourceLimit { kind: LimitKind, reason: LimitReason, detail: String }kind is Memory / Processes / Cpu (which limit), reason is Invalid / Unsupported / Unenforceable (why), and detail is a human-readable message replacing the old message field 1:1.

Fix a match on the old field — Error::ResourceLimit { message } => ... no longer compiles (the field is gone, and the variant is #[non_exhaustive] besides, see above). Either match with .. and read detail, or — better — classify without parsing English at all via the new accessors: Error::limit_kind(&self) -> Option<LimitKind> / Error::limit_reason(&self) -> Option<LimitReason>. reason reflects a real backend signal: Unsupported when the platform has no whole-tree containment mechanism at all, Unenforceable when a capable mechanism exists but this particular request was rejected, Invalid for a nonsensical value caught before the OS is ever touched.

Field and builder renames

BeforeAfterNotes
Error::OutputTooLarge { line_limit, byte_limit, .. }Error::OutputTooLarge { max_lines, max_bytes, .. }matches the OutputBufferPolicy knobs they report
ResourceLimits::memory_maxResourceLimits::max_memorylimits feature; matches max_processes's word order
ProcessGroupOptions::memory_max(bytes)ProcessGroupOptions::max_memory(bytes)same rename, on the builder

ProcessResult::output_contains_any takes any string collection

Was fn output_contains_any(&self, needles: &[&str]) -> bool, now fn output_contains_any(&self, needles: impl IntoIterator<Item = impl AsRef<str>>) -> bool. Backward-compatible for most call sites — a &[&str] slice still works, since &[&str] implements IntoIterator<Item = &&str> and &&str: AsRef<str> — and a bare array literal that previously needed an explicit & (output_contains_any(&["a", "b"])) now also works without it (output_contains_any(["a", "b"])); a Vec<String> or any other IntoIterator<Item = impl AsRef<str>> now passes directly too, where it previously needed collecting into a slice first. These are all additive, not breaking. The one call shape that stops compiling: a single bare &str passed where &[&str] previously coerced through an implicit reborrow — wrap it in an array instead: output_contains_any([needle]).

Removed: two deprecated aliases and a duplicated field

  • ProcessGroup::terminate_all (deprecated since 1.1.0) is removed — use ProcessGroup::kill_all.
  • RunProfile::avg_cpu (deprecated since 1.1.0) is removed — use RunProfile::avg_cpu_cores.
  • RunProfile::exit_code (the field) is removed — it duplicated outcome.code(); use RunProfile::code() -> Option<i32> instead of profile.exit_code.

Encoding and StreamExt move into processkit::prelude

The flat crate-root re-exports of two upstream 0.x dependencies' vocabulary types — Encoding (from encoding_rs) and StreamExt (from tokio-stream) — moved behind a new processkit::prelude module.

Fix — update the use:

BeforeAfter
use processkit::Encoding;use processkit::prelude::Encoding;
use processkit::StreamExt;use processkit::prelude::StreamExt;

This keeps use processkit::* from pulling in a StreamExt trait that collides with futures::StreamExt, and contains a future 0.x major bump of either dependency to the prelude module instead of the whole crate surface.

output_bytes now enforces the OutputBufferPolicy byte cap

output_bytes now honors OutputBufferPolicy::max_bytes on its raw stdout capture — previously the byte ceiling applied only to the line-pumped stderr path, not to output_bytes's own stdout collection. Affected if you set with_max_bytes(..) (or OutputBufferPolicy::bounded(..)/fail_loud(..)) and called output_bytes expecting an unbounded stdout capture regardless: past the cap, OverflowMode::Error now errors with Error::OutputTooLarge (max_lines: None — raw bytes have no lines), and the drop modes now bound the retained bytes to a head/tail and set ProcessResult::truncated. The default (no byte cap, OutputBufferPolicy::unbounded()) is unchanged — an output_bytes call with the default policy still captures everything, same as before.

RecordReplayRunner cassettes no longer match on cwd

(record feature.) A cassette entry recorded from one absolute working directory now replays against the same invocation run from a different one, instead of Error::CassetteMissing — the main blocker to recording on one machine and replaying on another (a dev box vs. CI vs. a teammate's checkout). cwd is still stored on each entry, verbatim, for visibility; it just no longer discriminates two otherwise-identical recorded runs. The on-disk cassette format revision bumped to 3, but this is not a compatibility gate — a cassette written by a previous build still loads and replays fine, unchanged. The only behavior change: if your test suite relied on two entries that differed only in cwd to answer two different invocations, they now collide and the first-recorded entry answers for both — give such entries a distinguishing argument instead.

1.0.0 (from 0.11.x)

A few breaking changes, all caught by the compiler — if it builds after the bump, you're done.

OutputLine.text is now an accessor

OutputLine (the per-line payload of RunningProcess::output_events) no longer exposes text as a public field — read it via line.text() -> &str (or line.into_text() -> String to take ownership). This frees the line representation to evolve. Fix: line.textline.text().

Error::ResourceLimit is now a struct variant

Error::ResourceLimit(String) became Error::ResourceLimit { message: String } (parity with the other rich variants, room for structured detail later). Fix a match Error::ResourceLimit(m)Error::ResourceLimit { message: m }. (Only relevant with the limits feature.)

This was the 1.0 shape only, kept here for historical context. 2.0 restructured the variant again — the current fields are { kind: LimitKind, reason: LimitReason, detail: String }, not { message: String }; see the 2.0.0 section below.

The text-capture verb is renamed outputoutput_string

The verb that runs to completion and returns the full ProcessResult<String> is now spelled output_string on every layer, matching output_bytes (and the spelling Command/Pipeline/RunningProcess already used). Two reasons: the same operation no longer has two names depending on the type, and a bare output clashed with std::process::Command::output, which returns bytes — the explicit name removes that footgun.

Affected if you call ProcessRunner::output, CliClient::output, the free fn processkit::output, or implement a custom ProcessRunner / use MockRunner. The symptom is a build error like "no method named output" / "cannot find function output in crate processkit".

Fix — rename the calls (mechanical):

BeforeAfter
runner.output(&cmd) / client.output(args)runner.output_string(&cmd) / client.output_string(args)
processkit::output(prog, args)processkit::output_string(prog, args)
impl ProcessRunner { async fn output(..) }async fn output_string(..) (the required method)
mock.expect_output()mock.expect_output_string()

output_bytes is unchanged, and Command/Pipeline/RunningProcess callers need no change (those already used output_string).

0.11.0 (from 0.10.x)

Two breaking changes, both small and caught by the compiler — if it builds after the bump, you're done. Plus one internal fix that needs no action.

1. stats is now opt-in — a Cargo.toml change

The default feature set is now just process-control; stats is no longer on by default. (It is the one feature carrying an extra build dependency — the Windows ProcessStatus FFI used solely for the peak-memory readout — and it gates a specialized metrics surface the core never needs.)

Affected if you use any metrics API: ProcessGroup::stats / ProcessGroupStats, RunningProcess::cpu_time / peak_memory_bytes, or RunProfile / RunningProcess::profile. The symptom is a build error like "no method named stats / cpu_time / peak_memory_bytes / profile" or "cannot find type ProcessGroupStats / RunProfile".

Fix — add the feature:

[dependencies]
processkit = { version = "0.11", features = ["stats"] }

If you already enable limits, do nothinglimits still implies stats.

If you don't use metrics: nothing to do. Your default build is now slightly leaner (no Windows ProcessStatus dependency).

2. OutputEvent carries OutputLine — a code change

Affects only callers of RunningProcess::output_events (the ordered lifecycle+output event stream). The per-line payload changed from a bare String to a #[non_exhaustive] OutputLine struct with a public text field.

Before:

#![allow(unused)]
fn main() {
use processkit::OutputEvent;

while let Some(ev) = events.next().await {
    match ev {
        OutputEvent::Stdout(s) => println!("out: {s}"),
        OutputEvent::Stderr(s) => eprintln!("err: {s}"),
        _ => {}
    }
}
}

After — read line.text (in 1.0 this becomes line.text(); see the 1.0.0 section above):

#![allow(unused)]
fn main() {
match ev {
    OutputEvent::Stdout(line) => println!("out: {}", line.text),
    OutputEvent::Stderr(line) => eprintln!("err: {}", line.text),
    _ => {}
}
}

Or, when you don't care which stream produced the line, use the new accessor:

#![allow(unused)]
fn main() {
if let Some(text) = ev.text() {
    println!("{text}");
}
}

OutputLine is #[non_exhaustive]: you receive it from the crate and read its fields — you don't construct it, and a match on it should use ... The change exists to reserve room for per-line metadata (e.g. a timestamp or a monotonic line index) in a later release without another break.

3. Cancel-precedence fix ("Issue 7") — no action

A run that reaps on its own is no longer at risk of being misreported as Err(Cancelled) by a cancellation token that fires in the narrow window between the reap and the disposition check. This is an internal correctness fix with no public-API change. If you carried a workaround that tolerated a spurious Cancelled on a self-completing run, you can remove it.

Verify the upgrade

cargo update -p processkit
cargo build      # both breaking changes are compiler-caught
cargo test

Upgrading from older than 0.10

The jumps below 0.10 predate this guide. Read the dated sections of the CHANGELOG for each minor you cross — every breaking entry there is marked Breaking and carries its own migration note. Notable recent non-breaking additions you gain along the way: Command::checked / run_unit (0.10.2) and the record-cassette symlink/Display-injection hardening (0.10.2).

What's next

‹ docs index

ProcessKit is published for Rust as processkit on crates.io. Its first cross-language binding is also available for Python as processkit-py, with a separate overview on this site, a full guide set, and a source repository.

The longer-term plan is to bring the same approach — kernel-backed whole-tree containment, honest error semantics, and testable seams — to Go, F#, and Kotlin. Each implementation will follow the same philosophy and receive its own language-appropriate documentation.