Merge nucleic/frosty-ember-raven-da8d into dev
This commit is contained in:
+126
-4
@@ -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<String>,
|
||||
cwd: String,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
pipeline: Option<PipelineRef>,
|
||||
},
|
||||
#[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<String>,
|
||||
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<brush_core::gate::RedirectRecord>,
|
||||
pipeline: Option<brush_core::gate::PipelineSlot>,
|
||||
/// 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<Event>,
|
||||
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<u64> = 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<String>, rx: mpsc::Receiver<Msg>) {
|
||||
let mut buf: Vec<Event> = Vec::new();
|
||||
let mut first_at: Option<Instant> = None;
|
||||
@@ -736,11 +832,31 @@ fn flusher(cfg: Config, shell_id: String, parent_shell: Option<String>, 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<String>, 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);
|
||||
|
||||
Reference in New Issue
Block a user