perf(#28): parallelize verify_integrity() with buffer_unordered(16)

Replace sequential blob-by-blob SHA-256 verification with
futures::stream::buffer_unordered(16) to hash up to 16 blobs
concurrently.

Each hash_file() already runs on spawn_blocking, so 16 concurrent
verifications saturate both disk I/O queue and CPU cores.

Before: 10K blobs × 13ms = ~130s (NVMe) — 1 core, 1 I/O in flight
After:  10K blobs / 16 concurrency = ~8s — 16 cores, 16 I/O in flight
This commit is contained in:
Dionisio
2026-02-23 23:43:59 +01:00
parent 44113eabdc
commit 12c6914d5d
+38 -13
View File
@@ -24,6 +24,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; use futures::Stream;
use sha2::{Digest, Sha256}; use sha2::{Digest, Sha256};
use sqlx::PgPool; use sqlx::PgPool;
@@ -640,8 +641,13 @@ impl DedupService {
// ── Maintenance ────────────────────────────────────────────── // ── Maintenance ──────────────────────────────────────────────
/// 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
/// `buffer_unordered`, saturating both disk I/O and CPU cores.
pub async fn verify_integrity(&self) -> Result<Vec<String>, DomainError> { pub async fn verify_integrity(&self) -> Result<Vec<String>, DomainError> {
let mut corrupted = Vec::new(); /// Max blobs verified concurrently. Each spawns a blocking
/// thread for SHA-256 so this also caps blocking-pool pressure.
const VERIFY_CONCURRENCY: usize = 16;
let rows = sqlx::query_as::<_, (String, i64)>( let rows = sqlx::query_as::<_, (String, i64)>(
"SELECT hash, size FROM storage.blobs ORDER BY hash", "SELECT hash, size FROM storage.blobs ORDER BY hash",
@@ -652,43 +658,62 @@ impl DedupService {
DomainError::internal_error("Dedup", format!("Failed to list blobs: {}", e)) DomainError::internal_error("Dedup", format!("Failed to list blobs: {}", e))
})?; })?;
for (hash, expected_size) in &rows { let total = rows.len();
let blob_path = self.blob_path(hash); let blob_root = self.blob_root.clone();
let corrupted: Vec<String> = stream::iter(rows)
.map(move |(hash, expected_size)| {
let blob_root = blob_root.clone();
async move {
let prefix = &hash[0..2];
let blob_path =
blob_root.join(prefix).join(format!("{}.blob", hash));
let mut issues = Vec::new();
// Check file exists // Check file exists
if !blob_path.exists() { if !blob_path.exists() {
corrupted.push(format!("{}: file missing on disk", hash)); issues.push(format!("{}: file missing on disk", hash));
continue; return issues;
} }
// Verify hash // Verify hash
match Self::hash_file(&blob_path).await { match Self::hash_file(&blob_path).await {
Ok(actual_hash) => { Ok(actual_hash) => {
if actual_hash != *hash { if actual_hash != hash {
corrupted issues.push(format!(
.push(format!("{}: hash mismatch (actual: {})", hash, actual_hash)); "{}: hash mismatch (actual: {})",
hash, actual_hash,
));
} }
} }
Err(e) => { Err(e) => {
corrupted.push(format!("{}: read error ({})", hash, e)); issues.push(format!("{}: read error ({})", hash, e));
} }
} }
// Check size // Check size
if let Ok(file_meta) = fs::metadata(&blob_path).await if let Ok(file_meta) = fs::metadata(&blob_path).await
&& file_meta.len() != *expected_size as u64 && file_meta.len() != expected_size as u64
{ {
corrupted.push(format!( issues.push(format!(
"{}: size mismatch (expected: {}, actual: {})", "{}: size mismatch (expected: {}, actual: {})",
hash, hash,
expected_size, expected_size,
file_meta.len() file_meta.len(),
)); ));
} }
issues
} }
})
.buffer_unordered(VERIFY_CONCURRENCY)
.flat_map(stream::iter)
.collect()
.await;
if corrupted.is_empty() { if corrupted.is_empty() {
tracing::info!("Integrity check passed for {} blobs", rows.len()); tracing::info!("Integrity check passed for {} blobs", total);
} else { } else {
tracing::warn!("Integrity check found {} issues", corrupted.len()); tracing::warn!("Integrity check found {} issues", corrupted.len());
} }