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
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
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
//! Per-chunk compression with an adaptive "is this worth it?" gate.
//!
//! Every chunk is compressed independently. That costs a little ratio versus a
//! single solid stream, and buys three things this engine needs: chunks can be
//! processed on any worker in any order, a resumed transfer can skip individual
//! chunks, and one corrupt chunk cannot poison the ones after it.

use crate::config::{CompressionConfig, CompressionMode};
use crate::error::{Error, Result};

/// Compression algorithm for a chunk payload.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
pub enum Algorithm {
    None = 0,
    /// Very high throughput (GB/s per core), modest ratio. Use when the CPU,
    /// not the link, is the constraint.
    Lz4 = 1,
    /// Better ratio at a few hundred MB/s per core. The default.
    Zstd = 2,
    /// Lossless predictive coder for uncompressed PCM audio. Roughly doubles
    /// what zstd achieves on `.wav`, which zstd barely compresses at all.
    Pcm = 3,
}

impl Default for Algorithm {
    fn default() -> Self {
        if cfg!(feature = "zstd-codec") {
            Algorithm::Zstd
        } else if cfg!(feature = "lz4-codec") {
            Algorithm::Lz4
        } else {
            Algorithm::None
        }
    }
}

impl Algorithm {
    pub fn from_u8(v: u8) -> Result<Self> {
        match v {
            0 => Ok(Algorithm::None),
            1 => Ok(Algorithm::Lz4),
            2 => Ok(Algorithm::Zstd),
            3 => Ok(Algorithm::Pcm),
            other => Err(Error::Compress(format!("unknown algorithm id {other}"))),
        }
    }

    pub fn available(self) -> bool {
        match self {
            Algorithm::None => true,
            Algorithm::Lz4 => cfg!(feature = "lz4-codec"),
            Algorithm::Zstd => cfg!(feature = "zstd-codec"),
            // No third-party dependency, so always available.
            Algorithm::Pcm => true,
        }
    }
}

/// What the encoder decided for one chunk.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Encoded {
    pub algorithm: Algorithm,
    /// Size before compression; the decoder needs it to size its output buffer.
    pub raw_len: usize,
}

/// Reusable codec state for one worker.
///
/// Two things live here that must not be rebuilt per chunk:
///
/// * **zstd contexts.** A `CCtx` owns the match tables — megabytes at level 3.
///   The convenience `zstd::bulk::compress_to_buffer` allocates and frees one on
///   every call, which on 1 MiB chunks costs more than the compression itself.
/// * **A scratch buffer.** zstd needs a destination sized to its worst case.
///   Growing a `Vec` for that means zeroing ~1 MiB per chunk that the codec is
///   about to overwrite anyway. One long-lived buffer pays that once.
pub struct Codec {
    /// Only consulted by zstd; without that feature there is no level to track.
    #[cfg_attr(not(feature = "zstd-codec"), allow(dead_code))]
    level: i32,
    scratch: Vec<u8>,
    #[cfg(feature = "zstd-codec")]
    zc: Option<zstd::bulk::Compressor<'static>>,
    #[cfg(feature = "zstd-codec")]
    zd: Option<zstd::bulk::Decompressor<'static>>,
}

impl Default for Codec {
    fn default() -> Self {
        Self::new()
    }
}

impl Codec {
    pub fn new() -> Self {
        Self {
            level: i32::MIN,
            scratch: Vec::new(),
            #[cfg(feature = "zstd-codec")]
            zc: None,
            #[cfg(feature = "zstd-codec")]
            zd: None,
        }
    }

    /// Compress `input`, appending the result to `out`.
    ///
    /// Returns the algorithm actually used, which may be `None` even when a
    /// codec was requested: if the result did not beat `min_gain`, the raw bytes
    /// are written instead. The caller must record the returned algorithm on the
    /// wire.
    pub fn compress_into(
        &mut self,
        cfg: &CompressionConfig,
        hint: FileHint,
        input: &[u8],
        out: &mut Vec<u8>,
    ) -> Result<Encoded> {
        let raw_len = input.len();
        let algo = select_algorithm(cfg, hint, input);

        if algo == Algorithm::None {
            out.extend_from_slice(input);
            return Ok(Encoded {
                algorithm: Algorithm::None,
                raw_len,
            });
        }

        let produced = self.run(algo, input, cfg.level)?;

        // The gate is applied to the real output, not the probe's guess. An
        // expensive miss costs CPU but never costs bytes on the wire.
        if !worth_it(produced, raw_len, cfg.min_gain) {
            out.extend_from_slice(input);
            return Ok(Encoded {
                algorithm: Algorithm::None,
                raw_len,
            });
        }

        out.extend_from_slice(&self.scratch[..produced]);
        Ok(Encoded {
            algorithm: algo,
            raw_len,
        })
    }

    /// Compress the payload already sitting in `buf` after `prefix` header bytes.
    ///
    /// The sender reads each chunk straight into its frame buffer, so the common
    /// case — a chunk that will not be compressed, which is every chunk of a
    /// `.flac` or `.mp4` — costs zero copies. Only a chunk that actually wins
    /// pays one, and then only of its compressed size.
    pub fn compress_in_place(
        &mut self,
        cfg: &CompressionConfig,
        hint: FileHint,
        buf: &mut Vec<u8>,
        prefix: usize,
    ) -> Result<Encoded> {
        let raw_len = buf.len() - prefix;

        // Uncompressed audio goes to the coder that models it. zstd manages
        // about 1.05x here, below the gate, so without this branch a `.wav`
        // ships raw.
        // `CompressionMode::Off` means off, including this coder — it runs
        // ahead of `select_algorithm`, so it has to honour the mode itself.
        if let Some(fmt) = hint
            .audio
            .filter(|_| cfg.audio_codec && cfg.mode != CompressionMode::Off)
        {
            self.scratch.clear();
            if let Some(n) =
                super::pcm::encode(&fmt, hint.chunk_offset, &buf[prefix..], &mut self.scratch)
            {
                if worth_it(n, raw_len, cfg.min_gain) {
                    buf.truncate(prefix);
                    buf.extend_from_slice(&self.scratch[..n]);
                    return Ok(Encoded {
                        algorithm: Algorithm::Pcm,
                        raw_len,
                    });
                }
            }
        }

        let algo = select_algorithm(cfg, hint, &buf[prefix..]);
        if algo == Algorithm::None {
            return Ok(Encoded {
                algorithm: Algorithm::None,
                raw_len,
            });
        }

        let produced = self.run_from(algo, buf, prefix, cfg.level)?;
        if !worth_it(produced, raw_len, cfg.min_gain) {
            return Ok(Encoded {
                algorithm: Algorithm::None,
                raw_len,
            });
        }

        buf.truncate(prefix);
        buf.extend_from_slice(&self.scratch[..produced]);
        Ok(Encoded {
            algorithm: algo,
            raw_len,
        })
    }

    /// Decompress one chunk. `raw_len` comes from the frame header and is
    /// validated against the decoded length, so a lying sender cannot make us
    /// over-allocate beyond the frame limit the caller already enforced.
    pub fn decompress_into(
        &mut self,
        algo: Algorithm,
        raw_len: usize,
        input: &[u8],
        out: &mut Vec<u8>,
    ) -> Result<()> {
        match algo {
            Algorithm::None => {
                if input.len() != raw_len {
                    return Err(Error::Compress(format!(
                        "raw chunk length {} does not match declared {}",
                        input.len(),
                        raw_len
                    )));
                }
                out.extend_from_slice(input);
                Ok(())
            }
            Algorithm::Zstd => self.decompress_zstd(input, raw_len, out),
            Algorithm::Lz4 => self.decompress_lz4(input, raw_len, out),
            Algorithm::Pcm => {
                let before = out.len();
                super::pcm::decode(input, out)?;
                if out.len() - before != raw_len {
                    out.truncate(before);
                    return Err(Error::Compress(format!(
                        "pcm produced {} bytes, header declared {raw_len}",
                        out.len() - before
                    )));
                }
                Ok(())
            }
        }
    }

    /// Grow the scratch buffer to `n`, keeping it initialised across calls so
    /// no chunk ever pays to zero it.
    #[cfg_attr(not(feature = "lz4-codec"), allow(dead_code))]
    fn scratch_at_least(&mut self, n: usize) -> &mut [u8] {
        if self.scratch.len() < n {
            self.scratch.resize(n, 0);
        }
        &mut self.scratch[..n]
    }

    fn run(&mut self, algo: Algorithm, input: &[u8], level: i32) -> Result<usize> {
        match algo {
            Algorithm::Zstd => self.compress_zstd(input, level),
            Algorithm::Lz4 => self.compress_lz4(input),
            // Handled ahead of general-purpose selection, since it needs the
            // chunk's file offset rather than just its bytes.
            Algorithm::Pcm | Algorithm::None => unreachable!("handled by the caller"),
        }
    }

    /// Same as [`Codec::run`], for input borrowed out of `buf`.
    fn run_from(
        &mut self,
        algo: Algorithm,
        buf: &[u8],
        prefix: usize,
        level: i32,
    ) -> Result<usize> {
        self.run(algo, &buf[prefix..], level)
    }
}

/// Did compression save enough to be worth sending compressed?
#[inline]
fn worth_it(produced: usize, raw_len: usize, min_gain: f32) -> bool {
    let gain = 1.0 - (produced as f32 / raw_len.max(1) as f32);
    gain >= min_gain
}

thread_local! {
    /// One codec per thread. The engine runs its CPU work on a fixed rayon
    /// pool, so a thread-local *is* a per-worker context: no pool bookkeeping,
    /// no lock, and the zstd tables stay hot from one chunk to the next.
    static TLS_CODEC: std::cell::RefCell<Codec> = std::cell::RefCell::new(Codec::new());
}

/// Compress using the calling thread's [`Codec`].
pub fn compress_into(
    cfg: &CompressionConfig,
    hint: FileHint,
    input: &[u8],
    out: &mut Vec<u8>,
) -> Result<Encoded> {
    TLS_CODEC.with(|c| c.borrow_mut().compress_into(cfg, hint, input, out))
}

/// Decompress using the calling thread's [`Codec`].
pub fn decompress_into(
    algo: Algorithm,
    raw_len: usize,
    input: &[u8],
    out: &mut Vec<u8>,
) -> Result<()> {
    TLS_CODEC.with(|c| c.borrow_mut().decompress_into(algo, raw_len, input, out))
}

/// Run `f` with the calling thread's codec.
pub fn with_codec<R>(f: impl FnOnce(&mut Codec) -> R) -> R {
    TLS_CODEC.with(|c| f(&mut c.borrow_mut()))
}

// ---------------------------------------------------------------------------
// Selection
// ---------------------------------------------------------------------------

/// What we know about a chunk before looking at its bytes.
#[derive(Debug, Clone, Copy, Default)]
pub struct FileHint {
    /// Extension matched the incompressible list, so skip the probe entirely.
    pub known_incompressible: bool,
    /// The file is uncompressed PCM with this layout, so the audio coder
    /// applies. Detected once per file from its container header.
    pub audio: Option<super::pcm::AudioFormat>,
    /// Where this chunk starts in the file. The audio coder needs it to find
    /// frame boundaries, since chunks do not respect them.
    pub chunk_offset: u64,
}

fn select_algorithm(cfg: &CompressionConfig, hint: FileHint, input: &[u8]) -> Algorithm {
    match cfg.mode {
        CompressionMode::Off => return Algorithm::None,
        CompressionMode::Always => {
            return if cfg.algorithm.available() {
                cfg.algorithm
            } else {
                Algorithm::None
            }
        }
        CompressionMode::Adaptive => {}
    }

    if !cfg.algorithm.available() || hint.known_incompressible {
        return Algorithm::None;
    }
    // Tiny chunks are all framing overhead; the codec cannot win.
    if input.len() < 1024 {
        return Algorithm::None;
    }
    if looks_incompressible(input, cfg.probe_bytes) {
        return Algorithm::None;
    }
    cfg.algorithm
}

/// Cheap entropy probe: a byte histogram over a sample, scored by Shannon
/// entropy. Encrypted and entropy-coded data sits at ~8.0 bits/byte; text and
/// PCM audio sit well below. Costs ~O(sample) with no allocation, versus a
/// trial compression that would cost a full codec pass.
///
/// This is a filter, not an oracle — anything it lets through still has to
/// clear the real `min_gain` check after compressing.
fn looks_incompressible(input: &[u8], probe_bytes: usize) -> bool {
    let n = probe_bytes.min(input.len());
    if n < 256 {
        return false;
    }
    // Sample the head; for a chunk this is representative and stays in L1/L2.
    let sample = &input[..n];

    let mut hist = [0u32; 256];
    for &b in sample {
        hist[b as usize] += 1;
    }

    let len = n as f32;
    let mut entropy = 0.0f32;
    for &c in hist.iter() {
        if c != 0 {
            let p = c as f32 / len;
            entropy -= p * p.log2();
        }
    }

    // 7.8 bits/byte leaves headroom for high-entropy-but-compressible inputs
    // (e.g. base64 of random data is ~6.0, dense binaries ~7.2).
    entropy > 7.8
}

/// Does this extension name a format that is already entropy-coded?
pub fn is_incompressible_extension(cfg: &CompressionConfig, path: &str) -> bool {
    let ext = match path.rsplit_once('.') {
        Some((_, e)) if !e.is_empty() && e.len() <= 12 => e,
        _ => return false,
    };
    let lower = ext.to_ascii_lowercase();
    cfg.incompressible_extensions.contains(&lower)
}

// ---------------------------------------------------------------------------
// Backends
// ---------------------------------------------------------------------------

impl Codec {
    #[cfg(feature = "zstd-codec")]
    fn compress_zstd(&mut self, input: &[u8], level: i32) -> Result<usize> {
        if self.zc.is_none() {
            self.zc = Some(
                zstd::bulk::Compressor::new(level)
                    .map_err(|e| Error::Compress(format!("zstd context: {e}")))?,
            );
            self.level = level;
        }
        if self.level != level {
            self.zc
                .as_mut()
                .expect("just built")
                .set_compression_level(level)
                .map_err(|e| Error::Compress(format!("zstd level: {e}")))?;
            self.level = level;
        }
        let bound = zstd::zstd_safe::compress_bound(input.len());
        if self.scratch.len() < bound {
            self.scratch.resize(bound, 0);
        }
        let (zc, scratch) = (
            self.zc.as_mut().expect("just built"),
            &mut self.scratch[..bound],
        );
        zc.compress_to_buffer(input, scratch)
            .map_err(|e| Error::Compress(format!("zstd: {e}")))
    }

    #[cfg(not(feature = "zstd-codec"))]
    fn compress_zstd(&mut self, _input: &[u8], _level: i32) -> Result<usize> {
        Err(Error::Compress("zstd support not compiled in".into()))
    }

    #[cfg(feature = "zstd-codec")]
    fn decompress_zstd(&mut self, input: &[u8], raw_len: usize, out: &mut Vec<u8>) -> Result<()> {
        if self.zd.is_none() {
            self.zd = Some(
                zstd::bulk::Decompressor::new()
                    .map_err(|e| Error::Compress(format!("zstd context: {e}")))?,
            );
        }
        let before = out.len();
        out.resize(before + raw_len, 0);
        let written = self
            .zd
            .as_mut()
            .expect("just built")
            .decompress_to_buffer(input, &mut out[before..])
            .map_err(|e| Error::Compress(format!("zstd decode: {e}")))?;
        if written != raw_len {
            out.truncate(before);
            return Err(Error::Compress(format!(
                "zstd produced {written} bytes, header declared {raw_len}"
            )));
        }
        Ok(())
    }

    #[cfg(not(feature = "zstd-codec"))]
    fn decompress_zstd(
        &mut self,
        _input: &[u8],
        _raw_len: usize,
        _out: &mut Vec<u8>,
    ) -> Result<()> {
        Err(Error::Compress(
            "peer used zstd but zstd support is not compiled in".into(),
        ))
    }

    #[cfg(feature = "lz4-codec")]
    fn compress_lz4(&mut self, input: &[u8]) -> Result<usize> {
        let bound = lz4_flex::block::get_maximum_output_size(input.len());
        let dst = self.scratch_at_least(bound);
        lz4_flex::block::compress_into(input, dst).map_err(|e| Error::Compress(format!("lz4: {e}")))
    }

    #[cfg(not(feature = "lz4-codec"))]
    fn compress_lz4(&mut self, _input: &[u8]) -> Result<usize> {
        Err(Error::Compress("lz4 support not compiled in".into()))
    }

    #[cfg(feature = "lz4-codec")]
    fn decompress_lz4(&mut self, input: &[u8], raw_len: usize, out: &mut Vec<u8>) -> Result<()> {
        let before = out.len();
        out.resize(before + raw_len, 0);
        let written = lz4_flex::block::decompress_into(input, &mut out[before..])
            .map_err(|e| Error::Compress(format!("lz4 decode: {e}")))?;
        if written != raw_len {
            out.truncate(before);
            return Err(Error::Compress(format!(
                "lz4 produced {written} bytes, header declared {raw_len}"
            )));
        }
        Ok(())
    }

    #[cfg(not(feature = "lz4-codec"))]
    fn decompress_lz4(&mut self, _input: &[u8], _raw_len: usize, _out: &mut Vec<u8>) -> Result<()> {
        Err(Error::Compress(
            "peer used lz4 but lz4 support is not compiled in".into(),
        ))
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    fn text_chunk() -> Vec<u8> {
        "the quick brown fox jumps over the lazy dog. "
            .repeat(4000)
            .into_bytes()
    }

    fn random_chunk(n: usize) -> Vec<u8> {
        // xorshift; high entropy, no dependency on a test RNG crate here.
        let mut s = 0x2545F4914F6CDD1Du64;
        (0..n)
            .map(|_| {
                s ^= s << 13;
                s ^= s >> 7;
                s ^= s << 17;
                (s >> 24) as u8
            })
            .collect()
    }

    #[test]
    fn roundtrip_all_algorithms() {
        let data = text_chunk();
        for algo in [Algorithm::None, Algorithm::Lz4, Algorithm::Zstd] {
            if !algo.available() {
                continue;
            }
            let cfg = CompressionConfig {
                mode: CompressionMode::Always,
                algorithm: algo,
                ..Default::default()
            };
            let mut enc = Vec::new();
            let e = compress_into(&cfg, FileHint::default(), &data, &mut enc).unwrap();
            let mut dec = Vec::new();
            decompress_into(e.algorithm, e.raw_len, &enc, &mut dec).unwrap();
            assert_eq!(dec, data, "roundtrip failed for {algo:?}");
        }
    }

    #[test]
    fn adaptive_skips_high_entropy_data() {
        let cfg = CompressionConfig::default();
        let data = random_chunk(256 * 1024);
        let mut enc = Vec::new();
        let e = compress_into(&cfg, FileHint::default(), &data, &mut enc).unwrap();
        assert_eq!(e.algorithm, Algorithm::None);
        // Incompressible input must never be inflated by the transfer.
        assert_eq!(enc.len(), data.len());
    }

    #[test]
    fn adaptive_compresses_text() {
        let cfg = CompressionConfig::default();
        let data = text_chunk();
        let mut enc = Vec::new();
        let e = compress_into(&cfg, FileHint::default(), &data, &mut enc).unwrap();
        assert_ne!(e.algorithm, Algorithm::None);
        assert!(enc.len() < data.len() / 2);
    }

    #[test]
    fn extension_hint_forces_raw() {
        let cfg = CompressionConfig::default();
        let data = text_chunk();
        let hint = FileHint {
            known_incompressible: true,
            ..Default::default()
        };
        let mut enc = Vec::new();
        let e = compress_into(&cfg, hint, &data, &mut enc).unwrap();
        assert_eq!(e.algorithm, Algorithm::None);
    }

    #[test]
    fn extension_matching() {
        let cfg = CompressionConfig::default();
        assert!(is_incompressible_extension(&cfg, "song.FLAC"));
        assert!(is_incompressible_extension(&cfg, "a/b/movie.mkv"));
        // Raw PCM must stay compressible; it is the main win for audio sets.
        assert!(!is_incompressible_extension(&cfg, "master.wav"));
        assert!(!is_incompressible_extension(&cfg, "notes.txt"));
        assert!(!is_incompressible_extension(&cfg, "no_extension"));
    }

    #[test]
    fn decompress_rejects_length_mismatch() {
        let cfg = CompressionConfig {
            mode: CompressionMode::Always,
            ..Default::default()
        };
        let data = text_chunk();
        let mut enc = Vec::new();
        let e = compress_into(&cfg, FileHint::default(), &data, &mut enc).unwrap();
        if e.algorithm == Algorithm::None {
            return;
        }
        let mut dec = Vec::new();
        // A sender lying about raw_len must be rejected, not trusted.
        assert!(decompress_into(e.algorithm, e.raw_len / 2, &enc, &mut dec).is_err());
    }
}