Skip to main content

mkit_cli/
progress.rs

1//! Honest transfer-progress reporting for `clone`/`push`/`pull`/`fetch`
2//! (#711).
3//!
4//! `mkit clone`/`push`/`pull`/`fetch` previously printed only a start
5//! banner and a final summary — the network transfer itself was silent.
6//! This module adds a lightweight, thread-local progress sink that the
7//! transfer call chain (`push_branch_with_depth` in
8//! `remote_dispatch::mod`, `unpack_downloaded_packs` in
9//! `remote_dispatch::packmap`) reports real, already-happened work to:
10//! objects staged into the outgoing pack, bytes handed to the transport,
11//! and objects unpacked from a downloaded pack.
12//!
13//! It deliberately never reports git's fabricated
14//! `Enumerating/Counting/Compressing objects` or `Total N (delta D)`
15//! lines — mkit's transport is one-object-per-pack and computes no
16//! cross-branch delta graph, so those numbers would be invented (see
17//! `docs/PARITY.md`'s "Human-facing output parity" section).
18//!
19//! ## Threading pattern
20//!
21//! Rather than adding a progress parameter to every function in the
22//! `push_all_with` → `push_branch_with_depth` → `pull_all` →
23//! `fetch_objects` call chain (touching dozens of existing call sites,
24//! including many integration tests that don't care about progress at
25//! all), this mirrors the pattern already used for interrupt handling:
26//! `crate::signal::is_shutdown()` is a global checkpoint polled inside
27//! the same loops. Here, [`report`] is the equivalent checkpoint —  a
28//! thread-local sink installed by [`start`] and torn down by the
29//! returned [`Guard`]'s `Drop`. When no sink is installed (the common
30//! case: every test that doesn't call [`start`], and any non-interactive
31//! run), `report` is a cheap thread-local check that does nothing.
32//! Concurrent callers (see `fetch_pull_lock_scope.rs`, which fetches
33//! from multiple threads) are unaffected: the sink is thread-local, so
34//! each thread has its own (absent, by default) reporter.
35//!
36//! ## Interactivity gating
37//!
38//! Mirrors `term::use_color_stderr`'s tty auto-detection: progress is
39//! shown only when stderr is a tty, unless overridden by an explicit
40//! `--quiet` flag (forces off) or the `MKIT_PROGRESS` env var
41//! (`always`/`never`/`auto`, mirroring `NO_COLOR`/`CLICOLOR_FORCE`'s
42//! override convention) — `always` is how the CLI integration tests
43//! observe progress lines over a piped (non-tty) stderr. Verification waiting
44//! is an exception: a piped, non-quiet push gets one start and one completion
45//! line because the wait can last for minutes.
46
47use mkit_core::hash::{Hash, to_hex};
48use mkit_transport_connect::{PendingEvent, UploadEvent};
49use std::cell::RefCell;
50use std::io::{IsTerminal, Write};
51
52/// One real, already-happened unit of transfer work. Never a projection
53/// or estimate.
54#[derive(Debug, Clone, Copy)]
55pub enum Event {
56    /// `count` objects were appended to the outgoing pack(s) (push side,
57    /// `build_and_upload_packs`'s `plan.raw` / `plan.deltas` loops).
58    ObjectsPacked(usize),
59    /// A finished pack (`bytes` long) was handed to
60    /// `Transport::upload_pack` — that pack's upload is complete. Fires
61    /// once per pack; a push whose plan exceeds a single pack's payload
62    /// cap fires this more than once, and `bytes` accumulates across
63    /// calls (issue #831) rather than reporting only the last pack.
64    PackUploaded(u64),
65    /// `count` objects were unpacked from one downloaded pack (pull/fetch
66    /// side, `unpack_downloaded_packs`) — real counts from the pack's own
67    /// [`mkit_core::pack::UnpackReport`].
68    ObjectsUnpacked(usize),
69    /// A split push (WP-1.17b) is starting advance `index` of `total`. The
70    /// counters restart for it and its label carries the step.
71    Step { index: usize, total: usize },
72}
73
74/// Objects between throttled stderr re-writes. The final event
75/// ([`Event::PackUploaded`], and [`Guard`]'s `Drop`) always emits
76/// regardless of this threshold, so the last line reflects the true
77/// final count even when it doesn't land on an interval boundary.
78const REPORT_INTERVAL: usize = 8;
79
80struct Reporter {
81    label: &'static str,
82    total: Option<usize>,
83    done: usize,
84    bytes: u64,
85    last_emit_done: usize,
86    emitted: bool,
87    step: Option<(usize, usize)>,
88}
89
90impl Reporter {
91    fn new(label: &'static str, total: Option<usize>) -> Self {
92        Self {
93            label,
94            total,
95            done: 0,
96            bytes: 0,
97            last_emit_done: 0,
98            emitted: false,
99            step: None,
100        }
101    }
102
103    fn record(&mut self, event: Event) {
104        match event {
105            Event::ObjectsPacked(n) | Event::ObjectsUnpacked(n) => {
106                self.done += n;
107                if self.done.saturating_sub(self.last_emit_done) >= REPORT_INTERVAL {
108                    self.emit();
109                }
110            }
111            Event::PackUploaded(bytes) => {
112                // Accumulate, not overwrite: a multi-pack push (#831)
113                // fires this once per pack, and the reported total must
114                // cover every pack uploaded so far, not just the last one.
115                self.bytes = self.bytes.saturating_add(bytes);
116                self.emit();
117            }
118            Event::Step { index, total } => {
119                self.finish();
120                self.done = 0;
121                self.bytes = 0;
122                self.last_emit_done = 0;
123                self.emitted = false;
124                self.step = Some((index, total));
125            }
126        }
127    }
128
129    fn label(&self) -> String {
130        match self.step {
131            Some((index, total)) => format!("{} (step {index}/{total})", self.label),
132            None => self.label.to_owned(),
133        }
134    }
135
136    fn emit(&mut self) {
137        self.last_emit_done = self.done;
138        self.emitted = true;
139        let label = self.label();
140        let mut stderr = std::io::stderr().lock();
141        let _ = match (self.total, self.bytes) {
142            (Some(total), 0) => write!(stderr, "\r{label}: {}/{total} objects", self.done),
143            (Some(total), bytes) => {
144                write!(
145                    stderr,
146                    "\r{label}: {}/{total} objects, {bytes} bytes",
147                    self.done
148                )
149            }
150            (None, 0) => write!(stderr, "\r{label}: {} objects", self.done),
151            (None, bytes) => write!(stderr, "\r{label}: {} objects, {bytes} bytes", self.done),
152        };
153        let _ = stderr.flush();
154    }
155
156    /// Force a final emit (bypassing the throttle) and move past the
157    /// self-overwriting `\r` line so later output isn't clobbered by it.
158    /// A no-op when nothing was ever reported (e.g. a no-op push).
159    fn finish(&mut self) {
160        if self.done == 0 && self.bytes == 0 {
161            return;
162        }
163        self.emit();
164        let mut stderr = std::io::stderr().lock();
165        let _ = writeln!(stderr, ", done.");
166    }
167}
168
169thread_local! {
170    /// `Some(quiet)` while a [`Guard`] is installed on this thread.
171    static QUIET: std::cell::Cell<Option<bool>> = const { std::cell::Cell::new(None) };
172    static REPORTER: RefCell<Option<Reporter>> = const { RefCell::new(None) };
173    static PENDING: RefCell<Option<PendingReporter>> = const { RefCell::new(None) };
174    static UPLOAD: RefCell<Option<UploadReporter>> = const { RefCell::new(None) };
175}
176
177struct UploadReporter {
178    quiet: bool,
179    interactive: bool,
180    active: bool,
181}
182
183impl UploadReporter {
184    fn render(&mut self, event: UploadEvent) -> Option<String> {
185        if self.quiet {
186            return None;
187        }
188        match event {
189            UploadEvent::PartsPlanned {
190                parts,
191                resumed,
192                saved_bytes,
193                bytes,
194            } => {
195                self.active = true;
196                if self.interactive {
197                    Some(format!(
198                        "\rUploading pack: part {resumed}/{parts} ({}/{} MiB), {resumed} resumed\x1b[K",
199                        saved_bytes / (1024 * 1024),
200                        bytes / (1024 * 1024)
201                    ))
202                } else {
203                    Some(format!(
204                        "Uploading pack: {parts} parts, {resumed} resumed.\n"
205                    ))
206                }
207            }
208            UploadEvent::PartSent {
209                index,
210                parts,
211                saved_bytes,
212                bytes,
213                resumed,
214            } if self.interactive => Some(format!(
215                "\rUploading pack: part {}/{parts} ({}/{} MiB), {resumed} resumed\x1b[K",
216                index + 1,
217                saved_bytes / (1024 * 1024),
218                bytes / (1024 * 1024)
219            )),
220            UploadEvent::Finished if self.active => {
221                self.active = false;
222                if self.interactive {
223                    Some("\rUploading pack: done.\x1b[K\n".to_owned())
224                } else {
225                    Some("Upload complete.\n".to_owned())
226                }
227            }
228            _ => None,
229        }
230    }
231}
232
233/// Render a multipart upload event on the current command's stderr sink.
234pub fn upload_event(event: UploadEvent) {
235    UPLOAD.with(|slot| {
236        if let Some(reporter) = slot.borrow_mut().as_mut()
237            && let Some(line) = reporter.render(event)
238        {
239            let mut stderr = std::io::stderr().lock();
240            let _ = stderr.write_all(line.as_bytes());
241            let _ = stderr.flush();
242        }
243    });
244}
245
246struct PendingReporter {
247    quiet: bool,
248    interactive: bool,
249    active: bool,
250}
251
252impl PendingReporter {
253    fn new(quiet: bool, mode: Option<&str>, is_tty: bool) -> Self {
254        Self {
255            quiet: quiet || mode == Some("never"),
256            interactive: mode == Some("always") || is_tty,
257            active: false,
258        }
259    }
260
261    fn render(&mut self, event: PendingEvent) -> Option<String> {
262        if self.quiet {
263            return None;
264        }
265        match event {
266            PendingEvent::Waiting { elapsed, .. } if self.interactive => {
267                self.active = true;
268                Some(format!(
269                    "\rWaiting for server verification: {}s\x1b[K",
270                    elapsed.as_secs()
271                ))
272            }
273            PendingEvent::Waiting { .. } if !self.active => {
274                self.active = true;
275                Some("Waiting for server verification...\n".to_owned())
276            }
277            PendingEvent::Finished { elapsed, succeeded } if self.active => {
278                self.active = false;
279                if self.interactive {
280                    let status = if succeeded { "done" } else { "stopped" };
281                    Some(format!(
282                        "\rWaiting for server verification: {}s, {status}.\x1b[K\n",
283                        elapsed.as_secs()
284                    ))
285                } else if succeeded {
286                    Some("Server verification complete.\n".to_owned())
287                } else {
288                    None
289                }
290            }
291            _ => None,
292        }
293    }
294}
295
296/// Report verification polling on the same caller thread as the transfer.
297pub fn pending_event(event: PendingEvent) {
298    PENDING.with(|slot| {
299        if let Some(reporter) = slot.borrow_mut().as_mut()
300            && let Some(line) = reporter.render(event)
301        {
302            let mut stderr = std::io::stderr().lock();
303            let _ = stderr.write_all(line.as_bytes());
304            let _ = stderr.flush();
305        }
306    });
307}
308
309/// End the current self-overwriting progress line before helper UX appears.
310pub fn suspend_for_admission() {
311    REPORTER.with(|slot| {
312        if slot
313            .borrow()
314            .as_ref()
315            .is_some_and(|reporter| reporter.emitted)
316        {
317            let mut stderr = std::io::stderr().lock();
318            let _ = writeln!(stderr);
319        }
320    });
321}
322
323/// RAII handle returned by [`start`]. Dropping it flushes a final
324/// progress line (if anything was reported) and uninstalls the
325/// thread-local sink, so a command can simply hold the guard for the
326/// duration of its transfer call and let scope-exit (including an early
327/// `return` on error) clean up.
328#[derive(Debug)]
329#[must_use = "dropping this immediately ends progress reporting"]
330pub struct Guard {
331    _private: (),
332}
333
334impl Drop for Guard {
335    fn drop(&mut self) {
336        QUIET.with(|quiet| quiet.set(None));
337        PENDING.with(|r| {
338            r.borrow_mut().take();
339        });
340        UPLOAD.with(|r| {
341            r.borrow_mut().take();
342        });
343        REPORTER.with(|r| {
344            if let Some(mut rep) = r.borrow_mut().take() {
345                rep.finish();
346            }
347        });
348    }
349}
350
351/// Install a thread-local progress reporter for the duration of the
352/// returned [`Guard`]. `enabled = false` installs no reporter, so
353/// [`report`] stays a cheap no-op — used when stderr isn't interactive
354/// or `--quiet` was passed (see [`should_report`]).
355///
356/// `total`, when known ahead of time (the push side plans its pack
357/// before building it), renders as `done/total`; `None` (the fetch/pull
358/// side, where the object count isn't known until each pack is
359/// downloaded) renders as a running count only — never a fabricated
360/// total.
361pub fn start(label: &'static str, total: Option<usize>, enabled: bool, quiet: bool) -> Guard {
362    let progress_mode = std::env::var("MKIT_PROGRESS").ok();
363    QUIET.with(|state| state.set(Some(quiet || progress_mode.as_deref() == Some("never"))));
364    PENDING.with(|r| {
365        *r.borrow_mut() = Some(PendingReporter::new(
366            quiet,
367            progress_mode.as_deref(),
368            std::io::stderr().is_terminal(),
369        ));
370    });
371    UPLOAD.with(|r| {
372        *r.borrow_mut() = Some(UploadReporter {
373            quiet: quiet || progress_mode.as_deref() == Some("never"),
374            interactive: progress_mode.as_deref() == Some("always")
375                || std::io::stderr().is_terminal(),
376            active: false,
377        });
378    });
379    REPORTER.with(|r| {
380        *r.borrow_mut() = if enabled {
381            Some(Reporter::new(label, total))
382        } else {
383            None
384        };
385    });
386    Guard { _private: () }
387}
388
389/// Report one real unit of already-completed transfer work to the
390/// current thread's installed reporter, if any. A no-op — a single
391/// thread-local check — when no [`Guard`] is active on this thread,
392/// which is the default for every caller that doesn't opt in (including
393/// every existing test that drives `push_branch_with_depth` /
394/// `push_all` / `pull_all` / `fetch_all` directly).
395pub fn report(event: Event) {
396    REPORTER.with(|r| {
397        // `try_borrow_mut` rather than `borrow_mut`: `report` is called
398        // from deep inside the transfer call chain and must never panic
399        // on a re-entrant borrow; silently dropping a progress tick is
400        // harmless (the running total is cosmetic), unlike the transfer
401        // itself.
402        if let Ok(mut slot) = r.try_borrow_mut()
403            && let Some(rep) = slot.as_mut()
404        {
405            rep.record(event);
406        }
407    });
408}
409
410/// The line a piped, non-quiet run gets when one advance of a split push has
411/// landed.
412fn step_line(index: usize, total: usize, head: &Hash) -> String {
413    format!(
414        "pushed step {index}/{total}: branch now at {}",
415        to_hex(head)
416    )
417}
418
419/// The line for one landed advance, or `None` when nothing should be printed:
420/// a quiet run, a single advance, or an interactive run, whose live progress
421/// line already carries the step.
422fn step_message(
423    quiet: bool,
424    interactive: bool,
425    index: usize,
426    total: usize,
427    head: &Hash,
428) -> Option<String> {
429    (!quiet && !interactive && total >= 2).then(|| step_line(index, total, head))
430}
431
432/// Report that advance `index` of `total` of a split push has landed. Outside
433/// a [`Guard`] (library callers, tests) nothing is printed.
434pub fn step_committed(index: usize, total: usize, head: &Hash) {
435    let Some(quiet) = QUIET.with(std::cell::Cell::get) else {
436        return;
437    };
438    let interactive = REPORTER.with(|r| r.try_borrow().is_ok_and(|slot| slot.is_some()));
439    if let Some(line) = step_message(quiet, interactive, index, total, head) {
440        let mut stderr = std::io::stderr().lock();
441        let _ = writeln!(stderr, "{line}");
442    }
443}
444
445/// Whether progress should be shown on stderr: not explicitly silenced
446/// (`--quiet` / `-q`), and either `MKIT_PROGRESS` forces a decision or
447/// stderr is a tty. Mirrors `term::use_color_stderr`'s
448/// `NO_COLOR`/`CLICOLOR_FORCE`-style override convention — `always` is
449/// how CLI integration tests observe progress lines over a piped
450/// (non-tty) stderr; `never` is an explicit opt-out distinct from
451/// `--quiet` (e.g. for scripting environments that set it once instead
452/// of threading `--quiet` through every call site).
453#[must_use]
454pub fn should_report(quiet: bool) -> bool {
455    if quiet {
456        return false;
457    }
458    match std::env::var("MKIT_PROGRESS").ok().as_deref() {
459        Some("always") => true,
460        Some("never") => false,
461        _ => std::io::stderr().is_terminal(),
462    }
463}
464
465#[cfg(test)]
466mod tests {
467    use super::*;
468
469    /// `report` with no active [`Guard`] must not panic and must not
470    /// touch stderr (there's no reporter to write through) — the
471    /// no-op path every existing push/pull/fetch integration test takes.
472    #[test]
473    fn report_without_guard_is_a_silent_no_op() {
474        report(Event::ObjectsPacked(1));
475        report(Event::PackUploaded(128));
476        report(Event::ObjectsUnpacked(3));
477    }
478
479    /// issue #831: a multi-pack push fires `PackUploaded` once per
480    /// pack. The reported byte total must accumulate across those
481    /// calls, not report only the last pack (the bug this test pins).
482    #[test]
483    fn pack_uploaded_accumulates_across_multiple_packs() {
484        let mut rep = Reporter::new("Writing objects", None);
485        rep.record(Event::PackUploaded(100));
486        assert_eq!(rep.bytes, 100);
487        rep.record(Event::PackUploaded(50));
488        assert_eq!(rep.bytes, 150, "second pack's bytes must add, not replace");
489        rep.record(Event::PackUploaded(25));
490        assert_eq!(rep.bytes, 175);
491    }
492
493    /// A disabled guard (`enabled: false`) installs no reporter, so
494    /// `report` inside its scope is still the no-op path.
495    #[test]
496    fn disabled_guard_installs_no_reporter() {
497        let guard = start("Writing objects", Some(4), false, true);
498        report(Event::ObjectsPacked(4));
499        drop(guard);
500    }
501
502    #[test]
503    fn pending_stderr_lines_for_piped_forced_and_quiet_modes() {
504        let wait = PendingEvent::Waiting {
505            elapsed: std::time::Duration::from_secs(42),
506            next: std::time::Duration::from_secs(1),
507        };
508        let done = PendingEvent::Finished {
509            elapsed: std::time::Duration::from_secs(43),
510            succeeded: true,
511        };
512        let mut piped = PendingReporter::new(false, None, false);
513        assert_eq!(
514            piped.render(wait).as_deref(),
515            Some("Waiting for server verification...\n")
516        );
517        assert_eq!(piped.render(wait), None);
518        assert_eq!(
519            piped.render(done).as_deref(),
520            Some("Server verification complete.\n")
521        );
522
523        let mut forced = PendingReporter::new(false, Some("always"), false);
524        assert_eq!(
525            forced.render(wait).as_deref(),
526            Some("\rWaiting for server verification: 42s\x1b[K")
527        );
528        assert_eq!(
529            forced.render(done).as_deref(),
530            Some("\rWaiting for server verification: 43s, done.\x1b[K\n")
531        );
532
533        let mut quiet = PendingReporter::new(true, Some("always"), false);
534        assert_eq!(quiet.render(wait), None);
535        assert_eq!(quiet.render(done), None);
536
537        let mut never = PendingReporter::new(false, Some("never"), true);
538        assert_eq!(never.render(wait), None);
539        assert_eq!(never.render(done), None);
540    }
541
542    #[test]
543    fn upload_progress_shows_resumed_bytes_on_tty_and_two_lines_when_piped() {
544        let start = UploadEvent::PartsPlanned {
545            parts: 12,
546            resumed: 4,
547            saved_bytes: 32 << 20,
548            bytes: 96 << 20,
549        };
550        let sent = UploadEvent::PartSent {
551            index: 4,
552            parts: 12,
553            saved_bytes: 40 << 20,
554            bytes: 96 << 20,
555            resumed: 4,
556        };
557        let mut tty = UploadReporter {
558            quiet: false,
559            interactive: true,
560            active: false,
561        };
562        assert!(tty.render(start).unwrap().contains("part 4/12 (32/96 MiB)"));
563        assert!(
564            tty.render(sent)
565                .unwrap()
566                .contains("part 5/12 (40/96 MiB), 4 resumed")
567        );
568        assert!(tty.render(UploadEvent::Finished).unwrap().contains("done."));
569
570        let mut piped = UploadReporter {
571            quiet: false,
572            interactive: false,
573            active: false,
574        };
575        assert_eq!(
576            piped.render(start).as_deref(),
577            Some("Uploading pack: 12 parts, 4 resumed.\n")
578        );
579        assert_eq!(piped.render(sent), None);
580        assert_eq!(
581            piped.render(UploadEvent::Finished).as_deref(),
582            Some("Upload complete.\n")
583        );
584    }
585
586    /// `should_report` precedence: `--quiet` wins outright, then
587    /// `MKIT_PROGRESS`, then tty-ness. Exercised via the pure
588    /// tty-independent branches only (quiet, and the env var forcing a
589    /// decision) — this process's stderr tty-ness varies by how tests
590    /// are invoked, so the `_ =>` fallthrough isn't asserted here.
591    #[test]
592    fn should_report_quiet_always_wins() {
593        assert!(!should_report(true));
594    }
595
596    /// A split push (WP-1.17b): the live progress line names the step, a
597    /// piped run gets one line per landed advance, a quiet run and a single
598    /// advance get none.
599    #[test]
600    fn split_push_output_in_tty_piped_and_quiet_modes() {
601        let mut reporter = Reporter::new("Writing objects", None);
602        assert_eq!(reporter.label(), "Writing objects");
603        reporter.record(Event::Step { index: 2, total: 5 });
604        assert_eq!(reporter.label(), "Writing objects (step 2/5)");
605        assert_eq!((reporter.done, reporter.bytes), (0, 0));
606
607        let head = [0xab; 32];
608        let hex = to_hex(&head);
609        assert_eq!(
610            step_message(false, false, 2, 5, &head),
611            Some(format!("pushed step 2/5: branch now at {hex}"))
612        );
613        assert_eq!(step_message(false, true, 2, 5, &head), None, "tty");
614        assert_eq!(step_message(true, false, 2, 5, &head), None, "quiet");
615        assert_eq!(step_message(false, false, 1, 1, &head), None, "unsplit");
616    }
617}