Merge nucleic/eager-glass-civet-f3ph into dev
This commit is contained in:
+274
-54
@@ -1999,71 +1999,291 @@ fn nash_stage_text(pipeline: &ast::Pipeline, index: usize) -> String {
|
|||||||
text
|
text
|
||||||
}
|
}
|
||||||
|
|
||||||
// nash: spawn a copier that moves bytes from `src` (producer's output) to `dst`
|
// nash: bytes moved per userspace copy on the portable tee path. 128 KiB rather than the 32 KiB
|
||||||
// (consumer's input), mirroring a bounded prefix, then reports the link to the
|
// this started with: the same stream costs a quarter of the read/write pairs, and a pipe link's
|
||||||
// gate (docs/NASH.md §5.3). The synchronous copy preserves pipe backpressure;
|
// tee cost is almost entirely syscalls and the context switches around them
|
||||||
// EOF on `src` or EPIPE on `dst` (consumer gone, e.g. `yes | head`) ends it and
|
// (docs/NASH_STREAM_PERF_PLAN.md §P2.2).
|
||||||
// closes `dst` so the consumer sees EOF.
|
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<Vec<std::sync::mpsc::Sender<NashTeeJob>>> =
|
||||||
|
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<std::sync::mpsc::Sender<NashTeeJob>> {
|
||||||
|
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<NashTeeJob>) -> 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<NashTeeJob>,
|
||||||
|
parked: &std::sync::mpsc::Sender<NashTeeJob>,
|
||||||
|
) {
|
||||||
|
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
|
// `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`.
|
// than passed; `from_text`/`to_text` are the two stages already rendered by `nash_stage_text`.
|
||||||
fn nash_spawn_pipe_tee(
|
fn nash_spawn_pipe_tee(
|
||||||
mut src: std::io::PipeReader,
|
src: std::io::PipeReader,
|
||||||
mut dst: std::io::PipeWriter,
|
dst: std::io::PipeWriter,
|
||||||
pipeline_id: u64,
|
pipeline_id: u64,
|
||||||
from_index: usize,
|
from_index: usize,
|
||||||
from_text: String,
|
from_text: String,
|
||||||
to_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::<NashTeeJob>();
|
||||||
|
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())
|
.name("nash-pipe-tee".into())
|
||||||
.spawn(move || {
|
.stack_size(NASH_TEE_STACK)
|
||||||
use std::io::{Read, Write};
|
.spawn(move || nash_tee_worker(&rx, &parked))
|
||||||
let mut buf = [0u8; 32 * 1024];
|
.is_ok()
|
||||||
let mut captured: Vec<u8> = Vec::new();
|
{
|
||||||
let mut total: u64 = 0;
|
let _ = tx.send(job);
|
||||||
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();
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 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<u8> = 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<u8>,
|
||||||
|
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<F: std::os::fd::AsRawFd>(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<F>(_pipe: &F) {}
|
||||||
|
|
||||||
// nash: record a regular-file redirect for post-run read-back (docs/NASH.md §5.2a).
|
// 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.
|
// The child still gets the real fd; this only notes what to read back afterward.
|
||||||
fn nash_record_file_redirect(
|
fn nash_record_file_redirect(
|
||||||
|
|||||||
Reference in New Issue
Block a user