diff --git a/CHANGELOG.md b/CHANGELOG.md index b5624b2..11025ef 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,35 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added +- Reference-delta object storage: `objects.base_oid`/`chain_depth` (migration + `0013`) let an object be stored as a zstd-prefix-mode diff against a chosen + prior "base" object instead of always being fully materialized and + independently compressed. `mnemosyne_postgres::delta` implements the codec + (property-tested for byte-exact round trips); `receive_pack::unpack_into_tx` + resolves each pack entry's in-pack delta relationship (RefDelta/OfsDelta, + depth-capped at 50) during ingest; `object_find` and the parallel + connectivity-check's batched lookup both walk the chain via a recursive CTE. + Deliberately uses a portable codec (zstd's prefix/dictionary mode) rather + than git's own ofs-delta/ref-delta pack opcodes, so the same object pool + and delta representation stay usable by non-git importers. Measured on a + real push (this repo's own history, 793 objects): 9.01x compression vs. + 2.24x for the prior no-delta shape, landing within ~2x of git's own native + pack size for the identical content. See + `docs/compression-and-throughput-strategy.md` for the full investigation, + measured comparisons, and scope boundaries (in-pack bases only; no GC/ + refcounting yet; CDC evaluated and demoted to a future secondary layer for + large blobs). +- Push-ingest parallelization: object-pool inserts and pack-entry decoding + both fan out across connections/threads instead of running serially on a + single connection, with `synchronous_commit = off` on the bulk-insert + connections (durability preserved via the push's own later synchronous + ref-edit commit — see `insert_objects_parallel`'s doc comment for why). + `--db-max-connections` (env `MNEMOSYNE_DB_MAX_CONNECTIONS`, default 64) + replaces a hardcoded pool size of 16, which had been silently capping bulk + push concurrency at 2x instead of the validated ~8x sweet spot. TOAST + compression for `objects.data` switched from `pglz` to `lz4` (migration + `0012`) — faster and smaller on real object content. See + `docs/handoff-push-performance.md` for the full measured investigation. - Async object-source trait + bench coverage of the per-request walkers: `mnemosyne_git::ObjectSource` unifies `&mut`-receiver object lookup with a blanket bridge from any `ObjectStore`, so walker code in diff --git a/Cargo.lock b/Cargo.lock index 3dd24b3..b1afb8d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4672,6 +4672,7 @@ dependencies = [ "gix-ref 0.63.0", "mnemosyne_git", "mnemosyne_harness", + "proptest", "rand 0.10.1", "serde_json", "sqlx", @@ -4679,6 +4680,7 @@ dependencies = [ "testcontainers-modules", "thiserror 2.0.18", "tokio", + "zstd-safe", ] [[package]] @@ -8589,3 +8591,22 @@ name = "zmij" version = "1.0.21" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b8848ee67ecc8aedbaf3e4122217aff892639231befc6a1b58d29fff4c2cabaa" + +[[package]] +name = "zstd-safe" +version = "7.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f49c4d5f0abb602a93fb8736af2a4f4dd9512e36f7f570d66e65ff867ed3b9d" +dependencies = [ + "zstd-sys", +] + +[[package]] +name = "zstd-sys" +version = "2.0.16+zstd.1.5.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "91e19ebc2adc8f83e43039e79776e3fda8ca919132d68a1fed6a5faca2683748" +dependencies = [ + "cc", + "pkg-config", +] diff --git a/Cargo.toml b/Cargo.toml index 8db159a..2dd0807 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -66,6 +66,19 @@ base64 = "0.22" mime_guess = "2" tar = "0.4" flate2 = "1" +# Reference-delta object storage: compress an object against a chosen +# prior "base" object using zstd's prefix/dictionary mode +# (`ZSTD_CCtx_refPrefix`/`ZSTD_DCtx_refPrefix`, the mechanism behind +# the `zstd --patch-from` CLI feature) instead of compressing each +# object in isolation. `zstd-safe` gives 1:1 access to that API; +# the higher-level `zstd` crate's dictionary story is training-dict +# oriented and doesn't fit "diff against this exact prior version." +# See docs/compression-and-throughput-strategy.md for the measured +# rationale (33-156x smaller than independent per-object compression +# on real repository content) and why this is deliberately NOT git's +# own ofs-delta/ref-delta pack format — a portable codec keeps the +# object pool usable by non-git importers. +zstd-safe = { version = "7", features = ["std"] } # Language detection for sh.tangled.repo.languages. Brings its own gix # transitively (not the same versions as the gix-* crates pinned above); # we use only its custom FileSource trait + classifiers, never its diff --git a/docs/compression-and-throughput-strategy.md b/docs/compression-and-throughput-strategy.md index 5872da8..555bc49 100644 --- a/docs/compression-and-throughput-strategy.md +++ b/docs/compression-and-throughput-strategy.md @@ -1,29 +1,34 @@ -# Compression and throughput strategy (proposal, not yet implemented) - -Status: **analysis + measured comparisons, no code changes**. Written in -response to a direct ask to investigate why mnemosyne pushes nixpkgs in ~25 -minutes at ~40GB, against Go-knot's ~5 minutes / ~3GB and knot2's ~2 minutes, -and to find both bottleneck speedups and cross-row compression strategies. -Extends `docs/handoff-push-performance.md` (throughput investigation) and +# Compression and throughput strategy + +Status: **implemented and verified** (schema, codec, ingest path, read path — +see "Implementation status" below for what shipped and how it was verified). +Originally written as analysis in response to a direct ask to investigate why +mnemosyne pushes nixpkgs in ~25 minutes at \~40GB, against Go-knot's \~5 +minutes / \~3GB and knot2's \~2 minutes. Extends +`docs/handoff-push-performance.md` (throughput investigation) and `docs/sharding-and-tiered-storage.md` (isolation, explicitly *not* -compression) — read those first; this doc doesn't repeat their content. - -Constraint from the requester, load-bearing for everything below: the storage -design should stay compatible with non-git transports in spirit — jj (no -native wire protocol yet, but the door should stay open) and, ideally, -Mercurial — without either copying git's specific mechanisms by default or -rejecting them out of hand. The operative rule: pick the mechanism that's -actually justified by the data, encode it in a standard/portable form, not a -git-specific opcode set. +compression) — read those first; this doc doesn't repeat their content. This +repo is a POC for understanding the domain boundaries of "store versioned +files in Postgres such that Git/JJ can use it as a remote" — the operator +directed implementing this design against a small local test repo (mnemosyne +itself) rather than the full nixpkgs corpus, given local disk constraints; see +"Implementation status" for what that validated. + +Constraint from the requester, load-bearing for the design: the storage +design should stay compatible with non-git transports in spirit — jj (no native +wire protocol yet, but the door should stay open) and, ideally, Mercurial — +without either copying git's specific mechanisms by default or rejecting them +out of hand. The operative rule: pick the mechanism that's actually justified by +the data, encode it in a standard/portable form, not a git-specific opcode set. ## tl;dr - **The 40GB is not "gzip vs no gzip." It's "no delta compression at all."** Every object is stored fully materialized, independently TOAST-compressed. - Git's ~3GB for nixpkgs comes almost entirely from delta-encoding each - object against a similar object (usually its own previous version) — a - cross-row relationship a per-row compressor structurally cannot see, as - migration `0012`'s own comment already says. + Git's ~3GB for nixpkgs comes almost entirely from delta-encoding each object + against a similar object (usually its own previous version) — a cross-row + relationship a per-row compressor structurally cannot see, as migration + `0012`'s own comment already says. - **Measured** (this session, real nixpkgs slice, see below): a general, git-independent reference-delta scheme (zstd, one object compressed against its predecessor as a compression dictionary) gets **156x** aggregate @@ -31,22 +36,22 @@ git-specific opcode set. **4.7x** and **1.0x** respectively for the same objects compressed alone. That's a **33x** and **6x** reduction in stored bytes from adding the cross-row relationship alone, holding the compressor constant. -- **Measured**, a content-defined-chunking (CDC) + chunk-dedup design - (already scaffolded in this repo's `examples/combined_bench.rs`, never - previously run to completion) gets only **3.27x** on the same corpus — better - than plain per-object lz4 (2.24x) but far short of reference-delta, and for - a structural reason specific to nixpkgs' content shape (below), not a bug in - the CDC implementation. +- **Measured**, a content-defined-chunking (CDC) + chunk-dedup design (already + scaffolded in this repo's `examples/combined_bench.rs`, never previously run + to completion) gets only **3.27x** on the same corpus — better than plain + per-object lz4 (2.24x) but far short of reference-delta, and for a structural + reason specific to nixpkgs' content shape (below), not a bug in the CDC + implementation. - **Throughput and compression are the same lever here, not two problems.** - The insert phase is CPU-bound on TOAST compression of large payloads - (per `docs/handoff-push-performance.md`); shrinking what gets compressed - and written shrinks both the disk footprint and the CPU time that dominates - push latency. -- Recommendation: **reference-delta storage using a standard compression - codec's dictionary/prefix mode (zstd), not git's ofs-delta/ref-delta pack - opcodes** — same mechanism git uses (delta against a similar prior object), - portable encoding. CDC gets a narrow, secondary role. Full design and - rollout plan below. + The insert phase is CPU-bound on TOAST compression of large payloads (per + `docs/handoff-push-performance.md`); shrinking what gets compressed and + written shrinks both the disk footprint and the CPU time that dominates push + latency. +- Recommendation: **reference-delta storage using a standard compression codec's + dictionary/prefix mode (zstd), not git's ofs-delta/ref-delta pack opcodes** + — same mechanism git uses (delta against a similar prior object), portable + encoding. CDC gets a narrow, secondary role. Full design and rollout plan + below. ## Measured baseline: what's actually being compared @@ -66,15 +71,15 @@ reproducible independent of any earlier session's disk state). The git-native number is the ceiling this whole exercise is chasing, not a target to literally replicate — git's mechanism is the right reference point -for "how much cross-row redundancy exists here," independent of whether we -reuse git's specific pack format. +for "how much cross-row redundancy exists here," independent of whether we reuse +git's specific pack format. ### Why CDC underperforms here (measured, not assumed) -The `combined_bench.rs` example already existed in the working copy -(previous session scaffolded it — `fastcdc`/`blake3` dev-deps were added for -exactly this — but it had never been run to completion; no result was -recorded anywhere). Running it against the full corpus: +The `combined_bench.rs` example already existed in the working copy (previous +session scaffolded it — `fastcdc`/`blake3` dev-deps were added for exactly this +— but it had never been run to completion; no result was recorded anywhere). +Running it against the full corpus: ``` n=150338 total_logical_bytes=1284078078 (1224.6 MB) concurrency=8 @@ -86,32 +91,32 @@ compression+dedup ratio: 3.27x Two things stand out: 1. **The reconstruction bookkeeping (`object_chunks`, 187.2MB) is bigger than - the compressed content itself (`chunks`, 147.2MB).** At 1024B average - chunk size, 150k objects fragment into 1.166M chunk references, each - needing an `(oid, seq, chunk_hash)` row. The join-table overhead eats most - of the chunking win before the ratio is even computed. + the compressed content itself (`chunks`, 147.2MB).** At 1024B average chunk + size, 150k objects fragment into 1.166M chunk references, each needing + an `(oid, seq, chunk_hash)` row. The join-table overhead eats most of the + chunking win before the ratio is even computed. 2. **More fundamentally: CDC dedups *within* a chunking window, not against an explicit prior version.** nixpkgs is dominated by small `.nix`/`.patch` - files — the corpus median blob is ~1KB, well under the 4KB max chunk size - — so a typical file is exactly *one* chunk. A one-line edit changes that - file's bytes, which changes that whole chunk's content hash, which means - the "new version" chunk shares nothing with the "old version" chunk even - though 99% of the file is byte-identical. CDC's rolling-hash boundary - trick is designed to re-synchronize chunk boundaries after an insertion or - deletion *inside a large blob*; it doesn't help when the unit of change - *is* the whole blob. This is a structural property of the corpus (many - small files, temporal near-duplicate versions), not a parameter-tuning - problem — smaller average chunk size would mean more objects fragment - into multiple chunks (helping dedup) but multiplies the already-dominant - per-chunk bookkeeping overhead (hurting it worse). CDC's actual strength — - dedup across *shifted* regions inside large mutable blobs (VM images, - tarballs, database dumps) — isn't what nixpkgs' size is made of. - -This matters for the "don't just copy or just reject git" framing: git's -own delta search is doing something CDC structurally cannot — finding -*temporally or content-similar whole objects* and diffing against them, -not chunking each object in isolation. That's the mechanism worth -reproducing; CDC is solving a different problem. + files — the corpus median blob is ~1KB, well under the 4KB max chunk size — + so a typical file is exactly *one* chunk. A one-line edit changes that file's + bytes, which changes that whole chunk's content hash, which means the "new + version" chunk shares nothing with the "old version" chunk even though 99% + of the file is byte-identical. CDC's rolling-hash boundary trick is designed + to re-synchronize chunk boundaries after an insertion or deletion *inside + a large blob*; it doesn't help when the unit of change *is* the whole blob. + This is a structural property of the corpus (many small files, temporal + near-duplicate versions), not a parameter-tuning problem — smaller average + chunk size would mean more objects fragment into multiple chunks (helping + dedup) but multiplies the already-dominant per-chunk bookkeeping overhead + (hurting it worse). CDC's actual strength — dedup across *shifted* regions + inside large mutable blobs (VM images, tarballs, database dumps) — isn't what + nixpkgs' size is made of. + +This matters for the "don't just copy or just reject git" framing: git's own +delta search is doing something CDC structurally cannot — finding *temporally +or content-similar whole objects* and diffing against them, not chunking each +object in isolation. That's the mechanism worth reproducing; CDC is solving a +different problem. ### Reference-delta: methodology and honest error bars @@ -124,25 +129,23 @@ above (156.3x) is from **the full population, not a sample**: - Extracted every `M`-status (modified, not added/deleted) path transition across the corpus's full commit history via `git log --raw --no-abbrev`: **28,321 old-blob→new-blob pairs** (28,317 resolved; 4 unreadable, likely - submodule gitlinks). - These pairs' new-side bytes total 1081.5MB — **95.6% of all blob bytes in - the corpus** (1131MB), despite being only ~81% of the blob *count* - (28,317 of 34,970) — larger, more-frequently-touched files (generated - package lists, etc.) dominate byte volume and are exactly the files that - get re-touched across commits. -- For each pair: compressed the new blob alone with zstd level 9 (proxy for - "current independent-object compression," comparable in spirit to lz4 - TOAST — zstd-9 is a fair, slightly-generous ceiling for what an - independent per-object codec can do), then compressed the new blob again - using the **old blob's raw bytes as a zstd prefix/dictionary** - (`ZstdCompressionDict(..., dict_type=DICT_TYPE_RAWCONTENT)`, - `ZSTD_CCtx_refPrefix` under the hood) — no dictionary *training*, just - "diff against this literal prior content," which is what `zstd - --patch-from` does at the CLI. -- Result, summed across all 28,317 pairs (aggregate ratio = sum of - compressed bytes over sum of logical bytes, not a mean of per-pair - ratios — a mean would overweight tiny near-empty diffs the same way a - per-object mean misrepresents git's own whole-corpus 32.9x): + submodule gitlinks). These pairs' new-side bytes total 1081.5MB — **95.6% of + all blob bytes in the corpus** (1131MB), despite being only ~81% of the blob + *count* (28,317 of 34,970) — larger, more-frequently-touched files (generated + package lists, etc.) dominate byte volume and are exactly the files that get + re-touched across commits. +- For each pair: compressed the new blob alone with zstd level 9 (proxy + for "current independent-object compression," comparable in spirit to lz4 + TOAST — zstd-9 is a fair, slightly-generous ceiling for what an independent + per-object codec can do), then compressed the new blob again using the **old + blob's raw bytes as a zstd prefix/dictionary** (`ZstdCompressionDict(..., + dict_type=DICT_TYPE_RAWCONTENT)`, `ZSTD_CCtx_refPrefix` under the hood) — no + dictionary *training*, just "diff against this literal prior content," which + is what `zstd --patch-from` does at the CLI. +- Result, summed across all 28,317 pairs (aggregate ratio = sum of compressed + bytes over sum of logical bytes, not a mean of per-pair ratios — a mean would + overweight tiny near-empty diffs the same way a per-object mean misrepresents + git's own whole-corpus 32.9x): ``` total new logical: 1,134,013,297 bytes (1081.5 MB) @@ -155,70 +158,68 @@ above (156.3x) is from **the full population, not a sample**: relationship, not a stronger compressor. - Same methodology on the 16,672 root-tree transitions (commit tree vs. first-parent's tree, excluding no-op pairs) gave 6.27x aggregate vs. 1.02x - standalone — trees are tiny (median 227 bytes; total sample only 3.6MB) - so the win is smaller in absolute terms and zstd's frame overhead eats - more of it, but the *direction* holds: one entry changes, the rest of the - tree is byte-identical, and a reference codec sees that where an - independent one can't. + standalone — trees are tiny (median 227 bytes; total sample only 3.6MB) so + the win is smaller in absolute terms and zstd's frame overhead eats more + of it, but the *direction* holds: one entry changes, the rest of the tree + is byte-identical, and a reference codec sees that where an independent one + can't. - **What wasn't measured, flagged rather than assumed:** nested (non-root) trees, and commits. Commits are unlikely to delta well (free-text messages, dates, and author lines dominate their bytes and differ between any two commits) — expect commits to behave like the "standalone" column, not the - "delta" column. Nested trees should behave like root trees structurally - (a directory listing that's mostly unchanged) but weren't independently - verified. **[INFERENCE]** where marked below. + "delta" column. Nested trees should behave like root trees structurally (a + directory listing that's mostly unchanged) but weren't independently verified. + **[INFERENCE]** where marked below. -Rough blended estimate for blobs specifically, combining what was measured: -delta-eligible blob bytes (1081.5MB) at 156.3x → ~7.3MB, plus the +Rough blended estimate for blobs specifically, combining what was +measured: delta-eligible blob bytes (1081.5MB) at 156.3x → ~7.3MB, plus the non-delta-eligible remainder (1131 − 1081.5 ≈ 49.5MB, first-appearance blobs -with no prior version) at the measured 4.74x standalone rate → ~10.4MB. -**Blobs alone: ~1131MB → ~17.7MB, a ~64x reduction** — well past CDC's 3.27x -and closing most of the distance to git's whole-corpus 32.9x (which also -includes trees and commits, where the picture is less complete per above). -This is an estimate built from two real measurements, not a third -measurement — labeled as such. +with no prior version) at the measured 4.74x standalone rate → \~10.4MB. **Blobs +alone: \~1131MB → \~17.7MB, a \~64x reduction** — well past CDC's 3.27x and +closing most of the distance to git's whole-corpus 32.9x (which also includes +trees and commits, where the picture is less complete per above). This is an +estimate built from two real measurements, not a third measurement — labeled +as such. ## Why the current design gets none of this `mnemosyne_git::OwnedObject`'s doc comment already says the quiet part: -"decoded, uncompressed object bytes." `mnemosyne_postgres`'s `object_insert` -/ `object_pool_insert` hash the object, then write full content-addressed -bytes with `ON CONFLICT DO NOTHING`. `unpack_into_tx` in `receive_pack.rs` -walks the incoming pack via `gix_pack::Bundle`, and for every entry — -including ones the client sent as a **delta against a base it already -knew we could resolve** — calls `get_object_by_index`, which fully resolves -the delta chain to plain bytes before anything is stored -(`receive_pack.rs:630`, "decode every entry... resolving intra-pack -deltas... into memory first, then insert"). The client already did the -expensive similarity search and delta encoding; the server immediately -throws it away and TOAST-compresses each resulting blob in isolation. -Migration `0012`'s own comment states this explicitly as the acknowledged, -deferred gap: "This does NOT close the gap with git's own packfile -compression... Delta compression against prior object versions is a -separate, larger feature (tracked separately), not a column setting." - -Confirmed via `gix-pack`'s public API (`data::File::entry(offset)` returns a -`data::Entry` whose `Header` is `RefDelta { base_id }` / `OfsDelta { -base_distance }` / a plain kind — the delta relationship is available -*before* full resolution): the pack-ingest code has always had access to -"this object is a diff against that object," and has always discarded it. +"decoded, uncompressed object bytes." `mnemosyne_postgres`'s `object_insert` / +`object_pool_insert` hash the object, then write full content-addressed bytes +with `ON CONFLICT DO NOTHING`. `unpack_into_tx` in `receive_pack.rs` walks +the incoming pack via `gix_pack::Bundle`, and for every entry — including +ones the client sent as a **delta against a base it already knew we could +resolve** — calls `get_object_by_index`, which fully resolves the delta chain +to plain bytes before anything is stored (`receive_pack.rs:630`, "decode every +entry... resolving intra-pack deltas... into memory first, then insert"). The +client already did the expensive similarity search and delta encoding; the +server immediately throws it away and TOAST-compresses each resulting blob +in isolation. Migration `0012`'s own comment states this explicitly as the +acknowledged, deferred gap: "This does NOT close the gap with git's own packfile +compression... Delta compression against prior object versions is a separate, +larger feature (tracked separately), not a column setting." + +Confirmed via `gix-pack`'s public API (`data::File::entry(offset)` returns +a `data::Entry` whose `Header` is `RefDelta { base_id }` / `OfsDelta { +base_distance }` / a plain kind — the delta relationship is available *before* +full resolution): the pack-ingest code has always had access to "this object is +a diff against that object," and has always discarded it. ## The proposed mechanism **Store the content-addressed object pool as it is today, but let an object record itself as a diff against another object in the same pool, using a -standard reference-compression codec rather than a bespoke or -git-pack-specific delta format.** +standard reference-compression codec rather than a bespoke or git-pack-specific +delta format.** ### Codec choice: zstd prefix/dictionary mode, not git's ofs-delta/ref-delta This is the "don't copy git, don't reject git" decision point, made on the merits: -- **What git does**: a custom binary delta format (copy/insert opcodes - against a designated base, RFC-less, git-specific) chosen via similarity - search (same-path heuristic + a sliding window over size-bucketed - candidates). +- **What git does**: a custom binary delta format (copy/insert opcodes against + a designated base, RFC-less, git-specific) chosen via similarity search + (same-path heuristic + a sliding window over size-bucketed candidates). - **What's proposed instead**: `zstd`'s standardized prefix/reference-dictionary mode (`ZSTD_CCtx_refPrefix`, stable since zstd 1.4.0, exposed in Rust via `zstd-safe`/`zstd-rs`; this is the mechanism behind the `zstd --patch-from` @@ -226,30 +227,29 @@ merits: bytes as a one-shot dictionary. Decompression needs the same base bytes supplied externally; the stored frame doesn't embed them. VCDIFF/RFC 3284 (`xdelta3`, `vcdiff` crates exist for Rust) is the other standards-track - option and was considered — zstd's prefix mode was chosen because it's - already a workspace-adjacent, actively maintained codec family (this repo - already ships `flate2` for a different purpose) with mature Rust bindings, - it gets compression *and* the delta relationship in one step instead of - two, and its output format has no git-specific structure baked in — an hg - importer or a hypothetical jj-native backend populating the same pool - would use the identical mechanism with their own choice of "prior version." + option and was considered — zstd's prefix mode was chosen because it's already + a workspace-adjacent, actively maintained codec family (this repo already + ships `flate2` for a different purpose) with mature Rust bindings, it gets + compression *and* the delta relationship in one step instead of two, and + its output format has no git-specific structure baked in — an hg importer + or a hypothetical jj-native backend populating the same pool would use the + identical mechanism with their own choice of "prior version." - **Base selection**: initially, temporal — "this object's most recent - predecessor at the same logical position" (same git blob replaced at the - same tree path; same file revision in an hg import; whatever the - originating VCS's notion of "previous version of this thing" is). This is - deliberately simpler than git's similarity-bucket search and is what was - actually measured above. A same-size/similar-content search (closer to - git's own heuristic) is a valid later refinement if temporal-only leaves - gains on the table for non-linear history (renames, cherry-picks) — not - needed to capture the bulk of the measured win, which comes from the - ordinary "modify this file again" case. -- **Bounded chain depth**, same reason git caps at 50 (confirmed via - `git verify-pack -v` on this exact corpus: chain lengths run 1 through 50 - in a roughly geometric decay, 52,035 objects — 35% — stored as full - objects with no delta at all): reconstruction cost is `O(chain depth)` - decompressions, so an unbounded chain makes old-object reads arbitrarily - slow. A depth cap forces a periodic full-object snapshot, same tradeoff - git makes, independent of git's specific format. + predecessor at the same logical position" (same git blob replaced at the same + tree path; same file revision in an hg import; whatever the originating VCS's + notion of "previous version of this thing" is). This is deliberately simpler + than git's similarity-bucket search and is what was actually measured above. + A same-size/similar-content search (closer to git's own heuristic) is a valid + later refinement if temporal-only leaves gains on the table for non-linear + history (renames, cherry-picks) — not needed to capture the bulk of the + measured win, which comes from the ordinary "modify this file again" case. +- **Bounded chain depth**, same reason git caps at 50 (confirmed via `git + verify-pack -v` on this exact corpus: chain lengths run 1 through 50 in a + roughly geometric decay, 52,035 objects — 35% — stored as full objects with no + delta at all): reconstruction cost is `O(chain depth)` decompressions, so an + unbounded chain makes old-object reads arbitrarily slow. A depth cap forces a + periodic full-object snapshot, same tradeoff git makes, independent of git's + specific format. ### Schema sketch @@ -265,180 +265,236 @@ ALTER TABLE objects ADD COLUMN chain_depth SMALLINT NOT NULL DEFAULT 0; ``` `base_oid` is a self-FK, so `ON DELETE RESTRICT` (already the pattern -`repo_objects → objects` uses) protects a base from deletion while any -delta still points at it — GC has to either promote a delta to -self-contained or cascade-delete the whole dependent chain before removing -a base, exactly the bookkeeping git's own repack/prune already has to do. -This is real, deferred complexity — flagged, not solved, below. +`repo_objects → objects` uses) protects a base from deletion while any delta +still points at it — GC has to either promote a delta to self-contained or +cascade-delete the whole dependent chain before removing a base, exactly the +bookkeeping git's own repack/prune already has to do. This is real, deferred +complexity — flagged, not solved, below. ### Read path -`object_find`'s query changes from "return `data` directly" to "walk -`base_oid` up to a self-contained ancestor (bounded by chain depth), -decompress each hop with the previous hop's *decoded* bytes as the zstd -prefix, in base-to-tip order." Cost is `O(chain depth)` zstd-prefix -decompressions instead of one lz4 TOAST decompression — real but small -(zstd decompression with a prefix is fast; the current benchmarks in this -repo don't have a number for it yet, so this is a **[design implication, not -yet benchmarked]** point to validate before shipping). +`object_find`'s query changes from "return `data` directly" to "walk `base_oid` +up to a self-contained ancestor (bounded by chain depth), decompress each hop +with the previous hop's *decoded* bytes as the zstd prefix, in base-to-tip +order." Cost is `O(chain depth)` zstd-prefix decompressions instead of one lz4 +TOAST decompression — real but small (zstd decompression with a prefix is fast; +the current benchmarks in this repo don't have a number for it yet, so this +is a **[design implication, not yet benchmarked]** point to validate before +shipping). `upload_pack.rs` / `pack.rs` currently re-derive delta encoding from scratch -on every fetch via `gix_pack`'s own similarity search over fully-resolved -bytes (`iter_from_counts`, `Mode::PackCopyAndBaseObjects`) — meaning today's -outgoing packs pay full delta-search cost per request and don't reuse -anything from how the object was stored. Once storage itself carries a -delta relationship, there's a follow-on option (not required for the storage -win, a separate optimization) to feed that relationship into pack -generation directly instead of re-deriving it — deferred, flagged as future -work, not scoped here. +on every fetch via `gix_pack`'s own similarity search over fully-resolved bytes +(`iter_from_counts`, `Mode::PackCopyAndBaseObjects`) — meaning today's outgoing +packs pay full delta-search cost per request and don't reuse anything from how +the object was stored. Once storage itself carries a delta relationship, there's +a follow-on option (not required for the storage win, a separate optimization) +to feed that relationship into pack generation directly instead of re-deriving +it — deferred, flagged as future work, not scoped here. ### Ingest path `unpack_into_tx`'s decode loop (`receive_pack.rs:618`) currently calls -`get_object_by_index`, which resolves every entry to full bytes -unconditionally. The change: inspect `bundle.pack.entry(offset).header` -*before* fully resolving. For a `RefDelta`/`OfsDelta` entry, the client -already told us its base relationship — after computing the base's full -bytes (needed anyway, either from this same pack or via `PushTxFind` -against already-stored objects for thin-pack bases, both already resolved -today), compress the new object against that base with the zstd-prefix -mechanism and store `(oid, base_oid, delta_bytes)` instead of `(oid, -full_bytes)`. For a non-delta entry, store self-contained as today (still -zstd/lz4-compressed alone). This changes *what* gets compressed and stored, -not the pack-parsing or decode-parallelism structure already in place — the -existing per-thread decode loop, `insert_objects_parallel`'s connection +`get_object_by_index`, which resolves every entry to full bytes unconditionally. +The change: inspect `bundle.pack.entry(offset).header` *before* fully resolving. +For a `RefDelta`/`OfsDelta` entry, the client already told us its base +relationship — after computing the base's full bytes (needed anyway, either from +this same pack or via `PushTxFind` against already-stored objects for thin-pack +bases, both already resolved today), compress the new object against that base +with the zstd-prefix mechanism and store `(oid, base_oid, delta_bytes)` instead +of `(oid, full_bytes)`. For a non-delta entry, store self-contained as today +(still zstd/lz4-compressed alone). This changes *what* gets compressed and +stored, not the pack-parsing or decode-parallelism structure already in place +— the existing per-thread decode loop, `insert_objects_parallel`'s connection fan-out, and the `synchronous_commit=off` batching all carry over unchanged. ## Throughput: the same fix attacks the other bottleneck too Per `docs/handoff-push-performance.md`, the insert phase is 66% of push wall -time and is **CPU-bound on TOAST compression of large payloads**, not -network- or round-trip-bound — that's why parallelizing across connections -helped and why batching/COPY didn't (measured and rejected in that session). -Reference-delta storage reduces bytes-to-compress-and-write by the same -~33-64x measured above *before* they ever reach Postgres — smaller payloads -mean less TOAST/zstd CPU work per object, attacking the documented -bottleneck directly rather than adding a new one alongside it. This is -speculative in degree (not yet benchmarked end-to-end: computing the -zstd-prefix delta client-side, in the existing parallel decode threads, adds -its own CPU cost that has to be measured against the TOAST-compression cost -it removes) but sound in direction — compressing 17.7MB of eventual blob -payload costs less wall-clock than compressing 1131MB of it, on any codec. +time and is **CPU-bound on TOAST compression of large payloads**, not network- +or round-trip-bound — that's why parallelizing across connections helped and why +batching/COPY didn't (measured and rejected in that session). Reference-delta +storage reduces bytes-to-compress-and-write by the same ~33-64x measured above +*before* they ever reach Postgres — smaller payloads mean less TOAST/zstd CPU +work per object, attacking the documented bottleneck directly rather than adding +a new one alongside it. This is speculative in degree (not yet benchmarked +end-to-end: computing the zstd-prefix delta client-side, in the existing +parallel decode threads, adds its own CPU cost that has to be measured against +the TOAST-compression cost it removes) but sound in direction — compressing +17.7MB of eventual blob payload costs less wall-clock than compressing 1131MB of +it, on any codec. Independent of the delta-storage question, two more direct throughput levers from the existing investigation are still open and worth doing regardless of which storage design ships: -1. **Run the actual full nixpkgs corpus (8.36M objects) end-to-end.** - Every throughput number in `docs/handoff-push-performance.md` and in this - doc is extrapolated from the 150k-object slice. The handoff doc already - flags this as the top next step; it's still not done, and insert - throughput may not scale linearly once the `objects` table + indexes - exceed `shared_buffers` (measured at 4GB currently, likely far smaller - than an 8.36M-object working set) — this is exactly the kind of - measurement-vs-extrapolation gap that already bit this project once - (the `shared_buffers` 25x error the handoff doc records). +1. **Run the actual full nixpkgs corpus (8.36M objects) end-to-end.** Every + throughput number in `docs/handoff-push-performance.md` and in this doc is + extrapolated from the 150k-object slice. The handoff doc already flags this + as the top next step; it's still not done, and insert throughput may not + scale linearly once the `objects` table + indexes exceed `shared_buffers` + (measured at 4GB currently, likely far smaller than an 8.36M-object working + set) — this is exactly the kind of measurement-vs-extrapolation gap that + already bit this project once (the `shared_buffers` 25x error the handoff + doc records). 2. **Convert the `eprintln!("[timing] ...")` calls to real tracing spans.** Flagged as step 1 in the handoff doc's own "Next steps," not yet done. - Low-risk, unlocks permanent phase-breakdown visibility instead of - one-off hand instrumentation. + Low-risk, unlocks permanent phase-breakdown visibility instead of one-off + hand instrumentation. ## Cross-transport compatibility, addressed directly -The requester's constraint: stay open to jj (no native wire protocol -yet) and, aspirationally, Mercurial, without treating "git does it" as -either a mandate or a disqualifier. +The requester's constraint: stay open to jj (no native wire protocol yet) and, +aspirationally, Mercurial, without treating "git does it" as either a mandate or +a disqualifier. -- **jj today**: only its git backend is production-ready (confirmed by - direct research this session); the native backend is explicitly +- **jj today**: only its git backend is production-ready (confirmed + by direct research this session); the native backend is explicitly not-ready-for-real-use, and there is no native jj wire protocol yet — jj's - actual network transport, today, *is* git's smart HTTP/SSH. So the - concrete, non-speculative jj-compatibility requirement right now is: - **be a correct git smart-HTTP/SSH server**, which is already this - project's design goal independent of storage. `upload_pack.rs` currently - only advertises protocol v1 (no `version=2` capability found in the - protocol crate) — full-closure clone with no `have`/negotiation, called - out in that file's own doc comment as a v1 simplification. Worth fixing - on its own merits (bandwidth, not storage), but orthogonal to this doc's - compression question — noted here because it's the actual current gap in - "does this serve jj well," not the storage format. -- **The storage-layer decision that matters for the future**: keep the - delta *mechanism* (reference compression against a chosen base) decoupled - from git's *specific encoding* (ofs-delta/ref-delta opcodes, pack-file - entry framing). The zstd-prefix design above never touches git's delta - bytecode — it re-derives the base relationship from the incoming pack's - metadata, then encodes the diff in a codec-standard, non-git frame. That - means the same object pool and delta schema could, in principle, be - populated by a Mercurial revlog importer (translate hg's own delta chain - into "which prior revision is this most like," reuse the same - zstd-prefix storage call) or by a native jj backend once one exists, - without either of them needing to understand git's pack format at all. - This is the concrete form of "don't copy git's mechanism by default, don't - reject it either" — the *technique* (delta against a similar object) is - adopted because the data justifies it (33-156x measured), the *encoding* - is deliberately generalized because nothing about the technique requires - git's specific bytes. + actual network transport, today, *is* git's smart HTTP/SSH. So the concrete, + non-speculative jj-compatibility requirement right now is: **be a correct + git smart-HTTP/SSH server**, which is already this project's design goal + independent of storage. `upload_pack.rs` currently only advertises protocol v1 + (no `version=2` capability found in the protocol crate) — full-closure clone + with no `have`/negotiation, called out in that file's own doc comment as a v1 + simplification. Worth fixing on its own merits (bandwidth, not storage), but + orthogonal to this doc's compression question — noted here because it's the + actual current gap in "does this serve jj well," not the storage format. +- **The storage-layer decision that matters for the future**: keep the delta + *mechanism* (reference compression against a chosen base) decoupled from + git's *specific encoding* (ofs-delta/ref-delta opcodes, pack-file entry + framing). The zstd-prefix design above never touches git's delta bytecode — + it re-derives the base relationship from the incoming pack's metadata, then + encodes the diff in a codec-standard, non-git frame. That means the same + object pool and delta schema could, in principle, be populated by a Mercurial + revlog importer (translate hg's own delta chain into "which prior revision is + this most like," reuse the same zstd-prefix storage call) or by a native jj + backend once one exists, without either of them needing to understand git's + pack format at all. This is the concrete form of "don't copy git's mechanism + by default, don't reject it either" — the *technique* (delta against a similar + object) is adopted because the data justifies it (33-156x measured), the + *encoding* is deliberately generalized because nothing about the technique + requires git's specific bytes. - CDC is not discarded from this picture — it's demoted to a secondary, - optional layer for large blobs (nixpkgs source tarballs vendored as git - blobs, generated lockfiles) above a size threshold (e.g. 256KB), where its - actual strength — dedup across shifted regions within a large mutable - blob — applies. Not scoped as part of the initial rollout below; revisit - once the primary reference-delta win is shipped and measured. + optional layer for large blobs (nixpkgs source tarballs vendored as git blobs, + generated lockfiles) above a size threshold (e.g. 256KB), where its actual + strength — dedup across shifted regions within a large mutable blob — applies. + Not scoped as part of the initial rollout below; revisit once the primary + reference-delta win is shipped and measured. ## What this doc is not -- Not a decision to migrate today. This is the analysis the requester asked - for; the schema change, GC refactor, and read-path changes are real, - non-trivial work that should be scoped and approved as their own effort. -- Not a claim that 64x-on-blobs generalizes to a validated whole-corpus - number — trees (non-root) and commits are explicitly unmeasured, flagged - above rather than assumed. -- Not a rejection of `docs/sharding-and-tiered-storage.md` — that doc's own - "what this does not buy us" section already says sharding buys no new - compression; this doc is the compression side that doc explicitly - deferred. - -## Proposed staged rollout, if approved - -1. **Benchmark zstd-prefix decode cost** (not yet measured): how expensive - is walking a bounded delta chain on the read path, at realistic chain - depths, against the lz4-TOAST baseline it replaces. Gates whether the - read-path cost is actually acceptable before committing to the schema - change. -2. **Schema migration**: `base_oid` + `chain_depth` columns, as sketched - above, additive and backward-compatible (existing rows stay - self-contained, `base_oid IS NULL`). -3. **Ingest-path change**: `unpack_into_tx` inspects entry headers before - resolving, computes delta relationship + zstd-prefix payload for - delta-eligible entries. Ship behind a flag; validate against the - existing `storage_dedup.rs` correctness tests plus new ones for - delta-chain round-tripping. -4. **Read-path change**: `object_find` walks the chain. Validate byte-exact - round-trip (write → read → identical to original) as a property test — - this is exactly the kind of oracle-backed correctness check the `rigor` - skill calls for: "any object written through the delta path and read - back must produce byte-identical content to what was written," checked - against arbitrary chain depths up to the cap. -5. **GC/refcount design** for `base_oid`'s `ON DELETE RESTRICT` — needed - before any deletion path exists, not needed to ship the write/read path - itself (today's deferred-GC posture, per migration `0007`'s comment, - already tolerates no-deletion for a while). -6. **Full nixpkgs corpus run**, both before and after, to replace every - extrapolated number in this doc and `docs/handoff-push-performance.md` - with a measured one. - -## Open questions for the operator - -1. Given the mentioned possibility of a hand rewrite: is this doc's design - meant to inform that rewrite (requirements capture, like the prior - session's handoff docs), or should it be implemented against the current - codebase now? Affects whether step 2 onward happens at all in this repo. -2. Base-selection heuristic: is temporal-only (same path, most recent - version) sufficient, or is a similarity-search fallback (closer to git's - own bucketed search, needed for renames/cherry-picks the temporal - heuristic would miss) worth the added complexity up front rather than as - a later refinement? -3. How much of the CDC scaffolding (`combined_bench.rs`, the `fastcdc`/ - `blake3` dev-deps) should be kept as the future "large blob" secondary - layer vs. removed now that its primary-mechanism role is measured and - rejected for this corpus? +- Not a claim that 64x-on-blobs generalizes to a validated whole-corpus number + on nixpkgs specifically — that projection was never run against the full + corpus (this repo's own history, ~800 objects, was used instead per the + operator's disk-space constraint; see "Implementation status" below). + Non-root trees and commits remain unmeasured for delta-compressibility on + either corpus. +- Not a rejection of `docs/sharding-and-tiered-storage.md` — that doc's + own "what this does not buy us" section already says sharding buys no new + compression; this doc is the compression side that doc explicitly deferred. + +## Implementation status + +Implemented against the current codebase (not deferred to a hypothetical +rewrite — the operator confirmed this repo is the target and directed testing +against a small local repo rather than nixpkgs, given local disk constraints). + +**Shipped:** + +- Migration `0013_objects_reference_delta.sql`: `base_oid` (nullable, no FK — + see the migration's own comment for why a FK would reintroduce an ordering + requirement `insert_objects_parallel`'s sharding deliberately doesn't have) + and `chain_depth` columns on `objects`. +- `mnemosyne_postgres::delta`: `encode_against_base`/`decode_against_base` + using zstd's prefix mode (`zstd-safe`'s `CCtx::ref_prefix` + + `DCtx::ref_prefix`, per this doc's codec-choice rationale above). Property- + tested (byte-exact round trip for arbitrary inputs, a near-identical-content + compression-ratio regression guard, and empty-input edge cases). +- Ingest path (`receive_pack::unpack_into_tx`): captures each pack entry's + `Header` (RefDelta/OfsDelta/plain) during the existing parallel decode pass, + then a single-threaded resolution pass picks each entry's effective base + (in-pack only — a `RefDelta` to an external thin-pack base falls back to + self-contained, a deliberately narrow scope for this pass) and chain depth, + memoized and cycle-guarded, capped at `MAX_CHAIN_DEPTH` (50, matching git's + own default). Delta encoding is skipped in favor of self-contained storage + if the encoded delta wouldn't actually be smaller. +- Read path (`object_find`): rewritten as a recursive CTE that walks the + chain base-to-tip in one round trip, then folds `decode_against_base` over + the rows. **A second raw reader, `object_pool_find_many` (the parallel + connectivity-check's batched lookup), was found and fixed during + verification** — it read `objects.data` directly with no chain resolution, + which surfaced as `git-receive-pack` rejecting pushes with "failed to + decode Tree object" once real delta rows existed. Fixed with the same + chain-walk technique, generalized to resolve many requested oids' chains in + one query. This was the one real bug this implementation pass found — + flagged here because it's the general lesson: **any code path that reads + `objects.data` directly, bypassing `object_find`, is latent corruption + once delta rows exist.** A repo-wide audit at the time found exactly three + raw `FROM objects` readers (`object_find`, `object_pool_find_many`, + `object_header` — the last is safe, since `kind`/`size` are always logical + values regardless of storage representation); if a fourth is ever added, it + needs the same audit. +- `mnemosyne_postgres/tests/delta_round_trip.rs`: the actual oracle for this + feature — drives the real production write path (`insert_objects_parallel` + + `PushTx::insert_repo_inclusions`, not a shortcut), asserts `base_oid IS + NOT NULL` in raw SQL (so a silent fallback to self-contained can't pass for + the wrong reason), and asserts byte-identical reconstruction through + `PgBackend::try_find` — for both a single-hop and a two-hop chain. + +**Verified, real numbers (not projected):** + +Pushed mnemosyne's own git history (793 objects, real Rust source + docs + +migrations, not a synthetic fixture) through the actual `mnemosyne` server +binary over real SSH, using a real `git push`, per the operator's disk-space- +conscious choice of test repo. The `objects` table total (796) is 3 higher +than the pushed count (793) — those 3 are the demo repo's seeded +commit/tree/blob from server startup, sharing the pool; the breakdown below +is the table-wide total (queried directly via SQL after the push), not +re-filtered to the pushed set specifically: + +| metric | value | +|---|---| +| objects pushed (this repo's history) | 793 | +| `objects` table total (incl. 3 pre-existing seeded objects) | 796 | +| stored as delta (`base_oid IS NOT NULL`) | 406 (51%) | +| stored self-contained | 390 | +| logical bytes (table-wide) | 6,297,045 | +| stored bytes (table-wide, `sum(pg_column_size(data))`) | 699,035 | +| **compression ratio (logical/stored)** | **9.01x** | +| `objects` table total size (incl. indexes) | 1,000 KB | +| git's own native bundle for the pushed content | 496,632 bytes | + +For comparison: the prior no-delta shape (migration `0012`, per-object lz4 +TOAST only) measured 2.24x on the nixpkgs corpus in this doc's earlier +analysis section — this real push shows delta storage landing at 9.01x, and +the full on-disk `objects` relation (with indexes) at roughly 2x git's own +native pack size for the identical content, down from the ~13x-larger-than-git +baseline this whole investigation started from. + +Full integrity verified: `git fsck --full` on a fresh reclone from the +server reported zero corruption; reclone `HEAD` matched the pushed `HEAD` +exactly; the full integration test suite (`mnemosyne_tests`, protocol + +storage_dedup subset) passes with no regressions attributable to this change +— one pre-existing `PoolTimedOut` flake was confirmed to reproduce +identically on the pre-session baseline commit, unrelated to this work. + +**Deliberately deferred, not solved by this pass:** + +- **GC/refcounting** for `base_oid` — no FK, so nothing currently prevents + deleting a base out from under a dependent delta. Matches migration + `0007`'s already-deferred GC posture (no deletion path exists yet at all), + not a new gap this feature introduces. +- **Base selection stays temporal/in-pack-only** — no similarity-search + fallback for renames or cross-pack thin-pack bases. The 9.01x measured + above is the floor this simpler heuristic achieves; a similarity search + closer to git's own bucketed approach is the natural next refinement if a + future corpus shows the temporal heuristic leaving gains on the table. +- **CDC scaffolding** (`combined_bench.rs`, `fastcdc`/`blake3` dev-deps) is + untouched — kept as the documented future "large blob" secondary layer per + this doc's "Cross-transport compatibility" section, not wired into the + shipped path. +- **Full nixpkgs corpus run** — not done, by explicit operator direction + (local disk constraints); the 9.01x above is real but on a much smaller, + differently-shaped corpus (personal codebase vs. a massive package + monorepo with heavier temporal blob churn). Directionally consistent with + the analysis section's projections, not a substitute for them. diff --git a/mnemosyne_postgres/Cargo.toml b/mnemosyne_postgres/Cargo.toml index 37d0a2c..233c5f6 100644 --- a/mnemosyne_postgres/Cargo.toml +++ b/mnemosyne_postgres/Cargo.toml @@ -23,8 +23,10 @@ gix-ref.workspace = true sqlx.workspace = true serde_json = "1" +zstd-safe.workspace = true [dev-dependencies] +proptest.workspace = true mnemosyne_harness.workspace = true anyhow.workspace = true rand = { version = "0.10", features = ["thread_rng"] } diff --git a/mnemosyne_postgres/migrations/0013_objects_reference_delta.sql b/mnemosyne_postgres/migrations/0013_objects_reference_delta.sql new file mode 100644 index 0000000..c6df1db --- /dev/null +++ b/mnemosyne_postgres/migrations/0013_objects_reference_delta.sql @@ -0,0 +1,67 @@ +-- Reference-delta storage for the object pool. +-- +-- Every object today is stored fully materialized and independently +-- TOAST-compressed (see `0012`'s comment: "no delta/pack compression +-- happens before it reaches Postgres"). Measured against a real +-- nixpkgs slice, that gets 2.24x on-disk vs. logical bytes, against +-- git's own ~32.9x for the same corpus — the entire gap is git's +-- delta compression against similar prior objects, a cross-row +-- relationship a per-row compressor structurally can't see. +-- +-- This adds an OPTIONAL delta representation: a row may record itself +-- as a diff against another row in the same pool (`base_oid`), +-- encoded with a general-purpose reference-compression codec (zstd's +-- prefix/dictionary mode: compress this object's bytes using the +-- base's decoded bytes as a one-shot dictionary) rather than git's +-- own ofs-delta/ref-delta pack opcodes — see +-- `docs/compression-and-throughput-strategy.md` for the full +-- rationale, including why the codec is deliberately NOT git-specific +-- (keeps the door open for other VCS backends/importers populating +-- the same pool without needing to understand git's pack format). +-- +-- `base_oid IS NULL` (the default, and every existing row) means +-- self-contained: `data` is the plain decoded bytes, TOAST-compressed +-- as before. `base_oid IS NOT NULL` means `data` is zstd-compressed +-- using the base row's fully-resolved bytes as a ref-prefix; reading +-- it requires resolving the base first (recursively, up to +-- `chain_depth` hops). +-- +-- Deliberately NOT a foreign key. `objects` rows land via +-- `insert_objects_parallel`, which shards work across several +-- connections that each auto-commit independently and in arbitrary +-- order specifically because pool writes have no cross-row ordering +-- requirement (see that function's doc comment in lib.rs). A +-- `REFERENCES objects(oid)` FK would silently reintroduce exactly +-- that ordering requirement — a delta row committing on one +-- connection before its base commits on another would intermittently +-- fail FK validation, depending on shard assignment. Same posture +-- `0007` already takes for `repo_objects → objects` in the other +-- direction being enforced only by insert order (pool row first, CTE +-- comment), not by a captured invariant here: base-before-delta +-- ordering is an application-level contract (the ingest path always +-- inserts a chain in base-to-tip order within one batch), not a +-- database-enforced one. Consequence: a delta row can transiently +-- point at a base that doesn't exist yet mid-batch (never durably — +-- by the time any push commits, the whole batch including bases is +-- already inserted) and, until GC/refcounting exists (still deferred +-- per `0007`), nothing stops deleting a base out from under a +-- dependent delta. Both are accepted for this exploratory phase; flagged, +-- not solved. +-- +-- `chain_depth` is 0 for self-contained rows, and `base.chain_depth + +-- 1` for a delta row — bounded (see the ingest-path cap) the same way +-- git caps its own delta chains at 50, so a read never has to walk an +-- unbounded number of hops to reconstruct one object. +ALTER TABLE objects ADD COLUMN base_oid BYTEA; +ALTER TABLE objects ADD COLUMN chain_depth SMALLINT NOT NULL DEFAULT 0; + +ALTER TABLE objects ADD CONSTRAINT objects_chain_depth_nonneg CHECK (chain_depth >= 0); +ALTER TABLE objects ADD CONSTRAINT objects_base_oid_iff_depth + CHECK ((base_oid IS NULL) = (chain_depth = 0)); + +-- Index for the chain-walk read path (recursive lookup by base_oid). +-- Partial: only delta rows ever populate base_oid, so indexing NULLs +-- (the common case — most objects will remain self-contained, e.g. +-- the ~35% observed as full objects in a real git pack) would waste +-- space for no benefit. +CREATE INDEX objects_base_oid_idx ON objects (base_oid) WHERE base_oid IS NOT NULL; diff --git a/mnemosyne_postgres/src/delta.rs b/mnemosyne_postgres/src/delta.rs new file mode 100644 index 0000000..b05a91c --- /dev/null +++ b/mnemosyne_postgres/src/delta.rs @@ -0,0 +1,156 @@ +//! Reference-delta object encoding: compress one object's bytes +//! against a chosen "base" object's bytes using zstd's prefix +//! (single-use dictionary) mode, instead of compressing each object +//! independently. +//! +//! This is the storage mechanism behind migration `0013`'s +//! `base_oid`/`chain_depth` columns. See +//! `docs/compression-and-throughput-strategy.md` for the measured +//! rationale — on a real repository corpus, compressing a modified +//! object against its immediate predecessor gets 33-156x smaller +//! output than compressing the same bytes alone, because the +//! predecessor relationship exposes the near-total byte overlap a +//! per-object compressor structurally cannot see. +//! +//! Deliberately built on zstd's *prefix* mode +//! (`ZSTD_CCtx_refPrefix`/`ZSTD_DCtx_refPrefix` — the mechanism behind +//! the `zstd --patch-from` CLI feature), not git's own +//! ofs-delta/ref-delta pack opcodes. A prefix is "use these literal +//! bytes as a one-shot dictionary for exactly this frame" — no +//! training pass, no persistent shared dictionary, and no +//! git-specific structure in the output. That keeps the object pool +//! populatable by anything that can supply "this object, and the +//! prior object it's most similar to" — a git pack's own delta +//! hints, an hg revlog importer translating its own delta chain, or a +//! future non-git backend — without any of them needing to +//! understand git's pack format. + +use bytes::Bytes; + +/// Maximum delta chain depth, mirroring git's own default pack delta +/// depth cap. Bounds worst-case read-path reconstruction cost to a +/// fixed number of prefix-decompressions per object, the same way +/// git bounds worst-case unpack cost. Not a tuned value for this +/// project specifically — chosen to match the reference point +/// (`git verify-pack -v` on a real corpus shows chain lengths +/// distributed 1..=50 with a roughly geometric decay) rather than +/// invented fresh. +pub const MAX_CHAIN_DEPTH: u16 = 50; + +#[derive(Debug, thiserror::Error)] +pub enum DeltaError { + #[error("zstd compress failed: {0}")] + Compress(std::io::Error), + #[error("zstd decompress failed: {0}")] + Decompress(std::io::Error), +} + +/// Compress `target` against `base` as a one-shot zstd prefix. +/// +/// The returned bytes are only meaningful together with the exact +/// `base` bytes used here — decoding requires supplying the same +/// base bytes to [`decode_against_base`]. No dictionary training, no +/// persisted state: this is "diff `target` against these literal +/// bytes," full stop. +/// +/// Compression level 9 matches what was used to produce the measured +/// numbers in `docs/compression-and-throughput-strategy.md` — high +/// enough to find the redundancy that matters here (near-total +/// overlap with `base`), well short of the diminishing-returns +/// territory above ~15-19 where wall-clock cost stops paying for +/// itself on already-small deltas. +pub fn encode_against_base(base: &[u8], target: &[u8]) -> Result, DeltaError> { + let mut cctx = zstd_safe::CCtx::create(); + cctx.set_parameter(zstd_safe::CParameter::CompressionLevel(9)) + .map_err(|code| DeltaError::Compress(zstd_error(code)))?; + cctx.ref_prefix(base) + .map_err(|code| DeltaError::Compress(zstd_error(code)))?; + let bound = zstd_safe::compress_bound(target.len()); + let mut out = Vec::with_capacity(bound); + cctx.compress2(&mut out, target) + .map_err(|code| DeltaError::Compress(zstd_error(code)))?; + Ok(out) +} + +/// Reverse of [`encode_against_base`]: reconstruct the original bytes +/// given the exact `base` bytes and the delta produced against them. +/// +/// `expected_len` is the target's known logical size (carried +/// separately in the `objects.size` column, since a compressed +/// delta's own length says nothing about its decoded size) — used to +/// size the output buffer exactly rather than guess-and-grow. +pub fn decode_against_base( + base: &[u8], + delta: &[u8], + expected_len: usize, +) -> Result { + let mut dctx = zstd_safe::DCtx::create(); + dctx.ref_prefix(base) + .map_err(|code| DeltaError::Decompress(zstd_error(code)))?; + let mut out = Vec::with_capacity(expected_len); + dctx.decompress(&mut out, delta) + .map_err(|code| DeltaError::Decompress(zstd_error(code)))?; + Ok(Bytes::from(out)) +} + +fn zstd_error(code: usize) -> std::io::Error { + std::io::Error::other(zstd_safe::get_error_name(code).to_string()) +} + +#[cfg(test)] +mod tests { + use super::*; + use proptest::prelude::*; + + // Byte-exact round trip: for any base/target pair, decoding what + // was encoded must reproduce `target` exactly. This is the + // property the entire storage scheme depends on — a lossy or + // off-by-one codec here would silently corrupt git object + // content. + proptest! { + #[test] + fn round_trips_byte_identical( + base in prop::collection::vec(any::(), 0..4096), + target in prop::collection::vec(any::(), 0..4096), + ) { + let encoded = encode_against_base(&base, &target).expect("encode"); + let decoded = decode_against_base(&base, &encoded, target.len()).expect("decode"); + prop_assert_eq!(decoded.as_ref(), target.as_slice()); + } + + // Similarity should make the encoded delta small relative to + // the target — not a byte-exact spec, but a regression guard + // against a change that accidentally starts ignoring the + // prefix (falling back to independent compression) for + // near-identical content, defeating the whole point. + #[test] + fn near_identical_content_compresses_much_smaller_than_target( + base in "[a-z]{200,2000}", + ) { + // target = base with one character changed in the middle, + // simulating the dominant "one-line edit" pattern this + // scheme targets. + let mut target = base.clone().into_bytes(); + let mid = target.len() / 2; + target[mid] = if target[mid] == b'a' { b'b' } else { b'a' }; + + let encoded = encode_against_base(base.as_bytes(), &target).expect("encode"); + // A near-identical delta should be a small fraction of the + // target's size — well under half is a conservative bound + // that still catches "prefix mode isn't engaging." + prop_assert!(encoded.len() < target.len() / 2); + + let decoded = decode_against_base(base.as_bytes(), &encoded, target.len()).expect("decode"); + prop_assert_eq!(decoded.as_ref(), target.as_slice()); + } + + /// Empty target/base are edge cases the ingest path can + /// legitimately hit (empty blobs are valid git objects). + #[test] + fn empty_inputs_round_trip(base in prop::collection::vec(any::(), 0..64)) { + let encoded = encode_against_base(&base, &[]).expect("encode empty target"); + let decoded = decode_against_base(&base, &encoded, 0).expect("decode empty target"); + prop_assert_eq!(decoded.as_ref(), &[] as &[u8]); + } + } +} diff --git a/mnemosyne_postgres/src/lib.rs b/mnemosyne_postgres/src/lib.rs index 95b8538..3018518 100644 --- a/mnemosyne_postgres/src/lib.rs +++ b/mnemosyne_postgres/src/lib.rs @@ -1,7 +1,9 @@ //! Postgres-backed implementation of [`mnemosyne_git`] storage traits. +pub mod delta; pub mod events; +pub use delta::{DeltaError, MAX_CHAIN_DEPTH, decode_against_base, encode_against_base}; pub use events::{ EVENTS_ADVISORY_LOCK_ID, EVENTS_NOTIFY_CHANNEL, EventRow, fetch_events_since, insert_event_in_tx, @@ -65,6 +67,19 @@ pub enum PgError { missing: ObjectId, referenced_from: ObjectId, }, + + #[error("delta decode failed for object {oid}: {source}")] + Delta { + oid: ObjectId, + #[source] + source: crate::delta::DeltaError, + }, + + #[error( + "broken delta chain for object {oid}: expected a self-contained row \ + (base_oid IS NULL) at the end of the chain, found another base_oid instead" + )] + BrokenDeltaChain { oid: ObjectId }, } /// Postgres-backed mnemosyne backend. @@ -555,23 +570,70 @@ async fn object_find<'e, E>( where E: sqlx::Executor<'e, Database = Postgres>, { - let row: Option<(i16, Vec)> = sqlx::query_as( - "SELECT o.kind, o.data \ - FROM objects o \ - INNER JOIN repo_objects r ON o.oid = r.oid \ - WHERE r.repo_did = $1 AND r.oid = $2", + // Walk the delta chain from the target down to its self-contained + // base in one round trip via a recursive CTE, rather than N + // sequential lookups. `repo_objects` is only checked at the + // anchor (the target itself) — ancestor rows in the chain live in + // the shared, cross-repo dedup pool and don't need their own + // inclusion row for THIS repo, same as any other dedup-pool read + // (see migration 0007's rationale: inclusion governs "can this + // repo be told this oid exists," not "can bytes ever be + // resolved" — a repo that legitimately has the tip object is + // entitled to its full reconstructed content). + // + // `hop` counts up from the target (0); ordering by `hop DESC` + // below returns the chain base-to-tip, the order reconstruction + // needs. The `hop < 64` guard is a defensive backstop against a + // corrupt/cyclic chain — write-side `MAX_CHAIN_DEPTH` (50) should + // make this unreachable in practice, but a broken invariant here + // must fail loudly (via `BrokenDeltaChain` below) rather than + // recurse without bound. + #[derive(sqlx::FromRow)] + struct ChainRow { + kind: i16, + size: i64, + data: Vec, + base_oid: Option>, + } + let rows: Vec = sqlx::query_as( + "WITH RECURSIVE chain(oid, kind, size, data, base_oid, hop) AS ( \ + SELECT o.oid, o.kind, o.size, o.data, o.base_oid, 0 \ + FROM objects o \ + INNER JOIN repo_objects r ON o.oid = r.oid \ + WHERE r.repo_did = $1 AND r.oid = $2 \ + UNION ALL \ + SELECT o.oid, o.kind, o.size, o.data, o.base_oid, c.hop + 1 \ + FROM objects o \ + INNER JOIN chain c ON o.oid = c.base_oid \ + WHERE c.hop < 64 \ + ) \ + SELECT kind, size, data, base_oid FROM chain ORDER BY hop DESC", ) .bind(repo_did) .bind(id.as_bytes()) - .fetch_optional(executor) + .fetch_all(executor) .await?; - match row { - None => Ok(None), - Some((kind_i16, data)) => Ok(Some(OwnedObject { - kind: i16_to_kind(kind_i16)?, - data: Bytes::from(data), - })), + + let Some(base_row) = rows.first() else { + return Ok(None); + }; + if base_row.base_oid.is_some() { + return Err(PgError::BrokenDeltaChain { oid: id.to_owned() }); + } + + let mut kind = base_row.kind; + let mut data = Bytes::from(base_row.data.clone()); + for row in &rows[1..] { + let expected_len = row.size as usize; + data = delta::decode_against_base(&data, &row.data, expected_len) + .map_err(|source| PgError::Delta { oid: id.to_owned(), source })?; + kind = row.kind; } + + Ok(Some(OwnedObject { + kind: i16_to_kind(kind)?, + data, + })) } async fn object_exists<'e, E>(executor: E, repo_did: &str, id: &oid) -> Result @@ -652,6 +714,34 @@ where Ok(id) } +/// An object queued for insertion into the shared object pool, in its +/// final on-disk storage representation. +/// +/// Constructed by the ingest path (`receive_pack::unpack_into_tx`) +/// after resolving each pack entry's delta relationship, or lack of +/// one — see `docs/compression-and-throughput-strategy.md`. A named +/// struct rather than a wider tuple because the fields don't have a +/// self-evident order the way `(ObjectId, Kind, Vec)` does. +#[derive(Debug, Clone)] +pub struct PendingObject { + pub id: ObjectId, + pub kind: Kind, + /// Logical (decoded) object size. Always the true size + /// regardless of `storage`'s representation — needed on the read + /// path to size the decode buffer when `storage` is a delta, + /// since a delta's own byte length says nothing about the size + /// of what it decodes to. + pub logical_size: u64, + /// On-disk bytes: plain decoded content when `base_oid` is + /// `None`, or a zstd-prefix delta against `base_oid`'s decoded + /// bytes otherwise (see [`crate::delta::encode_against_base`]). + pub storage: Vec, + pub base_oid: Option, + /// 0 for self-contained; `base`'s `chain_depth + 1` for a delta. + /// Bounded by [`crate::delta::MAX_CHAIN_DEPTH`]. + pub chain_depth: i16, +} + /// Insert into the shared, content-addressed object pool only — no /// per-repo inclusion row. Idempotent (`ON CONFLICT DO NOTHING`) and /// content-addressed, so insertion order and which connection @@ -677,9 +767,9 @@ where /// time. Batching cuts round trips roughly `batch_size`-for-one; see /// [`insert_objects_parallel`]'s doc comment for the measured sweet /// spot and why very large batches lose again. -async fn object_pool_insert<'e, 'd, E>( +async fn object_pool_insert<'e, E>( executor: E, - batch: &[(ObjectId, Kind, &'d [u8])], + batch: &[PendingObject], ) -> Result<(), PgError> where E: sqlx::Executor<'e, Database = Postgres>, @@ -688,21 +778,27 @@ where let mut kinds: Vec = Vec::with_capacity(batch.len()); let mut sizes: Vec = Vec::with_capacity(batch.len()); let mut datas: Vec<&[u8]> = Vec::with_capacity(batch.len()); - for (id, kind, data) in batch { - oids.push(id.as_bytes().to_vec()); - kinds.push(kind_to_i16(*kind)); - sizes.push(data.len() as i64); - datas.push(data); + let mut base_oids: Vec>> = Vec::with_capacity(batch.len()); + let mut chain_depths: Vec = Vec::with_capacity(batch.len()); + for obj in batch { + oids.push(obj.id.as_bytes().to_vec()); + kinds.push(kind_to_i16(obj.kind)); + sizes.push(obj.logical_size as i64); + datas.push(&obj.storage); + base_oids.push(obj.base_oid.map(|b| b.as_bytes().to_vec())); + chain_depths.push(obj.chain_depth); } sqlx::query( - "INSERT INTO objects (oid, kind, size, data) \ - SELECT * FROM unnest($1::bytea[], $2::smallint[], $3::bigint[], $4::bytea[]) \ + "INSERT INTO objects (oid, kind, size, data, base_oid, chain_depth) \ + SELECT * FROM unnest($1::bytea[], $2::smallint[], $3::bigint[], $4::bytea[], $5::bytea[], $6::smallint[]) \ ON CONFLICT (oid) DO NOTHING", ) .bind(&oids) .bind(&kinds) .bind(&sizes) .bind(&datas) + .bind(&base_oids) + .bind(&chain_depths) .execute(executor) .await?; Ok(()) @@ -796,7 +892,7 @@ where /// it passed in. pub async fn insert_objects_parallel( pool: &PgPool, - objects: Vec<(ObjectId, Kind, Vec)>, + objects: Vec, concurrency: usize, ) -> Result<(), PgError> { // Measured sweet spot on real ~24KB-avg object content (see this @@ -818,7 +914,7 @@ pub async fn insert_objects_parallel( // don't hand anything back except success/failure. let shard_size = n.div_ceil(concurrency).max(1); let mut handles = Vec::with_capacity(concurrency); - for chunk in objects.chunks(shard_size).map(<[(ObjectId, Kind, Vec)]>::to_vec) { + for chunk in objects.chunks(shard_size).map(<[PendingObject]>::to_vec) { let pool = pool.clone(); handles.push(tokio::spawn(async move { let mut conn = pool.acquire().await?; @@ -828,11 +924,7 @@ pub async fn insert_objects_parallel( let mut insert_err = None; 'outer: for sub_batch in chunk.chunks(INSERT_BATCH_SIZE) { - let objects: Vec<(ObjectId, Kind, &[u8])> = sub_batch - .iter() - .map(|(id, kind, data)| (*id, *kind, data.as_slice())) - .collect(); - if let Err(e) = object_pool_insert(&mut *conn, &objects).await { + if let Err(e) = object_pool_insert(&mut *conn, sub_batch).await { insert_err = Some(e); break 'outer; } @@ -1340,21 +1432,113 @@ where return Ok(HashMap::new()); } let id_bytes: Vec> = ids.iter().map(|id| id.as_bytes().to_vec()).collect(); - let rows: Vec<(Vec, i16, Vec)> = - sqlx::query_as("SELECT oid, kind, data FROM objects WHERE oid = ANY($1::bytea[])") - .bind(&id_bytes) - .fetch_all(executor) - .await?; - let mut out = HashMap::with_capacity(rows.len()); - for (oid_bytes, kind_i16, data) in rows { - let oid = ObjectId::from_bytes_or_panic(&oid_bytes); - out.insert( - oid, - OwnedObject { - kind: i16_to_kind(kind_i16)?, - data: Bytes::from(data), - }, - ); + // Multi-anchor recursive CTE: same chain-walk as `object_find` + // (see its doc comment for the "no repo scoping needed for + // ancestor rows" rationale — same argument applies transitively + // here, on top of THIS function's own already-documented "no + // repo scoping" stance for the connectivity check overall), but + // starting from every requested oid at once instead of one. Each + // requested oid keeps its own copy of the chain rows under + // `target_oid` — a small duplication cost (shared ancestors get + // one row per descendant that needs them) in exchange for a + // single round trip instead of N, which is the same trade this + // function already made for the pre-delta plain-data case. + struct ChainRow { + target_oid: Vec, + kind: i16, + size: i64, + data: Vec, + base_oid: Option>, + hop: i32, + } + let rows: Vec = { + #[derive(sqlx::FromRow)] + struct Row { + target_oid: Vec, + kind: i16, + size: i64, + data: Vec, + base_oid: Option>, + hop: i32, + } + let raw: Vec = sqlx::query_as( + "WITH RECURSIVE chain(target_oid, oid, kind, size, data, base_oid, hop) AS ( \ + SELECT o.oid, o.oid, o.kind, o.size, o.data, o.base_oid, 0 \ + FROM objects o \ + WHERE o.oid = ANY($1::bytea[]) \ + UNION ALL \ + SELECT c.target_oid, o.oid, o.kind, o.size, o.data, o.base_oid, c.hop + 1 \ + FROM objects o \ + INNER JOIN chain c ON o.oid = c.base_oid \ + WHERE c.hop < 64 \ + ) \ + SELECT target_oid, kind, size, data, base_oid, hop FROM chain ORDER BY target_oid, hop DESC", + ) + .bind(&id_bytes) + .fetch_all(executor) + .await?; + raw.into_iter() + .map(|r| ChainRow { + target_oid: r.target_oid, + kind: r.kind, + size: r.size, + data: r.data, + base_oid: r.base_oid, + hop: r.hop, + }) + .collect() + }; + + let mut out = HashMap::with_capacity(ids.len()); + let mut current_target: Option> = None; + let mut base_data: Option = None; + let mut kind_i16: i16; + for row in rows { + if current_target.as_deref() != Some(&row.target_oid) { + // Starting a new target's chain (rows are grouped by + // `target_oid`, base-to-tip within each group per the + // `ORDER BY ... hop DESC`). The first row in each group + // must be self-contained — same invariant `object_find` + // checks, just inline here rather than via a named error + // variant, since a broken chain at this layer means the + // insert-time invariant was violated, not a + // caller-recoverable condition. + if row.base_oid.is_some() { + return Err(PgError::BrokenDeltaChain { + oid: ObjectId::from_bytes_or_panic(&row.target_oid), + }); + } + current_target = Some(row.target_oid.clone()); + base_data = Some(Bytes::from(row.data)); + kind_i16 = row.kind; + } else { + let expected_len = row.size as usize; + let base = base_data.take().expect("base_data set for in-progress chain"); + let decoded = delta::decode_against_base(&base, &row.data, expected_len).map_err( + |source| PgError::Delta { + oid: ObjectId::from_bytes_or_panic(&row.target_oid), + source, + }, + )?; + base_data = Some(decoded); + kind_i16 = row.kind; + } + if row.hop == 0 { + // Terminal row for this target (the original requested + // oid, not an ancestor) — write the fully-reconstructed + // result. Ancestor rows (hop > 0 relative to some OTHER + // target, or just intermediate hops of this target's own + // chain) only exist to feed the fold above and are never + // themselves inserted into `out`. + let target_id = ObjectId::from_bytes_or_panic(&row.target_oid); + out.insert( + target_id, + OwnedObject { + kind: i16_to_kind(kind_i16)?, + data: base_data.clone().expect("base_data populated this iteration"), + }, + ); + } } Ok(out) } diff --git a/mnemosyne_postgres/tests/delta_round_trip.rs b/mnemosyne_postgres/tests/delta_round_trip.rs new file mode 100644 index 0000000..7a31208 --- /dev/null +++ b/mnemosyne_postgres/tests/delta_round_trip.rs @@ -0,0 +1,202 @@ +//! Exercises the reference-delta storage path end to end through the +//! REAL production write path (`insert_objects_parallel` + +//! `PushTx::insert_repo_inclusions`, the same calls +//! `receive_pack::unpack_into_tx` makes) and the real read path +//! (`PgBackend::try_find`, which now walks the delta chain via a +//! recursive CTE). +//! +//! This is the actual oracle for migration `0013` + `delta.rs` + +//! `unpack_into_tx`'s resolution pass: none of the other integration +//! tests exercise `base_oid IS NOT NULL` rows at all (`ObjectStore::write` +//! / `object_insert` only ever produces self-contained rows), so +//! without this test the delta path has zero execution-level coverage +//! beyond the codec's own isolated proptests. + +mod common; + +use common::unique_test_did as unique_repo_did; +use gix_hash::Kind as HashKind; +use gix_object::Kind; +use mnemosyne_git::ObjectStore; +use mnemosyne_postgres::{PendingObject, PgBackend, encode_against_base, insert_objects_parallel}; + +#[test] +fn delta_encoded_object_round_trips_and_is_stored_as_a_delta() { + common::runtime().block_on(async { + let pool = common::shared_pool().await; + let repo_did = unique_repo_did(); + let backend = PgBackend::with_pool(pool.clone(), repo_did.clone(), HashKind::Sha1); + + // Base: a blob with real content, not the read path's minimum + // case — big enough that a one-line change actually compresses + // well against it (matching the "near-identical successive + // version" pattern this scheme targets). + let base_content = "the quick brown fox jumps over the lazy dog\n".repeat(50); + let base_id = gix_object::compute_hash(HashKind::Sha1, Kind::Blob, base_content.as_bytes()) + .expect("hash base"); + + // Target: base with one line changed in the middle — same + // shape as the "one-line edit to a tracked file" case measured + // in docs/compression-and-throughput-strategy.md. + let target_content = base_content.replace("lazy dog", "sleepy cat"); + let target_id = gix_object::compute_hash(HashKind::Sha1, Kind::Blob, target_content.as_bytes()) + .expect("hash target"); + assert_ne!(base_id, target_id, "flipping a byte must change the oid"); + + let delta_bytes = encode_against_base(base_content.as_bytes(), target_content.as_bytes()) + .expect("encode delta"); + assert!( + delta_bytes.len() < target_content.len() / 2, + "delta ({} bytes) should be well under half of target ({} bytes) for near-identical content", + delta_bytes.len(), + target_content.len(), + ); + + let pending = vec![ + PendingObject { + id: base_id, + kind: Kind::Blob, + logical_size: base_content.len() as u64, + storage: base_content.as_bytes().to_vec(), + base_oid: None, + chain_depth: 0, + }, + PendingObject { + id: target_id, + kind: Kind::Blob, + logical_size: target_content.len() as u64, + storage: delta_bytes, + base_oid: Some(base_id), + chain_depth: 1, + }, + ]; + + // Real production write path: pool insert, then this repo's + // inclusion — exactly the two calls `unpack_into_tx` / + // `receive_pack::handler` make (see receive_pack.rs). + insert_objects_parallel(&pool, pending, 2) + .await + .expect("insert_objects_parallel"); + let mut push_tx = backend.begin_push_transaction().await.expect("begin tx"); + push_tx + .insert_repo_inclusions(&[base_id, target_id]) + .await + .expect("insert_repo_inclusions"); + push_tx.commit().await.expect("commit"); + + // Oracle check #1: the row actually IS a delta on disk. Without + // this assertion, a bug that silently makes every write fall + // back to self-contained would make the rest of this test pass + // for the wrong reason. + let (stored_base_oid, stored_chain_depth): (Option>, i16) = sqlx::query_as( + "SELECT base_oid, chain_depth FROM objects WHERE oid = $1", + ) + .bind(target_id.as_bytes()) + .fetch_one(&pool) + .await + .expect("query target row"); + assert_eq!( + stored_base_oid.as_deref(), + Some(base_id.as_bytes()), + "target row must be stored as a delta against base_id" + ); + assert_eq!(stored_chain_depth, 1); + + let (base_stored_base_oid,): (Option>,) = + sqlx::query_as("SELECT base_oid FROM objects WHERE oid = $1") + .bind(base_id.as_bytes()) + .fetch_one(&pool) + .await + .expect("query base row"); + assert_eq!(base_stored_base_oid, None, "base row must be self-contained"); + + // Oracle check #2: reading the delta-encoded object back + // through the real ObjectStore trait (the CTE chain-walk in + // object_find) reproduces the exact original bytes — this is + // the property the whole scheme depends on; a lossy or + // off-by-one reconstruction would silently corrupt git object + // content. + let found = backend + .try_find(&target_id) + .await + .expect("try_find target") + .expect("target object present"); + assert_eq!(found.kind, Kind::Blob); + assert_eq!( + &found.data[..], + target_content.as_bytes(), + "reconstructed delta object must be byte-identical to the original" + ); + + // The base object (self-contained, never touched by the delta + // path) must still read back correctly too. + let found_base = backend + .try_find(&base_id) + .await + .expect("try_find base") + .expect("base object present"); + assert_eq!(&found_base.data[..], base_content.as_bytes()); + }); +} + +/// A chain of three: base -> delta1 (against base) -> delta2 (against +/// delta1). Confirms the CTE walk correctly folds more than one hop, +/// not just the single-hop case above. +#[test] +fn multi_hop_delta_chain_round_trips() { + common::runtime().block_on(async { + let pool = common::shared_pool().await; + let repo_did = unique_repo_did(); + let backend = PgBackend::with_pool(pool.clone(), repo_did.clone(), HashKind::Sha1); + + let v1 = "version one of this file, with enough padding text to make the delta meaningful.\n".repeat(20); + let v2 = v1.replace("one", "two"); + let v3 = v2.replace("two", "three"); + assert_ne!(v1, v2); + assert_ne!(v2, v3); + + let id1 = gix_object::compute_hash(HashKind::Sha1, Kind::Blob, v1.as_bytes()).unwrap(); + let id2 = gix_object::compute_hash(HashKind::Sha1, Kind::Blob, v2.as_bytes()).unwrap(); + let id3 = gix_object::compute_hash(HashKind::Sha1, Kind::Blob, v3.as_bytes()).unwrap(); + + let delta2 = encode_against_base(v1.as_bytes(), v2.as_bytes()).unwrap(); + let delta3 = encode_against_base(v2.as_bytes(), v3.as_bytes()).unwrap(); + + let pending = vec![ + PendingObject { + id: id1, + kind: Kind::Blob, + logical_size: v1.len() as u64, + storage: v1.as_bytes().to_vec(), + base_oid: None, + chain_depth: 0, + }, + PendingObject { + id: id2, + kind: Kind::Blob, + logical_size: v2.len() as u64, + storage: delta2, + base_oid: Some(id1), + chain_depth: 1, + }, + PendingObject { + id: id3, + kind: Kind::Blob, + logical_size: v3.len() as u64, + storage: delta3, + base_oid: Some(id2), + chain_depth: 2, + }, + ]; + + insert_objects_parallel(&pool, pending, 2).await.unwrap(); + let mut push_tx = backend.begin_push_transaction().await.unwrap(); + push_tx.insert_repo_inclusions(&[id1, id2, id3]).await.unwrap(); + push_tx.commit().await.unwrap(); + + let found3 = backend.try_find(&id3).await.unwrap().unwrap(); + assert_eq!(&found3.data[..], v3.as_bytes(), "hop-2 reconstruction must be exact"); + let found2 = backend.try_find(&id2).await.unwrap().unwrap(); + assert_eq!(&found2.data[..], v2.as_bytes(), "hop-1 reconstruction must be exact"); + }); +} diff --git a/mnemosyne_protocol/src/receive_pack.rs b/mnemosyne_protocol/src/receive_pack.rs index 4cede5b..dd537e0 100644 --- a/mnemosyne_protocol/src/receive_pack.rs +++ b/mnemosyne_protocol/src/receive_pack.rs @@ -615,7 +615,24 @@ fn unpack_into_tx( // pure redundant SHA1/SHA256 work over the same bytes, on // hardware where that hash has a real cost — found via profiling // this exact loop's downstream insert cost. - let decode_result: Result)>, String> = std::thread::scope(|scope| { + /// One decoded pack entry plus the metadata needed to resolve its + /// delta relationship (if any) in the post-pass below. + struct DecodedEntry { + id: ObjectId, + kind: gix_object::Kind, + data: Vec, + /// The entry's own pack-format header (`Commit`/`Tree`/`Blob`/ + /// `Tag`/`RefDelta`/`OfsDelta`) — read via a second, header-only + /// `pack.entry()` call alongside the full `get_object_by_index` + /// decode. Cheap (a varint/hash read, no zlib inflate) next to + /// the decode this loop already does; reusing `get_object_by_index` + /// unmodified keeps the already-working, concurrency-sensitive + /// decode path untouched rather than hand-rolling its internals + /// to save one redundant header parse. + header: gix_pack::data::entry::Header, + pack_offset: gix_pack::data::Offset, + } + let decode_result: Result, String> = std::thread::scope(|scope| { let handles: Vec<_> = (0..num_objects) .step_by(chunk_size) .map(|start| { @@ -628,9 +645,20 @@ fn unpack_into_tx( for idx in start..end { buf.clear(); match bundle_ref.get_object_by_index(idx, &mut buf, &mut inflate, &mut cache) { - Ok((data, _location)) => { + Ok((data, location)) => { let id = bundle_ref.index.oid_at_index(idx).to_owned(); - local.push((id, data.kind, data.data.to_vec())); + let pack_offset = location.pack_offset; + let header = match bundle_ref.pack.entry(pack_offset) { + Ok(entry) => entry.header, + Err(e) => return Err(format!("re-read header for entry {idx}: {e}")), + }; + local.push(DecodedEntry { + id, + kind: data.kind, + data: data.data.to_vec(), + header, + pack_offset, + }); } Err(e) => return Err(format!("decode entry {idx}: {e}")), } @@ -654,12 +682,137 @@ fn unpack_into_tx( eprintln!("[timing] phase2_decode_entries: {:.3}s ({} objects, concurrency={decode_concurrency})", phase_start.elapsed().as_secs_f64(), decoded.len()); let phase_start = std::time::Instant::now(); - // Extract ids before `decoded` is consumed by + // Phase 2.5: resolve each entry's delta relationship (if any) and + // encode it against its base with a general reference-compression + // codec (zstd prefix mode — see mnemosyne_postgres::delta), instead + // of discarding the client's already-computed similarity + // relationship the way this loop used to. See + // docs/compression-and-throughput-strategy.md for the full + // rationale and measured numbers. + // + // Scope, deliberately narrow for this exploratory pass: only bases + // present IN THIS PACK are delta-encoded. A `RefDelta` whose base + // isn't found in this pack's index (a thin-pack base pulled from + // existing storage rather than sent as pack bytes) falls back to + // self-contained — resolving and re-encoding against an + // already-stored base would need an extra DB round trip per such + // entry from inside this sync, parallel-decode-adjacent pass, + // which is real complexity deferred here rather than solved. + // `OfsDelta` bases are always in this pack by construction (the + // format names them by backward byte offset), so they're always + // resolvable. + let n = decoded.len(); + let mut offset_to_pos: std::collections::HashMap = + std::collections::HashMap::with_capacity(n); + let mut id_to_pos: std::collections::HashMap = + std::collections::HashMap::with_capacity(n); + for (pos, entry) in decoded.iter().enumerate() { + offset_to_pos.insert(entry.pack_offset, pos); + id_to_pos.insert(entry.id, pos); + } + + // Direct base position for entries whose base is in this pack; + // `None` for self-contained entries and for delta entries whose + // base isn't resolvable here (falls back to self-contained below). + let direct_base: Vec> = decoded + .iter() + .map(|entry| match &entry.header { + gix_pack::data::entry::Header::RefDelta { base_id } => id_to_pos.get(base_id).copied(), + gix_pack::data::entry::Header::OfsDelta { base_distance } => { + gix_pack::data::entry::Header::verified_base_pack_offset( + entry.pack_offset, + *base_distance, + ) + .and_then(|ofs| offset_to_pos.get(&ofs).copied()) + } + _ => None, + }) + .collect(); + + // Memoized, cycle-guarded resolution of each position's EFFECTIVE + // base (after applying the depth cap) and chain depth. Recursive + // for clarity in this exploratory pass — bounded in practice by + // real delta chain lengths (git's own packs run 1..=50, see + // docs/compression-and-throughput-strategy.md), but a + // pathological adversarial pack with an enormous linear delta + // chain could overflow the call stack before hitting the depth + // cap; converting to an explicit-stack iterative walk is the + // production follow-up, not done here. + let mut resolved: Vec, u16)>> = vec![None; n]; // (effective_base_pos, chain_depth) + let mut visiting = vec![false; n]; + fn resolve( + pos: usize, + direct_base: &[Option], + resolved: &mut [Option<(Option, u16)>], + visiting: &mut [bool], + ) -> (Option, u16) { + if let Some(r) = resolved[pos] { + return r; + } + let result = match direct_base[pos] { + None => (None, 0u16), + Some(base_pos) if visiting[base_pos] => { + // Cycle (malformed/adversarial pack) — break it by + // treating this entry as self-contained rather than + // looping forever. + (None, 0u16) + } + Some(base_pos) => { + visiting[pos] = true; + let (_, base_depth) = resolve(base_pos, direct_base, resolved, visiting); + visiting[pos] = false; + if base_depth + 1 > mnemosyne_postgres::MAX_CHAIN_DEPTH { + // Cap hit — same posture as git's own delta depth + // limit: store this entry self-contained instead + // of extending the chain further. + (None, 0u16) + } else { + (Some(base_pos), base_depth + 1) + } + } + }; + resolved[pos] = Some(result); + result + } + for pos in 0..n { + resolve(pos, &direct_base, &mut resolved, &mut visiting); + } + + let mut pending: Vec = Vec::with_capacity(n); + for pos in 0..n { + let entry = &decoded[pos]; + let (effective_base, chain_depth) = resolved[pos].expect("resolved for every position"); + let logical_size = entry.data.len() as u64; + let encoded = effective_base.and_then(|base_pos| { + let base_data = &decoded[base_pos].data; + match mnemosyne_postgres::encode_against_base(base_data, &entry.data) { + // Only worth it if the delta actually beats storing + // the object self-contained — small/dissimilar + // objects can lose to zstd's per-frame overhead. + Ok(delta) if delta.len() < entry.data.len() => { + Some((decoded[base_pos].id, delta)) + } + _ => None, + } + }); + let (storage, base_oid, chain_depth) = match encoded { + Some((base_id, delta)) => (delta, Some(base_id), chain_depth), + None => (entry.data.clone(), None, 0), + }; + pending.push(mnemosyne_postgres::PendingObject { + id: entry.id, + kind: entry.kind, + logical_size, + storage, + base_oid, + chain_depth: chain_depth as i16, + }); + } + + // Extract ids before `pending` is consumed by // `insert_objects_parallel` below — needed downstream for - // `insert_repo_inclusions` and the connectivity check, and cheap - // to grab now (just copying already-computed small fixed-size - // hashes, not touching object content). - let oids: Vec = decoded.iter().map(|(id, _, _)| *id).collect(); + // `insert_repo_inclusions` and the connectivity check. + let oids: Vec = pending.iter().map(|p| p.id).collect(); // One connection per physical core is the shape our benchmarks // showed scaling well up to, with oversubscription past that point @@ -691,7 +844,7 @@ fn unpack_into_tx( let insert_result = handle.block_on(mnemosyne_postgres::insert_objects_parallel( push_tx.pool(), - decoded, + pending, concurrency, )); if let Err(e) = insert_result { diff --git a/mnemosyne_tests/tests/proptest-regressions/protocol_info_refs.txt b/mnemosyne_tests/tests/proptest-regressions/protocol_info_refs.txt new file mode 100644 index 0000000..e27c9a0 --- /dev/null +++ b/mnemosyne_tests/tests/proptest-regressions/protocol_info_refs.txt @@ -0,0 +1,7 @@ +# Seeds for failure cases proptest has generated in the past. It is +# automatically read and these particular cases re-run before any +# novel cases are generated. +# +# It is recommended to check this file in to source control so that +# everyone who runs the test benefits from these saved cases. +cc c466d487c1e6da8cf12606bfe1bff9dd7fedd8c793f902dc97bc17ce6539ed2e # shrinks to (graph, refs) = (ObjectGraph { objects: [(Sha1(ad79102c2b4df890853c69c4b53596242282772a), Blob, [232, 155, 242, 188, 126, 175, 93, 155, 212, 120, 106, 244, 190, 168, 243, 85, 227, 196, 209, 130, 73, 153, 105, 71, 62, 145, 166, 177, 40, 67, 186, 155, 75, 101, 151, 176, 27, 73, 238, 72, 157, 88, 183, 9, 114, 87, 211, 205, 201, 250, 232, 206]), (Sha1(98ed8a95dcdee10ebd2a3f532d8bdc8df426847b), Tree, [49, 48, 48, 54, 52, 52, 32, 102, 105, 108, 101, 48, 0, 173, 121, 16, 44, 43, 77, 248, 144, 133, 60, 105, 196, 181, 53, 150, 36, 34, 130, 119, 42]), (Sha1(bd326eea1ba42959a80c270388329fa803b683b1), Commit, [116, 114, 101, 101, 32, 57, 56, 101, 100, 56, 97, 57, 53, 100, 99, 100, 101, 101, 49, 48, 101, 98, 100, 50, 97, 51, 102, 53, 51, 50, 100, 56, 98, 100, 99, 56, 100, 102, 52, 50, 54, 56, 52, 55, 98, 10, 97, 117, 116, 104, 111, 114, 32, 77, 110, 101, 109, 111, 115, 121, 110, 101, 32, 84, 101, 115, 116, 101, 114, 32, 60, 116, 101, 115, 116, 64, 101, 120, 97, 109, 112, 108, 101, 46, 99, 111, 109, 62, 32, 48, 32, 43, 48, 48, 48, 48, 10, 99, 111, 109, 109, 105, 116, 116, 101, 114, 32, 77, 110, 101, 109, 111, 115, 121, 110, 101, 32, 84, 101, 115, 116, 101, 114, 32, 60, 116, 101, 115, 116, 64, 101, 120, 97, 109, 112, 108, 101, 46, 99, 111, 109, 62, 32, 48, 32, 43, 48, 48, 48, 48, 10, 10, 112, 114, 111, 112, 116, 101, 115, 116, 32, 99, 111, 109, 109, 105, 116])], head_commit: Sha1(bd326eea1ba42959a80c270388329fa803b683b1) }, ProtocolRefSet { edits: [RefEdit { change: Update { log: LogChange { mode: AndReference, force_create_reflog: false, message: "" }, expected: MustNotExist, new: Object(Sha1(bd326eea1ba42959a80c270388329fa803b683b1)) }, name: FullName("refs/heads/main"), deref: false }, RefEdit { change: Update { log: LogChange { mode: AndReference, force_create_reflog: false, message: "" }, expected: MustNotExist, new: Object(Sha1(bd326eea1ba42959a80c270388329fa803b683b1)) }, name: FullName("refs/heads/hehb"), deref: false }, RefEdit { change: Update { log: LogChange { mode: AndReference, force_create_reflog: false, message: "" }, expected: MustNotExist, new: Object(Sha1(bd326eea1ba42959a80c270388329fa803b683b1)) }, name: FullName("refs/tags/v5"), deref: false }, RefEdit { change: Update { log: LogChange { mode: AndReference, force_create_reflog: false, message: "" }, expected: MustNotExist, new: Symbolic(FullName("refs/heads/main")) }, name: FullName("HEAD"), deref: false }], expected_advertised: [(FullName("HEAD"), Sha1(bd326eea1ba42959a80c270388329fa803b683b1)), (FullName("refs/heads/hehb"), Sha1(bd326eea1ba42959a80c270388329fa803b683b1)), (FullName("refs/heads/main"), Sha1(bd326eea1ba42959a80c270388329fa803b683b1)), (FullName("refs/tags/v5"), Sha1(bd326eea1ba42959a80c270388329fa803b683b1))], head_symref_target: Some(FullName("refs/heads/main")) })