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}