diff --git a/nash-observe/src/lib.rs b/nash-observe/src/lib.rs index 3b7838e..4610049 100644 --- a/nash-observe/src/lib.rs +++ b/nash-observe/src/lib.rs @@ -43,6 +43,10 @@ enum Event { cwd: String, exit_code: i32, duration_ms: u64, + /// Present only for a command running as a stage of `a | b | c`; omitted + /// otherwise so a standalone command's event keeps its old shape. + #[serde(skip_serializing_if = "Option::is_none")] + pipeline: Option, }, #[serde(rename = "cd", rename_all = "camelCase")] Cd { seq: u64, ts: u64, from: String, to: String }, @@ -98,6 +102,27 @@ enum Event { Dropped { seq: u64, ts: u64, count: u64 }, } +/// The pipeline stage an exec event belongs to (docs/NASH.md §5.1). `id` is the same +/// counter the `pipe` events carry, so the host can name the commands either side of +/// a link instead of showing bare stage indices. +#[derive(Serialize, Clone, Copy)] +#[serde(rename_all = "camelCase")] +struct PipelineRef { + id: u64, + index: u64, + len: u64, +} + +impl From for PipelineRef { + fn from(slot: brush_core::gate::PipelineSlot) -> Self { + Self { + id: slot.id, + index: slot.index as u64, + len: slot.len as u64, + } + } +} + #[derive(Serialize)] #[serde(rename_all = "camelCase")] struct Batch<'a> { @@ -190,6 +215,7 @@ struct PendingExec { cwd: PathBuf, started: Instant, redirects: Vec, + pipeline: Option, } struct Observer { @@ -368,6 +394,7 @@ impl brush_core::gate::Gate for RecordingGate { cwd: ev.cwd, started: Instant::now(), redirects: ev.redirects, + pipeline: ev.pipeline, }, ); } @@ -433,6 +460,7 @@ impl brush_core::gate::Gate for RecordingGate { cwd: cwd_str, exit_code: ev.exit_code, duration_ms: u64::try_from(pending.started.elapsed().as_millis()).unwrap_or(0), + pipeline: pending.pipeline.map(PipelineRef::from), }); // Data-flow read-back for this command's redirects (docs/NASH.md §5.2). diff --git a/nash/tests/observe.rs b/nash/tests/observe.rs index 4cc5198..aad4029 100644 --- a/nash/tests/observe.rs +++ b/nash/tests/observe.rs @@ -467,6 +467,93 @@ fn pipe_links_captured_with_indices() { assert_eq!(decode_preview(link1), "1:a\n2:b\n3:c\n"); } +/// Each stage's exec event carries the slot that joins it to the links either side, so a viewer +/// can name a link's producer and consumer (`printf … | grep …`) instead of showing bare indices. +#[test] +fn pipeline_stages_carry_their_slot_on_exec() { + let sink = Sink::start(); + let status = run_nash(&sink, "printf 'a\\nb\\n' | grep -n . | cat > /dev/null"); + assert_eq!(status.code(), Some(0)); + + let events = sink.events_until(|evs| { + ["printf", "grep", "cat"] + .iter() + .all(|c| find(evs, "exec", |e| &e["argv"][0] == c).is_some()) + && evs.iter().any(|e| e["kind"] == "pipe") + }); + + let slot = |cmd: &str| { + find(&events, "exec", |e| &e["argv"][0] == cmd) + .unwrap_or_else(|| panic!("{cmd} exec event missing"))["pipeline"] + .clone() + }; + for (cmd, index) in [("printf", 0), ("grep", 1), ("cat", 2)] { + assert_eq!(slot(cmd)["index"], index, "{cmd} is stage {index}"); + assert_eq!(slot(cmd)["len"], 3, "{cmd} sees a three-stage pipeline"); + } + // The stages and the links they connect share one id — that join is the whole point. + let pipeline_id = slot("printf")["id"].clone(); + assert!(pipeline_id.is_number()); + assert_eq!(slot("grep")["id"], pipeline_id); + assert_eq!(slot("cat")["id"], pipeline_id); + let pipe = find(&events, "pipe", |_| true).expect("pipe event"); + assert_eq!(pipe["pipelineId"], pipeline_id); +} + +/// A standalone command has no slot at all, so the field stays absent rather than reporting a +/// meaningless stage 0 of 1. +#[test] +fn standalone_command_carries_no_pipeline_slot() { + let sink = Sink::start(); + assert_eq!(run_nash(&sink, "echo alone").code(), Some(0)); + + let events = sink.events_until(|evs| find(evs, "exec", |e| e["argv"][0] == "echo").is_some()); + let echo = find(&events, "exec", |e| e["argv"][0] == "echo").expect("echo exec event"); + assert!(echo.get("pipeline").is_none(), "no pipeline field on a lone command"); +} + +/// A command substitution runs during a stage's *expansion*, before the stage's own command, so +/// it must not claim the stage — otherwise `grep "$(printf x)" … | wc -l` reports `printf x` as +/// what produced the piped bytes. +#[test] +fn command_substitution_does_not_claim_its_stage() { + let sink = Sink::start(); + let status = run_nash(&sink, "grep \"$(printf a)\" /etc/hostname | cat > /dev/null"); + assert_eq!(status.code(), Some(0)); + + let events = sink.events_until(|evs| { + find(evs, "exec", |e| e["argv"][0] == "grep").is_some() + && find(evs, "exec", |e| e["argv"][0] == "cat").is_some() + }); + let printf = find(&events, "exec", |e| e["argv"][0] == "printf").expect("printf exec event"); + assert!( + printf.get("pipeline").is_none(), + "the substituted command is an argument to stage 0, not stage 0" + ); + let grep = find(&events, "exec", |e| e["argv"][0] == "grep").expect("grep exec event"); + assert_eq!(grep["pipeline"]["index"], 0, "grep is what stage 0 actually ran"); +} + +/// A stage that is a compound command attributes the commands *inside* it to that same stage — +/// they really are what that stage ran — while a nested pipeline takes over with its own id. +#[test] +fn nested_commands_inherit_their_enclosing_stage() { + let sink = Sink::start(); + let status = run_nash(&sink, "{ echo one; echo two; } | cat > /dev/null"); + assert_eq!(status.code(), Some(0)); + + let events = sink.events_until(|evs| { + find(evs, "exec", |e| e["argv"][1] == "two").is_some() + && find(evs, "exec", |e| e["argv"][0] == "cat").is_some() + }); + for arg in ["one", "two"] { + let ev = find(&events, "exec", |e| e["argv"][1] == arg).expect("echo exec event"); + assert_eq!(ev["pipeline"]["index"], 0, "`echo {arg}` runs inside stage 0"); + } + let cat = find(&events, "exec", |e| e["argv"][0] == "cat").expect("cat exec event"); + assert_eq!(cat["pipeline"]["index"], 1); +} + #[test] fn pipe_sigpipe_consumer_exits_early() { // `yes | head` — the consumer closes after N lines; the producer must get