runsync-transfer 2026.1.0

High-throughput P2P file transfer engine: adaptive compression, end-to-end AEAD, parallel chunked pipeline over QUIC or any async transport.
Documentation
//! Progress and throughput accounting.
//!
//! Counters are plain atomics updated from every worker. They are `Relaxed`
//! because nothing branches on them — they exist to be read by an observer, and
//! paying for ordering on a per-chunk counter would show up in the throughput
//! numbers it is meant to measure.

use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};

#[derive(Debug, Default)]
pub struct Counters {
    /// Payload bytes as they exist on disk.
    pub logical_bytes: AtomicU64,
    /// Bytes handed to the transport, after compression and sealing.
    pub wire_bytes: AtomicU64,
    pub chunks: AtomicU64,
    pub chunks_compressed: AtomicU64,
    /// Chunks skipped because the receiver already had them.
    pub chunks_skipped: AtomicU64,
    /// Chunks that were entirely zero and crossed the wire as a flag.
    pub chunks_zero: AtomicU64,
    /// Chunks the receiver already had, sent as a reference rather than data.
    pub chunks_reused: AtomicU64,
    /// Chunks the compressor actually ran on, whether or not the result was
    /// kept. The gap between this and `chunks_compressed` is wasted CPU.
    pub compressor_runs: AtomicU64,
    pub files_completed: AtomicU64,
    pub files_total: AtomicU64,
    pub bytes_total: AtomicU64,
}

/// Shared, cloneable metrics handle.
#[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);
        }
    }

    /// A hole: counted as progress and as a chunk, but the payload never
    /// existed on the wire.
    #[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);
    }

    /// Record that the compressor ran on a chunk, regardless of the outcome.
    #[inline]
    pub fn compressor_ran(&self) {
        self.counters
            .compressor_runs
            .fetch_add(1, Ordering::Relaxed);
    }

    /// A chunk the far side already held: counted as progress, but the payload
    /// never crossed the wire.
    #[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);
        // Skipped bytes count as progress: from the caller's point of view the
        // file is that much closer to done, even though nothing crossed the wire.
        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,
        }
    }
}

/// A point-in-time view of a transfer.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Progress {
    /// Payload bytes transferred, measured as they exist on disk.
    pub logical_bytes: u64,
    /// Bytes actually put on the wire.
    pub wire_bytes: u64,
    pub bytes_total: u64,
    pub chunks: u64,
    pub chunks_compressed: u64,
    pub chunks_skipped: u64,
    pub chunks_zero: u64,
    /// Chunks the receiver already had; the payload never crossed the wire.
    pub chunks_reused: u64,
    /// Chunks handed to the compressor. Compare with `chunks_compressed`: the
    /// difference is compression work whose output was discarded.
    pub compressor_runs: u64,
    pub files_completed: u64,
    pub files_total: u64,
    pub elapsed: Duration,
}

impl Progress {
    /// Effective throughput in bytes/sec, measured against logical bytes. This
    /// is the number a user cares about: how fast their data moved, not how
    /// many packets it took.
    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
    }

    /// Throughput measured against bytes actually sent.
    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
    }

    /// Compression ratio achieved, logical / wire. 1.0 means no saving.
    pub fn compression_ratio(&self) -> f64 {
        if self.wire_bytes == 0 {
            return 1.0;
        }
        self.logical_bytes as f64 / self.wire_bytes as f64
    }

    /// Fraction of the transfer complete, 0.0..=1.0.
    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)
    }

    /// Estimated time remaining, from the average rate so far.
    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])
    }
}

/// Callback invoked periodically with a fresh snapshot.
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");
        // Units stop at PiB rather than inventing an exabyte suffix.
        assert_eq!(human_bytes(u64::MAX), "16384.00 PiB");
    }
}