Skip to main content

runsync_transfer/
metrics.rs

1//! Progress and throughput accounting.
2//!
3//! Counters are plain atomics updated from every worker. They are `Relaxed`
4//! because nothing branches on them — they exist to be read by an observer, and
5//! paying for ordering on a per-chunk counter would show up in the throughput
6//! numbers it is meant to measure.
7
8use std::sync::atomic::{AtomicU64, Ordering};
9use std::sync::Arc;
10use std::time::{Duration, Instant};
11
12#[derive(Debug, Default)]
13pub struct Counters {
14    /// Payload bytes as they exist on disk.
15    pub logical_bytes: AtomicU64,
16    /// Bytes handed to the transport, after compression and sealing.
17    pub wire_bytes: AtomicU64,
18    pub chunks: AtomicU64,
19    pub chunks_compressed: AtomicU64,
20    /// Chunks skipped because the receiver already had them.
21    pub chunks_skipped: AtomicU64,
22    /// Chunks that were entirely zero and crossed the wire as a flag.
23    pub chunks_zero: AtomicU64,
24    /// Chunks the receiver already had, sent as a reference rather than data.
25    pub chunks_reused: AtomicU64,
26    /// Chunks the compressor actually ran on, whether or not the result was
27    /// kept. The gap between this and `chunks_compressed` is wasted CPU.
28    pub compressor_runs: AtomicU64,
29    pub files_completed: AtomicU64,
30    pub files_total: AtomicU64,
31    pub bytes_total: AtomicU64,
32}
33
34/// Shared, cloneable metrics handle.
35#[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    /// A hole: counted as progress and as a chunk, but the payload never
72    /// existed on the wire.
73    #[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    /// Record that the compressor ran on a chunk, regardless of the outcome.
80    #[inline]
81    pub fn compressor_ran(&self) {
82        self.counters
83            .compressor_runs
84            .fetch_add(1, Ordering::Relaxed);
85    }
86
87    /// A chunk the far side already held: counted as progress, but the payload
88    /// never crossed the wire.
89    #[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        // Skipped bytes count as progress: from the caller's point of view the
99        // file is that much closer to done, even though nothing crossed the wire.
100        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/// A point-in-time view of a transfer.
139#[derive(Debug, Clone, Copy, PartialEq, Eq)]
140pub struct Progress {
141    /// Payload bytes transferred, measured as they exist on disk.
142    pub logical_bytes: u64,
143    /// Bytes actually put on the wire.
144    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    /// Chunks the receiver already had; the payload never crossed the wire.
151    pub chunks_reused: u64,
152    /// Chunks handed to the compressor. Compare with `chunks_compressed`: the
153    /// difference is compression work whose output was discarded.
154    pub compressor_runs: u64,
155    pub files_completed: u64,
156    pub files_total: u64,
157    pub elapsed: Duration,
158}
159
160impl Progress {
161    /// Effective throughput in bytes/sec, measured against logical bytes. This
162    /// is the number a user cares about: how fast their data moved, not how
163    /// many packets it took.
164    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    /// Throughput measured against bytes actually sent.
173    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    /// Compression ratio achieved, logical / wire. 1.0 means no saving.
182    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    /// Fraction of the transfer complete, 0.0..=1.0.
190    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    /// Estimated time remaining, from the average rate so far.
202    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
250/// Callback invoked periodically with a fresh snapshot.
251pub 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        // Units stop at PiB rather than inventing an exabyte suffix.
324        assert_eq!(human_bytes(u64::MAX), "16384.00 PiB");
325    }
326}