use std::io::{self, Write};
use std::sync::mpsc::{self, Receiver, RecvTimeoutError};
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::thread::{self, JoinHandle};
use std::time::{Duration, Instant};
use fdu_core::query::SizeMetric;
use fdu_core::{Progress, ProgressPhase, ProgressSnapshot};
use crate::progress_line::{FrameFacts, Phase, ProgressPlan, render_frame};
pub(crate) const ERASE_LINE: &str = "\r\x1b[2K";
const FIRST_FRAME_DELAY: Duration = Duration::from_millis(500);
const REDRAW_INTERVAL: Duration = Duration::from_millis(100);
const FALLBACK_WIDTH: usize = 80;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) struct Timing {
pub delay: Duration,
pub tick: Duration,
}
impl Default for Timing {
fn default() -> Self {
Self { delay: FIRST_FRAME_DELAY, tick: REDRAW_INTERVAL }
}
}
pub(crate) struct ProgressIo {
pub out: Box<dyn Write + Send>,
pub timing: Timing,
pub width: fn() -> usize,
pub interrupt: fn(Arc<Line>, bool),
}
impl ProgressIo {
pub(crate) fn for_process() -> Self {
Self {
out: Box::new(io::stderr()),
timing: Timing::default(),
width: stderr_width,
interrupt: crate::interrupt::install,
}
}
#[cfg(test)]
pub(crate) fn inert() -> Self {
let never = Duration::from_secs(60 * 60);
Self {
out: Box::new(io::sink()),
timing: Timing { delay: never, tick: never },
width: || 80,
interrupt: |_, _| {},
}
}
}
pub(crate) fn stderr_width() -> usize {
terminal_size::terminal_size_of(io::stderr())
.map_or(FALLBACK_WIDTH, |(terminal_size::Width(width), _)| usize::from(width))
}
pub(crate) struct LineState {
pub stopped: bool,
pub frame_on_screen: bool,
pub out: Box<dyn Write + Send>,
}
impl LineState {
pub(crate) fn draw(&mut self, frame: &str) {
let written = self
.out
.write_all(ERASE_LINE.as_bytes())
.and_then(|()| self.out.write_all(frame.as_bytes()))
.and_then(|()| self.out.flush());
match written {
Ok(()) => self.frame_on_screen = true,
Err(_) => self.stopped = true,
}
}
pub(crate) fn stop(&mut self) {
self.stopped = true;
if self.frame_on_screen {
self.frame_on_screen = false;
let _ = self.out.write_all(ERASE_LINE.as_bytes()).and_then(|()| self.out.flush());
}
}
}
pub(crate) struct Line {
state: Mutex<LineState>,
}
impl Line {
pub(crate) fn new(out: Box<dyn Write + Send>) -> Arc<Self> {
Arc::new(Self {
state: Mutex::new(LineState { stopped: false, frame_on_screen: false, out }),
})
}
pub(crate) fn lock(&self) -> MutexGuard<'_, LineState> {
self.state.lock().unwrap_or_else(PoisonError::into_inner)
}
}
pub(crate) fn frame_facts(snapshot: &ProgressSnapshot, size: SizeMetric) -> Option<FrameFacts> {
let phase = match snapshot.phase {
ProgressPhase::Starting => return None,
ProgressPhase::Loading => Phase::Loading,
ProgressPhase::Scanning => Phase::Scanning,
ProgressPhase::Revalidating => Phase::Revalidating,
ProgressPhase::Analyzing => Phase::Analyzing,
ProgressPhase::Saving => Phase::Saving,
ProgressPhase::Indexing => Phase::Indexing,
ProgressPhase::Summarizing => Phase::Summarizing,
};
Some(FrameFacts {
phase,
directories: snapshot.directories,
files: snapshot.files,
bytes: match size {
SizeMetric::Apparent => snapshot.bytes,
SizeMetric::Allocated => snapshot.allocated,
},
analysis: snapshot.analysis,
})
}
pub(crate) struct Ticker {
line: Arc<Line>,
stop: Option<mpsc::Sender<()>>,
thread: Option<JoinHandle<()>>,
}
impl Ticker {
pub(crate) fn start(
plan: ProgressPlan,
progress: Progress,
started: Instant,
io: ProgressIo,
) -> Self {
let line = Line::new(io.out);
let color = plan.color;
let (stop, stopped) = mpsc::channel();
let thread = thread::Builder::new()
.name("fdu-progress".to_string())
.spawn({
let line = Arc::clone(&line);
move || {
redraw_until_stopped(
&line, &plan, &progress, started, &stopped, io.timing, io.width,
);
}
})
.ok();
(io.interrupt)(Arc::clone(&line), color);
Self { line, stop: Some(stop), thread }
}
#[cfg(test)]
pub(crate) fn line(&self) -> &Arc<Line> {
&self.line
}
pub(crate) fn stop(&mut self) {
drop(self.stop.take());
if let Some(thread) = self.thread.take() {
let _ = thread.join();
}
self.line.lock().stop();
}
}
impl Drop for Ticker {
fn drop(&mut self) {
self.stop();
}
}
fn redraw_until_stopped(
line: &Line,
plan: &ProgressPlan,
progress: &Progress,
started: Instant,
stopped: &Receiver<()>,
timing: Timing,
width: fn() -> usize,
) {
let mut wait = timing.delay;
let mut spinner_step = 0;
loop {
match stopped.recv_timeout(wait) {
Err(RecvTimeoutError::Timeout) => {}
Ok(()) | Err(RecvTimeoutError::Disconnected) => return,
}
wait = timing.tick;
let Some(facts) = frame_facts(&progress.snapshot(), plan.size) else {
continue;
};
let frame =
render_frame(&plan.root, &facts, started.elapsed(), spinner_step, width(), plan.color);
let mut state = line.lock();
if state.stopped {
return;
}
state.draw(&frame);
if state.stopped {
return;
}
spinner_step = spinner_step.wrapping_add(1);
}
}
#[cfg(test)]
#[derive(Clone, Default)]
pub(crate) struct SharedBuffer(Arc<Mutex<Vec<u8>>>);
#[cfg(test)]
impl SharedBuffer {
pub(crate) fn contents(&self) -> Vec<u8> {
self.0.lock().unwrap_or_else(PoisonError::into_inner).clone()
}
pub(crate) fn text(&self) -> String {
String::from_utf8(self.contents()).expect("the buffer holds UTF-8")
}
}
#[cfg(test)]
impl Write for SharedBuffer {
fn write(&mut self, buffer: &[u8]) -> io::Result<usize> {
self.0.lock().unwrap_or_else(PoisonError::into_inner).extend_from_slice(buffer);
Ok(buffer.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
#[cfg(test)]
pub(crate) fn scanned(root: &std::path::Path) -> Progress {
use std::time::SystemTime;
use fdu_core::content::AnalysisSet;
use fdu_core::query::{Basis, Delivery, Query, Request};
use fdu_core::{CachePolicy, ScanConfig, prepare_report_with_progress};
let progress = Progress::new();
let request = Request::new(
Basis {
root: root.to_path_buf(),
scope: ScanConfig::default().into(),
content: AnalysisSet::NONE,
},
Query::default(),
SystemTime::now(),
);
let delivery = Delivery::new(CachePolicy::Off, None);
let (_, pending, _) =
prepare_report_with_progress(&request, &delivery, &progress).expect("a report");
pending.join().expect("no save to fail");
assert_eq!(progress.snapshot().phase, ProgressPhase::Summarizing);
progress
}
#[cfg(test)]
mod tests {
use super::*;
const ROOT: &str = "~/wrk/github";
fn plan() -> ProgressPlan {
ProgressPlan {
draw: true,
root: ROOT.to_string(),
color: false,
size: SizeMetric::Allocated,
}
}
fn io(out: &SharedBuffer, timing: Timing) -> ProgressIo {
ProgressIo { out: Box::new(out.clone()), timing, width: || 100, interrupt: |_, _| {} }
}
fn erases(bytes: &[u8]) -> usize {
bytes.windows(ERASE_LINE.len()).filter(|window| *window == ERASE_LINE.as_bytes()).count()
}
fn wait_for(out: &SharedBuffer, condition: impl Fn(&[u8]) -> bool) {
let deadline = Instant::now() + Duration::from_secs(10);
while !condition(&out.contents()) {
assert!(Instant::now() < deadline, "the ticker never got there");
thread::sleep(Duration::from_millis(1));
}
}
#[test]
fn a_fresh_handle_gives_no_frame_and_every_phase_maps_across() {
assert_eq!(frame_facts(&Progress::new().snapshot(), SizeMetric::Allocated), None);
let snapshot = ProgressSnapshot {
phase: ProgressPhase::Analyzing,
directories: 3,
files: 40,
bytes: 4_096,
allocated: 8_192,
analysis: Some((2, 5)),
};
assert_eq!(
frame_facts(&snapshot, SizeMetric::Allocated),
Some(FrameFacts {
phase: Phase::Analyzing,
directories: 3,
files: 40,
bytes: 8_192,
analysis: Some((2, 5)),
})
);
for (engine, frame) in [
(ProgressPhase::Loading, Phase::Loading),
(ProgressPhase::Scanning, Phase::Scanning),
(ProgressPhase::Revalidating, Phase::Revalidating),
(ProgressPhase::Analyzing, Phase::Analyzing),
(ProgressPhase::Saving, Phase::Saving),
(ProgressPhase::Indexing, Phase::Indexing),
(ProgressPhase::Summarizing, Phase::Summarizing),
] {
let snapshot = ProgressSnapshot { phase: engine, ..snapshot };
assert_eq!(
frame_facts(&snapshot, SizeMetric::Allocated).map(|facts| facts.phase),
Some(frame)
);
}
}
#[test]
fn the_bytes_follow_the_answers_size_metric() {
let snapshot = ProgressSnapshot {
phase: ProgressPhase::Scanning,
directories: 1,
files: 1,
bytes: 8_796_093_022_208,
allocated: 41_107_456,
analysis: None,
};
let bytes = |size| frame_facts(&snapshot, size).map(|facts| facts.bytes);
assert_eq!(bytes(SizeMetric::Allocated), Some(41_107_456));
assert_eq!(bytes(SizeMetric::Apparent), Some(8_796_093_022_208));
}
#[test]
fn a_run_that_finishes_inside_the_delay_writes_nothing_and_stops_at_once() {
let root = tempfile::tempdir().expect("tempdir");
let out = SharedBuffer::default();
let timing = Timing { delay: Duration::from_secs(60), tick: Duration::from_secs(60) };
let started = Instant::now();
let mut ticker = Ticker::start(plan(), scanned(root.path()), started, io(&out, timing));
ticker.stop();
assert!(started.elapsed() < Duration::from_secs(10), "stopping waited out the delay");
assert!(out.contents().is_empty(), "{:?}", out.text());
assert!(ticker.line().lock().stopped);
assert!(!ticker.line().lock().frame_on_screen);
}
#[test]
fn the_first_bytes_are_an_erase_and_a_frame_and_the_last_are_an_erase() {
let root = tempfile::tempdir().expect("tempdir");
let out = SharedBuffer::default();
let timing = Timing { delay: Duration::ZERO, tick: Duration::from_millis(1) };
let mut ticker =
Ticker::start(plan(), scanned(root.path()), Instant::now(), io(&out, timing));
wait_for(&out, |bytes| erases(bytes) >= 3);
ticker.stop();
assert!(ticker.line().lock().stopped);
assert!(!ticker.line().lock().frame_on_screen);
let text = out.text();
let frames: Vec<&str> = text.split(ERASE_LINE).collect();
assert_eq!(frames[0], "", "the first bytes are the erase sequence");
assert!(
frames[1].starts_with(&format!(
"⠋ {ROOT} Summarizing 0 files · 1 dirs · 0 B · "
)),
"{:?}",
frames[1]
);
assert!(frames[2].starts_with('⠙'), "the spinner advances one cell per redraw");
assert!(frames[3].starts_with('⠹'));
assert_eq!(frames[frames.len() - 1], "", "the last bytes are the erase sequence");
assert!(text.ends_with(ERASE_LINE));
assert!(!text.contains("\x1b[?25l"), "the cursor is never hidden");
let before = out.contents();
ticker.stop();
drop(ticker);
assert_eq!(out.contents(), before);
}
#[test]
fn a_stopped_ticker_can_be_stopped_again_and_is_stopped_by_drop() {
let root = tempfile::tempdir().expect("tempdir");
let out = SharedBuffer::default();
let timing = Timing { delay: Duration::ZERO, tick: Duration::from_millis(1) };
let ticker = Ticker::start(plan(), scanned(root.path()), Instant::now(), io(&out, timing));
wait_for(&out, |bytes| !bytes.is_empty());
let line = Arc::clone(ticker.line());
drop(ticker);
assert!(line.lock().stopped);
assert!(!line.lock().frame_on_screen);
assert!(out.text().ends_with(ERASE_LINE));
let before = out.contents();
thread::sleep(Duration::from_millis(20));
assert_eq!(out.contents(), before, "a dropped ticker draws nothing more");
}
struct FailsAfterOneFrame {
accepted: SharedBuffer,
writes: usize,
}
impl Write for FailsAfterOneFrame {
fn write(&mut self, buffer: &[u8]) -> io::Result<usize> {
if self.writes >= 2 {
return Err(io::Error::other("stderr went away"));
}
self.writes += 1;
self.accepted.write(buffer)
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
#[test]
fn after_a_failed_write_nothing_more_is_written() {
let root = tempfile::tempdir().expect("tempdir");
let accepted = SharedBuffer::default();
let out = FailsAfterOneFrame { accepted: accepted.clone(), writes: 0 };
let timing = Timing { delay: Duration::ZERO, tick: Duration::from_millis(1) };
let io = ProgressIo { out: Box::new(out), timing, width: || 100, interrupt: |_, _| {} };
let mut ticker = Ticker::start(plan(), scanned(root.path()), Instant::now(), io);
wait_for(&accepted, |bytes| !bytes.is_empty());
let deadline = Instant::now() + Duration::from_secs(10);
while !ticker.line().lock().stopped {
assert!(Instant::now() < deadline, "the failed write never stopped the ticker");
thread::sleep(Duration::from_millis(1));
}
let first_frame = accepted.text();
assert!(
first_frame.starts_with(&format!("{ERASE_LINE}⠋ {ROOT} Summarizing")),
"{first_frame:?}"
);
assert_eq!(first_frame.matches(ERASE_LINE).count(), 1, "exactly one frame got through");
ticker.stop();
assert_eq!(accepted.text(), first_frame, "the erase at the stop point was refused too");
assert!(!ticker.line().lock().frame_on_screen);
}
#[test]
fn the_shipped_timing_is_the_plan_s_and_the_width_falls_back_to_eighty() {
let timing = Timing::default();
assert_eq!(timing.delay, Duration::from_millis(500));
assert_eq!(timing.tick, Duration::from_millis(100));
let width = stderr_width();
assert!(width >= 1, "a width of {width} fits nothing");
if !io::IsTerminal::is_terminal(&io::stderr()) {
assert_eq!(width, FALLBACK_WIDTH);
}
}
#[test]
fn a_new_line_shows_nothing_and_stop_without_a_frame_writes_nothing() {
let out = SharedBuffer::default();
let line = Line::new(Box::new(out.clone()));
{
let mut state = line.lock();
assert!(!state.stopped);
assert!(!state.frame_on_screen);
state.stop();
assert!(state.stopped);
}
assert!(out.contents().is_empty());
let line = Line::new(Box::new(out.clone()));
line.lock().draw("⠋ a frame");
assert!(line.lock().frame_on_screen);
line.lock().stop();
assert_eq!(out.text(), format!("{ERASE_LINE}⠋ a frame{ERASE_LINE}"));
line.lock().stop();
assert_eq!(out.text(), format!("{ERASE_LINE}⠋ a frame{ERASE_LINE}"), "idempotent");
}
}