diff --git a/nash-observe/src/lib.rs b/nash-observe/src/lib.rs index 87d0a3b..817f69e 100644 --- a/nash-observe/src/lib.rs +++ b/nash-observe/src/lib.rs @@ -23,6 +23,14 @@ use serde::Serialize; const FLUSH_MAX_EVENTS: usize = 200; const FLUSH_MAX_AGE: Duration = Duration::from_millis(500); const CHANNEL_BOUND: usize = 4096; +/// How long a command must have been running before nash announces it as still in flight +/// (docs/NASH.md §5.1). Everything shorter is reported once, on exit: a shell that emitted a +/// start event for every `echo` would double the exec traffic to say "running" about commands +/// that were already over before the batch carrying it was flushed. +const ANNOUNCE_AFTER: Duration = Duration::from_millis(400); +/// How often the flusher re-checks for a command that has crossed `ANNOUNCE_AFTER`. Only ticks +/// while an un-announced command is actually in flight; an idle shell parks as it always did. +const ANNOUNCE_TICK: Duration = Duration::from_millis(250); 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). @@ -35,10 +43,34 @@ const PREVIEW_CAP: usize = 64 * 1024; #[derive(Serialize)] #[serde(tag = "kind", rename_all = "camelCase")] enum Event { + /// A command that has been running longer than [`ANNOUNCE_AFTER`] and has not returned + /// (docs/NASH.md §5.1). Emitted at most once per command, from the flusher's sweep rather + /// than the gate, so a short command costs no extra event. `token` is the gate token the + /// matching [`Event::Exec`] carries, which is what lets a host retire the in-flight row + /// instead of showing the command twice. + #[serde(rename = "exec-start", rename_all = "camelCase")] + ExecStart { + seq: u64, + ts: u64, + token: u64, + /// How long the command had *already* been running when this went out. A host that showed + /// elapsed time from the moment it heard about the command would under-report it by the + /// announce threshold plus however long the batch sat in the flusher; this is what lets + /// the row count from when the command actually started. + elapsed_ms: u64, + argv: Vec, + cwd: String, + #[serde(skip_serializing_if = "Option::is_none")] + pipeline: Option, + }, #[serde(rename = "exec", rename_all = "camelCase")] Exec { seq: u64, ts: u64, + /// The gate token this command ran under — the join key for the `exec-start` that may + /// have announced it (see [`Event::ExecStart`]). Always present; a command that returned + /// before `ANNOUNCE_AFTER` simply has no start event to join. + token: u64, argv: Vec, cwd: String, exit_code: i32, @@ -241,6 +273,11 @@ impl Config { enum Msg { Event(Event), FlushSync(mpsc::SyncSender<()>), + /// A command just started, so the flusher must re-arm its timeout for the in-flight sweep + /// (§5.1). Without it the thread stays parked on its idle wait — a shell whose only activity + /// is one long command sends nothing to wake it, which is exactly the case being announced. + /// Carries nothing and is never batched: a nudge, not an event. + Wake, } struct PendingExec { @@ -249,6 +286,9 @@ struct PendingExec { started: Instant, redirects: Vec, pipeline: Option, + /// Whether an [`Event::ExecStart`] has already gone out for this command. The sweep sets it + /// so a command that outlives many ticks is announced exactly once. + announced: bool, } struct Observer { @@ -456,9 +496,14 @@ impl brush_core::gate::Gate for RecordingGate { started: Instant::now(), redirects: ev.redirects, pipeline: ev.pipeline, + announced: false, }, ); } + // Best-effort by design: a dropped nudge (full channel) costs at most a late + // announcement, and a full channel means events are flowing anyway — which wakes the + // flusher just the same. + let _ = obs.tx.try_send(Msg::Wake); Some(token) })) .ok() @@ -517,6 +562,7 @@ impl brush_core::gate::Gate for RecordingGate { obs.send(Event::Exec { seq: exec_seq, ts, + token, argv: pending.argv, cwd: cwd_str, exit_code: ev.exit_code, @@ -710,6 +756,56 @@ fn deliver(cfg: &Config, shell_id: &str, body: &[u8]) { // Flusher thread // --------------------------------------------------------------------------- +/// What one sweep of the pending map found: the commands to announce as still running, and +/// whether anything is still too young to announce (which is what keeps the flusher ticking). +struct Sweep { + announce: Vec, + waiting: bool, +} + +/// Announce every command that has been in flight longer than [`ANNOUNCE_AFTER`] without +/// returning (docs/NASH.md §5.1). +/// +/// Runs on the flusher thread, never on a command's own path, so a slow or contended sweep can +/// only delay an event — never a command. Events are returned rather than sent through the +/// channel the flusher is itself draining: pushing them straight into the outgoing buffer keeps +/// a start event ahead of the exit event it belongs to (a command whose `on_exit` has already +/// removed it from the map is no longer here to announce). +fn sweep_running() -> Sweep { + let mut sweep = Sweep { announce: Vec::new(), waiting: false }; + let Some(obs) = OBSERVER.get() else { return sweep }; + let Ok(mut pending) = obs.pending.lock() else { return sweep }; + // Oldest first, so a batch that announces several commands at once still numbers them in the + // order they started (a `HashMap` would otherwise hand them over in an arbitrary one). + let mut ripe: Vec = pending + .iter() + .filter(|(_, exec)| !exec.announced) + .filter_map(|(token, exec)| { + if exec.started.elapsed() >= ANNOUNCE_AFTER { + Some(*token) + } else { + sweep.waiting = true; + None + } + }) + .collect(); + ripe.sort_unstable(); + for token in ripe { + let Some(exec) = pending.get_mut(&token) else { continue }; + exec.announced = true; + sweep.announce.push(Event::ExecStart { + seq: obs.seq(), + ts: now_millis(), + token, + elapsed_ms: u64::try_from(exec.started.elapsed().as_millis()).unwrap_or(0), + argv: exec.argv.clone(), + cwd: exec.cwd.to_string_lossy().into_owned(), + pipeline: exec.pipeline.map(PipelineRef::from), + }); + } + sweep +} + 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; @@ -736,11 +832,31 @@ fn flusher(cfg: Config, shell_id: String, parent_shell: Option, rx: mpsc }; loop { - let timeout = match first_at { + // Commands that have outlived `ANNOUNCE_AFTER` go out with (or ahead of) whatever else + // is buffered, so a host learns about a long build while it is still building. + let sweep = std::panic::catch_unwind(sweep_running) + .unwrap_or(Sweep { announce: Vec::new(), waiting: false }); + if !sweep.announce.is_empty() { + if buf.is_empty() { + first_at = Some(Instant::now()); + } + buf.extend(sweep.announce); + if buf.len() >= FLUSH_MAX_EVENTS { + flush(&mut buf); + first_at = None; + } + } + + let flush_due = match first_at { Some(t) => FLUSH_MAX_AGE.saturating_sub(t.elapsed()), None => Duration::from_secs(3600), }; - match rx.recv_timeout(timeout) { + // Tick only while something is in flight and not yet announced. Once every running + // command has been announced — and whenever none is — the thread parks as it always did. + // The tick shortens the *wait*, never the flush: waking early to re-sweep must not turn + // the locked 500 ms batching cadence into a 250 ms one for every shell running a command. + let ticking = sweep.waiting && ANNOUNCE_TICK < flush_due; + match rx.recv_timeout(if ticking { ANNOUNCE_TICK } else { flush_due }) { Ok(Msg::Event(ev)) => { if buf.is_empty() { first_at = Some(Instant::now()); @@ -756,9 +872,15 @@ fn flusher(cfg: Config, shell_id: String, parent_shell: Option, rx: mpsc first_at = None; let _ = ack.send(()); } + // Nothing to do but go round again: the sweep at the top of the loop is the work, and + // the timeout it computes is the point of having been woken. + Ok(Msg::Wake) => {} Err(mpsc::RecvTimeoutError::Timeout) => { - flush(&mut buf); - first_at = None; + // A tick is just "go re-sweep"; the batch is not due yet. + if !ticking { + flush(&mut buf); + first_at = None; + } } Err(mpsc::RecvTimeoutError::Disconnected) => { flush(&mut buf); diff --git a/nash/tests/observe.rs b/nash/tests/observe.rs index fc02e3d..2bf1c0a 100644 --- a/nash/tests/observe.rs +++ b/nash/tests/observe.rs @@ -165,8 +165,9 @@ fn header(buf: &[u8], name: &str) -> Option { }) } -fn run_nash(sink: &Sink, script: &str) -> std::process::ExitStatus { - Command::new(env!("CARGO_BIN_EXE_nash")) +fn nash_command(sink: &Sink, script: &str) -> Command { + let mut command = Command::new(env!("CARGO_BIN_EXE_nash")); + command .args(["-c", script]) .env("NUCLEIC_SHELL_SOCKET", sink.socket()) .env("NUCLEIC_HOOK_TOKEN", "test-token") @@ -174,9 +175,12 @@ fn run_nash(sink: &Sink, script: &str) -> std::process::ExitStatus { .env("NUCLEIC_SHELL_ENVIRONMENT_KIND", "agent") .env("NUCLEIC_SHELL_ENVIRONMENT_ID", "agent") .env("NUCLEIC_SHELL_ENVIRONMENT_LABEL", "Agent") - .env_remove("NUCLEIC_SHELL_PARENT") - .status() - .unwrap() + .env_remove("NUCLEIC_SHELL_PARENT"); + command +} + +fn run_nash(sink: &Sink, script: &str) -> std::process::ExitStatus { + nash_command(sink, script).status().unwrap() } fn find<'a>( @@ -208,6 +212,54 @@ fn exec_events_flow() { let false_ev = find(&events, "exec", |e| e["argv"][0] == "false").expect("false exec event"); assert_eq!(false_ev["exitCode"], 1); + + // A command that returns promptly is reported once, on exit. Announcing every command as + // "running" would double the exec traffic to say so about commands already long over. + assert!( + find(&events, "exec-start", |_| true).is_none(), + "fast commands must not be announced as running" + ); +} + +/// A command that has not returned is reported while it runs (docs/NASH.md §5.1) — otherwise the +/// host hears nothing at all during exactly the commands worth watching, and a five-minute build +/// first appears once it is over. +#[test] +fn long_running_command_announced_before_it_returns() { + let sink = Sink::start(); + // Spawned, not run to completion: the point of the event is that it arrives *mid-command*. + let mut child = nash_command(&sink, "sleep 3").spawn().unwrap(); + + let events = sink.events_until(|evs| find(evs, "exec-start", |e| e["argv"][0] == "sleep").is_some()); + let start = find(&events, "exec-start", |e| e["argv"][0] == "sleep") + .expect("exec-start while the command is still running"); + assert!( + child.try_wait().unwrap().is_none(), + "the announcement arrived only after the command had already returned" + ); + assert_eq!(start["argv"][1], "3"); + assert!(start["cwd"].is_string()); + // How long it had *already* been running, so a host counts elapsed time from the command's + // own start rather than from the moment the batch reached it. + let elapsed = start["elapsedMs"].as_u64().expect("elapsedMs on exec-start"); + assert!((400..3_000).contains(&elapsed), "elapsedMs was {elapsed}"); + // Nothing about how it ended — that is what the exec event is for. + assert!(start.get("exitCode").is_none()); + assert!(start.get("durationMs").is_none()); + let token = start["token"].as_u64().expect("gate token on exec-start"); + + child.wait().unwrap(); + let rest = sink.events_until(|evs| find(evs, "exec", |e| e["argv"][0] == "sleep").is_some()); + let done = find(&rest, "exec", |e| e["argv"][0] == "sleep").expect("exec event on return"); + // The join a host folds the two into one row on: same command run, same token. + assert_eq!(done["token"].as_u64(), Some(token)); + assert_eq!(done["exitCode"], 0); + // Announced once, however many sweeps the command outlived — a row that re-announced itself + // every quarter second would be a new row every quarter second. + let announcements = events.iter().chain(rest.iter()).filter(|e| { + e["kind"] == "exec-start" && e["argv"][0] == "sleep" + }); + assert_eq!(announcements.count(), 1); } #[test]