use mkit_core::hash::{Hash, to_hex};
use mkit_transport_connect::{PendingEvent, UploadEvent};
use std::cell::RefCell;
use std::io::{IsTerminal, Write};
#[derive(Debug, Clone, Copy)]
pub enum Event {
ObjectsPacked(usize),
PackUploaded(u64),
ObjectsUnpacked(usize),
Step { index: usize, total: usize },
}
const REPORT_INTERVAL: usize = 8;
struct Reporter {
label: &'static str,
total: Option<usize>,
done: usize,
bytes: u64,
last_emit_done: usize,
emitted: bool,
step: Option<(usize, usize)>,
}
impl Reporter {
fn new(label: &'static str, total: Option<usize>) -> Self {
Self {
label,
total,
done: 0,
bytes: 0,
last_emit_done: 0,
emitted: false,
step: None,
}
}
fn record(&mut self, event: Event) {
match event {
Event::ObjectsPacked(n) | Event::ObjectsUnpacked(n) => {
self.done += n;
if self.done.saturating_sub(self.last_emit_done) >= REPORT_INTERVAL {
self.emit();
}
}
Event::PackUploaded(bytes) => {
self.bytes = self.bytes.saturating_add(bytes);
self.emit();
}
Event::Step { index, total } => {
self.finish();
self.done = 0;
self.bytes = 0;
self.last_emit_done = 0;
self.emitted = false;
self.step = Some((index, total));
}
}
}
fn label(&self) -> String {
match self.step {
Some((index, total)) => format!("{} (step {index}/{total})", self.label),
None => self.label.to_owned(),
}
}
fn emit(&mut self) {
self.last_emit_done = self.done;
self.emitted = true;
let label = self.label();
let mut stderr = std::io::stderr().lock();
let _ = match (self.total, self.bytes) {
(Some(total), 0) => write!(stderr, "\r{label}: {}/{total} objects", self.done),
(Some(total), bytes) => {
write!(
stderr,
"\r{label}: {}/{total} objects, {bytes} bytes",
self.done
)
}
(None, 0) => write!(stderr, "\r{label}: {} objects", self.done),
(None, bytes) => write!(stderr, "\r{label}: {} objects, {bytes} bytes", self.done),
};
let _ = stderr.flush();
}
fn finish(&mut self) {
if self.done == 0 && self.bytes == 0 {
return;
}
self.emit();
let mut stderr = std::io::stderr().lock();
let _ = writeln!(stderr, ", done.");
}
}
thread_local! {
static QUIET: std::cell::Cell<Option<bool>> = const { std::cell::Cell::new(None) };
static REPORTER: RefCell<Option<Reporter>> = const { RefCell::new(None) };
static PENDING: RefCell<Option<PendingReporter>> = const { RefCell::new(None) };
static UPLOAD: RefCell<Option<UploadReporter>> = const { RefCell::new(None) };
}
struct UploadReporter {
quiet: bool,
interactive: bool,
active: bool,
}
impl UploadReporter {
fn render(&mut self, event: UploadEvent) -> Option<String> {
if self.quiet {
return None;
}
match event {
UploadEvent::PartsPlanned {
parts,
resumed,
saved_bytes,
bytes,
} => {
self.active = true;
if self.interactive {
Some(format!(
"\rUploading pack: part {resumed}/{parts} ({}/{} MiB), {resumed} resumed\x1b[K",
saved_bytes / (1024 * 1024),
bytes / (1024 * 1024)
))
} else {
Some(format!(
"Uploading pack: {parts} parts, {resumed} resumed.\n"
))
}
}
UploadEvent::PartSent {
index,
parts,
saved_bytes,
bytes,
resumed,
} if self.interactive => Some(format!(
"\rUploading pack: part {}/{parts} ({}/{} MiB), {resumed} resumed\x1b[K",
index + 1,
saved_bytes / (1024 * 1024),
bytes / (1024 * 1024)
)),
UploadEvent::Finished if self.active => {
self.active = false;
if self.interactive {
Some("\rUploading pack: done.\x1b[K\n".to_owned())
} else {
Some("Upload complete.\n".to_owned())
}
}
_ => None,
}
}
}
pub fn upload_event(event: UploadEvent) {
UPLOAD.with(|slot| {
if let Some(reporter) = slot.borrow_mut().as_mut()
&& let Some(line) = reporter.render(event)
{
let mut stderr = std::io::stderr().lock();
let _ = stderr.write_all(line.as_bytes());
let _ = stderr.flush();
}
});
}
struct PendingReporter {
quiet: bool,
interactive: bool,
active: bool,
}
impl PendingReporter {
fn new(quiet: bool, mode: Option<&str>, is_tty: bool) -> Self {
Self {
quiet: quiet || mode == Some("never"),
interactive: mode == Some("always") || is_tty,
active: false,
}
}
fn render(&mut self, event: PendingEvent) -> Option<String> {
if self.quiet {
return None;
}
match event {
PendingEvent::Waiting { elapsed, .. } if self.interactive => {
self.active = true;
Some(format!(
"\rWaiting for server verification: {}s\x1b[K",
elapsed.as_secs()
))
}
PendingEvent::Waiting { .. } if !self.active => {
self.active = true;
Some("Waiting for server verification...\n".to_owned())
}
PendingEvent::Finished { elapsed, succeeded } if self.active => {
self.active = false;
if self.interactive {
let status = if succeeded { "done" } else { "stopped" };
Some(format!(
"\rWaiting for server verification: {}s, {status}.\x1b[K\n",
elapsed.as_secs()
))
} else if succeeded {
Some("Server verification complete.\n".to_owned())
} else {
None
}
}
_ => None,
}
}
}
pub fn pending_event(event: PendingEvent) {
PENDING.with(|slot| {
if let Some(reporter) = slot.borrow_mut().as_mut()
&& let Some(line) = reporter.render(event)
{
let mut stderr = std::io::stderr().lock();
let _ = stderr.write_all(line.as_bytes());
let _ = stderr.flush();
}
});
}
pub fn suspend_for_admission() {
REPORTER.with(|slot| {
if slot
.borrow()
.as_ref()
.is_some_and(|reporter| reporter.emitted)
{
let mut stderr = std::io::stderr().lock();
let _ = writeln!(stderr);
}
});
}
#[derive(Debug)]
#[must_use = "dropping this immediately ends progress reporting"]
pub struct Guard {
_private: (),
}
impl Drop for Guard {
fn drop(&mut self) {
QUIET.with(|quiet| quiet.set(None));
PENDING.with(|r| {
r.borrow_mut().take();
});
UPLOAD.with(|r| {
r.borrow_mut().take();
});
REPORTER.with(|r| {
if let Some(mut rep) = r.borrow_mut().take() {
rep.finish();
}
});
}
}
pub fn start(label: &'static str, total: Option<usize>, enabled: bool, quiet: bool) -> Guard {
let progress_mode = std::env::var("MKIT_PROGRESS").ok();
QUIET.with(|state| state.set(Some(quiet || progress_mode.as_deref() == Some("never"))));
PENDING.with(|r| {
*r.borrow_mut() = Some(PendingReporter::new(
quiet,
progress_mode.as_deref(),
std::io::stderr().is_terminal(),
));
});
UPLOAD.with(|r| {
*r.borrow_mut() = Some(UploadReporter {
quiet: quiet || progress_mode.as_deref() == Some("never"),
interactive: progress_mode.as_deref() == Some("always")
|| std::io::stderr().is_terminal(),
active: false,
});
});
REPORTER.with(|r| {
*r.borrow_mut() = if enabled {
Some(Reporter::new(label, total))
} else {
None
};
});
Guard { _private: () }
}
pub fn report(event: Event) {
REPORTER.with(|r| {
if let Ok(mut slot) = r.try_borrow_mut()
&& let Some(rep) = slot.as_mut()
{
rep.record(event);
}
});
}
fn step_line(index: usize, total: usize, head: &Hash) -> String {
format!(
"pushed step {index}/{total}: branch now at {}",
to_hex(head)
)
}
fn step_message(
quiet: bool,
interactive: bool,
index: usize,
total: usize,
head: &Hash,
) -> Option<String> {
(!quiet && !interactive && total >= 2).then(|| step_line(index, total, head))
}
pub fn step_committed(index: usize, total: usize, head: &Hash) {
let Some(quiet) = QUIET.with(std::cell::Cell::get) else {
return;
};
let interactive = REPORTER.with(|r| r.try_borrow().is_ok_and(|slot| slot.is_some()));
if let Some(line) = step_message(quiet, interactive, index, total, head) {
let mut stderr = std::io::stderr().lock();
let _ = writeln!(stderr, "{line}");
}
}
#[must_use]
pub fn should_report(quiet: bool) -> bool {
if quiet {
return false;
}
match std::env::var("MKIT_PROGRESS").ok().as_deref() {
Some("always") => true,
Some("never") => false,
_ => std::io::stderr().is_terminal(),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn report_without_guard_is_a_silent_no_op() {
report(Event::ObjectsPacked(1));
report(Event::PackUploaded(128));
report(Event::ObjectsUnpacked(3));
}
#[test]
fn pack_uploaded_accumulates_across_multiple_packs() {
let mut rep = Reporter::new("Writing objects", None);
rep.record(Event::PackUploaded(100));
assert_eq!(rep.bytes, 100);
rep.record(Event::PackUploaded(50));
assert_eq!(rep.bytes, 150, "second pack's bytes must add, not replace");
rep.record(Event::PackUploaded(25));
assert_eq!(rep.bytes, 175);
}
#[test]
fn disabled_guard_installs_no_reporter() {
let guard = start("Writing objects", Some(4), false, true);
report(Event::ObjectsPacked(4));
drop(guard);
}
#[test]
fn pending_stderr_lines_for_piped_forced_and_quiet_modes() {
let wait = PendingEvent::Waiting {
elapsed: std::time::Duration::from_secs(42),
next: std::time::Duration::from_secs(1),
};
let done = PendingEvent::Finished {
elapsed: std::time::Duration::from_secs(43),
succeeded: true,
};
let mut piped = PendingReporter::new(false, None, false);
assert_eq!(
piped.render(wait).as_deref(),
Some("Waiting for server verification...\n")
);
assert_eq!(piped.render(wait), None);
assert_eq!(
piped.render(done).as_deref(),
Some("Server verification complete.\n")
);
let mut forced = PendingReporter::new(false, Some("always"), false);
assert_eq!(
forced.render(wait).as_deref(),
Some("\rWaiting for server verification: 42s\x1b[K")
);
assert_eq!(
forced.render(done).as_deref(),
Some("\rWaiting for server verification: 43s, done.\x1b[K\n")
);
let mut quiet = PendingReporter::new(true, Some("always"), false);
assert_eq!(quiet.render(wait), None);
assert_eq!(quiet.render(done), None);
let mut never = PendingReporter::new(false, Some("never"), true);
assert_eq!(never.render(wait), None);
assert_eq!(never.render(done), None);
}
#[test]
fn upload_progress_shows_resumed_bytes_on_tty_and_two_lines_when_piped() {
let start = UploadEvent::PartsPlanned {
parts: 12,
resumed: 4,
saved_bytes: 32 << 20,
bytes: 96 << 20,
};
let sent = UploadEvent::PartSent {
index: 4,
parts: 12,
saved_bytes: 40 << 20,
bytes: 96 << 20,
resumed: 4,
};
let mut tty = UploadReporter {
quiet: false,
interactive: true,
active: false,
};
assert!(tty.render(start).unwrap().contains("part 4/12 (32/96 MiB)"));
assert!(
tty.render(sent)
.unwrap()
.contains("part 5/12 (40/96 MiB), 4 resumed")
);
assert!(tty.render(UploadEvent::Finished).unwrap().contains("done."));
let mut piped = UploadReporter {
quiet: false,
interactive: false,
active: false,
};
assert_eq!(
piped.render(start).as_deref(),
Some("Uploading pack: 12 parts, 4 resumed.\n")
);
assert_eq!(piped.render(sent), None);
assert_eq!(
piped.render(UploadEvent::Finished).as_deref(),
Some("Upload complete.\n")
);
}
#[test]
fn should_report_quiet_always_wins() {
assert!(!should_report(true));
}
#[test]
fn split_push_output_in_tty_piped_and_quiet_modes() {
let mut reporter = Reporter::new("Writing objects", None);
assert_eq!(reporter.label(), "Writing objects");
reporter.record(Event::Step { index: 2, total: 5 });
assert_eq!(reporter.label(), "Writing objects (step 2/5)");
assert_eq!((reporter.done, reporter.bytes), (0, 0));
let head = [0xab; 32];
let hex = to_hex(&head);
assert_eq!(
step_message(false, false, 2, 5, &head),
Some(format!("pushed step 2/5: branch now at {hex}"))
);
assert_eq!(step_message(false, true, 2, 5, &head), None, "tty");
assert_eq!(step_message(true, false, 2, 5, &head), None, "quiet");
assert_eq!(step_message(false, false, 1, 1, &head), None, "unsplit");
}
}