diff --git a/CHANGELOG.md b/CHANGELOG.md
index 5afd9cd..05193aa 100644
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -112,8 +112,9 @@ layout is unchanged (golden length still 886).
wording, and symlink/empty-directory quick-checks.
- `--delete-before`'s phase-0 late-file divergence remains (rsync's pre-scan
fixes the file list before the data pass).
-- A single file larger than 256 MiB cannot be streamed in the default path
- (a general whole-file limit, not basis-specific).
+- A whole-file sender that cannot stream its codec (lz4's one-shot block
+ format) or a `--append`/delta source above the bound still buffers; the
+ default zstd/zlib and the uncompressed paths stream (see #318).
- `--stats` byte totals and `--msgs2stderr` stay documented divergences.
### Security
diff --git a/HANDOFF.md b/HANDOFF.md
index e2931dc..4f0f02e 100644
--- a/HANDOFF.md
+++ b/HANDOFF.md
@@ -235,7 +235,9 @@
and symlink/empty-dir quick-check feedback.
- **`--delete-before` phase-0 keep-set** (rsync fixes the file list before
the data pass; FastSync keeps its pre-scan snapshot race).
- - **>256 MiB single-file streaming** (B4, the general whole-file limit).
+ - ~~**>256 MiB single-file streaming** (B4, the general whole-file limit).~~
+ Closed by #318: the whole-file payload, the basis read/verify and the fuzzy
+ basis are streamed through bounded buffers (lz4/append remain buffered).
- **Wire native-size framing:** lengths are native `size_t` and the protocol
assumes homogeneous word size/endianness — document or move to fixed-width
framing.
diff --git a/README.md b/README.md
index af2240f..7efdd70 100644
--- a/README.md
+++ b/README.md
@@ -206,9 +206,9 @@ This produces `./build/client` and `./build/server`. `compile_commands.json` is
| `--incremental` | Skip files unchanged since last transfer (size + mtime). Auto-enables `--preserve`. Incompatible with `--chunk-serialization`. |
| `--existing` | Skip files not already present at the destination; update existing files normally. |
| `--ignore-existing` | Skip files that already exist on the receiver; like rsync it does not apply to directories or symlinks. |
-| `--compare-dest
` | Extra comparison basis: unchanged files are not transferred (requires/implies `--incremental`; a basis MISS above the 256 MiB whole-file payload bound is refused — see [`RSYNC_COMPAT.md`](RSYNC_COMPAT.md)) |
-| `--copy-dest ` | Like `--compare-dest`, but copies the unchanged file from DIR into the destination (same basis-size caveat; see [`RSYNC_COMPAT.md`](RSYNC_COMPAT.md)) |
-| `--link-dest ` | Like `--copy-dest`, but hard-links the unchanged file from DIR (repeatable; earlier DIRs win; same basis-size caveat; see [`RSYNC_COMPAT.md`](RSYNC_COMPAT.md)) |
+| `--compare-dest ` | Extra comparison basis: unchanged files are not transferred (requires/implies `--incremental`; a basis of any size is supported, streamed in bounded chunks — see [`RSYNC_COMPAT.md`](RSYNC_COMPAT.md)) |
+| `--copy-dest ` | Like `--compare-dest`, but copies the unchanged file from DIR into the destination (the copy is streamed, so a basis of any size works; see [`RSYNC_COMPAT.md`](RSYNC_COMPAT.md)) |
+| `--link-dest ` | Like `--copy-dest`, but hard-links the unchanged file from DIR (repeatable; earlier DIRs win; a basis of any size works; see [`RSYNC_COMPAT.md`](RSYNC_COMPAT.md)) |
| `--verify-basis` | FastSync-only: require a basis hit (`--compare-dest`/`--copy-dest`/`--link-dest`) to match the source by whole-file digest instead of trusting the size+mtime quick-check (default matches rsync) |
| `--delete` | Delete files on receiver not present in source (default timing: delete-during, matching rsync, so destination space is freed progressively). Scoped to the synchronized directories, so `--files-from` subsets are safe |
| `--delete-before` | Delete extras before the transfer starts (implies `--delete`) |
@@ -300,6 +300,7 @@ transfer is never aborted.
| `FASTSYNC_SOURCE_DIR` | — | Source directory fallback |
| `FASTSYNC_DEST_DIR` | — | Destination directory fallback |
| `FASTSYNC_SAVE_TO_DISK` | `false` | Disk persistence fallback |
+| `FASTSYNC_MAX_WHOLE_FILE_SIZE` | `268435456` | Receiver-only test hook: a byte count that lowers the whole-file streaming bound. Payloads above it are streamed through a bounded buffer. Values are clamped to the 256 MiB protocol ceiling, so it can only lower, never raise, the bound. |
## Implementation Details
@@ -580,9 +581,9 @@ remote SSH argv is already built injection-safe.
| `--files-from ` | Read the source file list from FILE (paths relative to the source root). |
| `-0, --from0` | Treat entries in `--files-from` files as NUL-delimited instead of newline-delimited. |
| `--delay-updates` | Put updated files into place only at the end of the transfer (the fixed `.fastsync-stage` staging name diverges from rsync; see [`RSYNC_COMPAT.md`](RSYNC_COMPAT.md)). |
-| `--compare-dest ` | Extra comparison basis: unchanged files are not transferred (requires/implies `--incremental`; a basis MISS above the 256 MiB whole-file payload bound is refused — see [`RSYNC_COMPAT.md`](RSYNC_COMPAT.md)). |
-| `--copy-dest ` | Like `--compare-dest`, but copies the unchanged file from DIR into the destination (same basis-size caveat; see [`RSYNC_COMPAT.md`](RSYNC_COMPAT.md)). |
-| `--link-dest ` | Like `--copy-dest`, but hard-links the unchanged file from DIR (repeatable; earlier DIRs win; same basis-size caveat; see [`RSYNC_COMPAT.md`](RSYNC_COMPAT.md)). |
+| `--compare-dest ` | Extra comparison basis: unchanged files are not transferred (requires/implies `--incremental`; a basis of any size is supported, streamed in bounded chunks — see [`RSYNC_COMPAT.md`](RSYNC_COMPAT.md)). |
+| `--copy-dest ` | Like `--compare-dest`, but copies the unchanged file from DIR into the destination (the copy is streamed, so a basis of any size works; see [`RSYNC_COMPAT.md`](RSYNC_COMPAT.md)). |
+| `--link-dest ` | Like `--copy-dest`, but hard-links the unchanged file from DIR (repeatable; earlier DIRs win; a basis of any size works; see [`RSYNC_COMPAT.md`](RSYNC_COMPAT.md)). |
| `--verify-basis` | FastSync-only: require a basis hit to match the source by whole-file digest instead of trusting the size+mtime quick-check (default matches rsync). |
| `--preallocate` | Allocate destination file space up front (fail-fast on a full disk). |
| `--append` | Resume a shorter destination by appending only its tail (prefix not verified; requires `--incremental`). |
diff --git a/RSYNC_COMPAT.md b/RSYNC_COMPAT.md
index 677c8aa..6ccc8b2 100644
--- a/RSYNC_COMPAT.md
+++ b/RSYNC_COMPAT.md
@@ -60,7 +60,7 @@ matrix is **111 ✅ / 13 ⚠️ / 33 ❌ = 157**.
**Parity cycle 2.29 (on `feat/parity-2.29`; `PROTOCOL_VERSION` stays 2.28.0 — no wire change was needed).** Five independent residuals were closed and four rows moved to ✅:
- **Scanner order.** The sequential scanner now buffers and sorts each directory's inspected entries (non-directories ascending, then directories ascending) and walks them depth-first, reproducing rsync 3.4.1's flist order. This makes the `--info=name` transfer order, the `--delete-during`/`--delete-delay`/`-n` would-delete order, and the partial-`--max-delete` survivor set byte-identical to rsync (`test_parity_order.py`). `--threads` has no rsync analogue and stays unordered.
- **Delete timing.** The complete per-directory plan set is transmitted before the first data frame, so a mid-transfer abort has already removed every planned extra like rsync's generator; `-d/--dirs` uses the same per-directory plans (shielded untraversed subdirectories) instead of the end-of-transfer commit (`test_delete_boundary_parity.py`). `-n`/`--delete`/`--del`/`--delete-delay` move ⚠️ → ✅.
-- **Basis relative-DIR.** A relative `--compare-dest`/`--copy-dest`/`--link-dest` DIR resolves against the destination directory with the transfer-relative name appended, exactly like rsync (`test_parity_basis_fuzzy.py`); the >256 MiB basis-MISS limit remains (a general whole-file limit, not basis-specific).
+- **Basis relative-DIR.** A relative `--compare-dest`/`--copy-dest`/`--link-dest` DIR resolves against the destination directory with the transfer-relative name appended, exactly like rsync (`test_parity_basis_fuzzy.py`); the >256 MiB basis-MISS limit is closed (#318 streams a MISS above the bound, `test_large_stream.py`).
- **Fuzzy eligibility.** The `-y/--fuzzy` candidate search no longer inherits the ordinary delta engine's 16 KiB minimum or 10× ratio bound, so an oversized or sub-16-KiB sibling is reused as rsync reuses it (`test_parity_basis_fuzzy.py`).
- **Output partials.** `--info=mount`/`--info=stats`, the `--stats` `dir:` breakdown under `-r`, and real `--debug` output for `flist`/`del`/`hash`/`deltasum`/`recv`/`filter`/`send` were added (`test_parity_info_mount_stats.py`, `test_output_parity.py`, `test_parity_debug.py`); those rows stay ⚠️ for their remaining documented residuals. `--delete-before`'s phase-0 late-file divergence and the `--progress` root/ancestor/symlink feedback remain open (they need a receiver→sender event channel), and the >256 MiB single-file streaming limit (B4) was not addressed. The matrix is now **120 ✅ / 10 ⚠️ / 27 ❌ = 157**.
@@ -703,10 +703,10 @@ targets verbatim, matching rsync.
|------|-------------------|-----------------|-------|
| `--checksum` | Skip based on checksum | ✅ Parity | `-c`/`--checksum` compares per-file whole-file content digests to skip unchanged files. **As of protocol 2.23.0 the short `-c` implies the checksum quick-check**, so a plain `-c` run verifies content rather than only affecting the `--incremental` handshake. The digest algorithm is `xxh128` by default (protocol 2.26.0's negotiated default) and is selectable via `--checksum-choice`/`--cc` (`xxh128`/`xxh3`/`xxh64`/`xxhash`/`md5`/`md4`/`sha1`/`none`/`auto`, plus rsync's two-name form) and `--checksum-seed=NUM` (see those rows) |
| `--checksum-choice=STR`, `--cc=STR` | Choose checksum algorithm | ✅ Parity | Real algorithm selection for the per-file whole-file digest used by the `--incremental`/`--checksum` handshake and basis-dir verification. **Protocol 2.26.0 accepts rsync 3.4.1's full set** — `xxh128` (the negotiated default), `xxh3`, `xxh64`, `xxhash`, `md5`, `md4`, `sha1`, `none`, `auto`, and the two-name `transfer,pre-transfer` form — with rsync's exit-4 rejection of an unknown name and of `none` on the transfer side when `--checksum` is on. `--cc=ALG` and space forms both parse. The algorithm id and seed cross the wire; the receiver hashes its old file with the same algorithm+seed and the per-file `STATUS_CHECK` handshake carries a bounded digest pinned to the negotiated length. `checksum_digest_file` now streams **every** supported algorithm (md4 via the self-contained RFC 1320 code, sha1/md5 via EVP, none as an empty digest), so the streaming path matches its contract, and `--out-format %C` uses the selected **transfer** half of a two-name choice and renders each algorithm byte-for-byte like rsync (xxh128 high-then-low, xxh64/xxh3 big-endian, md5/md4/sha1 standard hex, none a blank 2-char column) — differential-tested across all algorithms. **Track-3b finding — the block-checksum residual is not observable.** rsync applies the choice to the block checksum on its wire too, while FastSync selects only the whole-file comparison digest and keeps the delta BLOCK strong checksum fixed at xxHash32 (`DeltaBlockSig`, `delta_signature_create_seeded`). Because FastSync does not interoperate with rsync on the wire, only the compared surface matters, and a pre-seeded delta differential against rsync 3.4.1 (`--no-whole-file -B8192 --stats --out-format=%c|%C %n` vs `--incremental --delta`) shows the choice does not move it: across `xxh64`, `xxh128`, `xxh3`, `md5`, `md4`, `sha1` and both two-name orders the destination tree is byte-identical, `Matched data`/`Literal data`/`Total transferred file size` are unchanged (and equal to rsync's with the block size pinned), the `%c` block-checksum token is invariant (rsync `16 + 6·ceil(size/block)`, already documented under `--out-format`; FastSync its own basis-read counter), and the exit code is 0. The negotiated algorithm is visible only in `%C`, which applies it to the whole-file transfer digest and matches rsync byte-for-byte. A false block match would require the 4-byte adler32 AND the 4-byte xxHash32 to collide; at the 256 MiB maximum with the 1 KiB minimum block size the expected false matches are ≤2⁻¹⁸, and the choice cannot change this because FastSync's block strong sum is fixed. `auto` now consults `RSYNC_CHECKSUM_LIST` (rsync's whitespace-separated preference list; unknown names skipped, first supported wins, all-unknown is exit 4) before the compiled-in order; because both peers run the identical build this deterministic resolution needs no rsync peer probe, and an explicit `--cc` still wins. The list is differential-tested through `--out-format %C` (byte-identical digests to rsync for md5/sha1/xxh3) |
-| `--compare-dest=DIR` | Compare dest files relative to DIR | ⚠️ Caveat | DIR is a receiver-side basis; protocol 2.26.0 uses an absolute path verbatim (rsync semantics) and resolves a relative path below the destination root (`..` components are rejected, `//` collapsed and trailing `/` dropped) — note rsync resolves a relative DIR against the destination directory while FastSync resolves it below the receive root and appends the mirrored source path, so the same relative spelling addresses a different tree (use an absolute DIR for exact parity). On the receiver's per-file check (implies `--incremental`) an exact match is rsync's metadata quick-check: same size and mtime (unless `--size-only`; `-I` disables matching), with NO content digest required by default (track 5a). A match suppresses the data transfer. The FastSync-only `--verify-basis` restores the stricter whole-file content equality. compare-dest never copies: it only skips a file the destination does **not** already hold (sparse destination, rsync parity), and is consulted before the normal delta/full paths. Repeatable; searched in command-line order, first match wins. Differential-tested against rsync 3.4.1 (`compare_dest`, and `test_verify_basis_restores_strict_content`). **Relative-DIR parity (parity-2.29):** a relative DIR now resolves against the destination directory with the file's transfer-relative name appended, exactly like rsync 3.4.1, instead of FastSync's source-mirrored wire path (the historical spelling stays as a fallback; differential `test_parity_basis_fuzzy.py`). Residual: a basis MISS above the 256 MiB whole-file payload bound is refused up front (FastSync's general whole-file limit, not basis-specific); rsync applies basis dirs to arbitrary sizes. Wire: a basis-count field plus the `verify_basis` bool are present on the config frame (protocol 2.9.0/2.28.0) |
-| `--copy-dest=DIR` | Include copies of unchanged files | ⚠️ Caveat | Same basis rules as `--compare-dest`, but an exact match materializes a **local copy** of the DIR file into the destination (via the atomic temp+rename store path, so `--existing`/`--ignore-existing`/`--update`/`--backup`/`--delay-updates` all still apply) instead of transferring data. Track 5a re-applies the SOURCE attributes on the copy (rsync's "copy then fix attributes"): the sender transmits the source metadata with the basis check frame, so the copy's mode/uid/gid/mtime match the source rather than the basis inode (differential `copy_dest` compares modes). The copy streams the basis file through a bounded buffer, so a basis larger than the whole-file payload bound still materializes. Repeatable; command-line order = priority. Requires `--incremental` (implied); incompatible with `-s`. Wire: protocol 2.9.0 |
-| `--link-dest=DIR` | Hardlink to files when unchanged | ⚠️ Caveat | Same basis rules as `--copy-dest`, but an exact match installs an atomic **hard link** to the DIR file (temp hard link + rename) so no data or disk space is used; where the link is impossible (basis on another filesystem, filesystem refuses links) it falls back cleanly to a byte-identical local copy (streamed from the basis, so an over-limit basis still works), never a corrupt/partial file. `--delay-updates` stages the link and publishes by rename, so the final entry stays a real hard link. Repeatable (searched in command-line order, first match wins). Differential-tested against rsync 3.4.1 (`link_dest`). Inherent shared-inode semantics (identical to rsync): a link keeps the basis inode's own mode/uid/gid and mtime — metadata is never written through the shared inode (that would mutate the basis file), so a later `--inplace` run that rewrites such a destination path **will mutate the basis snapshot** through the shared inode (use `--copy-dest` when the destination must stay independently writable); protocol 2.26.0 re-links an already up-to-date destination file to the basis; a `--remove-source-files` source satisfied by a basis dir is treated as skipped and therefore **retained** (never removed); basis dirs are excluded from `--delete`. **Relative-DIR parity (parity-2.29):** a relative DIR resolves against the destination directory with the transfer-relative name appended, exactly like rsync 3.4.1 (differential `test_parity_basis_fuzzy.py`). Residual: a basis MISS above the 256 MiB whole-file payload bound is refused (FastSync's general whole-file limit). Requires `--incremental` (implied); incompatible with `-s`. Wire: protocol 2.9.0 |
-| `-y`, `--fuzzy`, `--no-fuzzy` | Find similar file for basis | ⚠️ Caveat | `-y/--fuzzy` is a pure bandwidth optimization on the existing receiver-driven delta path: when a file must be transferred and the destination holds no usable content at the exact path (file absent, or the destination file is outside the delta engine's size bounds), the receiver searches the SAME destination directory for an existing regular file whose basename is similar to the incoming name and uses it as the delta basis, so the sender transmits only the differences instead of the whole file. The output is always byte-exact regardless of which (or whether any) basis is chosen. Decision location: the receiver performs the candidate search inside `receive_incremental_check` and sends the normal `STATUS_DELTA_SIGNATURE`; the sender never learns the basis was a different file, so no new frame type or sender logic was needed — only the config frame grew a `fuzzy` boolean, so `PROTOCOL_VERSION` was bumped **2.8.0 → 2.9.0** (peers must match). Similarity heuristic (a deterministic port of rsync 3.4.1's matcher — `util1.c` `fuzzy_distance`/`find_filename_suffix` plus `generator.c find_fuzzy`'s exact size+mtime pass — documented precisely): candidates are the target's sibling entries in its destination directory, opened `O_NOFOLLOW`/`AT_SYMLINK_NOFOLLOW` under the confined root (symlinks never followed; nothing outside the destination root is ever read or hashed); dotfiles, directories, the target's own name, and the `.fastsync-stage`/temp scratch names are excluded; like the ordinary delta path, the block signature the receiver transmits is derived from on-disk content it may not otherwise send, so a negotiated `--fuzzy` run exposes the destination's sibling files (at block granularity) to the sender as a known-plaintext oracle — the same information class as the normal delta handshake over the file being replaced; the size gate is the delta engine's own bounds (both files ≥ 16 KiB, ≤ `--delta-max`, ratio ≤ 10×); rsync's fuzzy matcher is not tied to a delta size bound and empirically reuses a basis well outside FastSync's window (a 64 KiB source against a repeated-content sibling from 0.25× to 10000×, and files as small as 300 B), so candidate ELIGIBILITY — and hence the chosen basis — can differ even though the name heuristic is the same; protocol 2.26.0 uses rsync's weighted-Levenshtein name/suffix distance plus an exact size+mtime pass and reads a single best candidate; the tie-break (smallest size gap, then lexical name) is deterministic where rsync leaves equal distances to its file-list order; the directory scan is capped at 4096 entries so a pathological directory cannot stall a transfer. When fuzzy applies: only to files the receiver would otherwise send whole — the destination's own file is always preferred as the delta basis when it exists and fits the delta size bounds, so fuzzy does NOT replace an existing-but-different destination basis; FastSync's 10× delta size-ratio bound means an existing destination file that is too far away in size still lets the fuzzy search run. When no similar candidate exists the transfer falls back to the normal whole-file transfer. rsync-divergence note: the name matching is rsync's own rule; the residual is eligibility bounded by FastSync's delta engine, so no name-matcher port can widen it. Because FastSync's delta machinery is off by default (rsync's is on), `--fuzzy` implies `--incremental` + `--delta` (unless `--whole-file`/`-W` or an explicit `--no-delta` switched delta off, in which case fuzzy is inert — matching rsync where `--whole-file` makes fuzzy irrelevant). Unlike the basis-dir options, `--fuzzy` honors an explicit `--no-incremental` (it does not force the handshake back on); an explicit `--no-incremental` also suppresses the delta implication so no invalid `--delta requires --incremental` config results. `--no-fuzzy` negates it. All surrounding semantics are untouched: a fuzzy-reconstructed file is stored as a normal file, so `--remove-source-files`, itemize/`-i`, `--stats`, `--backup`, `--delay-updates`, `--existing`/`--ignore-existing`/`--update` behave exactly as for a whole-file transfer (the fuzzy delta does not skip the file). **Reclassified Caveat (track 5b):** the output is always byte-exact regardless of the basis, and the name rule is rsync's, **Eligibility parity (parity-2.29):** the fuzzy candidate search no longer inherits the ordinary delta engine's 16 KiB minimum or 10× size-ratio bound, so an oversized or sub-16-KiB sibling is now reused exactly as rsync reuses it — the chosen basis (and therefore `Matched data`/`Literal data`/`Total transferred file size`) matches rsync where the block size is pinned (differential `test_parity_basis_fuzzy.py`; the old boundary tests were flipped to assert both tools reuse the basis). Residual: the whole-basis buffer cap (`MAX_RECEIVE_WHOLE_FILE_SIZE`, 256 MiB) and the deterministic tie-break where rsync leaves equal distances to its file-list order. |
+| `--compare-dest=DIR` | Compare dest files relative to DIR | ⚠️ Caveat | DIR is a receiver-side basis; protocol 2.26.0 uses an absolute path verbatim (rsync semantics) and resolves a relative path below the destination root (`..` components are rejected, `//` collapsed and trailing `/` dropped) — note rsync resolves a relative DIR against the destination directory while FastSync resolves it below the receive root and appends the mirrored source path, so the same relative spelling addresses a different tree (use an absolute DIR for exact parity). On the receiver's per-file check (implies `--incremental`) an exact match is rsync's metadata quick-check: same size and mtime (unless `--size-only`; `-I` disables matching), with NO content digest required by default (track 5a). A match suppresses the data transfer. The FastSync-only `--verify-basis` restores the stricter whole-file content equality. compare-dest never copies: it only skips a file the destination does **not** already hold (sparse destination, rsync parity), and is consulted before the normal delta/full paths. Repeatable; searched in command-line order, first match wins. Differential-tested against rsync 3.4.1 (`compare_dest`, and `test_verify_basis_restores_strict_content`). **Relative-DIR parity (parity-2.29):** a relative DIR now resolves against the destination directory with the file's transfer-relative name appended, exactly like rsync 3.4.1, instead of FastSync's source-mirrored wire path (the historical spelling stays as a fallback; differential `test_parity_basis_fuzzy.py`). **Streaming (parity-2.29/#318):** the over-256 MiB basis-MISS residual is closed — a basis MISS above the whole-file bound is no longer refused; it falls through to the streaming whole-file transfer, so basis dirs now apply to files of arbitrary size (`test_large_stream.py`). Residual: a relative DIR additionally probes the historical mirror-appended spelling as a fallback, where rsync probes only the destination-relative name. Wire: a basis-count field plus the `verify_basis` bool are present on the config frame (protocol 2.9.0/2.28.0) |
+| `--copy-dest=DIR` | Include copies of unchanged files | ⚠️ Caveat | Same basis rules as `--compare-dest`, but an exact match materializes a **local copy** of the DIR file into the destination (via the atomic temp+rename store path, so `--existing`/`--ignore-existing`/`--update`/`--backup`/`--delay-updates` all still apply) instead of transferring data. Track 5a re-applies the SOURCE attributes on the copy (rsync's "copy then fix attributes"): the sender transmits the source metadata with the basis check frame, so the copy's mode/uid/gid/mtime match the source rather than the basis inode (differential `copy_dest` compares modes). The copy streams the basis file through a bounded buffer, so a basis larger than the whole-file payload bound still materializes, and a basis MISS above the bound now falls through to the streaming whole-file transfer (the historical over-256 MiB MISS refusal is closed; `test_large_stream.py`). Repeatable; command-line order = priority. Requires `--incremental` (implied); incompatible with `-s`. Wire: protocol 2.9.0 |
+| `--link-dest=DIR` | Hardlink to files when unchanged | ⚠️ Caveat | Same basis rules as `--copy-dest`, but an exact match installs an atomic **hard link** to the DIR file (temp hard link + rename) so no data or disk space is used; where the link is impossible (basis on another filesystem, filesystem refuses links) it falls back cleanly to a byte-identical local copy (streamed from the basis, so an over-limit basis still works), never a corrupt/partial file. `--delay-updates` stages the link and publishes by rename, so the final entry stays a real hard link. Repeatable (searched in command-line order, first match wins). Differential-tested against rsync 3.4.1 (`link_dest`). Inherent shared-inode semantics (identical to rsync): a link keeps the basis inode's own mode/uid/gid and mtime — metadata is never written through the shared inode (that would mutate the basis file), so a later `--inplace` run that rewrites such a destination path **will mutate the basis snapshot** through the shared inode (use `--copy-dest` when the destination must stay independently writable); protocol 2.26.0 re-links an already up-to-date destination file to the basis; a `--remove-source-files` source satisfied by a basis dir is treated as skipped and therefore **retained** (never removed); basis dirs are excluded from `--delete`. **Relative-DIR parity (parity-2.29):** a relative DIR resolves against the destination directory with the transfer-relative name appended, exactly like rsync 3.4.1 (differential `test_parity_basis_fuzzy.py`). **Streaming (parity-2.29/#318):** the over-256 MiB basis-MISS residual is closed — a MISS above the bound now falls through to the streaming whole-file transfer (`test_large_stream.py`). Requires `--incremental` (implied); incompatible with `-s`. Wire: protocol 2.9.0 |
+| `-y`, `--fuzzy`, `--no-fuzzy` | Find similar file for basis | ⚠️ Caveat | `-y/--fuzzy` is a pure bandwidth optimization on the existing receiver-driven delta path: when a file must be transferred and the destination holds no usable content at the exact path (file absent, or the destination file is outside the delta engine's size bounds), the receiver searches the SAME destination directory for an existing regular file whose basename is similar to the incoming name and uses it as the delta basis, so the sender transmits only the differences instead of the whole file. The output is always byte-exact regardless of which (or whether any) basis is chosen. Decision location: the receiver performs the candidate search inside `receive_incremental_check` and sends the normal `STATUS_DELTA_SIGNATURE`; the sender never learns the basis was a different file, so no new frame type or sender logic was needed — only the config frame grew a `fuzzy` boolean, so `PROTOCOL_VERSION` was bumped **2.8.0 → 2.9.0** (peers must match). Similarity heuristic (a deterministic port of rsync 3.4.1's matcher — `util1.c` `fuzzy_distance`/`find_filename_suffix` plus `generator.c find_fuzzy`'s exact size+mtime pass — documented precisely): candidates are the target's sibling entries in its destination directory, opened `O_NOFOLLOW`/`AT_SYMLINK_NOFOLLOW` under the confined root (symlinks never followed; nothing outside the destination root is ever read or hashed); dotfiles, directories, the target's own name, and the `.fastsync-stage`/temp scratch names are excluded; like the ordinary delta path, the block signature the receiver transmits is derived from on-disk content it may not otherwise send, so a negotiated `--fuzzy` run exposes the destination's sibling files (at block granularity) to the sender as a known-plaintext oracle — the same information class as the normal delta handshake over the file being replaced; the size gate is the delta engine's own bounds (both files ≥ 16 KiB, ≤ `--delta-max`, ratio ≤ 10×); rsync's fuzzy matcher is not tied to a delta size bound and empirically reuses a basis well outside FastSync's window (a 64 KiB source against a repeated-content sibling from 0.25× to 10000×, and files as small as 300 B), so candidate ELIGIBILITY — and hence the chosen basis — can differ even though the name heuristic is the same; protocol 2.26.0 uses rsync's weighted-Levenshtein name/suffix distance plus an exact size+mtime pass and reads a single best candidate; the tie-break (smallest size gap, then lexical name) is deterministic where rsync leaves equal distances to its file-list order; the directory scan is capped at 4096 entries so a pathological directory cannot stall a transfer. When fuzzy applies: only to files the receiver would otherwise send whole — the destination's own file is always preferred as the delta basis when it exists and fits the delta size bounds, so fuzzy does NOT replace an existing-but-different destination basis; FastSync's 10× delta size-ratio bound means an existing destination file that is too far away in size still lets the fuzzy search run. When no similar candidate exists the transfer falls back to the normal whole-file transfer. rsync-divergence note: the name matching is rsync's own rule; the residual is eligibility bounded by FastSync's delta engine, so no name-matcher port can widen it. Because FastSync's delta machinery is off by default (rsync's is on), `--fuzzy` implies `--incremental` + `--delta` (unless `--whole-file`/`-W` or an explicit `--no-delta` switched delta off, in which case fuzzy is inert — matching rsync where `--whole-file` makes fuzzy irrelevant). Unlike the basis-dir options, `--fuzzy` honors an explicit `--no-incremental` (it does not force the handshake back on); an explicit `--no-incremental` also suppresses the delta implication so no invalid `--delta requires --incremental` config results. `--no-fuzzy` negates it. All surrounding semantics are untouched: a fuzzy-reconstructed file is stored as a normal file, so `--remove-source-files`, itemize/`-i`, `--stats`, `--backup`, `--delay-updates`, `--existing`/`--ignore-existing`/`--update` behave exactly as for a whole-file transfer (the fuzzy delta does not skip the file). **Reclassified Caveat (track 5b):** the output is always byte-exact regardless of the basis, and the name rule is rsync's, **Eligibility parity (parity-2.29):** the fuzzy candidate search no longer inherits the ordinary delta engine's 16 KiB minimum or 10× size-ratio bound, so an oversized or sub-16-KiB sibling is now reused exactly as rsync reuses it — the chosen basis (and therefore `Matched data`/`Literal data`/`Total transferred file size`) matches rsync where the block size is pinned (differential `test_parity_basis_fuzzy.py`; the old boundary tests were flipped to assert both tools reuse the basis). Residual: the deterministic tie-break where rsync leaves equal distances to its file-list order. (The whole-basis buffer cap — `MAX_RECEIVE_WHOLE_FILE_SIZE`, 256 MiB — is closed: #318 opens the sibling and streams its block signature and the delta reconstruction from the descriptor in bounded chunks, so a basis of arbitrary size is reused without materializing it; `test_large_stream.py`.) |
## 12. Compression
diff --git a/src/client/client_send.c b/src/client/client_send.c
index 354b86f..27f7745 100644
--- a/src/client/client_send.c
+++ b/src/client/client_send.c
@@ -877,6 +877,23 @@ static bool source_is_regular_file(const File* file) {
return stat(file->path, &st) == 0 && S_ISREG(st.st_mode);
}
+/* True when an over-threshold source will be sent by STREAMING from disk rather
+ * than loaded into memory: either the zero-copy sendfile path (no compression)
+ * or the sender-side streaming compressor (zstd/zlib, when this file is not on
+ * the --skip-compress list). Otherwise the loader must materialize it. */
+static bool loader_can_stream(const Config* config, const File* file) {
+ if (!config || !file || !file->data || file->data->size <= STREAM_THRESHOLD)
+ return false;
+ if (!config->use_compression)
+ return true;
+ if (!compression_stream_compress_supported(compression_get_algo()) ||
+ config->compression_level <= 0)
+ return false;
+ int skip_count = config->skip_compress_set ? config->skip_compress_count : -1;
+ return !compression_should_skip_with_suffixes(file->path, config->skip_compress_suffixes,
+ skip_count);
+}
+
static int send_chunk_with_removal(Client* client, Chunk* chunk, Config* config,
ArrayList* remove_sources, TransferStats* stats) {
/* --stderr=client: this is a frame boundary, so forward any diagnostics the
@@ -1423,7 +1440,7 @@ static int load_files_multithreaded(void* pipeline_context) {
if (!context->config->use_sendfile) {
for (int i = 0; i < chunk->element_count; i++) {
File* f = chunk->items[i];
- if (f->data->size > STREAM_THRESHOLD && !context->config->use_compression)
+ if (loader_can_stream(context->config, f))
continue;
if (!file_load_data(f)) {
log_message(LOG_LEVEL_ERROR, "Failed to load file data");
@@ -1812,7 +1829,7 @@ static bool send_files_run(Config* config, SendFilesState* state) {
bool load_ok = true;
for (int i = 0; i < current_chunk->element_count; i++) {
File* f = current_chunk->items[i];
- if (f->data->size > STREAM_THRESHOLD && !config->use_compression)
+ if (loader_can_stream(config, f))
continue;
if (!file_load_data(f)) {
log_message(LOG_LEVEL_ERROR, "Failed to load file data");
diff --git a/src/shared/compression.c b/src/shared/compression.c
index a3d7d66..150b2f8 100644
--- a/src/shared/compression.c
+++ b/src/shared/compression.c
@@ -3,6 +3,7 @@
#include "log.h"
#include "protocol.h"
#include "utils.h"
+#include
#include
#include
#include
@@ -716,3 +717,379 @@ Data* data_decompress_limited(Data* compressed_data, size_t maximum_size) {
Data* data_decompress(Data* compressed_data) {
return data_decompress_limited(compressed_data, MAX_DECOMPRESSED_SIZE);
}
+
+/* ---- streaming decompression ---- */
+
+#define STREAM_DECOMPRESS_OUT_CHUNK (256 * 1024)
+
+struct CompressionStreamDecompressor {
+ CompressionAlgo algo;
+ unsigned long long expected_out;
+ unsigned long long total;
+ int out_fd;
+ unsigned char* out_buf;
+ ZSTD_DCtx* dctx;
+ z_stream zs;
+ bool zs_initialized;
+ bool failed;
+};
+
+static bool stream_write_all(int fd, const void* data, size_t size) {
+ const unsigned char* p = data;
+ size_t done = 0;
+ while (done < size) {
+ ssize_t n = write(fd, p + done, size - done);
+ if (n < 0 && errno == EINTR)
+ continue;
+ if (n <= 0)
+ return false;
+ done += (size_t)n;
+ }
+ return true;
+}
+
+CompressionStreamDecompressor*
+compression_stream_decompressor_create(CompressionAlgo algo, unsigned long long expected_out) {
+ CompressionStreamDecompressor* d = calloc(1, sizeof(*d));
+ if (!d)
+ return NULL;
+ d->algo = algo;
+ d->expected_out = expected_out;
+ d->out_fd = -1;
+ d->out_buf = malloc(STREAM_DECOMPRESS_OUT_CHUNK);
+ if (!d->out_buf) {
+ free(d);
+ return NULL;
+ }
+ if (algo == COMPRESSION_ALGO_ZSTD) {
+ d->dctx = ZSTD_createDCtx();
+ if (!d->dctx) {
+ free(d->out_buf);
+ free(d);
+ return NULL;
+ }
+ } else if (algo == COMPRESSION_ALGO_ZLIB || algo == COMPRESSION_ALGO_ZLIBX) {
+ if (inflateInit(&d->zs) != Z_OK) {
+ free(d->out_buf);
+ free(d);
+ return NULL;
+ }
+ d->zs_initialized = true;
+ } else if (algo != COMPRESSION_ALGO_NONE) {
+ /* lz4's block format cannot be decompressed incrementally. */
+ free(d->out_buf);
+ free(d);
+ return NULL;
+ }
+ return d;
+}
+
+static bool stream_emit(CompressionStreamDecompressor* d, const void* buf, size_t len) {
+ if (len == 0)
+ return true;
+ if (d->expected_out != 0 && (d->total > d->expected_out || len > d->expected_out - d->total)) {
+ d->failed = true;
+ return false;
+ }
+ if (!stream_write_all(d->out_fd, buf, len)) {
+ d->failed = true;
+ return false;
+ }
+ d->total += len;
+ return true;
+}
+
+static bool stream_feed_none(CompressionStreamDecompressor* d, const void* in, size_t in_len,
+ bool* done) {
+ if (!stream_emit(d, in, in_len))
+ return false;
+ /* NONE has no end marker; the caller knows the frame length. */
+ *done = true;
+ return true;
+}
+
+static bool stream_feed_zstd(CompressionStreamDecompressor* d, const void* in, size_t in_len,
+ bool* done) {
+ ZSTD_inBuffer input = {in, in_len, 0};
+ while (input.pos < input.size) {
+ ZSTD_outBuffer output = {d->out_buf, STREAM_DECOMPRESS_OUT_CHUNK, 0};
+ size_t ret = ZSTD_decompressStream(d->dctx, &output, &input);
+ if (ZSTD_isError(ret)) {
+ d->failed = true;
+ return false;
+ }
+ if (!stream_emit(d, d->out_buf, output.pos))
+ return false;
+ if (ret == 0) {
+ *done = true;
+ /* Trailing bytes after a complete frame are malformed; stop consuming. */
+ if (input.pos < input.size) {
+ d->failed = true;
+ return false;
+ }
+ return true;
+ }
+ }
+ return true;
+}
+
+static bool stream_feed_zlib(CompressionStreamDecompressor* d, const void* in, size_t in_len,
+ bool* done) {
+ d->zs.next_in = (Bytef*)in;
+ d->zs.avail_in = (uInt)in_len;
+ while (d->zs.avail_in > 0) {
+ d->zs.next_out = d->out_buf;
+ d->zs.avail_out = STREAM_DECOMPRESS_OUT_CHUNK;
+ int rc = inflate(&d->zs, Z_NO_FLUSH);
+ if (rc != Z_OK && rc != Z_STREAM_END && rc != Z_BUF_ERROR) {
+ d->failed = true;
+ return false;
+ }
+ size_t produced = STREAM_DECOMPRESS_OUT_CHUNK - d->zs.avail_out;
+ if (!stream_emit(d, d->out_buf, produced))
+ return false;
+ if (rc == Z_STREAM_END) {
+ *done = true;
+ return d->zs.avail_in == 0;
+ }
+ if (rc == Z_BUF_ERROR && produced == 0) {
+ /* Need more input. */
+ break;
+ }
+ }
+ return true;
+}
+
+bool compression_stream_decompressor_feed(CompressionStreamDecompressor* d, const void* in,
+ size_t in_len, int out_fd, bool* done) {
+ if (!d || d->failed)
+ return false;
+ d->out_fd = out_fd;
+ if (done)
+ *done = false;
+ switch (d->algo) {
+ case COMPRESSION_ALGO_NONE:
+ return stream_feed_none(d, in, in_len, done);
+ case COMPRESSION_ALGO_ZSTD:
+ return stream_feed_zstd(d, in, in_len, done);
+ case COMPRESSION_ALGO_ZLIB:
+ case COMPRESSION_ALGO_ZLIBX:
+ return stream_feed_zlib(d, in, in_len, done);
+ case COMPRESSION_ALGO_LZ4:
+ break;
+ }
+ d->failed = true;
+ return false;
+}
+
+unsigned long long compression_stream_decompressor_total(const CompressionStreamDecompressor* d) {
+ return d ? d->total : 0;
+}
+
+void compression_stream_decompressor_destroy(CompressionStreamDecompressor* d) {
+ if (!d)
+ return;
+ if (d->dctx)
+ ZSTD_freeDCtx(d->dctx);
+ if (d->zs_initialized)
+ inflateEnd(&d->zs);
+ free(d->out_buf);
+ free(d);
+}
+
+/* ---- streaming compression ---- */
+
+struct CompressionStreamCompressor {
+ CompressionAlgo algo;
+ int level;
+ ZSTD_CCtx* cctx;
+ z_stream zs;
+ bool zs_initialized;
+ unsigned char* out_buf;
+ bool failed;
+};
+
+bool compression_stream_compress_supported(CompressionAlgo algo) {
+ return algo == COMPRESSION_ALGO_ZSTD || algo == COMPRESSION_ALGO_ZLIB ||
+ algo == COMPRESSION_ALGO_ZLIBX;
+}
+
+CompressionStreamCompressor* compression_stream_compressor_create(CompressionAlgo algo, int level,
+ int threads) {
+ (void)threads;
+ if (!compression_algo_valid((int)algo) || algo == COMPRESSION_ALGO_LZ4)
+ return NULL;
+ CompressionStreamCompressor* c = calloc(1, sizeof(*c));
+ if (!c)
+ return NULL;
+ c->algo = algo;
+ c->level = level;
+ c->out_buf = malloc(STREAM_DECOMPRESS_OUT_CHUNK);
+ if (!c->out_buf) {
+ free(c);
+ return NULL;
+ }
+ if (algo == COMPRESSION_ALGO_ZSTD) {
+ c->cctx = ZSTD_createCCtx();
+ if (!c->cctx) {
+ free(c->out_buf);
+ free(c);
+ return NULL;
+ }
+ } else if (algo == COMPRESSION_ALGO_ZLIB || algo == COMPRESSION_ALGO_ZLIBX) {
+ if (deflateInit(&c->zs, level < 1 ? Z_DEFAULT_COMPRESSION : level) != Z_OK) {
+ free(c->out_buf);
+ free(c);
+ return NULL;
+ }
+ c->zs_initialized = true;
+ }
+ return c;
+}
+
+bool compression_stream_compressor_begin(CompressionStreamCompressor* c,
+ unsigned long long raw_size, int out_fd) {
+ if (!c || c->failed)
+ return false;
+ unsigned char hdr[1 + 4];
+ size_t hdr_len = 1;
+ hdr[0] = (unsigned char)c->algo;
+ if (c->algo == COMPRESSION_ALGO_ZLIB || c->algo == COMPRESSION_ALGO_ZLIBX) {
+ uint32_t size32 = raw_size > UINT32_MAX ? UINT32_MAX : (uint32_t)raw_size;
+ for (int i = 0; i < 4; i++)
+ hdr[1 + i] = (uint8_t)((size32 >> (8 * i)) & 0xff);
+ hdr_len = 5;
+ }
+ if (c->algo == COMPRESSION_ALGO_ZSTD) {
+ /* Pledge the source size and force the frame content-size field so the
+ receiver can decide whether to stream from the frame header alone. */
+ if (ZSTD_isError(ZSTD_CCtx_setPledgedSrcSize(c->cctx, raw_size)) ||
+ ZSTD_isError(ZSTD_CCtx_setParameter(c->cctx, ZSTD_c_compressionLevel, c->level)) ||
+ ZSTD_isError(ZSTD_CCtx_setParameter(c->cctx, ZSTD_c_contentSizeFlag, 1))) {
+ c->failed = true;
+ return false;
+ }
+ if (ZSTD_isError(ZSTD_CCtx_setParameter(c->cctx, ZSTD_c_checksumFlag, 0))) {
+ c->failed = true;
+ return false;
+ }
+ }
+ if (!stream_write_all(out_fd, hdr, hdr_len)) {
+ c->failed = true;
+ return false;
+ }
+ return true;
+}
+
+static bool stream_compress_zlib(CompressionStreamCompressor* c, const void* in, size_t in_len,
+ int out_fd, int flush) {
+ c->zs.next_in = (Bytef*)in;
+ c->zs.avail_in = (uInt)in_len;
+ do {
+ c->zs.next_out = c->out_buf;
+ c->zs.avail_out = STREAM_DECOMPRESS_OUT_CHUNK;
+ int rc = deflate(&c->zs, flush);
+ if (rc != Z_OK && rc != Z_STREAM_END && rc != Z_BUF_ERROR) {
+ c->failed = true;
+ return false;
+ }
+ size_t produced = STREAM_DECOMPRESS_OUT_CHUNK - c->zs.avail_out;
+ if (!stream_write_all(out_fd, c->out_buf, produced)) {
+ c->failed = true;
+ return false;
+ }
+ if (rc == Z_STREAM_END)
+ return true;
+ if (rc == Z_BUF_ERROR && produced == 0)
+ break;
+ } while (c->zs.avail_in > 0 || flush == Z_FINISH);
+ return true;
+}
+
+bool compression_stream_compressor_feed(CompressionStreamCompressor* c, const void* in,
+ size_t in_len, int out_fd) {
+ if (!c || c->failed)
+ return false;
+ if (c->algo == COMPRESSION_ALGO_NONE)
+ return stream_write_all(out_fd, in, in_len);
+ if (c->algo == COMPRESSION_ALGO_ZSTD) {
+ ZSTD_inBuffer input = {in, in_len, 0};
+ while (input.pos < input.size) {
+ ZSTD_outBuffer output = {c->out_buf, STREAM_DECOMPRESS_OUT_CHUNK, 0};
+ size_t ret = ZSTD_compressStream2(c->cctx, &output, &input, ZSTD_e_continue);
+ if (ZSTD_isError(ret)) {
+ c->failed = true;
+ return false;
+ }
+ if (!stream_write_all(out_fd, c->out_buf, output.pos)) {
+ c->failed = true;
+ return false;
+ }
+ if (output.pos == 0 && input.pos < input.size)
+ break; /* avoid spinning; zstd buffers the rest internally */
+ }
+ return true;
+ }
+ return stream_compress_zlib(c, in, in_len, out_fd, Z_NO_FLUSH);
+}
+
+bool compression_stream_compressor_finish(CompressionStreamCompressor* c, int out_fd) {
+ if (!c || c->failed)
+ return false;
+ if (c->algo == COMPRESSION_ALGO_NONE)
+ return true;
+ if (c->algo == COMPRESSION_ALGO_ZSTD) {
+ size_t ret;
+ do {
+ ZSTD_inBuffer input = {NULL, 0, 0};
+ ZSTD_outBuffer output = {c->out_buf, STREAM_DECOMPRESS_OUT_CHUNK, 0};
+ ret = ZSTD_compressStream2(c->cctx, &output, &input, ZSTD_e_end);
+ if (ZSTD_isError(ret)) {
+ c->failed = true;
+ return false;
+ }
+ if (!stream_write_all(out_fd, c->out_buf, output.pos)) {
+ c->failed = true;
+ return false;
+ }
+ } while (ret > 0);
+ return true;
+ }
+ return stream_compress_zlib(c, NULL, 0, out_fd, Z_FINISH);
+}
+
+void compression_stream_compressor_destroy(CompressionStreamCompressor* c) {
+ if (!c)
+ return;
+ if (c->cctx)
+ ZSTD_freeCCtx(c->cctx);
+ if (c->zs_initialized)
+ deflateEnd(&c->zs);
+ free(c->out_buf);
+ free(c);
+}
+
+unsigned long long compression_peek_frame_content_size(const void* buf, size_t len) {
+ if (!buf || len < 1)
+ return 0;
+ const uint8_t* p = buf;
+ uint8_t codec = p[0];
+ if (!compression_algo_valid(codec))
+ return 0;
+ if (codec == (uint8_t)COMPRESSION_ALGO_NONE)
+ return len - 1;
+ if (codec == (uint8_t)COMPRESSION_ALGO_ZSTD) {
+ if (len < 2)
+ return 0;
+ unsigned long long size = ZSTD_getFrameContentSize(p + 1, len - 1);
+ if (size == ZSTD_CONTENTSIZE_ERROR || size == ZSTD_CONTENTSIZE_UNKNOWN)
+ return 0;
+ return size;
+ }
+ if (len < 1 + LZ4_SIZE_PREFIX_LEN)
+ return 0;
+ uint32_t raw = 0;
+ for (int i = 0; i < LZ4_SIZE_PREFIX_LEN; i++)
+ raw |= (uint32_t)p[1 + i] << (8 * i);
+ return raw;
+}
diff --git a/src/shared/compression.h b/src/shared/compression.h
index ca874b8..4590b97 100644
--- a/src/shared/compression.h
+++ b/src/shared/compression.h
@@ -81,6 +81,54 @@ Data* data_compress_with_threads(Data* data_to_compress, int compression_level,
Data* data_decompress(Data* compressed_data);
bool compression_should_skip_with_suffixes(const char* path, char* const* suffixes, int count);
+/* Streaming decompression for a payload too large to hold in memory. The
+ * caller consumes the frame's leading codec byte (and, for lz4/zlib/zlibx, the
+ * 4-byte little-endian raw-size prefix) and then feeds the remaining frame
+ * bytes in bounded chunks; decompressed output is written straight to `out_fd`
+ * so neither the compressed nor the decompressed image is ever materialized.
+ * Only zstd (the default), zlib/zlibx and none support streaming; lz4's block
+ * format is one-shot, so its stream decompressor reports failure and the caller
+ * falls back (the whole-buffer path keeps its existing bound). */
+typedef struct CompressionStreamDecompressor CompressionStreamDecompressor;
+
+CompressionStreamDecompressor*
+compression_stream_decompressor_create(CompressionAlgo algo, unsigned long long expected_out);
+/* Feed one chunk. Returns false on a malformed frame, an I/O error, or when the
+ * total output would exceed `expected_out` (when non-zero). *done is set once
+ * the frame end has been reached. */
+bool compression_stream_decompressor_feed(CompressionStreamDecompressor* d, const void* in,
+ size_t in_len, int out_fd, bool* done);
+unsigned long long compression_stream_decompressor_total(const CompressionStreamDecompressor* d);
+void compression_stream_decompressor_destroy(CompressionStreamDecompressor* d);
+
+/* Peek the logical (decompressed) size from the leading bytes of a compressed
+ * frame (codec byte + header), returning 0 when it cannot be determined from
+ * the supplied prefix. Used to decide whether a frame must take the streaming
+ * path before its body is read. */
+unsigned long long compression_peek_frame_content_size(const void* buf, size_t len);
+
+/* Streaming compression (sender side). Compresses a source in bounded chunks
+ * into `out_fd` as one self-describing frame (codec byte, the lz4/zlib raw-size
+ * prefix, then the codec stream), so a whole file can be compressed without
+ * materializing it in memory. zstd/zlib/zlibx/none are supported; lz4's block
+ * format is one-shot, so its create() returns NULL and the caller keeps the
+ * buffered path. `raw_size` is the known source length (used for the zlib
+ * prefix and, for zstd, the frame content-size field). */
+typedef struct CompressionStreamCompressor CompressionStreamCompressor;
+
+/* True when `algo` can be stream-compressed (zstd/zlib/zlibx; lz4's block format
+ * is one-shot). Used by the sender to decide whether an over-threshold source
+ * may stay unloaded. */
+bool compression_stream_compress_supported(CompressionAlgo algo);
+CompressionStreamCompressor* compression_stream_compressor_create(CompressionAlgo algo, int level,
+ int threads);
+bool compression_stream_compressor_begin(CompressionStreamCompressor* c,
+ unsigned long long raw_size, int out_fd);
+bool compression_stream_compressor_feed(CompressionStreamCompressor* c, const void* in,
+ size_t in_len, int out_fd);
+bool compression_stream_compressor_finish(CompressionStreamCompressor* c, int out_fd);
+void compression_stream_compressor_destroy(CompressionStreamCompressor* c);
+
/* Release the calling thread's cached zstd contexts (compressor, decompressor
* and scratch buffer). The cache is thread-local and is also released
* automatically when a worker thread exits (via a C11 tss destructor) and for
diff --git a/src/shared/delta.c b/src/shared/delta.c
index dc5ea8a..cdc6f5a 100644
--- a/src/shared/delta.c
+++ b/src/shared/delta.c
@@ -1,10 +1,12 @@
#include "delta.h"
#include "log.h"
#include "protocol.h"
+#include
#include
#include
#include
#include
+#include
#define XXH_STATIC_LINKING_ONLY
#define XXH_IMPLEMENTATION
@@ -82,6 +84,64 @@ DeltaSignature* delta_signature_create_seeded(const void* old_file_data, uint64_
return sig;
}
+/* Bounded read of exactly `len` bytes at `off`; retries on EINTR. */
+static bool pread_all(int fd, void* buf, size_t len, uint64_t off) {
+ uint8_t* p = buf;
+ size_t done = 0;
+ while (done < len) {
+ ssize_t n = pread(fd, p + done, len - done, (off_t)(off + done));
+ if (n < 0 && errno == EINTR)
+ continue;
+ if (n <= 0)
+ return false;
+ done += (size_t)n;
+ }
+ return true;
+}
+
+DeltaSignature* delta_signature_create_fd_seeded(int fd, uint64_t old_file_size,
+ uint32_t block_size, uint32_t seed) {
+ if (fd < 0 || old_file_size == 0 || block_size == 0 || block_size > DELTA_BLOCK_SIZE_MAX ||
+ old_file_size > UINT32_MAX * (uint64_t)block_size)
+ return NULL;
+ uint32_t block_count = (uint32_t)((old_file_size + block_size - 1) / block_size);
+ /* Bound the signature's own memory (block_count * sizeof(DeltaBlockSig)). */
+ if (block_count == 0 || block_count > MAX_DELTA_BLOCKS)
+ return NULL;
+ DeltaSignature* sig = protocol_alloc(sizeof(DeltaSignature));
+ if (!sig)
+ return NULL;
+ sig->file_size = old_file_size;
+ sig->block_size = block_size;
+ sig->block_count = block_count;
+ sig->blocks = protocol_alloc((size_t)block_count * sizeof(DeltaBlockSig));
+ if (!sig->blocks) {
+ free(sig);
+ return NULL;
+ }
+ uint8_t* block = malloc(block_size);
+ if (!block) {
+ free(sig->blocks);
+ free(sig);
+ return NULL;
+ }
+ for (uint32_t i = 0; i < block_count; i++) {
+ uint64_t offset = (uint64_t)i * block_size;
+ uint32_t len =
+ (uint32_t)((old_file_size - offset < block_size) ? (old_file_size - offset) : block_size);
+ if (!pread_all(fd, block, len, offset)) {
+ free(block);
+ free(sig->blocks);
+ free(sig);
+ return NULL;
+ }
+ sig->blocks[i].adler32 = delta_adler32(block, len);
+ sig->blocks[i].xxhash = delta_xxhash32_seeded(block, len, seed);
+ }
+ free(block);
+ return sig;
+}
+
Data* delta_signature_serialize(const DeltaSignature* sig) {
if (!sig)
return NULL;
@@ -687,6 +747,130 @@ void* delta_apply(const void* old_data, uint64_t old_size, const Delta* delta,
return output;
}
+void* delta_apply_fd(int src_fd, uint64_t old_size, const Delta* delta, uint32_t block_size) {
+ if (!delta || block_size == 0 || block_size > DELTA_BLOCK_SIZE_MAX ||
+ (delta->instruction_count > 0 && !delta->instructions) || delta->new_file_size == 0 ||
+ delta->new_file_size > SIZE_MAX)
+ return NULL;
+
+ void* output = protocol_alloc((size_t)delta->new_file_size);
+ if (!output)
+ return NULL;
+
+ uint8_t* out = (uint8_t*)output;
+ uint64_t out_pos = 0;
+
+ for (uint32_t i = 0; i < delta->instruction_count; i++) {
+ if (delta->instructions[i].type == DELTA_INSTR_BLOCK_MATCH) {
+ uint64_t src_offset = (uint64_t)delta->instructions[i].match.block_index * block_size;
+ if (src_offset > UINT64_MAX - delta->instructions[i].match.block_offset) {
+ free(output);
+ return NULL;
+ }
+ src_offset += delta->instructions[i].match.block_offset;
+ uint32_t len = delta->instructions[i].match.length;
+
+ if (src_offset > old_size || (uint64_t)len > old_size - src_offset ||
+ out_pos > delta->new_file_size || (uint64_t)len > delta->new_file_size - out_pos) {
+ free(output);
+ return NULL;
+ }
+ if (!pread_all(src_fd, out + out_pos, len, src_offset)) {
+ free(output);
+ return NULL;
+ }
+ out_pos += len;
+ } else if (delta->instructions[i].type == DELTA_INSTR_LITERAL) {
+ uint32_t len = delta->instructions[i].literal.length;
+ if (out_pos > delta->new_file_size || (uint64_t)len > delta->new_file_size - out_pos) {
+ free(output);
+ return NULL;
+ }
+ memcpy(out + out_pos, delta->instructions[i].literal.data, len);
+ out_pos += len;
+ } else {
+ free(output);
+ return NULL;
+ }
+ }
+
+ if (out_pos != delta->new_file_size) {
+ free(output);
+ return NULL;
+ }
+ return output;
+}
+
+bool delta_apply_to_fd(const void* old_data, int src_fd, uint64_t old_size, const Delta* delta,
+ uint32_t block_size, int dst_fd) {
+ if (!delta || (old_data == NULL && src_fd < 0) || block_size == 0 ||
+ block_size > DELTA_BLOCK_SIZE_MAX || (delta->instruction_count > 0 && !delta->instructions))
+ return false;
+ const int chunk = 1 << 20;
+ uint8_t* buf = malloc((size_t)chunk);
+ if (!buf)
+ return false;
+ uint64_t out_pos = 0;
+ bool ok = true;
+ for (uint32_t i = 0; i < delta->instruction_count && ok; i++) {
+ uint64_t src_offset = 0;
+ uint64_t len = 0;
+ const uint8_t* lit = NULL;
+ if (delta->instructions[i].type == DELTA_INSTR_BLOCK_MATCH) {
+ src_offset = (uint64_t)delta->instructions[i].match.block_index * block_size;
+ if (src_offset > UINT64_MAX - delta->instructions[i].match.block_offset ||
+ src_offset + delta->instructions[i].match.block_offset > old_size) {
+ ok = false;
+ break;
+ }
+ src_offset += delta->instructions[i].match.block_offset;
+ len = delta->instructions[i].match.length;
+ if (len > old_size - src_offset) {
+ ok = false;
+ break;
+ }
+ } else if (delta->instructions[i].type == DELTA_INSTR_LITERAL) {
+ lit = delta->instructions[i].literal.data;
+ len = delta->instructions[i].literal.length;
+ } else {
+ ok = false;
+ break;
+ }
+ if (out_pos > delta->new_file_size || len > delta->new_file_size - out_pos) {
+ ok = false;
+ break;
+ }
+ uint64_t done = 0;
+ while (ok && done < len) {
+ size_t want = (len - done) < (uint64_t)chunk ? (size_t)(len - done) : (size_t)chunk;
+ if (lit) {
+ memcpy(buf, lit + done, want);
+ } else if (old_data) {
+ memcpy(buf, (const uint8_t*)old_data + src_offset + done, want);
+ } else if (!pread_all(src_fd, buf, want, src_offset + done)) {
+ ok = false;
+ break;
+ }
+ const uint8_t* p = buf;
+ size_t written = 0;
+ while (written < want) {
+ ssize_t n = write(dst_fd, p + written, want - written);
+ if (n < 0 && errno == EINTR)
+ continue;
+ if (n <= 0) {
+ ok = false;
+ break;
+ }
+ written += (size_t)n;
+ }
+ done += want;
+ }
+ out_pos += len;
+ }
+ free(buf);
+ return ok && out_pos == delta->new_file_size;
+}
+
void delta_destroy(Delta* delta) {
if (!delta)
return;
diff --git a/src/shared/delta.h b/src/shared/delta.h
index d395a90..58b9718 100644
--- a/src/shared/delta.h
+++ b/src/shared/delta.h
@@ -61,6 +61,13 @@ DeltaSignature* delta_signature_create(const void* old_file_data, uint64_t old_f
* identical to the unseeded function. */
DeltaSignature* delta_signature_create_seeded(const void* old_file_data, uint64_t old_file_size,
uint32_t block_size, uint32_t seed);
+/* Streaming equivalent of delta_signature_create_seeded: reads the basis blocks
+ * from `fd` in bounded chunks, so an arbitrarily large basis can be signed
+ * without materializing it. The signature itself is bounded (MAX_DELTA_BLOCKS
+ * entries); an over-large basis returns NULL and the caller falls back to a
+ * whole-file transfer. */
+DeltaSignature* delta_signature_create_fd_seeded(int fd, uint64_t old_file_size,
+ uint32_t block_size, uint32_t seed);
Data* delta_signature_serialize(const DeltaSignature* sig);
DeltaSignature* delta_signature_deserialize(const Data* data);
void delta_signature_destroy(DeltaSignature* sig);
@@ -75,6 +82,15 @@ Delta* delta_compute_seeded(const void* new_file_data, uint64_t new_file_size,
Data* delta_serialize(const Delta* delta);
Delta* delta_deserialize(const Data* data);
void* delta_apply(const void* old_data, uint64_t old_size, const Delta* delta, uint32_t block_size);
+/* Streaming equivalent of delta_apply: matched blocks are read from `src_fd` as
+ * they are emitted, so the basis never has to be resident. */
+void* delta_apply_fd(int src_fd, uint64_t old_size, const Delta* delta, uint32_t block_size);
+/* Fully streamed reconstruction: matched blocks come from `old_data` (when
+ * non-NULL) or are read from `src_fd`, and the reconstructed bytes are written
+ * straight to `dst_fd` in bounded chunks, so a reconstructed file larger than
+ * memory is never materialized. */
+bool delta_apply_to_fd(const void* old_data, int src_fd, uint64_t old_size, const Delta* delta,
+ uint32_t block_size, int dst_fd);
void delta_destroy(Delta* delta);
bool delta_should_attempt(uint64_t old_size, uint64_t new_size, uint64_t max_file_size);
diff --git a/src/shared/file.c b/src/shared/file.c
index de06c09..b9cafc8 100644
--- a/src/shared/file.c
+++ b/src/shared/file.c
@@ -16,6 +16,8 @@
#include "data.h"
#include "checksum.h"
+#include "chunk.h"
+#include "compression.h"
#include "delta.h"
#include "file.h"
#include "file_store.h"
@@ -209,6 +211,7 @@ File* file_create(const char* path) {
file->dir_time_only = false;
file->basis_link = NULL;
file->basis_copy = NULL;
+ file->data_spool = false;
file->link_group = 0;
file->link_first = false;
file->hardlink_target = NULL;
@@ -238,8 +241,11 @@ void file_destroy(void* item) {
file->send_path = NULL;
free(file->basis_link);
file->basis_link = NULL;
+ if (file->data_spool && file->basis_copy)
+ unlink(file->basis_copy);
free(file->basis_copy);
file->basis_copy = NULL;
+ file->data_spool = false;
free(file->hardlink_target);
file->hardlink_target = NULL;
free(file->symlink_target);
@@ -381,6 +387,261 @@ size_t file_content_to_buffer(File* file) {
return bytes_read;
}
+/* ---- Streamed whole-file payload receive ----
+ *
+ * A whole-file data frame is a uint64 length followed by that many bytes. When
+ * the logical payload is at or below the receiver's streaming bound the
+ * historical charged whole-buffer path is kept (decompressing in one shot for
+ * `-z`). Above the bound the frame is read in bounded chunks and written into
+ * a spool temp file next to the destination, decompressing incrementally for
+ * zstd/zlib (lz4's block format cannot be streamed and keeps the buffered path).
+ * The spool is then installed by the existing bounded-buffer basis-copy path
+ * (file_copy_basis_stream_attrs), so the atomic temp+rename store, --partial/
+ * --partial-dir, --delay-updates, --preallocate and metadata/xattr application
+ * are all reused unchanged. Only bounded buffers (64 KiB read chunk + the
+ * decompressor's 256 KiB output window) are ever live. */
+
+#define FILE_PAYLOAD_READ_CHUNK (64 * 1024)
+#define FILE_PAYLOAD_PEEK_MAX 32
+
+/* Create a confined spool temp file in the destination's directory. Returns an
+ * open write fd and an owned absolute path, or -1 (errno set). The parent walk
+ * creates missing directories exactly as a normal store would. */
+static int file_spool_create(const char* dest_path, char** out_spool_path) {
+ *out_spool_path = NULL;
+ char* leaf = NULL;
+ int dirfd = file_open_secure_parent(dest_path, &leaf, true);
+ free(leaf);
+ if (dirfd < 0)
+ return -1;
+ char name[64];
+ for (unsigned int i = 0; i < 100; i++) {
+ snprintf(name, sizeof(name), ".fastsync-spool.%ld.%llu", (long)getpid(), next_temp_sequence());
+ int fd = openat(dirfd, name, O_WRONLY | O_CREAT | O_EXCL | O_CLOEXEC | O_NOFOLLOW, 0600);
+ if (fd < 0) {
+ if (errno != EEXIST)
+ break;
+ continue;
+ }
+ char* copy = str_dup(dest_path);
+ char* dir = copy ? dirname(copy) : NULL;
+ size_t need = dir ? strlen(dir) + 1 + strlen(name) + 1 : 0;
+ char* full = need ? malloc(need) : NULL;
+ if (!full) {
+ free(copy);
+ close(fd);
+ unlinkat(dirfd, name, 0);
+ close(dirfd);
+ return -1;
+ }
+ snprintf(full, need, "%s/%s", dir, name);
+ free(copy);
+ close(dirfd);
+ *out_spool_path = full;
+ return fd;
+ }
+ close(dirfd);
+ return -1;
+}
+
+int file_spool_for_payload(const char* dest_path, char** out_spool_path) {
+ return file_spool_create(dest_path, out_spool_path);
+}
+
+bool file_receive_payload(int fd, bool compress, unsigned long long expected_size,
+ const char* dest_path, unsigned long long stream_limit, Data** out_buffer,
+ char** out_spool, unsigned long long* out_size) {
+ if (out_buffer)
+ *out_buffer = NULL;
+ if (out_spool)
+ *out_spool = NULL;
+ if (out_size)
+ *out_size = 0;
+ if (!out_buffer || !out_spool || !out_size || !dest_path)
+ return false;
+
+ unsigned long long frame_size = 0;
+ if (!receive_n_data(fd, &frame_size, sizeof(frame_size)))
+ return false;
+ /* A frame larger than MAX_DATA_PAYLOAD_SIZE is allowed only on the streamed
+ path, which never materializes it; the buffered path's
+ receive_data_alloc/body enforces the historical bound itself. Only a size
+ unrepresentable on this platform is rejected here. */
+ if (frame_size > SIZE_MAX) {
+ log_message(LOG_LEVEL_ERROR, "Data size %llu is not representable", frame_size);
+ return false;
+ }
+
+ unsigned char prefix[FILE_PAYLOAD_PEEK_MAX];
+ size_t prefix_len = 0;
+ bool stream;
+ if (expected_size != 0) {
+ /* The check frame already told us the logical size: no need to peek. */
+ stream = expected_size > stream_limit;
+ } else if (!compress) {
+ stream = frame_size > stream_limit;
+ } else {
+ /* Non-incremental compressed frame with unknown logical size: peek the
+ header to decide, so a small compressed frame that expands past the bound
+ still streams instead of allocating the whole logical image. */
+ size_t want = frame_size < sizeof(prefix) ? (size_t)frame_size : sizeof(prefix);
+ if (want > 0 && !receive_n_data(fd, prefix, want))
+ return false;
+ prefix_len = want;
+ unsigned long long logical = compression_peek_frame_content_size(prefix, prefix_len);
+ stream = logical != 0 ? logical > stream_limit : frame_size > stream_limit;
+ }
+
+ if (!stream) {
+ Data* frame;
+ if (prefix_len > 0) {
+ frame = receive_data_alloc(fd, frame_size);
+ if (!frame)
+ return false;
+ memcpy(frame->data, prefix, prefix_len);
+ if (frame_size > prefix_len &&
+ !receive_n_data(fd, (char*)frame->data + prefix_len, (size_t)(frame_size - prefix_len))) {
+ data_destroy(frame);
+ return false;
+ }
+ } else {
+ frame = receive_data_body(fd, frame_size);
+ if (!frame)
+ return false;
+ }
+ if (compress) {
+ unsigned long long bound =
+ expected_size != 0 ? expected_size : (unsigned long long)MAX_RECEIVE_WHOLE_FILE_SIZE;
+ Data* uncompressed = data_decompress_limited(frame, (size_t)bound);
+ ProtocolSession* owner = frame->owner;
+ data_destroy(frame);
+ if (!uncompressed)
+ return false;
+ if (expected_size != 0 && uncompressed->size != expected_size) {
+ data_destroy(uncompressed);
+ return false;
+ }
+ if (!data_charge_session(uncompressed, owner, uncompressed->size)) {
+ data_destroy(uncompressed);
+ return false;
+ }
+ *out_buffer = uncompressed;
+ *out_size = uncompressed->size;
+ return true;
+ }
+ *out_buffer = frame;
+ *out_size = frame->size;
+ return true;
+ }
+
+ /* Streamed path. */
+ int spool_fd = file_spool_create(dest_path, out_spool);
+ if (spool_fd < 0)
+ return false;
+
+ CompressionStreamDecompressor* dec = NULL;
+ size_t header_len = 0;
+ unsigned long long learned_size = expected_size;
+ if (compress) {
+ unsigned char hdr[1 + 4];
+ size_t hdr_have = 0;
+ if (prefix_len >= 1) {
+ hdr[0] = prefix[0];
+ hdr_have = 1;
+ } else {
+ if (!receive_n_data(fd, hdr, 1))
+ goto stream_fail;
+ hdr_have = 1;
+ }
+ CompressionAlgo algo = (CompressionAlgo)hdr[0];
+ if (!compression_algo_valid((int)algo) || algo == COMPRESSION_ALGO_LZ4)
+ goto stream_fail;
+ header_len = (algo == COMPRESSION_ALGO_ZLIB || algo == COMPRESSION_ALGO_ZLIBX) ? 5 : 1;
+ while (hdr_have < header_len) {
+ size_t need = header_len - hdr_have;
+ if (prefix_len > hdr_have) {
+ size_t avail = prefix_len - hdr_have;
+ size_t take = avail < need ? avail : need;
+ memcpy(hdr + hdr_have, prefix + hdr_have, take);
+ hdr_have += take;
+ } else {
+ if (!receive_n_data(fd, hdr + hdr_have, need))
+ goto stream_fail;
+ hdr_have = header_len;
+ }
+ }
+ if (header_len == 5) {
+ uint32_t raw = 0;
+ for (int i = 0; i < 4; i++)
+ raw |= (uint32_t)hdr[1 + i] << (8 * i);
+ if (expected_size != 0 && raw != expected_size)
+ goto stream_fail;
+ if (learned_size == 0)
+ learned_size = raw;
+ }
+ dec = compression_stream_decompressor_create(algo, learned_size);
+ if (!dec)
+ goto stream_fail;
+ }
+
+ if (frame_size < header_len)
+ goto stream_fail;
+ {
+ unsigned long long body_remaining = frame_size - header_len;
+ if (prefix_len > header_len) {
+ size_t avail = prefix_len - header_len;
+ if (compress) {
+ bool done = false;
+ if (!compression_stream_decompressor_feed(dec, prefix + header_len, avail, spool_fd, &done))
+ goto stream_fail;
+ } else if (!write_all(spool_fd, prefix + header_len, avail)) {
+ goto stream_fail;
+ }
+ body_remaining -= avail;
+ }
+ unsigned char buf[FILE_PAYLOAD_READ_CHUNK];
+ while (body_remaining > 0) {
+ size_t want = body_remaining < sizeof(buf) ? (size_t)body_remaining : (size_t)sizeof(buf);
+ if (!receive_n_data(fd, buf, want))
+ goto stream_fail;
+ if (compress) {
+ bool done = false;
+ if (!compression_stream_decompressor_feed(dec, buf, want, spool_fd, &done))
+ goto stream_fail;
+ } else if (!write_all(spool_fd, buf, want)) {
+ goto stream_fail;
+ }
+ body_remaining -= want;
+ }
+ }
+
+ {
+ unsigned long long total = compress ? compression_stream_decompressor_total(dec) : frame_size;
+ if (learned_size != 0 && total != learned_size)
+ goto stream_fail;
+ *out_size = total;
+ }
+ if (dec)
+ compression_stream_decompressor_destroy(dec);
+ if (close(spool_fd) != 0) {
+ spool_fd = -1;
+ goto stream_fail;
+ }
+ return true;
+
+stream_fail:
+ if (dec)
+ compression_stream_decompressor_destroy(dec);
+ if (spool_fd >= 0)
+ close(spool_fd);
+ if (*out_spool) {
+ unlink(*out_spool);
+ free(*out_spool);
+ *out_spool = NULL;
+ }
+ return false;
+}
+
/* ---- Secure filesystem primitives ---- */
bool file_path_exists_secure(const char* path) {
diff --git a/src/shared/file.h b/src/shared/file.h
index dc02584..f573cf9 100644
--- a/src/shared/file.h
+++ b/src/shared/file.h
@@ -177,6 +177,24 @@ bool file_copy_basis_stream_attrs(const char* path, const char* basis_path,
const FileMetadata* metadata, FileAttrPolicy policy, bool update,
bool use_fsync, const FileXattrList* xattrs, bool fake_super,
const char* temp_dir);
+/* Receive one length-prefixed whole-file data frame, streaming the payload
+ * through a bounded buffer when it (or its known logical size) exceeds
+ * `stream_limit`. On success exactly one of *out_buffer / *out_spool is set:
+ * - *out_buffer: the historical charged whole-buffer Data (caller destroys);
+ * - *out_spool: an owned temp path holding the payload, installed through the
+ * File's basis_copy field with File.data_spool set so file_destroy unlinks
+ * it. The destination policy/metadata/atomic-store handling is then the
+ * existing bounded-buffer basis install (file_copy_basis_stream_attrs).
+ * `expected_size` (0 = unknown) is the logical size from the check frame;
+ * `dest_path` locates the spool next to the destination; `compress` selects
+ * incremental decompression. Returns false on any framing/I/O/size error. */
+bool file_receive_payload(int fd, bool compress, unsigned long long expected_size,
+ const char* dest_path, unsigned long long stream_limit, Data** out_buffer,
+ char** out_spool, unsigned long long* out_size);
+/* Create a confined spool temp file next to `dest_path`; returns an open write
+ * fd and an owned absolute path (to be installed via File.basis_copy with
+ * File.data_spool set, and unlinked by file_destroy). */
+int file_spool_for_payload(const char* dest_path, char** out_spool_path);
/* Protocol 2.28.0 receiver-stat variants: like the two above but additionally
* report through `dirs_created` (when non-NULL) how many parent directories the
* confined secure walk had to create that lie strictly below `count_floor` (a
diff --git a/src/shared/file_receive.c b/src/shared/file_receive.c
index be49127..763b221 100644
--- a/src/shared/file_receive.c
+++ b/src/shared/file_receive.c
@@ -58,36 +58,41 @@ File* file_receive(const Config* config, int file_descriptor) {
file_destroy(file);
return NULL;
}
- Data* file_data = receive_data_limited(file_descriptor, MAX_RECEIVE_WHOLE_FILE_SIZE);
- if (file_data == NULL) {
+ bool compress =
+ config->use_compression && !compression_should_skip_with_suffixes(
+ file->path, config->skip_compress_suffixes,
+ config->skip_compress_set ? config->skip_compress_count : -1);
+ char* dest_path = path_cat(config->receive_root_directory, file->path);
+ if (!dest_path) {
file_destroy(file);
return NULL;
}
- if (config->use_compression &&
- !compression_should_skip_with_suffixes(file->path, config->skip_compress_suffixes,
- config->skip_compress_set ? config->skip_compress_count
- : -1)) {
- Data* file_data_uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_WHOLE_FILE_SIZE);
- ProtocolSession* owner = file_data->owner;
- data_destroy(file_data);
- if (file_data_uncompressed == NULL) {
+ Data* buffer = NULL;
+ char* spool = NULL;
+ unsigned long long size = 0;
+ bool ok = file_receive_payload(file_descriptor, compress, 0, dest_path,
+ protocol_whole_file_receive_limit(), &buffer, &spool, &size);
+ free(dest_path);
+ if (!ok) {
+ file_destroy(file);
+ return NULL;
+ }
+ if (spool) {
+ Data* reserved = data_create_reserve((size_t)size);
+ if (!reserved) {
+ unlink(spool);
+ free(spool);
file_destroy(file);
return NULL;
}
- if (!data_charge_session(file_data_uncompressed, owner, file_data_uncompressed->size)) {
- data_destroy(file_data_uncompressed);
- file_destroy(file);
- return NULL;
- }
- if (file_data_uncompressed->size > MAX_FILE_DATA_SIZE) {
- data_destroy(file_data_uncompressed);
- file_destroy(file);
- return NULL;
- }
- file_data = file_data_uncompressed;
+ data_destroy(file->data);
+ file->data = reserved;
+ file->basis_copy = spool;
+ file->data_spool = true;
+ } else {
+ data_destroy(file->data);
+ file->data = buffer;
}
- data_destroy(file->data);
- file->data = file_data;
return file;
}
diff --git a/src/shared/file_send.c b/src/shared/file_send.c
index 70f0e95..f62af4b 100644
--- a/src/shared/file_send.c
+++ b/src/shared/file_send.c
@@ -44,12 +44,138 @@ bool file_send_single_calls(File* file, int file_descriptor, bool use_metadata,
send_path, NULL, -1, 0, false);
}
+/* Stream-compress a whole source file into a temp file, then transmit it as the
+ * normal length-prefixed data frame. Used when the source was too large to
+ * load (File.data.data == NULL): the source is read in bounded chunks through a
+ * streaming codec, so neither the raw nor the compressed image is held in
+ * memory. The compressed length must be known before the frame is sent (the
+ * wire is length-prefixed), so the stream lands in a private temp file first. */
+bool file_send_compressed_stream_with_skip(File* file, int file_descriptor, bool use_metadata,
+ int compression_level, bool send_path,
+ char* const* skip_suffixes, int skip_count,
+ int compression_threads, bool send_xattrs) {
+ /* The caller has already decided this file compresses; the skip list is part
+ of the public signature for symmetry with the buffered path. */
+ (void)skip_suffixes;
+ (void)skip_count;
+ if (!file || !file->path || !file->data || file->data->size == 0 || file->data->data != NULL)
+ return false;
+ CompressionAlgo algo = compression_get_algo();
+ CompressionStreamCompressor* compressor =
+ compression_stream_compressor_create(algo, compression_level, compression_threads);
+ if (!compressor) {
+ log_message(LOG_LEVEL_ERROR, "streaming compression is not available for this codec; "
+ "the source was not loaded for the buffered path");
+ return false;
+ }
+
+ int src = file_open_for_read(file->path);
+ if (src < 0) {
+ compression_stream_compressor_destroy(compressor);
+ return false;
+ }
+ struct stat src_st;
+ if (fstat(src, &src_st) != 0 || !S_ISREG(src_st.st_mode) ||
+ (unsigned long long)src_st.st_size < file->data->size) {
+ close(src);
+ compression_stream_compressor_destroy(compressor);
+ return false;
+ }
+
+ const char* tmpdir = getenv("TMPDIR");
+ if (!tmpdir || tmpdir[0] == '\0')
+ tmpdir = "/tmp";
+ size_t tmplen = strlen(tmpdir) + strlen("/fastsync-z-XXXXXX") + 1;
+ char* tmpl = malloc(tmplen);
+ if (!tmpl) {
+ close(src);
+ compression_stream_compressor_destroy(compressor);
+ return false;
+ }
+ snprintf(tmpl, tmplen, "%s/fastsync-z-XXXXXX", tmpdir);
+ int tmp_fd = mkstemp(tmpl);
+ if (tmp_fd < 0) {
+ log_perror("Could not create compression temp file");
+ free(tmpl);
+ close(src);
+ compression_stream_compressor_destroy(compressor);
+ return false;
+ }
+ unlink(tmpl);
+ free(tmpl);
+
+ bool ok = compression_stream_compressor_begin(compressor, file->data->size, tmp_fd);
+ unsigned char buf[64 * 1024];
+ unsigned long long remaining = file->data->size;
+ while (ok && remaining > 0) {
+ size_t want = remaining < sizeof(buf) ? (size_t)remaining : sizeof(buf);
+ ssize_t got = read(src, buf, want);
+ if (got <= 0) {
+ ok = false;
+ break;
+ }
+ if (!compression_stream_compressor_feed(compressor, buf, (size_t)got, tmp_fd))
+ ok = false;
+ remaining -= (unsigned long long)got;
+ }
+ if (ok)
+ ok = compression_stream_compressor_finish(compressor, tmp_fd);
+ close(src);
+ compression_stream_compressor_destroy(compressor);
+ if (!ok) {
+ close(tmp_fd);
+ return false;
+ }
+ struct stat tmp_st;
+ if (fstat(tmp_fd, &tmp_st) != 0 || tmp_st.st_size < 0) {
+ close(tmp_fd);
+ return false;
+ }
+ unsigned long long compressed_size = (unsigned long long)tmp_st.st_size;
+ if (lseek(tmp_fd, 0, SEEK_SET) == (off_t)-1) {
+ close(tmp_fd);
+ return false;
+ }
+
+ ok = true;
+ if (send_path && !send_wire_str(file_descriptor, file_wire_path(file)))
+ ok = false;
+ if (ok && use_metadata && !metadata_send(file_descriptor, file->metadata))
+ ok = false;
+ if (ok && send_xattrs && !xattr_send(file_descriptor, file ? file->xattrs : NULL))
+ ok = false;
+ if (ok && !send_n_data(file_descriptor, &compressed_size, sizeof(compressed_size)))
+ ok = false;
+ unsigned long long left = compressed_size;
+ while (ok && left > 0) {
+ size_t want = left < sizeof(buf) ? (size_t)left : sizeof(buf);
+ ssize_t got = read(tmp_fd, buf, want);
+ if (got <= 0 || !send_n_data(file_descriptor, buf, (size_t)got)) {
+ ok = false;
+ break;
+ }
+ left -= (unsigned long long)got;
+ }
+ close(tmp_fd);
+ return ok;
+}
+
bool file_send_single_calls_with_skip(File* file, int file_descriptor, bool use_metadata,
int compression_level, bool send_path,
char* const* skip_suffixes, int skip_count,
int compression_threads, bool send_xattrs) {
- if (!file || !file->path || !file->data || (file->data->size != 0 && !file->data->data))
+ if (!file || !file->path || !file->data)
return false;
+ /* A source too large to load is streamed: compression streams through a
+ * temp file, no compression must have taken the sendfile path instead. */
+ if (file->data->size != 0 && file->data->data == NULL) {
+ if (compression_level <= 0 ||
+ compression_should_skip_with_suffixes(file->path, skip_suffixes, skip_count))
+ return false;
+ return file_send_compressed_stream_with_skip(file, file_descriptor, use_metadata,
+ compression_level, send_path, skip_suffixes,
+ skip_count, compression_threads, send_xattrs);
+ }
const Data* data_to_send = file->data;
Data* compressed_data = NULL;
if (compression_level > 0 &&
diff --git a/src/shared/file_send.h b/src/shared/file_send.h
index b48d33c..3d919f0 100644
--- a/src/shared/file_send.h
+++ b/src/shared/file_send.h
@@ -13,6 +13,12 @@ bool file_send_single_calls_with_skip(File* file, int file_descriptor, bool use_
int compression_level, bool send_path,
char* const* skip_suffixes, int skip_count,
int compression_threads, bool send_xattrs);
+/* Stream-compress an unloaded whole source file into a temp file and send it as
+ * the usual data frame (see file_send.c). */
+bool file_send_compressed_stream_with_skip(File* file, int file_descriptor, bool use_metadata,
+ int compression_level, bool send_path,
+ char* const* skip_suffixes, int skip_count,
+ int compression_threads, bool send_xattrs);
bool file_send_sendfile(File* file, int file_descriptor, bool use_metadata, int compression_level,
bool send_path);
bool file_send_sendfile_with_skip(File* file, int file_descriptor, bool use_metadata,
diff --git a/src/shared/file_types.h b/src/shared/file_types.h
index 1a3bd82..4383bab 100644
--- a/src/shared/file_types.h
+++ b/src/shared/file_types.h
@@ -62,6 +62,11 @@ typedef struct {
* This lets a basis larger than any whole-file bound materialize without
* buffering it; the source metadata on `metadata` is applied afterwards. */
char* basis_copy;
+ /* Receiver-only. When set, `basis_copy` points at a receiver-created spool
+ * temp file holding a STREAMED whole-file payload (rather than a --copy-dest
+ * basis). file_destroy unlinks it after the install consumes it, so an
+ * over-limit file leaves no scratch behind. */
+ bool data_spool;
/* --hard-links (-H), sender + receiver wire state. link_group is a run-local
* id shared by every member of one source inode (0 = not part of a group).
* The FIRST member (link_first == true) carries its data on the wire and is
diff --git a/src/shared/incremental_check.c b/src/shared/incremental_check.c
index 98d4eb0..8fd4bbf 100644
--- a/src/shared/incremental_check.c
+++ b/src/shared/incremental_check.c
@@ -44,16 +44,26 @@ bool receive_file_xattrs(File* file, int fd, const Config* config) {
return true;
}
+static bool receive_file_payload_into(File* file, int fd, const Config* config,
+ const char* dest_path, unsigned long long expected_size);
+
static File* receive_delta_file(int fd, const Config* config, const char* check_path,
- void* old_data, unsigned long long old_size, bool* failed) {
- if (!old_data) {
- free(old_data); /* defensive: old_data is always non-NULL today */
+ void* old_data, unsigned long long old_size,
+ unsigned long long expected_size, int basis_fd, bool* failed) {
+ if (!old_data && basis_fd < 0) {
+ free(old_data); /* defensive: a basis source is always provided today */
*failed = true;
return NULL;
}
- DeltaSignature* sig = delta_signature_create_seeded(old_data, old_size, config->delta_block_size,
- (uint32_t)config->checksum_seed);
+ /* The basis is either an in-memory snapshot (the destination file, bounded) or
+ * a confined descriptor (a --fuzzy sibling, possibly larger than memory) that
+ * is signed/applied in bounded chunks. */
+ DeltaSignature* sig =
+ old_data ? delta_signature_create_seeded(old_data, old_size, config->delta_block_size,
+ (uint32_t)config->checksum_seed)
+ : delta_signature_create_fd_seeded(basis_fd, old_size, config->delta_block_size,
+ (uint32_t)config->checksum_seed);
if (!sig) {
free(old_data);
*failed = true;
@@ -130,7 +140,7 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_
}
uint64_t new_size = delta->new_file_size;
- if (new_size > MAX_RECEIVE_WHOLE_FILE_SIZE || new_size > SIZE_MAX) {
+ if (new_size > SIZE_MAX) {
delta_destroy(delta);
free(old_data);
delta_signature_destroy(sig);
@@ -149,19 +159,57 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_
else if (delta->instructions[k].type == DELTA_INSTR_LITERAL)
literal += delta->instructions[k].literal.length;
}
- void* new_data = delta_apply(old_data, old_size, delta, config->delta_block_size);
- delta_destroy(delta);
- if (!new_data) {
- free(old_data);
- delta_signature_destroy(sig);
- *failed = true;
- return NULL;
+ /* A reconstructed file above the streaming bound is written into a spool
+ temp file through delta_apply_to_fd; a smaller one keeps the historical
+ in-memory reconstruction. */
+ void* new_data = NULL;
+ char* spool = NULL;
+ if (new_size > protocol_whole_file_receive_limit()) {
+ char* dest_path = path_cat(config->receive_root_directory, check_path);
+ int spool_fd = dest_path ? file_spool_for_payload(dest_path, &spool) : -1;
+ free(dest_path);
+ if (spool_fd < 0) {
+ delta_destroy(delta);
+ free(old_data);
+ delta_signature_destroy(sig);
+ send_status(fd, STATUS_ERROR);
+ *failed = true;
+ return NULL;
+ }
+ bool applied = delta_apply_to_fd(old_data, basis_fd, old_size, delta,
+ config->delta_block_size, spool_fd);
+ if (close(spool_fd) != 0)
+ applied = false;
+ delta_destroy(delta);
+ if (!applied) {
+ unlink(spool);
+ free(spool);
+ free(old_data);
+ delta_signature_destroy(sig);
+ send_status(fd, STATUS_ERROR);
+ *failed = true;
+ return NULL;
+ }
+ } else {
+ new_data = old_data ? delta_apply(old_data, old_size, delta, config->delta_block_size)
+ : delta_apply_fd(basis_fd, old_size, delta, config->delta_block_size);
+ delta_destroy(delta);
+ if (!new_data) {
+ free(old_data);
+ delta_signature_destroy(sig);
+ *failed = true;
+ return NULL;
+ }
}
File* file = file_create(check_path);
if (!file) {
free(new_data);
+ if (spool) {
+ unlink(spool);
+ free(spool);
+ }
free(old_data);
delta_signature_destroy(sig);
*failed = true;
@@ -176,6 +224,10 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_
if (!meta_ok) {
file_destroy(file);
free(new_data);
+ if (spool) {
+ unlink(spool);
+ free(spool);
+ }
free(old_data);
delta_signature_destroy(sig);
*failed = true;
@@ -185,23 +237,45 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_
if (!receive_file_xattrs(file, fd, config)) {
file_destroy(file);
free(new_data);
+ if (spool) {
+ unlink(spool);
+ free(spool);
+ }
free(old_data);
delta_signature_destroy(sig);
*failed = true;
return NULL;
}
- Data* replacement = data_create(new_data, (size_t)new_size);
- if (replacement == NULL) {
- file_destroy(file);
- free(old_data);
- delta_signature_destroy(sig);
- send_status(fd, STATUS_ERROR);
- *failed = true;
- return NULL;
+ if (spool) {
+ Data* reserved = data_create_reserve((size_t)new_size);
+ if (reserved == NULL) {
+ unlink(spool);
+ free(spool);
+ file_destroy(file);
+ free(old_data);
+ delta_signature_destroy(sig);
+ send_status(fd, STATUS_ERROR);
+ *failed = true;
+ return NULL;
+ }
+ data_destroy(file->data);
+ file->data = reserved;
+ file->basis_copy = spool;
+ file->data_spool = true;
+ } else {
+ Data* replacement = data_create(new_data, (size_t)new_size);
+ if (replacement == NULL) {
+ file_destroy(file);
+ free(old_data);
+ delta_signature_destroy(sig);
+ send_status(fd, STATUS_ERROR);
+ *failed = true;
+ return NULL;
+ }
+ data_destroy(file->data);
+ file->data = replacement;
}
- data_destroy(file->data);
- file->data = replacement;
free(old_data);
delta_signature_destroy(sig);
@@ -233,44 +307,19 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_
return NULL;
}
- Data* file_data = receive_data_limited(fd, MAX_RECEIVE_WHOLE_FILE_SIZE);
- if (file_data == NULL) {
+ char* dest_path = path_cat(config->receive_root_directory, check_path);
+ if (!dest_path) {
file_destroy(file);
*failed = true;
return NULL;
}
-
- if (config->use_compression &&
- !compression_should_skip_with_suffixes(
- file->path, config->skip_compress_suffixes,
- config->skip_compress_set ? config->skip_compress_count : -1)) {
- Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_WHOLE_FILE_SIZE);
- ProtocolSession* owner = file_data->owner;
- data_destroy(file_data);
- if (uncompressed == NULL) {
- file_destroy(file);
- *failed = true;
- return NULL;
- }
- if (!data_charge_session(uncompressed, owner, uncompressed->size)) {
- data_destroy(uncompressed);
- file_destroy(file);
- send_status(fd, STATUS_ERROR);
- *failed = true;
- return NULL;
- }
- if (uncompressed->size > MAX_FILE_DATA_SIZE) {
- data_destroy(uncompressed);
- file_destroy(file);
- send_status(fd, STATUS_ERROR);
- *failed = true;
- return NULL;
- }
- file_data = uncompressed;
+ bool payload_ok = receive_file_payload_into(file, fd, config, dest_path, expected_size);
+ free(dest_path);
+ if (!payload_ok) {
+ file_destroy(file);
+ *failed = true;
+ return NULL;
}
-
- data_destroy(file->data);
- file->data = file_data;
return file;
}
@@ -640,28 +689,29 @@ static bool fuzzy_candidate_better(const FuzzyCandidate* cand, const FuzzyCandid
return strcmp(cand->name, best->name) < 0;
}
-/* Search the destination directory that will contain `check_path` for a
- * similar regular file usable as a --fuzzy delta basis and return its full
- * content in a malloc'd (protocol_alloc) buffer. Returns NULL (with *out_size
- * = 0) when no candidate qualifies, which means the caller performs the normal
- * whole-file transfer. */
-static void* fuzzy_basis_find_and_load(const Config* config, const char* check_path,
- unsigned long long check_size, time_t check_mtime,
- long check_mtime_nsec, unsigned long long* out_size) {
+/* Search the destination directory that will contain `check_path` for a similar
+ * regular file usable as a --fuzzy delta basis and return an open, confined
+ * read descriptor to it (with *out_size set). Returns -1 (with *out_size 0)
+ * when no candidate qualifies, which means the caller performs the normal
+ * whole-file transfer. The basis is signed/applied by streaming its descriptor,
+ * so no whole-basis buffer is ever needed and its size is not capped. */
+static int fuzzy_basis_find_and_open(const Config* config, const char* check_path,
+ unsigned long long check_size, time_t check_mtime,
+ long check_mtime_nsec, unsigned long long* out_size) {
*out_size = 0;
if (!config || !config->receive_root_directory || !config->fuzzy || !config->use_delta ||
- !check_path || check_size > MAX_RECEIVE_WHOLE_FILE_SIZE)
- return NULL;
+ !check_path)
+ return -1;
char* full_path = path_cat(config->receive_root_directory, check_path);
if (!full_path)
- return NULL;
+ return -1;
char* leaf = NULL;
int dir_fd = file_open_secure_parent(full_path, &leaf, false);
if (dir_fd < 0 || !leaf) {
free(leaf);
free(full_path);
- return NULL;
+ return -1;
}
size_t target_len = strlen(leaf);
/* A target basename longer than FUZZY_NAME_LIMIT can never pass the name gate
@@ -670,7 +720,7 @@ static void* fuzzy_basis_find_and_load(const Config* config, const char* check_p
close(dir_fd);
free(leaf);
free(full_path);
- return NULL;
+ return -1;
}
int scanfd = dup(dir_fd);
@@ -678,7 +728,7 @@ static void* fuzzy_basis_find_and_load(const Config* config, const char* check_p
close(dir_fd);
free(leaf);
free(full_path);
- return NULL;
+ return -1;
}
DIR* dir = fdopendir(scanfd);
if (!dir) {
@@ -686,7 +736,7 @@ static void* fuzzy_basis_find_and_load(const Config* config, const char* check_p
close(dir_fd);
free(leaf);
free(full_path);
- return NULL;
+ return -1;
}
/* The weighted-distance scratch row is allocated once per scan (not once per
@@ -697,7 +747,7 @@ static void* fuzzy_basis_find_and_load(const Config* config, const char* check_p
close(dir_fd);
free(leaf);
free(full_path);
- return NULL;
+ return -1;
}
int fname_suf_len = 0;
const char* fname_suf = fuzzy_find_suffix(leaf, (int)target_len, &fname_suf_len);
@@ -727,7 +777,7 @@ static void* fuzzy_basis_find_and_load(const Config* config, const char* check_p
if (fstatat(dir_fd, name, &st, AT_SYMLINK_NOFOLLOW) != 0 || !S_ISREG(st.st_mode))
continue;
unsigned long long cand_size = (unsigned long long)st.st_size;
- if (cand_size == 0 || cand_size > MAX_RECEIVE_WHOLE_FILE_SIZE)
+ if (cand_size == 0)
continue;
long cand_nsec = 0;
#ifdef __linux__
@@ -771,7 +821,7 @@ static void* fuzzy_basis_find_and_load(const Config* config, const char* check_p
if (exact.name[0])
best = exact;
- void* basis = NULL;
+ int basis_fd = -1;
if (best.name[0]) {
/* O_NONBLOCK: a name raced to a FIFO between the fstatat gate and this open
would otherwise block the receive thread forever on open(2); with it the
@@ -781,29 +831,55 @@ static void* fuzzy_basis_find_and_load(const Config* config, const char* check_p
if (fd >= 0) {
struct stat st;
if (fstat(fd, &st) == 0 && S_ISREG(st.st_mode) &&
- (unsigned long long)st.st_size == best.size && best.size <= SIZE_MAX) {
- basis = protocol_alloc((size_t)best.size);
- if (basis) {
- size_t got = 0;
- while (got < (size_t)best.size) {
- ssize_t n = read(fd, (char*)basis + got, (size_t)best.size - got);
- if (n <= 0) {
- free(basis);
- basis = NULL;
- break;
- }
- got += (size_t)n;
- }
- }
+ (unsigned long long)st.st_size == best.size) {
+ basis_fd = fd;
+ } else {
+ close(fd);
}
- close(fd);
}
}
close(dir_fd);
free(full_path);
- if (basis)
+ if (basis_fd >= 0)
*out_size = best.size;
- return basis;
+ return basis_fd;
+}
+
+/* Receive one whole-file data frame into `file`. A payload at or below the
+ * receiver's streaming bound keeps the historical charged whole-buffer path; a
+ * larger one is streamed into a spool temp file (decompressing incrementally)
+ * and installed through the File's basis_copy field. `expected_size` is the
+ * logical size from the check frame (0 when unknown, e.g. the non-incremental
+ * path). */
+static bool receive_file_payload_into(File* file, int fd, const Config* config,
+ const char* dest_path, unsigned long long expected_size) {
+ bool compress =
+ config->use_compression && !compression_should_skip_with_suffixes(
+ file->path, config->skip_compress_suffixes,
+ config->skip_compress_set ? config->skip_compress_count : -1);
+ Data* buffer = NULL;
+ char* spool = NULL;
+ unsigned long long size = 0;
+ if (!file_receive_payload(fd, compress, expected_size, dest_path,
+ protocol_whole_file_receive_limit(), &buffer, &spool, &size)) {
+ return false;
+ }
+ if (spool) {
+ Data* reserved = data_create_reserve((size_t)size);
+ if (!reserved) {
+ unlink(spool);
+ free(spool);
+ return false;
+ }
+ data_destroy(file->data);
+ file->data = reserved;
+ file->basis_copy = spool;
+ file->data_spool = true;
+ } else {
+ data_destroy(file->data);
+ file->data = buffer;
+ }
+ return true;
}
/* Read the remainder of a full-file transfer after the receiver has already
@@ -811,7 +887,8 @@ static void* fuzzy_basis_find_and_load(const Config* config, const char* check_p
* data frame, and return an owned File. Shared by the plain full-transfer path
* and the --append-verify prefix-mismatch fallback (a clean full transfer
* instead of a corrupt prefix+tail blend). */
-static File* receive_full_file(int fd, const Config* config, const char* path) {
+static File* receive_full_file(int fd, const Config* config, const char* path,
+ unsigned long long expected_size) {
File* file = file_create(path);
if (!file)
return NULL;
@@ -827,36 +904,17 @@ static File* receive_full_file(int fd, const Config* config, const char* path) {
file_destroy(file);
return NULL;
}
- Data* file_data = receive_data_limited(fd, MAX_RECEIVE_WHOLE_FILE_SIZE);
- if (file_data == NULL) {
+ char* dest_path = path_cat(config->receive_root_directory, path);
+ if (!dest_path) {
file_destroy(file);
return NULL;
}
- if (config->use_compression &&
- !compression_should_skip_with_suffixes(file->path, config->skip_compress_suffixes,
- config->skip_compress_set ? config->skip_compress_count
- : -1)) {
- Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_WHOLE_FILE_SIZE);
- ProtocolSession* owner = file_data->owner;
- data_destroy(file_data);
- if (uncompressed == NULL) {
- file_destroy(file);
- return NULL;
- }
- if (!data_charge_session(uncompressed, owner, uncompressed->size)) {
- data_destroy(uncompressed);
- file_destroy(file);
- return NULL;
- }
- if (uncompressed->size > MAX_FILE_DATA_SIZE) {
- data_destroy(uncompressed);
- file_destroy(file);
- return NULL;
- }
- file_data = uncompressed;
+ bool ok = receive_file_payload_into(file, fd, config, dest_path, expected_size);
+ free(dest_path);
+ if (!ok) {
+ file_destroy(file);
+ return NULL;
}
- data_destroy(file->data);
- file->data = file_data;
return file;
}
@@ -974,13 +1032,12 @@ static IncrementalCheckOutcome incremental_check_receive_request(IncrementalChec
return INCREMENTAL_ERROR;
}
- /* A basis-configured run may materialize a file larger than the whole-file
- payload bound: a basis hit is streamed from the basis path (bounded
- buffers), so the check size is not itself an allocation. Every other
- path (delta/append/full) still applies MAX_RECEIVE_WHOLE_FILE_SIZE, and a
- miss simply falls through to the normal transfer with its own bound. */
- if (!config_has_basis(config) && state->check_size > MAX_RECEIVE_WHOLE_FILE_SIZE) {
- send_error_detail(fd, "check size exceeds receiver limit");
+ /* No file-size refusal: a whole-file payload larger than the historical
+ whole-file bound is streamed through a bounded buffer (see
+ file_receive_payload). Only a size that cannot be represented on this
+ platform is rejected. */
+ if (state->check_size > SIZE_MAX) {
+ send_error_detail(fd, "check size is not representable");
return INCREMENTAL_ERROR;
}
@@ -1411,7 +1468,7 @@ static IncrementalCheckOutcome incremental_check_try_append_resume(IncrementalCh
close(state->old_fd);
state->old_fd = -1;
}
- *out_file = receive_full_file(fd, config, check_path);
+ *out_file = receive_full_file(fd, config, check_path, check_size);
return INCREMENTAL_FILE;
}
@@ -1516,8 +1573,9 @@ static IncrementalCheckOutcome incremental_check_try_delta(IncrementalCheckState
bool try_delta, File** out_file) {
if (try_delta && state->old_data != NULL) {
bool delta_failed = false;
- File* delta_file = receive_delta_file(state->fd, state->config, state->check_path,
- state->old_data, state->old_size, &delta_failed);
+ File* delta_file =
+ receive_delta_file(state->fd, state->config, state->check_path, state->old_data,
+ state->old_size, state->check_size, -1, &delta_failed);
state->old_data = NULL; /* receive_delta_file consumes the snapshot on every path */
if (delta_file) {
*out_file = delta_file;
@@ -1540,14 +1598,14 @@ static IncrementalCheckOutcome incremental_check_try_fuzzy(IncrementalCheckState
if (!config->fuzzy || !config->use_delta)
return INCREMENTAL_CONTINUE;
unsigned long long fuzzy_size = 0;
- void* fuzzy_basis = fuzzy_basis_find_and_load(config, state->check_path, state->check_size,
- (time_t)state->check_mtime,
- (long)state->check_mtime_nsec, &fuzzy_size);
- if (fuzzy_basis != NULL) {
+ int fuzzy_fd = fuzzy_basis_find_and_open(config, state->check_path, state->check_size,
+ (time_t)state->check_mtime,
+ (long)state->check_mtime_nsec, &fuzzy_size);
+ if (fuzzy_fd >= 0) {
bool fuzzy_failed = false;
- File* fuzzy_file = receive_delta_file(state->fd, config, state->check_path, fuzzy_basis,
- fuzzy_size, &fuzzy_failed);
- fuzzy_basis = NULL; /* receive_delta_file consumes the buffer on every path */
+ File* fuzzy_file = receive_delta_file(state->fd, config, state->check_path, NULL, fuzzy_size,
+ state->check_size, fuzzy_fd, &fuzzy_failed);
+ close(fuzzy_fd);
if (fuzzy_file) {
*out_file = fuzzy_file;
return INCREMENTAL_FILE;
@@ -1555,7 +1613,6 @@ static IncrementalCheckOutcome incremental_check_try_fuzzy(IncrementalCheckState
if (fuzzy_failed)
return INCREMENTAL_ERROR;
}
- free(fuzzy_basis);
return INCREMENTAL_CONTINUE;
}
@@ -1567,7 +1624,7 @@ static File* incremental_check_receive_full(IncrementalCheckState* state) {
close(state->old_fd);
state->old_fd = -1;
}
- return receive_full_file(state->fd, state->config, state->check_path);
+ return receive_full_file(state->fd, state->config, state->check_path, state->check_size);
}
/* Core implementation. `would_transfer` (may be NULL) is set true only on the
diff --git a/src/shared/protocol.c b/src/shared/protocol.c
index 035c5fc..a6733f0 100644
--- a/src/shared/protocol.c
+++ b/src/shared/protocol.c
@@ -38,6 +38,27 @@ static atomic_ullong io_bytes_read = 0;
static unsigned long long global_bwlimit(void);
+/* Runtime whole-file receive bound (see protocol.h). Resolved once; an
+ * override can only LOWER the ceiling, never raise it above the protocol
+ * constant, so the wire/security bound is unchanged. A parse failure or a
+ * non-positive value leaves the default in place. */
+unsigned long long protocol_whole_file_receive_limit(void) {
+ static atomic_ullong cached = 0;
+ unsigned long long value = atomic_load_explicit(&cached, memory_order_relaxed);
+ if (value != 0)
+ return value;
+ value = MAX_RECEIVE_WHOLE_FILE_SIZE;
+ const char* env = getenv("FASTSYNC_MAX_WHOLE_FILE_SIZE");
+ if (env && env[0] != '\0') {
+ char* end = NULL;
+ unsigned long long parsed = strtoull(env, &end, 10);
+ if (end && *end == '\0' && parsed > 0 && parsed < value)
+ value = parsed;
+ }
+ atomic_store_explicit(&cached, value, memory_order_relaxed);
+ return value;
+}
+
/* ------------------------------------------------------------------------- *
* Transport vtable implementations.
*
@@ -769,17 +790,11 @@ bool protocol_send_data(ProtocolSession* session, const Data* data) {
return true;
}
-Data* protocol_receive_data_limited(ProtocolSession* session, unsigned long long maximum_size) {
+Data* protocol_receive_data_alloc(ProtocolSession* session, unsigned long long size) {
if (!session)
return NULL;
- unsigned long long size = 0;
- if (!protocol_receive_n_data(session, &size, sizeof(unsigned long long)))
+ if (size > MAX_DATA_PAYLOAD_SIZE)
return NULL;
- if (size > MAX_DATA_PAYLOAD_SIZE || size > maximum_size) {
- log_message(LOG_LEVEL_ERROR, "Data size %llu exceeds maximum %llu", size,
- (unsigned long long)MAX_DATA_PAYLOAD_SIZE);
- return NULL;
- }
if (size > SIZE_MAX)
return NULL;
size_t allocation_size = size == 0 ? 1 : (size_t)size;
@@ -794,12 +809,6 @@ Data* protocol_receive_data_limited(ProtocolSession* session, unsigned long long
protocol_release_memory_for_session(session, allocation_size);
return NULL;
}
- if (!protocol_receive_n_data(session, data, (size_t)size)) {
- free(data);
- protocol_release_memory_for_session(session, allocation_size);
- return NULL;
- }
- log_debug_message(LOG_DEBUG_PROTO, "Received %llu data", size);
Data* result = data_create(data, (size_t)size);
if (!result) {
protocol_release_memory_for_session(session, allocation_size);
@@ -810,6 +819,32 @@ Data* protocol_receive_data_limited(ProtocolSession* session, unsigned long long
return result;
}
+Data* protocol_receive_data_body(ProtocolSession* session, unsigned long long size) {
+ Data* result = protocol_receive_data_alloc(session, size);
+ if (!result)
+ return NULL;
+ if (!protocol_receive_n_data(session, result->data, (size_t)size)) {
+ data_destroy(result);
+ return NULL;
+ }
+ log_debug_message(LOG_DEBUG_PROTO, "Received %llu data", size);
+ return result;
+}
+
+Data* protocol_receive_data_limited(ProtocolSession* session, unsigned long long maximum_size) {
+ if (!session)
+ return NULL;
+ unsigned long long size = 0;
+ if (!protocol_receive_n_data(session, &size, sizeof(unsigned long long)))
+ return NULL;
+ if (size > MAX_DATA_PAYLOAD_SIZE || size > maximum_size) {
+ log_message(LOG_LEVEL_ERROR, "Data size %llu exceeds maximum %llu", size,
+ (unsigned long long)MAX_DATA_PAYLOAD_SIZE);
+ return NULL;
+ }
+ return protocol_receive_data_body(session, size);
+}
+
bool protocol_send_int(ProtocolSession* session, int data) {
if (!protocol_send_n_data(session, &data, sizeof(int)))
return false;
@@ -1114,6 +1149,12 @@ Data* receive_data(int fd) {
Data* receive_data_limited(int fd, unsigned long long maximum_size) {
return protocol_receive_data_limited(legacy_session(fd, -1), maximum_size);
}
+Data* receive_data_body(int fd, unsigned long long size) {
+ return protocol_receive_data_body(legacy_session(fd, -1), size);
+}
+Data* receive_data_alloc(int fd, unsigned long long size) {
+ return protocol_receive_data_alloc(legacy_session(fd, -1), size);
+}
bool send_int(int fd, int data) {
return protocol_send_int(legacy_session(-1, fd), data);
}
diff --git a/src/shared/protocol.h b/src/shared/protocol.h
index 4a04def..101ca73 100644
--- a/src/shared/protocol.h
+++ b/src/shared/protocol.h
@@ -32,6 +32,16 @@
/* Maximum allowed data payload size for receive_data (whole-file bound) */
#define MAX_DATA_PAYLOAD_SIZE MAX_RECEIVE_WHOLE_FILE_SIZE
+/* Runtime whole-file receive bound. It defaults to MAX_RECEIVE_WHOLE_FILE_SIZE
+ * and exists so the test suite can lower the ceiling (via the
+ * FASTSYNC_MAX_WHOLE_FILE_SIZE environment variable, a byte count) and exercise
+ * the streaming path with a small, fast transfer. A payload at or below the
+ * bound keeps the historical whole-buffer path; a larger one is streamed
+ * through a bounded buffer. The value is resolved once per process and never
+ * exceeds the compile-time ceiling, so a malicious environment cannot raise it
+ * beyond the protocol limit. */
+unsigned long long protocol_whole_file_receive_limit(void);
+
/* Maximum chunk size (64 MB) — prevents unbounded allocation from the wire */
#define MAX_CHUNK_SIZE (64ULL * 1024 * 1024)
/* Files larger than this are not kept fully in memory while loading: the
@@ -370,6 +380,15 @@ char* receive_str_redacted(int file_descriptor);
bool send_data(int file_descriptor, const Data* data);
Data* receive_data(int file_descriptor);
Data* receive_data_limited(int file_descriptor, unsigned long long maximum_size);
+/* Read exactly `size` bytes as a charged Data body. The length-prefixed
+ * receive_data_limited() reads the header itself; this variant is for callers
+ * that must inspect the declared size (and possibly stream the body instead)
+ * before allocating. `size` must already be within MAX_DATA_PAYLOAD_SIZE. */
+Data* receive_data_body(int file_descriptor, unsigned long long size);
+/* Allocate (and charge) a `size`-byte Data body without reading it; the caller
+ * fills `result->data` itself. Used when a frame's leading bytes must be
+ * inspected before the rest of the body is read. */
+Data* receive_data_alloc(int file_descriptor, unsigned long long size);
bool send_int(int file_descriptor, int data);
bool receive_int(int file_descriptor, int* data);
bool send_status(int file_descriptor, Status status);
diff --git a/tests/integration/common.py b/tests/integration/common.py
index b9eb14b..ebb2ab2 100644
--- a/tests/integration/common.py
+++ b/tests/integration/common.py
@@ -34,7 +34,7 @@ class ServerManager:
self._proc = None
self._port = None
- def start(self, extra_args=None):
+ def start(self, extra_args=None, env=None):
self.stop()
self._port = _find_free_port()
# Plain TCP is intentionally explicit in the server; integration tests
@@ -42,7 +42,11 @@ class ServerManager:
cmd = SERVER_CMD + ["-p", str(self._port), "--allow-unauthenticated"]
if extra_args:
cmd += extra_args
- self._proc = subprocess.Popen(cmd, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)
+ proc_env = dict(os.environ)
+ if env:
+ proc_env.update(env)
+ self._proc = subprocess.Popen(cmd, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL,
+ env=proc_env)
_wait_for_port(self._port, timeout=5)
def stop(self):
diff --git a/tests/integration/test_large_stream.py b/tests/integration/test_large_stream.py
new file mode 100644
index 0000000..cbed046
--- /dev/null
+++ b/tests/integration/test_large_stream.py
@@ -0,0 +1,222 @@
+"""#318: whole-file streaming above the receiver's 256 MiB ceiling.
+
+The receiver's historical whole-file bound (``MAX_RECEIVE_WHOLE_FILE_SIZE``,
+256 MiB) refused any single-file payload above it. The transfer engine now
+streams such a payload (and the basis read/verify/hash) through a bounded buffer
+and spools it to a temp file, so arbitrarily large single files transfer without
+being materialized in memory.
+
+To exercise the streaming path deterministically and quickly, these tests lower
+the receiver bound with the test-only ``FASTSYNC_MAX_WHOLE_FILE_SIZE`` hook (it
+can only lower, never raise, the protocol ceiling) and transfer a file a few
+times larger than the lowered bound. A real >256 MiB transfer is covered once,
+unmarked, so it runs in the full suite but not the fast PR gate.
+"""
+import hashlib
+import os
+import random
+import shutil
+import sys
+
+import pytest
+
+sys.path.insert(0, os.path.dirname(__file__))
+from common import ( # noqa: E402
+ ServerManager,
+ TEST_DATA_DIR,
+ clean_dir,
+ get_dest_received_dir,
+ run_client,
+)
+
+LOW_BOUND = 1024 * 1024
+FILE_SIZE = 3 * 1024 * 1024
+OLD_MTIME = 1_500_000_000
+
+
+@pytest.fixture(scope="module")
+def small_bound_server():
+ """A server whose whole-file streaming bound is 1 MiB."""
+ server = ServerManager()
+ server.start(extra_args=["--allow-super"],
+ env={"FASTSYNC_MAX_WHOLE_FILE_SIZE": str(LOW_BOUND)})
+ yield server
+ server.stop()
+
+
+def _payload(n):
+ rng = random.Random(0xC0FFEE)
+ return rng.randbytes(n)
+
+
+def _write(path, data, mtime=None):
+ os.makedirs(os.path.dirname(path), exist_ok=True)
+ with open(path, "wb") as fh:
+ fh.write(data)
+ if mtime is not None:
+ os.utime(path, (mtime, mtime))
+
+
+def _resolved(dest, source, rel):
+ return os.path.join(get_dest_received_dir(dest, source), rel)
+
+
+def _no_spool_leftovers(dest):
+ leftovers = []
+ for root, _dirs, files in os.walk(dest):
+ leftovers += [os.path.join(root, f) for f in files if ".fastsync-spool." in f]
+ return leftovers
+
+
+class TestStreamedWholeFile:
+ """A file above the (lowered) bound transfers correctly in every mode."""
+
+ @pytest.mark.parametrize(
+ "flags",
+ [
+ ["-a"],
+ ["-a", "--incremental"],
+ ["-a", "-z"],
+ ["-a", "--incremental", "-z"],
+ ["-a", "--threads", "--incremental"],
+ ["-a", "--inplace"],
+ ["-a", "--partial"],
+ ],
+ )
+ def test_above_bound_transfers(self, small_bound_server, flags):
+ tag = "_".join(f.strip("-") for f in flags) or "default"
+ source = os.path.join(TEST_DATA_DIR, f"stream_src_{tag}")
+ dest = os.path.join(TEST_DATA_DIR, f"stream_dst_{tag}")
+ clean_dir(source)
+ clean_dir(dest)
+ data = _payload(FILE_SIZE)
+ _write(os.path.join(source, "big.bin"), data, OLD_MTIME)
+ if "--inplace" in flags:
+ # --inplace only matters when the destination already exists.
+ _write(_resolved(dest, source, "big.bin"), b"stale", OLD_MTIME)
+
+ result, _ = run_client(source, dest, flags=flags, port=small_bound_server.port)
+ assert result.returncode == 0, (result.stderr or result.stdout)[:400]
+
+ got = os.path.join(get_dest_received_dir(dest, source), "big.bin")
+ assert os.path.exists(got), "streamed file was not written"
+ with open(got, "rb") as fh:
+ assert fh.read() == data, "streamed file content mismatch"
+ assert _no_spool_leftovers(dest) == [], "a spool temp file leaked"
+
+
+class TestStreamedBasis:
+ """A basis above the bound is streamed, not refused (compare/copy/link)."""
+
+ def _seed(self, dest, source, data):
+ clean_dir(source)
+ clean_dir(dest)
+ _write(os.path.join(source, "big.bin"), data, OLD_MTIME)
+ # FastSync resolves a relative basis DIR against the destination and
+ # appends the transfer-relative name.
+ _write(os.path.join(dest, "basis", "big.bin"), data, OLD_MTIME)
+
+ def test_compare_dest_above_bound(self, small_bound_server):
+ source = os.path.join(TEST_DATA_DIR, "sbasis_cmp_src")
+ dest = os.path.join(TEST_DATA_DIR, "sbasis_cmp_dst")
+ data = _payload(FILE_SIZE)
+ self._seed(dest, source, data)
+ result, _ = run_client(source, dest,
+ flags=["-a", "--compare-dest=basis", "--incremental"],
+ port=small_bound_server.port)
+ assert result.returncode == 0, (result.stderr or result.stdout)[:400]
+ # compare-dest never copies: an already-present destination stays sparse.
+ assert not os.path.exists(_resolved(dest, source, "big.bin"))
+
+ def test_copy_dest_above_bound(self, small_bound_server):
+ source = os.path.join(TEST_DATA_DIR, "sbasis_cpy_src")
+ dest = os.path.join(TEST_DATA_DIR, "sbasis_cpy_dst")
+ data = _payload(FILE_SIZE)
+ self._seed(dest, source, data)
+ result, _ = run_client(source, dest,
+ flags=["-a", "--copy-dest=basis", "--incremental"],
+ port=small_bound_server.port)
+ assert result.returncode == 0, (result.stderr or result.stdout)[:400]
+ got = _resolved(dest, source, "big.bin")
+ assert os.path.exists(got)
+ with open(got, "rb") as fh:
+ assert fh.read() == data
+ assert os.stat(got).st_ino != os.stat(os.path.join(dest, "basis", "big.bin")).st_ino
+ assert _no_spool_leftovers(dest) == []
+
+ def test_link_dest_above_bound(self, small_bound_server):
+ source = os.path.join(TEST_DATA_DIR, "sbasis_lnk_src")
+ dest = os.path.join(TEST_DATA_DIR, "sbasis_lnk_dst")
+ data = _payload(FILE_SIZE)
+ self._seed(dest, source, data)
+ result, _ = run_client(source, dest,
+ flags=["-a", "--link-dest=basis", "--incremental"],
+ port=small_bound_server.port)
+ assert result.returncode == 0, (result.stderr or result.stdout)[:400]
+ got = _resolved(dest, source, "big.bin")
+ assert os.path.exists(got)
+ with open(got, "rb") as fh:
+ assert fh.read() == data
+ assert os.stat(got).st_ino == os.stat(os.path.join(dest, "basis", "big.bin")).st_ino
+
+
+class TestFuzzyAboveBound:
+ """-y/--fuzzy reuses a basis above the bound by streaming its signature."""
+
+ def test_fuzzy_oversized_sibling(self, small_bound_server):
+ source = os.path.join(TEST_DATA_DIR, "sfuzzy_src")
+ dest = os.path.join(TEST_DATA_DIR, "sfuzzy_dst")
+ clean_dir(source)
+ clean_dir(dest)
+ base = _payload(FILE_SIZE)
+ sibling = bytearray(base)
+ sibling[FILE_SIZE // 2:FILE_SIZE // 2 + 4096] = bytes(
+ (b + 1) % 256 for b in sibling[FILE_SIZE // 2:FILE_SIZE // 2 + 4096])
+ _write(os.path.join(source, "report_v2.txt"), base)
+ _write(os.path.join(get_dest_received_dir(dest, source), "report_v1.txt"), bytes(sibling))
+ result, _ = run_client(
+ source, dest,
+ flags=["-a", "--incremental", "--delta", "--fuzzy", "--stats"],
+ port=small_bound_server.port)
+ assert result.returncode == 0, (result.stderr or result.stdout)[:400]
+ got = _resolved(dest, source, "report_v2.txt")
+ with open(got, "rb") as fh:
+ assert fh.read() == base, "fuzzy reconstruction mismatch"
+ assert _no_spool_leftovers(dest) == []
+
+
+class TestRealLargeFile:
+ """A real >256 MiB transfer, run only in the full (non-PR-gate) suite."""
+
+ def test_real_300mib_transfer(self, shared_server):
+ source = os.path.join(TEST_DATA_DIR, "real_large_src")
+ dest = os.path.join(TEST_DATA_DIR, "real_large_dst")
+ clean_dir(source)
+ clean_dir(dest)
+ n = 300 * 1024 * 1024
+ # Deterministic, compressible pattern written in bounded chunks.
+ chunk = bytes(range(256)) * 4096
+ digest = hashlib.sha256()
+ with open(os.path.join(source, "big.bin"), "wb") as fh:
+ written = 0
+ while written < n:
+ piece = chunk[: min(len(chunk), n - written)]
+ fh.write(piece)
+ digest.update(piece)
+ written += len(piece)
+
+ result, _ = run_client(source, dest, flags=["-a", "--incremental"],
+ port=shared_server.port)
+ assert result.returncode == 0, (result.stderr or result.stdout)[:400]
+ got = os.path.join(get_dest_received_dir(dest, source), "big.bin")
+ assert os.path.getsize(got) == n
+ got_digest = hashlib.sha256()
+ with open(got, "rb") as fh:
+ while True:
+ block = fh.read(1 << 20)
+ if not block:
+ break
+ got_digest.update(block)
+ assert got_digest.hexdigest() == digest.hexdigest()
+ shutil.rmtree(source, ignore_errors=True)
+ shutil.rmtree(dest, ignore_errors=True)
diff --git a/tests/test_file.c b/tests/test_file.c
index 7958687..15c3233 100644
--- a/tests/test_file.c
+++ b/tests/test_file.c
@@ -5,6 +5,8 @@
#include "file.h"
#include "file_receive.h"
#include "data.h"
+#include "compression.h"
+#include "delta.h"
#include "config.h"
#include "charset.h"
#include "utils.h"
@@ -2465,7 +2467,160 @@ static void test_manifest_would_delete_protects_absolute_basis() {
free(extra);
}
+/* #318: a whole-file payload above the streaming bound must be written to a
+ * spool temp file in bounded chunks, not materialized in memory. Exercises the
+ * raw and zstd-compressed paths and asserts the exact bytes land in the spool. */
+static void test_file_receive_payload_streams(void) {
+ const char* dir = "test_file_stream_tmp";
+ char dest_path[512];
+ snprintf(dest_path, sizeof(dest_path), "%s/out.bin", dir);
+ mkdir(dir, 0777);
+
+ const unsigned long long stream_limit = 4096;
+ size_t size = 20000;
+ unsigned char* payload = malloc(size);
+ EXPECT_NOT_NULL(payload);
+ for (size_t i = 0; i < size; i++)
+ payload[i] = (unsigned char)((i * 7 + 3) & 0xff);
+
+ /* Raw (uncompressed) streamed payload. */
+ {
+ int p[2];
+ EXPECT_EQ_INT(pipe(p), 0);
+ unsigned long long hdr = size;
+ EXPECT_TRUE(send_n_data(p[1], &hdr, sizeof(hdr)));
+ EXPECT_TRUE(send_n_data(p[1], payload, size));
+ Data* buffer = NULL;
+ char* spool = NULL;
+ unsigned long long out_size = 0;
+ EXPECT_TRUE(file_receive_payload(p[0], false, size, dest_path, stream_limit, &buffer, &spool,
+ &out_size));
+ EXPECT_NULL(buffer);
+ EXPECT_NOT_NULL(spool);
+ EXPECT_TRUE(out_size == size);
+ if (spool) {
+ FILE* fh = fopen(spool, "rb");
+ EXPECT_NOT_NULL(fh);
+ if (fh) {
+ unsigned char* got = malloc(size);
+ EXPECT_TRUE(fread(got, 1, size, fh) == size);
+ EXPECT_EQ_INT(memcmp(got, payload, size), 0);
+ free(got);
+ fclose(fh);
+ }
+ unlink(spool);
+ free(spool);
+ }
+ close(p[0]);
+ close(p[1]);
+ }
+
+ /* zstd-compressed payload whose logical size exceeds the bound: the frame is
+ * decompressed incrementally straight into the spool. */
+ {
+ unsigned char* copy = malloc(size);
+ EXPECT_NOT_NULL(copy);
+ memcpy(copy, payload, size);
+ Data* raw = data_create(copy, size); /* data_create takes ownership of copy */
+ Data* compressed = data_compress_codec(raw, COMPRESSION_ALGO_ZSTD, 3, 0);
+ data_destroy(raw);
+ EXPECT_NOT_NULL(compressed);
+ if (compressed) {
+ int p[2];
+ EXPECT_EQ_INT(pipe(p), 0);
+ EXPECT_TRUE(send_data(p[1], compressed));
+ Data* buffer = NULL;
+ char* spool = NULL;
+ unsigned long long out_size = 0;
+ EXPECT_TRUE(file_receive_payload(p[0], true, size, dest_path, stream_limit, &buffer, &spool,
+ &out_size));
+ EXPECT_NULL(buffer);
+ EXPECT_NOT_NULL(spool);
+ EXPECT_TRUE(out_size == size);
+ if (spool) {
+ FILE* fh = fopen(spool, "rb");
+ EXPECT_NOT_NULL(fh);
+ if (fh) {
+ unsigned char* got = malloc(size);
+ EXPECT_TRUE(fread(got, 1, size, fh) == size);
+ EXPECT_EQ_INT(memcmp(got, payload, size), 0);
+ free(got);
+ fclose(fh);
+ }
+ unlink(spool);
+ free(spool);
+ }
+ close(p[0]);
+ close(p[1]);
+ data_destroy(compressed);
+ }
+ }
+
+ free(payload);
+ unlink(dest_path);
+ rmdir(dir);
+}
+
+/* #318: the fd-based delta helpers must match the in-memory ones and stream the
+ * reconstruction to a descriptor without allocating the whole output. */
+static void test_delta_stream_helpers(void) {
+ const unsigned char basis[] = {0x11, 0x22, 0x33, 0x44};
+ const char* basis_path = "test_file_delta_basis.bin";
+ int bfd = open(basis_path, O_RDWR | O_CREAT | O_TRUNC, 0600);
+ EXPECT_TRUE(bfd >= 0);
+ EXPECT_TRUE(write(bfd, basis, sizeof(basis)) == (ssize_t)sizeof(basis));
+ EXPECT_TRUE(lseek(bfd, 0, SEEK_SET) == 0);
+
+ DeltaSignature* fd_sig = delta_signature_create_fd_seeded(bfd, sizeof(basis), 4, 0);
+ DeltaSignature* mem_sig = delta_signature_create_seeded(basis, sizeof(basis), 4, 0);
+ EXPECT_NOT_NULL(fd_sig);
+ EXPECT_NOT_NULL(mem_sig);
+ if (fd_sig && mem_sig) {
+ EXPECT_TRUE(fd_sig->block_count == mem_sig->block_count);
+ for (uint32_t i = 0; i < fd_sig->block_count && i < mem_sig->block_count; i++) {
+ EXPECT_TRUE(fd_sig->blocks[i].adler32 == mem_sig->blocks[i].adler32);
+ EXPECT_TRUE(fd_sig->blocks[i].xxhash == mem_sig->blocks[i].xxhash);
+ }
+ }
+ delta_signature_destroy(fd_sig);
+ delta_signature_destroy(mem_sig);
+
+ DeltaInstruction instrs[2];
+ instrs[0].type = DELTA_INSTR_BLOCK_MATCH;
+ instrs[0].match.block_index = 0;
+ instrs[0].match.block_offset = 0;
+ instrs[0].match.length = 4;
+ instrs[1].type = DELTA_INSTR_LITERAL;
+ instrs[1].literal.data = (uint8_t*)"XY";
+ instrs[1].literal.length = 2;
+ Delta delta;
+ delta.new_file_size = 6;
+ delta.instruction_count = 2;
+ delta.instructions = instrs;
+ delta.delta_size = 0;
+
+ int out[2];
+ EXPECT_EQ_INT(pipe(out), 0);
+ EXPECT_TRUE(delta_apply_to_fd(NULL, bfd, sizeof(basis), &delta, 4, out[1]));
+ close(out[1]);
+ unsigned char got[6] = {0};
+ size_t total = 0;
+ while (total < sizeof(got)) {
+ ssize_t n = read(out[0], got + total, sizeof(got) - total);
+ if (n <= 0)
+ break;
+ total += (size_t)n;
+ }
+ EXPECT_TRUE(total == sizeof(got));
+ EXPECT_TRUE(memcmp(got, "\x11\x22\x33\x44XY", 6) == 0);
+ close(out[0]);
+ close(bfd);
+ unlink(basis_path);
+}
+
void test_file() {
+ test_file_receive_payload_streams();
+ test_delta_stream_helpers();
test_file_create();
test_file_special_rdev_valid();
test_file_destroy_null();