amont_runtime/live.rs
1//! One check, one block — and while it runs, one line: per-check output
2//! capture plus a live progress region for the concurrent stage.
3//!
4//! Twenty checks used to print straight to inherited stdio from their own
5//! threads, so two failing linters shuffled their lines together and the
6//! reader un-shuffled them by hand — the dispatcher's roll-up existed partly
7//! to apologise for it. Now every check writes into its own slot, and a
8//! completed check's output reaches stdout as ONE locked write: contiguous,
9//! whatever the other nineteen were doing.
10//!
11//! Three writers feed a slot:
12//!
13//! 1. The check's own thread, through [`say`] — which is what
14//! `common::ok/fail/warn` call. A thread with no slot installed (commit-msg,
15//! `amont install`, the dispatcher itself) prints directly, exactly as
16//! before; nothing outside a stage changes.
17//! 2. A captured child's reader threads, through [`Stage::append_raw`] —
18//! they are not the check's thread, so the thread-local cannot carry the
19//! routing; the `Arc` is captured before the spawn instead.
20//! 3. Nobody else. The dispatcher's own lines (skips, pins, the roll-up)
21//! happen strictly before or after the fan-out and stay direct.
22//!
23//! Order across checks is COMPLETION order — deterministic per block, not
24//! per stage, which is the same nondeterminism the interleaved version had
25//! without the shuffling. `amont.progress false` switches the whole
26//! mechanism off and restores raw streaming for anyone who wants to watch a
27//! tool write in real time.
28//!
29//! # The region
30//!
31//! When stderr is a real terminal ([`watching`]) the stage also paints a
32//! live region UNDER the finished blocks: one line per running check —
33//! braille spinner, name, elapsed — repainted every 80ms by a ticker
34//! thread, shrinking as checks finish, gone without a trace when the stage
35//! ends. Blocks go to stdout, the region to stderr; both feed one tty, and
36//! every write to either happens under the same [`Stage::out`] lock, so a
37//! block never tears a repaint in half. Piped, redirected, `TERM=dumb`, or
38//! CI: [`watching`] is false, no ticker starts, and the region costs
39//! nothing — which is also why the test suite (piped stdio throughout)
40//! exercises capture but never the paint.
41
42use std::cell::RefCell;
43use std::io::{IsTerminal, Write};
44use std::sync::atomic::{AtomicBool, Ordering};
45use std::sync::{Arc, Mutex, Weak};
46use std::time::Instant;
47
48/// The fleet spinner's frames (progress.rs) — cycled by elapsed time, so a
49/// frame needs no state beyond the clock.
50const FRAMES: [char; 10] = ['⠋', '⠙', '⠹', '⠸', '⠼', '⠴', '⠦', '⠧', '⠇', '⠏'];
51
52/// The region never grows past this many check lines; the rest fold into
53/// one `… and N more`. Twelve is the whole default fleet on one screen.
54const MAX_LINES: usize = 12;
55
56/// One check's place in the stage.
57struct Slot {
58 /// Sanitised at [`Stage::begin`]: a manifest-declared name is
59 /// repo-derived text and the region writes it to a live terminal.
60 name: String,
61 /// Restamped by [`Stage::enter`], so a serial stage (pre-push) times
62 /// each check from its own start, not the stage's.
63 started: Instant,
64 /// The last byte or line that landed in `buf` — what the region's
65 /// `quiet` figure and the heartbeat's `last output` read.
66 last_output: Instant,
67 /// Elapsed seconds at which the non-tty heartbeat next speaks.
68 next_beat: u64,
69 buf: Vec<u8>,
70 /// Entered and not yet finished — the region shows exactly these.
71 running: bool,
72 done: bool,
73 /// The command this check is waiting on, while it runs: its output
74 /// clock and what its CPU is doing — the same object the kill decision
75 /// reads, so the displays can never disagree with it (ADR-0008).
76 activity: Option<Arc<crate::hooks::common::Activity>>,
77}
78
79/// A running stage: the slots, and the one lock every terminal write inside
80/// the stage goes through.
81pub struct Stage {
82 slots: Mutex<Vec<Slot>>,
83 /// Serialises block emission and region repaints; the value is how many
84 /// region lines are currently painted (what an erase must remove).
85 out: Mutex<usize>,
86 /// Painting at all? [`enabled`] && [`watching`], decided once at begin.
87 live: bool,
88 /// Is this the PUSH stage? Read by the heartbeat, which has something to
89 /// say about a long gate there and nothing to say about one at commit
90 /// time — see [`beat_line`]. Derived from the names, which already
91 /// carry the trigger.
92 on_push: bool,
93 stop: AtomicBool,
94}
95
96thread_local! {
97 /// Where [`say`] routes on THIS thread: a stage and a slot index.
98 static SINK: RefCell<Option<(Arc<Stage>, usize)>> = const { RefCell::new(None) };
99}
100
101impl Stage {
102 /// A stage over `names`, in dispatch order. Does nothing visible until
103 /// checks start entering (the region) or finishing (the blocks).
104 pub fn begin(settings: &crate::config::Settings, names: &[&str]) -> Arc<Stage> {
105 let now = Instant::now();
106 let stage = Arc::new(Stage {
107 slots: Mutex::new(
108 names
109 .iter()
110 .map(|n| Slot {
111 // Every name in a stage carries the stage's own
112 // prefix ("pre-commit-clippy"); the region drops it
113 // — twelve identical prefixes say nothing.
114 name: crate::ui::sanitize(
115 n.strip_prefix("pre-commit-")
116 .or_else(|| n.strip_prefix("pre-push-"))
117 .unwrap_or(n),
118 ),
119 started: now,
120 last_output: now,
121 next_beat: HEARTBEAT_SECS,
122 buf: Vec::new(),
123 running: false,
124 done: false,
125 activity: None,
126 })
127 .collect(),
128 ),
129 out: Mutex::new(0),
130 live: enabled(settings) && watching(),
131 // The names arrive fully qualified and the loop above has
132 // already had to strip the trigger to display them, so the
133 // stage can answer this without dispatch passing anything in.
134 on_push: names.iter().any(|n| n.starts_with("pre-push-")),
135 stop: AtomicBool::new(false),
136 });
137 if stage.live {
138 // The ticker holds a Weak: the stage dropping is what ends it,
139 // so a paint can never outlive the region's owner.
140 let weak = Arc::downgrade(&stage);
141 let own = settings.for_thread();
142 let _ = std::thread::Builder::new()
143 .name("amont-live".into())
144 .spawn(move || tick(own, weak));
145 } else if enabled(settings) {
146 // Nobody is watching a terminal — an agent, CI, a pipe — and a
147 // captured check shows nothing until it finishes. The heartbeat
148 // is the one line a minute that says it is alive, which is the
149 // difference between "wait" and "kill it" for whoever is on the
150 // other end of the pipe.
151 let weak = Arc::downgrade(&stage);
152 let own = settings.for_thread();
153 let _ = std::thread::Builder::new()
154 .name("amont-heartbeat".into())
155 .spawn(move || heartbeat(own, weak));
156 }
157 stage
158 }
159
160 /// Route this thread's [`say`] calls into slot `idx` until the guard
161 /// drops. Installed by the dispatcher around each `check.run`. Also
162 /// starts the slot's clock and puts it in the region.
163 pub fn enter(self: &Arc<Stage>, idx: usize) -> SinkGuard {
164 {
165 let mut slots = self.slots.lock().unwrap_or_else(|p| p.into_inner());
166 if let Some(slot) = slots.get_mut(idx) {
167 slot.running = true;
168 slot.started = Instant::now();
169 slot.last_output = slot.started;
170 slot.next_beat = HEARTBEAT_SECS;
171 }
172 }
173 SINK.with(|s| *s.borrow_mut() = Some((Arc::clone(self), idx)));
174 SinkGuard
175 }
176
177 /// Show slot `idx`'s spawned command in the displays while the returned
178 /// guard lives. A check that runs several commands in turn attaches each
179 /// one; between them the slot falls back to its own output clock.
180 pub fn attach(
181 self: &Arc<Stage>,
182 idx: usize,
183 activity: Arc<crate::hooks::common::Activity>,
184 ) -> AttachGuard {
185 let mut slots = self.slots.lock().unwrap_or_else(|p| p.into_inner());
186 if let Some(slot) = slots.get_mut(idx) {
187 slot.activity = Some(activity);
188 }
189 AttachGuard {
190 stage: Arc::clone(self),
191 idx,
192 }
193 }
194
195 /// Append raw bytes (a captured child's output) to slot `idx`.
196 pub fn append_raw(&self, idx: usize, bytes: &[u8]) {
197 let mut slots = self.slots.lock().unwrap_or_else(|p| p.into_inner());
198 if let Some(slot) = slots.get_mut(idx) {
199 if !slot.done {
200 slot.buf.extend_from_slice(bytes);
201 slot.last_output = Instant::now();
202 }
203 }
204 }
205
206 fn append_line(&self, idx: usize, line: &str) {
207 let mut slots = self.slots.lock().unwrap_or_else(|p| p.into_inner());
208 if let Some(slot) = slots.get_mut(idx) {
209 if !slot.done {
210 slot.buf.extend_from_slice(line.as_bytes());
211 slot.buf.push(b'\n');
212 slot.last_output = Instant::now();
213 }
214 }
215 }
216
217 /// The check is over: emit everything it said as ONE contiguous write,
218 /// with the region lifted out of the way first and repainted after —
219 /// blocks pile up above, spinners stay below.
220 ///
221 /// Called by the dispatcher after `check.run` returns (still on the
222 /// check's thread, so a torn-down thread cannot strand a buffer — the
223 /// same `catch_unwind` that feeds the dead-check outcome runs first).
224 pub fn finish(&self, settings: &crate::config::Settings, idx: usize) {
225 let block = {
226 let mut slots = self.slots.lock().unwrap_or_else(|p| p.into_inner());
227 let Some(slot) = slots.get_mut(idx) else {
228 return;
229 };
230 slot.done = true;
231 slot.running = false;
232 std::mem::take(&mut slot.buf)
233 };
234 if block.is_empty() && !self.live {
235 return;
236 }
237 let mut drawn = self.out.lock().unwrap_or_else(|p| p.into_inner());
238 if !block.is_empty() {
239 if *drawn > 0 {
240 let mut err = std::io::stderr().lock();
241 let _ = write!(err, "\x1b[{}A\x1b[J", *drawn);
242 let _ = err.flush();
243 *drawn = 0;
244 }
245 let stdout = std::io::stdout();
246 let mut handle = stdout.lock();
247 let _ = handle.write_all(&block);
248 let _ = handle.flush();
249 }
250 self.repaint(settings, &mut drawn);
251 }
252
253 /// Erase and redraw the region in one stderr write. Lock order is
254 /// `out` → `slots`, everywhere — never the reverse.
255 fn repaint(&self, settings: &crate::config::Settings, drawn: &mut usize) {
256 if !self.live {
257 return;
258 }
259 let entries: Vec<Row> = {
260 let slots = self.slots.lock().unwrap_or_else(|p| p.into_inner());
261 let now = Instant::now();
262 slots
263 .iter()
264 .filter(|s| s.running && !s.done)
265 .map(|s| row_of(s, now, now.duration_since(s.started).as_secs_f64()))
266 .collect()
267 };
268 let text = region(&entries, term_width(), budgets(settings));
269 let mut paint = String::new();
270 if *drawn > 0 {
271 paint.push_str(&format!("\x1b[{}A\x1b[J", *drawn));
272 }
273 paint.push_str(&text);
274 if paint.is_empty() {
275 return;
276 }
277 let mut err = std::io::stderr().lock();
278 let _ = err.write_all(paint.as_bytes());
279 let _ = err.flush();
280 *drawn = text.matches('\n').count();
281 }
282}
283
284impl Drop for Stage {
285 /// The stage's end erases whatever the region still shows — a Block
286 /// verdict, a panic on the dispatcher path, anything: no spinner junk
287 /// above the roll-up. (`get_mut`: dropping proves no other thread holds
288 /// the stage, so the locks are free.)
289 fn drop(&mut self) {
290 self.stop.store(true, Ordering::Relaxed);
291 if !self.live {
292 return;
293 }
294 let drawn = self.out.get_mut().unwrap_or_else(|p| p.into_inner());
295 if *drawn > 0 {
296 let mut err = std::io::stderr().lock();
297 let _ = write!(err, "\x1b[{}A\x1b[J", *drawn);
298 let _ = err.flush();
299 *drawn = 0;
300 }
301 }
302}
303
304/// The ticker: repaint every 80ms until the stage drops or tells it to
305/// stop. Holds only a `Weak`, so it can never keep a finished stage alive.
306/// Owns its `Settings` (see [`crate::config::Settings::for_thread`]): a
307/// spawned thread is `'static`, and the budgets must be read lazily, not
308/// pre-resolved at the spawn.
309fn tick(settings: crate::config::Settings, weak: Weak<Stage>) {
310 loop {
311 std::thread::sleep(std::time::Duration::from_millis(80));
312 let Some(stage) = weak.upgrade() else { return };
313 if stage.stop.load(Ordering::Relaxed) {
314 return;
315 }
316 let mut drawn = stage.out.lock().unwrap_or_else(|p| p.into_inner());
317 stage.repaint(&settings, &mut drawn);
318 }
319}
320
321/// One running check, as the region and the heartbeat see it.
322#[derive(Debug, Clone)]
323pub struct Row {
324 pub name: String,
325 /// Seconds since the check entered.
326 pub elapsed: f64,
327 /// Seconds since it last wrote anything.
328 pub quiet: f64,
329 /// Seconds it has been silent AND idle on CPU — what the silence budget
330 /// is judged against. Equal to `quiet` when CPU is not sampled.
331 pub still: f64,
332 pub cpu: RowCpu,
333}
334
335/// What a running check's CPU is doing, as far as the displays may say.
336#[derive(Debug, Clone, Copy, PartialEq, Eq)]
337pub enum RowCpu {
338 /// No spawned command is attached (an in-process check, or between two
339 /// commands): nothing to say.
340 None,
341 /// A command is attached but its CPU is not sampled here
342 /// (`amont.idleCpuCredit false`, no silence budget, or the platform):
343 /// silence alone counts.
344 NotSampled,
345 /// Sampled, but nothing fresh to report (not quiet long enough yet, or
346 /// the last snapshot was incomplete).
347 Unmeasured,
348 /// Measurably working, at this many thousandths of a core.
349 Busy(u32),
350 /// Measured under the busy threshold.
351 Idle,
352}
353
354/// A slot as a [`Row`], reading its attached command's clocks when there is
355/// one. The quiet figure is the more recent of the slot's own lines and the
356/// command's bytes (a captured command writes to one and not the other).
357fn row_of(s: &Slot, now: Instant, elapsed: f64) -> Row {
358 use crate::hooks::common::{CpuState, BUSY_MILLI_CORES};
359 let slot_quiet = now.duration_since(s.last_output).as_secs_f64();
360 let (quiet, still, cpu) = match &s.activity {
361 None => (slot_quiet, slot_quiet, RowCpu::None),
362 Some(a) => {
363 let quiet = slot_quiet.min(a.quiet_for().as_secs_f64());
364 let still = quiet.min(a.still_for().as_secs_f64());
365 let cpu = match (a.cpu_state(), a.fresh_rate()) {
366 (CpuState::Off, _) => RowCpu::NotSampled,
367 (_, Some(r)) if r >= BUSY_MILLI_CORES => RowCpu::Busy(r),
368 (_, Some(_)) => RowCpu::Idle,
369 (_, None) => RowCpu::Unmeasured,
370 };
371 (quiet, still, cpu)
372 }
373 };
374 Row {
375 name: s.name.clone(),
376 elapsed,
377 quiet,
378 still,
379 cpu,
380 }
381}
382
383/// The two clocks, as the region annotates them: `(idle, ceiling)` in
384/// seconds, `0` for off.
385#[derive(Debug, Clone, Copy)]
386pub struct Budgets {
387 pub idle: u64,
388 pub ceiling: u64,
389}
390
391fn budgets(settings: &crate::config::Settings) -> Budgets {
392 Budgets {
393 idle: crate::hooks::common::idle_timeout(settings),
394 ceiling: crate::hooks::common::check_timeout(settings),
395 }
396}
397
398/// How long a check must be quiet before the region says so. A test suite
399/// pauses this long between crates without anything being wrong; past it,
400/// the reader wants to know the silence is being counted.
401const QUIET_NOTE_SECS: f64 = 30.0;
402
403/// The non-tty heartbeat's period: one line a minute per running check.
404const HEARTBEAT_SECS: u64 = 60;
405
406/// Elapsed time in a fixed six-column figure: ` 3.2s` under a minute,
407/// `8m12s` and `1h02m` above, so the column stays aligned as the suite
408/// crosses the minute.
409fn elapsed_column(secs: f64) -> String {
410 if secs < 60.0 {
411 format!("{secs:>5.1}s")
412 } else {
413 format!("{:>6}", crate::hooks::common::human_secs(secs as u64))
414 }
415}
416
417/// The region's text: one `⠹ name 12.3s` line per running check, capped at
418/// [`MAX_LINES`] plus a `… and N more` overflow line. Pure — the ticker is
419/// a thin shell around this, and the tests drive it directly.
420///
421/// Two annotations, each only when it carries news: `· quiet 45s/2m` once
422/// a check has been silent past [`QUIET_NOTE_SECS`] (with the silence
423/// budget it is counting toward, when there is one), and `· 48m/60m` once
424/// elapsed passes 80% of the ceiling — the cliff, shown before the fall.
425fn region(entries: &[Row], width: usize, budgets: Budgets) -> String {
426 if entries.is_empty() {
427 return String::new();
428 }
429 let pad = entries
430 .iter()
431 .take(MAX_LINES)
432 .map(|r| r.name.chars().count())
433 .max()
434 .unwrap_or(0);
435 let mut out = String::new();
436 for row in entries.iter().take(MAX_LINES) {
437 let frame = FRAMES[((row.elapsed * 10.0) as usize) % FRAMES.len()];
438 let name = &row.name;
439 let mut line = format!("{frame} {name:<pad$} {}", elapsed_column(row.elapsed));
440 if row.quiet >= QUIET_NOTE_SECS {
441 let quiet = crate::hooks::common::human_secs(row.quiet as u64);
442 match row.cpu {
443 // Working: no countdown — no kill is coming — just how hard.
444 RowCpu::Busy(m) => line.push_str(&format!(
445 " · quiet {quiet} · {}",
446 crate::hooks::common::cores(m)
447 )),
448 _ if budgets.idle == 0 => line.push_str(&format!(" · quiet {quiet}")),
449 // The countdown counts what the kill decision counts: the
450 // still-time, which only differs from the silence once CPU
451 // work has pushed it back.
452 RowCpu::Idle | RowCpu::Unmeasured if row.quiet - row.still >= 1.0 => {
453 line.push_str(&format!(
454 " · quiet {quiet} · idle {}/{}",
455 crate::hooks::common::human_secs(row.still as u64),
456 crate::hooks::common::human_secs(budgets.idle)
457 ))
458 }
459 _ => line.push_str(&format!(
460 " · quiet {quiet}/{}",
461 crate::hooks::common::human_secs(budgets.idle)
462 )),
463 }
464 }
465 if budgets.ceiling > 0 && row.elapsed >= 0.8 * budgets.ceiling as f64 {
466 line.push_str(&format!(
467 " · {}/{}",
468 crate::hooks::common::human_secs(row.elapsed as u64),
469 crate::hooks::common::human_secs(budgets.ceiling)
470 ));
471 }
472 if line.chars().count() > width {
473 out.extend(line.chars().take(width));
474 } else {
475 out.push_str(&line);
476 }
477 out.push('\n');
478 }
479 if entries.len() > MAX_LINES {
480 out.push_str(&format!("… and {} more\n", entries.len() - MAX_LINES));
481 }
482 out
483}
484
485/// The heartbeat: once a minute, for each check still running, one plain
486/// line on stderr — elapsed, and how long since it last said anything.
487/// Not a region: nothing is erased or repainted, because nobody is looking
488/// at a cursor; whoever reads this reads a log.
489///
490/// The first beat for a check also names the two budgets, once, so the
491/// reader can tell how far it is from being killed without opening the
492/// docs. Written under the same `out` lock as the blocks, so a beat never
493/// lands inside one.
494/// Owns its `Settings` for the same reason [`tick`] does.
495fn heartbeat(settings: crate::config::Settings, weak: Weak<Stage>) {
496 loop {
497 std::thread::sleep(std::time::Duration::from_secs(1));
498 let Some(stage) = weak.upgrade() else { return };
499 if stage.stop.load(Ordering::Relaxed) {
500 return;
501 }
502 let due: Vec<(Row, bool)> = {
503 let mut slots = stage.slots.lock().unwrap_or_else(|p| p.into_inner());
504 let now = Instant::now();
505 let mut due = Vec::new();
506 for s in slots.iter_mut().filter(|s| s.running && !s.done) {
507 let elapsed = now.duration_since(s.started).as_secs();
508 if elapsed >= s.next_beat {
509 let first = s.next_beat == HEARTBEAT_SECS;
510 s.next_beat += HEARTBEAT_SECS;
511 due.push((row_of(s, now, elapsed as f64), first));
512 }
513 }
514 due
515 };
516 if due.is_empty() {
517 continue;
518 }
519 let text: String = due
520 .iter()
521 .map(|(row, first)| beat_line(row, *first, budgets(&settings), stage.on_push))
522 .collect();
523 let _guard = stage.out.lock().unwrap_or_else(|p| p.into_inner());
524 let mut err = std::io::stderr().lock();
525 let _ = err.write_all(text.as_bytes());
526 let _ = err.flush();
527 }
528}
529
530/// One heartbeat line. Pure, for the tests.
531///
532/// On the FIRST beat of a PUSH gate it also names something no other part of
533/// the system is placed to explain. `git push` opens its connection to the
534/// remote, reads the remote refs — which is where the `pre-push` hook's own
535/// stdin comes from — and only then calls the hook. The connection is
536/// therefore already open and goes idle for exactly as long as the gate
537/// runs, and a remote may close it before the gate finishes. git then
538/// reports `Connection reset by peer`, which reads as a network fault and
539/// says nothing about the seven minutes that caused it.
540///
541/// The note does NOT recommend ssh keepalive, and that omission is
542/// deliberate: `ServerAliveInterval 60` was already in force on the machine
543/// where this was diagnosed, and GitHub reset the connection anyway.
544/// Whatever the remote is measuring, it is not packets. Recommending it
545/// would be a confident instruction to change a setting that is probably
546/// already on and cannot help, so the note says so and points at the thing
547/// that does work.
548///
549/// Only on a first beat, so it is said once; only on a push, so a commit
550/// gate never hears it. A first beat is a check that has already run a full
551/// minute, which is the population at risk — no threshold to invent.
552fn beat_line(row: &Row, first: bool, budgets: Budgets, on_push: bool) -> String {
553 use crate::hooks::common::human_secs;
554 // The prefix is byte-for-byte what it always was: log readers grep it.
555 // What CPU sampling adds goes after it.
556 let mut line = format!(
557 " … {} still running: {}, last output {} ago",
558 row.name,
559 human_secs(row.elapsed as u64),
560 human_secs(row.quiet as u64)
561 );
562 match row.cpu {
563 RowCpu::Busy(m) => line.push_str(&format!(", busy {}", crate::hooks::common::cores(m))),
564 RowCpu::Idle => line.push_str(&format!(", CPU idle {}", human_secs(row.still as u64))),
565 RowCpu::Unmeasured => line.push_str(", CPU unmeasured"),
566 RowCpu::None | RowCpu::NotSampled => {}
567 }
568 if first {
569 let idle = match budgets.idle {
570 0 => "off".to_string(),
571 s => human_secs(s),
572 };
573 let ceiling = match budgets.ceiling {
574 0 => "off".to_string(),
575 s => human_secs(s),
576 };
577 match row.cpu {
578 RowCpu::Busy(_) | RowCpu::Idle | RowCpu::Unmeasured if budgets.idle > 0 => {
579 line.push_str(&format!(
580 " (killed after {idle} with no output and under 0.1 core of CPU, or \
581 {ceiling} in total — amont.idleTimeout / amont.timeout)"
582 ))
583 }
584 _ => line.push_str(&format!(
585 " (killed after {idle} of silence or {ceiling} in total — amont.idleTimeout / amont.timeout)"
586 )),
587 }
588 if row.cpu == RowCpu::NotSampled && budgets.idle > 0 {
589 line.push_str("; CPU not sampled here, silence alone counts");
590 }
591 if on_push {
592 // `concat!`, not a `\`-continued literal: a continuation keeps
593 // the next line's indentation, which turns the message into runs
594 // of spaces. Each line is its own literal and the newlines are
595 // written down, so what is here is what a reader sees.
596 line.push_str(concat!(
597 "\n git opened its connection to the remote before calling this",
598 "\n gate, and it stays idle until the gate finishes. A remote may",
599 "\n close it first — GitHub does — and the push then fails with",
600 "\n \"Connection reset by peer\", naming the network rather than the",
601 "\n wait. ssh keepalive does not prevent this.",
602 "\n Declaring this check at pre-commit moves it off the push path —",
603 "\n see \"Moving a gate entry earlier\" in the docs.",
604 ));
605 }
606 }
607 line.push('\n');
608 line
609}
610
611/// `$COLUMNS` when it is exported and sane, else a conservative 80 — the
612/// region's lines are short and an ioctl is not worth its portability. 80,
613/// not wider: shells rarely export `COLUMNS`, and a region line longer than
614/// the real terminal wraps, which breaks the erase arithmetic.
615///
616/// `pub` is now wider than it needs to be — the out-of-crate caller that
617/// justified it, `amont-agent`, is its own project and carries its own copy.
618/// Left public rather than narrowed in the same change that removed it.
619pub fn term_width() -> usize {
620 std::env::var("COLUMNS")
621 .ok()
622 .and_then(|c| c.parse::<usize>().ok())
623 .filter(|w| *w >= 20)
624 .unwrap_or(80)
625}
626
627/// Emits slot `idx`'s block when dropped — however the check's closure
628/// exits, a panic included: the partial output of a check that died still
629/// reaches the reader, above the dead-check verdict the runner fills in.
630pub struct FinishOnDrop<'a> {
631 stage: &'a Stage,
632 idx: usize,
633 /// Carried, because `Drop` takes no arguments and the finish paint
634 /// needs the budgets. Same lifetime as the stage it belongs to.
635 settings: &'a crate::config::Settings,
636}
637
638impl<'a> FinishOnDrop<'a> {
639 pub fn new(
640 settings: &'a crate::config::Settings,
641 stage: &'a Stage,
642 idx: usize,
643 ) -> FinishOnDrop<'a> {
644 FinishOnDrop {
645 stage,
646 idx,
647 settings,
648 }
649 }
650}
651
652impl Drop for FinishOnDrop<'_> {
653 fn drop(&mut self) {
654 self.stage.finish(self.settings, self.idx);
655 }
656}
657
658/// Uninstalls the thread's sink on drop, whatever path the check took out.
659pub struct SinkGuard;
660
661impl Drop for SinkGuard {
662 fn drop(&mut self) {
663 SINK.with(|s| *s.borrow_mut() = None);
664 }
665}
666
667/// Detaches a command from its slot's displays when dropped. See
668/// [`Stage::attach`].
669pub struct AttachGuard {
670 stage: Arc<Stage>,
671 idx: usize,
672}
673
674impl Drop for AttachGuard {
675 fn drop(&mut self) {
676 let mut slots = self.stage.slots.lock().unwrap_or_else(|p| p.into_inner());
677 if let Some(slot) = slots.get_mut(self.idx) {
678 slot.activity = None;
679 }
680 }
681}
682
683/// The sink installed on THIS thread, if any — how a child-capture helper on
684/// the check's own thread learns where the reader threads should append.
685pub fn current_sink() -> Option<(Arc<Stage>, usize)> {
686 SINK.with(|s| s.borrow().clone())
687}
688
689/// One line of check output, wherever it should go.
690///
691/// THE funnel: `common::ok/fail/warn` call this, so a check's helper prints
692/// land in its slot during a stage and on stdout everywhere else. `line` is
693/// taken without a trailing newline, exactly like `println!`.
694pub fn say(line: &str) {
695 let routed = SINK.with(|s| {
696 s.borrow().as_ref().map(|(stage, idx)| {
697 stage.append_line(*idx, line);
698 })
699 });
700 if routed.is_none() {
701 println!("{line}");
702 }
703}
704
705/// `println!`, stage-aware: formats and routes through [`say`]. What every
706/// direct print inside a CHECK BODY becomes — a line printed raw from a
707/// check thread bypasses the slot and interleaves, which is the bug this
708/// module exists to close.
709#[macro_export]
710macro_rules! say {
711 ($($arg:tt)*) => {
712 $crate::live::say(&format!($($arg)*))
713 };
714}
715
716/// Should a check's SUCCESS line be swallowed?
717///
718/// A hook that passes says one line per check, and on a clean run that is the
719/// entire output: fourteen lines to say nothing happened. At a terminal those
720/// lines are the reassurance that the gate ran. Captured — an agent's tool
721/// result, a CI log — they are re-read on every later turn of the session and
722/// say no more the tenth time than the first.
723///
724/// So the setting names WHO is reading, not how loud to be:
725///
726/// - `auto` (default) — quiet when nobody is watching, verbose at a terminal.
727/// - `never` — every check says it passed, whoever is reading.
728/// - `always` — quiet everywhere.
729///
730/// `auto` is the default because the reader it costs nothing is the one at a
731/// terminal: `watching()` is true there, so a person sees exactly what they
732/// saw before. The reader it saves is the one who cannot skim — a captured
733/// log, an agent's tool result — and that reader was paying for fourteen
734/// lines of nothing on every turn of a session. A default that is free for
735/// one audience and compounding for the other is not a neutral default.
736///
737/// Only the success lines go. A failure, a warning, a check that could not
738/// run, a repaired file, and the blocked summary are printed under every
739/// setting: quiet is about the uneventful path, and nothing else.
740pub fn quiet(settings: &crate::config::Settings) -> bool {
741 *settings.quiet.get_or_init(|| {
742 decide(
743 crate::config::enumerated_or(settings, "amont.quiet", QUIET_VALUES, "auto"),
744 watching(),
745 )
746 })
747}
748
749pub const QUIET_VALUES: &[&str] = &["never", "auto", "always"];
750
751/// Pure, so the three-way decision is testable without a terminal or a config.
752fn decide(setting: &str, watching: bool) -> bool {
753 match setting {
754 "always" => true,
755 "auto" => !watching,
756 // `never`. A value `enumerated_or` rejected never reaches here — it
757 // complains and hands back the default, which is now `auto`.
758 _ => false,
759 }
760}
761
762/// Whether the capture mechanism is on at all. `amont.progress false` is the
763/// escape hatch back to raw streaming — one knob, read once.
764pub fn enabled(settings: &crate::config::Settings) -> bool {
765 *settings
766 .progress
767 .get_or_init(|| crate::config::boolean_or(settings, "amont.progress", true))
768}
769
770/// Is anyone watching? True only when stderr is a real terminal that speaks
771/// VT: not piped, not redirected, not `TERM=dumb` — and on Windows only
772/// with `TERM` actually set, because bare conhost may not interpret the
773/// cursor codes the region depends on. This is the paint gate; capture
774/// ([`enabled`]) does not consult it.
775pub fn watching() -> bool {
776 static WATCHING: std::sync::OnceLock<bool> = std::sync::OnceLock::new();
777 *WATCHING.get_or_init(|| {
778 if !std::io::stderr().is_terminal() {
779 return false;
780 }
781 match std::env::var("TERM") {
782 Ok(term) => term != "dumb",
783 Err(_) => !cfg!(windows),
784 }
785 })
786}
787
788#[cfg(test)]
789mod tests {
790 use super::*;
791
792 fn test_settings() -> crate::config::Settings {
793 crate::config::Settings::default()
794 }
795
796 #[test]
797 fn quiet_asks_who_is_reading() {
798 assert!(!decide("never", true));
799 assert!(!decide("never", false));
800 assert!(always_and_auto_agree_at_a_terminal());
801 assert!(decide("auto", false), "captured: nobody is watching");
802 assert!(decide("always", true));
803 assert!(decide("always", false));
804 // An unreadable value has already been reported by `enumerated_or`,
805 // which hands back the default; silence is never assumed.
806 assert!(!decide("shhh", false));
807 }
808
809 fn always_and_auto_agree_at_a_terminal() -> bool {
810 !decide("auto", true) && decide("always", true)
811 }
812
813 /// The atomicity contract at the unit level: two threads writing
814 /// interleaved lines into their own slots come out as two contiguous
815 /// buffers, whatever the scheduler did.
816 #[test]
817 fn slots_do_not_share_a_buffer() {
818 let stage = Stage::begin(&test_settings(), &["a", "b"]);
819 std::thread::scope(|scope| {
820 for idx in 0..2 {
821 let stage = Arc::clone(&stage);
822 scope.spawn(move || {
823 let _guard = stage.enter(idx);
824 for i in 0..50 {
825 say(&format!("check-{idx} line-{i}"));
826 std::thread::yield_now();
827 }
828 });
829 }
830 });
831 let slots = stage.slots.lock().unwrap();
832 for idx in 0..2 {
833 let text = String::from_utf8(slots[idx].buf.clone()).unwrap();
834 assert_eq!(text.lines().count(), 50);
835 assert!(
836 text.lines()
837 .all(|l| l.starts_with(&format!("check-{idx} "))),
838 "a foreign line landed in slot {idx}"
839 );
840 }
841 }
842
843 /// A thread with no sink prints; its lines never land in anyone's slot.
844 #[test]
845 fn no_sink_means_no_capture() {
846 let stage = Stage::begin(&test_settings(), &["a"]);
847 say("goes to stdout, not to a slot");
848 let slots = stage.slots.lock().unwrap();
849 assert!(slots[0].buf.is_empty());
850 }
851
852 /// After finish, late writes are dropped rather than stranded — a child
853 /// reader thread that outlives its check must not corrupt a later block.
854 #[test]
855 fn a_finished_slot_takes_no_more_writes() {
856 let stage = Stage::begin(&test_settings(), &["a"]);
857 stage.append_raw(0, b"before\n");
858 stage.finish(&test_settings(), 0);
859 stage.append_raw(0, b"after\n");
860 let slots = stage.slots.lock().unwrap();
861 assert!(slots[0].buf.is_empty(), "a write landed after finish");
862 }
863
864 /// A repo-derived check name cannot smuggle control bytes onto a live
865 /// terminal: sanitised at begin, once, for every later paint.
866 #[test]
867 fn a_slot_name_is_sanitised_at_begin() {
868 let stage = Stage::begin(&test_settings(), &["evil\u{1b}[2Jname\rhere"]);
869 let slots = stage.slots.lock().unwrap();
870 assert!(!slots[0].name.contains('\u{1b}'), "{:?}", slots[0].name);
871 assert!(!slots[0].name.contains('\r'), "{:?}", slots[0].name);
872 }
873
874 /// Region names drop the stage's own prefix — it is the same twelve
875 /// characters on every line.
876 #[test]
877 fn a_slot_name_drops_the_stage_prefix() {
878 let stage = Stage::begin(
879 &test_settings(),
880 &["pre-commit-clippy", "pre-push-run-tests", "bare"],
881 );
882 let slots = stage.slots.lock().unwrap();
883 assert_eq!(slots[0].name, "clippy");
884 assert_eq!(slots[1].name, "run-tests");
885 assert_eq!(slots[2].name, "bare");
886 }
887
888 fn row(name: &str, elapsed: f64) -> Row {
889 Row {
890 name: name.into(),
891 elapsed,
892 quiet: 0.0,
893 still: 0.0,
894 cpu: RowCpu::None,
895 }
896 }
897
898 const B: Budgets = Budgets {
899 idle: 120,
900 ceiling: 3600,
901 };
902
903 /// The spinner frame comes from the clock: different elapsed, different
904 /// frame; same elapsed, same frame.
905 #[test]
906 fn frames_advance_with_time() {
907 let a = region(&[row("clippy", 0.0)], 80, B);
908 let b = region(&[row("clippy", 0.1)], 80, B);
909 let c = region(&[row("clippy", 1.0)], 80, B);
910 assert_ne!(a.chars().next(), b.chars().next());
911 assert_eq!(a.chars().next(), c.chars().next(), "10 frames per second");
912 }
913
914 /// Names pad to a column so the elapsed figures align — across the
915 /// minute mark too, where the figure changes shape.
916 #[test]
917 fn region_lines_align() {
918 let text = region(&[row("a", 0.0), row("longer-name", 0.0)], 80, B);
919 let widths: Vec<usize> = text.lines().map(|l| l.chars().count()).collect();
920 assert_eq!(widths[0], widths[1], "{text:?}");
921 let text = region(&[row("a", 3.2), row("b", 492.0)], 80, B);
922 let widths: Vec<usize> = text.lines().map(|l| l.chars().count()).collect();
923 assert_eq!(widths[0], widths[1], "{text:?}");
924 assert!(text.contains("8m12s"), "{text:?}");
925 }
926
927 /// Thirteen running checks paint as twelve lines and one overflow.
928 #[test]
929 fn region_caps_and_counts_the_rest() {
930 let entries: Vec<Row> = (0..13).map(|i| row(&format!("check-{i}"), 0.0)).collect();
931 let text = region(&entries, 80, B);
932 assert_eq!(text.lines().count(), MAX_LINES + 1);
933 assert!(text.ends_with("… and 1 more\n"), "{text:?}");
934 }
935
936 /// A narrow terminal truncates rather than wraps — a wrapped region
937 /// line would break the erase arithmetic.
938 #[test]
939 fn region_respects_width() {
940 let text = region(&[row("a-name-much-longer-than-the-terminal", 0.0)], 20, B);
941 assert!(text.lines().all(|l| l.chars().count() <= 20), "{text:?}");
942 }
943
944 /// No running checks, no region — not even a blank line.
945 #[test]
946 fn an_empty_region_is_empty() {
947 assert_eq!(region(&[], 80, B), "");
948 }
949
950 /// Silence is annotated only once it is news, and names the budget it
951 /// counts toward — a check that just paused between crates says
952 /// nothing extra.
953 #[test]
954 fn a_quiet_check_shows_its_silence_against_the_budget() {
955 let mut r = row("cargo-test", 300.0);
956 r.quiet = 5.0;
957 assert!(!region(&[r.clone()], 80, B).contains("quiet"));
958 r.quiet = 45.0;
959 let text = region(&[r.clone()], 80, B);
960 assert!(text.contains("quiet 45s/2m00s"), "{text:?}");
961 let off = Budgets { idle: 0, ..B };
962 let text = region(&[r], 80, off);
963 assert!(
964 text.contains("quiet 45s") && !text.contains('/'),
965 "{text:?}"
966 );
967 }
968
969 /// A silent check that is working shows how hard, with no countdown —
970 /// no kill is coming; one whose CPU work pushed the still-time back
971 /// counts down the still-time, which is what the kill decision uses.
972 /// Both fit an 80-column terminal with a longish name.
973 #[test]
974 fn a_quiet_busy_check_shows_cores_and_an_idle_one_counts_down_the_still_time() {
975 let mut r = row("vitest-workspace", 240.0);
976 r.quiet = 130.0;
977 r.still = 130.0;
978 r.cpu = RowCpu::Busy(3900);
979 let busy = region(&[r.clone()], 80, B);
980 assert!(busy.contains("· quiet 2m10s · ~3.9 cores"), "{busy:?}");
981 assert!(
982 !busy.contains("/2m00s"),
983 "no countdown while busy: {busy:?}"
984 );
985
986 r.cpu = RowCpu::Idle;
987 r.still = 40.0;
988 let idle = region(&[r.clone()], 80, B);
989 assert!(idle.contains("· quiet 2m10s · idle 40s/2m00s"), "{idle:?}");
990
991 r.cpu = RowCpu::Unmeasured;
992 r.still = 130.0;
993 let plain = region(&[r], 80, B);
994 assert!(plain.contains("· quiet 2m10s/2m00s"), "{plain:?}");
995
996 for text in [busy, idle, plain] {
997 assert!(text.lines().all(|l| l.chars().count() <= 80), "{text:?}");
998 }
999 }
1000
1001 /// The heartbeat's prefix is unchanged — log readers grep it — and the
1002 /// CPU state rides after it.
1003 #[test]
1004 fn a_heartbeat_appends_the_cpu_state_after_an_unchanged_prefix() {
1005 let mut r = row("vitest", 240.0);
1006 r.quiet = 130.0;
1007 r.still = 40.0;
1008 let prefix = " … vitest still running: 4m00s, last output 2m10s ago";
1009 for (cpu, suffix) in [
1010 (RowCpu::Busy(3900), ", busy ~3.9 cores\n"),
1011 (RowCpu::Idle, ", CPU idle 40s\n"),
1012 (RowCpu::Unmeasured, ", CPU unmeasured\n"),
1013 (RowCpu::None, "\n"),
1014 (RowCpu::NotSampled, "\n"),
1015 ] {
1016 r.cpu = cpu;
1017 assert_eq!(beat_line(&r, false, B, false), format!("{prefix}{suffix}"));
1018 }
1019 }
1020
1021 /// The first beat states the rule that actually applies to this check.
1022 #[test]
1023 fn the_first_beat_states_the_rule_in_force() {
1024 let mut r = row("vitest", 60.0);
1025 r.cpu = RowCpu::Unmeasured;
1026 let sampled = beat_line(&r, true, B, false);
1027 assert!(
1028 flat(&sampled).contains(
1029 "killed after 2m00s with no output and under 0.1 core of CPU, or 1h00m in total"
1030 ),
1031 "{sampled:?}"
1032 );
1033 r.cpu = RowCpu::NotSampled;
1034 let not = beat_line(&r, true, B, false);
1035 assert!(
1036 not.contains("2m00s of silence or 1h00m in total"),
1037 "{not:?}"
1038 );
1039 assert!(
1040 not.contains("CPU not sampled here, silence alone counts"),
1041 "{not:?}"
1042 );
1043 }
1044
1045 /// The ceiling appears once a check is 80% of the way to it — the cliff,
1046 /// shown before the fall — and never for a disabled ceiling.
1047 #[test]
1048 fn the_ceiling_shows_only_when_it_is_near() {
1049 assert!(!region(&[row("cargo-test", 1000.0)], 80, B).contains("/1h00m"));
1050 let text = region(&[row("cargo-test", 3000.0)], 80, B);
1051 assert!(text.contains("50m00s/1h00m"), "{text:?}");
1052 let off = Budgets { ceiling: 0, ..B };
1053 assert!(!region(&[row("cargo-test", 3000.0)], 80, off).contains("/"));
1054 }
1055
1056 /// The heartbeat says how long, how quiet, and — the first time — the
1057 /// budgets, so a reader at the far end of a pipe can tell "wait" from
1058 /// "kill it" without the docs.
1059 #[test]
1060 fn a_heartbeat_names_the_budgets_once() {
1061 let mut r = row("cargo-test", 60.0);
1062 r.quiet = 2.0;
1063 let first = beat_line(&r, true, B, false);
1064 assert!(
1065 first.contains("cargo-test still running: 1m00s"),
1066 "{first:?}"
1067 );
1068 assert!(first.contains("last output 2s ago"), "{first:?}");
1069 assert!(
1070 first.contains("2m00s of silence or 1h00m in total"),
1071 "{first:?}"
1072 );
1073 assert!(first.contains("amont.idleTimeout"), "{first:?}");
1074 let later = beat_line(&r, false, B, false);
1075 assert!(!later.contains("amont.idleTimeout"), "{later:?}");
1076 let off = beat_line(
1077 &r,
1078 true,
1079 Budgets {
1080 idle: 0,
1081 ceiling: 0,
1082 },
1083 false,
1084 );
1085 assert!(off.contains("off of silence or off in total"), "{off:?}");
1086 }
1087
1088 /// The message, with newlines and indentation flattened.
1089 ///
1090 /// The note is wrapped for a terminal, so a literal substring can fall
1091 /// across a line break — asserting on `"may close it first"` failed for
1092 /// no better reason than that `may` ended a line. These tests are about
1093 /// what the message SAYS; re-wrapping it should not break them.
1094 fn flat(line: &str) -> String {
1095 line.split_whitespace().collect::<Vec<_>>().join(" ")
1096 }
1097
1098 /// A long PUSH gate is told what it is sitting on; a commit gate is not.
1099 ///
1100 /// The three negatives matter as much as the positive. Said on every
1101 /// beat it would be nagging; said at commit time it would be false —
1102 /// there is no connection open — and a future refactor that wires
1103 /// `on_push` to a constant would show up here and nowhere else.
1104 #[test]
1105 fn a_long_push_gate_is_told_what_it_is_sitting_on() {
1106 let r = row("cargo-test", 60.0);
1107
1108 let pushing = flat(&beat_line(&r, true, B, true));
1109 assert!(
1110 pushing.contains("A remote may close it first"),
1111 "{pushing:?}"
1112 );
1113 assert!(pushing.contains("Connection reset by peer"), "{pushing:?}");
1114 assert!(
1115 pushing.contains("Moving a gate entry earlier"),
1116 "{pushing:?}"
1117 );
1118
1119 // Once, not every minute.
1120 let later = flat(&beat_line(&r, false, B, true));
1121 assert!(!later.contains("close it first"), "{later:?}");
1122
1123 // Never at commit time: nothing is waiting on a socket there.
1124 let committing = flat(&beat_line(&r, true, B, false));
1125 assert!(!committing.contains("close it first"), "{committing:?}");
1126 }
1127
1128 /// The advice that does NOT appear, and must not come back.
1129 ///
1130 /// `ServerAliveInterval 60` is the obvious suggestion and it is wrong:
1131 /// it was already in force on the machine where this failure was
1132 /// diagnosed, and the remote reset the connection regardless. Telling
1133 /// every amont user to set it would be confident, actionable and
1134 /// useless. This test exists so that a future reader who has the same
1135 /// obvious idea meets an argument instead of a blank.
1136 #[test]
1137 fn the_push_note_does_not_recommend_ssh_keepalive() {
1138 let r = row("cargo-test", 60.0);
1139 let pushing = flat(&beat_line(&r, true, B, true));
1140 assert!(
1141 !pushing.contains("ServerAlive"),
1142 "keepalive was already on when this failed; recommending it \
1143 would be useless advice: {pushing:?}"
1144 );
1145 assert!(
1146 pushing.contains("ssh keepalive does not prevent this"),
1147 "say so, rather than leaving the reader to try it: {pushing:?}"
1148 );
1149 }
1150}