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
use crate::codec::compress::Algorithm;
use crate::error::{Error, Result};
use std::time::Duration;

/// How the engine decides whether to compress a given chunk.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CompressionMode {
    /// Never compress. Lowest CPU; best when the link is faster than the CPU
    /// (loopback, 25GbE) or every input is already compressed media.
    Off,
    /// Always compress with the configured algorithm.
    Always,
    /// Decide per chunk from a cheap entropy probe, and per file from its
    /// extension. This is the default and the right answer for mixed payloads.
    Adaptive,
}

/// Tuning for the compression stage.
#[derive(Debug, Clone)]
pub struct CompressionConfig {
    pub mode: CompressionMode,
    pub algorithm: Algorithm,
    /// zstd level. Ignored by LZ4. 1..=9 is the useful range for transfers;
    /// above ~6 the compressor becomes the bottleneck before the network does.
    pub level: i32,
    /// In `Adaptive` mode a chunk is sent raw unless compression saves at least
    /// this fraction. 0.06 means "must shrink by 6% to be worth it".
    pub min_gain: f32,
    /// Bytes sampled from the head of a chunk for the entropy probe.
    pub probe_bytes: usize,
    /// Extensions (lowercase, no dot) that skip compression entirely.
    pub incompressible_extensions: Vec<String>,
    /// Use the dedicated lossless coder for uncompressed PCM audio.
    ///
    /// zstd manages ~1.05x on `.wav`, below `min_gain`, so with this off the
    /// engine ships raw audio untouched. Costs a container-header sniff per
    /// file, and applies only to files that really are PCM.
    pub audio_codec: bool,
}

impl Default for CompressionConfig {
    fn default() -> Self {
        Self {
            mode: CompressionMode::Adaptive,
            algorithm: Algorithm::default(),
            level: 3,
            min_gain: 0.06,
            probe_bytes: 16 * 1024,
            incompressible_extensions: default_incompressible_extensions(),
            audio_codec: true,
        }
    }
}

/// Formats that are already entropy-coded. Compressing these burns CPU to make
/// the payload very slightly larger. FLAC and ALAC are lossless *codecs* but
/// still entropy-coded, so they belong here; WAV/AIFF/PCM do not — raw PCM
/// compresses well and is deliberately absent from this list.
pub fn default_incompressible_extensions() -> Vec<String> {
    [
        // audio
        "flac", "mp3", "aac", "m4a", "ogg", "opus", "wma", "ape", "alac", "dsf", "dff",
        // video
        "mp4", "mkv", "mov", "avi", "webm", "m4v", "mpg", "mpeg", "ts", "m2ts", "wmv", "flv",
        // images
        "jpg", "jpeg", "png", "gif", "webp", "heic", "heif", "avif", "jxl",
        // archives / already-compressed containers
        "zst", "gz", "bz2", "xz", "lz4", "7z", "zip", "rar", "br", "zipx", "cab",
        // packages & disk images that are internally compressed
        "whl", "jar", "apk", "crate", "deb", "rpm", "dmg", "appimage",
    ]
    .iter()
    .map(|s| s.to_string())
    .collect()
}

/// End-to-end payload confidentiality, layered *inside* whatever the transport
/// already provides.
#[derive(Clone)]
pub enum Secrecy {
    /// No extra layer. Payloads are protected only by the transport (QUIC/TLS
    /// 1.3). Correct when both endpoints terminate their own TLS and you trust
    /// every hop; wrong if traffic crosses a relay you do not control.
    TransportOnly,
    /// Ephemeral X25519 with a pre-shared key mixed into the KDF. Both sides
    /// must hold the same 32-byte secret; it authenticates the exchange and
    /// gives forward secrecy for recorded traffic.
    Psk([u8; 32]),
    /// Static X25519 identity with the peer's public key pinned, plus an
    /// ephemeral share. Mutual authentication without a PSK.
    Static {
        our_secret: [u8; 32],
        peer_public: [u8; 32],
    },
}

impl std::fmt::Debug for Secrecy {
    // Never let key material reach a log line.
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        match self {
            Secrecy::TransportOnly => f.write_str("TransportOnly"),
            Secrecy::Psk(_) => f.write_str("Psk(<redacted>)"),
            Secrecy::Static { .. } => f.write_str("Static { <redacted> }"),
        }
    }
}

/// AEAD used for the end-to-end layer.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Cipher {
    /// Pick AES-GCM where the CPU has AES instructions, ChaCha20 otherwise.
    Auto,
    Aes256Gcm,
    ChaCha20Poly1305,
}

#[derive(Debug, Clone)]
pub struct Config {
    /// Payload bytes per chunk before compression. Chunks are the unit of
    /// parallelism, resume, and AEAD sealing.
    pub chunk_size: usize,
    /// Concurrent data streams. Each is an independent QUIC unidirectional
    /// stream, so head-of-line blocking is per stream, not per transfer.
    pub streams: usize,
    /// CPU workers for compress/encrypt (sender) and decrypt/decompress
    /// (receiver). Defaults to the core count.
    pub workers: usize,
    /// Chunks allowed in flight per stream between the CPU stage and the wire.
    /// This is what bounds memory: peak ≈ streams * queue_depth * chunk_size.
    pub queue_depth: usize,
    pub compression: CompressionConfig,
    pub secrecy: Secrecy,
    pub cipher: Cipher,
    /// Verify each file's BLAKE3 hash on the receiver after the last chunk lands.
    pub verify_hashes: bool,
    /// Write a sidecar state file so an interrupted transfer resumes instead of
    /// restarting.
    pub resume: bool,
    /// Ask the filesystem to reserve space up front. Avoids fragmentation and
    /// surfaces ENOSPC before the first byte crosses the network.
    pub preallocate: bool,
    /// Reuse blocks the receiver already has from an older copy of a file.
    ///
    /// The receiver hashes whatever is already at the destination and tells the
    /// sender; the sender recognises matching chunks and sends 28 bytes instead
    /// of a megabyte. Changing one byte of a large file then costs one chunk,
    /// not the whole file. Costs one read of the existing copy on the receiver.
    pub delta: bool,
    /// Remember chunk hashes between runs, keyed by size and modification time.
    ///
    /// Without it, delta sync re-reads and re-hashes every file on both ends
    /// every time. With it, an unchanged file costs a `stat`. The trade is that
    /// a file edited within the timestamp's resolution *and* left the same
    /// length would go unnoticed; the receiver's hash check still catches it as
    /// a failed transfer rather than a corrupt file.
    pub trust_mtime: bool,
    /// Detect all-zero chunks and send them as a flag instead of as data.
    /// Turns a sparse or preallocated file — VM images, database files,
    /// preallocated media containers — into a transfer proportional to the data
    /// it actually holds, and reproduces the holes on the far side. The scan
    /// costs one pass over memory the sender has already read.
    pub sparse: bool,
    /// Ceiling on the bytes of chunk hashes offered for one file's reuse index.
    /// Past this the index costs more to announce than it can save.
    pub chunk_hash_budget: usize,
    /// Cap on a single frame's payload. Rejects hostile length prefixes.
    pub max_frame_bytes: usize,
    /// Ceiling on entries in one manifest.
    pub max_manifest_entries: usize,
    pub handshake_timeout: Duration,
    /// Preserve mtime and unix permission bits on received files.
    pub preserve_metadata: bool,
}

impl Default for Config {
    fn default() -> Self {
        let workers = num_cpus::get().max(1);
        Self {
            chunk_size: 1024 * 1024,
            // More streams than cores keeps the wire busy while workers are
            // mid-chunk, without the scheduling cost of a stream per chunk.
            streams: (workers * 2).clamp(4, 32),
            workers,
            queue_depth: 4,
            compression: CompressionConfig::default(),
            secrecy: Secrecy::TransportOnly,
            cipher: Cipher::Auto,
            verify_hashes: true,
            resume: true,
            preallocate: true,
            delta: true,
            trust_mtime: true,
            sparse: true,
            chunk_hash_budget: 8 * 1024 * 1024,
            max_frame_bytes: 64 * 1024 * 1024,
            max_manifest_entries: 4_000_000,
            handshake_timeout: Duration::from_secs(30),
            preserve_metadata: true,
        }
    }
}

impl Config {
    /// Saturate a fast link with large files: bigger chunks, cheap compression.
    pub fn throughput() -> Self {
        Self {
            chunk_size: 4 * 1024 * 1024,
            queue_depth: 6,
            compression: CompressionConfig {
                algorithm: Algorithm::Lz4,
                ..CompressionConfig::default()
            },
            ..Self::default()
        }
    }

    /// Minimise bytes on the wire for a slow or metered link.
    pub fn bandwidth_saving() -> Self {
        Self {
            compression: CompressionConfig {
                mode: CompressionMode::Always,
                algorithm: Algorithm::Zstd,
                level: 9,
                ..CompressionConfig::default()
            },
            ..Self::default()
        }
    }

    pub fn with_chunk_size(mut self, n: usize) -> Self {
        self.chunk_size = n;
        self
    }
    pub fn with_streams(mut self, n: usize) -> Self {
        self.streams = n;
        self
    }
    pub fn with_workers(mut self, n: usize) -> Self {
        self.workers = n;
        self
    }
    pub fn with_secrecy(mut self, s: Secrecy) -> Self {
        self.secrecy = s;
        self
    }
    pub fn with_compression(mut self, c: CompressionConfig) -> Self {
        self.compression = c;
        self
    }
    pub fn without_compression(mut self) -> Self {
        self.compression.mode = CompressionMode::Off;
        self
    }

    /// Worst-case resident bytes for the chunk pipeline, excluding OS cache.
    pub fn memory_budget(&self) -> usize {
        // Each in-flight slot holds a plaintext chunk and its encoded form.
        self.streams * self.queue_depth * self.chunk_size * 2
    }

    pub(crate) fn validate(&self) -> Result<()> {
        if self.chunk_size < 4096 {
            return Err(Error::Config("chunk_size must be at least 4 KiB".into()));
        }
        if self.chunk_size > self.max_frame_bytes / 2 {
            return Err(Error::Config(
                "chunk_size must leave headroom under max_frame_bytes for incompressible expansion"
                    .into(),
            ));
        }
        if self.streams == 0 {
            return Err(Error::Config("streams must be >= 1".into()));
        }
        if self.workers == 0 {
            return Err(Error::Config("workers must be >= 1".into()));
        }
        if self.queue_depth == 0 {
            return Err(Error::Config("queue_depth must be >= 1".into()));
        }
        if !(0.0..1.0).contains(&self.compression.min_gain) {
            return Err(Error::Config("min_gain must be in [0, 1)".into()));
        }
        Ok(())
    }
}