use std::collections::HashMap;
use std::path::PathBuf;
use std::time::{Duration, Instant};
use crate::progress::Progress;
const SAMPLE_WINDOW: Duration = Duration::from_millis(500);
const EWMA_ALPHA: f64 = 0.3;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Regime {
Directory,
SmallFile,
LargeFile,
}
#[derive(Debug, Clone, Default)]
struct Rate {
ewma: Option<f64>,
pending_work: f64,
pending_secs: f64,
}
impl Rate {
fn add_work(&mut self, work: f64) {
self.pending_work += work;
}
fn add_time(&mut self, secs: f64) {
self.pending_secs += secs;
if self.pending_secs >= SAMPLE_WINDOW.as_secs_f64() {
let sample = self.pending_work / self.pending_secs;
self.ewma = Some(match self.ewma {
Some(previous) => previous * (1.0 - EWMA_ALPHA) + sample * EWMA_ALPHA,
None => sample,
});
self.pending_work = 0.0;
self.pending_secs = 0.0;
}
}
fn per_sec(&self) -> Option<f64> {
if let Some(ewma) = self.ewma {
if ewma > 0.0 {
return Some(ewma);
}
}
if self.pending_secs > 0.0 && self.pending_work > 0.0 {
return Some(self.pending_work / self.pending_secs);
}
None
}
}
#[derive(Debug, Clone)]
pub struct EtaEstimator {
small_file_threshold: u64,
directories_remaining: usize,
small_files_remaining: usize,
large_bytes_remaining: u64,
directory_rate: Rate,
small_file_rate: Rate,
large_file_rate: Rate,
overall_byte_rate: Rate,
small_in_flight: usize,
large_in_flight: usize,
directories_in_flight: bool,
credited_bytes: HashMap<PathBuf, u64>,
last_event: Option<Instant>,
awaiting_planned_start: bool,
}
impl Default for EtaEstimator {
fn default() -> Self {
Self::new()
}
}
impl EtaEstimator {
pub fn new() -> Self {
Self {
small_file_threshold: 0,
directories_remaining: 0,
small_files_remaining: 0,
large_bytes_remaining: 0,
directory_rate: Rate::default(),
small_file_rate: Rate::default(),
large_file_rate: Rate::default(),
overall_byte_rate: Rate::default(),
small_in_flight: 0,
large_in_flight: 0,
directories_in_flight: false,
credited_bytes: HashMap::new(),
last_event: None,
awaiting_planned_start: false,
}
}
pub fn observe(&mut self, progress: &Progress) {
self.observe_at(progress, Instant::now());
}
fn observe_at(&mut self, progress: &Progress, now: Instant) {
let was_in_flight = self.in_flight_regimes();
match progress {
Progress::Planned {
directories,
small_files,
large_bytes,
small_file_threshold,
..
} => {
self.small_file_threshold = *small_file_threshold;
self.directories_remaining = *directories;
self.small_files_remaining = *small_files;
self.large_bytes_remaining = *large_bytes;
self.awaiting_planned_start = true;
}
Progress::DirectoriesStarted { total } => {
self.directories_remaining = *total;
self.directories_in_flight = *total > 0;
}
Progress::DirectoryCompleted { .. } | Progress::DirectoryFailed { .. } => {
self.directory_rate.add_work(1.0);
self.directories_remaining = self.directories_remaining.saturating_sub(1);
if self.directories_remaining == 0 {
self.directories_in_flight = false;
}
}
Progress::Started { entries_total, .. } => {
self.directories_in_flight = false;
if self.awaiting_planned_start {
self.awaiting_planned_start = false;
} else {
self.small_files_remaining = *entries_total;
self.large_bytes_remaining = 0;
self.small_file_threshold = u64::MAX;
}
}
Progress::EntryStarted { entry } => {
if self.regime_for(entry.size) == Regime::LargeFile {
self.large_in_flight += 1;
} else {
self.small_in_flight += 1;
}
}
Progress::EntryCompleted { entry } | Progress::EntryFailed { entry } => {
if self.regime_for(entry.size) == Regime::LargeFile {
let outstanding = entry
.size
.saturating_sub(self.credited_bytes.remove(&entry.path).unwrap_or(0));
self.large_file_rate.add_work(outstanding as f64);
self.overall_byte_rate.add_work(outstanding as f64);
self.large_bytes_remaining =
self.large_bytes_remaining.saturating_sub(outstanding);
self.large_in_flight = self.large_in_flight.saturating_sub(1);
} else {
self.overall_byte_rate.add_work(entry.size as f64);
self.small_file_rate.add_work(1.0);
self.small_files_remaining = self.small_files_remaining.saturating_sub(1);
self.small_in_flight = self.small_in_flight.saturating_sub(1);
}
}
Progress::EntryProgress {
entry,
bytes_copied,
} => {
let credited = self.credited_bytes.entry(entry.path.clone()).or_insert(0);
let delta = bytes_copied.saturating_sub(*credited);
if delta > 0 {
*credited = *bytes_copied;
self.large_file_rate.add_work(delta as f64);
self.overall_byte_rate.add_work(delta as f64);
self.large_bytes_remaining = self.large_bytes_remaining.saturating_sub(delta);
}
}
}
if let Some(last) = self.last_event {
let elapsed = now.saturating_duration_since(last).as_secs_f64();
for regime in &was_in_flight {
self.rate_mut(*regime).add_time(elapsed);
}
if was_in_flight
.iter()
.any(|r| matches!(r, Regime::SmallFile | Regime::LargeFile))
{
self.overall_byte_rate.add_time(elapsed);
}
}
self.last_event = Some(now);
}
pub fn estimate(&self) -> Option<Duration> {
let seconds_for = |remaining: f64, rate: &Rate| -> Option<f64> {
if remaining <= 0.0 {
return Some(0.0);
}
Some(remaining / rate.per_sec()?)
};
let directories = seconds_for(self.directories_remaining as f64, &self.directory_rate)?;
let small = seconds_for(self.small_files_remaining as f64, &self.small_file_rate)?;
let large_rate = self
.large_file_rate
.per_sec()
.or_else(|| self.overall_byte_rate.per_sec());
let seconds_for_bytes = |bytes: u64| -> Option<f64> {
match (bytes, large_rate) {
(0, _) => Some(0.0),
(bytes, Some(rate)) => Some(bytes as f64 / rate),
(_, None) => None,
}
};
let large = seconds_for_bytes(self.large_bytes_remaining)?;
Duration::try_from_secs_f64(directories + small.max(large)).ok()
}
pub fn bytes_per_sec(&self) -> Option<f64> {
self.large_file_rate.per_sec()
}
fn in_flight_regimes(&self) -> Vec<Regime> {
let mut regimes = Vec::with_capacity(3);
if self.directories_in_flight {
regimes.push(Regime::Directory);
}
if self.small_in_flight > 0 {
regimes.push(Regime::SmallFile);
}
if self.large_in_flight > 0 {
regimes.push(Regime::LargeFile);
}
regimes
}
fn regime_for(&self, size: u64) -> Regime {
if size <= self.small_file_threshold {
Regime::SmallFile
} else {
Regime::LargeFile
}
}
fn rate_mut(&mut self, regime: Regime) -> &mut Rate {
match regime {
Regime::Directory => &mut self.directory_rate,
Regime::SmallFile => &mut self.small_file_rate,
Regime::LargeFile => &mut self.large_file_rate,
}
}
}
#[cfg(test)]
mod tests {
use std::path::PathBuf;
use crate::profiler::Entry;
use super::*;
fn entry(size: u64) -> Entry {
Entry {
path: PathBuf::from("x"),
relative_path: PathBuf::from("x"),
size,
modified: None,
}
}
fn planned(directories: usize, small_files: usize, large_bytes: u64) -> Progress {
Progress::Planned {
directories,
small_files,
small_bytes: small_files as u64,
large_files: usize::from(large_bytes > 0),
large_bytes,
small_file_threshold: 1024,
}
}
fn started() -> Progress {
Progress::Started {
bytes_total: Some(0),
entries_total: 0,
}
}
fn replay(estimator: &mut EtaEstimator, clock: &mut Instant, script: &[(Progress, f64)]) {
for (event, delay) in script {
*clock += Duration::from_secs_f64(*delay);
estimator.observe_at(event, *clock);
}
}
#[test]
fn no_estimate_before_anything_completes() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
replay(
&mut eta,
&mut clock,
&[(planned(0, 10, 0), 0.0), (started(), 0.0)],
);
assert_eq!(eta.estimate(), None);
}
#[test]
fn estimates_zero_when_nothing_is_outstanding() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
replay(
&mut eta,
&mut clock,
&[(planned(0, 0, 0), 0.0), (started(), 0.0)],
);
assert_eq!(eta.estimate(), Some(Duration::ZERO));
}
#[test]
fn small_files_are_estimated_per_file_not_per_byte() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
let mut script = vec![(planned(0, 100, 0), 0.0), (started(), 0.0)];
for i in 0..10 {
let size = if i % 2 == 0 { 10 } else { 1000 };
script.push((Progress::EntryStarted { entry: entry(size) }, 0.0));
script.push((Progress::EntryCompleted { entry: entry(size) }, 0.1));
}
replay(&mut eta, &mut clock, &script);
let estimate = eta.estimate().unwrap().as_secs_f64();
assert!(
(estimate - 9.0).abs() < 0.5,
"expected ~9s for 90 files at 10 files/sec, got {estimate}"
);
}
#[test]
fn large_files_are_estimated_per_byte() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
let big = 10_000_u64;
let mut script = vec![(planned(0, 0, big * 10), 0.0), (started(), 0.0)];
for _ in 0..4 {
script.push((Progress::EntryStarted { entry: entry(big) }, 0.0));
script.push((Progress::EntryCompleted { entry: entry(big) }, 1.0));
}
replay(&mut eta, &mut clock, &script);
let estimate = eta.estimate().unwrap().as_secs_f64();
assert!(
(estimate - 6.0).abs() < 0.5,
"expected ~6s for 60_000 bytes at 10_000 B/s, got {estimate}"
);
}
#[test]
fn directory_pre_pass_is_estimated_before_any_file_work_is_known() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
let mut script = vec![
(planned(100, 0, 0), 0.0),
(Progress::DirectoriesStarted { total: 100 }, 0.0),
];
for _ in 0..20 {
script.push((
Progress::DirectoryCompleted {
path: PathBuf::from("d"),
},
0.1,
));
}
replay(&mut eta, &mut clock, &script);
let estimate = eta.estimate().unwrap().as_secs_f64();
assert!(
(estimate - 8.0).abs() < 0.5,
"expected ~8s for 80 directories at 10/sec, got {estimate}"
);
}
#[test]
fn directory_and_file_costs_are_summed_not_conflated() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
let mut script = vec![
(planned(30, 30, 0), 0.0),
(Progress::DirectoriesStarted { total: 30 }, 0.0),
];
for _ in 0..10 {
script.push((
Progress::DirectoryCompleted {
path: PathBuf::from("d"),
},
0.1,
));
}
script.push((started(), 0.0));
for _ in 0..10 {
script.push((Progress::EntryStarted { entry: entry(10) }, 0.0));
script.push((Progress::EntryCompleted { entry: entry(10) }, 0.2));
}
replay(&mut eta, &mut clock, &script);
let estimate = eta.estimate().unwrap().as_secs_f64();
assert!(
(estimate - 6.0).abs() < 0.7,
"expected ~6s (2s of directories + 4s of files), got {estimate}"
);
}
#[test]
fn failed_entries_count_as_progress() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
let mut script = vec![(planned(0, 20, 0), 0.0), (started(), 0.0)];
for _ in 0..10 {
script.push((Progress::EntryStarted { entry: entry(10) }, 0.0));
script.push((Progress::EntryFailed { entry: entry(10) }, 0.1));
}
replay(&mut eta, &mut clock, &script);
let estimate = eta.estimate().unwrap().as_secs_f64();
assert!((estimate - 1.0).abs() < 0.3, "expected ~1s, got {estimate}");
}
#[test]
fn started_without_planned_is_modelled_as_a_metadata_only_phase() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
let mut script = vec![(
Progress::Started {
bytes_total: None,
entries_total: 100,
},
0.0,
)];
for _ in 0..10 {
script.push((
Progress::EntryStarted {
entry: entry(5_000_000),
},
0.0,
));
script.push((
Progress::EntryCompleted {
entry: entry(5_000_000),
},
0.1,
));
}
replay(&mut eta, &mut clock, &script);
let estimate = eta.estimate().unwrap().as_secs_f64();
assert!(
(estimate - 9.0).abs() < 0.5,
"expected ~9s for 90 deletions at 10/sec, got {estimate}"
);
}
#[test]
fn a_later_phase_resets_remaining_work_without_discarding_learned_rates() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
let mut script = vec![(planned(0, 10, 0), 0.0), (started(), 0.0)];
for _ in 0..10 {
script.push((Progress::EntryStarted { entry: entry(10) }, 0.0));
script.push((Progress::EntryCompleted { entry: entry(10) }, 0.1));
}
replay(&mut eta, &mut clock, &script);
assert_eq!(eta.estimate(), Some(Duration::ZERO));
replay(
&mut eta,
&mut clock,
&[(
Progress::Started {
bytes_total: None,
entries_total: 5,
},
0.0,
)],
);
let estimate = eta.estimate().unwrap().as_secs_f64();
assert!(
(estimate - 0.5).abs() < 0.2,
"expected ~0.5s for 5 entries at 10/sec, got {estimate}"
);
}
#[test]
fn slowdown_is_tracked_rather_than_averaged_away() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
let mut script = vec![(planned(0, 1000, 0), 0.0), (started(), 0.0)];
for _ in 0..100 {
script.push((Progress::EntryStarted { entry: entry(10) }, 0.0));
script.push((Progress::EntryCompleted { entry: entry(10) }, 0.01));
}
replay(&mut eta, &mut clock, &script);
let fast = eta.estimate().unwrap();
let mut script = Vec::new();
for _ in 0..100 {
script.push((Progress::EntryStarted { entry: entry(10) }, 0.0));
script.push((Progress::EntryCompleted { entry: entry(10) }, 0.1));
}
replay(&mut eta, &mut clock, &script);
let slow = eta.estimate().unwrap();
assert!(
slow > fast * 3,
"estimate should track the slowdown: {fast:?} -> {slow:?}"
);
}
#[test]
fn outstanding_large_files_are_costed_from_overall_byte_throughput() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
let plan = Progress::Planned {
directories: 0,
small_files: 100,
small_bytes: 10_000_000,
large_files: 1,
large_bytes: 100_000_000,
small_file_threshold: 1_000_000,
};
let mut script = vec![
(plan, 0.0),
(started(), 0.0),
(
Progress::EntryStarted {
entry: entry(100_000_000),
},
0.0,
),
];
for _ in 0..50 {
script.push((
Progress::EntryStarted {
entry: entry(100_000),
},
0.0,
));
script.push((
Progress::EntryCompleted {
entry: entry(100_000),
},
0.1,
));
}
replay(&mut eta, &mut clock, &script);
let estimate = eta
.estimate()
.expect("an unfinished large file must not suppress the estimate")
.as_secs_f64();
assert!(
(estimate - 100.0).abs() < 15.0,
"expected ~100s dominated by the outstanding large file, got {estimate}"
);
}
#[test]
fn a_measured_large_file_rate_supersedes_the_overall_fallback() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
let plan = Progress::Planned {
directories: 0,
small_files: 0,
small_bytes: 0,
large_files: 3,
large_bytes: 300_000_000,
small_file_threshold: 1_000_000,
};
replay(&mut eta, &mut clock, &[(plan, 0.0), (started(), 0.0)]);
replay(
&mut eta,
&mut clock,
&[
(
Progress::EntryStarted {
entry: entry(100_000_000),
},
0.0,
),
(
Progress::EntryCompleted {
entry: entry(100_000_000),
},
1.0,
),
],
);
assert_eq!(eta.bytes_per_sec(), Some(100_000_000.0));
let estimate = eta.estimate().unwrap().as_secs_f64();
assert!(
(estimate - 2.0).abs() < 0.3,
"expected ~2s for the remaining 200MB at 100MB/s, got {estimate}"
);
}
#[test]
fn a_single_large_file_is_estimated_from_in_flight_samples() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
let plan = Progress::Planned {
directories: 0,
small_files: 0,
small_bytes: 0,
large_files: 1,
large_bytes: 1_000_000_000,
small_file_threshold: 1_000_000,
};
let big = entry(1_000_000_000);
let mut script = vec![
(plan, 0.0),
(started(), 0.0),
(Progress::EntryStarted { entry: big.clone() }, 0.0),
];
for i in 1..=4 {
script.push((
Progress::EntryProgress {
entry: big.clone(),
bytes_copied: i * 25_000_000,
},
0.25,
));
}
replay(&mut eta, &mut clock, &script);
let estimate = eta
.estimate()
.expect("in-flight samples must produce an estimate")
.as_secs_f64();
assert!(
(estimate - 9.0).abs() < 1.0,
"expected ~9s for the outstanding 900MB at 100MB/s, got {estimate}"
);
assert_eq!(eta.bytes_per_sec(), Some(100_000_000.0));
}
#[test]
fn sampled_bytes_are_not_counted_again_on_completion() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
let plan = Progress::Planned {
directories: 0,
small_files: 0,
small_bytes: 0,
large_files: 2,
large_bytes: 200_000_000,
small_file_threshold: 1_000_000,
};
let big = entry(100_000_000);
replay(
&mut eta,
&mut clock,
&[
(plan, 0.0),
(started(), 0.0),
(Progress::EntryStarted { entry: big.clone() }, 0.0),
(
Progress::EntryProgress {
entry: big.clone(),
bytes_copied: 25_000_000,
},
0.25,
),
(
Progress::EntryProgress {
entry: big.clone(),
bytes_copied: 50_000_000,
},
0.25,
),
(
Progress::EntryProgress {
entry: big.clone(),
bytes_copied: 75_000_000,
},
0.25,
),
(
Progress::EntryProgress {
entry: big.clone(),
bytes_copied: 100_000_000,
},
0.25,
),
(Progress::EntryCompleted { entry: big }, 0.0),
],
);
assert_eq!(eta.bytes_per_sec(), Some(100_000_000.0));
let estimate = eta.estimate().unwrap().as_secs_f64();
assert!(
(estimate - 1.0).abs() < 0.2,
"expected ~1s for the remaining 100MB at 100MB/s, got {estimate}"
);
}
#[test]
fn out_of_order_and_surplus_completions_do_not_panic() {
let mut eta = EtaEstimator::new();
let mut clock = Instant::now();
replay(
&mut eta,
&mut clock,
&[
(Progress::EntryCompleted { entry: entry(10) }, 0.1),
(
Progress::DirectoryCompleted {
path: PathBuf::from("d"),
},
0.1,
),
(planned(0, 1, 0), 0.0),
(started(), 0.0),
(Progress::EntryCompleted { entry: entry(10) }, 0.1),
(Progress::EntryCompleted { entry: entry(10) }, 0.1),
],
);
assert_eq!(eta.estimate(), Some(Duration::ZERO));
}
}