mkit_core/transfer.rs
1//! Transfer-layer helpers: delta-aware pack planning and the packlist
2//! discovery wire format.
3//!
4//! These sit ABOVE the on-disk pack format (SPEC-PACKFILE / SPEC-DELTA)
5//! and below the transport. The push path uses [`plan_pack`] to decide
6//! which objects of a ref's closure to send, and how — raw, or as a delta
7//! against a base the remote already holds. The plan is then serialised
8//! with [`crate::pack::PackWriter`] into one or more packs — split when
9//! the plan's payload exceeds a single pack's cap (issue #831; the CLI's
10//! `remote_dispatch::build_and_upload_packs` owns the splitting) — each
11//! keyed by its own BLAKE3 digest.
12//!
13//! Because a delta-encoded pack is keyed by the pack digest (not by the
14//! reconstructed object's hash), the fetch side can no longer find it by
15//! walking object hashes. Each push records a [`PackListNode`] — this
16//! module's small versioned wire format — holding the pack(s) it added plus
17//! a `prev` pointer to the previous node, forming a per-branch chain. The
18//! push side advertises the chain head through a `refs/mkit/packmap/<branch>`
19//! metadata ref whose value is the BLAKE3 of the head node (itself stored as
20//! a pack object). The fetch side reads that ref, walks the chain, and
21//! unpacks every pack oldest-first. Chaining keeps each push O(1) on the
22//! wire rather than re-uploading the whole history's list.
23//!
24//! Content addressing is unchanged: delta is a transfer encoding only. The
25//! reconstructed object's id (BLAKE3 of its canonical bytes, or the merkle
26//! BMT root for a `Tree`/`ChunkedBlob` — see `crate::merkle`) is
27//! re-verified by `PackReader` before storing.
28
29use std::collections::{BTreeSet, HashMap, HashSet};
30
31use crate::delta;
32use crate::hash::{self, Hash};
33use crate::object::{Object, ObjectType};
34use crate::store::{ObjectStore, StoreError};
35
36// ---------------------------------------------------------------------------
37// PackList wire format
38// ---------------------------------------------------------------------------
39
40/// ASCII magic at the start of every packlist node ("mkit pack list").
41pub const PACKLIST_MAGIC: &[u8; 4] = b"MKPL";
42/// Current packlist version. Readers reject anything else so a future
43/// format change is a loud error, not a silent misparse.
44pub const PACKLIST_VERSION: u8 = 1;
45/// Hard cap on packs recorded in a single node — a normal push records
46/// one; refuse to allocate unboundedly on a malformed blob.
47pub const PACKLIST_MAX_ENTRIES: u32 = 1_000_000;
48
49/// The fixed guard header preceding the codec body: `[4B magic][1B
50/// version]`. The body (`prev` pointer + pack list) is encoded with
51/// `commonware-codec` and is variable-length.
52const PACKLIST_HEADER_LEN: usize = 4 + 1;
53
54/// A single node in a branch's packlist chain.
55///
56/// The push path appends one node per push: `prev` points at the previous
57/// node (the current packmap value before this push), and `packs` holds the
58/// pack(s) this push added. The full ordered pack set for a branch is the
59/// chain walked oldest-first — see `remote_dispatch::fetch_pack_chain`.
60///
61/// Chaining keeps each push O(1) on the wire (read a 32-byte pointer, write
62/// a ~64-byte node) instead of re-uploading the whole history's list.
63///
64/// A node's `prev` is reset to `None` (rather than linked to the prior head)
65/// when a push re-baselines (#406) — proactively bounding chain depth, or
66/// reactively escaping a broken chain. Either way the packs the superseded
67/// chain referenced become unreachable from the new head node: they are not
68/// deleted here (this module only ever writes new nodes/packs, never
69/// deletes), so they linger as orphaned storage on the remote until a
70/// server-side sweep reclaims them (tracked as makechain#849). A stale
71/// pack lingering harmlessly is safe; nothing on the fetch side ever
72/// resolves it once no live chain points to it.
73#[derive(Debug, Clone, PartialEq, Eq)]
74pub struct PackListNode {
75 /// Previous node's key, or `None` for the first node of a branch.
76 pub prev: Option<Hash>,
77 /// Pack keys added by this node, in apply order.
78 pub packs: Vec<Hash>,
79}
80
81/// Errors decoding a [`PackListNode`] blob.
82#[derive(Debug, thiserror::Error, PartialEq, Eq)]
83pub enum PackListError {
84 #[error("packlist is shorter than the {PACKLIST_HEADER_LEN}-byte magic+version header")]
85 TooShort,
86 #[error("packlist magic is not \"MKPL\"")]
87 InvalidMagic,
88 #[error("packlist version {0} is not supported (v1 only)")]
89 UnsupportedVersion(u8),
90 #[error("packlist pack count exceeds the {PACKLIST_MAX_ENTRIES} cap")]
91 TooManyEntries(u32),
92 /// The codec body (prev pointer / pack list) is malformed, truncated,
93 /// or has trailing bytes.
94 #[error("packlist body is malformed (bad codec payload or trailing bytes)")]
95 Malformed,
96}
97
98/// Serialise one packlist node: the `MKPL`/version guard header followed by
99/// the `prev`/`packs` body encoded with `commonware-codec` (idiomatic
100/// `Option` + `Vec`).
101///
102/// # Errors
103///
104/// [`PackListError::TooManyEntries`] if `packs` exceeds the cap.
105pub fn encode_packlist(prev: Option<Hash>, packs: &[Hash]) -> Result<Vec<u8>, PackListError> {
106 use commonware_codec::Write;
107 let count = u32::try_from(packs.len()).map_err(|_| PackListError::TooManyEntries(u32::MAX))?;
108 if count > PACKLIST_MAX_ENTRIES {
109 return Err(PackListError::TooManyEntries(count));
110 }
111 let mut out = Vec::new();
112 out.extend_from_slice(PACKLIST_MAGIC);
113 out.push(PACKLIST_VERSION);
114 // Body via commonware-codec: an `Option<Hash>` then a `Vec<Hash>`
115 // (length-prefixed). `Vec<u8>` is a `bytes::BufMut`.
116 prev.write(&mut out);
117 packs.to_vec().write(&mut out);
118 Ok(out)
119}
120
121/// Parse a packlist node blob.
122///
123/// # Errors
124///
125/// Returns the matching [`PackListError`] for a short buffer, wrong magic,
126/// or unknown version (the explicit guard header), or
127/// [`PackListError::Malformed`] / [`PackListError::TooManyEntries`] from
128/// the `commonware-codec` body (over-cap pack list, truncation, or
129/// trailing bytes).
130pub fn decode_packlist(bytes: &[u8]) -> Result<PackListNode, PackListError> {
131 use bytes::Buf as _;
132 use commonware_codec::{ReadExt, ReadRangeExt};
133
134 if bytes.len() < PACKLIST_HEADER_LEN {
135 return Err(PackListError::TooShort);
136 }
137 if &bytes[..4] != PACKLIST_MAGIC.as_slice() {
138 return Err(PackListError::InvalidMagic);
139 }
140 let version = bytes[4];
141 if version != PACKLIST_VERSION {
142 return Err(PackListError::UnsupportedVersion(version));
143 }
144 // Body: `Option<Hash>` then a length-capped `Vec<Hash>`. `&[u8]` is a
145 // `bytes::Buf`; the `RangeCfg` enforces the entry cap at decode time.
146 let mut buf: &[u8] = &bytes[PACKLIST_HEADER_LEN..];
147 let prev = <Option<Hash>>::read(&mut buf).map_err(|_| PackListError::Malformed)?;
148 let packs = <Vec<Hash>>::read_range(&mut buf, 0..=PACKLIST_MAX_ENTRIES as usize)
149 .map_err(|_| PackListError::Malformed)?;
150 // Trailing bytes after the declared body are a malformed packlist.
151 if buf.has_remaining() {
152 return Err(PackListError::Malformed);
153 }
154 Ok(PackListNode { prev, packs })
155}
156
157// ---------------------------------------------------------------------------
158// Delta base selection
159// ---------------------------------------------------------------------------
160
161/// Cap on tree-diff recursion depth when pairing chunked blobs across two
162/// commits. Bounds the walk on adversarial / pathologically nested trees.
163const MAX_TREE_DEPTH: usize = 64;
164
165/// Choose, for each changed `FastCDC` chunk or small blob in `new_tip`, a
166/// base object from `old_tip` to delta against.
167///
168/// (SPEC-DELTA §5 is informative — any base the size-gate later accepts is
169/// fine.) Diff the two commits' trees by path; where the same path is a
170/// [`crate::object::ChunkedBlob`] on both sides with a different manifest
171/// hash, pair each changed new chunk against a base in the old manifest —
172/// by same-index when the chunk counts match (an in-place edit didn't shift
173/// boundaries) or by content similarity when they differ (an insert/delete
174/// did, via `pair_chunks`). Identical chunks (present byte-for-byte in the
175/// old manifest) are skipped — they dedup for free and never need a delta.
176/// Where the same path is a plain [`crate::object::Blob`] on both sides
177/// (a changed small file, under `CHUNK_THRESHOLD`), the old blob is paired
178/// as the new blob's delta base directly — no similarity search, just
179/// "same path, old version" (#646a). A path whose kind changed between the
180/// two blob/chunked-blob variants, or that has no old-side counterpart at
181/// all (added path, or a rename), is left unpaired.
182///
183/// The returned map is `new_hash -> base_hash`. Every base is reachable
184/// from `old_tip`, so a remote that already holds `old_tip` holds the
185/// base. The caller still gates on whether the delta actually saves bytes
186/// ([`plan_pack`]).
187///
188/// # Errors
189///
190/// Propagates [`StoreError`] other than `ObjectNotFound`, which is treated
191/// as "no pairing available" (a partial local history is not fatal here —
192/// the worst case is that we send the chunk raw).
193pub fn select_chunk_delta_bases(
194 store: &ObjectStore,
195 new_tip: Hash,
196 old_tip: Hash,
197) -> Result<HashMap<Hash, Hash>, StoreError> {
198 let mut out = HashMap::new();
199 let (Some(new_tree), Some(old_tree)) = (tip_tree(store, new_tip)?, tip_tree(store, old_tip)?)
200 else {
201 return Ok(out);
202 };
203 pair_trees(store, new_tree, old_tree, 0, &mut out)?;
204 Ok(out)
205}
206
207/// Resolve a commit/remix tip to its root tree hash; `None` for any other
208/// object kind (a tip that is not a commit has no tree to diff).
209fn tip_tree(store: &ObjectStore, tip: Hash) -> Result<Option<Hash>, StoreError> {
210 match store.read_object(&tip) {
211 Ok(Object::Commit(c)) => Ok(Some(c.tree_hash)),
212 Ok(Object::Remix(r)) => Ok(Some(r.tree_hash)),
213 // A non-commit tip, or one we don't have, simply has no tree to diff.
214 Ok(_) | Err(StoreError::ObjectNotFound(_)) => Ok(None),
215 Err(e) => Err(e),
216 }
217}
218
219/// Merge-join two sorted trees by entry name, recursing into matching
220/// subtrees and pairing chunks for matching chunked-blob entries.
221fn pair_trees(
222 store: &ObjectStore,
223 new_tree: Hash,
224 old_tree: Hash,
225 depth: usize,
226 out: &mut HashMap<Hash, Hash>,
227) -> Result<(), StoreError> {
228 if depth > MAX_TREE_DEPTH || new_tree == old_tree {
229 return Ok(());
230 }
231 let (Some(Object::Tree(new_t)), Some(Object::Tree(old_t))) = (
232 read_optional(store, new_tree)?,
233 read_optional(store, old_tree)?,
234 ) else {
235 return Ok(());
236 };
237
238 // Both trees are lex-sorted by name (SPEC-OBJECTS §4); two-pointer
239 // merge to find entries present under the same name on both sides.
240 let (mut i, mut j) = (0usize, 0usize);
241 while i < new_t.entries.len() && j < old_t.entries.len() {
242 let ne = &new_t.entries[i];
243 let oe = &old_t.entries[j];
244 match ne.name.cmp(&oe.name) {
245 std::cmp::Ordering::Less => i += 1,
246 std::cmp::Ordering::Greater => j += 1,
247 std::cmp::Ordering::Equal => {
248 if ne.object_hash != oe.object_hash {
249 pair_entry(store, ne.object_hash, oe.object_hash, depth, out)?;
250 }
251 i += 1;
252 j += 1;
253 }
254 }
255 }
256 Ok(())
257}
258
259/// Dispatch a changed same-named entry: recurse into subtrees, pair
260/// chunked blobs, pair small blobs against their same-path prior version
261/// (#646a), ignore everything else (in particular, a `Blob<->ChunkedBlob`
262/// kind mismatch — a file crossing `CHUNK_THRESHOLD` between the two
263/// commits — is deliberately left unpaired here).
264fn pair_entry(
265 store: &ObjectStore,
266 new_hash: Hash,
267 old_hash: Hash,
268 depth: usize,
269 out: &mut HashMap<Hash, Hash>,
270) -> Result<(), StoreError> {
271 match (store.read_object(&new_hash), store.read_object(&old_hash)) {
272 (Ok(Object::Tree(_)), Ok(Object::Tree(_))) => {
273 pair_trees(store, new_hash, old_hash, depth + 1, out)
274 }
275 (Ok(Object::ChunkedBlob(new_cb)), Ok(Object::ChunkedBlob(old_cb))) => {
276 pair_chunks(store, &new_cb.chunks, &old_cb.chunks, out)
277 }
278 // A changed same-path small blob: the old blob at this path is a
279 // natural delta base for the new one (#646a). Never overwrite an
280 // existing pairing for `new_hash`, mirroring `pair_chunks`'
281 // determinism guarantee.
282 (Ok(Object::Blob(_)), Ok(Object::Blob(_))) => {
283 out.entry(new_hash).or_insert(old_hash);
284 Ok(())
285 }
286 // A missing object or a kind mismatch (file became a chunked blob,
287 // etc.) just yields no pairing for this entry.
288 (Err(StoreError::ObjectNotFound(_)), _) | (_, Err(StoreError::ObjectNotFound(_))) => Ok(()),
289 (Err(e), _) | (_, Err(e)) => Err(e),
290 _ => Ok(()),
291 }
292}
293
294/// Window size for content-similarity sketching. Matches the delta
295/// encoder's match-block size, so a shared 16-byte run is a shared feature.
296const FEATURE_WINDOW: usize = 16;
297/// Number of min-hash "super-features" sampled per chunk. Small keeps the
298/// index cheap; the delta size-gate in [`plan_pack`] is the final judge of a
299/// pick, so a few features are enough signal to rank candidates.
300const FEATURE_COUNT: usize = 4;
301
302/// Pair each changed new chunk against a base chunk in the old manifest.
303/// Skips chunks already present byte-identically in the old manifest (those
304/// dedup for free) and never overwrites an existing pairing, so the result
305/// is deterministic.
306///
307/// Two regimes:
308///
309/// * **Equal chunk counts** ⇒ `FastCDC` boundaries didn't shift (an in-place,
310/// same-length edit), so same-index pairing lands the new chunk on its own
311/// prior version. This is the common case and reads no chunk content.
312/// * **Differing counts** ⇒ an insert/delete shifted indices, so position is
313/// unreliable. We match by **content similarity**: index every old chunk's
314/// min-hash super-features, then pick, for each changed new chunk, the old
315/// chunk sharing the most features (deterministic tiebreak: lowest hash),
316/// falling back to the clamped same-index chunk when there is no overlap.
317/// This reads the changed file's chunk content (`O(file size)`), but only
318/// when there actually was a shift — and it trades those local reads for a
319/// much smaller upload than re-sending the shifted chunks raw.
320fn pair_chunks(
321 store: &ObjectStore,
322 new_chunks: &[Hash],
323 old_chunks: &[Hash],
324 out: &mut HashMap<Hash, Hash>,
325) -> Result<(), StoreError> {
326 if old_chunks.is_empty() {
327 return Ok(());
328 }
329 let old_set: HashSet<&Hash> = old_chunks.iter().collect();
330
331 // Fast path: no boundary shift → same-index pairing, no content reads.
332 if new_chunks.len() == old_chunks.len() {
333 for (j, nj) in new_chunks.iter().enumerate() {
334 if old_set.contains(nj) || out.contains_key(nj) {
335 continue;
336 }
337 if old_chunks[j] != *nj {
338 out.insert(*nj, old_chunks[j]);
339 }
340 }
341 return Ok(());
342 }
343
344 // Shifted: index old chunks by their super-features once. An unreadable
345 // old chunk (absent / non-blob) contributes no features and can't be a
346 // base — it is simply absent from the index.
347 let mut feature_index: HashMap<u64, Vec<Hash>> = HashMap::new();
348 for oc in old_chunks {
349 if let Some(content) = chunk_bytes(store, oc)? {
350 for f in chunk_features(&content) {
351 feature_index.entry(f).or_default().push(*oc);
352 }
353 }
354 }
355
356 for nj in new_chunks {
357 if old_set.contains(nj) || out.contains_key(nj) {
358 continue;
359 }
360 // An unreadable new chunk has no content signal — skip pairing it
361 // (it's sent raw) rather than guessing a base.
362 let Some(content) = chunk_bytes(store, nj)? else {
363 continue;
364 };
365 // Count shared features per candidate old chunk.
366 let mut votes: HashMap<Hash, u32> = HashMap::new();
367 for f in &chunk_features(&content) {
368 if let Some(cands) = feature_index.get(f) {
369 for cand in cands {
370 *votes.entry(*cand).or_default() += 1;
371 }
372 }
373 }
374 // Pair ONLY on genuine content overlap: most shared features wins,
375 // ties → lowest hash (deterministic, order-independent). No overlap →
376 // no base; the chunk is sent raw (a dissimilar base would be rejected
377 // by the size-gate anyway). No position fallback after a shift.
378 if let Some((base, _)) = votes
379 .into_iter()
380 .filter(|(cand, _)| cand != nj)
381 .max_by(|x, y| x.1.cmp(&y.1).then_with(|| y.0.cmp(&x.0)))
382 {
383 out.insert(*nj, base);
384 }
385 }
386 Ok(())
387}
388
389/// Read a chunk blob's CONTENT bytes (not its serialized object form).
390///
391/// Returns `None` when the object is absent or not a blob — there is no
392/// content to compare, so the caller skips pairing that chunk rather than
393/// inventing a base. A genuine read/IO error propagates, preserving the
394/// [`select_chunk_delta_bases`] contract (only `ObjectNotFound` is swallowed).
395fn chunk_bytes(store: &ObjectStore, h: &Hash) -> Result<Option<Vec<u8>>, StoreError> {
396 match store.read_object(h) {
397 Ok(Object::Blob(b)) => Ok(Some(b.data)),
398 Ok(_) | Err(StoreError::ObjectNotFound(_)) => Ok(None),
399 Err(e) => Err(e),
400 }
401}
402
403/// Up to [`FEATURE_COUNT`] shift-resistant "super-features" for a chunk: the
404/// smallest distinct `FEATURE_WINDOW`-byte FNV-1a window hashes. Two chunks
405/// that share content share these min-hashes with high probability even when
406/// the content is shifted, making them a cheap similarity key. Chunks shorter
407/// than the window yield none (such a chunk has no similarity key, so it is
408/// left unpaired and sent raw).
409fn chunk_features(bytes: &[u8]) -> Vec<u64> {
410 if bytes.len() < FEATURE_WINDOW {
411 return Vec::new();
412 }
413 // Keep the K smallest DISTINCT window hashes, ascending.
414 let mut best: Vec<u64> = Vec::with_capacity(FEATURE_COUNT + 1);
415 for w in bytes.windows(FEATURE_WINDOW) {
416 let h = fnv1a(w);
417 if best.len() == FEATURE_COUNT && h >= best[FEATURE_COUNT - 1] {
418 continue;
419 }
420 if let Err(pos) = best.binary_search(&h) {
421 best.insert(pos, h);
422 best.truncate(FEATURE_COUNT);
423 }
424 }
425 best
426}
427
428/// FNV-1a 64-bit over a fixed window — this module's own min-hash feature
429/// key; unrelated to (and independent of) `delta::block_hash`'s algorithm,
430/// which only needs to agree with itself between the base-index build and
431/// the scan, not with any hash used here.
432fn fnv1a(block: &[u8]) -> u64 {
433 let mut h: u64 = 0xcbf2_9ce4_8422_2325;
434 for &b in block {
435 h ^= u64::from(b);
436 h = h.wrapping_mul(0x0000_0001_0000_01b3);
437 }
438 h
439}
440
441/// Read an object, mapping a missing object to `None` so callers can treat
442/// "absent" the same as "not the kind I wanted" without a hard error.
443fn read_optional(store: &ObjectStore, h: Hash) -> Result<Option<Object>, StoreError> {
444 match store.read_object(&h) {
445 Ok(o) => Ok(Some(o)),
446 Err(StoreError::ObjectNotFound(_)) => Ok(None),
447 Err(e) => Err(e),
448 }
449}
450
451// ---------------------------------------------------------------------------
452// Pack planning
453// ---------------------------------------------------------------------------
454
455/// One delta entry in a [`PackPlan`]: a target object encoded against a
456/// base the remote already holds.
457#[derive(Debug, Clone)]
458pub struct PlannedDelta {
459 /// BLAKE3 of the reconstructed (canonical) object — its storage id.
460 pub target: Hash,
461 /// BLAKE3 of the base object the delta is applied against.
462 pub base: Hash,
463 /// SPEC-DELTA instruction stream.
464 pub stream: Vec<u8>,
465}
466
467/// A deterministic plan for the pack(s) a push uploads for one ref —
468/// serialised into a single pack when it fits under the payload cap, or
469/// split across several when it doesn't (issue #831).
470///
471/// Entries are pre-ordered for [`crate::pack::PackWriter`]: all `raw`
472/// objects first (non-blobs before blobs, each group in BLAKE3 order),
473/// then `deltas`. Delta bases are external — they live in `old_tip`'s
474/// closure, which the remote already holds and earlier packs already
475/// delivered, NEVER in an entry introduced earlier in this same plan —
476/// so no in-pack base ordering is required (SPEC-PACKFILE §4), and a
477/// left-to-right split of this sequence across multiple packs is always
478/// safe.
479#[derive(Debug, Clone, Default)]
480pub struct PackPlan {
481 /// Objects to send verbatim, already ordered non-blobs-then-blobs.
482 pub raw: Vec<Hash>,
483 /// Objects to send as a delta against an already-present base.
484 pub deltas: Vec<PlannedDelta>,
485 /// `true` when the pack needs no externally-resolved base — i.e. this is
486 /// a full-closure push (no usable `old_tip`), so the pack reconstructs
487 /// the ref's whole closure on its own. The push path normally **appends**
488 /// this pack to the branch's packlist chain; it uses `self_contained`
489 /// to decide whether a push may **reset** to a fresh chain, in two
490 /// cases:
491 ///
492 /// * Reactively, when the prior chain is unreadable — a self-contained
493 /// pack is the one kind that can safely escape a broken chain.
494 /// * Proactively (#406): when a *healthy* chain has simply grown past
495 /// the re-baseline depth threshold, the push side calls [`plan_pack`]
496 /// with `old_tip: None` specifically to force `self_contained: true`,
497 /// then resets the chain to bound its depth. This reuses the atomic
498 /// head+packmap advance (#408) the reactive reset already depends on
499 /// — no separate mechanism was needed once that landed.
500 pub self_contained: bool,
501}
502
503impl PackPlan {
504 /// Total objects the plan transfers (raw + delta).
505 #[must_use]
506 pub fn object_count(&self) -> usize {
507 self.raw.len() + self.deltas.len()
508 }
509
510 /// `true` when the plan carries nothing — the remote already holds the
511 /// ref's whole closure.
512 #[must_use]
513 pub fn is_empty(&self) -> bool {
514 self.raw.is_empty() && self.deltas.is_empty()
515 }
516}
517
518/// Plan the pack for pushing `new_tip` to a remote whose current tip is
519/// `old_tip` (`None` for a first push or a remote we cannot diff against).
520///
521/// Computes the send-set as `closure(new_tip) \ closure(old_tip)` — so
522/// objects the remote already holds are never re-sent (this is the
523/// identical-object dedup, preferred over delta) — then delta-encodes
524/// changed chunks against same-path prior chunks where that actually
525/// saves bytes, falling back to raw otherwise.
526///
527/// Delta candidates are encoded one at a time, sequentially — see
528/// [`plan_pack_with`] for a version that lets a caller fan that step out
529/// across a thread pool.
530///
531/// # Errors
532///
533/// Propagates [`StoreError`] from reading the local closure. A missing
534/// `old_tip` closure is treated as "diff unavailable" (full-closure,
535/// all-raw push) rather than an error.
536pub fn plan_pack(
537 store: &ObjectStore,
538 new_tip: Hash,
539 old_tip: Option<Hash>,
540) -> Result<PackPlan, StoreError> {
541 plan_pack_with(store, new_tip, old_tip, |store, candidates| {
542 candidates
543 .iter()
544 .map(|&c| encode_delta_candidate(store, c))
545 .collect()
546 })
547}
548
549/// One delta-encoding candidate handed in a batch to [`plan_pack_with`]'s
550/// `encode_deltas` callback: encode `target` against `base` if the
551/// remote is already known to hold `base`.
552#[derive(Debug, Clone, Copy)]
553pub struct DeltaCandidate {
554 /// BLAKE3 of the blob to encode.
555 pub target: Hash,
556 /// BLAKE3 of the base object to encode `target` against.
557 pub base: Hash,
558}
559
560/// Encode one [`DeltaCandidate`] against `store`. Returns `None` when the
561/// delta stream would not actually shrink the payload — [`plan_pack`] and
562/// [`plan_pack_with`] both fall back to sending `target` raw in that case.
563///
564/// Reads both `target`'s and `base`'s bytes and runs [`delta::encode`]'s
565/// block-hash-table build plus greedy scan — both CPU-bound and
566/// independent of every other candidate, which is why this is exposed as
567/// its own function: it's the unit [`plan_pack_with`] callers fan out per
568/// candidate rather than per whole batch.
569///
570/// Multiple candidates commonly share the same `base` (e.g. several
571/// chunks of one file all diffed against the same prior chunk) — each
572/// call here re-reads it independently. [`encode_delta_candidate_with_base`]
573/// is the same unit for a caller that has already deduplicated those
574/// reads across a batch.
575///
576/// # Errors
577///
578/// Propagates [`StoreError`] from either read.
579pub fn encode_delta_candidate(
580 store: &ObjectStore,
581 candidate: DeltaCandidate,
582) -> Result<Option<PlannedDelta>, StoreError> {
583 let target_bytes = store.read(&candidate.target)?;
584 let base_bytes = store.read(&candidate.base)?;
585 Ok(try_delta(
586 candidate.target,
587 candidate.base,
588 &target_bytes,
589 &base_bytes,
590 ))
591}
592
593/// [`encode_delta_candidate`], but for a caller that has already read
594/// `candidate.base`'s bytes — e.g. `mkit-cli`'s rayon fan-out, which
595/// reads each *distinct* base in `plan_pack_with`'s candidate batch once
596/// up front (a `HashMap<Hash, Vec<u8>>` keyed by base hash) instead of
597/// leaving `N` candidates sharing a base to each independently
598/// `store.read` it — redundant I/O that a parallel fan-out turns into
599/// concurrent redundant I/O and concurrent redundant in-memory copies of
600/// the same bytes, not just serialized-and-short-lived like the
601/// sequential default pays.
602///
603/// # Errors
604///
605/// Propagates [`StoreError`] from reading `target`'s bytes.
606pub fn encode_delta_candidate_with_base(
607 store: &ObjectStore,
608 candidate: DeltaCandidate,
609 base_bytes: &[u8],
610) -> Result<Option<PlannedDelta>, StoreError> {
611 let target_bytes = store.read(&candidate.target)?;
612 Ok(try_delta(
613 candidate.target,
614 candidate.base,
615 &target_bytes,
616 base_bytes,
617 ))
618}
619
620/// A blob's slot in [`plan_pack_with`]'s final (BLAKE3-ordered)
621/// `blob_raw`/`deltas` split: either an immediate raw hash, or a
622/// placeholder pointing at its index in the candidates slice handed to
623/// `encode_deltas` — filled in once that callback returns.
624enum BlobSlot {
625 Raw(Hash),
626 Pending(usize),
627}
628
629/// [`plan_pack`], but with an explicit `encode_deltas` batch callback in
630/// place of the built-in sequential loop over [`encode_delta_candidate`].
631///
632/// `mkit-core` has no thread-pool dependency of its own — it stays usable
633/// from wasm targets, which have no OS threads — so the fan-out decision
634/// lives with the caller instead of living in this function. `mkit-cli`'s
635/// native, rayon-backed push path passes a parallel `encode_deltas`
636/// (mirroring the `prepare_raw_batch`/`prepare_delta_batch` compression
637/// fan-out already in `remote_dispatch/mod.rs`, but for the diffing step
638/// itself rather than the post-plan zstd pass); [`plan_pack`] passes a
639/// plain sequential loop.
640///
641/// `encode_deltas` receives every delta candidate for this plan in one
642/// batch, in the order they were discovered (BLAKE3/send-set order), and
643/// MUST return exactly one `Option<PlannedDelta>` per input candidate, in the same
644/// order — a `None` for a given candidate falls back to sending that
645/// blob raw, matching [`encode_delta_candidate`]'s own "smaller or raw"
646/// rule.
647///
648/// # Errors
649///
650/// Propagates [`StoreError`] from the local closure walk or from
651/// `encode_deltas`. Also returns [`StoreError::DeltaBatchLengthMismatch`]
652/// if `encode_deltas` returns a `Vec` of a different length than the
653/// candidate slice it was given — a contract violation by the caller, not
654/// a normal runtime condition, but reported as a typed error rather than
655/// a panic since a push is a network-facing operation where a clean error
656/// is a better failure mode than crashing the process.
657pub fn plan_pack_with(
658 store: &ObjectStore,
659 new_tip: Hash,
660 old_tip: Option<Hash>,
661 encode_deltas: impl FnOnce(
662 &ObjectStore,
663 &[DeltaCandidate],
664 ) -> Result<Vec<Option<PlannedDelta>>, StoreError>,
665) -> Result<PackPlan, StoreError> {
666 let new_set = crate::ops::reachable_objects(store, &new_tip)?;
667
668 // What the remote already holds = the closure of its current tip, if we
669 // can compute it locally. A partial local history (some referenced
670 // object pruned) degrades to "send everything" rather than failing.
671 let mut remote_set = BTreeSet::new();
672 let mut base_map = HashMap::new();
673 let mut have_old = false;
674 if let Some(o) = old_tip
675 && store.contains(&o)
676 {
677 match crate::ops::reachable_objects(store, &o) {
678 Ok(s) => {
679 remote_set = s;
680 base_map = select_chunk_delta_bases(store, new_tip, o)?;
681 have_old = true;
682 }
683 Err(StoreError::ObjectNotFound(_)) => {}
684 Err(e) => return Err(e),
685 }
686 }
687
688 let send: BTreeSet<Hash> = new_set.difference(&remote_set).copied().collect();
689
690 // Partition the send-set, iterating in BTreeSet (BLAKE3) order so the
691 // plan — and therefore the pack bytes — is deterministic.
692 let mut non_blob_raw = Vec::new();
693 let mut blob_slots = Vec::new();
694 let mut candidates: Vec<DeltaCandidate> = Vec::new();
695
696 for h in &send {
697 // Classify via the cheap 6-byte prologue check (`object_type`)
698 // instead of a full read+verify of the object's bytes — the send-set
699 // classification only needs the type tag, not the content (INV-14).
700 // Bytes are read lazily, only for the subset that are actually
701 // delta candidates, by `encode_deltas` below.
702 let is_blob = store.object_type(h)? == ObjectType::Blob;
703 if !is_blob {
704 non_blob_raw.push(*h);
705 continue;
706 }
707
708 // Only blobs (FastCDC chunks) are delta candidates, and only against
709 // a base the remote actually holds.
710 match base_map.get(h) {
711 Some(&base) if remote_set.contains(&base) => {
712 blob_slots.push(BlobSlot::Pending(candidates.len()));
713 candidates.push(DeltaCandidate { target: *h, base });
714 }
715 _ => blob_slots.push(BlobSlot::Raw(*h)),
716 }
717 }
718
719 let mut delta_results = encode_deltas(store, &candidates)?;
720 if delta_results.len() != candidates.len() {
721 return Err(StoreError::DeltaBatchLengthMismatch {
722 expected: candidates.len(),
723 actual: delta_results.len(),
724 });
725 }
726
727 let mut blob_raw = Vec::with_capacity(blob_slots.len());
728 let mut deltas = Vec::with_capacity(candidates.len());
729 for slot in blob_slots {
730 match slot {
731 BlobSlot::Raw(h) => blob_raw.push(h),
732 BlobSlot::Pending(idx) => match delta_results[idx].take() {
733 Some(planned) => deltas.push(planned),
734 None => blob_raw.push(candidates[idx].target),
735 },
736 }
737 }
738
739 let mut raw = non_blob_raw;
740 raw.extend(blob_raw);
741
742 Ok(PackPlan {
743 raw,
744 deltas,
745 self_contained: !have_old,
746 })
747}
748
749/// Encode `target` against `base` and return the delta only if it is
750/// strictly smaller on the wire than sending the target raw. The per-entry
751/// frame (SPEC-PACKFILE §2) is identical for raw and delta, so only the
752/// payloads differ: a delta payload is `base_hash (HASH_LEN) + stream`
753/// (SPEC-PACKFILE §3.2) versus the raw object bytes. Compare those.
754///
755/// Takes `base_bytes` already read — both callers ([`encode_delta_candidate`]
756/// and [`encode_delta_candidate_with_base`]) read `base` themselves,
757/// from different sourcing strategies, so this pure comparison step
758/// doesn't need to know which.
759fn try_delta(
760 target: Hash,
761 base: Hash,
762 target_bytes: &[u8],
763 base_bytes: &[u8],
764) -> Option<PlannedDelta> {
765 let Ok(stream) = delta::encode(base_bytes, target_bytes) else {
766 // Over-u32 inputs can't happen here (object cap < 4 GiB), but treat
767 // any encode failure as "send raw" rather than propagating.
768 return None;
769 };
770 if hash::HASH_LEN + stream.len() < target_bytes.len() {
771 Some(PlannedDelta {
772 target,
773 base,
774 stream,
775 })
776 } else {
777 None
778 }
779}
780
781// =========================================================================
782// Tests
783// =========================================================================
784
785#[cfg(test)]
786mod tests {
787 use super::*;
788 use crate::chunker::{ChunkIterator, FastCdc};
789 use crate::object::{Blob, ChunkedBlob, Commit, EntryMode, Identity, Tree, TreeEntry};
790 use crate::pack::{PackReader, PackWriter};
791 use crate::serialize;
792 use tempfile::TempDir;
793
794 fn store() -> (TempDir, ObjectStore) {
795 let d = TempDir::new().unwrap();
796 let s = ObjectStore::init(&crate::layout::RepoLayout::single(d.path())).unwrap();
797 (d, s)
798 }
799
800 fn put(s: &ObjectStore, obj: &Object) -> Hash {
801 s.write(&serialize::serialize(obj).unwrap()).unwrap()
802 }
803
804 /// Store `data` as `FastCDC` chunks + a `ChunkedBlob` manifest, mirroring
805 /// the worktree large-file path. Returns the manifest hash.
806 fn put_chunked(s: &ObjectStore, data: &[u8]) -> Hash {
807 let chunks: Vec<Hash> = ChunkIterator::new(FastCdc::v1(), data)
808 .map(|b| {
809 put(
810 s,
811 &Object::Blob(Blob {
812 data: data[b.offset..b.offset + b.length].to_vec(),
813 }),
814 )
815 })
816 .collect();
817 put(
818 s,
819 &Object::ChunkedBlob(ChunkedBlob {
820 total_size: data.len() as u64,
821 chunk_size: 0,
822 chunks,
823 }),
824 )
825 }
826
827 fn commit_with_file(s: &ObjectStore, file_hash: Hash, parents: Vec<Hash>, msg: &str) -> Hash {
828 commit_with_named_file(s, b"big.bin", file_hash, parents, msg)
829 }
830
831 /// Like [`commit_with_file`], but with a caller-chosen tree entry name —
832 /// needed to build trees with more than one entry, or to control whether
833 /// two commits share a path.
834 fn commit_with_named_file(
835 s: &ObjectStore,
836 name: &[u8],
837 file_hash: Hash,
838 parents: Vec<Hash>,
839 msg: &str,
840 ) -> Hash {
841 let tree = put(
842 s,
843 &Object::Tree(Tree {
844 entries: vec![TreeEntry {
845 name: name.to_vec(),
846 mode: EntryMode::Blob,
847 object_hash: file_hash,
848 }],
849 }),
850 );
851 put(
852 s,
853 &Object::Commit(Commit::new_unannotated(
854 tree,
855 parents,
856 Identity::ed25519([7; 32]),
857 [0; 32],
858 msg.as_bytes().to_vec(),
859 msg.len() as u64,
860 [0; 64],
861 )),
862 )
863 }
864
865 /// A >1 MiB pseudo-random buffer (deterministic) that `FastCDC` splits
866 /// into several chunks. Splitmix64 keeps it dependency-free and
867 /// reproducible.
868 fn big_buffer() -> Vec<u8> {
869 let mut data = vec![0u8; 2 * 1024 * 1024];
870 let mut state: u64 = 0x1234_5678_9abc_def0;
871 for chunk in data.chunks_mut(8) {
872 state = state.wrapping_add(0x9e37_79b9_7f4a_7c15);
873 let mut z = state;
874 z = (z ^ (z >> 30)).wrapping_mul(0xbf58_476d_1ce4_e5b9);
875 z = (z ^ (z >> 27)).wrapping_mul(0x94d0_49bb_1331_11eb);
876 z ^= z >> 31;
877 let bytes = z.to_le_bytes();
878 let n = chunk.len();
879 chunk.copy_from_slice(&bytes[..n]);
880 }
881 data
882 }
883
884 #[test]
885 fn packlist_node_roundtrip_with_prev() {
886 let prev = [0x11u8; 32];
887 let packs = vec![[1u8; 32], [2u8; 32], [0xABu8; 32]];
888 let bytes = encode_packlist(Some(prev), &packs).unwrap();
889 assert_eq!(&bytes[..4], PACKLIST_MAGIC);
890 assert_eq!(bytes[4], PACKLIST_VERSION);
891 let node = decode_packlist(&bytes).unwrap();
892 assert_eq!(node.prev, Some(prev));
893 assert_eq!(node.packs, packs);
894 }
895
896 #[test]
897 fn packlist_node_roundtrip_first_node() {
898 // First node of a branch: no predecessor, one pack.
899 let bytes = encode_packlist(None, &[[7u8; 32]]).unwrap();
900 let node = decode_packlist(&bytes).unwrap();
901 assert_eq!(node.prev, None);
902 assert_eq!(node.packs, vec![[7u8; 32]]);
903 }
904
905 #[test]
906 fn packlist_rejects_bad_magic_version_body_and_length() {
907 let good = encode_packlist(Some([0x11u8; 32]), &[[9u8; 32]]).unwrap();
908
909 // Magic and version are the explicit guard header (precise errors).
910 let mut bad_magic = good.clone();
911 bad_magic[0] = b'X';
912 assert_eq!(
913 decode_packlist(&bad_magic),
914 Err(PackListError::InvalidMagic)
915 );
916
917 let mut bad_ver = good.clone();
918 bad_ver[4] = 2;
919 assert_eq!(
920 decode_packlist(&bad_ver),
921 Err(PackListError::UnsupportedVersion(2))
922 );
923
924 // Truncated codec body → Malformed.
925 let mut short = good.clone();
926 short.pop();
927 assert_eq!(decode_packlist(&short), Err(PackListError::Malformed));
928
929 // Trailing byte after the declared body → Malformed.
930 let mut long = good;
931 long.push(0);
932 assert_eq!(decode_packlist(&long), Err(PackListError::Malformed));
933
934 // Shorter than the magic+version guard header → TooShort.
935 assert_eq!(decode_packlist(&[0u8; 3]), Err(PackListError::TooShort));
936 }
937
938 #[test]
939 fn packlist_decode_rejects_over_cap_pack_list() {
940 // A REAL over-cap body: `encode_packlist` refuses to build one,
941 // so assemble the node by hand with the same codec — the guard
942 // header followed by `Option<Hash>` + a `Vec<Hash>` holding
943 // PACKLIST_MAX_ENTRIES + 1 entries (~32 MiB). The decode-side
944 // RangeCfg must reject it as Malformed, not allocate it.
945 use commonware_codec::Write as _;
946 let over_cap = PACKLIST_MAX_ENTRIES as usize + 1;
947 let mut node = Vec::with_capacity(PACKLIST_HEADER_LEN + 8 + 32 * over_cap);
948 node.extend_from_slice(PACKLIST_MAGIC);
949 node.push(PACKLIST_VERSION);
950 let prev: Option<Hash> = None;
951 prev.write(&mut node);
952 vec![[1u8; 32]; over_cap].write(&mut node);
953 assert_eq!(decode_packlist(&node), Err(PackListError::Malformed));
954
955 // Sanity: the identical construction at exactly the cap is
956 // accepted, so the failure above is the cap check itself, not
957 // an artifact of the hand-rolled encoding.
958 let mut at_cap = Vec::with_capacity(PACKLIST_HEADER_LEN + 8 + 32 * (over_cap - 1));
959 at_cap.extend_from_slice(PACKLIST_MAGIC);
960 at_cap.push(PACKLIST_VERSION);
961 let prev: Option<Hash> = None;
962 prev.write(&mut at_cap);
963 vec![[1u8; 32]; over_cap - 1].write(&mut at_cap);
964 let decoded = decode_packlist(&at_cap).expect("at-cap list must decode");
965 assert_eq!(decoded.packs.len(), PACKLIST_MAX_ENTRIES as usize);
966 }
967
968 #[test]
969 fn plan_first_push_is_self_contained_all_raw() {
970 let (_d, s) = store();
971 let file = put_chunked(&s, &big_buffer());
972 let c1 = commit_with_file(&s, file, vec![], "v1");
973
974 let plan = plan_pack(&s, c1, None).unwrap();
975 assert!(plan.self_contained);
976 assert!(plan.deltas.is_empty(), "no base on a first push");
977 // commit + tree + manifest + every chunk must be present.
978 let full = crate::ops::reachable_objects(&s, &c1).unwrap();
979 assert_eq!(plan.object_count(), full.len());
980 }
981
982 // =================================================================
983 // object_type() short-circuit for plan_pack's type check (#636 /
984 // INV-14) — classification must use the cheap type-prologue check
985 // rather than a full read+verify, while objects actually selected as
986 // delta candidates still get a real, verified read.
987 // =================================================================
988
989 /// Flip a byte well past the 6-byte type prologue so the object's
990 /// on-disk content no longer matches its BLAKE3 hash, while its type
991 /// tag stays intact and `object_type()` still reads it correctly.
992 fn corrupt_payload_byte(s: &ObjectStore, h: &Hash) {
993 use std::fs::OpenOptions;
994 use std::io::{Read, Seek, SeekFrom, Write};
995 let path = s.path_for(h);
996 let mut f = OpenOptions::new()
997 .read(true)
998 .write(true)
999 .open(&path)
1000 .unwrap();
1001 f.seek(SeekFrom::Start(6)).unwrap();
1002 let mut byte = [0u8; 1];
1003 f.read_exact(&mut byte).unwrap();
1004 f.seek(SeekFrom::Start(6)).unwrap();
1005 f.write_all(&[byte[0] ^ 0xFF]).unwrap();
1006 f.sync_all().unwrap();
1007 }
1008
1009 #[test]
1010 fn plan_pack_first_push_tolerates_corrupted_raw_blob_content() {
1011 // First push: every blob in the send-set is `raw` (no delta base
1012 // exists yet), so classification never needs to read a blob's
1013 // content — only its type. A blob whose payload has bit-rotted
1014 // must therefore still plan cleanly and land in `raw`; the actual
1015 // pack build is where its (now-failing) integrity check belongs.
1016 let (_d, s) = store();
1017 let file = put_chunked(&s, &big_buffer());
1018 let c1 = commit_with_file(&s, file, vec![], "v1");
1019
1020 let full = crate::ops::reachable_objects(&s, &c1).unwrap();
1021 let blob_hash = *full
1022 .iter()
1023 .find(|h| s.object_type(h).unwrap() == ObjectType::Blob)
1024 .expect("chunked big_buffer must contain at least one blob chunk");
1025 corrupt_payload_byte(&s, &blob_hash);
1026
1027 // Sanity: a full verified read of this object now fails.
1028 assert!(matches!(
1029 s.read(&blob_hash),
1030 Err(StoreError::HashMismatch { .. })
1031 ));
1032
1033 let plan = plan_pack(&s, c1, None).unwrap();
1034 assert!(plan.raw.contains(&blob_hash));
1035 assert_eq!(plan.object_count(), full.len());
1036 }
1037
1038 #[test]
1039 fn plan_pack_second_push_still_verifies_delta_candidate_bytes() {
1040 // A blob that IS a delta candidate (paired with a base the remote
1041 // holds) must still get a real, verified read — the short-circuit
1042 // only removes the *classification* read, not the read needed to
1043 // actually build a delta.
1044 let (_d, s) = store();
1045 let v1 = big_buffer();
1046 let file1 = put_chunked(&s, &v1);
1047 let c1 = commit_with_file(&s, file1, vec![], "v1");
1048
1049 let mut v2 = v1.clone();
1050 for k in 0..16 {
1051 v2[900_000 + k] ^= 0xFF;
1052 }
1053 let file2 = put_chunked(&s, &v2);
1054 let c2 = commit_with_file(&s, file2, vec![c1], "v2");
1055
1056 let bases = select_chunk_delta_bases(&s, c2, c1).unwrap();
1057 let target = *bases.keys().next().expect("expected a paired chunk");
1058 corrupt_payload_byte(&s, &target);
1059
1060 let err = plan_pack(&s, c2, Some(c1)).unwrap_err();
1061 assert!(matches!(err, StoreError::HashMismatch { .. }));
1062 }
1063
1064 #[test]
1065 fn plan_second_push_deltas_changed_chunk_and_skips_unchanged() {
1066 let (_d, s) = store();
1067 let v1 = big_buffer();
1068 let file1 = put_chunked(&s, &v1);
1069 let c1 = commit_with_file(&s, file1, vec![], "v1");
1070
1071 // Edit a small in-place region (same length → stable boundaries).
1072 let mut v2 = v1.clone();
1073 for k in 0..16 {
1074 v2[900_000 + k] ^= 0xFF;
1075 }
1076 let file2 = put_chunked(&s, &v2);
1077 let c2 = commit_with_file(&s, file2, vec![c1], "v2");
1078
1079 let plan = plan_pack(&s, c2, Some(c1)).unwrap();
1080 assert!(
1081 !plan.self_contained,
1082 "second push diffs against the prior tip"
1083 );
1084
1085 // At least one changed chunk should delta-compress.
1086 assert!(
1087 !plan.deltas.is_empty(),
1088 "expected a delta for the edited chunk"
1089 );
1090
1091 // The send-set must exclude every object the v1 closure already
1092 // holds (identical-chunk dedup), so it is far smaller than the full
1093 // v2 closure.
1094 let v2_full = crate::ops::reachable_objects(&s, &c2).unwrap();
1095 assert!(
1096 plan.object_count() < v2_full.len(),
1097 "unchanged chunks must not be re-sent"
1098 );
1099
1100 // The pack must reconstruct, bit-for-bit, against a store seeded
1101 // with the v1 closure (what the remote/fetcher already holds).
1102 assert_pack_reconstructs(&s, &plan, c1, c2);
1103 }
1104
1105 /// Build the planned pack, replay it into a fresh store pre-seeded with
1106 /// `old_tip`'s closure, and assert every `new_tip` object reconstructs
1107 /// with a matching hash.
1108 fn assert_pack_reconstructs(src: &ObjectStore, plan: &PackPlan, old_tip: Hash, new_tip: Hash) {
1109 let mut w = PackWriter::new();
1110 for h in &plan.raw {
1111 let bytes = src.read(h).unwrap();
1112 w.push_raw(*h, &bytes).unwrap();
1113 }
1114 for d in &plan.deltas {
1115 w.push_delta(&d.base, &d.stream).unwrap();
1116 }
1117 let pack = w.finish().unwrap();
1118
1119 // Seed the destination with what the remote already has.
1120 let (_d2, dst) = store();
1121 for h in crate::ops::reachable_objects(src, &old_tip).unwrap() {
1122 dst.write(&src.read(&h).unwrap()).unwrap();
1123 }
1124
1125 PackReader::read(&pack, &dst).unwrap();
1126
1127 for h in crate::ops::reachable_objects(src, &new_tip).unwrap() {
1128 // `dst.read` already re-verifies the object addresses to `h`
1129 // (merkle root for Tree/ChunkedBlob, BLAKE3 otherwise); assert
1130 // the dispatched id explicitly via `Object::id`.
1131 let got = dst.read(&h).unwrap();
1132 assert_eq!(
1133 crate::serialize::deserialize(&got).unwrap().id().unwrap(),
1134 h,
1135 "reconstructed object must address to its id"
1136 );
1137 assert_eq!(got, src.read(&h).unwrap());
1138 }
1139 }
1140
1141 #[test]
1142 fn select_bases_pairs_only_changed_chunks() {
1143 let (_d, s) = store();
1144 let v1 = big_buffer();
1145 let file1 = put_chunked(&s, &v1);
1146 let c1 = commit_with_file(&s, file1, vec![], "v1");
1147
1148 let mut v2 = v1.clone();
1149 v2[1_000_000] ^= 0xFF;
1150 let file2 = put_chunked(&s, &v2);
1151 let c2 = commit_with_file(&s, file2, vec![c1], "v2");
1152
1153 let bases = select_chunk_delta_bases(&s, c2, c1).unwrap();
1154 assert!(!bases.is_empty());
1155 // Every paired target is a chunk that is new in v2, and every base
1156 // is a chunk that v1 held.
1157 let v1_set = crate::ops::reachable_objects(&s, &c1).unwrap();
1158 for (target, base) in &bases {
1159 assert!(!v1_set.contains(target), "target should be new");
1160 assert!(v1_set.contains(base), "base must be present at old tip");
1161 }
1162 }
1163
1164 /// Deterministic, feature-rich chunk content (xorshift so 16-byte window
1165 /// hashes are well-distributed rather than colliding on flat data).
1166 fn chunk_content(seed: u8, len: usize) -> Vec<u8> {
1167 let mut state = u64::from(seed).wrapping_mul(0x9E37_79B9_7F4A_7C15) | 1;
1168 let mut v = vec![0u8; len];
1169 for byte in &mut v {
1170 state ^= state << 13;
1171 state ^= state >> 7;
1172 state ^= state << 17;
1173 *byte = (state & 0xff) as u8;
1174 }
1175 v
1176 }
1177
1178 fn put_blob(s: &ObjectStore, content: Vec<u8>) -> Hash {
1179 put(s, &Object::Blob(Blob { data: content }))
1180 }
1181
1182 /// Build a `ChunkedBlob` from explicit chunk contents (bypassing
1183 /// `FastCDC`) so a test can control chunk boundaries and shifts precisely.
1184 fn put_chunked_from(s: &ObjectStore, chunks: &[Vec<u8>]) -> Hash {
1185 let total: usize = chunks.iter().map(Vec::len).sum();
1186 let chunk_ids: Vec<Hash> = chunks.iter().map(|c| put_blob(s, c.clone())).collect();
1187 put(
1188 s,
1189 &Object::ChunkedBlob(ChunkedBlob {
1190 total_size: total as u64,
1191 chunk_size: 0,
1192 chunks: chunk_ids,
1193 }),
1194 )
1195 }
1196
1197 #[test]
1198 #[allow(clippy::many_single_char_names)] // a..e keep the chunk-shift table compact
1199 fn content_aware_pairing_picks_similar_base_after_a_shift() {
1200 let (_d, s) = store();
1201 let (a, b, c, d, e) = (
1202 chunk_content(1, 400),
1203 chunk_content(2, 400),
1204 chunk_content(3, 400),
1205 chunk_content(4, 400),
1206 chunk_content(5, 400),
1207 );
1208 let mut e_mod = e.clone();
1209 e_mod[200] ^= 0xFF; // a one-byte edit to E
1210
1211 // Old file = [A,B,C,D,E]; new file = [B,C,D,E'] — A deleted (so every
1212 // surviving chunk's index shifts down by one) and E modified. Counts
1213 // differ, so the content-similarity path runs.
1214 let old_file = put_chunked_from(&s, &[a, b.clone(), c.clone(), d.clone(), e.clone()]);
1215 let new_file = put_chunked_from(&s, &[b, c, d.clone(), e_mod.clone()]);
1216 let c_old = commit_with_file(&s, old_file, vec![], "old");
1217 let c_new = commit_with_file(&s, new_file, vec![c_old], "new");
1218
1219 let bases = select_chunk_delta_bases(&s, c_new, c_old).unwrap();
1220
1221 let e_mod_id = put_blob(&s, e_mod);
1222 let e_id = put_blob(&s, e);
1223 let d_id = put_blob(&s, d);
1224
1225 // E' is the only changed chunk. It sits at new index 3; the old chunk
1226 // at index 3 is D, so position-clamp pairing would mispair E'→D.
1227 // Content similarity must instead pair E'→E (its real prior version).
1228 assert_eq!(
1229 bases.get(&e_mod_id),
1230 Some(&e_id),
1231 "E' must delta against E (content match), not the wrong-index D"
1232 );
1233 assert_ne!(bases.get(&e_mod_id), Some(&d_id));
1234 }
1235
1236 #[test]
1237 #[allow(clippy::many_single_char_names)] // a..d + z keep the chunk table compact
1238 fn shifted_pairing_skips_a_chunk_with_no_content_overlap() {
1239 // After a shift (differing counts), a brand-new chunk dissimilar to
1240 // every old chunk must NOT be paired — no position/clamp fallback.
1241 let (_d, s) = store();
1242 let (a, b, c, d) = (
1243 chunk_content(1, 400),
1244 chunk_content(2, 400),
1245 chunk_content(3, 400),
1246 chunk_content(4, 400),
1247 );
1248 let z = chunk_content(99, 400); // unrelated to a/b/c/d
1249 // old = [A,B,C,D]; new = [B,C,Z] — counts differ (content path), and Z
1250 // is new + dissimilar to every old chunk.
1251 let old_file = put_chunked_from(&s, &[a, b.clone(), c.clone(), d]);
1252 let new_file = put_chunked_from(&s, &[b, c, z.clone()]);
1253 let c_old = commit_with_file(&s, old_file, vec![], "old");
1254 let c_new = commit_with_file(&s, new_file, vec![c_old], "new");
1255
1256 let bases = select_chunk_delta_bases(&s, c_new, c_old).unwrap();
1257 let z_id = put_blob(&s, z);
1258 assert!(
1259 !bases.contains_key(&z_id),
1260 "a dissimilar new chunk must be left unpaired (sent raw), not clamped to a base"
1261 );
1262 }
1263
1264 #[test]
1265 fn equal_count_edit_still_uses_fast_index_path() {
1266 // Same chunk count ⇒ no shift ⇒ same-index pairing (no content reads).
1267 let (_d, s) = store();
1268 let (a, b, c) = (
1269 chunk_content(10, 400),
1270 chunk_content(11, 400),
1271 chunk_content(12, 400),
1272 );
1273 let mut b_mod = b.clone();
1274 b_mod[50] ^= 0xFF;
1275 let old_file = put_chunked_from(&s, &[a.clone(), b.clone(), c.clone()]);
1276 let new_file = put_chunked_from(&s, &[a, b_mod.clone(), c]);
1277 let c_old = commit_with_file(&s, old_file, vec![], "old");
1278 let c_new = commit_with_file(&s, new_file, vec![c_old], "new");
1279
1280 let bases = select_chunk_delta_bases(&s, c_new, c_old).unwrap();
1281 let b_mod_id = put_blob(&s, b_mod);
1282 let b_id = put_blob(&s, b);
1283 assert_eq!(
1284 bases.get(&b_mod_id),
1285 Some(&b_id),
1286 "in-place edit pairs by index"
1287 );
1288 }
1289
1290 // =================================================================
1291 // #646a — small (non-chunked) blob delta pairing. `pair_entry` must
1292 // pair a changed `Blob<->Blob` same-path entry (previously a no-op),
1293 // using the old blob as the delta base for the new one, without
1294 // touching cross-path or `Blob<->ChunkedBlob` transitions.
1295 // =================================================================
1296
1297 #[test]
1298 fn small_blob_edit_is_delta_paired_against_same_path_prior_version() {
1299 let (_d, s) = store();
1300 // Small, well under CHUNK_THRESHOLD: a single Blob object, not a
1301 // ChunkedBlob manifest. Large enough (many unchanged lines around a
1302 // single mid-file edit) that a delta's fixed per-instruction
1303 // overhead (SPEC-DELTA header + two COPYs + one INSERT + the
1304 // HASH_LEN base pointer) is comfortably smaller than the raw blob.
1305 let mut v1 = Vec::new();
1306 for i in 0..40 {
1307 v1.extend_from_slice(format!("unchanged line number {i:02}\n").as_bytes());
1308 }
1309 let mut v2 = v1.clone();
1310 // Replace exactly one line in the middle with a different one,
1311 // keeping everything before/after byte-identical to v1.
1312 let needle = b"unchanged line number 20\n".to_vec();
1313 let pos = v2
1314 .windows(needle.len())
1315 .position(|w| w == needle.as_slice())
1316 .unwrap();
1317 v2.splice(
1318 pos..pos + needle.len(),
1319 b"THIS LINE WAS EDITED\n".iter().copied(),
1320 );
1321
1322 let blob1 = put_blob(&s, v1.clone());
1323 let blob2 = put_blob(&s, v2.clone());
1324 let c1 = commit_with_named_file(&s, b"small.txt", blob1, vec![], "v1");
1325 let c2 = commit_with_named_file(&s, b"small.txt", blob2, vec![c1], "v2");
1326
1327 let bases = select_chunk_delta_bases(&s, c2, c1).unwrap();
1328 assert_eq!(
1329 bases.get(&blob2),
1330 Some(&blob1),
1331 "same-path small-blob edit must pair against the prior version"
1332 );
1333
1334 // The pairing must actually be usable by plan_pack: it should
1335 // produce a delta for blob2, and that delta must be strictly
1336 // smaller on the wire than sending blob2 raw.
1337 let plan = plan_pack(&s, c2, Some(c1)).unwrap();
1338 let planned = plan
1339 .deltas
1340 .iter()
1341 .find(|d| d.target == blob2)
1342 .expect("expected a planned delta for the edited small blob");
1343 assert_eq!(planned.base, blob1);
1344 assert!(
1345 hash::HASH_LEN + planned.stream.len() < v2.len(),
1346 "delta payload must beat sending the new blob raw"
1347 );
1348 assert!(
1349 !plan.raw.contains(&blob2),
1350 "the edited blob should not also be sent raw"
1351 );
1352
1353 assert_pack_reconstructs(&s, &plan, c1, c2);
1354 }
1355
1356 #[test]
1357 fn plan_pack_with_reports_a_typed_error_when_encode_deltas_miscounts() {
1358 // Same fixture as `small_blob_edit_is_delta_paired_against_same_path_prior_version`
1359 // — one delta candidate — but with a broken `encode_deltas` that
1360 // drops it instead of returning one result per candidate. Pins
1361 // `plan_pack_with`'s length-check contract: a caller bug here must
1362 // surface as `StoreError::DeltaBatchLengthMismatch`, not a panic.
1363 let (_d, s) = store();
1364 let mut v1 = Vec::new();
1365 for i in 0..40 {
1366 v1.extend_from_slice(format!("unchanged line number {i:02}\n").as_bytes());
1367 }
1368 let mut v2 = v1.clone();
1369 let needle = b"unchanged line number 20\n".to_vec();
1370 let pos = v2
1371 .windows(needle.len())
1372 .position(|w| w == needle.as_slice())
1373 .unwrap();
1374 v2.splice(
1375 pos..pos + needle.len(),
1376 b"THIS LINE WAS EDITED\n".iter().copied(),
1377 );
1378 let blob1 = put_blob(&s, v1);
1379 let blob2 = put_blob(&s, v2);
1380 let c1 = commit_with_named_file(&s, b"small.txt", blob1, vec![], "v1");
1381 let c2 = commit_with_named_file(&s, b"small.txt", blob2, vec![c1], "v2");
1382
1383 let bases = select_chunk_delta_bases(&s, c2, c1).unwrap();
1384 assert_eq!(
1385 bases.get(&blob2),
1386 Some(&blob1),
1387 "sanity: expects a delta candidate"
1388 );
1389
1390 let err = plan_pack_with(&s, c2, Some(c1), |_store, candidates| {
1391 assert_eq!(
1392 candidates.len(),
1393 1,
1394 "sanity: exactly one candidate expected"
1395 );
1396 Ok(Vec::new()) // drops the one candidate's result — the bug under test
1397 })
1398 .unwrap_err();
1399 assert!(
1400 matches!(
1401 err,
1402 StoreError::DeltaBatchLengthMismatch {
1403 expected: 1,
1404 actual: 0
1405 }
1406 ),
1407 "expected DeltaBatchLengthMismatch, got {err:?}"
1408 );
1409 }
1410
1411 #[test]
1412 fn small_blob_pairing_is_strictly_same_path() {
1413 // Two unrelated paths, each holding a small blob, where the "new"
1414 // path's content is a near-duplicate of the "old" path's content.
1415 // Despite the content similarity, pairing is same-path only — a
1416 // blob at a path with no old-side counterpart must never be paired
1417 // against a differently-named entry's prior blob.
1418 let (_d, s) = store();
1419 let a_v1 = b"shared prefix content for path a\n".to_vec();
1420 let mut a_v2 = a_v1.clone();
1421 a_v2.push(b'!');
1422
1423 let blob_a1 = put_blob(&s, a_v1.clone());
1424 let blob_a2 = put_blob(&s, a_v2);
1425 let c1 = commit_with_named_file(&s, b"a.txt", blob_a1, vec![], "v1");
1426
1427 // c2 changes "a.txt" AND introduces a brand-new path "b.txt" whose
1428 // content is near-identical to a.txt's OLD content (but it is a
1429 // distinct path with no old-side entry at all).
1430 let blob_b = put_blob(&s, a_v1);
1431 let tree2 = put(
1432 &s,
1433 &Object::Tree(Tree {
1434 entries: vec![
1435 TreeEntry {
1436 name: b"a.txt".to_vec(),
1437 mode: EntryMode::Blob,
1438 object_hash: blob_a2,
1439 },
1440 TreeEntry {
1441 name: b"b.txt".to_vec(),
1442 mode: EntryMode::Blob,
1443 object_hash: blob_b,
1444 },
1445 ],
1446 }),
1447 );
1448 let c2 = put(
1449 &s,
1450 &Object::Commit(Commit::new_unannotated(
1451 tree2,
1452 vec![c1],
1453 Identity::ed25519([7; 32]),
1454 [0; 32],
1455 b"v2".to_vec(),
1456 2,
1457 [0; 64],
1458 )),
1459 );
1460
1461 let bases = select_chunk_delta_bases(&s, c2, c1).unwrap();
1462 assert_eq!(
1463 bases.get(&blob_a2),
1464 Some(&blob_a1),
1465 "a.txt's edit still pairs against its own prior version"
1466 );
1467 assert!(
1468 !bases.contains_key(&blob_b),
1469 "b.txt has no old-side counterpart and must not be paired at all, \
1470 even against a content-similar blob from a different path"
1471 );
1472 }
1473}
1474
1475#[cfg(test)]
1476mod wire_conformance {
1477 use super::*;
1478
1479 /// Conformance pin: the exact on-wire bytes of a `PackListNode`. Locks
1480 /// the `commonware-codec` body encoding (`MKPL` + version guard, then a
1481 /// codec `Option<Hash>` and a length-prefixed `Vec<Hash>`) so a
1482 /// commonware version bump that silently changes the encoding is caught
1483 /// here rather than in the field. Regenerate deliberately only on an
1484 /// intentional format change (and bump `PACKLIST_VERSION`).
1485 #[test]
1486 fn packlist_wire_format_is_pinned() {
1487 let bytes = encode_packlist(Some([0x11u8; 32]), &[[0x22u8; 32], [0x33u8; 32]]).unwrap();
1488 let expected = "4d4b504c0101\
1489 1111111111111111111111111111111111111111111111111111111111111111\
1490 02\
1491 2222222222222222222222222222222222222222222222222222222222222222\
1492 3333333333333333333333333333333333333333333333333333333333333333";
1493 assert_eq!(
1494 hash::to_hex_bytes(&bytes),
1495 expected,
1496 "PackListNode wire format drifted (commonware-codec change?)"
1497 );
1498 let node = decode_packlist(&bytes).unwrap();
1499 assert_eq!(node.prev, Some([0x11u8; 32]));
1500 assert_eq!(node.packs, vec![[0x22u8; 32], [0x33u8; 32]]);
1501 }
1502}