Merge nucleic/sleek-thistle-egret-fyej into dev
This commit is contained in:
@@ -8,6 +8,7 @@
|
||||
//! any further changes to this crate.
|
||||
|
||||
use std::path::PathBuf;
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::sync::OnceLock;
|
||||
|
||||
/// Details about a simple command about to be executed (builtin, function, or
|
||||
@@ -62,6 +63,24 @@ pub struct ExecEnd {
|
||||
pub exit_code: i32,
|
||||
}
|
||||
|
||||
/// Data observed flowing across one pipe link `a | b` (docs/NASH.md §5.3).
|
||||
/// Reported once the link drains (producer closed or consumer went away).
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct PipeEvent {
|
||||
/// Correlates the links of one pipeline invocation.
|
||||
pub pipeline_id: u64,
|
||||
/// Zero-based index of the producing stage.
|
||||
pub from_index: usize,
|
||||
/// Zero-based index of the consuming stage.
|
||||
pub to_index: usize,
|
||||
/// Total bytes that crossed the link.
|
||||
pub total_bytes: u64,
|
||||
/// A bounded prefix of the bytes (observer applies its own caps/redaction).
|
||||
pub captured: Vec<u8>,
|
||||
/// Whether `total_bytes` exceeded `captured`.
|
||||
pub truncated: bool,
|
||||
}
|
||||
|
||||
/// The gate's decision for a command about to execute.
|
||||
pub enum Verdict {
|
||||
/// Run the command; the token is echoed back via `on_exit`.
|
||||
@@ -87,6 +106,9 @@ pub trait Gate: Send + Sync {
|
||||
/// nash: called with the (trimmed) output of a command substitution
|
||||
/// `$(…)` / backticks (docs/NASH.md §5.3). Default no-op.
|
||||
fn on_cmdsub(&self, _output: &str) {}
|
||||
/// nash: called once a pipe link `a | b` drains (docs/NASH.md §5.3).
|
||||
/// Invoked from the tee copier thread. Default no-op.
|
||||
fn on_pipe(&self, _ev: PipeEvent) {}
|
||||
}
|
||||
|
||||
struct AllowAll;
|
||||
@@ -111,3 +133,16 @@ pub fn set_gate(gate: Box<dyn Gate>) {
|
||||
pub(crate) fn gate() -> &'static dyn Gate {
|
||||
GATE.get().map_or(&ALLOW_ALL as &dyn Gate, AsRef::as_ref)
|
||||
}
|
||||
|
||||
/// Whether a real (non-default) gate is installed. Pipe teeing — which adds a
|
||||
/// thread + extra pipe per link — is only done when this is true (docs/NASH.md §5.3).
|
||||
pub(crate) fn is_active() -> bool {
|
||||
GATE.get().is_some()
|
||||
}
|
||||
|
||||
static PIPELINE_COUNTER: AtomicU64 = AtomicU64::new(1);
|
||||
|
||||
/// A fresh id correlating the links of one pipeline invocation.
|
||||
pub(crate) fn next_pipeline_id() -> u64 {
|
||||
PIPELINE_COUNTER.fetch_add(1, Ordering::Relaxed)
|
||||
}
|
||||
|
||||
@@ -460,7 +460,28 @@ async fn spawn_pipeline_processes(
|
||||
// Create pipes to use between commands, but only bother doing so if there's more than one
|
||||
// command.
|
||||
if pipeline_len > 1 {
|
||||
for _ in 0..(pipeline_len - 1) {
|
||||
// nash: when observation is active, interpose a tee on each link so the
|
||||
// bytes flowing `a | b` are captured (docs/NASH.md §5.3). Off otherwise —
|
||||
// no extra pipe/thread on the unobserved fast path.
|
||||
let nash_pipeline_id = if crate::gate::is_active() {
|
||||
Some(crate::gate::next_pipeline_id())
|
||||
} else {
|
||||
None
|
||||
};
|
||||
for link in 0..(pipeline_len - 1) {
|
||||
if let Some(pipeline_id) = nash_pipeline_id {
|
||||
// Producer writes `a_write`; a copier moves bytes to `b_write`;
|
||||
// consumer reads `b_read`. Slots stay identical to the untapped
|
||||
// case, so the pop-ordering below is unchanged.
|
||||
let (a_read, a_write) = std::io::pipe()?;
|
||||
let (b_read, b_write) = std::io::pipe()?;
|
||||
// Link `link` connects producer stage (len-2-link) → consumer (len-1-link).
|
||||
let from_index = pipeline_len - 2 - link;
|
||||
nash_spawn_pipe_tee(a_read, b_write, pipeline_id, from_index, from_index + 1);
|
||||
pipe_readers.push(Some(openfiles::OpenFile::PipeReader(b_read)));
|
||||
pipe_writers.push(Some(openfiles::OpenFile::PipeWriter(a_write)));
|
||||
continue;
|
||||
}
|
||||
let (reader, writer) = std::io::pipe()?;
|
||||
pipe_readers.push(Some(reader.into()));
|
||||
pipe_writers.push(Some(writer.into()));
|
||||
@@ -1888,6 +1909,69 @@ pub(crate) async fn setup_redirect(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// nash: bounded prefix captured per pipe link before reporting (docs/NASH.md §5.3).
|
||||
const NASH_PIPE_CAPTURE_CAP: usize = 64 * 1024;
|
||||
|
||||
// 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.
|
||||
fn nash_spawn_pipe_tee(
|
||||
mut src: std::io::PipeReader,
|
||||
mut dst: std::io::PipeWriter,
|
||||
pipeline_id: u64,
|
||||
from_index: usize,
|
||||
to_index: usize,
|
||||
) {
|
||||
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<u8> = 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,
|
||||
total_bytes: total,
|
||||
captured,
|
||||
truncated,
|
||||
});
|
||||
// Now close the write end so the consumer sees EOF.
|
||||
drop(dst);
|
||||
})
|
||||
.ok();
|
||||
}
|
||||
|
||||
// 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(
|
||||
|
||||
Reference in New Issue
Block a user