From 13a2f20558b96a254d7e4964adf3f3dbe175ac49 Mon Sep 17 00:00:00 2001 From: Edouard Vanbelle Date: Fri, 4 Sep 2026 21:43:49 +0200 Subject: [PATCH] fix(dedup): make the chunk reap guard registry-driven too MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit GC phase 2 already had the right shape — `ref_count <= 0 AND NOT EXISTS(manifest lists it) AND NOT EXISTS(file points at it)` — so unlike phase 1 before 6dc045ea, a stale counter could only delay collection there, never delete live bytes. What it did not have is any connection to `BlobReferenceRegistry`: the two cross-checks named `storage.chunk_manifests` and `storage.files` literally. That is correct today and one source away from not being. Both `content_derived_blobs` and `file_attached_blobs` return None at RefLevel::Chunk, so the registry's chunk union is exactly manifests + legacy files. The moment anything contributes at that level — a legacy whole-file derived blob, or file_versions when versioning lands — phase 2 misses it and reaps referenced bytes. That is precisely the failure the registry was built to prevent, and precisely what the phase 1 comment warns about while phase 2 sat unfixed. ## Why this is additive, not a swap `no_reference_predicate` is assembled from fragments designed for COUNTING, and FilesReferenceSource's chunk-level fragment deliberately excludes files whose blob_hash has a manifest — otherwise a single-chunk blob, whose file hash and lone chunk hash are the same BLAKE3, would be counted at both levels. Correct for a recompute; too narrow for a reap guard. Concretely: a `storage.blobs` row keyed by a MULTI-chunk file's hash is not a member of its own manifest's chunk_hashes, and such rows exist transiently while `rechunk` migrates a legacy blob. Replacing the hardcoded guards with the registry predicate would have satisfied "unreferenced" for that row while a live storage.files row still pointed at it — reaping it mid-migration. So the guards stay and the registry predicate is ANDed on top. Adding a conjunct can only spare more rows, never reap more, so this cannot regress; what it buys is that a future chunk-level source is honoured automatically. ## Also: EXISTS instead of COUNT in the hot path ChunksReferenceSource had no `ref_exists_sql` override, so the trait default wrapped its counting fragment as `(SELECT COUNT(*) …) > 0`. That now runs per candidate row inside the reap guard, and a heavily-deduplicated chunk is exactly where counting every referrer is most expensive and least necessary. FilesReferenceSource already carried this override for the same reason; ChunksReferenceSource now does too. Semantically identical, so no golden-test drift beyond the shape. ## Tests `blob_reap_statement_is_stable` pins the assembled statement, and `empty_registry_refuses_to_build_blob_reap_statement` mirrors the manifest builder's loud failure on a wiring bug. `a_new_chunk_level_source_reaches_the_blob_reap_statement` is the one that earns its keep: since no shipped source contributes at chunk level, a golden test alone would not notice the registry conjunct being dropped. It registers a synthetic source and asserts the fragment appears. Verified 921 passed / 0 failed on a clean database, and again on a second consecutive run against the same one. Co-Authored-By: Claude Opus 5 (1M context) --- .../repositories/pg/blob_reference_sources.rs | 25 +++ src/infrastructure/services/dedup_service.rs | 209 +++++++++++++++--- 2 files changed, 206 insertions(+), 28 deletions(-) diff --git a/src/infrastructure/repositories/pg/blob_reference_sources.rs b/src/infrastructure/repositories/pg/blob_reference_sources.rs index e9534c82..c2b2fa39 100644 --- a/src/infrastructure/repositories/pg/blob_reference_sources.rs +++ b/src/infrastructure/repositories/pg/blob_reference_sources.rs @@ -76,6 +76,27 @@ fn files_exists_sql(level: RefLevel, outer_hash_expr: &str) -> Option { } } +/// Short-circuiting existence form of [`chunks_ref_sql`]. +/// +/// Same motivation as [`files_exists_sql`], and it now matters more: this +/// fragment sits in `dedup_gc`'s **phase-2 reap guard**, evaluated per +/// candidate blob row. Without the override the trait default wraps the +/// counting form as `(SELECT COUNT(*) …) > 0`, which scans every manifest +/// listing the chunk before comparing — a heavily-deduplicated chunk is +/// exactly the case where that is most expensive and least necessary. +fn chunks_exists_sql(level: RefLevel, outer_hash_expr: &str) -> Option { + match level { + RefLevel::Chunk => { + let m = MANIFEST_ALIAS; + Some(format!( + "EXISTS (SELECT 1 FROM storage.chunk_manifests {m} \ + WHERE {outer_hash_expr} = ANY({m}.chunk_hashes))" + )) + } + RefLevel::Manifest => None, + } +} + /// Fragment for [`ChunksReferenceSource`]. See [`files_ref_sql`]. fn chunks_ref_sql(level: RefLevel, outer_hash_expr: &str) -> Option { match level { @@ -280,6 +301,10 @@ impl BlobReferenceSource for ChunksReferenceSource { chunks_ref_sql(level, outer_hash_expr) } + fn ref_exists_sql(&self, level: RefLevel, outer_hash_expr: &str) -> Option { + chunks_exists_sql(level, outer_hash_expr) + } + async fn count_references(&self, blob_hash: &str) -> Result { let n: i64 = sqlx::query_scalar( "SELECT COUNT(*) FROM storage.chunk_manifests WHERE $1 = ANY(chunk_hashes)", diff --git a/src/infrastructure/services/dedup_service.rs b/src/infrastructure/services/dedup_service.rs index 850b8e3a..f5faab14 100644 --- a/src/infrastructure/services/dedup_service.rs +++ b/src/infrastructure/services/dedup_service.rs @@ -478,6 +478,74 @@ async fn populate_integrity_blob_sizes<'a>( /// true for every row and this statement would delete every manifest in the /// database. `DedupService::new` always registers `FilesReferenceSource`, so /// the only way to reach this is to pass a deliberately empty registry. +/// Build the chunk/blob reap statement (GC phase 2) from the registered +/// reference sources. +/// +/// Unlike [`manifest_reap_sql`], the registry predicate here is **added to** +/// the hardcoded guards rather than replacing them. That asymmetry is +/// deliberate and the reason this was not a mechanical swap. +/// +/// `no_reference_predicate` is built from fragments designed for *counting*, +/// and `FilesReferenceSource`'s chunk-level fragment deliberately excludes +/// files whose `blob_hash` has a manifest — otherwise a single-chunk blob, +/// where the file hash and its lone chunk hash are the same BLAKE3, would be +/// counted at both levels. Correct for a recompute; too narrow for a reap +/// guard. A `storage.blobs` row keyed by a MULTI-chunk file's hash — which +/// exists transiently while `rechunk` migrates a legacy blob, and is not a +/// member of its own manifest's `chunk_hashes` — would satisfy the registry's +/// "unreferenced" test while a live `storage.files` row still points at it. +/// Swapping the guards out would have reaped it mid-migration. +/// +/// So the statement keeps `NOT EXISTS (manifest lists it as a chunk)` and +/// `NOT EXISTS (any file points at it)`, and ANDs the registry predicate on +/// top. Adding a conjunct can only ever spare more rows, never reap more, so +/// this cannot regress; what it buys is that a future source contributing at +/// [`RefLevel::Chunk`] is honoured automatically instead of being silently +/// missed — the same failure that made Phase 1's hardcoded cross-check +/// dangerous. +/// +/// Today the registry adds nothing operationally: +/// `content_derived_blobs` and `file_attached_blobs` both return `None` at +/// `RefLevel::Chunk`, so its union is exactly manifests + legacy files. The +/// point is what happens when that stops being true. +/// +/// `$1` is the batch limit, `$2` the grace window in seconds. +/// +/// # Panics +/// +/// If no source contributes at [`RefLevel::Chunk`]. Same reasoning as +/// [`manifest_reap_sql`]: a missing predicate must be loud rather than +/// silently degrading to "nothing references anything". +fn blob_reap_sql(registry: &BlobReferenceRegistry) -> String { + let unreferenced = registry + .no_reference_predicate(RefLevel::Chunk, "b.hash") + .expect( + "no chunk-level blob reference source registered: the reap \ + predicate would lose its registry cross-check", + ); + + format!( + "DELETE FROM storage.blobs + WHERE ctid = ANY( + SELECT b.ctid FROM storage.blobs b + WHERE b.ref_count <= 0 + AND (b.orphaned_at IS NULL + OR b.orphaned_at < now() - ($2::int * interval '1 second')) + AND NOT EXISTS ( + SELECT 1 FROM storage.chunk_manifests m + WHERE m.chunk_hashes @> ARRAY[b.hash::text] + ) + AND NOT EXISTS ( + SELECT 1 FROM storage.files f + WHERE f.blob_hash = b.hash + ) + AND {unreferenced} + LIMIT $1 + ) + RETURNING hash, size" + ) +} + fn manifest_reap_sql(registry: &BlobReferenceRegistry) -> String { let orphaned = registry .no_reference_predicate(RefLevel::Manifest, "m.file_hash") @@ -526,6 +594,10 @@ pub struct DedupService { /// Kept as a field so `garbage_collect` runs a fixed statement rather /// than assembling SQL inside a delete loop — see `manifest_reap_sql`. manifest_reap_sql: String, + /// The chunk/blob reap statement (GC phase 2), same treatment — see + /// [`blob_reap_sql`], including why its registry predicate is additive + /// rather than a replacement for the hardcoded guards. + blob_reap_sql: String, } impl DedupService { @@ -548,6 +620,7 @@ impl DedupService { manifest_cache: Self::build_manifest_cache(), reference_registry: registry.clone(), manifest_reap_sql: manifest_reap_sql(®istry), + blob_reap_sql: blob_reap_sql(®istry), } } @@ -581,6 +654,7 @@ impl DedupService { /// entirely — see `docs/plan/derived-blobs.md`. pub fn with_reference_registry(mut self, registry: Arc) -> Self { self.manifest_reap_sql = manifest_reap_sql(®istry); + self.blob_reap_sql = blob_reap_sql(®istry); self.reference_registry = registry; self } @@ -1033,6 +1107,7 @@ impl DedupService { manifest_cache: Self::build_manifest_cache(), reference_registry: stub_registry.clone(), manifest_reap_sql: manifest_reap_sql(&stub_registry), + blob_reap_sql: blob_reap_sql(&stub_registry), } } @@ -3230,41 +3305,29 @@ impl DedupService { // NULL orphaned_at — a pre-migration row or a path that never // stamped it; those are safe to take immediately), AND // • no manifest still lists it as a chunk, AND - // • no file still points at it directly (legacy whole-file blob). + // • no file still points at it directly (legacy whole-file blob), + // AND + // • no registered reference source claims it at the chunk level. // - // The two NOT EXISTS guards mirror Phase 1's file cross-check: a stale - // ref_count = 0 on still-referenced content can then only delay - // collection, never delete live bytes. The grace window keeps a + // The NOT EXISTS guards mean a stale ref_count = 0 on still-referenced + // content can only delay collection, never delete live bytes — unlike + // Phase 1 before `manifest_reap_sql` dropped its ref_count arm, this + // phase always had that property. The registry conjunct is additive + // (see `blob_reap_sql`): it cannot reap anything the hardcoded guards + // would have spared, it just stops a future chunk-level source from + // being missed. The grace window keeps a // concurrent uploader that is about to pin a just-orphaned chunk from // racing the row-delete → file-unlink gap (see GC_ORPHAN_GRACE_SECS). // The ctid snapshot already protects against a pin that commits DURING // the DELETE (the pin rewrites the row's ctid, so it drops out of the // set); grace covers the remaining post-commit unlink window. loop { - let batch: Vec<(String, i64)> = sqlx::query_as( - "DELETE FROM storage.blobs - WHERE ctid = ANY( - SELECT b.ctid FROM storage.blobs b - WHERE b.ref_count <= 0 - AND (b.orphaned_at IS NULL - OR b.orphaned_at < now() - ($2::int * interval '1 second')) - AND NOT EXISTS ( - SELECT 1 FROM storage.chunk_manifests m - WHERE m.chunk_hashes @> ARRAY[b.hash::text] - ) - AND NOT EXISTS ( - SELECT 1 FROM storage.files f - WHERE f.blob_hash = b.hash - ) - LIMIT $1 - ) - RETURNING hash, size", - ) - .bind(BATCH_SIZE) - .bind(grace_secs as i32) - .fetch_all(self.maintenance_pool.as_ref()) - .await - .map_err(|e| DomainError::internal_error("Dedup", format!("GC blobs: {e}")))?; + let batch: Vec<(String, i64)> = sqlx::query_as(&self.blob_reap_sql) + .bind(BATCH_SIZE) + .bind(grace_secs as i32) + .fetch_all(self.maintenance_pool.as_ref()) + .await + .map_err(|e| DomainError::internal_error("Dedup", format!("GC blobs: {e}")))?; if batch.is_empty() { break; @@ -3889,6 +3952,96 @@ mod tests { fn empty_registry_refuses_to_build_reap_statement() { let _ = manifest_reap_sql(&BlobReferenceRegistry::new()); } + + #[test] + #[should_panic(expected = "no chunk-level blob reference source")] + fn empty_registry_refuses_to_build_blob_reap_statement() { + let _ = blob_reap_sql(&BlobReferenceRegistry::new()); + } + + /// Golden test for GC phase 2, same purpose as the manifest one. + /// + /// Note what this pins that the manifest statement does not: the two + /// hardcoded `NOT EXISTS` guards **and** the registry predicate, ANDed. + /// The registry fragment is not a replacement here — see `blob_reap_sql` + /// for why substituting it would reap a legacy blob row mid-rechunk. + #[tokio::test] + async fn blob_reap_statement_is_stable() { + let sql = DedupService::new_stub().blob_reap_sql; + let expected = r#"DELETE FROM storage.blobs + WHERE ctid = ANY( + SELECT b.ctid FROM storage.blobs b + WHERE b.ref_count <= 0 + AND (b.orphaned_at IS NULL + OR b.orphaned_at < now() - ($2::int * interval '1 second')) + AND NOT EXISTS ( + SELECT 1 FROM storage.chunk_manifests m + WHERE m.chunk_hashes @> ARRAY[b.hash::text] + ) + AND NOT EXISTS ( + SELECT 1 FROM storage.files f + WHERE f.blob_hash = b.hash + ) + AND NOT (EXISTS (SELECT 1 FROM storage.files cnt_f WHERE cnt_f.blob_hash = b.hash AND NOT EXISTS (SELECT 1 FROM storage.chunk_manifests cnt_m WHERE cnt_m.file_hash = cnt_f.blob_hash)) + OR EXISTS (SELECT 1 FROM storage.chunk_manifests cnt_m WHERE b.hash = ANY(cnt_m.chunk_hashes))) + LIMIT $1 + ) + RETURNING hash, size"#; + assert_eq!(sql, expected, "blob reap statement changed:\n{sql}"); + } + + /// The reason phase 2 became registry-driven at all. + /// + /// Today no source contributes at [`RefLevel::Chunk`] beyond files and + /// manifests, so the registry conjunct is operationally redundant and a + /// golden test alone would not notice if it stopped being wired up. This + /// registers a synthetic chunk-level source and asserts its fragment + /// reaches the statement — which is what stops a future + /// `content_derived_blobs`-style table from being silently missed the way + /// Phase 1's hardcoded cross-check missed them. + #[tokio::test] + async fn a_new_chunk_level_source_reaches_the_blob_reap_statement() { + use crate::application::ports::blob_reference_ports::BlobReferenceSource; + + struct FakeChunkSource; + + #[async_trait::async_trait] + impl BlobReferenceSource for FakeChunkSource { + fn source_name(&self) -> &'static str { + "fake_chunk_source" + } + fn ref_count_sql(&self, level: RefLevel, outer: &str) -> Option { + self.ref_exists_sql(level, outer) + } + fn ref_exists_sql(&self, level: RefLevel, outer: &str) -> Option { + match level { + RefLevel::Chunk => Some(format!( + "EXISTS (SELECT 1 FROM storage.zzz_fake WHERE blob_hash = {outer})" + )), + RefLevel::Manifest => None, + } + } + async fn count_references(&self, _hash: &str) -> Result { + Ok(0) + } + async fn list_referenced_blobs( + &self, + _cursor: Option>, + _limit: usize, + ) -> Result<(Vec, Option>), DomainError> { + Ok((Vec::new(), None)) + } + } + + let mut registry = BlobReferenceRegistry::new(); + registry.register(Arc::new(FakeChunkSource)); + let sql = blob_reap_sql(®istry); + + assert!( + sql.contains("storage.zzz_fake"), + "a chunk-level source must reach the phase-2 reap guard:\n{sql}" + ); + } use std::collections::HashSet; use tempfile::NamedTempFile;