1use std::sync::atomic::{AtomicU64, Ordering};
9use std::sync::Arc;
10use std::time::{Duration, Instant};
11
12#[derive(Debug, Default)]
13pub struct Counters {
14 pub logical_bytes: AtomicU64,
16 pub wire_bytes: AtomicU64,
18 pub chunks: AtomicU64,
19 pub chunks_compressed: AtomicU64,
20 pub chunks_skipped: AtomicU64,
22 pub chunks_zero: AtomicU64,
24 pub chunks_reused: AtomicU64,
26 pub compressor_runs: AtomicU64,
29 pub files_completed: AtomicU64,
30 pub files_total: AtomicU64,
31 pub bytes_total: AtomicU64,
32}
33
34#[derive(Clone)]
36pub struct Metrics {
37 counters: Arc<Counters>,
38 start: Instant,
39}
40
41impl Default for Metrics {
42 fn default() -> Self {
43 Self::new()
44 }
45}
46
47impl Metrics {
48 pub fn new() -> Self {
49 Self {
50 counters: Arc::new(Counters::default()),
51 start: Instant::now(),
52 }
53 }
54
55 pub fn set_totals(&self, files: u64, bytes: u64) {
56 self.counters.files_total.store(files, Ordering::Relaxed);
57 self.counters.bytes_total.store(bytes, Ordering::Relaxed);
58 }
59
60 #[inline]
61 pub fn chunk_done(&self, logical: u64, wire: u64, compressed: bool) {
62 let c = &self.counters;
63 c.logical_bytes.fetch_add(logical, Ordering::Relaxed);
64 c.wire_bytes.fetch_add(wire, Ordering::Relaxed);
65 c.chunks.fetch_add(1, Ordering::Relaxed);
66 if compressed {
67 c.chunks_compressed.fetch_add(1, Ordering::Relaxed);
68 }
69 }
70
71 #[inline]
74 pub fn chunk_zero(&self, logical: u64, wire: u64) {
75 self.counters.chunks_zero.fetch_add(1, Ordering::Relaxed);
76 self.chunk_done(logical, wire, false);
77 }
78
79 #[inline]
81 pub fn compressor_ran(&self) {
82 self.counters
83 .compressor_runs
84 .fetch_add(1, Ordering::Relaxed);
85 }
86
87 #[inline]
90 pub fn chunk_reused(&self, logical: u64, wire: u64) {
91 self.counters.chunks_reused.fetch_add(1, Ordering::Relaxed);
92 self.chunk_done(logical, wire, false);
93 }
94
95 #[inline]
96 pub fn chunk_skipped(&self, logical: u64) {
97 self.counters.chunks_skipped.fetch_add(1, Ordering::Relaxed);
98 self.counters
101 .logical_bytes
102 .fetch_add(logical, Ordering::Relaxed);
103 }
104
105 #[inline]
106 pub fn file_done(&self) {
107 self.counters
108 .files_completed
109 .fetch_add(1, Ordering::Relaxed);
110 }
111
112 pub fn elapsed(&self) -> Duration {
113 self.start.elapsed()
114 }
115
116 pub fn snapshot(&self) -> Progress {
117 let c = &self.counters;
118 let elapsed = self.start.elapsed();
119 let logical = c.logical_bytes.load(Ordering::Relaxed);
120 let wire = c.wire_bytes.load(Ordering::Relaxed);
121 Progress {
122 logical_bytes: logical,
123 wire_bytes: wire,
124 bytes_total: c.bytes_total.load(Ordering::Relaxed),
125 chunks: c.chunks.load(Ordering::Relaxed),
126 chunks_compressed: c.chunks_compressed.load(Ordering::Relaxed),
127 chunks_skipped: c.chunks_skipped.load(Ordering::Relaxed),
128 chunks_zero: c.chunks_zero.load(Ordering::Relaxed),
129 chunks_reused: c.chunks_reused.load(Ordering::Relaxed),
130 compressor_runs: c.compressor_runs.load(Ordering::Relaxed),
131 files_completed: c.files_completed.load(Ordering::Relaxed),
132 files_total: c.files_total.load(Ordering::Relaxed),
133 elapsed,
134 }
135 }
136}
137
138#[derive(Debug, Clone, Copy, PartialEq, Eq)]
140pub struct Progress {
141 pub logical_bytes: u64,
143 pub wire_bytes: u64,
145 pub bytes_total: u64,
146 pub chunks: u64,
147 pub chunks_compressed: u64,
148 pub chunks_skipped: u64,
149 pub chunks_zero: u64,
150 pub chunks_reused: u64,
152 pub compressor_runs: u64,
155 pub files_completed: u64,
156 pub files_total: u64,
157 pub elapsed: Duration,
158}
159
160impl Progress {
161 pub fn throughput(&self) -> f64 {
165 let s = self.elapsed.as_secs_f64();
166 if s <= 0.0 {
167 return 0.0;
168 }
169 self.logical_bytes as f64 / s
170 }
171
172 pub fn wire_throughput(&self) -> f64 {
174 let s = self.elapsed.as_secs_f64();
175 if s <= 0.0 {
176 return 0.0;
177 }
178 self.wire_bytes as f64 / s
179 }
180
181 pub fn compression_ratio(&self) -> f64 {
183 if self.wire_bytes == 0 {
184 return 1.0;
185 }
186 self.logical_bytes as f64 / self.wire_bytes as f64
187 }
188
189 pub fn fraction(&self) -> f64 {
191 if self.bytes_total == 0 {
192 return if self.files_total > 0 && self.files_completed >= self.files_total {
193 1.0
194 } else {
195 0.0
196 };
197 }
198 (self.logical_bytes as f64 / self.bytes_total as f64).min(1.0)
199 }
200
201 pub fn eta(&self) -> Option<Duration> {
203 let rate = self.throughput();
204 if rate <= 0.0 || self.bytes_total == 0 {
205 return None;
206 }
207 let left = self.bytes_total.saturating_sub(self.logical_bytes);
208 Some(Duration::from_secs_f64(left as f64 / rate))
209 }
210}
211
212impl std::fmt::Display for Progress {
213 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
214 write!(
215 f,
216 "{}/{} files, {} of {} ({:.1}%), {}/s wire {}/s, ratio {:.2}x, {} chunks ({} zero, {} skipped, {} compressed of {} tried), {:.1}s",
217 self.files_completed,
218 self.files_total,
219 human_bytes(self.logical_bytes),
220 human_bytes(self.bytes_total),
221 self.fraction() * 100.0,
222 human_bytes(self.throughput() as u64),
223 human_bytes(self.wire_throughput() as u64),
224 self.compression_ratio(),
225 self.chunks,
226 self.chunks_zero,
227 self.chunks_skipped,
228 self.chunks_compressed,
229 self.compressor_runs,
230 self.elapsed.as_secs_f64(),
231 )
232 }
233}
234
235pub fn human_bytes(n: u64) -> String {
236 const UNITS: [&str; 6] = ["B", "KiB", "MiB", "GiB", "TiB", "PiB"];
237 let mut v = n as f64;
238 let mut i = 0;
239 while v >= 1024.0 && i < UNITS.len() - 1 {
240 v /= 1024.0;
241 i += 1;
242 }
243 if i == 0 {
244 format!("{n} B")
245 } else {
246 format!("{v:.2} {}", UNITS[i])
247 }
248}
249
250pub type ProgressFn = Arc<dyn Fn(Progress) + Send + Sync>;
252
253#[cfg(test)]
254mod tests {
255 use super::*;
256
257 #[test]
258 fn counters_accumulate_across_threads() {
259 let m = Metrics::new();
260 m.set_totals(4, 4_000_000);
261 let threads: Vec<_> = (0..4)
262 .map(|_| {
263 let m = m.clone();
264 std::thread::spawn(move || {
265 for _ in 0..1000 {
266 m.chunk_done(1000, 500, true);
267 }
268 m.file_done();
269 })
270 })
271 .collect();
272 for t in threads {
273 t.join().unwrap();
274 }
275 let p = m.snapshot();
276 assert_eq!(p.chunks, 4000);
277 assert_eq!(p.logical_bytes, 4_000_000);
278 assert_eq!(p.wire_bytes, 2_000_000);
279 assert_eq!(p.chunks_compressed, 4000);
280 assert_eq!(p.files_completed, 4);
281 assert_eq!(p.compression_ratio(), 2.0);
282 assert_eq!(p.fraction(), 1.0);
283 }
284
285 #[test]
286 fn skipped_chunks_still_count_as_progress() {
287 let m = Metrics::new();
288 m.set_totals(1, 1000);
289 m.chunk_skipped(1000);
290 let p = m.snapshot();
291 assert_eq!(p.chunks_skipped, 1);
292 assert_eq!(p.wire_bytes, 0);
293 assert_eq!(p.fraction(), 1.0);
294 }
295
296 #[test]
297 fn empty_transfer_does_not_divide_by_zero() {
298 let m = Metrics::new();
299 let p = m.snapshot();
300 assert_eq!(p.fraction(), 0.0);
301 assert_eq!(p.compression_ratio(), 1.0);
302 assert!(p.eta().is_none());
303 m.set_totals(1, 0);
304 m.file_done();
305 assert_eq!(m.snapshot().fraction(), 1.0);
306 }
307
308 #[test]
309 fn fraction_is_clamped() {
310 let m = Metrics::new();
311 m.set_totals(1, 100);
312 m.chunk_done(500, 500, false);
313 assert_eq!(m.snapshot().fraction(), 1.0);
314 }
315
316 #[test]
317 fn human_bytes_formats_sensibly() {
318 assert_eq!(human_bytes(0), "0 B");
319 assert_eq!(human_bytes(512), "512 B");
320 assert_eq!(human_bytes(1024), "1.00 KiB");
321 assert_eq!(human_bytes(1536), "1.50 KiB");
322 assert_eq!(human_bytes(100 * 1024 * 1024 * 1024), "100.00 GiB");
323 assert_eq!(human_bytes(u64::MAX), "16384.00 PiB");
325 }
326}