1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
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(())
}
}