diff --git a/migrations/20260802000000_blob_gc_grace.sql b/migrations/20260802000000_blob_gc_grace.sql new file mode 100644 index 00000000..daa93dcf --- /dev/null +++ b/migrations/20260802000000_blob_gc_grace.sql @@ -0,0 +1,47 @@ +-- Garbage-collection safety for orphaned blobs. +-- +-- The dedup GC deletes a blob row (committed) and then unlinks the backing +-- file. A concurrent uploader of identical content can re-reference a chunk in +-- that window. Two mechanisms make the sweep safe: +-- (a) garbage_collect() never collects a blob still referenced by a manifest +-- (chunk) or a file (legacy whole-file blob) — cross-checks backed by +-- idx_chunk_manifests_chunk_hashes_gin and idx_files_blob_hash. A stale +-- ref_count = 0 on live content can then only delay collection, never +-- delete it. +-- (b) garbage_collect() never collects a blob that became unreferenced only +-- moments ago — the grace period below, mirroring git's gc.pruneExpire, +-- so a writer about to pin a just-orphaned chunk cannot race the sweep. +-- +-- `orphaned_at` records when ref_count last reached 0. NULL means the row is +-- referenced (ref_count > 0) or predates this column. + +ALTER TABLE storage.blobs ADD COLUMN IF NOT EXISTS orphaned_at TIMESTAMPTZ; + +-- Existing orphans start their grace window now, so applying this migration +-- never triggers an immediate sweep of content a writer might still be racing. +UPDATE storage.blobs + SET orphaned_at = now() + WHERE ref_count <= 0 AND orphaned_at IS NULL; + +-- GC scan index: orphan rows ordered by when they became collectible. Replaces +-- the old ref_count-only partial index (the GC now also filters on orphaned_at). +DROP INDEX IF EXISTS storage.idx_blobs_orphaned; +CREATE INDEX IF NOT EXISTS idx_blobs_gc_eligible + ON storage.blobs (orphaned_at) WHERE ref_count = 0; + +-- Stamp orphaned_at when a file delete drops a blob's ref_count to 0, so the +-- grace window starts at the moment of orphaning. No-op for multi-chunk files +-- whose file_hash is not itself a storage.blobs row. +CREATE OR REPLACE FUNCTION storage.decrement_blob_ref() +RETURNS trigger AS $$ +BEGIN + UPDATE storage.blobs + SET ref_count = GREATEST(ref_count - 1, 0), + orphaned_at = CASE WHEN GREATEST(ref_count - 1, 0) = 0 THEN now() ELSE orphaned_at END + WHERE hash = OLD.blob_hash; + RETURN OLD; +END; +$$ LANGUAGE plpgsql; + +COMMENT ON COLUMN storage.blobs.orphaned_at IS + 'When ref_count last reached 0; GC waits a grace period past this before deleting (NULL = referenced or pre-migration)'; diff --git a/src/application/services/trash_service.rs b/src/application/services/trash_service.rs index 935f294e..589f4ae9 100644 --- a/src/application/services/trash_service.rs +++ b/src/application/services/trash_service.rs @@ -568,9 +568,11 @@ impl TrashUseCase for TrashService { } } Err(e) => { - // Check if the file is not found - in that case, we can continue - // because we still want to remove the item from the trash index - if format!("{}", e).contains("not found") { + // File already gone — still remove the trash index + // entry. Match on the typed error kind, not the + // message text, so a reworded message can't + // silently turn this into a hard failure. + if e.kind == ErrorKind::NotFound { info!( "File not found, may already have been deleted: {}", file_id @@ -608,8 +610,9 @@ impl TrashUseCase for TrashService { info!("Successfully deleted folder permanently: {}", folder_id); } Err(e) => { - // Check if the folder is not found - in that case, we can continue - if format!("{}", e).contains("not found") { + // Folder already gone — still remove the trash + // index entry. Typed-kind match (see file branch). + if e.kind == ErrorKind::NotFound { info!( "Folder not found, may already have been deleted: {}", folder_id diff --git a/src/infrastructure/services/dedup_service.rs b/src/infrastructure/services/dedup_service.rs index 9115463f..1219d093 100644 --- a/src/infrastructure/services/dedup_service.rs +++ b/src/infrastructure/services/dedup_service.rs @@ -161,7 +161,9 @@ impl IngestGuard { ) { if !pinned.is_empty() && let Err(e) = sqlx::query( - "UPDATE storage.blobs SET ref_count = GREATEST(ref_count - 1, 0) + "UPDATE storage.blobs + SET ref_count = GREATEST(ref_count - 1, 0), + orphaned_at = CASE WHEN GREATEST(ref_count - 1, 0) = 0 THEN now() ELSE orphaned_at END WHERE hash = ANY($1)", ) .bind(&pinned) @@ -190,8 +192,8 @@ impl IngestGuard { ); } if let Err(e) = sqlx::query( - "INSERT INTO storage.blobs (hash, size, ref_count) - SELECT h, s, 0 FROM UNNEST($1::text[], $2::bigint[]) AS t(h, s) + "INSERT INTO storage.blobs (hash, size, ref_count, orphaned_at) + SELECT h, s, 0, now() FROM UNNEST($1::text[], $2::bigint[]) AS t(h, s) ON CONFLICT (hash) DO NOTHING", ) .bind(&hashes) @@ -380,6 +382,16 @@ impl DedupService { /// ~9 MiB regardless of file size. const FLUSH_MAX_BYTES: usize = 8 * 1024 * 1024; + /// Grace period (seconds) a blob must stay orphaned (`ref_count = 0`) + /// before [`garbage_collect`](Self::garbage_collect) may physically delete + /// it. Mirrors git's `gc.pruneExpire`: content that became unreferenced + /// only moments ago is never reaped, so a concurrent uploader about to pin + /// a just-orphaned chunk — or a delta-upload client that registered loose + /// chunks at `ref_count = 0` and is about to commit their manifest — cannot + /// race the sweep. Must comfortably exceed the longest plausible gap + /// between registering a chunk and referencing it (any in-flight upload). + const GC_ORPHAN_GRACE_SECS: i64 = 60 * 60; // 1 hour + /// Store content with CDC deduplication, straight from a byte stream — /// the single write path for every upload surface (REST multipart, /// WebDAV PUT, NextCloud PUT, chunked-upload assembly, WOPI PutFile). @@ -738,8 +750,8 @@ impl DedupService { let sizes: Vec = new_rows.iter().map(|(_, s)| *s).collect(); self.backend.sync_blobs(&hashes).await?; sqlx::query( - "INSERT INTO storage.blobs (hash, size, ref_count) - SELECT h, s, 0 FROM UNNEST($1::text[], $2::bigint[]) AS t(h, s) + "INSERT INTO storage.blobs (hash, size, ref_count, orphaned_at) + SELECT h, s, 0, now() FROM UNNEST($1::text[], $2::bigint[]) AS t(h, s) ON CONFLICT (hash) DO NOTHING", ) .bind(&hashes) @@ -923,7 +935,7 @@ impl DedupService { "INSERT INTO storage.blobs (hash, size, ref_count) SELECT h, s, 1 FROM UNNEST($1::text[], $2::bigint[]) AS t(h, s) ON CONFLICT (hash) DO UPDATE - SET ref_count = storage.blobs.ref_count + 1", + SET ref_count = storage.blobs.ref_count + 1, orphaned_at = NULL", ) .bind(&new_hashes) .bind(&new_sizes) @@ -971,7 +983,7 @@ impl DedupService { // session's reference NOW; hashes not returned don't exist and are // ours to write. let pinned: HashSet = sqlx::query_scalar::<_, String>( - "UPDATE storage.blobs SET ref_count = ref_count + 1 + "UPDATE storage.blobs SET ref_count = ref_count + 1, orphaned_at = NULL WHERE hash = ANY($1) RETURNING hash", ) @@ -1123,7 +1135,9 @@ impl DedupService { // Legacy blob let rows_affected = - sqlx::query("UPDATE storage.blobs SET ref_count = ref_count + 1 WHERE hash = $1") + sqlx::query( + "UPDATE storage.blobs SET ref_count = ref_count + 1, orphaned_at = NULL WHERE hash = $1", + ) .bind(hash) .execute(self.pool.as_ref()) .await @@ -1810,9 +1824,13 @@ impl DedupService { /// Garbage collect orphaned manifests and blobs. /// - /// Phase 1: Delete manifests with ref_count = 0, then decrement - /// chunk ref_counts for their chunks. - /// Phase 2: Delete blobs (chunks + legacy) with ref_count = 0. + /// Phase 1: Delete manifests with ref_count = 0 (or no referencing file), + /// then decrement chunk ref_counts for their chunks. + /// Phase 2: Delete blobs (chunks + legacy) that are unreferenced + /// (ref_count = 0), no longer listed by any manifest or file, and have + /// been orphaned for at least [`GC_ORPHAN_GRACE_SECS`](Self::GC_ORPHAN_GRACE_SECS). + /// The grace window and reference cross-checks together make the sweep safe + /// against a concurrent uploader re-referencing a just-orphaned chunk. pub async fn garbage_collect(&self) -> Result<(u64, u64), DomainError> { const BATCH_SIZE: i64 = 500; @@ -1855,9 +1873,12 @@ impl DedupService { // single-chunk file case where the PG file-delete trigger already // decremented blobs.ref_count (because file_hash == chunk_hash); // without the clamp this would underflow the CHECK constraint. + // Stamp orphaned_at so chunks freed here get the same GC grace + // window as any other newly-orphaned blob. sqlx::query( "UPDATE storage.blobs - SET ref_count = GREATEST(ref_count - 1, 0) + SET ref_count = GREATEST(ref_count - 1, 0), + orphaned_at = CASE WHEN GREATEST(ref_count - 1, 0) = 0 THEN now() ELSE orphaned_at END WHERE hash = ANY($1)", ) .bind(chunk_hashes) @@ -1880,17 +1901,44 @@ impl DedupService { } // ── Phase 2: GC orphaned blobs/chunks ──────────────────── + // A blob row is collectible only when ALL of these hold: + // • ref_count <= 0, AND + // • it has been orphaned for at least GC_ORPHAN_GRACE_SECS (or has a + // 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). + // + // 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 + // 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 ctid FROM storage.blobs - WHERE ref_count <= 0 + 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] + ) + AND NOT EXISTS ( + SELECT 1 FROM storage.files f + WHERE f.blob_hash = b.hash + ) LIMIT $1 ) RETURNING hash, size", ) .bind(BATCH_SIZE) + .bind(Self::GC_ORPHAN_GRACE_SECS as i32) .fetch_all(self.maintenance_pool.as_ref()) .await .map_err(|e| DomainError::internal_error("Dedup", format!("GC blobs: {e}")))?; @@ -1898,15 +1946,33 @@ impl DedupService { if batch.is_empty() { break; } + let n = batch.len(); - for (hash, size) in &batch { - if let Err(e) = self.backend.delete_blob(hash).await { - tracing::warn!("Failed to delete orphan blob {hash}: {e}"); - } + // The rows are already gone, so a concurrent re-upload of identical + // content recreates both row and file (durability before + // visibility); the grace window above keeps that race vanishingly + // narrow. Unlink the backing files with bounded fan-out so a large + // sweep doesn't serialise on a slow (e.g. S3) backend. + let backend = self.backend.clone(); + let deleted: Vec<(String, i64)> = stream::iter(batch) + .map(|(hash, size)| { + let backend = backend.clone(); + async move { + if let Err(e) = backend.delete_blob(&hash).await { + tracing::warn!("Failed to delete orphan blob {hash}: {e}"); + } + (hash, size) + } + }) + .buffer_unordered(Self::CHUNK_UPLOAD_CONCURRENCY) + .collect() + .await; + + for (hash, size) in &deleted { self.fire_blob_hooks(hash); total_bytes += *size as u64; } - total_deleted += batch.len() as u64; + total_deleted += n as u64; tokio::task::yield_now().await; } @@ -2256,7 +2322,9 @@ impl DedupService { return; } if let Err(e) = sqlx::query( - "UPDATE storage.blobs SET ref_count = GREATEST(ref_count - 1, 0) + "UPDATE storage.blobs + SET ref_count = GREATEST(ref_count - 1, 0), + orphaned_at = CASE WHEN GREATEST(ref_count - 1, 0) = 0 THEN now() ELSE orphaned_at END WHERE hash = ANY($1)", ) .bind(chunk_hashes) @@ -3299,6 +3367,103 @@ mod delta_upload_integration_tests { cleanup(&pool, &file_hash, file_id, &[fresh_hash]).await; } + // ── Garbage collection: grace window + reference cross-checks ─ + #[tokio::test] + async fn garbage_collect_honours_grace_window_and_references() { + let pool = test_pool().await; + let dir = TempDir::new().unwrap(); + let svc = local_svc(&pool, &dir).await; + let user = seed_user(&pool).await; + + // (A) An aged orphan (orphaned well past the grace window) with no + // references → must be collected (row + backing file). + // (B) A freshly orphaned blob (orphaned_at = now()) → must survive: a + // concurrent uploader could still be about to pin it. + let aged = blake3::hash(format!("aged-{}", Uuid::new_v4()).as_bytes()) + .to_hex() + .to_string(); + let fresh = blake3::hash(format!("fresh-{}", Uuid::new_v4()).as_bytes()) + .to_hex() + .to_string(); + for h in [&aged, &fresh] { + svc.backend() + .put_blob_from_bytes_unsynced(h, Bytes::from_static(b"xyz")) + .await + .expect("write blob"); + } + svc.backend() + .sync_blobs(&[aged.clone(), fresh.clone()]) + .await + .expect("sync"); + sqlx::query( + "INSERT INTO storage.blobs (hash, size, ref_count, orphaned_at) VALUES + ($1, 3, 0, now() - interval '2 hours'), + ($2, 3, 0, now())", + ) + .bind(&aged) + .bind(&fresh) + .execute(pool.as_ref()) + .await + .expect("seed orphans"); + + // (C) A chunk still listed by a live file's manifest, but whose + // blobs.ref_count has drifted to 0 and aged past the grace window. + // The manifest cross-check must keep it (and its bytes) alive — a + // stale ref_count must never delete referenced content. + let data = content(3 * 1024 * 1024, 71); + let (file_hash, owned_chunks, file_id) = + seed_owned_content(&svc, &pool, user, &data, "gc").await; + let referenced = owned_chunks[0].clone(); + sqlx::query( + "UPDATE storage.blobs + SET ref_count = 0, orphaned_at = now() - interval '2 hours' + WHERE hash = $1", + ) + .bind(&referenced) + .execute(pool.as_ref()) + .await + .expect("drift referenced chunk"); + + let (deleted, _bytes) = svc.garbage_collect().await.expect("gc"); + assert!(deleted >= 1, "the aged orphan must be collected"); + + // Aged orphan fully gone. + assert!( + blob_ref(&pool, &aged).await.is_none(), + "aged orphan row removed" + ); + assert!( + !svc.backend().blob_exists(&aged).await.unwrap(), + "aged orphan file unlinked" + ); + // Fresh orphan preserved by the grace window. + assert_eq!( + blob_ref(&pool, &fresh).await, + Some(0), + "fresh orphan survives the grace window" + ); + assert!( + svc.backend().blob_exists(&fresh).await.unwrap(), + "fresh orphan bytes kept" + ); + // Referenced chunk preserved by the manifest cross-check despite ref 0. + assert_eq!( + blob_ref(&pool, &referenced).await, + Some(0), + "referenced chunk row kept" + ); + assert!( + svc.backend().blob_exists(&referenced).await.unwrap(), + "referenced chunk bytes kept" + ); + + let _ = sqlx::query("DELETE FROM storage.blobs WHERE hash = ANY($1)") + .bind(vec![aged, fresh]) + .execute(pool.as_ref()) + .await; + cleanup(&pool, &file_hash, file_id, &[]).await; + } + // ── Verification read ──────────────────────────────────────── #[tokio::test] async fn hash_chunk_sequence_recomputes_and_validates_sizes() { diff --git a/src/infrastructure/services/search_index/content_index_worker.rs b/src/infrastructure/services/search_index/content_index_worker.rs index bcc74781..74d3649f 100644 --- a/src/infrastructure/services/search_index/content_index_worker.rs +++ b/src/infrastructure/services/search_index/content_index_worker.rs @@ -56,6 +56,11 @@ const PREVIEW_BYTES: usize = 16 * 1024; /// 1.5 s interval). const ORPHAN_SWEEP_TICKS: u64 = 2400; +/// Backoff before the supervisor restarts the drain loop after an abnormal +/// exit (a panic). Long enough that a tight crash-loop can't busy-spin, short +/// enough that indexing resumes promptly. +const WORKER_RESTART_BACKOFF_SECS: u64 = 5; + pub struct ContentIndexWorker { maintenance_pool: Arc, dedup: Arc, @@ -86,49 +91,77 @@ impl ContentIndexWorker { } } - /// Spawn the indexing loop. Fire-and-forget: the loop logs and survives - /// every error (an exited loop would silently freeze the index while the - /// queue grows), and the first drain runs immediately to absorb rows left - /// over from a previous run or the migration backfill. + /// Spawn the indexing loop, supervised. The drain loop logs and survives + /// every *operational* error (a failed drain just retries next tick), but a + /// panic in the loop body would otherwise kill the task and silently freeze + /// the index while the dirty queue grows unbounded. The supervisor restarts + /// the loop after a panic (with backoff) so indexing self-heals. The first + /// drain runs immediately to absorb rows left over from a previous run or + /// the migration backfill. #[instrument(skip(self))] pub fn start(self, needs_reseed: bool) { info!( "Starting content-index worker (every {}ms, batch {}, reseed: {})", self.interval_ms, DRAIN_BATCH, needs_reseed ); + let worker = Arc::new(self); tokio::spawn(async move { - if let Err(e) = self.prepare(needs_reseed).await { + // Reseed/version cleanup runs once, not on every restart. + if let Err(e) = worker.prepare(needs_reseed).await { error!("Content-index prepare failed (continuing with queue as-is): {e}"); } - let mut ticker = - tokio::time::interval(std::time::Duration::from_millis(self.interval_ms)); - ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); - let mut ticks: u64 = 0; + // run_loop() never returns under normal operation, so any exit is + // abnormal: a panic surfaces as a JoinError; a plain return would + // be a logic bug. Either way, log loudly and restart. loop { - ticker.tick().await; - for _ in 0..MAX_BATCHES_PER_TICK { - match self.drain_once().await { - Ok(0) => break, - Ok(drained) => { - debug!("Content-index drain: processed {drained} queue row(s)"); - if drained < DRAIN_BATCH as usize { - break; - } - } - Err(e) => { - error!("Content-index drain failed (queue preserved, will retry): {e}"); + let w = worker.clone(); + match tokio::spawn(async move { w.run_loop().await }).await { + Ok(()) => error!( + "Content-index drain loop returned unexpectedly; \ + restarting in {WORKER_RESTART_BACKOFF_SECS}s" + ), + Err(e) if e.is_panic() => error!( + "Content-index drain loop panicked ({e}); \ + restarting in {WORKER_RESTART_BACKOFF_SECS}s" + ), + Err(_) => return, // task cancelled — runtime shutting down + } + tokio::time::sleep(std::time::Duration::from_secs(WORKER_RESTART_BACKOFF_SECS)) + .await; + } + }); + } + + /// The perpetual drain loop. Extracted from [`start`](Self::start) so the + /// supervisor can run it in a child task and restart it after a panic. + async fn run_loop(&self) { + let mut ticker = tokio::time::interval(std::time::Duration::from_millis(self.interval_ms)); + ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + let mut ticks: u64 = 0; + loop { + ticker.tick().await; + for _ in 0..MAX_BATCHES_PER_TICK { + match self.drain_once().await { + Ok(0) => break, + Ok(drained) => { + debug!("Content-index drain: processed {drained} queue row(s)"); + if drained < DRAIN_BATCH as usize { break; } } - } - - ticks += 1; - if ticks.is_multiple_of(ORPHAN_SWEEP_TICKS) { - self.sweep_orphaned_text().await; + Err(e) => { + error!("Content-index drain failed (queue preserved, will retry): {e}"); + break; + } } } - }); + + ticks += 1; + if ticks.is_multiple_of(ORPHAN_SWEEP_TICKS) { + self.sweep_orphaned_text().await; + } + } } /// Spawn the discard-only janitor used when content search is DISABLED: