diff --git a/Cargo.lock b/Cargo.lock index 752a9f6..d7fe137 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1240,6 +1240,12 @@ dependencies = [ "either", ] +[[package]] +name = "itoa" +version = "1.0.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682" + [[package]] name = "js-sys" version = "0.3.103" @@ -1364,12 +1370,20 @@ dependencies = [ name = "nash" version = "0.1.0" dependencies = [ + "brush-parser", "brush-shell", + "nash-observe", + "serde_json", ] [[package]] name = "nash-observe" version = "0.1.0" +dependencies = [ + "brush-core", + "serde", + "serde_json", +] [[package]] name = "nibble_vec" @@ -1912,6 +1926,19 @@ dependencies = [ "syn", ] +[[package]] +name = "serde_json" +version = "1.0.150" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e8014e44b4736ed0538adeecded0fce2a272f22dc9578a7eb6b2d9993c74cfb9" +dependencies = [ + "itoa", + "memchr", + "serde", + "serde_core", + "zmij", +] + [[package]] name = "serde_spanned" version = "1.1.1" @@ -2998,3 +3025,9 @@ dependencies = [ "quote", "syn", ] + +[[package]] +name = "zmij" +version = "1.0.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "29666d0abbfad1e3dc4dcf6144730dd3a3ab225bbbdac83319345b1b44ccfc1b" diff --git a/corpus/replay.py b/corpus/replay.py index ec4c3b1..69c2e8e 100644 --- a/corpus/replay.py +++ b/corpus/replay.py @@ -7,7 +7,12 @@ separately and does NOT count against parity (error-message wording legitimately differs between shells); a stderr-only difference is reported for review. Usage: replay.py [--nash PATH] [--bash PATH] [--corpus PATH] [--json PATH] + [--observe-spool DIR] Parity target (M0/M3 gate): >= 99%. + +--observe-spool DIR turns nash observation ON during the replay (spool +transport; the env is set identically for bash, where it is inert), proving +the gate hooks don't perturb command semantics (docs/NASH.md §10.1). """ import argparse @@ -115,8 +120,17 @@ def main(): default=os.path.join(os.path.dirname(os.path.abspath(__file__)), "corpus.jsonl"), ) ap.add_argument("--json", help="also write the full report to this path") + ap.add_argument("--observe-spool", help="enable nash observation, spooling to this dir") args = ap.parse_args() + observe_env = None + if args.observe_spool: + os.makedirs(args.observe_spool, exist_ok=True) + observe_env = { + "NUCLEIC_SHELL_SPOOL": args.observe_spool, + "NUCLEIC_SESSION_ID": "corpus-replay", + } + entries = [] with open(args.corpus) as f: for line in f: @@ -126,8 +140,8 @@ def main(): results = [] for entry in entries: - b = run_one(args.bash, entry["cmd"]) - n = run_one(args.nash, entry["cmd"]) + b = run_one(args.bash, entry["cmd"], extra_env=observe_env) + n = run_one(args.nash, entry["cmd"], extra_env=observe_env) verdict = classify(b, n) results.append({"id": entry["id"], "cmd": entry["cmd"], "verdict": verdict, "bash": b, "nash": n}) diff --git a/nash-observe/Cargo.toml b/nash-observe/Cargo.toml index cc59fd5..b5b5313 100644 --- a/nash-observe/Cargo.toml +++ b/nash-observe/Cargo.toml @@ -5,3 +5,8 @@ edition = "2021" description = "nash-observe — event model, taps, batching, and transport for nash" license = "MIT" publish = false + +[dependencies] +brush-core = { path = "../../third_party/brush/brush-core" } +serde = { version = "1", features = ["derive"] } +serde_json = "1" diff --git a/nash-observe/src/lib.rs b/nash-observe/src/lib.rs index b172924..3914a73 100644 --- a/nash-observe/src/lib.rs +++ b/nash-observe/src/lib.rs @@ -1,5 +1,497 @@ -//! nash-observe — event model, redirection/pipe taps, batching, spool, and -//! transport for nash (see docs/NASH.md §4–§6). +//! nash-observe — the observation layer for nash (docs/NASH.md §4–§6). //! -//! M0 placeholder: this crate gains the `Gate` implementation (allow-all), -//! the shell-batch event schema, and the unix-socket HTTP transport in M1. +//! Implements the allow-all [`brush_core::gate::Gate`], the shell-batch event +//! model, near-live batching (200 events / 500 ms / exit), and the transport +//! chain: HTTP over a unix socket (`NUCLEIC_SHELL_SOCKET`), then HTTP over TCP +//! (`NUCLEIC_SHELL_HOOK_URL`), then a JSONL disk spool (`NUCLEIC_SHELL_SPOOL`). +//! +//! Failure doctrine (docs/NASH.md §5.5): nothing in this crate may affect a +//! command's outcome. Hook bodies are panic-caught, the event channel drops +//! (with a `dropped` counter) rather than blocks, and transport is +//! fire-and-forget with short timeouts and spool fallback. + +use std::collections::HashMap; +use std::io::{Read, Write}; +use std::path::{Path, PathBuf}; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::mpsc; +use std::sync::{Mutex, OnceLock}; +use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; + +use serde::Serialize; + +const FLUSH_MAX_EVENTS: usize = 200; +const FLUSH_MAX_AGE: Duration = Duration::from_millis(500); +const CHANNEL_BOUND: usize = 4096; +const IO_TIMEOUT: Duration = Duration::from_millis(1500); +const EXIT_FLUSH_TIMEOUT: Duration = Duration::from_millis(2000); + +// --------------------------------------------------------------------------- +// Event model (docs/NASH.md §6.1) +// --------------------------------------------------------------------------- + +#[derive(Serialize)] +#[serde(tag = "kind", rename_all = "camelCase")] +enum Event { + #[serde(rename = "exec", rename_all = "camelCase")] + Exec { + seq: u64, + ts: u64, + argv: Vec, + cwd: String, + exit_code: i32, + duration_ms: u64, + }, + #[serde(rename = "cd", rename_all = "camelCase")] + Cd { seq: u64, ts: u64, from: String, to: String }, + #[serde(rename = "export", rename_all = "camelCase")] + Export { + seq: u64, + ts: u64, + op: String, + names: Vec, + }, + #[serde(rename = "fallback", rename_all = "camelCase")] + Fallback { + seq: u64, + ts: u64, + reason: String, + input: String, + }, + #[serde(rename = "dropped", rename_all = "camelCase")] + Dropped { seq: u64, ts: u64, count: u64 }, +} + +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +struct Batch<'a> { + #[serde(rename = "type")] + batch_type: &'static str, + source: &'static str, + session_id: Option<&'a str>, + shell_id: &'a str, + parent_shell: Option<&'a str>, + events: &'a [Event], +} + +// --------------------------------------------------------------------------- +// Config +// --------------------------------------------------------------------------- + +struct Config { + socket: Option, + hook_url: Option, + token: Option, + session_id: Option, + spool: Option, +} + +struct HookUrl { + host: String, + port: u16, + path: String, +} + +fn parse_hook_url(url: &str) -> Option { + let rest = url.strip_prefix("http://")?; + let (authority, path) = match rest.find('/') { + Some(idx) => (&rest[..idx], rest[idx..].to_string()), + None => (rest, "/shell-event".to_string()), + }; + let (host, port) = match authority.rsplit_once(':') { + Some((h, p)) => (h.to_string(), p.parse().ok()?), + None => (authority.to_string(), 80), + }; + Some(HookUrl { host, port, path }) +} + +impl Config { + fn from_env() -> Option { + if std::env::var("NUCLEIC_SHELL_CAPTURE").as_deref() == Ok("off") { + return None; + } + let socket = std::env::var_os("NUCLEIC_SHELL_SOCKET").map(PathBuf::from); + let hook_url = std::env::var("NUCLEIC_SHELL_HOOK_URL") + .ok() + .and_then(|u| parse_hook_url(&u)); + let spool = std::env::var_os("NUCLEIC_SHELL_SPOOL").map(PathBuf::from); + if socket.is_none() && hook_url.is_none() && spool.is_none() { + return None; + } + Some(Self { + socket, + hook_url, + token: std::env::var("NUCLEIC_HOOK_TOKEN").ok(), + session_id: std::env::var("NUCLEIC_SESSION_ID").ok(), + spool, + }) + } +} + +// --------------------------------------------------------------------------- +// Observer state +// --------------------------------------------------------------------------- + +enum Msg { + Event(Event), + FlushSync(mpsc::SyncSender<()>), +} + +struct PendingExec { + argv: Vec, + cwd: PathBuf, + started: Instant, +} + +struct Observer { + tx: mpsc::SyncSender, + next_token: AtomicU64, + next_seq: AtomicU64, + dropped: AtomicU64, + pending: Mutex>, +} + +static OBSERVER: OnceLock = OnceLock::new(); + +impl Observer { + fn seq(&self) -> u64 { + self.next_seq.fetch_add(1, Ordering::Relaxed) + } + + fn send(&self, event: Event) { + match self.tx.try_send(Msg::Event(event)) { + Ok(()) => { + // Piggyback a dropped-counter event whenever capacity returns. + let dropped = self.dropped.swap(0, Ordering::Relaxed); + if dropped > 0 { + let _ = self.tx.try_send(Msg::Event(Event::Dropped { + seq: self.seq(), + ts: now_millis(), + count: dropped, + })); + } + } + Err(_) => { + self.dropped.fetch_add(1, Ordering::Relaxed); + } + } + } +} + +fn now_millis() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|d| u64::try_from(d.as_millis()).unwrap_or(0)) + .unwrap_or(0) +} + +// --------------------------------------------------------------------------- +// Gate implementation (allow-all; docs/NASH.md §4.2) +// --------------------------------------------------------------------------- + +struct RecordingGate; + +impl brush_core::gate::Gate for RecordingGate { + fn on_exec(&self, ev: brush_core::gate::ExecStart) -> brush_core::gate::Verdict { + let token = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let obs = OBSERVER.get()?; + let token = obs.next_token.fetch_add(1, Ordering::Relaxed); + if let Ok(mut pending) = obs.pending.lock() { + pending.insert( + token, + PendingExec { + argv: ev.argv, + cwd: ev.cwd, + started: Instant::now(), + }, + ); + } + Some(token) + })) + .ok() + .flatten() + .unwrap_or(0); + brush_core::gate::Verdict::Allow(token) + } + + fn on_exit(&self, token: u64, ev: brush_core::gate::ExecEnd) { + let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let Some(obs) = OBSERVER.get() else { return }; + let Some(pending) = obs.pending.lock().ok().and_then(|mut p| p.remove(&token)) + else { + return; + }; + let ts = now_millis(); + let cwd_str = pending.cwd.to_string_lossy().into_owned(); + + // Derived, names-only state-change events (docs/NASH.md §5.3–5.4). + match pending.argv.first().map(String::as_str) { + Some("cd") if ev.exit_code == 0 => { + let to = pending + .argv + .get(1) + .map_or_else(|| "~".to_string(), Clone::clone); + obs.send(Event::Cd { + seq: obs.seq(), + ts, + from: cwd_str.clone(), + to, + }); + } + Some(op @ ("export" | "declare" | "typeset" | "unset" | "readonly" | "local")) + if ev.exit_code == 0 => + { + let names: Vec = pending + .argv + .iter() + .skip(1) + .filter(|a| !a.starts_with('-')) + .map(|a| a.split('=').next().unwrap_or(a).to_string()) + .collect(); + if !names.is_empty() { + obs.send(Event::Export { + seq: obs.seq(), + ts, + op: op.to_string(), + names, + }); + } + } + _ => {} + } + + obs.send(Event::Exec { + seq: obs.seq(), + ts, + argv: pending.argv, + cwd: cwd_str, + exit_code: ev.exit_code, + duration_ms: u64::try_from(pending.started.elapsed().as_millis()).unwrap_or(0), + }); + })); + } +} + +// --------------------------------------------------------------------------- +// Transport +// --------------------------------------------------------------------------- + +/// Set once the host rejects our token; we stop posting and spool instead. +static UNAUTHORIZED: AtomicBool = AtomicBool::new(false); + +fn http_request(path: &str, token: Option<&str>, body: &[u8]) -> Vec { + 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", + body.len() + ); + if let Some(token) = token { + req.push_str(&format!("Authorization: Bearer {token}\r\n")); + } + req.push_str("\r\n"); + let mut bytes = req.into_bytes(); + bytes.extend_from_slice(body); + 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() +} + +fn check_status(status: Option) -> Result<(), ()> { + match status { + Some(code) if (200..300).contains(&code) => Ok(()), + Some(401) => { + UNAUTHORIZED.store(true, Ordering::Relaxed); + Err(()) + } + // Treat unreadable/other responses as delivered-at-best-effort: the + // host 202s known tokens, so anything else is not retryable. + Some(_) | None => Err(()), + } +} + +fn post_unix(socket: &Path, token: Option<&str>, body: &[u8]) -> Result<(), ()> { + let stream = std::os::unix::net::UnixStream::connect(socket).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("/shell-event", token, body)) + .map_err(|_| ())?; + check_status(read_status(&mut stream)) +} + +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)) +} + +fn spool(dir: &Path, shell_id: &str, body: &[u8]) { + let _ = std::fs::create_dir_all(dir); + let path = dir.join(format!("{shell_id}.jsonl")); + if let Ok(mut f) = std::fs::OpenOptions::new().create(true).append(true).open(path) { + let _ = f.write_all(body); + let _ = f.write_all(b"\n"); + } +} + +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; + } + } + } + if let Some(dir) = &cfg.spool { + spool(dir, shell_id, body); + } +} + +// --------------------------------------------------------------------------- +// Flusher thread +// --------------------------------------------------------------------------- + +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 flush = |buf: &mut Vec| { + if buf.is_empty() { + return; + } + let batch = Batch { + batch_type: "shell-batch", + source: "nash", + session_id: cfg.session_id.as_deref(), + shell_id: &shell_id, + parent_shell: parent_shell.as_deref(), + events: buf, + }; + if let Ok(body) = serde_json::to_vec(&batch) { + deliver(&cfg, &shell_id, &body); + } + buf.clear(); + }; + + loop { + let timeout = match first_at { + Some(t) => FLUSH_MAX_AGE.saturating_sub(t.elapsed()), + None => Duration::from_secs(3600), + }; + match rx.recv_timeout(timeout) { + Ok(Msg::Event(ev)) => { + if buf.is_empty() { + first_at = Some(Instant::now()); + } + buf.push(ev); + if buf.len() >= FLUSH_MAX_EVENTS { + flush(&mut buf); + first_at = None; + } + } + Ok(Msg::FlushSync(ack)) => { + flush(&mut buf); + first_at = None; + let _ = ack.send(()); + } + Err(mpsc::RecvTimeoutError::Timeout) => { + flush(&mut buf); + first_at = None; + } + Err(mpsc::RecvTimeoutError::Disconnected) => { + flush(&mut buf); + return; + } + } + } +} + +// --------------------------------------------------------------------------- +// Public API +// --------------------------------------------------------------------------- + +/// Reads `NUCLEIC_SHELL_*` config from the environment; when observation is +/// configured, installs the recording gate into brush-core and starts the +/// flusher thread. Returns whether observation is active. +/// +/// Also threads shell lineage: reads `NUCLEIC_SHELL_PARENT` as this shell's +/// parent and re-exports it as this shell's own id for child shells. +pub fn install_from_env() -> bool { + let Some(cfg) = Config::from_env() else { + return false; + }; + + let shell_id = format!("{:x}-{:x}", std::process::id(), now_millis()); + let parent_shell = std::env::var("NUCLEIC_SHELL_PARENT").ok(); + // SAFETY: called from `main` before any other threads exist. + unsafe { std::env::set_var("NUCLEIC_SHELL_PARENT", &shell_id) }; + + let (tx, rx) = mpsc::sync_channel(CHANNEL_BOUND); + let observer = Observer { + tx, + next_token: AtomicU64::new(1), + next_seq: AtomicU64::new(1), + dropped: AtomicU64::new(0), + pending: Mutex::new(HashMap::new()), + }; + if OBSERVER.set(observer).is_err() { + return false; + } + + std::thread::Builder::new() + .name("nash-observe-flush".into()) + .spawn(move || flusher(cfg, shell_id, parent_shell, rx)) + .ok(); + + brush_core::gate::set_gate(Box::new(RecordingGate)); + true +} + +/// Flushes any buffered events, waiting up to a short deadline. Intended for +/// the process-exit hook; a no-op when observation is not installed. +pub fn flush_sync() { + let Some(obs) = OBSERVER.get() else { return }; + let (ack_tx, ack_rx) = mpsc::sync_channel(1); + if obs.tx.try_send(Msg::FlushSync(ack_tx)).is_ok() { + let _ = ack_rx.recv_timeout(EXIT_FLUSH_TIMEOUT); + } +} + +/// Records a compat-fallback event (docs/NASH.md §4.1): nash is about to +/// re-exec the real bash because it could not handle `input`. +pub fn report_fallback(reason: &str, input: &str) { + let Some(obs) = OBSERVER.get() else { return }; + let truncated: String = input.chars().take(4096).collect(); + obs.send(Event::Fallback { + seq: obs.seq(), + ts: now_millis(), + reason: reason.to_string(), + input: truncated, + }); + flush_sync(); +} diff --git a/nash/Cargo.toml b/nash/Cargo.toml index bedbd46..6629004 100644 --- a/nash/Cargo.toml +++ b/nash/Cargo.toml @@ -12,3 +12,8 @@ path = "src/main.rs" [dependencies] brush-shell = { path = "../../third_party/brush/brush-shell" } +brush-parser = { path = "../../third_party/brush/brush-parser" } +nash-observe = { path = "../nash-observe" } + +[dev-dependencies] +serde_json = "1" diff --git a/nash/src/main.rs b/nash/src/main.rs index c235190..7785888 100644 --- a/nash/src/main.rs +++ b/nash/src/main.rs @@ -1,13 +1,95 @@ //! nash — the Nucleic agent shell (see docs/NASH.md). //! -//! M0: a transparent embedding of the vendored brush shell, reusing its full -//! bash-compatible CLI (`-c`, scripts, stdin, set flags). The observation layer -//! (`nash-observe`) and the brush-core `Gate` patch land in M1; until then this -//! binary must behave byte-for-byte like upstream brush. +//! A thin binary around the vendored brush fork: installs the `nash-observe` +//! recording gate (when `NUCLEIC_SHELL_*` config is present), registers the +//! exit-time event flush, applies the compat fallback for `-c` input brush +//! cannot parse, and otherwise defers entirely to brush's bash-compatible CLI. + +use std::os::unix::process::CommandExt; fn main() { // docs/NASH.md §2: let scripts/tests detect they are running under nash. // SAFETY: called before any other threads exist. unsafe { std::env::set_var("NUCLEIC_NASH", "1") }; + + // Kill switch (docs/NASH.md §4.1): behave as the real bash, verbatim argv. + if std::env::var("NUCLEIC_NASH_DISABLE").as_deref() == Ok("1") { + exec_real_bash(); + // exec failed; fall through and run as nash rather than break the command. + } + + let observing = nash_observe::install_from_env(); + if observing { + brush_shell::entry::set_exit_hook(Box::new(nash_observe::flush_sync)); + } + + // Compat fallback (docs/NASH.md §4.1): if brush can't parse `-c` input + // that may still be valid bash, record a fallback event and re-exec the + // preserved real bash with identical argv so the command still runs. + if let Some(command) = dash_c_command() { + if let Err(parse_err) = try_parse(&command) { + if observing { + nash_observe::report_fallback(&format!("parse-error: {parse_err}"), &command); + } + exec_real_bash(); + // exec failed (no real bash available): let brush surface its own + // parse error so the user sees a real diagnostic. + } + } + brush_shell::entry::run(); } + +/// Extracts the command string from a bash-style `-c` invocation, handling +/// combined short options (`-lc `, `-elc `, …). +fn dash_c_command() -> Option { + let args: Vec = std::env::args().skip(1).collect(); + let mut iter = args.iter(); + while let Some(arg) = iter.next() { + if arg == "--" { + return None; // operands only from here on; no -c seen + } + if let Some(flags) = arg.strip_prefix('-') { + if !flags.is_empty() + && flags.chars().all(|c| c.is_ascii_alphanumeric()) + && flags.ends_with('c') + { + return iter.next().cloned(); + } + } + } + None +} + +fn try_parse(command: &str) -> Result<(), String> { + let reader = std::io::BufReader::new(command.as_bytes()); + let mut parser = brush_parser::Parser::new(reader, &brush_parser::ParserOptions::default()); + parser.parse_program().map(|_| ()).map_err(|e| e.to_string()) +} + +/// Replaces this process with the preserved real bash, keeping argv intact. +/// Returns only if every candidate exec fails. +fn exec_real_bash() { + let mut candidates: Vec = Vec::new(); + if let Ok(explicit) = std::env::var("NUCLEIC_REAL_BASH") { + candidates.push(explicit); + } + candidates.push("/usr/bin/bash.real".into()); + candidates.push("/bin/bash".into()); + + let self_path = std::env::current_exe().ok(); + let args: Vec = std::env::args().skip(1).collect(); + for candidate in candidates { + let path = std::path::Path::new(&candidate); + // Never re-exec ourselves (e.g. /bin/bash diverted to nash). + if !path.is_file() + || self_path + .as_deref() + .and_then(|p| path.canonicalize().ok().map(|c| c == p)) + .unwrap_or(false) + { + continue; + } + let _ = std::process::Command::new(path).args(&args).exec(); + } +} diff --git a/nash/tests/observe.rs b/nash/tests/observe.rs new file mode 100644 index 0000000..6a98c2f --- /dev/null +++ b/nash/tests/observe.rs @@ -0,0 +1,296 @@ +//! End-to-end observation tests: run the real nash binary against a fake +//! Nucleic control-socket sink and assert on the shell-batch events it posts +//! (docs/NASH.md §6.1, §10). + +use std::io::{Read, Write}; +use std::os::unix::net::UnixListener; +use std::process::Command; +use std::sync::mpsc; +use std::time::Duration; + +struct Sink { + dir: std::path::PathBuf, + rx: mpsc::Receiver<(Option, serde_json::Value)>, +} + +impl Sink { + /// Starts a unix-socket HTTP sink; each posted batch is parsed and sent + /// through the channel along with its Authorization header. + fn start() -> Self { + let dir = std::env::temp_dir().join(format!("nash-test-{}-{:?}", std::process::id(), std::time::Instant::now())); + std::fs::create_dir_all(&dir).unwrap(); + let socket_path = dir.join("control.sock"); + let listener = UnixListener::bind(&socket_path).unwrap(); + let (tx, rx) = mpsc::channel(); + std::thread::spawn(move || { + for stream in listener.incoming() { + let Ok(mut stream) = stream else { break }; + 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. + 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; + } + } + Err(_) => break, + } + } + }); + } + }); + Self { dir, rx } + } + + fn socket(&self) -> std::path::PathBuf { + self.dir.join("control.sock") + } + + /// Collects events from batches until `pred` is satisfied or ~5s elapse. + fn events_until( + &self, + pred: impl Fn(&[serde_json::Value]) -> bool, + ) -> Vec { + let mut events = Vec::new(); + let deadline = std::time::Instant::now() + Duration::from_secs(5); + while std::time::Instant::now() < deadline { + match self.rx.recv_timeout(Duration::from_millis(250)) { + Ok((auth, batch)) => { + assert_eq!( + auth.as_deref(), + Some("Bearer test-token"), + "batch must carry the bearer token" + ); + assert_eq!(batch["type"], "shell-batch"); + assert_eq!(batch["source"], "nash"); + assert_eq!(batch["sessionId"], "sess-1"); + assert!(batch["shellId"].is_string()); + events.extend(batch["events"].as_array().cloned().unwrap_or_default()); + if pred(&events) { + break; + } + } + Err(_) => { + if pred(&events) { + break; + } + } + } + } + events + } +} + +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()?; + let body = &buf[headers_end..]; + (body.len() >= len).then(|| &body[..len]) +} + +fn header(buf: &[u8], name: &str) -> Option { + let text = String::from_utf8_lossy(buf); + text.lines() + .take_while(|l| !l.is_empty()) + .find_map(|l| { + let (k, v) = l.split_once(':')?; + k.trim().eq_ignore_ascii_case(name).then(|| v.trim().to_string()) + }) +} + +fn run_nash(sink: &Sink, script: &str) -> std::process::ExitStatus { + Command::new(env!("CARGO_BIN_EXE_nash")) + .args(["-c", script]) + .env("NUCLEIC_SHELL_SOCKET", sink.socket()) + .env("NUCLEIC_HOOK_TOKEN", "test-token") + .env("NUCLEIC_SESSION_ID", "sess-1") + .env_remove("NUCLEIC_SHELL_PARENT") + .status() + .unwrap() +} + +fn find<'a>( + events: &'a [serde_json::Value], + kind: &str, + pred: impl Fn(&serde_json::Value) -> bool, +) -> Option<&'a serde_json::Value> { + events + .iter() + .find(|e| e["kind"] == kind && pred(e)) +} + +#[test] +fn exec_events_flow() { + let sink = Sink::start(); + let status = run_nash(&sink, "echo observed; false"); + assert_eq!(status.code(), Some(1)); + + let events = sink.events_until(|evs| { + find(evs, "exec", |e| e["argv"][0] == "echo").is_some() + && find(evs, "exec", |e| e["argv"][0] == "false").is_some() + }); + + let echo = find(&events, "exec", |e| e["argv"][0] == "echo").expect("echo exec event"); + assert_eq!(echo["argv"][1], "observed"); + assert_eq!(echo["exitCode"], 0); + assert!(echo["cwd"].is_string()); + assert!(echo["durationMs"].is_u64()); + + let false_ev = find(&events, "exec", |e| e["argv"][0] == "false").expect("false exec event"); + assert_eq!(false_ev["exitCode"], 1); +} + +#[test] +fn external_command_exit_codes() { + let sink = Sink::start(); + let status = run_nash(&sink, "/bin/echo ext; /usr/bin/env true; sh -c 'exit 7'"); + assert_eq!(status.code(), Some(7)); + + let events = sink.events_until(|evs| { + find(evs, "exec", |e| e["argv"][0] == "sh").is_some() + }); + let ext = find(&events, "exec", |e| e["argv"][0] == "/bin/echo").expect("external exec event"); + assert_eq!(ext["exitCode"], 0); + let sh = find(&events, "exec", |e| e["argv"][0] == "sh").expect("sh exec event"); + assert_eq!(sh["exitCode"], 7); +} + +#[test] +fn cd_and_export_events() { + let sink = Sink::start(); + let status = run_nash(&sink, "cd /tmp && export NASH_E2E=1 OTHER=2 && unset OTHER"); + assert_eq!(status.code(), Some(0)); + + let events = sink.events_until(|evs| { + find(evs, "cd", |_| true).is_some() + && find(evs, "export", |e| e["op"] == "unset").is_some() + }); + + let cd = find(&events, "cd", |_| true).expect("cd event"); + assert_eq!(cd["to"], "/tmp"); + assert!(cd["from"].is_string()); + + let export = find(&events, "export", |e| e["op"] == "export").expect("export event"); + let names: Vec<_> = export["names"].as_array().unwrap().iter().collect(); + assert_eq!(names.len(), 2); + assert_eq!(export["names"][0], "NASH_E2E"); + assert_eq!(export["names"][1], "OTHER"); + // Names only — values must never appear anywhere in the event. + assert!(!serde_json::to_string(export).unwrap().contains('=')); + + let unset = find(&events, "export", |e| e["op"] == "unset").expect("unset event"); + assert_eq!(unset["names"][0], "OTHER"); +} + +#[test] +fn pipeline_members_each_reported() { + let sink = Sink::start(); + let status = run_nash(&sink, "printf 'a\\nb\\n' | wc -l | cat > /dev/null"); + assert_eq!(status.code(), Some(0)); + + let events = sink.events_until(|evs| { + ["printf", "wc", "cat"] + .iter() + .all(|c| find(evs, "exec", |e| &e["argv"][0] == c).is_some()) + }); + for cmd in ["printf", "wc", "cat"] { + let ev = find(&events, "exec", |e| &e["argv"][0] == cmd) + .unwrap_or_else(|| panic!("{cmd} exec event missing")); + assert_eq!(ev["exitCode"], 0, "{cmd} should exit 0"); + } +} + +#[test] +fn command_not_found_reported() { + let sink = Sink::start(); + let status = run_nash(&sink, "definitely-not-a-command-xyz 2>/dev/null; true"); + assert_eq!(status.code(), Some(0)); + + let events = sink.events_until(|evs| { + find(evs, "exec", |e| e["argv"][0] == "definitely-not-a-command-xyz").is_some() + }); + let ev = find(&events, "exec", |e| e["argv"][0] == "definitely-not-a-command-xyz") + .expect("not-found exec event"); + assert_eq!(ev["exitCode"], 127); +} + +#[test] +fn parse_failure_falls_back_to_bash() { + let sink = Sink::start(); + let status = Command::new(env!("CARGO_BIN_EXE_nash")) + .args(["-c", "if then fi"]) + .env("NUCLEIC_SHELL_SOCKET", sink.socket()) + .env("NUCLEIC_HOOK_TOKEN", "test-token") + .env("NUCLEIC_SESSION_ID", "sess-1") + .env("NUCLEIC_REAL_BASH", "/bin/bash") + .env_remove("NUCLEIC_SHELL_PARENT") + .status() + .unwrap(); + // bash also rejects this input; what matters is bash's diagnostic + code. + assert_eq!(status.code(), Some(2)); + + let events = sink.events_until(|evs| find(evs, "fallback", |_| true).is_some()); + let fb = find(&events, "fallback", |_| true).expect("fallback event"); + assert_eq!(fb["input"], "if then fi"); + assert!(fb["reason"].as_str().unwrap().starts_with("parse-error")); +} + +#[test] +fn silent_without_config() { + // No NUCLEIC_SHELL_* config: nash must run normally and post nothing. + let out = Command::new(env!("CARGO_BIN_EXE_nash")) + .args(["-c", "echo quiet"]) + .env_remove("NUCLEIC_SHELL_SOCKET") + .env_remove("NUCLEIC_SHELL_HOOK_URL") + .env_remove("NUCLEIC_SHELL_SPOOL") + .output() + .unwrap(); + assert_eq!(out.status.code(), Some(0)); + assert_eq!(String::from_utf8_lossy(&out.stdout), "quiet\n"); +} + +#[test] +fn spool_fallback_when_socket_absent() { + let dir = std::env::temp_dir().join(format!("nash-spool-{}", std::process::id())); + let _ = std::fs::remove_dir_all(&dir); + let status = Command::new(env!("CARGO_BIN_EXE_nash")) + .args(["-c", "echo spooled"]) + .env("NUCLEIC_SHELL_SOCKET", "/nonexistent/control.sock") + .env("NUCLEIC_SHELL_SPOOL", &dir) + .env("NUCLEIC_SESSION_ID", "sess-1") + .env_remove("NUCLEIC_SHELL_PARENT") + .status() + .unwrap(); + assert_eq!(status.code(), Some(0)); + + let mut found = false; + for entry in std::fs::read_dir(&dir).expect("spool dir created") { + let content = std::fs::read_to_string(entry.unwrap().path()).unwrap(); + for line in content.lines() { + let batch: serde_json::Value = serde_json::from_str(line).unwrap(); + if batch["events"] + .as_array() + .unwrap() + .iter() + .any(|e| e["kind"] == "exec" && e["argv"][0] == "echo") + { + found = true; + } + } + } + assert!(found, "exec event spooled to disk"); + let _ = std::fs::remove_dir_all(&dir); +}