fix(dedup): count chunk ref_count per distinct chunk, not per occurrence
store_chunks bumped storage.blobs.ref_count once per chunk *occurrence* (it looped over the full chunk list, duplicates included), but remove_manifest_reference decrements once per *distinct* chunk (WHERE hash = ANY(chunk_hashes) matches each row a single time). For any file that repeats a chunk -- zero-filled regions in disk/VM images, repeated document structures, concatenated archives -- storing added +N while deleting removed -1, so the blob's ref_count never returned to 0 and the chunk was never garbage-collected: a permanent storage leak. Count per distinct chunk on the store side too, matching deletion. This also makes it faster: - existing chunks: one batched `UPDATE ... WHERE hash = ANY($1)` instead of one UPDATE per occurrence; - a brand-new chunk repeated within a file is read, uploaded and INSERTed once instead of once per occurrence. The manifest still stores the full per-occurrence chunk sequence (needed to reassemble the file). Forward fix: blobs already over-counted by the old path stay over-counted (a reconcile/verify pass could recompute them), but the bias is upward (leak), so no data is ever deleted early. CDC tests pass (12); fmt + clippy clean. https://claude.ai/code/session_01UtfkS3nZF1vrF5jNAps6wV
This commit is contained in:
@@ -459,8 +459,12 @@ impl DedupService {
|
|||||||
/// exist in `storage.blobs`.
|
/// exist in `storage.blobs`.
|
||||||
/// Phase 1: Reads only *new* chunks from the source file (the biggest
|
/// Phase 1: Reads only *new* chunks from the source file (the biggest
|
||||||
/// I/O saving for versioned files where most chunks are unchanged).
|
/// I/O saving for versioned files where most chunks are unchanged).
|
||||||
/// Phase 2: Parallel operations — uploads new chunks, bumps ref_count
|
/// Uploads each new chunk and bumps `ref_count` for chunks that already
|
||||||
/// for existing ones — with up to [`CHUNK_UPLOAD_CONCURRENCY`] in flight.
|
/// exist, with up to [`CHUNK_UPLOAD_CONCURRENCY`] uploads in flight.
|
||||||
|
///
|
||||||
|
/// `ref_count` is incremented once per *distinct* chunk (one reference per
|
||||||
|
/// manifest), staying symmetric with `remove_manifest_reference` so a file
|
||||||
|
/// that repeats a chunk cannot over-count and leak the blob forever.
|
||||||
async fn store_chunks(
|
async fn store_chunks(
|
||||||
&self,
|
&self,
|
||||||
source_path: &Path,
|
source_path: &Path,
|
||||||
@@ -469,20 +473,22 @@ impl DedupService {
|
|||||||
let pool = &self.pool;
|
let pool = &self.pool;
|
||||||
let backend = &self.backend;
|
let backend = &self.backend;
|
||||||
|
|
||||||
// ── Phase 0: Batch-check which chunks already exist ──────
|
// ── Phase 0: de-duplicate chunk hashes, then batch-check existence ──
|
||||||
let unique_hashes: Vec<String> = {
|
// A single file can legitimately repeat the same chunk many times
|
||||||
let mut seen = std::collections::HashSet::new();
|
// (zero-filled regions in disk/VM images, repeated document structures,
|
||||||
chunks
|
// concatenated archives). `ref_count` is tracked per *distinct* chunk —
|
||||||
.iter()
|
// one reference per manifest — to stay symmetric with
|
||||||
.filter_map(|c| {
|
// `remove_manifest_reference`, which decrements via
|
||||||
if seen.insert(c.hash.as_str()) {
|
// `WHERE hash = ANY(chunk_hashes)` (matching each row once). Counting
|
||||||
Some(c.hash.clone())
|
// per-occurrence here would over-increment and leak the blob forever.
|
||||||
} else {
|
// Keep the first occurrence of each hash so new chunks know where to
|
||||||
None
|
// read their bytes.
|
||||||
}
|
let mut seen = std::collections::HashSet::new();
|
||||||
})
|
let unique_chunks: Vec<&ChunkMeta> = chunks
|
||||||
.collect()
|
.iter()
|
||||||
};
|
.filter(|c| seen.insert(c.hash.as_str()))
|
||||||
|
.collect();
|
||||||
|
let unique_hashes: Vec<String> = unique_chunks.iter().map(|c| c.hash.clone()).collect();
|
||||||
|
|
||||||
let existing_hashes: std::collections::HashSet<String> =
|
let existing_hashes: std::collections::HashSet<String> =
|
||||||
sqlx::query_scalar::<_, String>("SELECT hash FROM storage.blobs WHERE hash = ANY($1)")
|
sqlx::query_scalar::<_, String>("SELECT hash FROM storage.blobs WHERE hash = ANY($1)")
|
||||||
@@ -498,98 +504,86 @@ impl DedupService {
|
|||||||
.into_iter()
|
.into_iter()
|
||||||
.collect();
|
.collect();
|
||||||
|
|
||||||
// ── Phase 1+2 (fused): upload NEW chunks just-in-time ────
|
// ── Phase 1: bump ref_count for every existing chunk in ONE query ──
|
||||||
// Read each new chunk by positioned I/O immediately before its
|
// (was one UPDATE per occurrence — now a single batched round-trip).
|
||||||
// upload, instead of first materializing every new chunk's *data* in
|
let existing: Vec<String> = unique_hashes
|
||||||
// a Vec. Peak heap for file content is bounded to
|
|
||||||
// ~CHUNK_UPLOAD_CONCURRENCY × CDC_MAX_CHUNK (≈ 8 MiB) — proportional
|
|
||||||
// to the chunk size, never the file size, so storing a large
|
|
||||||
// brand-new file no longer spikes RAM. Existing chunks skip all disk
|
|
||||||
// I/O and just bump ref_count.
|
|
||||||
//
|
|
||||||
// We first collect *owned* per-chunk metadata (hash + offset + length
|
|
||||||
// + existence flag — no file data) so the stream below does not borrow
|
|
||||||
// the `chunks` parameter across an `.await` (which would make this
|
|
||||||
// future non-`Send` and break the upload handlers).
|
|
||||||
let chunk_ops: Vec<(String, u64, usize, bool)> = chunks
|
|
||||||
.iter()
|
.iter()
|
||||||
.map(|chunk| {
|
.filter(|h| existing_hashes.contains(*h))
|
||||||
let exists = existing_hashes.contains(&chunk.hash);
|
.cloned()
|
||||||
(
|
.collect();
|
||||||
chunk.hash.clone(),
|
if !existing.is_empty() {
|
||||||
chunk.offset as u64,
|
sqlx::query("UPDATE storage.blobs SET ref_count = ref_count + 1 WHERE hash = ANY($1)")
|
||||||
chunk.length,
|
.bind(&existing)
|
||||||
exists,
|
.execute(pool.as_ref())
|
||||||
)
|
.await
|
||||||
})
|
.map_err(|e| {
|
||||||
|
DomainError::internal_error("Dedup", format!("Failed to bump ref_count: {}", e))
|
||||||
|
})?;
|
||||||
|
}
|
||||||
|
|
||||||
|
// ── Phase 2: upload each NEW chunk once, concurrently ──────────────
|
||||||
|
// Read each new chunk by positioned I/O immediately before its upload
|
||||||
|
// instead of materializing every chunk's *data* up front. Peak heap for
|
||||||
|
// file content stays bounded to ~CHUNK_UPLOAD_CONCURRENCY × CDC_MAX_CHUNK
|
||||||
|
// (≈ 8 MiB) — proportional to the chunk size, never the file size.
|
||||||
|
//
|
||||||
|
// Owned metadata (hash + offset + length, no file data) so the stream
|
||||||
|
// below does not borrow `chunks` across an `.await` (which would make
|
||||||
|
// this future non-`Send` and break the upload handlers).
|
||||||
|
let new_ops: Vec<(String, u64, usize)> = unique_chunks
|
||||||
|
.iter()
|
||||||
|
.filter(|c| !existing_hashes.contains(&c.hash))
|
||||||
|
.map(|c| (c.hash.clone(), c.offset as u64, c.length))
|
||||||
.collect();
|
.collect();
|
||||||
|
|
||||||
let source = Arc::new(std::fs::File::open(source_path).map_err(|e| {
|
let source = Arc::new(std::fs::File::open(source_path).map_err(|e| {
|
||||||
DomainError::internal_error("Dedup", format!("Failed to open source file: {}", e))
|
DomainError::internal_error("Dedup", format!("Failed to open source file: {}", e))
|
||||||
})?);
|
})?);
|
||||||
|
|
||||||
let results: Vec<Result<(), DomainError>> = stream::iter(chunk_ops)
|
let results: Vec<Result<(), DomainError>> = stream::iter(new_ops)
|
||||||
.map(|(hash, offset, length, exists)| {
|
.map(|(hash, offset, length)| {
|
||||||
let source = source.clone();
|
let source = source.clone();
|
||||||
let pool = pool.clone();
|
let pool = pool.clone();
|
||||||
let backend = backend.clone();
|
let backend = backend.clone();
|
||||||
async move {
|
async move {
|
||||||
if exists {
|
// Positioned read of just this chunk (≤ CDC_MAX_CHUNK) off
|
||||||
// Existing chunk: bump ref_count, no disk I/O.
|
// the async runtime, then upload.
|
||||||
sqlx::query(
|
let bytes = tokio::task::spawn_blocking(move || {
|
||||||
"UPDATE storage.blobs
|
use std::os::unix::fs::FileExt;
|
||||||
SET ref_count = ref_count + 1
|
let mut buf = vec![0u8; length];
|
||||||
WHERE hash = $1",
|
source.read_exact_at(&mut buf, offset)?;
|
||||||
)
|
Ok::<Vec<u8>, std::io::Error>(buf)
|
||||||
.bind(&hash)
|
})
|
||||||
.execute(pool.as_ref())
|
.await
|
||||||
.await
|
.map_err(|e| {
|
||||||
.map_err(|e| {
|
DomainError::internal_error("Dedup", format!("Read task failed: {}", e))
|
||||||
DomainError::internal_error(
|
})?
|
||||||
"Dedup",
|
.map_err(|e| {
|
||||||
format!("Failed to bump ref_count: {}", e),
|
DomainError::internal_error("Dedup", format!("Failed to read chunk: {}", e))
|
||||||
)
|
})?;
|
||||||
})?;
|
|
||||||
} else {
|
|
||||||
// New chunk: positioned read of just this chunk
|
|
||||||
// (≤ CDC_MAX_CHUNK) off the async runtime, then upload.
|
|
||||||
let bytes = tokio::task::spawn_blocking(move || {
|
|
||||||
use std::os::unix::fs::FileExt;
|
|
||||||
let mut buf = vec![0u8; length];
|
|
||||||
source.read_exact_at(&mut buf, offset)?;
|
|
||||||
Ok::<Vec<u8>, std::io::Error>(buf)
|
|
||||||
})
|
|
||||||
.await
|
|
||||||
.map_err(|e| {
|
|
||||||
DomainError::internal_error("Dedup", format!("Read task failed: {}", e))
|
|
||||||
})?
|
|
||||||
.map_err(|e| {
|
|
||||||
DomainError::internal_error(
|
|
||||||
"Dedup",
|
|
||||||
format!("Failed to read chunk: {}", e),
|
|
||||||
)
|
|
||||||
})?;
|
|
||||||
|
|
||||||
backend
|
backend
|
||||||
.put_blob_from_bytes(&hash, Bytes::from(bytes))
|
.put_blob_from_bytes(&hash, Bytes::from(bytes))
|
||||||
.await?;
|
.await?;
|
||||||
sqlx::query(
|
// ON CONFLICT covers a concurrent uploader inserting the
|
||||||
"INSERT INTO storage.blobs (hash, size, ref_count)
|
// same brand-new chunk between the existence check above and
|
||||||
VALUES ($1, $2, 1)
|
// this INSERT.
|
||||||
ON CONFLICT (hash) DO UPDATE
|
sqlx::query(
|
||||||
SET ref_count = storage.blobs.ref_count + 1",
|
"INSERT INTO storage.blobs (hash, size, ref_count)
|
||||||
|
VALUES ($1, $2, 1)
|
||||||
|
ON CONFLICT (hash) DO UPDATE
|
||||||
|
SET ref_count = storage.blobs.ref_count + 1",
|
||||||
|
)
|
||||||
|
.bind(&hash)
|
||||||
|
.bind(length as i64)
|
||||||
|
.execute(pool.as_ref())
|
||||||
|
.await
|
||||||
|
.map_err(|e| {
|
||||||
|
DomainError::internal_error(
|
||||||
|
"Dedup",
|
||||||
|
format!("Failed to upsert chunk: {}", e),
|
||||||
)
|
)
|
||||||
.bind(&hash)
|
})?;
|
||||||
.bind(length as i64)
|
|
||||||
.execute(pool.as_ref())
|
|
||||||
.await
|
|
||||||
.map_err(|e| {
|
|
||||||
DomainError::internal_error(
|
|
||||||
"Dedup",
|
|
||||||
format!("Failed to upsert chunk: {}", e),
|
|
||||||
)
|
|
||||||
})?;
|
|
||||||
}
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
@@ -597,13 +591,12 @@ impl DedupService {
|
|||||||
.collect()
|
.collect()
|
||||||
.await;
|
.await;
|
||||||
|
|
||||||
// All operations must succeed. Order preservation is not needed
|
|
||||||
// here — chunk_hashes/chunk_sizes are derived from the input
|
|
||||||
// `chunks` slice which keeps the original CDC order.
|
|
||||||
for result in results {
|
for result in results {
|
||||||
result?;
|
result?;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// chunk_hashes/chunk_sizes keep the full per-occurrence CDC sequence —
|
||||||
|
// the manifest needs every chunk, in order, to reassemble the file.
|
||||||
let chunk_hashes: Vec<String> = chunks.iter().map(|c| c.hash.clone()).collect();
|
let chunk_hashes: Vec<String> = chunks.iter().map(|c| c.hash.clone()).collect();
|
||||||
let chunk_sizes: Vec<u64> = chunks.iter().map(|c| c.length as u64).collect();
|
let chunk_sizes: Vec<u64> = chunks.iter().map(|c| c.length as u64).collect();
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user