Merge nucleic/eager-glass-civet-f3ph into dev

This commit is contained in:
2026-07-29 04:03:22 -07:00
parent 85e458ef04
commit 2cbff0667d
3 changed files with 897 additions and 163 deletions
+601 -118
View File
@@ -18,6 +18,7 @@ use std::sync::mpsc;
use std::sync::{Mutex, OnceLock};
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use serde::ser::SerializeMap;
use serde::Serialize;
const FLUSH_MAX_EVENTS: usize = 200;
@@ -35,11 +36,144 @@ const IO_TIMEOUT: Duration = Duration::from_millis(1500);
const EXIT_FLUSH_TIMEOUT: Duration = Duration::from_millis(2000);
/// Per-redirection / per-cmdsub preview cap in bytes (docs/NASH.md §5.4).
const PREVIEW_CAP: usize = 64 * 1024;
/// Captured payload one batch may carry, in raw bytes (docs/NASH.md §5.4's per-batch cap).
/// Beyond it the events keep their counts and hashes and lose their previews — a pipeline storm
/// can otherwise put 200 × 64 KiB of preview in one body, which the host then has to parse,
/// base64-decode and hold (docs/NASH_STREAM_PERF_PLAN.md §P4.2).
const BATCH_PREVIEW_BUDGET: usize = 1024 * 1024;
// ---------------------------------------------------------------------------
// Event model (docs/NASH.md §6.1)
// ---------------------------------------------------------------------------
/// The captured half of a data-flow event (`redirect` / `cmdsub` / `pipe`) — everything the wire
/// calls `bytes`, `truncated`, `hash` and `previewB64`.
///
/// It exists so that **none of that encoding happens where the bytes were captured**
/// (docs/NASH_STREAM_PERF_PLAN.md §P3). The interpreter thread, the tee thread and the command's
/// own exit path hand over raw bytes; redaction, base64 and the FNV hash run inside
/// [`Serialize`] — which is only ever called from the flusher thread, off every path a command
/// waits on. The wire shape is unchanged.
enum Capture {
/// Nothing was captured: a special file (`/dev/null`, a tty, a socket), or a target that
/// could not be read. Counts as zero, and carries no hash — an empty hash is how the host
/// tells "no payload" from "a payload that happened to be empty".
None,
/// A file range that has not been read back yet; the flusher opens and reads it
/// (docs/NASH_STREAM_PERF_PLAN.md §P3.2). `offset`/`total` are fixed at the moment the
/// command exited, so a later command appending to the same file cannot widen this event's
/// range — only the *reading* is deferred, never the measurement.
Deferred {
path: PathBuf,
offset: u64,
total: u64,
},
/// Bytes in hand, raw and already capped to [`PREVIEW_CAP`].
Bytes {
bytes: u64,
truncated: bool,
data: Vec<u8>,
},
/// A preview dropped by the per-batch budget (docs/NASH.md §5.4): counts and hash survive,
/// and `truncated` is true because the preview no longer stands for the bytes.
Counted { bytes: u64, hash: String },
}
impl Capture {
/// Bytes in hand, capped to the preview cap.
fn from_bytes(total: u64, data: &[u8]) -> Self {
let capped = data.len().min(PREVIEW_CAP);
Self::Bytes {
bytes: total,
truncated: total > capped as u64,
data: data[..capped].to_vec(),
}
}
/// The payload this will preview, in raw bytes — what the per-batch budget is spent on.
const fn preview_len(&self) -> usize {
match self {
Self::Bytes { data, .. } => data.len(),
// The budget is spent after the deferred reads are resolved, so this arm is only
// reached for a range that could not be read; charge its bounded worst case anyway
// rather than letting an unresolved capture spend nothing.
Self::Deferred { total, .. } => {
if *total > PREVIEW_CAP as u64 {
PREVIEW_CAP
} else {
*total as usize
}
}
Self::None | Self::Counted { .. } => 0,
}
}
/// Read back a deferred range (flusher thread only). A range that has since become
/// unreadable degrades to [`Capture::None`] — the event still reports what it observed.
fn resolve(&mut self) {
let Self::Deferred {
path,
offset,
total,
} = self
else {
return;
};
*self = read_range(path, *offset, *total).map_or(Self::None, |(total, data)| Self::Bytes {
bytes: total,
truncated: total > data.len() as u64,
data,
});
}
/// Drop the preview, keeping the counts and the hash (the per-batch budget).
fn count_only(&mut self) {
if let Self::Bytes { bytes, data, .. } = self {
*self = Self::Counted {
bytes: *bytes,
hash: fnv1a_hex(data),
};
}
}
}
impl Serialize for Capture {
/// Flattened into its event, so the four fields sit where they always have. This is where the
/// hashing, redaction and base64 of every preview nash sends actually happen.
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
let mut map = serializer.serialize_map(Some(4))?;
match self {
// A range left unresolved would mean the flusher skipped it; report it as the
// metadata-only event it effectively is rather than inventing a payload.
Self::None | Self::Deferred { .. } => {
map.serialize_entry("bytes", &0u64)?;
map.serialize_entry("truncated", &false)?;
map.serialize_entry("hash", "")?;
map.serialize_entry("previewB64", "")?;
}
Self::Bytes {
bytes,
truncated,
data,
} => {
map.serialize_entry("bytes", bytes)?;
map.serialize_entry("truncated", truncated)?;
// Over the captured prefix: the full stream is never buffered for a pipe, and this
// is a did-these-bytes-repeat marker rather than an integrity claim (NASH.md §5.2).
map.serialize_entry("hash", &fnv1a_hex(data))?;
map.serialize_entry("previewB64", &preview_b64(data))?;
}
Self::Counted { bytes, hash } => {
map.serialize_entry("bytes", bytes)?;
map.serialize_entry("truncated", &true)?;
map.serialize_entry("hash", hash)?;
map.serialize_entry("previewB64", "")?;
}
}
map.end()
}
}
#[derive(Serialize)]
#[serde(tag = "kind", rename_all = "camelCase")]
enum Event {
@@ -104,19 +238,17 @@ enum Event {
op: String,
fd: u32,
target: Option<String>,
bytes: u64,
truncated: bool,
hash: String,
preview_b64: String,
/// Usually [`Capture::Deferred`] when it goes into the channel and [`Capture::Bytes`] by
/// the time it is serialized — the file read happens on the flusher (§P3.2).
#[serde(flatten)]
capture: Capture,
},
#[serde(rename = "cmdsub", rename_all = "camelCase")]
Cmdsub {
seq: u64,
ts: u64,
bytes: u64,
truncated: bool,
hash: String,
preview_b64: String,
#[serde(flatten)]
capture: Capture,
},
#[serde(rename = "pipe", rename_all = "camelCase")]
Pipe {
@@ -128,12 +260,15 @@ enum Event {
/// The commands either side of this link, as written — brush-core renders them from the
/// AST at spawn, so they name a compound stage correctly and arrive with the very first
/// link rather than waiting on the stages to exit.
///
/// Stage text is source, so it can carry a literal credential the same way argv can;
/// masking runs at serialization time like every other outbound string (NASH.md §5.4).
#[serde(serialize_with = "serialize_redacted")]
from_text: String,
#[serde(serialize_with = "serialize_redacted")]
to_text: String,
bytes: u64,
truncated: bool,
hash: String,
preview_b64: String,
#[serde(flatten)]
capture: Capture,
},
#[serde(rename = "dropped", rename_all = "camelCase")]
Dropped { seq: u64, ts: u64, count: u64 },
@@ -406,26 +541,48 @@ fn looks_secret(word: &str) -> bool {
false
}
/// The captured (bytes-total, capped-preview, truncated) for one redirect record.
fn read_back(record: &brush_core::gate::RedirectRecord) -> Option<(u64, Vec<u8>, bool)> {
/// What one redirect record captured, measured on the command's own path but **not yet read**
/// (docs/NASH_STREAM_PERF_PLAN.md §P3.2).
///
/// The split is deliberate. The *range* — is this a regular file at all, and which bytes did this
/// command put there — is only true at the moment the command exited: `echo a > f; echo b >> f`
/// would otherwise have the truncating write report both lines, because by the time a batch
/// flushes the file has grown. So the exit path pays one `stat` and nothing else; opening the
/// file, reading up to 64 KiB of it, hashing, redacting and base64-encoding all move to the
/// flusher, which is where the bytes were always going anyway.
fn measure_readback(record: &brush_core::gate::RedirectRecord) -> Capture {
use brush_core::gate::RedirectReadback;
if let Some(inline) = &record.inline {
// Heredoc / here-string: the body was known at setup, so there is nothing to read.
let full = inline.as_bytes();
let capped = full.len().min(PREVIEW_CAP);
return Some((full.len() as u64, full[..capped].to_vec(), full.len() > capped));
return Capture::from_bytes(full.len() as u64, full);
}
let path = record.path.as_ref()?;
let Some(path) = record.path.as_ref() else {
return Capture::None;
};
// Only read back regular files; skip /dev/null, ttys, fifos, sockets, devices.
let meta = std::fs::metadata(path).ok()?;
let Ok(meta) = std::fs::metadata(path) else {
return Capture::None;
};
if !meta.is_file() {
return None;
return Capture::None;
}
let size = meta.len();
let (offset, total) = match record.readback {
RedirectReadback::Append => (record.size_before, size.saturating_sub(record.size_before)),
RedirectReadback::Truncate | RedirectReadback::Input => (0, size),
RedirectReadback::Inline => return None,
RedirectReadback::Inline => return Capture::None,
};
Capture::Deferred {
path: path.clone(),
offset,
total,
}
}
/// Read a measured range back, capped to [`PREVIEW_CAP`]. Returns the range's full length (which
/// may exceed what was read) and the bytes. Runs on the flusher thread only.
fn read_range(path: &Path, offset: u64, total: u64) -> Option<(u64, Vec<u8>)> {
let mut file = std::fs::File::open(path).ok()?;
if offset > 0 {
use std::io::Seek;
@@ -435,7 +592,7 @@ fn read_back(record: &brush_core::gate::RedirectRecord) -> Option<(u64, Vec<u8>,
let mut buf = vec![0u8; want];
let n = std::io::Read::read(&mut file, &mut buf).unwrap_or(0);
buf.truncate(n);
Some((total, buf, total > n as u64))
Some((total, buf))
}
/// Resolve a redirection target against the directory its command ran in, so the
@@ -471,11 +628,24 @@ fn preview_b64(data: &[u8]) -> String {
let redacted = redact(&String::from_utf8_lossy(data));
b64(redacted.as_bytes())
} else {
let hex: String = data.iter().take(1024).map(|b| format!("{b:02x}")).collect();
b64(format!("<binary> {hex}").as_bytes())
let mut hex = String::with_capacity(1024 * 2 + 9);
hex.push_str("<binary> ");
for byte in data.iter().take(1024) {
hex.push(HEX[usize::from(byte >> 4)] as char);
hex.push(HEX[usize::from(byte & 0xf)] as char);
}
b64(hex.as_bytes())
}
}
const HEX: &[u8; 16] = b"0123456789abcdef";
/// Mask a string on its way out (docs/NASH.md §5.4). Used where redaction is deferred to
/// serialization — the flusher thread — rather than done where the string was captured.
fn serialize_redacted<S: serde::Serializer>(text: &str, serializer: S) -> Result<S::Ok, S::Error> {
serializer.serialize_str(&redact(text))
}
// ---------------------------------------------------------------------------
// Gate implementation (allow-all; docs/NASH.md §4.2)
// ---------------------------------------------------------------------------
@@ -570,7 +740,8 @@ impl brush_core::gate::Gate for RecordingGate {
pipeline: pending.pipeline.map(PipelineRef::from),
});
// Data-flow read-back for this command's redirects (docs/NASH.md §5.2).
// Data-flow read-back for this command's redirects (docs/NASH.md §5.2). Only the
// *range* is settled here (one stat); the read and the encoding ride to the flusher.
for record in &pending.redirects {
// The host resolves a redirect target against the session's worktree to decide
// which file was edited (and which lock that is), so report the path absolutely:
@@ -580,22 +751,6 @@ impl brush_core::gate::Gate for RecordingGate {
.path
.as_ref()
.map(|p| absolute_path(p, &pending.cwd).to_string_lossy().into_owned());
let Some((total, data, truncated)) = read_back(record) else {
// Special file / unreadable — emit metadata only.
obs.send(Event::Redirect {
seq: obs.seq(),
ts,
cmd_seq: exec_seq,
op: record.op.clone(),
fd: record.fd,
target,
bytes: 0,
truncated: false,
hash: String::new(),
preview_b64: String::new(),
});
continue;
};
obs.send(Event::Redirect {
seq: obs.seq(),
ts,
@@ -603,10 +758,7 @@ impl brush_core::gate::Gate for RecordingGate {
op: record.op.clone(),
fd: record.fd,
target,
bytes: total,
truncated,
hash: fnv1a_hex(&data),
preview_b64: preview_b64(&data),
capture: measure_readback(record),
});
}
}));
@@ -616,14 +768,13 @@ impl brush_core::gate::Gate for RecordingGate {
let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let Some(obs) = OBSERVER.get() else { return };
let full = output.as_bytes();
let capped = full.len().min(PREVIEW_CAP);
// Copies the capped prefix and nothing else: this runs on the interpreter thread, in
// the middle of the expansion the shell is waiting on, and `$(cat big-file)` is a
// substitution whose *whole* result used to be hashed here (§P3.1).
obs.send(Event::Cmdsub {
seq: obs.seq(),
ts: now_millis(),
bytes: full.len() as u64,
truncated: full.len() > capped,
hash: fnv1a_hex(full),
preview_b64: preview_b64(&full[..capped]),
capture: Capture::from_bytes(full.len() as u64, full),
});
}));
}
@@ -632,21 +783,21 @@ impl brush_core::gate::Gate for RecordingGate {
let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let Some(obs) = OBSERVER.get() else { return };
let capped = ev.captured.len().min(PREVIEW_CAP);
let mut data = ev.captured;
data.truncate(capped);
obs.send(Event::Pipe {
seq: obs.seq(),
ts: now_millis(),
pipeline_id: ev.pipeline_id,
from_index: ev.from_index as u64,
to_index: ev.to_index as u64,
// Stage text is source, so it can carry a literal credential the same way argv
// can — every outbound string goes through the same masking (docs/NASH.md §5.4).
from_text: redact(&ev.from_text),
to_text: redact(&ev.to_text),
bytes: ev.total_bytes,
truncated: ev.truncated || ev.captured.len() > capped,
// Hash over the captured prefix (the full stream isn't buffered).
hash: fnv1a_hex(&ev.captured),
preview_b64: preview_b64(&ev.captured[..capped]),
from_text: ev.from_text,
to_text: ev.to_text,
capture: Capture::Bytes {
bytes: ev.total_bytes,
truncated: ev.truncated || ev.total_bytes > capped as u64,
data,
},
});
}));
}
@@ -660,8 +811,12 @@ impl brush_core::gate::Gate for RecordingGate {
static UNAUTHORIZED: AtomicBool = AtomicBool::new(false);
fn http_request(path: &str, token: Option<&str>, body: &[u8]) -> Vec<u8> {
// Keep-alive, deliberately (docs/NASH_STREAM_PERF_PLAN.md §P4.1): a shell posts a batch every
// 500 ms for as long as it lives, and `Connection: close` made every one of them a fresh
// connect + accept + teardown on both sides. The host's route already answers with
// `Connection: keep-alive` and a `Content-Length`, which is what makes reuse framable.
let mut req = format!(
"POST {path} HTTP/1.1\r\nHost: nucleic\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n",
"POST {path} HTTP/1.1\r\nHost: nucleic\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: keep-alive\r\n",
body.len()
);
if let Some(token) = token {
@@ -673,55 +828,185 @@ fn http_request(path: &str, token: Option<&str>, body: &[u8]) -> Vec<u8> {
bytes
}
fn read_status<R: Read>(mut stream: R) -> Option<u16> {
let mut buf = [0u8; 64];
let n = stream.read(&mut buf).ok()?;
let text = String::from_utf8_lossy(&buf[..n]);
let status = text.split_whitespace().nth(1)?;
status.parse().ok()
/// The outcome of one post on a live connection.
enum Posted {
/// The host took the batch and the connection can carry another.
Ok,
/// The host took the batch but the connection cannot be reused (no `Content-Length` to frame
/// the next response by, or the host asked to close).
OkAndClose,
/// The batch was not delivered; the connection is dropped. `retry` is false when a fresh
/// connection would not help (an outright refusal), true for an I/O failure — most often a
/// kept-open socket the host has since closed, which is nobody's fault and worth one retry.
Failed { retry: bool },
}
fn check_status(status: Option<u16>) -> Result<(), ()> {
match status {
Some(code) if (200..300).contains(&code) => Ok(()),
Some(401) => {
UNAUTHORIZED.store(true, Ordering::Relaxed);
Err(())
/// Anything the transport can post over: the unix socket and the TCP fallback both qualify.
trait Wire: Read + Write {}
impl<T: Read + Write> Wire for T {}
/// A connection kept open across batches (docs/NASH_STREAM_PERF_PLAN.md §P4.1). Lives on the
/// flusher thread, which is the only thing that posts, so it needs no locking.
#[derive(Default)]
struct Connection {
/// The live connection and the request path the transport it came from wants.
open: Option<(Box<dyn Wire>, String)>,
}
impl Connection {
/// Post one batch, reconnecting once if the kept-open connection turned out to be dead.
fn post(&mut self, cfg: &Config, body: &[u8]) -> Result<(), ()> {
for _ in 0..2 {
if self.open.is_none() {
self.open = connect(cfg);
}
let Some((stream, path)) = self.open.as_mut() else {
return Err(()); // nothing to connect to; the caller spools
};
let request = http_request(path, cfg.token.as_deref(), body);
match post_on(stream.as_mut(), &request) {
Posted::Ok => return Ok(()),
Posted::OkAndClose => {
self.open = None;
return Ok(());
}
Posted::Failed { retry } => {
self.open = None;
if !retry {
return Err(());
}
}
}
}
// Treat unreadable/other responses as delivered-at-best-effort: the
// host 202s known tokens, so anything else is not retryable.
Some(_) | None => Err(()),
Err(())
}
}
fn post_unix(socket: &Path, token: Option<&str>, body: &[u8]) -> Result<(), ()> {
let stream = std::os::unix::net::UnixStream::connect(socket).map_err(|_| ())?;
/// Open the configured transport: the unix socket first, then the TCP hook.
fn connect(cfg: &Config) -> Option<(Box<dyn Wire>, String)> {
if let Some(socket) = &cfg.socket {
if let Ok(stream) = std::os::unix::net::UnixStream::connect(socket) {
let _ = stream.set_write_timeout(Some(IO_TIMEOUT));
let _ = stream.set_read_timeout(Some(IO_TIMEOUT));
return Some((Box::new(stream), "/shell-event".to_string()));
}
}
let url = cfg.hook_url.as_ref()?;
let addr = std::net::ToSocketAddrs::to_socket_addrs(&(url.host.as_str(), url.port))
.ok()?
.next()?;
let stream = std::net::TcpStream::connect_timeout(&addr, IO_TIMEOUT).ok()?;
let _ = stream.set_write_timeout(Some(IO_TIMEOUT));
let _ = stream.set_read_timeout(Some(IO_TIMEOUT));
let mut stream = stream;
stream
.write_all(&http_request("/shell-event", token, body))
.map_err(|_| ())?;
check_status(read_status(&mut stream))
// Batches are one write each and latency is not what this path optimizes, but a delayed ACK
// waiting on Nagle would hold the *response* — and with it the next batch — for milliseconds.
let _ = stream.set_nodelay(true);
Some((Box::new(stream), url.path.clone()))
}
fn post_tcp(url: &HookUrl, token: Option<&str>, body: &[u8]) -> Result<(), ()> {
let addr = (url.host.as_str(), url.port);
let stream = std::net::TcpStream::connect_timeout(
&std::net::ToSocketAddrs::to_socket_addrs(&addr)
.ok()
.and_then(|mut a| a.next())
.ok_or(())?,
IO_TIMEOUT,
)
.map_err(|_| ())?;
let _ = stream.set_write_timeout(Some(IO_TIMEOUT));
let _ = stream.set_read_timeout(Some(IO_TIMEOUT));
let mut stream = stream;
stream
.write_all(&http_request(&url.path, token, body))
.map_err(|_| ())?;
check_status(read_status(&mut stream))
/// Write one request and read its whole response, so the connection is left framed for the next.
fn post_on(stream: &mut dyn Wire, request: &[u8]) -> Posted {
if stream.write_all(request).is_err() {
return Posted::Failed { retry: true };
}
let Some(response) = read_response(stream) else {
return Posted::Failed { retry: true };
};
match response.status {
401 => {
UNAUTHORIZED.store(true, Ordering::Relaxed);
Posted::Failed { retry: false }
}
code if (200..300).contains(&code) => {
if response.reusable {
Posted::Ok
} else {
Posted::OkAndClose
}
}
// The host 202s known tokens, so anything else is not retryable.
_ => Posted::Failed { retry: false },
}
}
struct Response {
status: u16,
/// Whether the body was fully consumed and the peer means to keep the connection.
reusable: bool,
}
/// Read a response's head, then drain exactly its `Content-Length` body — the whole point being
/// that the next response starts where this one ends. A response we cannot frame (no length,
/// chunked, `Connection: close`) is still reported, but its connection is not reused.
fn read_response(stream: &mut dyn Read) -> Option<Response> {
let mut buf = Vec::with_capacity(256);
let mut chunk = [0u8; 512];
let head_end = loop {
if let Some(at) = find_headers_end(&buf) {
break at;
}
// A response head this long is not one of ours; give up rather than read forever.
if buf.len() > 8192 {
return None;
}
match stream.read(&mut chunk) {
Ok(0) => return None,
Ok(n) => buf.extend_from_slice(&chunk[..n]),
Err(ref e) if e.kind() == std::io::ErrorKind::Interrupted => (),
Err(_) => return None,
}
};
let head = String::from_utf8_lossy(&buf[..head_end]);
let mut lines = head.split("\r\n");
let status: u16 = lines.next()?.split_whitespace().nth(1)?.parse().ok()?;
let mut length: Option<usize> = None;
let mut close = false;
for line in lines {
let Some((name, value)) = line.split_once(':') else {
continue;
};
match name.trim().to_ascii_lowercase().as_str() {
"content-length" => length = value.trim().parse().ok(),
"connection" => close = value.trim().eq_ignore_ascii_case("close"),
_ => (),
}
}
let Some(length) = length else {
return Some(Response {
status,
reusable: false,
});
};
// Drain the body so the socket is positioned at the next response.
let mut remaining = length.saturating_sub(buf.len() - (head_end + 4));
while remaining > 0 {
let want = remaining.min(chunk.len());
match stream.read(&mut chunk[..want]) {
Ok(0) => {
return Some(Response {
status,
reusable: false,
})
}
Ok(n) => remaining -= n,
Err(ref e) if e.kind() == std::io::ErrorKind::Interrupted => (),
Err(_) => {
return Some(Response {
status,
reusable: false,
})
}
}
}
Some(Response {
status,
reusable: !close,
})
}
/// Offset of the `\r\n\r\n` that ends a response head.
fn find_headers_end(buf: &[u8]) -> Option<usize> {
buf.windows(4).position(|w| w == b"\r\n\r\n")
}
fn spool(dir: &Path, shell_id: &str, body: &[u8]) {
@@ -733,19 +1018,9 @@ fn spool(dir: &Path, shell_id: &str, body: &[u8]) {
}
}
fn deliver(cfg: &Config, shell_id: &str, body: &[u8]) {
let token = cfg.token.as_deref();
if !UNAUTHORIZED.load(Ordering::Relaxed) {
if let Some(socket) = &cfg.socket {
if post_unix(socket, token, body).is_ok() {
return;
}
}
if let Some(url) = &cfg.hook_url {
if post_tcp(url, token, body).is_ok() {
return;
}
}
fn deliver(cfg: &Config, shell_id: &str, conn: &mut Connection, body: &[u8]) {
if !UNAUTHORIZED.load(Ordering::Relaxed) && conn.post(cfg, body).is_ok() {
return;
}
if let Some(dir) = &cfg.spool {
spool(dir, shell_id, body);
@@ -806,14 +1081,46 @@ fn sweep_running() -> Sweep {
sweep
}
/// The captured half of an event, for the flusher's two passes over a batch.
fn event_capture(event: &mut Event) -> Option<&mut Capture> {
match event {
Event::Redirect { capture, .. }
| Event::Cmdsub { capture, .. }
| Event::Pipe { capture, .. } => Some(capture),
_ => None,
}
}
/// Settle a batch's payloads before it is serialized — the flusher-thread half of §P3/§P4.2.
///
/// Two passes, in this order: read back every deferred file range, then spend the per-batch
/// preview budget over what that leaves. Events keep their place and their counts either way; only
/// previews are ever given up, oldest kept first, so a storm of huge pipes loses its tail rather
/// than the whole batch losing its head.
fn settle(events: &mut [Event]) {
let mut spent = 0usize;
for event in events {
let Some(capture) = event_capture(event) else {
continue;
};
capture.resolve();
spent = spent.saturating_add(capture.preview_len());
if spent > BATCH_PREVIEW_BUDGET {
capture.count_only();
}
}
}
fn flusher(cfg: Config, shell_id: String, parent_shell: Option<String>, rx: mpsc::Receiver<Msg>) {
let mut buf: Vec<Event> = Vec::new();
let mut first_at: Option<Instant> = None;
let mut conn = Connection::default();
let flush = |buf: &mut Vec<Event>| {
let flush = |buf: &mut Vec<Event>, conn: &mut Connection| {
if buf.is_empty() {
return;
}
settle(buf);
let batch = Batch {
batch_type: "shell-batch",
source: "nash",
@@ -826,7 +1133,7 @@ fn flusher(cfg: Config, shell_id: String, parent_shell: Option<String>, rx: mpsc
events: buf,
};
if let Ok(body) = serde_json::to_vec(&batch) {
deliver(&cfg, &shell_id, &body);
deliver(&cfg, &shell_id, conn, &body);
}
buf.clear();
};
@@ -842,7 +1149,7 @@ fn flusher(cfg: Config, shell_id: String, parent_shell: Option<String>, rx: mpsc
}
buf.extend(sweep.announce);
if buf.len() >= FLUSH_MAX_EVENTS {
flush(&mut buf);
flush(&mut buf, &mut conn);
first_at = None;
}
}
@@ -863,12 +1170,12 @@ fn flusher(cfg: Config, shell_id: String, parent_shell: Option<String>, rx: mpsc
}
buf.push(ev);
if buf.len() >= FLUSH_MAX_EVENTS {
flush(&mut buf);
flush(&mut buf, &mut conn);
first_at = None;
}
}
Ok(Msg::FlushSync(ack)) => {
flush(&mut buf);
flush(&mut buf, &mut conn);
first_at = None;
let _ = ack.send(());
}
@@ -878,12 +1185,12 @@ fn flusher(cfg: Config, shell_id: String, parent_shell: Option<String>, rx: mpsc
Err(mpsc::RecvTimeoutError::Timeout) => {
// A tick is just "go re-sweep"; the batch is not due yet.
if !ticking {
flush(&mut buf);
flush(&mut buf, &mut conn);
first_at = None;
}
}
Err(mpsc::RecvTimeoutError::Disconnected) => {
flush(&mut buf);
flush(&mut buf, &mut conn);
return;
}
}
@@ -976,3 +1283,179 @@ pub fn report_fallback(reason: &str, input: &str) {
});
flush_sync();
}
// ---------------------------------------------------------------------------
// Tests
// ---------------------------------------------------------------------------
#[cfg(test)]
mod tests {
use super::*;
/// Serialize one event the way the flusher does, as JSON.
fn json(event: &Event) -> serde_json::Value {
serde_json::to_value(event).expect("event serializes")
}
fn pipe(bytes: u64, data: Vec<u8>) -> Event {
Event::Pipe {
seq: 1,
ts: 1,
pipeline_id: 1,
from_index: 0,
to_index: 1,
from_text: "a".to_string(),
to_text: "b".to_string(),
capture: Capture::Bytes {
bytes,
truncated: bytes > data.len() as u64,
data,
},
}
}
/// The four payload fields sit where they always have, and an empty hash still means "nothing
/// was captured" rather than "the hash of nothing" — the host reads it that way.
#[test]
fn capture_serializes_the_wire_shape() {
let event = json(&pipe(5, b"hello".to_vec()));
assert_eq!(event["kind"], "pipe");
assert_eq!(event["bytes"], 5);
assert_eq!(event["truncated"], false);
assert_eq!(event["previewB64"], b64(b"hello"));
assert_eq!(event["hash"], fnv1a_hex(b"hello"));
let empty = json(&Event::Cmdsub {
seq: 1,
ts: 1,
capture: Capture::None,
});
assert_eq!(empty["bytes"], 0);
assert_eq!(empty["hash"], "");
assert_eq!(empty["previewB64"], "");
}
/// Redaction and base64 run at serialization time now, so they must still actually run —
/// on the preview *and* on the stage labels a link carries.
#[test]
fn serialization_redacts_previews_and_stage_text() {
let secret = "ghp_0123456789abcdefghij";
let event = json(&Event::Pipe {
seq: 1,
ts: 1,
pipeline_id: 1,
from_index: 0,
to_index: 1,
from_text: format!("curl -H {secret}"),
to_text: "sh".to_string(),
capture: Capture::from_bytes(secret.len() as u64, secret.as_bytes()),
});
assert!(!event["fromText"].as_str().unwrap().contains(secret));
assert_eq!(event["fromText"], "curl -H «redacted»");
assert_eq!(event["previewB64"], b64("«redacted»".as_bytes()));
}
/// A capture is capped to the preview cap where it is taken, and says so.
#[test]
fn capture_caps_and_marks_truncation() {
let big = vec![b'x'; PREVIEW_CAP + 10];
let Capture::Bytes {
bytes,
truncated,
data,
} = Capture::from_bytes(big.len() as u64, &big)
else {
panic!("expected captured bytes");
};
assert_eq!(bytes, big.len() as u64);
assert_eq!(data.len(), PREVIEW_CAP);
assert!(truncated);
}
/// The per-batch budget (docs/NASH.md §5.4): once a batch has carried its allowance of
/// payload, later events keep their counts and hashes and give up their previews.
#[test]
fn batch_preview_budget_drops_previews_not_events() {
let payload = vec![b'z'; PREVIEW_CAP];
let over = BATCH_PREVIEW_BUDGET / PREVIEW_CAP + 2;
let mut events: Vec<Event> = (0..over)
.map(|_| pipe(PREVIEW_CAP as u64, payload.clone()))
.collect();
settle(&mut events);
let rendered: Vec<serde_json::Value> = events.iter().map(json).collect();
// Every event survives — only previews are given up, and the earliest keep theirs.
assert_eq!(rendered.len(), over);
assert_ne!(rendered[0]["previewB64"], "");
let last = rendered.last().unwrap();
assert_eq!(last["previewB64"], "");
assert_eq!(last["bytes"], PREVIEW_CAP);
assert_eq!(last["hash"], fnv1a_hex(&payload));
assert_eq!(last["truncated"], true);
let kept = rendered
.iter()
.filter(|e| e["previewB64"] != "")
.count();
assert_eq!(kept, BATCH_PREVIEW_BUDGET / PREVIEW_CAP);
}
/// A redirect's bytes are read by the flusher, from the range measured when the command
/// exited — an append that lands afterwards must not widen it.
#[test]
fn deferred_readback_reads_the_measured_range() {
let dir = std::env::temp_dir().join(format!("nash-observe-{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join("out.txt");
std::fs::write(&path, b"first\n").unwrap();
let mut event = Event::Redirect {
seq: 1,
ts: 1,
cmd_seq: 1,
op: ">".to_string(),
fd: 1,
target: Some(path.to_string_lossy().into_owned()),
capture: Capture::Deferred {
path: path.clone(),
offset: 0,
total: 6,
},
};
// The file grows before the batch flushes; the event still reports what its command wrote.
std::fs::write(&path, b"first\nsecond\n").unwrap();
settle(std::slice::from_mut(&mut event));
let rendered = json(&event);
assert_eq!(rendered["bytes"], 6);
assert_eq!(rendered["previewB64"], b64(b"first\n"));
assert_eq!(rendered["truncated"], false);
// A target that has since vanished degrades to metadata only rather than failing.
std::fs::remove_file(&path).unwrap();
let mut gone = Event::Redirect {
seq: 2,
ts: 1,
cmd_seq: 1,
op: ">".to_string(),
fd: 1,
target: None,
capture: Capture::Deferred {
path,
offset: 0,
total: 6,
},
};
settle(std::slice::from_mut(&mut gone));
assert_eq!(json(&gone)["bytes"], 0);
assert_eq!(json(&gone)["hash"], "");
let _ = std::fs::remove_dir_all(&dir);
}
/// Binary payloads are previewed as a hex head (now table-driven, §P3.3) — same bytes as the
/// `format!`-per-byte encoding it replaced.
#[test]
fn binary_previews_as_hex() {
let preview = preview_b64(&[0xff, 0x00, 0x0a]);
assert_eq!(preview, b64(b"<binary> ff000a"));
}
}