From 2cbff0667ddee7838e67c30ef3d34debcb119343 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 --- corpus/overhead.py | 235 +++++++++++-- nash-observe/src/lib.rs | 719 +++++++++++++++++++++++++++++++++------- nash/tests/observe.rs | 106 +++++- 3 files changed, 897 insertions(+), 163 deletions(-) diff --git a/corpus/overhead.py b/corpus/overhead.py index 8c9bbc9..af20603 100644 --- a/corpus/overhead.py +++ b/corpus/overhead.py @@ -1,22 +1,44 @@ #!/usr/bin/env python3 -"""Corpus wall-clock overhead harness for nash (docs/NASH.md §11, M2/M3 gate: <3%). +"""Corpus overhead harness for nash (docs/NASH.md §11, M2/M3 gate: <3%). Replays the corpus under bash and nash in alternating rounds and compares total -wall-clock per shell. Only the shell subprocess is timed (fixture seeding and -filesystem snapshots are outside the clock). Rounds alternate shell order so -cache/thermal drift cancels; the reported figure uses the median round total. +wall-clock **and CPU time** per shell. Only the shell subprocess is measured +(fixture seeding and filesystem snapshots are outside both clocks). Rounds +alternate shell order so cache/thermal drift cancels; the reported figures use +the median round total. + +Two measurements, because they answer different questions +(docs/NASH_STREAM_PERF_PLAN.md §2): + + * **wall** is what a single command waits for — the M3 gate. + * **cpu** (user+sys of the shell and everything it spawned, via `rusage`) is + what nash actually spends. The pipe tee runs on its own thread, so on idle + hardware it can add CPU while *costing no wall time at all* — and a + wall-only harness would report that as free. It is not: the deployed box + runs many agents at once, where that CPU comes out of everyone's clock. + +Beyond the corpus, `--stream-mb` runs a **throughput** case — hundreds of MB +across one and two pipe links — which is the shape the tee is optimized for and +the one the corpus (a few KB per command) cannot see. Usage: overhead.py [--nash PATH] [--bash PATH] [--corpus PATH] [--rounds N] - [--observe-spool DIR] [--gate PCT] [--json PATH] + [--observe-spool DIR] [--gate PCT] [--stream-mb MB] + [--stream-gate PCT] [--json PATH] [--allow-nash-baseline] --observe-spool enables nash observation (spool transport) so the measured configuration is the deployed one; the env is set identically for bash, where it is inert. + +The baseline shell must be a *real* bash: on a Nucleic-managed box `/bin/bash` +IS nash (docs/NASH.md §7), so the default baseline is `/usr/bin/bash.real` and +a baseline that reports `NUCLEIC_NASH=1` is refused outright — benchmarking +nash against itself reports ~0% overhead and means nothing. """ import argparse import json import os +import resource import shutil import statistics import subprocess @@ -26,12 +48,30 @@ import time from replay import FIXTURES, TIMEOUT_S, seed +#: Where a Nucleic-managed image keeps the real bash after the nash divert. +DEFAULT_BASH = "/usr/bin/bash.real" -def run_timed(shell, cmd, extra_env): - parent = tempfile.mkdtemp(prefix="nash-overhead-") - workdir = os.path.join(parent, "workspace") - os.makedirs(workdir) - seed(workdir) + +def child_cpu(): + """User+sys seconds of every child this process has reaped.""" + usage = resource.getrusage(resource.RUSAGE_CHILDREN) + return usage.ru_utime + usage.ru_stime + + +def run_timed(shell, cmd, extra_env, workdir=None): + """Run one command, returning (wall_seconds, cpu_seconds). + + CPU comes from the delta of `RUSAGE_CHILDREN`, which only counts children + already reaped — `subprocess.run` waits, and the harness runs one child at a + time, so the delta is exactly this command's shell and its descendants + (including nash's tee and flusher threads). + """ + parent = None + if workdir is None: + parent = tempfile.mkdtemp(prefix="nash-overhead-") + workdir = os.path.join(parent, "workspace") + os.makedirs(workdir) + seed(workdir) env = { "PATH": "/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin", "HOME": workdir, @@ -42,6 +82,7 @@ def run_timed(shell, cmd, extra_env): } if extra_env: env.update(extra_env) + cpu_before = child_cpu() start = time.monotonic() try: subprocess.run( @@ -54,24 +95,136 @@ def run_timed(shell, cmd, extra_env): except subprocess.TimeoutExpired: pass elapsed = time.monotonic() - start - shutil.rmtree(parent, ignore_errors=True) - return elapsed + cpu = child_cpu() - cpu_before + if parent: + shutil.rmtree(parent, ignore_errors=True) + return elapsed, cpu + + +def is_nash(shell): + """Whether `shell` is nash wearing another name (docs/NASH.md §2). + + Probed with a scrubbed environment on purpose: the harness itself is very + likely running *under* nash, which exports `NUCLEIC_NASH=1` to everything it + spawns — inheriting that would make every shell look like nash. Only a shell + that sets the variable for itself answers 1 here. + """ + try: + out = subprocess.run( + [shell, "-c", 'printf %s "${NUCLEIC_NASH-}"'], + capture_output=True, + timeout=30, + text=True, + env={"PATH": "/usr/bin:/bin"}, + ) + except (OSError, subprocess.SubprocessError): + return False + return out.stdout.strip() == "1" + + +def stream_cases(megabytes): + """Throughput scripts: the same payload across one link and across two. + + `/dev/zero → /dev/null` deliberately: the point is the cost of *carrying* + bytes across a tapped link, so neither end should be doing work of its own. + """ + count = megabytes * 1024 * 1024 + return [ + (f"1 link ({megabytes} MB)", f"head -c {count} /dev/zero | cat > /dev/null"), + ( + f"2 links ({megabytes} MB)", + f"head -c {count} /dev/zero | cat | cat > /dev/null", + ), + ] + + +def measure(shell, cmds, extra_env): + """Total (wall, cpu) for one pass over `cmds`.""" + wall = cpu = 0.0 + for cmd in cmds: + w, c = run_timed(shell, cmd, extra_env) + wall += w + cpu += c + return wall, cpu + + +def percent(nash, bash): + """nash's cost over bash's, in percent. Infinite-safe for a zero baseline.""" + return 100.0 * (nash / bash - 1.0) if bash > 0 else float("nan") + + +def report(label, bash_vals, nash_vals, gate=None): + """Print (and return) one comparison's medians and percentages.""" + bash_wall = statistics.median(w for w, _ in bash_vals) + bash_cpu = statistics.median(c for _, c in bash_vals) + nash_wall = statistics.median(w for w, _ in nash_vals) + nash_cpu = statistics.median(c for _, c in nash_vals) + result = { + "bash_wall_s": bash_wall, + "bash_cpu_s": bash_cpu, + "nash_wall_s": nash_wall, + "nash_cpu_s": nash_cpu, + "wall_percent": percent(nash_wall, bash_wall), + "cpu_percent": percent(nash_cpu, bash_cpu), + } + print(f"\n{label}") + print(f" bash: {bash_wall:.3f}s wall / {bash_cpu:.3f}s cpu") + print(f" nash: {nash_wall:.3f}s wall / {nash_cpu:.3f}s cpu") + suffix = f" (gate: <{gate:.1f}%)" if gate is not None else "" + print( + f" overhead: {result['wall_percent']:+.2f}% wall / " + f"{result['cpu_percent']:+.2f}% cpu{suffix}" + ) + return result def main(): ap = argparse.ArgumentParser() ap.add_argument("--nash", default=os.environ.get("NASH_BIN", "nash")) - ap.add_argument("--bash", default="/bin/bash") + ap.add_argument( + "--bash", + default=DEFAULT_BASH, + help=f"baseline shell — must not be nash (default {DEFAULT_BASH})", + ) ap.add_argument( "--corpus", default=os.path.join(os.path.dirname(os.path.abspath(__file__)), "corpus.jsonl"), ) ap.add_argument("--rounds", type=int, default=5) ap.add_argument("--observe-spool", help="enable nash observation, spooling to this dir") - ap.add_argument("--gate", type=float, default=3.0, help="max overhead percent") + ap.add_argument("--gate", type=float, default=3.0, help="max corpus wall overhead percent") + ap.add_argument( + "--stream-mb", + type=int, + default=256, + help="payload per throughput case (0 disables the stream cases)", + ) + ap.add_argument( + "--stream-gate", + type=float, + default=50.0, + help="max stream CPU overhead percent (the tee's own cost)", + ) + ap.add_argument( + "--allow-nash-baseline", + action="store_true", + help="benchmark against a nash baseline anyway (produces meaningless numbers)", + ) ap.add_argument("--json", help="also write results to this path") args = ap.parse_args() + # A baseline that is itself nash makes every number here ~0% and hides whatever + # regressed. On a Nucleic box that is the *default* state of /bin/bash, so this is + # a refusal rather than a warning (docs/NASH_STREAM_PERF_PLAN.md §regression guard). + if is_nash(args.bash) and not args.allow_nash_baseline: + print( + f"error: baseline shell {args.bash} reports NUCLEIC_NASH=1 — it IS nash.\n" + f" Point --bash at the real bash ({DEFAULT_BASH} on a Nucleic image),\n" + " or pass --allow-nash-baseline if you really mean to compare nash to nash.", + file=sys.stderr, + ) + return 2 + observe_env = None if args.observe_spool: os.makedirs(args.observe_spool, exist_ok=True) @@ -82,6 +235,7 @@ def main(): with open(args.corpus) as f: cmds = [json.loads(line)["cmd"] for line in f if line.strip()] + streams = stream_cases(args.stream_mb) if args.stream_mb > 0 else [] # Warm-up: one untimed pass per shell (page cache, binary load). for shell in (args.bash, args.nash): @@ -89,22 +243,33 @@ def main(): run_timed(shell, cmd, observe_env) totals = {"bash": [], "nash": []} + stream_totals = {name: {"bash": [], "nash": []} for name, _ in streams} for round_no in range(args.rounds): order = [("bash", args.bash), ("nash", args.nash)] if round_no % 2: order.reverse() for name, shell in order: - total = sum(run_timed(shell, cmd, observe_env) for cmd in cmds) - totals[name].append(total) - print(f"round {round_no + 1} {name}: {total:.3f}s", flush=True) + wall, cpu = measure(shell, cmds, observe_env) + totals[name].append((wall, cpu)) + print(f"round {round_no + 1} {name}: {wall:.3f}s wall / {cpu:.3f}s cpu", flush=True) + for case, script in streams: + stream_totals[case][name].append(run_timed(shell, script, observe_env)) - bash_med = statistics.median(totals["bash"]) - nash_med = statistics.median(totals["nash"]) - overhead = 100.0 * (nash_med / bash_med - 1.0) + corpus = report( + f"corpus ({len(cmds)} cmds x {args.rounds} rounds)", + totals["bash"], + totals["nash"], + gate=args.gate, + ) + stream_results = { + case: report(f"stream {case}", vals["bash"], vals["nash"], gate=args.stream_gate) + for case, vals in stream_totals.items() + } - print(f"\ncorpus: {len(cmds)} cmds x {args.rounds} rounds") - print(f"bash median: {bash_med:.3f}s nash median: {nash_med:.3f}s") - print(f"overhead: {overhead:+.2f}% (gate: <{args.gate:.1f}%)") + corpus_ok = corpus["wall_percent"] < args.gate + # The stream cases are gated on CPU: the tee's cost is a copier thread, which idle + # hardware hides from wall-clock entirely. + stream_ok = all(r["cpu_percent"] < args.stream_gate for r in stream_results.values()) if args.json: with open(args.json, "w") as f: @@ -112,20 +277,30 @@ def main(): { "cmds": len(cmds), "rounds": args.rounds, - "totals": totals, - "bash_median_s": bash_med, - "nash_median_s": nash_med, - "overhead_percent": overhead, + "totals": {k: [{"wall_s": w, "cpu_s": c} for w, c in v] for k, v in totals.items()}, + "corpus": corpus, + "streams": stream_results, + "stream_mb": args.stream_mb, "gate_percent": args.gate, + "stream_gate_percent": args.stream_gate, "observed": bool(args.observe_spool), + "bash": args.bash, + "nash": args.nash, + # Kept for readers of the pre-CPU schema. + "bash_median_s": corpus["bash_wall_s"], + "nash_median_s": corpus["nash_wall_s"], + "overhead_percent": corpus["wall_percent"], }, f, indent=1, ) - gate = overhead < args.gate - print(f"M3 overhead gate (<{args.gate:.1f}%): {'PASS' if gate else 'FAIL'}") - return 0 if gate else 1 + print(f"\nM3 overhead gate (corpus wall <{args.gate:.1f}%): {'PASS' if corpus_ok else 'FAIL'}") + if streams: + print( + f"stream gate (cpu <{args.stream_gate:.1f}%): {'PASS' if stream_ok else 'FAIL'}" + ) + return 0 if corpus_ok and stream_ok else 1 if __name__ == "__main__": diff --git a/nash-observe/src/lib.rs b/nash-observe/src/lib.rs index 817f69e..6ea75ef 100644 --- a/nash-observe/src/lib.rs +++ b/nash-observe/src/lib.rs @@ -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, + }, + /// 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(&self, serializer: S) -> Result { + 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, - 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, 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)> { 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, 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!(" {hex}").as_bytes()) + let mut hex = String::with_capacity(1024 * 2 + 9); + hex.push_str(" "); + 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(text: &str, serializer: S) -> Result { + 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 { + // 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 { bytes } -fn read_status(mut stream: R) -> Option { - 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) -> 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 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, 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, 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 { + 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 = 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 { + 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, rx: mpsc::Receiver) { let mut buf: Vec = Vec::new(); let mut first_at: Option = None; + let mut conn = Connection::default(); - let flush = |buf: &mut Vec| { + let flush = |buf: &mut Vec, 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, 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, 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, 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, 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) -> 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 = (0..over) + .map(|_| pipe(PREVIEW_CAP as u64, payload.clone())) + .collect(); + settle(&mut events); + + let rendered: Vec = 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" ff000a")); + } +} diff --git a/nash/tests/observe.rs b/nash/tests/observe.rs index 2bf1c0a..0d7776f 100644 --- a/nash/tests/observe.rs +++ b/nash/tests/observe.rs @@ -48,6 +48,9 @@ fn scratch_dir(prefix: &str) -> std::path::PathBuf { struct Sink { dir: std::path::PathBuf, rx: mpsc::Receiver<(Option, serde_json::Value)>, + /// Connections accepted so far — how the keep-alive test tells one reused connection from + /// one per batch (docs/NASH_STREAM_PERF_PLAN.md §P4.1). + connections: std::sync::Arc, } impl Sink { @@ -59,37 +62,52 @@ impl Sink { let listener = UnixListener::bind(&socket_path) .unwrap_or_else(|e| panic!("bind {} ({} bytes): {e}", socket_path.display(), socket_path.as_os_str().len())); let (tx, rx) = mpsc::channel(); + let connections = std::sync::Arc::new(AtomicU32::new(0)); + let accepted = std::sync::Arc::clone(&connections); std::thread::spawn(move || { for stream in listener.incoming() { let Ok(mut stream) = stream else { break }; + accepted.fetch_add(1, Ordering::Relaxed); let tx = tx.clone(); std::thread::spawn(move || { let mut buf = Vec::new(); let mut chunk = [0u8; 4096]; - // Read until headers + declared body length are complete. + // Read until headers + declared body length are complete, answer, and keep + // reading: nash reuses one connection for every batch of a shell's life + // (docs/NASH_STREAM_PERF_PLAN.md §P4.1), exactly as the host's route does. loop { - match stream.read(&mut chunk) { - Ok(0) => break, - Ok(n) => { - buf.extend_from_slice(&chunk[..n]); - if let Some(body) = full_body(&buf) { - let auth = header(&buf, "authorization"); - if let Ok(json) = serde_json::from_slice(body) { - let _ = tx.send((auth, json)); - } - let _ = stream.write_all( - b"HTTP/1.1 202 Accepted\r\nContent-Length: 0\r\n\r\n", - ); - break; + if let Some(request) = request_len(&buf) { + let auth = header(&buf, "authorization"); + if let Some(body) = full_body(&buf) { + if let Ok(json) = serde_json::from_slice(body) { + let _ = tx.send((auth, json)); } } + if stream + .write_all( + b"HTTP/1.1 202 Accepted\r\nContent-Length: 0\r\nConnection: keep-alive\r\n\r\n", + ) + .is_err() + { + break; + } + buf.drain(..request); + continue; + } + match stream.read(&mut chunk) { + Ok(0) => break, + Ok(n) => buf.extend_from_slice(&chunk[..n]), Err(_) => break, } } }); } }); - Self { dir, rx } + Self { + dir, + rx, + connections, + } } fn socket(&self) -> std::path::PathBuf { @@ -148,6 +166,14 @@ impl Drop for Sink { } } +/// Total bytes of the first complete request in `buf` (head + body), or None while it is still +/// arriving — what lets one connection carry the next batch after this one is answered. +fn request_len(buf: &[u8]) -> Option { + let headers_end = buf.windows(4).position(|w| w == b"\r\n\r\n")? + 4; + let len: usize = header(buf, "content-length")?.parse().ok()?; + (buf.len() >= headers_end + len).then_some(headers_end + len) +} + fn full_body(buf: &[u8]) -> Option<&[u8]> { let headers_end = buf.windows(4).position(|w| w == b"\r\n\r\n")? + 4; let len: usize = header(buf, "content-length")?.parse().ok()?; @@ -869,3 +895,53 @@ fn capture_off_silences_nash_without_a_policy() { "capture=off with no policy posts nothing" ); } + +/// A shell posts a batch every 500 ms for as long as it lives, so the connection is opened once +/// and kept (docs/NASH_STREAM_PERF_PLAN.md §P4.1) rather than reconnected per batch. Two commands +/// either side of a flush window are two batches; one connection has to carry both. +#[test] +fn one_connection_carries_every_batch() { + let sink = Sink::start(); + let status = run_nash(&sink, "echo one; sleep 0.9; echo two"); + assert_eq!(status.code(), Some(0)); + + let events = sink.events_until(|evs| { + find(evs, "exec", |e| e["argv"][0] == "echo" && e["argv"][1] == "two").is_some() + }); + // Both commands were reported… + assert!(find(&events, "exec", |e| e["argv"][1] == "one").is_some()); + assert!(find(&events, "exec", |e| e["argv"][1] == "two").is_some()); + // …and the shell only ever dialled the sink once. + assert_eq!( + sink.connections.load(Ordering::Relaxed), + 1, + "batches must share one kept-open connection" + ); +} + +/// The budget's end-to-end half (docs/NASH.md §5.4): a burst of fat pipes in one batch keeps every +/// event and every byte count, and gives up only the previews past the allowance. +#[test] +fn batch_preview_budget_keeps_counts_and_drops_previews() { + let sink = Sink::start(); + // 24 links × 64 KiB of payload each — comfortably past the 1 MiB per-batch preview budget, + // and fast enough that they all land in one 500 ms flush. + let status = run_nash( + &sink, + "for i in $(seq 1 24); do head -c 65536 /dev/zero | cat > /dev/null; done", + ); + assert_eq!(status.code(), Some(0)); + + let events = sink.events_until(|evs| evs.iter().filter(|e| e["kind"] == "pipe").count() >= 24); + let pipes: Vec<_> = events.iter().filter(|e| e["kind"] == "pipe").collect(); + assert!(pipes.len() >= 24, "every link is still reported: {}", pipes.len()); + // Nothing loses its count or its fingerprint… + for pipe in &pipes { + assert_eq!(pipe["bytes"], 65536); + assert_ne!(pipe["hash"], ""); + } + // …and the ones past the budget are the ones without a preview. + let dropped = pipes.iter().filter(|p| p["previewB64"] == "").count(); + assert!(dropped > 0, "the budget dropped no previews"); + assert!(dropped < pipes.len(), "the budget dropped every preview"); +}