Skip to main content

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}