diff --git a/nash-observe/src/lib.rs b/nash-observe/src/lib.rs index 6f47daf..6377d76 100644 --- a/nash-observe/src/lib.rs +++ b/nash-observe/src/lib.rs @@ -82,6 +82,18 @@ enum Event { hash: String, preview_b64: String, }, + #[serde(rename = "pipe", rename_all = "camelCase")] + Pipe { + seq: u64, + ts: u64, + pipeline_id: u64, + from_index: u64, + to_index: u64, + bytes: u64, + truncated: bool, + hash: String, + preview_b64: String, + }, #[serde(rename = "dropped", rename_all = "camelCase")] Dropped { seq: u64, ts: u64, count: u64 }, } @@ -460,6 +472,25 @@ impl brush_core::gate::Gate for RecordingGate { }); })); } + + fn on_pipe(&self, ev: brush_core::gate::PipeEvent) { + let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let Some(obs) = OBSERVER.get() else { return }; + let capped = ev.captured.len().min(PREVIEW_CAP); + 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, + 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]), + }); + })); + } } // --------------------------------------------------------------------------- diff --git a/nash/tests/observe.rs b/nash/tests/observe.rs index 6d494d5..d8f2506 100644 --- a/nash/tests/observe.rs +++ b/nash/tests/observe.rs @@ -374,6 +374,65 @@ fn credential_shapes_redacted_in_previews() { assert_eq!(r["bytes"], 48); } +#[test] +fn pipe_links_captured_with_indices() { + let sink = Sink::start(); + let status = run_nash(&sink, "printf 'a\\nb\\nc\\n' | grep -n . | cat > /dev/null"); + assert_eq!(status.code(), Some(0)); + + let events = sink.events_until(|evs| { + evs.iter().filter(|e| e["kind"] == "pipe").count() >= 2 + }); + let pipes: Vec<_> = events.iter().filter(|e| e["kind"] == "pipe").collect(); + assert_eq!(pipes.len(), 2, "two links in a three-stage pipeline"); + // All links share one pipeline id. + assert_eq!(pipes[0]["pipelineId"], pipes[1]["pipelineId"]); + + let link0 = pipes.iter().find(|e| e["fromIndex"] == 0).expect("link 0->1"); + assert_eq!(link0["toIndex"], 1); + assert_eq!(decode_preview(link0), "a\nb\nc\n"); + assert_eq!(link0["bytes"], 6); + + let link1 = pipes.iter().find(|e| e["fromIndex"] == 1).expect("link 1->2"); + assert_eq!(link1["toIndex"], 2); + assert_eq!(decode_preview(link1), "1:a\n2:b\n3:c\n"); +} + +#[test] +fn pipe_sigpipe_consumer_exits_early() { + // `yes | head` — the consumer closes after N lines; the producer must get + // SIGPIPE/EPIPE and the whole thing must terminate (not hang), with the link + // still reported. + let sink = Sink::start(); + let status = run_nash(&sink, "yes | head -3 > /dev/null"); + assert_eq!(status.code(), Some(0)); + let events = sink.events_until(|evs| evs.iter().any(|e| e["kind"] == "pipe")); + assert!(events.iter().any(|e| e["kind"] == "pipe"), "link reported despite early consumer exit"); +} + +#[test] +fn pipe_preserves_binary_stream() { + // Byte-for-byte integrity across a tapped pipe: md5 through nash must match bash. + let dir = std::env::temp_dir().join(format!("nash-pipebin-{}", std::process::id())); + std::fs::create_dir_all(&dir).unwrap(); + let src = dir.join("rand.bin"); + let data: Vec = (0..100_000u32).map(|i| (i.wrapping_mul(2654435761) >> 16) as u8).collect(); + std::fs::write(&src, &data).unwrap(); + let script = format!("cat {} | cat | md5sum", src.display()); + + let sink = Sink::start(); + let nash_out = Command::new(env!("CARGO_BIN_EXE_nash")) + .args(["-c", &script]) + .env("NUCLEIC_SHELL_SOCKET", sink.socket()) + .env("NUCLEIC_SESSION_ID", "sess-1") + .env_remove("NUCLEIC_SHELL_PARENT") + .output() + .unwrap(); + let bash_out = Command::new("/bin/bash").args(["-c", &script]).output().unwrap(); + assert_eq!(nash_out.stdout, bash_out.stdout, "md5 through tapped pipe must match bash"); + let _ = std::fs::remove_dir_all(&dir); +} + #[test] fn redirect_preserves_seek_semantics() { // A seeking writer (dd with seek=) must see a real, seekable fd — the file