From 613c9825f4c6267c13d04d1e5d226cd92ab54173 Mon Sep 17 00:00:00 2001 From: Nucleic Date: Wed, 29 Jul 2026 04:03:22 -0700 Subject: [PATCH] Merge nucleic/eager-glass-civet-f3ph into dev --- brush-core/src/interp.rs | 328 ++++++++++++++++++++++++++++++++------- 1 file changed, 274 insertions(+), 54 deletions(-) diff --git a/brush-core/src/interp.rs b/brush-core/src/interp.rs index 9cb0091..8ae8812 100644 --- a/brush-core/src/interp.rs +++ b/brush-core/src/interp.rs @@ -1999,71 +1999,291 @@ fn nash_stage_text(pipeline: &ast::Pipeline, index: usize) -> String { text } -// nash: spawn a copier that moves bytes from `src` (producer's output) to `dst` -// (consumer's input), mirroring a bounded prefix, then reports the link to the -// gate (docs/NASH.md §5.3). The synchronous copy preserves pipe backpressure; -// EOF on `src` or EPIPE on `dst` (consumer gone, e.g. `yes | head`) ends it and -// closes `dst` so the consumer sees EOF. +// nash: bytes moved per userspace copy on the portable tee path. 128 KiB rather than the 32 KiB +// this started with: the same stream costs a quarter of the read/write pairs, and a pipe link's +// tee cost is almost entirely syscalls and the context switches around them +// (docs/NASH_STREAM_PERF_PLAN.md §P2.2). +const NASH_TEE_BUF: usize = 128 * 1024; + +// nash: how large both ends of a tapped link are grown to, where the platform allows it. A tee +// puts *two* pipes where the shell asked for one, so a producer that outruns the default 64 KiB +// buffer pays for it twice; growing them cuts the wakeups on both sides. +#[cfg(target_os = "linux")] +const NASH_TEE_PIPE_SIZE: libc::c_int = 256 * 1024; + +// nash: bytes asked of one `splice` call. The kernel moves at most a pipe-buffer's worth per call +// regardless, so this only has to be comfortably larger than the pipe. +#[cfg(target_os = "linux")] +const NASH_TEE_SPLICE_CHUNK: usize = 1024 * 1024; + +// nash: how many finished tee threads stay parked waiting for the next pipeline. Thread creation +// is the bulk of the ~0.13 ms a tapped pipeline costs to *set up*, and shell-heavy builds run +// pipelines in the thousands (docs/NASH_STREAM_PERF_PLAN.md §P2.3). A shell runs its pipelines +// mostly one at a time, so a small pool covers the common case; a burst of concurrent pipelines +// simply spawns beyond it and lets the extra workers retire when they finish. +const NASH_TEE_POOL_MAX_IDLE: usize = 4; + +// nash: stack for a tee worker. It copies through a heap buffer and hands the gate an already-built +// event, so its own frames are shallow — a quarter of the 2 MiB default is ample and maps less. +const NASH_TEE_STACK: usize = 512 * 1024; + +// nash: one link waiting to be copied — what a tee worker is handed. +struct NashTeeJob { + src: std::io::PipeReader, + dst: std::io::PipeWriter, + pipeline_id: u64, + from_index: usize, + from_text: String, + to_text: String, +} + +// nash: senders of the tee workers currently parked, newest last. +static NASH_TEE_POOL: std::sync::Mutex>> = + std::sync::Mutex::new(Vec::new()); + +// nash: the process the parked workers belong to. A `fork` copies the pool's *memory* but none of +// its threads, so a child that inherited it would hand jobs to workers that do not exist — and a +// dropped job silently breaks the pipeline it was tapping. Checked (and reset) on every dispatch. +static NASH_TEE_POOL_PID: std::sync::atomic::AtomicU32 = std::sync::atomic::AtomicU32::new(0); + +// nash: take a parked worker, if this process has one to spare. +// +// `try_lock`, never `lock`: this runs on the shell's own pipeline-setup path, so waiting on the +// pool would put observation in front of a command — and after a `fork` the mutex can be left +// locked by a thread that no longer exists. Failing to take a worker just means spawning one. +fn nash_tee_take_idle() -> Option> { + let mut pool = NASH_TEE_POOL.try_lock().ok()?; + let pid = std::process::id(); + if NASH_TEE_POOL_PID.swap(pid, std::sync::atomic::Ordering::Relaxed) != pid { + pool.clear(); // inherited across a fork: those workers are not in this process + } + pool.pop() +} + +// nash: park a worker that just finished a link. Returns whether it was taken — a worker the pool +// has no room for (or cannot reach) simply exits. +fn nash_tee_park(tx: &std::sync::mpsc::Sender) -> bool { + let Ok(mut pool) = NASH_TEE_POOL.try_lock() else { + return false; + }; + if pool.len() >= NASH_TEE_POOL_MAX_IDLE { + return false; + } + pool.push(tx.clone()); + true +} + +// nash: a parked worker's loop — copy a link, park again, wait for the next one. +fn nash_tee_worker( + rx: &std::sync::mpsc::Receiver, + parked: &std::sync::mpsc::Sender, +) { + while let Ok(job) = rx.recv() { + nash_run_pipe_tee(job); + if !nash_tee_park(parked) { + return; + } + } +} + +// nash: hand a copier the link `src` (producer's output) → `dst` (consumer's input): it mirrors a +// bounded prefix, then reports the link to the gate (docs/NASH.md §5.3). The synchronous copy +// preserves pipe backpressure; EOF on `src` or EPIPE on `dst` (consumer gone, e.g. `yes | head`) +// ends it and closes `dst` so the consumer sees EOF. // `to_index` is always `from_index + 1` — a link joins adjacent stages — so it is derived rather // than passed; `from_text`/`to_text` are the two stages already rendered by `nash_stage_text`. fn nash_spawn_pipe_tee( - mut src: std::io::PipeReader, - mut dst: std::io::PipeWriter, + src: std::io::PipeReader, + dst: std::io::PipeWriter, pipeline_id: u64, from_index: usize, from_text: String, to_text: String, ) { - std::thread::Builder::new() + let mut job = NashTeeJob { + src, + dst, + pipeline_id, + from_index, + from_text, + to_text, + }; + // A worker that retired between being parked and being handed this job returns it, so try the + // next one down rather than losing the link. + while let Some(tx) = nash_tee_take_idle() { + match tx.send(job) { + Ok(()) => return, + Err(std::sync::mpsc::SendError(returned)) => job = returned, + } + } + let (tx, rx) = std::sync::mpsc::channel::(); + let parked = tx.clone(); + // If the spawn fails there is nobody to copy this link and the job is dropped, closing both + // ends — the same outcome an unspawnable tee thread has always had, in a shell that is out of + // threads either way. + if std::thread::Builder::new() .name("nash-pipe-tee".into()) - .spawn(move || { - use std::io::{Read, Write}; - let mut buf = [0u8; 32 * 1024]; - let mut captured: Vec = Vec::new(); - let mut total: u64 = 0; - let mut truncated = false; - loop { - let n = match src.read(&mut buf) { - Ok(0) => break, // producer closed → done - Ok(n) => n, - Err(ref e) if e.kind() == std::io::ErrorKind::Interrupted => continue, - Err(_) => break, - }; - total += n as u64; - if captured.len() < NASH_PIPE_CAPTURE_CAP { - let room = NASH_PIPE_CAPTURE_CAP - captured.len(); - let take = room.min(n); - captured.extend_from_slice(&buf[..take]); - if take < n { - truncated = true; - } - } else { - truncated = true; - } - // Pass through; if the consumer went away, stop (its EOF is enough). - if dst.write_all(&buf[..n]).is_err() { - break; - } - } - // Report BEFORE closing `dst`: dropping it unblocks the consumer, which - // lets the shell reach its exit-flush — so the event must already be - // enqueued to avoid a lost-event race on short-lived shells. - crate::gate::gate().on_pipe(crate::gate::PipeEvent { - pipeline_id, - from_index, - to_index: from_index + 1, - from_text, - to_text, - total_bytes: total, - captured, - truncated, - }); - // Now close the write end so the consumer sees EOF. - drop(dst); - }) - .ok(); + .stack_size(NASH_TEE_STACK) + .spawn(move || nash_tee_worker(&rx, &parked)) + .is_ok() + { + let _ = tx.send(job); + } } +// nash: copy one link through to its consumer, then report it. +fn nash_run_pipe_tee(job: NashTeeJob) { + let NashTeeJob { + mut src, + mut dst, + pipeline_id, + from_index, + from_text, + to_text, + } = job; + nash_grow_pipe(&src); + nash_grow_pipe(&dst); + + let mut captured: Vec = Vec::new(); + let mut total: u64 = 0; + // Only the prefix is ever copied through userspace; the rest of the stream — which is all of + // it, for the transfers that actually cost something — is handed to the kernel. + if nash_tee_prefix(&mut src, &mut dst, &mut captured, &mut total) { + nash_tee_passthrough(&mut src, &mut dst, &mut total); + } + + // Report BEFORE closing `dst`: dropping it unblocks the consumer, which + // lets the shell reach its exit-flush — so the event must already be + // enqueued to avoid a lost-event race on short-lived shells. + crate::gate::gate().on_pipe(crate::gate::PipeEvent { + pipeline_id, + from_index, + to_index: from_index + 1, + from_text, + to_text, + total_bytes: total, + truncated: total > captured.len() as u64, + captured, + }); + // Now close the write end so the consumer sees EOF. + drop(dst); +} + +// nash: copy until the capture cap is reached, mirroring what goes past. Returns whether the link +// is still open (false on EOF, or once either end gives up). +fn nash_tee_prefix( + src: &mut std::io::PipeReader, + dst: &mut std::io::PipeWriter, + captured: &mut Vec, + total: &mut u64, +) -> bool { + use std::io::{Read, Write}; + let mut buf = vec![0u8; NASH_TEE_BUF]; + while captured.len() < NASH_PIPE_CAPTURE_CAP { + let n = match src.read(&mut buf) { + Ok(0) => return false, // producer closed → done + Ok(n) => n, + Err(ref e) if e.kind() == std::io::ErrorKind::Interrupted => continue, + Err(_) => return false, + }; + *total += n as u64; + let room = NASH_PIPE_CAPTURE_CAP - captured.len(); + captured.extend_from_slice(&buf[..room.min(n)]); + // Pass through; if the consumer went away, stop (its EOF is enough). + if dst.write_all(&buf[..n]).is_err() { + return false; + } + } + true +} + +// nash: the portable passthrough — the same read/write loop as the prefix phase, minus the +// capture. Used on platforms without `splice`, and as the fallback when the kernel refuses it. +fn nash_tee_copy(src: &mut std::io::PipeReader, dst: &mut std::io::PipeWriter, total: &mut u64) { + use std::io::{Read, Write}; + let mut buf = vec![0u8; NASH_TEE_BUF]; + loop { + let n = match src.read(&mut buf) { + Ok(0) => return, + Ok(n) => n, + Err(ref e) if e.kind() == std::io::ErrorKind::Interrupted => continue, + Err(_) => return, + }; + *total += n as u64; + if dst.write_all(&buf[..n]).is_err() { + return; + } + } +} + +// nash: move the rest of the link pipe-to-pipe inside the kernel (docs/NASH_STREAM_PERF_PLAN.md +// §P2.1). Past the captured prefix the tee has nothing to look at, so there is no reason for the +// bytes to enter this process at all — and the byte count comes back from the same call that +// moves them. +#[cfg(target_os = "linux")] +fn nash_tee_passthrough( + src: &mut std::io::PipeReader, + dst: &mut std::io::PipeWriter, + total: &mut u64, +) { + use std::os::fd::AsRawFd; + let (in_fd, out_fd) = (src.as_raw_fd(), dst.as_raw_fd()); + loop { + // SAFETY: both descriptors are pipes owned by this job for the duration of the call, and + // the offsets are null because pipes have none. + let moved = unsafe { + libc::splice( + in_fd, + std::ptr::null_mut(), + out_fd, + std::ptr::null_mut(), + NASH_TEE_SPLICE_CHUNK, + libc::SPLICE_F_MOVE, + ) + }; + if moved == 0 { + return; // producer closed → done + } + if moved > 0 { + *total += u64::try_from(moved).unwrap_or(0); + continue; + } + let err = std::io::Error::last_os_error(); + match err.raw_os_error() { + Some(libc::EINTR) => (), + // A kernel (or a seccomp policy) that will not splice these descriptors: finish the + // link the portable way rather than dropping it. + Some(libc::EINVAL | libc::ENOSYS | libc::EPERM) => { + nash_tee_copy(src, dst, total); + return; + } + // Consumer gone (EPIPE) or a link that broke: its EOF is enough. + _ => return, + } + } +} + +#[cfg(not(target_os = "linux"))] +fn nash_tee_passthrough( + src: &mut std::io::PipeReader, + dst: &mut std::io::PipeWriter, + total: &mut u64, +) { + nash_tee_copy(src, dst, total); +} + +// nash: grow a tapped pipe's buffer, best-effort. The kernel caps unprivileged growth at +// `/proc/sys/fs/pipe-max-size` and simply refuses anything larger, which costs one failed syscall +// per link and leaves the default in place. +#[cfg(target_os = "linux")] +fn nash_grow_pipe(pipe: &F) { + // SAFETY: `fcntl` with `F_SETPIPE_SZ` takes an int and only reads the descriptor's pipe size. + let _ = unsafe { libc::fcntl(pipe.as_raw_fd(), libc::F_SETPIPE_SZ, NASH_TEE_PIPE_SIZE) }; +} + +#[cfg(not(target_os = "linux"))] +fn nash_grow_pipe(_pipe: &F) {} + // nash: record a regular-file redirect for post-run read-back (docs/NASH.md §5.2a). // The child still gets the real fd; this only notes what to read back afterward. fn nash_record_file_redirect(