Skip to main content

runsync_transfer/
config.rs

1use crate::codec::compress::Algorithm;
2use crate::error::{Error, Result};
3use std::time::Duration;
4
5/// How the engine decides whether to compress a given chunk.
6#[derive(Debug, Clone, Copy, PartialEq, Eq)]
7pub enum CompressionMode {
8    /// Never compress. Lowest CPU; best when the link is faster than the CPU
9    /// (loopback, 25GbE) or every input is already compressed media.
10    Off,
11    /// Always compress with the configured algorithm.
12    Always,
13    /// Decide per chunk from a cheap entropy probe, and per file from its
14    /// extension. This is the default and the right answer for mixed payloads.
15    Adaptive,
16}
17
18/// Tuning for the compression stage.
19#[derive(Debug, Clone)]
20pub struct CompressionConfig {
21    pub mode: CompressionMode,
22    pub algorithm: Algorithm,
23    /// zstd level. Ignored by LZ4. 1..=9 is the useful range for transfers;
24    /// above ~6 the compressor becomes the bottleneck before the network does.
25    pub level: i32,
26    /// In `Adaptive` mode a chunk is sent raw unless compression saves at least
27    /// this fraction. 0.06 means "must shrink by 6% to be worth it".
28    pub min_gain: f32,
29    /// Bytes sampled from the head of a chunk for the entropy probe.
30    pub probe_bytes: usize,
31    /// Extensions (lowercase, no dot) that skip compression entirely.
32    pub incompressible_extensions: Vec<String>,
33    /// Use the dedicated lossless coder for uncompressed PCM audio.
34    ///
35    /// zstd manages ~1.05x on `.wav`, below `min_gain`, so with this off the
36    /// engine ships raw audio untouched. Costs a container-header sniff per
37    /// file, and applies only to files that really are PCM.
38    pub audio_codec: bool,
39}
40
41impl Default for CompressionConfig {
42    fn default() -> Self {
43        Self {
44            mode: CompressionMode::Adaptive,
45            algorithm: Algorithm::default(),
46            level: 3,
47            min_gain: 0.06,
48            probe_bytes: 16 * 1024,
49            incompressible_extensions: default_incompressible_extensions(),
50            audio_codec: true,
51        }
52    }
53}
54
55/// Formats that are already entropy-coded. Compressing these burns CPU to make
56/// the payload very slightly larger. FLAC and ALAC are lossless *codecs* but
57/// still entropy-coded, so they belong here; WAV/AIFF/PCM do not — raw PCM
58/// compresses well and is deliberately absent from this list.
59pub fn default_incompressible_extensions() -> Vec<String> {
60    [
61        // audio
62        "flac", "mp3", "aac", "m4a", "ogg", "opus", "wma", "ape", "alac", "dsf", "dff",
63        // video
64        "mp4", "mkv", "mov", "avi", "webm", "m4v", "mpg", "mpeg", "ts", "m2ts", "wmv", "flv",
65        // images
66        "jpg", "jpeg", "png", "gif", "webp", "heic", "heif", "avif", "jxl",
67        // archives / already-compressed containers
68        "zst", "gz", "bz2", "xz", "lz4", "7z", "zip", "rar", "br", "zipx", "cab",
69        // packages & disk images that are internally compressed
70        "whl", "jar", "apk", "crate", "deb", "rpm", "dmg", "appimage",
71    ]
72    .iter()
73    .map(|s| s.to_string())
74    .collect()
75}
76
77/// End-to-end payload confidentiality, layered *inside* whatever the transport
78/// already provides.
79#[derive(Clone)]
80pub enum Secrecy {
81    /// No extra layer. Payloads are protected only by the transport (QUIC/TLS
82    /// 1.3). Correct when both endpoints terminate their own TLS and you trust
83    /// every hop; wrong if traffic crosses a relay you do not control.
84    TransportOnly,
85    /// Ephemeral X25519 with a pre-shared key mixed into the KDF. Both sides
86    /// must hold the same 32-byte secret; it authenticates the exchange and
87    /// gives forward secrecy for recorded traffic.
88    Psk([u8; 32]),
89    /// Static X25519 identity with the peer's public key pinned, plus an
90    /// ephemeral share. Mutual authentication without a PSK.
91    Static {
92        our_secret: [u8; 32],
93        peer_public: [u8; 32],
94    },
95}
96
97impl std::fmt::Debug for Secrecy {
98    // Never let key material reach a log line.
99    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
100        match self {
101            Secrecy::TransportOnly => f.write_str("TransportOnly"),
102            Secrecy::Psk(_) => f.write_str("Psk(<redacted>)"),
103            Secrecy::Static { .. } => f.write_str("Static { <redacted> }"),
104        }
105    }
106}
107
108/// AEAD used for the end-to-end layer.
109#[derive(Debug, Clone, Copy, PartialEq, Eq)]
110pub enum Cipher {
111    /// Pick AES-GCM where the CPU has AES instructions, ChaCha20 otherwise.
112    Auto,
113    Aes256Gcm,
114    ChaCha20Poly1305,
115}
116
117#[derive(Debug, Clone)]
118pub struct Config {
119    /// Payload bytes per chunk before compression. Chunks are the unit of
120    /// parallelism, resume, and AEAD sealing.
121    pub chunk_size: usize,
122    /// Concurrent data streams. Each is an independent QUIC unidirectional
123    /// stream, so head-of-line blocking is per stream, not per transfer.
124    pub streams: usize,
125    /// CPU workers for compress/encrypt (sender) and decrypt/decompress
126    /// (receiver). Defaults to the core count.
127    pub workers: usize,
128    /// Chunks allowed in flight per stream between the CPU stage and the wire.
129    /// This is what bounds memory: peak ≈ streams * queue_depth * chunk_size.
130    pub queue_depth: usize,
131    pub compression: CompressionConfig,
132    pub secrecy: Secrecy,
133    pub cipher: Cipher,
134    /// Verify each file's BLAKE3 hash on the receiver after the last chunk lands.
135    pub verify_hashes: bool,
136    /// Write a sidecar state file so an interrupted transfer resumes instead of
137    /// restarting.
138    pub resume: bool,
139    /// Ask the filesystem to reserve space up front. Avoids fragmentation and
140    /// surfaces ENOSPC before the first byte crosses the network.
141    pub preallocate: bool,
142    /// Reuse blocks the receiver already has from an older copy of a file.
143    ///
144    /// The receiver hashes whatever is already at the destination and tells the
145    /// sender; the sender recognises matching chunks and sends 28 bytes instead
146    /// of a megabyte. Changing one byte of a large file then costs one chunk,
147    /// not the whole file. Costs one read of the existing copy on the receiver.
148    pub delta: bool,
149    /// Remember chunk hashes between runs, keyed by size and modification time.
150    ///
151    /// Without it, delta sync re-reads and re-hashes every file on both ends
152    /// every time. With it, an unchanged file costs a `stat`. The trade is that
153    /// a file edited within the timestamp's resolution *and* left the same
154    /// length would go unnoticed; the receiver's hash check still catches it as
155    /// a failed transfer rather than a corrupt file.
156    pub trust_mtime: bool,
157    /// Detect all-zero chunks and send them as a flag instead of as data.
158    /// Turns a sparse or preallocated file — VM images, database files,
159    /// preallocated media containers — into a transfer proportional to the data
160    /// it actually holds, and reproduces the holes on the far side. The scan
161    /// costs one pass over memory the sender has already read.
162    pub sparse: bool,
163    /// Ceiling on the bytes of chunk hashes offered for one file's reuse index.
164    /// Past this the index costs more to announce than it can save.
165    pub chunk_hash_budget: usize,
166    /// Cap on a single frame's payload. Rejects hostile length prefixes.
167    pub max_frame_bytes: usize,
168    /// Ceiling on entries in one manifest.
169    pub max_manifest_entries: usize,
170    pub handshake_timeout: Duration,
171    /// Preserve mtime and unix permission bits on received files.
172    pub preserve_metadata: bool,
173}
174
175impl Default for Config {
176    fn default() -> Self {
177        let workers = num_cpus::get().max(1);
178        Self {
179            chunk_size: 1024 * 1024,
180            // More streams than cores keeps the wire busy while workers are
181            // mid-chunk, without the scheduling cost of a stream per chunk.
182            streams: (workers * 2).clamp(4, 32),
183            workers,
184            queue_depth: 4,
185            compression: CompressionConfig::default(),
186            secrecy: Secrecy::TransportOnly,
187            cipher: Cipher::Auto,
188            verify_hashes: true,
189            resume: true,
190            preallocate: true,
191            delta: true,
192            trust_mtime: true,
193            sparse: true,
194            chunk_hash_budget: 8 * 1024 * 1024,
195            max_frame_bytes: 64 * 1024 * 1024,
196            max_manifest_entries: 4_000_000,
197            handshake_timeout: Duration::from_secs(30),
198            preserve_metadata: true,
199        }
200    }
201}
202
203impl Config {
204    /// Saturate a fast link with large files: bigger chunks, cheap compression.
205    pub fn throughput() -> Self {
206        Self {
207            chunk_size: 4 * 1024 * 1024,
208            queue_depth: 6,
209            compression: CompressionConfig {
210                algorithm: Algorithm::Lz4,
211                ..CompressionConfig::default()
212            },
213            ..Self::default()
214        }
215    }
216
217    /// Minimise bytes on the wire for a slow or metered link.
218    pub fn bandwidth_saving() -> Self {
219        Self {
220            compression: CompressionConfig {
221                mode: CompressionMode::Always,
222                algorithm: Algorithm::Zstd,
223                level: 9,
224                ..CompressionConfig::default()
225            },
226            ..Self::default()
227        }
228    }
229
230    pub fn with_chunk_size(mut self, n: usize) -> Self {
231        self.chunk_size = n;
232        self
233    }
234    pub fn with_streams(mut self, n: usize) -> Self {
235        self.streams = n;
236        self
237    }
238    pub fn with_workers(mut self, n: usize) -> Self {
239        self.workers = n;
240        self
241    }
242    pub fn with_secrecy(mut self, s: Secrecy) -> Self {
243        self.secrecy = s;
244        self
245    }
246    pub fn with_compression(mut self, c: CompressionConfig) -> Self {
247        self.compression = c;
248        self
249    }
250    pub fn without_compression(mut self) -> Self {
251        self.compression.mode = CompressionMode::Off;
252        self
253    }
254
255    /// Worst-case resident bytes for the chunk pipeline, excluding OS cache.
256    pub fn memory_budget(&self) -> usize {
257        // Each in-flight slot holds a plaintext chunk and its encoded form.
258        self.streams * self.queue_depth * self.chunk_size * 2
259    }
260
261    pub(crate) fn validate(&self) -> Result<()> {
262        if self.chunk_size < 4096 {
263            return Err(Error::Config("chunk_size must be at least 4 KiB".into()));
264        }
265        if self.chunk_size > self.max_frame_bytes / 2 {
266            return Err(Error::Config(
267                "chunk_size must leave headroom under max_frame_bytes for incompressible expansion"
268                    .into(),
269            ));
270        }
271        if self.streams == 0 {
272            return Err(Error::Config("streams must be >= 1".into()));
273        }
274        if self.workers == 0 {
275            return Err(Error::Config("workers must be >= 1".into()));
276        }
277        if self.queue_depth == 0 {
278            return Err(Error::Config("queue_depth must be >= 1".into()));
279        }
280        if !(0.0..1.0).contains(&self.compression.min_gain) {
281            return Err(Error::Config("min_gain must be in [0, 1)".into()));
282        }
283        Ok(())
284    }
285}