use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
#[derive(Debug, Default)]
pub struct Counters {
pub logical_bytes: AtomicU64,
pub wire_bytes: AtomicU64,
pub chunks: AtomicU64,
pub chunks_compressed: AtomicU64,
pub chunks_skipped: AtomicU64,
pub chunks_zero: AtomicU64,
pub chunks_reused: AtomicU64,
pub compressor_runs: AtomicU64,
pub files_completed: AtomicU64,
pub files_total: AtomicU64,
pub bytes_total: AtomicU64,
}
#[derive(Clone)]
pub struct Metrics {
counters: Arc<Counters>,
start: Instant,
}
impl Default for Metrics {
fn default() -> Self {
Self::new()
}
}
impl Metrics {
pub fn new() -> Self {
Self {
counters: Arc::new(Counters::default()),
start: Instant::now(),
}
}
pub fn set_totals(&self, files: u64, bytes: u64) {
self.counters.files_total.store(files, Ordering::Relaxed);
self.counters.bytes_total.store(bytes, Ordering::Relaxed);
}
#[inline]
pub fn chunk_done(&self, logical: u64, wire: u64, compressed: bool) {
let c = &self.counters;
c.logical_bytes.fetch_add(logical, Ordering::Relaxed);
c.wire_bytes.fetch_add(wire, Ordering::Relaxed);
c.chunks.fetch_add(1, Ordering::Relaxed);
if compressed {
c.chunks_compressed.fetch_add(1, Ordering::Relaxed);
}
}
#[inline]
pub fn chunk_zero(&self, logical: u64, wire: u64) {
self.counters.chunks_zero.fetch_add(1, Ordering::Relaxed);
self.chunk_done(logical, wire, false);
}
#[inline]
pub fn compressor_ran(&self) {
self.counters
.compressor_runs
.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn chunk_reused(&self, logical: u64, wire: u64) {
self.counters.chunks_reused.fetch_add(1, Ordering::Relaxed);
self.chunk_done(logical, wire, false);
}
#[inline]
pub fn chunk_skipped(&self, logical: u64) {
self.counters.chunks_skipped.fetch_add(1, Ordering::Relaxed);
self.counters
.logical_bytes
.fetch_add(logical, Ordering::Relaxed);
}
#[inline]
pub fn file_done(&self) {
self.counters
.files_completed
.fetch_add(1, Ordering::Relaxed);
}
pub fn elapsed(&self) -> Duration {
self.start.elapsed()
}
pub fn snapshot(&self) -> Progress {
let c = &self.counters;
let elapsed = self.start.elapsed();
let logical = c.logical_bytes.load(Ordering::Relaxed);
let wire = c.wire_bytes.load(Ordering::Relaxed);
Progress {
logical_bytes: logical,
wire_bytes: wire,
bytes_total: c.bytes_total.load(Ordering::Relaxed),
chunks: c.chunks.load(Ordering::Relaxed),
chunks_compressed: c.chunks_compressed.load(Ordering::Relaxed),
chunks_skipped: c.chunks_skipped.load(Ordering::Relaxed),
chunks_zero: c.chunks_zero.load(Ordering::Relaxed),
chunks_reused: c.chunks_reused.load(Ordering::Relaxed),
compressor_runs: c.compressor_runs.load(Ordering::Relaxed),
files_completed: c.files_completed.load(Ordering::Relaxed),
files_total: c.files_total.load(Ordering::Relaxed),
elapsed,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Progress {
pub logical_bytes: u64,
pub wire_bytes: u64,
pub bytes_total: u64,
pub chunks: u64,
pub chunks_compressed: u64,
pub chunks_skipped: u64,
pub chunks_zero: u64,
pub chunks_reused: u64,
pub compressor_runs: u64,
pub files_completed: u64,
pub files_total: u64,
pub elapsed: Duration,
}
impl Progress {
pub fn throughput(&self) -> f64 {
let s = self.elapsed.as_secs_f64();
if s <= 0.0 {
return 0.0;
}
self.logical_bytes as f64 / s
}
pub fn wire_throughput(&self) -> f64 {
let s = self.elapsed.as_secs_f64();
if s <= 0.0 {
return 0.0;
}
self.wire_bytes as f64 / s
}
pub fn compression_ratio(&self) -> f64 {
if self.wire_bytes == 0 {
return 1.0;
}
self.logical_bytes as f64 / self.wire_bytes as f64
}
pub fn fraction(&self) -> f64 {
if self.bytes_total == 0 {
return if self.files_total > 0 && self.files_completed >= self.files_total {
1.0
} else {
0.0
};
}
(self.logical_bytes as f64 / self.bytes_total as f64).min(1.0)
}
pub fn eta(&self) -> Option<Duration> {
let rate = self.throughput();
if rate <= 0.0 || self.bytes_total == 0 {
return None;
}
let left = self.bytes_total.saturating_sub(self.logical_bytes);
Some(Duration::from_secs_f64(left as f64 / rate))
}
}
impl std::fmt::Display for Progress {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"{}/{} files, {} of {} ({:.1}%), {}/s wire {}/s, ratio {:.2}x, {} chunks ({} zero, {} skipped, {} compressed of {} tried), {:.1}s",
self.files_completed,
self.files_total,
human_bytes(self.logical_bytes),
human_bytes(self.bytes_total),
self.fraction() * 100.0,
human_bytes(self.throughput() as u64),
human_bytes(self.wire_throughput() as u64),
self.compression_ratio(),
self.chunks,
self.chunks_zero,
self.chunks_skipped,
self.chunks_compressed,
self.compressor_runs,
self.elapsed.as_secs_f64(),
)
}
}
pub fn human_bytes(n: u64) -> String {
const UNITS: [&str; 6] = ["B", "KiB", "MiB", "GiB", "TiB", "PiB"];
let mut v = n as f64;
let mut i = 0;
while v >= 1024.0 && i < UNITS.len() - 1 {
v /= 1024.0;
i += 1;
}
if i == 0 {
format!("{n} B")
} else {
format!("{v:.2} {}", UNITS[i])
}
}
pub type ProgressFn = Arc<dyn Fn(Progress) + Send + Sync>;
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn counters_accumulate_across_threads() {
let m = Metrics::new();
m.set_totals(4, 4_000_000);
let threads: Vec<_> = (0..4)
.map(|_| {
let m = m.clone();
std::thread::spawn(move || {
for _ in 0..1000 {
m.chunk_done(1000, 500, true);
}
m.file_done();
})
})
.collect();
for t in threads {
t.join().unwrap();
}
let p = m.snapshot();
assert_eq!(p.chunks, 4000);
assert_eq!(p.logical_bytes, 4_000_000);
assert_eq!(p.wire_bytes, 2_000_000);
assert_eq!(p.chunks_compressed, 4000);
assert_eq!(p.files_completed, 4);
assert_eq!(p.compression_ratio(), 2.0);
assert_eq!(p.fraction(), 1.0);
}
#[test]
fn skipped_chunks_still_count_as_progress() {
let m = Metrics::new();
m.set_totals(1, 1000);
m.chunk_skipped(1000);
let p = m.snapshot();
assert_eq!(p.chunks_skipped, 1);
assert_eq!(p.wire_bytes, 0);
assert_eq!(p.fraction(), 1.0);
}
#[test]
fn empty_transfer_does_not_divide_by_zero() {
let m = Metrics::new();
let p = m.snapshot();
assert_eq!(p.fraction(), 0.0);
assert_eq!(p.compression_ratio(), 1.0);
assert!(p.eta().is_none());
m.set_totals(1, 0);
m.file_done();
assert_eq!(m.snapshot().fraction(), 1.0);
}
#[test]
fn fraction_is_clamped() {
let m = Metrics::new();
m.set_totals(1, 100);
m.chunk_done(500, 500, false);
assert_eq!(m.snapshot().fraction(), 1.0);
}
#[test]
fn human_bytes_formats_sensibly() {
assert_eq!(human_bytes(0), "0 B");
assert_eq!(human_bytes(512), "512 B");
assert_eq!(human_bytes(1024), "1.00 KiB");
assert_eq!(human_bytes(1536), "1.50 KiB");
assert_eq!(human_bytes(100 * 1024 * 1024 * 1024), "100.00 GiB");
assert_eq!(human_bytes(u64::MAX), "16384.00 PiB");
}
}