perf(issue#4): stream dedup verify_integrity & garbage_collect

- Replace fetch_all() with fetch() streaming cursor in verify_integrity
  so memory stays O(batch=16) instead of O(total_blobs)
- Replace fetch_all() with fetch() streaming cursor in garbage_collect
  so memory stays O(1) instead of O(orphans)
- Add TryStreamExt import for try_next() on cursors

Eliminates OOM risk with millions of blobs — RAM usage is now constant
regardless of table size.
This commit is contained in:
Dionisio
2026-02-24 10:45:38 +01:00
parent b0235e05c8
commit b628a3a166
+42 -19
View File
@@ -25,7 +25,7 @@
use async_trait::async_trait; use async_trait::async_trait;
use bytes::Bytes; use bytes::Bytes;
use futures::stream::{self, StreamExt}; use futures::stream::{self, StreamExt};
use futures::Stream; use futures::{Stream, TryStreamExt};
use sha2::{Digest, Sha256}; use sha2::{Digest, Sha256};
use sqlx::PgPool; use sqlx::PgPool;
use std::path::{Path, PathBuf}; use std::path::{Path, PathBuf};
@@ -642,26 +642,42 @@ impl DedupService {
/// Verify integrity of all blobs (PG index vs filesystem). /// Verify integrity of all blobs (PG index vs filesystem).
/// ///
/// Hashes up to `VERIFY_CONCURRENCY` blobs in parallel using /// Uses a **streaming cursor** (`fetch()`) so memory stays O(batch)
/// `buffer_unordered`, saturating both disk I/O and CPU cores. /// instead of O(total_blobs). Blobs are verified in micro-batches
/// of `VERIFY_CONCURRENCY` using `buffer_unordered`.
pub async fn verify_integrity(&self) -> Result<Vec<String>, DomainError> { pub async fn verify_integrity(&self) -> Result<Vec<String>, DomainError> {
/// Max blobs verified concurrently. Each spawns a blocking /// Max blobs verified concurrently. Each spawns a blocking
/// thread for SHA-256 so this also caps blocking-pool pressure. /// thread for SHA-256 so this also caps blocking-pool pressure.
const VERIFY_CONCURRENCY: usize = 16; const VERIFY_CONCURRENCY: usize = 16;
let rows = sqlx::query_as::<_, (String, i64)>( let mut row_stream = sqlx::query_as::<_, (String, i64)>(
"SELECT hash, size FROM storage.blobs ORDER BY hash", "SELECT hash, size FROM storage.blobs ORDER BY hash",
) )
.fetch_all(self.pool.as_ref()) .fetch(self.pool.as_ref());
.await
.map_err(|e| { let mut total = 0usize;
let mut corrupted = Vec::<String>::new();
let mut batch = Vec::with_capacity(VERIFY_CONCURRENCY);
loop {
let maybe_row = row_stream.try_next().await.map_err(|e| {
DomainError::internal_error("Dedup", format!("Failed to list blobs: {}", e)) DomainError::internal_error("Dedup", format!("Failed to list blobs: {}", e))
})?; })?;
let total = rows.len(); let is_done = maybe_row.is_none();
let blob_root = self.blob_root.clone();
let corrupted: Vec<String> = stream::iter(rows) if let Some(row) = maybe_row {
total += 1;
batch.push(row);
}
// Flush when batch is full or we've exhausted the cursor
if batch.len() >= VERIFY_CONCURRENCY || (is_done && !batch.is_empty()) {
let blob_root = self.blob_root.clone();
let current_batch =
std::mem::replace(&mut batch, Vec::with_capacity(VERIFY_CONCURRENCY));
let issues: Vec<String> = stream::iter(current_batch)
.map(move |(hash, expected_size)| { .map(move |(hash, expected_size)| {
let blob_root = blob_root.clone(); let blob_root = blob_root.clone();
async move { async move {
@@ -712,6 +728,14 @@ impl DedupService {
.collect() .collect()
.await; .await;
corrupted.extend(issues);
}
if is_done {
break;
}
}
if corrupted.is_empty() { if corrupted.is_empty() {
tracing::info!("Integrity check passed for {} blobs", total); tracing::info!("Integrity check passed for {} blobs", total);
} else { } else {
@@ -724,26 +748,25 @@ impl DedupService {
/// Garbage collect orphaned blobs (ref_count = 0). /// Garbage collect orphaned blobs (ref_count = 0).
/// ///
/// Uses `DELETE … RETURNING` for an atomic "find and remove" operation. /// Uses `DELETE … RETURNING` for an atomic "find and remove" operation.
/// Results are streamed so memory stays O(1) even with millions of orphans.
pub async fn garbage_collect(&self) -> Result<(u64, u64), DomainError> { pub async fn garbage_collect(&self) -> Result<(u64, u64), DomainError> {
let orphans = sqlx::query_as::<_, (String, i64)>( let mut orphan_stream = sqlx::query_as::<_, (String, i64)>(
"DELETE FROM storage.blobs WHERE ref_count = 0 RETURNING hash, size", "DELETE FROM storage.blobs WHERE ref_count = 0 RETURNING hash, size",
) )
.fetch_all(self.pool.as_ref()) .fetch(self.pool.as_ref());
.await
.map_err(|e| {
DomainError::internal_error("Dedup", format!("Failed to garbage collect: {}", e))
})?;
let mut deleted_count = 0u64; let mut deleted_count = 0u64;
let mut deleted_bytes = 0u64; let mut deleted_bytes = 0u64;
for (hash, size) in &orphans { while let Some((hash, size)) = orphan_stream.try_next().await.map_err(|e| {
let blob_path = self.blob_path(hash); DomainError::internal_error("Dedup", format!("Failed to garbage collect: {}", e))
})? {
let blob_path = self.blob_path(&hash);
if let Err(e) = fs::remove_file(&blob_path).await { if let Err(e) = fs::remove_file(&blob_path).await {
tracing::warn!("Failed to delete orphaned blob file {}: {}", hash, e); tracing::warn!("Failed to delete orphaned blob file {}: {}", hash, e);
} }
deleted_count += 1; deleted_count += 1;
deleted_bytes += *size as u64; deleted_bytes += size as u64;
} }
if deleted_count > 0 { if deleted_count > 0 {