1use mkit_core::hash::{Hash, to_hex};
48use mkit_transport_connect::{PendingEvent, UploadEvent};
49use std::cell::RefCell;
50use std::io::{IsTerminal, Write};
51
52#[derive(Debug, Clone, Copy)]
55pub enum Event {
56 ObjectsPacked(usize),
59 PackUploaded(u64),
65 ObjectsUnpacked(usize),
69 Step { index: usize, total: usize },
72}
73
74const 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 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 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 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
233pub 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
296pub 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
309pub 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#[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
351pub 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
389pub fn report(event: Event) {
396 REPORTER.with(|r| {
397 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
410fn 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
419fn 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
432pub 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#[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 #[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 #[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 #[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 #[test]
592 fn should_report_quiet_always_wins() {
593 assert!(!should_report(true));
594 }
595
596 #[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}