2026-02-14 19:30:49 +01:00
|
|
|
|
//! Content-Addressable Storage with Deduplication (PostgreSQL-backed)
|
2026-02-14 01:29:34 +01:00
|
|
|
|
//!
|
|
|
|
|
|
//! Implements hash-based deduplication to eliminate redundant file storage.
|
2026-03-01 21:47:39 +01:00
|
|
|
|
//! Files are stored by their BLAKE3 hash, and multiple references can point
|
2026-02-14 01:29:34 +01:00
|
|
|
|
//! to the same physical blob.
|
|
|
|
|
|
//!
|
|
|
|
|
|
//! Architecture:
|
|
|
|
|
|
//! ```text
|
|
|
|
|
|
//! ┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
|
2026-02-14 19:30:49 +01:00
|
|
|
|
//! │ storage.files │────▶│ storage.blobs │────▶│ Blob Store │
|
|
|
|
|
|
//! │ (references) │ │ (PG dedup index)│ │ (.blobs/ on FS) │
|
2026-02-14 01:29:34 +01:00
|
|
|
|
//! └─────────────────┘ └─────────────────┘ └─────────────────┘
|
|
|
|
|
|
//! ```
|
|
|
|
|
|
//!
|
2026-02-14 19:30:49 +01:00
|
|
|
|
//! The dedup index lives in PostgreSQL (`storage.blobs`) — no in-memory
|
2026-02-25 23:31:51 +01:00
|
|
|
|
//! HashMap, no JSON file, no WAL.
|
|
|
|
|
|
//!
|
2026-04-11 16:17:47 +02:00
|
|
|
|
//! **Write-first strategy** (store_from_file):
|
2026-02-25 23:31:51 +01:00
|
|
|
|
//! 1. Write/move the blob file to disk *before* touching PostgreSQL.
|
|
|
|
|
|
//! 2. Single `INSERT … ON CONFLICT … RETURNING ref_count` upsert
|
|
|
|
|
|
//! (~2-4 ms) — no explicit transaction, no `SELECT FOR UPDATE`.
|
|
|
|
|
|
//! 3. PG connection is never held during disk I/O.
|
|
|
|
|
|
//!
|
|
|
|
|
|
//! `remove_reference` retains `SELECT … FOR UPDATE` inside a short
|
|
|
|
|
|
//! transaction because it must atomically decide whether to delete the
|
|
|
|
|
|
//! row *and* the blob file.
|
2026-02-14 19:30:49 +01:00
|
|
|
|
//!
|
2026-02-14 01:29:34 +01:00
|
|
|
|
//! Benefits:
|
2026-02-14 19:30:49 +01:00
|
|
|
|
//! - ACID durability — crash-safe, zero orphaned index entries
|
2026-02-25 23:31:51 +01:00
|
|
|
|
//! - PG connections never blocked by disk I/O (write-first)
|
2026-02-14 01:29:34 +01:00
|
|
|
|
//! - 30-50% storage reduction typical
|
|
|
|
|
|
//! - Faster uploads for existing content (instant dedup)
|
|
|
|
|
|
|
|
|
|
|
|
use bytes::Bytes;
|
2026-02-23 23:43:59 +01:00
|
|
|
|
use futures::stream::{self, StreamExt};
|
2026-02-24 10:45:38 +01:00
|
|
|
|
use futures::{Stream, TryStreamExt};
|
2026-03-01 21:47:39 +01:00
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
use sqlx::PgPool;
|
2026-02-14 01:29:34 +01:00
|
|
|
|
use std::path::{Path, PathBuf};
|
2026-02-15 17:53:25 +01:00
|
|
|
|
use std::pin::Pin;
|
2026-02-14 01:29:34 +01:00
|
|
|
|
use std::sync::Arc;
|
|
|
|
|
|
use tokio::fs::{self, File};
|
2026-02-23 00:51:46 +01:00
|
|
|
|
use tokio::io::{AsyncReadExt, AsyncSeekExt};
|
2026-02-15 17:53:25 +01:00
|
|
|
|
use tokio_util::io::ReaderStream;
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
|
|
|
|
|
use crate::application::ports::dedup_ports::{
|
|
|
|
|
|
BlobMetadataDto, DedupPort, DedupResultDto, DedupStatsDto,
|
|
|
|
|
|
};
|
|
|
|
|
|
use crate::domain::errors::{DomainError, ErrorKind};
|
|
|
|
|
|
|
2026-02-23 00:51:46 +01:00
|
|
|
|
/// Chunk size for streaming file reads (256 KB)
|
2026-02-15 17:53:25 +01:00
|
|
|
|
const STREAM_CHUNK_SIZE: usize = 256 * 1024;
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
/// Content-Addressable Storage Service (PostgreSQL-backed)
|
2026-02-14 01:29:34 +01:00
|
|
|
|
pub struct DedupService {
|
2026-02-14 19:30:49 +01:00
|
|
|
|
/// Root directory for blob storage on the filesystem
|
2026-02-14 01:29:34 +01:00
|
|
|
|
blob_root: PathBuf,
|
|
|
|
|
|
/// Root directory for temporary files during upload
|
|
|
|
|
|
temp_root: PathBuf,
|
2026-02-24 19:28:00 +01:00
|
|
|
|
/// PostgreSQL connection pool (dedup index in `storage.blobs`) — primary,
|
2026-04-11 16:17:47 +02:00
|
|
|
|
/// used by request-path operations (store_from_file, etc.).
|
2026-02-14 19:30:49 +01:00
|
|
|
|
pool: Arc<PgPool>,
|
2026-02-24 19:28:00 +01:00
|
|
|
|
/// Isolated maintenance pool for long-running operations
|
|
|
|
|
|
/// (verify_integrity, garbage_collect) that must never starve the primary.
|
|
|
|
|
|
maintenance_pool: Arc<PgPool>,
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-03-02 04:46:13 +01:00
|
|
|
|
/// Compile-time lookup table for the 256 two-digit lowercase hex prefixes ("00"…"ff").
|
|
|
|
|
|
/// Avoids a `format!("{:02x}", i)` allocation on every iteration of `initialize()`.
|
|
|
|
|
|
static HEX_PREFIXES: [&str; 256] = [
|
|
|
|
|
|
"00", "01", "02", "03", "04", "05", "06", "07", "08", "09", "0a", "0b", "0c", "0d", "0e", "0f",
|
|
|
|
|
|
"10", "11", "12", "13", "14", "15", "16", "17", "18", "19", "1a", "1b", "1c", "1d", "1e", "1f",
|
|
|
|
|
|
"20", "21", "22", "23", "24", "25", "26", "27", "28", "29", "2a", "2b", "2c", "2d", "2e", "2f",
|
|
|
|
|
|
"30", "31", "32", "33", "34", "35", "36", "37", "38", "39", "3a", "3b", "3c", "3d", "3e", "3f",
|
|
|
|
|
|
"40", "41", "42", "43", "44", "45", "46", "47", "48", "49", "4a", "4b", "4c", "4d", "4e", "4f",
|
|
|
|
|
|
"50", "51", "52", "53", "54", "55", "56", "57", "58", "59", "5a", "5b", "5c", "5d", "5e", "5f",
|
|
|
|
|
|
"60", "61", "62", "63", "64", "65", "66", "67", "68", "69", "6a", "6b", "6c", "6d", "6e", "6f",
|
|
|
|
|
|
"70", "71", "72", "73", "74", "75", "76", "77", "78", "79", "7a", "7b", "7c", "7d", "7e", "7f",
|
|
|
|
|
|
"80", "81", "82", "83", "84", "85", "86", "87", "88", "89", "8a", "8b", "8c", "8d", "8e", "8f",
|
|
|
|
|
|
"90", "91", "92", "93", "94", "95", "96", "97", "98", "99", "9a", "9b", "9c", "9d", "9e", "9f",
|
|
|
|
|
|
"a0", "a1", "a2", "a3", "a4", "a5", "a6", "a7", "a8", "a9", "aa", "ab", "ac", "ad", "ae", "af",
|
|
|
|
|
|
"b0", "b1", "b2", "b3", "b4", "b5", "b6", "b7", "b8", "b9", "ba", "bb", "bc", "bd", "be", "bf",
|
|
|
|
|
|
"c0", "c1", "c2", "c3", "c4", "c5", "c6", "c7", "c8", "c9", "ca", "cb", "cc", "cd", "ce", "cf",
|
|
|
|
|
|
"d0", "d1", "d2", "d3", "d4", "d5", "d6", "d7", "d8", "d9", "da", "db", "dc", "dd", "de", "df",
|
|
|
|
|
|
"e0", "e1", "e2", "e3", "e4", "e5", "e6", "e7", "e8", "e9", "ea", "eb", "ec", "ed", "ee", "ef",
|
|
|
|
|
|
"f0", "f1", "f2", "f3", "f4", "f5", "f6", "f7", "f8", "f9", "fa", "fb", "fc", "fd", "fe", "ff",
|
|
|
|
|
|
];
|
|
|
|
|
|
|
2026-02-14 01:29:34 +01:00
|
|
|
|
impl DedupService {
|
2026-02-14 19:30:49 +01:00
|
|
|
|
/// Create a new dedup service backed by PostgreSQL.
|
2026-02-24 19:28:00 +01:00
|
|
|
|
///
|
|
|
|
|
|
/// * `pool` — primary pool for request-path operations.
|
|
|
|
|
|
/// * `maintenance_pool` — isolated pool for verify_integrity / garbage_collect.
|
|
|
|
|
|
pub fn new(storage_root: &Path, pool: Arc<PgPool>, maintenance_pool: Arc<PgPool>) -> Self {
|
2026-02-14 01:29:34 +01:00
|
|
|
|
let blob_root = storage_root.join(".blobs");
|
|
|
|
|
|
let temp_root = storage_root.join(".dedup_temp");
|
|
|
|
|
|
|
|
|
|
|
|
Self {
|
|
|
|
|
|
blob_root,
|
|
|
|
|
|
temp_root,
|
2026-02-14 19:30:49 +01:00
|
|
|
|
pool,
|
2026-02-24 19:28:00 +01:00
|
|
|
|
maintenance_pool,
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-03-04 14:02:15 +01:00
|
|
|
|
/// Creates a stub instance for testing — never hits PG or the filesystem.
|
2026-03-04 21:40:38 +01:00
|
|
|
|
#[cfg(any(test, feature = "integration_tests"))]
|
2026-03-04 14:02:15 +01:00
|
|
|
|
pub fn new_stub() -> Self {
|
|
|
|
|
|
let stub_pool = Arc::new(
|
|
|
|
|
|
sqlx::pool::PoolOptions::<sqlx::Postgres>::new()
|
|
|
|
|
|
.max_connections(1)
|
|
|
|
|
|
.connect_lazy("postgres://invalid:5432/none")
|
|
|
|
|
|
.unwrap(),
|
|
|
|
|
|
);
|
|
|
|
|
|
Self {
|
|
|
|
|
|
blob_root: std::path::PathBuf::from("/tmp/oxicloud_stub_blobs"),
|
|
|
|
|
|
temp_root: std::path::PathBuf::from("/tmp/oxicloud_stub_temp"),
|
|
|
|
|
|
pool: stub_pool.clone(),
|
|
|
|
|
|
maintenance_pool: stub_pool,
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
/// Initialize the service (create blob directories on the filesystem).
|
|
|
|
|
|
pub async fn initialize(&self) -> Result<(), DomainError> {
|
2026-02-14 01:29:34 +01:00
|
|
|
|
// Create directories
|
2026-02-14 19:30:49 +01:00
|
|
|
|
fs::create_dir_all(&self.blob_root)
|
|
|
|
|
|
.await
|
|
|
|
|
|
.map_err(DomainError::from)?;
|
|
|
|
|
|
fs::create_dir_all(&self.temp_root)
|
|
|
|
|
|
.await
|
|
|
|
|
|
.map_err(DomainError::from)?;
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
|
|
|
|
|
// Create hash prefix directories (00-ff)
|
2026-03-02 04:46:13 +01:00
|
|
|
|
for prefix in &HEX_PREFIXES {
|
|
|
|
|
|
fs::create_dir_all(self.blob_root.join(prefix))
|
2026-02-14 19:30:49 +01:00
|
|
|
|
.await
|
|
|
|
|
|
.map_err(DomainError::from)?;
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
// Log existing blob stats from PG
|
|
|
|
|
|
let count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM storage.blobs")
|
|
|
|
|
|
.fetch_one(self.pool.as_ref())
|
|
|
|
|
|
.await
|
|
|
|
|
|
.unwrap_or(0);
|
|
|
|
|
|
|
|
|
|
|
|
let total_bytes: i64 =
|
|
|
|
|
|
sqlx::query_scalar("SELECT COALESCE(SUM(size), 0) FROM storage.blobs")
|
|
|
|
|
|
.fetch_one(self.pool.as_ref())
|
|
|
|
|
|
.await
|
|
|
|
|
|
.unwrap_or(0);
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
|
|
|
|
|
tracing::info!(
|
2026-02-14 19:30:49 +01:00
|
|
|
|
"Dedup service initialized (PostgreSQL-backed): {} blobs, {} bytes stored",
|
|
|
|
|
|
count,
|
|
|
|
|
|
total_bytes
|
2026-02-14 01:29:34 +01:00
|
|
|
|
);
|
|
|
|
|
|
|
|
|
|
|
|
Ok(())
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
// ── Path helpers ─────────────────────────────────────────────
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
/// Get the blob path for a given hash.
|
2026-02-14 01:29:34 +01:00
|
|
|
|
pub fn blob_path(&self, hash: &str) -> PathBuf {
|
|
|
|
|
|
let prefix = &hash[0..2];
|
|
|
|
|
|
self.blob_root.join(prefix).join(format!("{}.blob", hash))
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
// ── Hash helpers ─────────────────────────────────────────────
|
|
|
|
|
|
|
2026-03-01 21:47:39 +01:00
|
|
|
|
/// Calculate BLAKE3 hash of a file (~5× faster than SHA-256).
|
2026-02-23 00:51:46 +01:00
|
|
|
|
///
|
|
|
|
|
|
/// Runs entirely on `spawn_blocking` with synchronous I/O so the Tokio
|
2026-03-03 11:18:40 +01:00
|
|
|
|
/// worker threads are never blocked by CPU-bound hashing.
|
|
|
|
|
|
///
|
2026-03-06 22:14:43 +01:00
|
|
|
|
/// Uses memory-mapped I/O (`update_mmap_rayon`) which avoids loading the
|
|
|
|
|
|
/// entire file into the heap. The OS pages in data on demand and BLAKE3
|
|
|
|
|
|
/// parallelises the computation across all available cores via rayon.
|
|
|
|
|
|
/// Peak RAM for a 500 MB file is only a few MB of active pages instead
|
|
|
|
|
|
/// of the full 500 MB.
|
2026-02-14 01:29:34 +01:00
|
|
|
|
pub async fn hash_file(path: &Path) -> std::io::Result<String> {
|
2026-02-23 00:51:46 +01:00
|
|
|
|
let path = path.to_path_buf();
|
|
|
|
|
|
tokio::task::spawn_blocking(move || {
|
2026-03-01 21:47:39 +01:00
|
|
|
|
let mut hasher = blake3::Hasher::new();
|
2026-03-06 22:14:43 +01:00
|
|
|
|
hasher.update_mmap_rayon(&path)?;
|
2026-03-01 21:47:39 +01:00
|
|
|
|
Ok(hasher.finalize().to_hex().to_string())
|
2026-02-23 00:51:46 +01:00
|
|
|
|
})
|
|
|
|
|
|
.await
|
|
|
|
|
|
.expect("hash_file: spawn_blocking task panicked")
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
// ── Core store operations ────────────────────────────────────
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
/// Store content with deduplication (streaming from file).
|
2026-02-25 23:31:51 +01:00
|
|
|
|
///
|
|
|
|
|
|
/// **Write-first strategy**: the source file is moved/copied to the
|
|
|
|
|
|
/// blob store *before* touching PostgreSQL, so the PG connection is
|
|
|
|
|
|
/// never held during disk I/O.
|
2026-02-15 17:53:25 +01:00
|
|
|
|
///
|
|
|
|
|
|
/// If `pre_computed_hash` is `Some`, the file will NOT be re-read for
|
2026-03-02 02:13:08 +01:00
|
|
|
|
/// BLAKE3 — saving one full sequential read (the biggest I/O win).
|
2026-02-14 01:29:34 +01:00
|
|
|
|
pub async fn store_from_file(
|
|
|
|
|
|
&self,
|
|
|
|
|
|
source_path: &Path,
|
|
|
|
|
|
content_type: Option<String>,
|
2026-02-15 17:53:25 +01:00
|
|
|
|
pre_computed_hash: Option<String>,
|
2026-02-14 19:30:49 +01:00
|
|
|
|
) -> Result<DedupResultDto, DomainError> {
|
2026-02-14 01:29:34 +01:00
|
|
|
|
let file_size = fs::metadata(source_path)
|
|
|
|
|
|
.await
|
2026-02-14 19:30:49 +01:00
|
|
|
|
.map_err(|e| {
|
2026-02-15 18:04:32 +01:00
|
|
|
|
DomainError::internal_error("Dedup", format!("Failed to get file metadata: {}", e))
|
2026-02-14 19:30:49 +01:00
|
|
|
|
})?
|
2026-02-14 01:29:34 +01:00
|
|
|
|
.len();
|
|
|
|
|
|
|
2026-02-15 17:53:25 +01:00
|
|
|
|
// Use pre-computed hash if available, otherwise calculate (streaming)
|
|
|
|
|
|
let hash = match pre_computed_hash {
|
|
|
|
|
|
Some(h) => h,
|
|
|
|
|
|
None => Self::hash_file(source_path)
|
|
|
|
|
|
.await
|
|
|
|
|
|
.map_err(DomainError::from)?,
|
|
|
|
|
|
};
|
2026-02-14 19:30:49 +01:00
|
|
|
|
|
2026-02-25 23:31:51 +01:00
|
|
|
|
let blob_path = self.blob_path(&hash);
|
2026-02-14 19:30:49 +01:00
|
|
|
|
|
2026-02-25 23:31:51 +01:00
|
|
|
|
// ── Phase 1: Move/place blob on disk (NO PG connection held) ─
|
|
|
|
|
|
//
|
|
|
|
|
|
// If the blob file already exists on disk, the source is simply
|
|
|
|
|
|
// deleted — the file content is identical by definition.
|
2026-03-06 22:27:43 +01:00
|
|
|
|
if fs::try_exists(&blob_path).await.unwrap_or(false) {
|
2026-02-25 23:31:51 +01:00
|
|
|
|
// Blob already on disk — discard the source file
|
|
|
|
|
|
let _ = fs::remove_file(source_path).await;
|
|
|
|
|
|
} else {
|
2026-03-02 00:12:33 +01:00
|
|
|
|
// Parent directory (xx/) guaranteed to exist — created by initialize()
|
2026-02-14 19:30:49 +01:00
|
|
|
|
|
2026-02-25 23:31:51 +01:00
|
|
|
|
// rename is atomic on the same filesystem. If source and blob
|
|
|
|
|
|
// dirs live on different filesystems (rare), this falls back to
|
|
|
|
|
|
// copy+delete which is slower but still correct.
|
|
|
|
|
|
if let Err(e) = fs::rename(source_path, &blob_path).await {
|
2026-03-05 17:28:50 -05:00
|
|
|
|
if e.raw_os_error() == Some(18) {
|
|
|
|
|
|
// EXDEV: cross-device link — fall back to copy+delete
|
|
|
|
|
|
fs::copy(source_path, &blob_path).await.map_err(|ce| {
|
|
|
|
|
|
DomainError::internal_error(
|
|
|
|
|
|
"Dedup",
|
|
|
|
|
|
format!("Failed to copy file to blob store: {}", ce),
|
|
|
|
|
|
)
|
|
|
|
|
|
})?;
|
|
|
|
|
|
let _ = fs::remove_file(source_path).await;
|
2026-03-06 22:27:43 +01:00
|
|
|
|
} else if fs::try_exists(&blob_path).await.unwrap_or(false) {
|
2026-03-05 17:28:50 -05:00
|
|
|
|
// Another writer may have placed the blob concurrently
|
2026-02-25 23:31:51 +01:00
|
|
|
|
let _ = fs::remove_file(source_path).await;
|
|
|
|
|
|
tracing::debug!("Blob file placed by concurrent writer: {}", e);
|
|
|
|
|
|
} else {
|
|
|
|
|
|
return Err(DomainError::internal_error(
|
|
|
|
|
|
"Dedup",
|
|
|
|
|
|
format!("Failed to move file to blob store: {}", e),
|
|
|
|
|
|
));
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-25 23:31:51 +01:00
|
|
|
|
// ── Phase 2: Single atomic upsert (~2-4 ms, no explicit TX) ─
|
|
|
|
|
|
let ref_count: i32 = sqlx::query_scalar(
|
2026-02-14 19:30:49 +01:00
|
|
|
|
"INSERT INTO storage.blobs (hash, size, ref_count, content_type)
|
|
|
|
|
|
VALUES ($1, $2, 1, $3)
|
2026-02-25 23:31:51 +01:00
|
|
|
|
ON CONFLICT (hash) DO UPDATE SET ref_count = storage.blobs.ref_count + 1
|
|
|
|
|
|
RETURNING ref_count",
|
2026-02-14 19:30:49 +01:00
|
|
|
|
)
|
|
|
|
|
|
.bind(&hash)
|
|
|
|
|
|
.bind(file_size as i64)
|
|
|
|
|
|
.bind(&content_type)
|
2026-02-25 23:31:51 +01:00
|
|
|
|
.fetch_one(self.pool.as_ref())
|
2026-02-14 19:30:49 +01:00
|
|
|
|
.await
|
|
|
|
|
|
.map_err(|e| {
|
2026-02-25 23:31:51 +01:00
|
|
|
|
DomainError::internal_error("Dedup", format!("Failed to upsert blob: {}", e))
|
2026-02-14 19:30:49 +01:00
|
|
|
|
})?;
|
|
|
|
|
|
|
2026-02-25 23:31:51 +01:00
|
|
|
|
if ref_count > 1 {
|
|
|
|
|
|
tracing::info!(
|
|
|
|
|
|
"DEDUP HIT (file): {} ({} bytes saved)",
|
|
|
|
|
|
&hash[..12],
|
|
|
|
|
|
file_size
|
|
|
|
|
|
);
|
|
|
|
|
|
Ok(DedupResultDto::ExistingBlob {
|
|
|
|
|
|
hash,
|
|
|
|
|
|
size: file_size,
|
|
|
|
|
|
blob_path,
|
|
|
|
|
|
saved_bytes: file_size,
|
|
|
|
|
|
})
|
|
|
|
|
|
} else {
|
|
|
|
|
|
tracing::info!("NEW BLOB (file): {} ({} bytes)", &hash[..12], file_size);
|
|
|
|
|
|
Ok(DedupResultDto::NewBlob {
|
|
|
|
|
|
hash,
|
|
|
|
|
|
size: file_size,
|
|
|
|
|
|
blob_path,
|
|
|
|
|
|
})
|
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
// ── Reference counting ───────────────────────────────────────
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
/// Check if a blob with the given hash exists in the PG index.
|
|
|
|
|
|
pub async fn blob_exists(&self, hash: &str) -> bool {
|
|
|
|
|
|
sqlx::query_scalar::<_, bool>("SELECT EXISTS(SELECT 1 FROM storage.blobs WHERE hash = $1)")
|
|
|
|
|
|
.bind(hash)
|
|
|
|
|
|
.fetch_one(self.pool.as_ref())
|
|
|
|
|
|
.await
|
|
|
|
|
|
.unwrap_or(false)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-03-05 13:15:34 +01:00
|
|
|
|
/// Returns `true` if `user_id` owns at least one (non-trashed) file that
|
|
|
|
|
|
/// references the blob identified by `hash`.
|
|
|
|
|
|
///
|
|
|
|
|
|
/// Used by the dedup API handlers to enforce per-user access control on
|
|
|
|
|
|
/// the content-addressed blob store.
|
|
|
|
|
|
pub async fn user_owns_blob_reference(&self, hash: &str, user_id: &str) -> bool {
|
|
|
|
|
|
sqlx::query_scalar::<_, bool>(
|
|
|
|
|
|
"SELECT EXISTS(SELECT 1 FROM storage.files WHERE blob_hash = $1 AND user_id = $2 AND NOT is_trashed)",
|
|
|
|
|
|
)
|
|
|
|
|
|
.bind(hash)
|
|
|
|
|
|
.bind(user_id)
|
|
|
|
|
|
.fetch_one(self.pool.as_ref())
|
|
|
|
|
|
.await
|
|
|
|
|
|
.unwrap_or(false)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
/// Get metadata for a blob from PostgreSQL.
|
|
|
|
|
|
pub async fn get_blob_metadata(&self, hash: &str) -> Option<BlobMetadataDto> {
|
|
|
|
|
|
let row = sqlx::query_as::<_, (String, i64, i32, Option<String>)>(
|
|
|
|
|
|
"SELECT hash, size, ref_count, content_type FROM storage.blobs WHERE hash = $1",
|
|
|
|
|
|
)
|
|
|
|
|
|
.bind(hash)
|
|
|
|
|
|
.fetch_optional(self.pool.as_ref())
|
|
|
|
|
|
.await
|
|
|
|
|
|
.ok()
|
|
|
|
|
|
.flatten()?;
|
|
|
|
|
|
|
|
|
|
|
|
Some(BlobMetadataDto {
|
|
|
|
|
|
hash: row.0,
|
|
|
|
|
|
size: row.1 as u64,
|
|
|
|
|
|
ref_count: row.2 as u32,
|
|
|
|
|
|
content_type: row.3,
|
|
|
|
|
|
})
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
/// Add a reference to a blob (increment ref_count).
|
|
|
|
|
|
pub async fn add_reference(&self, hash: &str) -> Result<(), DomainError> {
|
|
|
|
|
|
let rows_affected =
|
|
|
|
|
|
sqlx::query("UPDATE storage.blobs SET ref_count = ref_count + 1 WHERE hash = $1")
|
|
|
|
|
|
.bind(hash)
|
|
|
|
|
|
.execute(self.pool.as_ref())
|
|
|
|
|
|
.await
|
|
|
|
|
|
.map_err(|e| {
|
|
|
|
|
|
DomainError::internal_error(
|
|
|
|
|
|
"Dedup",
|
|
|
|
|
|
format!("Failed to increment ref_count: {}", e),
|
|
|
|
|
|
)
|
|
|
|
|
|
})?
|
|
|
|
|
|
.rows_affected();
|
|
|
|
|
|
|
|
|
|
|
|
if rows_affected == 0 {
|
|
|
|
|
|
return Err(DomainError::new(
|
|
|
|
|
|
ErrorKind::NotFound,
|
|
|
|
|
|
"Blob",
|
|
|
|
|
|
format!("Blob not found: {}", hash),
|
|
|
|
|
|
));
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
Ok(())
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
/// Remove a reference from a blob.
|
|
|
|
|
|
///
|
|
|
|
|
|
/// Uses a single transaction with `SELECT … FOR UPDATE` to atomically
|
|
|
|
|
|
/// decrement ref_count and delete the row + blob file if it reaches 0.
|
|
|
|
|
|
/// Returns `true` if the blob was deleted.
|
|
|
|
|
|
pub async fn remove_reference(&self, hash: &str) -> Result<bool, DomainError> {
|
|
|
|
|
|
let mut tx = self.pool.begin().await.map_err(|e| {
|
|
|
|
|
|
DomainError::internal_error("Dedup", format!("Failed to begin transaction: {}", e))
|
|
|
|
|
|
})?;
|
|
|
|
|
|
|
2026-04-11 16:17:47 +02:00
|
|
|
|
// Lock the row exclusively — prevents concurrent store_from_file from
|
2026-02-14 19:30:49 +01:00
|
|
|
|
// incrementing ref_count while we might be deleting
|
|
|
|
|
|
let row = sqlx::query_as::<_, (i32, i64)>(
|
|
|
|
|
|
"SELECT ref_count, size FROM storage.blobs WHERE hash = $1 FOR UPDATE",
|
|
|
|
|
|
)
|
|
|
|
|
|
.bind(hash)
|
|
|
|
|
|
.fetch_optional(&mut *tx)
|
|
|
|
|
|
.await
|
|
|
|
|
|
.map_err(|e| {
|
|
|
|
|
|
DomainError::internal_error("Dedup", format!("Failed to lock blob row: {}", e))
|
|
|
|
|
|
})?;
|
|
|
|
|
|
|
|
|
|
|
|
let Some((ref_count, _size)) = row else {
|
|
|
|
|
|
// Blob doesn't exist — nothing to do
|
|
|
|
|
|
tx.rollback().await.ok();
|
|
|
|
|
|
return Ok(false);
|
2026-02-14 01:29:34 +01:00
|
|
|
|
};
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
let new_ref_count = (ref_count - 1).max(0);
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
if new_ref_count == 0 {
|
|
|
|
|
|
// Last reference — delete row from PG
|
|
|
|
|
|
sqlx::query("DELETE FROM storage.blobs WHERE hash = $1")
|
|
|
|
|
|
.bind(hash)
|
|
|
|
|
|
.execute(&mut *tx)
|
|
|
|
|
|
.await
|
|
|
|
|
|
.map_err(|e| {
|
|
|
|
|
|
DomainError::internal_error(
|
|
|
|
|
|
"Dedup",
|
|
|
|
|
|
format!("Failed to delete blob row: {}", e),
|
|
|
|
|
|
)
|
|
|
|
|
|
})?;
|
|
|
|
|
|
|
|
|
|
|
|
tx.commit().await.map_err(|e| {
|
|
|
|
|
|
DomainError::internal_error("Dedup", format!("Failed to commit: {}", e))
|
|
|
|
|
|
})?;
|
|
|
|
|
|
|
|
|
|
|
|
// Delete blob file AFTER committing PG — the row is gone, so no
|
2026-04-11 16:17:47 +02:00
|
|
|
|
// concurrent store_from_file can resurrect a reference to this hash.
|
2026-02-14 19:30:49 +01:00
|
|
|
|
let blob_path = self.blob_path(hash);
|
|
|
|
|
|
if let Err(e) = fs::remove_file(&blob_path).await {
|
|
|
|
|
|
tracing::warn!("Failed to delete blob file {}: {}", hash, e);
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
tracing::info!("BLOB DELETED: {} (no more references)", &hash[..12]);
|
2026-02-14 01:29:34 +01:00
|
|
|
|
Ok(true)
|
|
|
|
|
|
} else {
|
2026-02-14 19:30:49 +01:00
|
|
|
|
// Still has references — just decrement
|
|
|
|
|
|
sqlx::query("UPDATE storage.blobs SET ref_count = $1 WHERE hash = $2")
|
|
|
|
|
|
.bind(new_ref_count)
|
|
|
|
|
|
.bind(hash)
|
|
|
|
|
|
.execute(&mut *tx)
|
|
|
|
|
|
.await
|
|
|
|
|
|
.map_err(|e| {
|
|
|
|
|
|
DomainError::internal_error(
|
|
|
|
|
|
"Dedup",
|
|
|
|
|
|
format!("Failed to decrement ref_count: {}", e),
|
|
|
|
|
|
)
|
|
|
|
|
|
})?;
|
|
|
|
|
|
|
|
|
|
|
|
tx.commit().await.map_err(|e| {
|
|
|
|
|
|
DomainError::internal_error("Dedup", format!("Failed to commit: {}", e))
|
|
|
|
|
|
})?;
|
|
|
|
|
|
|
|
|
|
|
|
tracing::debug!("Reference removed from blob {}", &hash[..12]);
|
2026-02-14 01:29:34 +01:00
|
|
|
|
Ok(false)
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
// ── Read operations ──────────────────────────────────────────
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
2026-02-15 17:53:25 +01:00
|
|
|
|
/// Stream blob content in 64 KB chunks — constant memory (~64 KB per stream).
|
|
|
|
|
|
///
|
|
|
|
|
|
/// A 1 GB file uses the same ~64 KB as a 1 KB file.
|
|
|
|
|
|
pub async fn read_blob_stream(
|
|
|
|
|
|
&self,
|
|
|
|
|
|
hash: &str,
|
|
|
|
|
|
) -> Result<Pin<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>>, DomainError>
|
|
|
|
|
|
{
|
|
|
|
|
|
let blob_path = self.blob_path(hash);
|
|
|
|
|
|
let file = File::open(&blob_path).await.map_err(|e| {
|
|
|
|
|
|
DomainError::new(
|
|
|
|
|
|
ErrorKind::NotFound,
|
|
|
|
|
|
"Blob",
|
|
|
|
|
|
format!("Failed to open blob {}: {}", hash, e),
|
|
|
|
|
|
)
|
|
|
|
|
|
})?;
|
2026-02-15 18:04:32 +01:00
|
|
|
|
Ok(Box::pin(ReaderStream::with_capacity(
|
|
|
|
|
|
file,
|
|
|
|
|
|
STREAM_CHUNK_SIZE,
|
|
|
|
|
|
)))
|
2026-02-15 17:53:25 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
/// Stream a byte range of a blob — only reads the requested portion.
|
|
|
|
|
|
///
|
|
|
|
|
|
/// Uses seek + take so a 1 MB range request on a 1 GB file only reads 1 MB.
|
|
|
|
|
|
pub async fn read_blob_range_stream(
|
|
|
|
|
|
&self,
|
|
|
|
|
|
hash: &str,
|
|
|
|
|
|
start: u64,
|
|
|
|
|
|
end: Option<u64>,
|
|
|
|
|
|
) -> Result<Pin<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>>, DomainError>
|
|
|
|
|
|
{
|
|
|
|
|
|
let blob_path = self.blob_path(hash);
|
|
|
|
|
|
let mut file = File::open(&blob_path).await.map_err(|e| {
|
|
|
|
|
|
DomainError::new(
|
|
|
|
|
|
ErrorKind::NotFound,
|
|
|
|
|
|
"Blob",
|
|
|
|
|
|
format!("Failed to open blob {}: {}", hash, e),
|
|
|
|
|
|
)
|
|
|
|
|
|
})?;
|
|
|
|
|
|
|
|
|
|
|
|
// Seek to the start position
|
|
|
|
|
|
file.seek(std::io::SeekFrom::Start(start))
|
|
|
|
|
|
.await
|
|
|
|
|
|
.map_err(|e| {
|
|
|
|
|
|
DomainError::internal_error("Blob", format!("Failed to seek in blob: {}", e))
|
|
|
|
|
|
})?;
|
|
|
|
|
|
|
|
|
|
|
|
// If an end is specified, limit the read with take()
|
|
|
|
|
|
if let Some(end_pos) = end {
|
|
|
|
|
|
let limit = end_pos.saturating_sub(start);
|
|
|
|
|
|
let limited = file.take(limit);
|
2026-02-15 18:04:32 +01:00
|
|
|
|
Ok(Box::pin(ReaderStream::with_capacity(
|
|
|
|
|
|
limited,
|
|
|
|
|
|
STREAM_CHUNK_SIZE,
|
|
|
|
|
|
)))
|
2026-02-15 17:53:25 +01:00
|
|
|
|
} else {
|
2026-02-15 18:04:32 +01:00
|
|
|
|
Ok(Box::pin(ReaderStream::with_capacity(
|
|
|
|
|
|
file,
|
|
|
|
|
|
STREAM_CHUNK_SIZE,
|
|
|
|
|
|
)))
|
2026-02-15 17:53:25 +01:00
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
/// Get the size of a blob without reading its content.
|
|
|
|
|
|
pub async fn blob_size(&self, hash: &str) -> Result<u64, DomainError> {
|
|
|
|
|
|
let blob_path = self.blob_path(hash);
|
|
|
|
|
|
let meta = fs::metadata(&blob_path).await.map_err(|e| {
|
|
|
|
|
|
DomainError::new(
|
|
|
|
|
|
ErrorKind::NotFound,
|
|
|
|
|
|
"Blob",
|
|
|
|
|
|
format!("Failed to stat blob {}: {}", hash, e),
|
|
|
|
|
|
)
|
|
|
|
|
|
})?;
|
|
|
|
|
|
Ok(meta.len())
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
// ── Statistics (computed from PG) ────────────────────────────
|
|
|
|
|
|
|
|
|
|
|
|
/// Get deduplication statistics by querying PostgreSQL.
|
|
|
|
|
|
pub async fn get_stats(&self) -> DedupStatsDto {
|
|
|
|
|
|
let row = sqlx::query_as::<_, (i64, i64, i64)>(
|
|
|
|
|
|
"SELECT
|
|
|
|
|
|
COUNT(*) AS total_blobs,
|
|
|
|
|
|
COALESCE(SUM(size), 0) AS total_bytes_stored,
|
|
|
|
|
|
COALESCE(SUM(size::BIGINT * ref_count), 0) AS total_bytes_referenced
|
|
|
|
|
|
FROM storage.blobs",
|
|
|
|
|
|
)
|
|
|
|
|
|
.fetch_one(self.pool.as_ref())
|
|
|
|
|
|
.await
|
|
|
|
|
|
.unwrap_or((0, 0, 0));
|
|
|
|
|
|
|
|
|
|
|
|
let total_blobs = row.0 as u64;
|
|
|
|
|
|
let total_bytes_stored = row.1 as u64;
|
|
|
|
|
|
let total_bytes_referenced = row.2 as u64;
|
|
|
|
|
|
let bytes_saved = total_bytes_referenced.saturating_sub(total_bytes_stored);
|
|
|
|
|
|
let dedup_ratio = if total_bytes_stored > 0 {
|
|
|
|
|
|
total_bytes_referenced as f64 / total_bytes_stored as f64
|
|
|
|
|
|
} else {
|
|
|
|
|
|
1.0
|
|
|
|
|
|
};
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
DedupStatsDto {
|
|
|
|
|
|
total_blobs,
|
|
|
|
|
|
total_bytes_stored,
|
|
|
|
|
|
total_bytes_referenced,
|
|
|
|
|
|
bytes_saved,
|
|
|
|
|
|
dedup_hits: 0, // Not tracked per-session — derive from SUM(ref_count - 1)
|
|
|
|
|
|
dedup_ratio,
|
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
// ── Maintenance ──────────────────────────────────────────────
|
|
|
|
|
|
|
|
|
|
|
|
/// Verify integrity of all blobs (PG index vs filesystem).
|
2026-02-23 23:43:59 +01:00
|
|
|
|
///
|
2026-02-24 10:45:38 +01:00
|
|
|
|
/// Uses a **streaming cursor** (`fetch()`) so memory stays O(batch)
|
|
|
|
|
|
/// instead of O(total_blobs). Blobs are verified in micro-batches
|
|
|
|
|
|
/// of `VERIFY_CONCURRENCY` using `buffer_unordered`.
|
2026-02-14 19:30:49 +01:00
|
|
|
|
pub async fn verify_integrity(&self) -> Result<Vec<String>, DomainError> {
|
2026-02-23 23:43:59 +01:00
|
|
|
|
/// Max blobs verified concurrently. Each spawns a blocking
|
2026-03-02 02:13:08 +01:00
|
|
|
|
/// thread for BLAKE3 so this also caps blocking-pool pressure.
|
2026-02-23 23:43:59 +01:00
|
|
|
|
const VERIFY_CONCURRENCY: usize = 16;
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
2026-02-24 10:45:38 +01:00
|
|
|
|
let mut row_stream = sqlx::query_as::<_, (String, i64)>(
|
2026-02-14 19:30:49 +01:00
|
|
|
|
"SELECT hash, size FROM storage.blobs ORDER BY hash",
|
|
|
|
|
|
)
|
2026-02-24 19:28:00 +01:00
|
|
|
|
.fetch(self.maintenance_pool.as_ref());
|
2026-02-24 10:45:38 +01:00
|
|
|
|
|
|
|
|
|
|
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))
|
|
|
|
|
|
})?;
|
|
|
|
|
|
|
|
|
|
|
|
let is_done = maybe_row.is_none();
|
2026-02-14 19:30:49 +01:00
|
|
|
|
|
2026-02-24 10:45:38 +01:00
|
|
|
|
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)| {
|
|
|
|
|
|
let blob_root = blob_root.clone();
|
|
|
|
|
|
async move {
|
|
|
|
|
|
let prefix = &hash[0..2];
|
2026-02-25 10:28:34 +01:00
|
|
|
|
let blob_path = blob_root.join(prefix).join(format!("{}.blob", hash));
|
2026-02-24 10:45:38 +01:00
|
|
|
|
|
|
|
|
|
|
let mut issues = Vec::new();
|
|
|
|
|
|
|
2026-02-24 13:06:40 +01:00
|
|
|
|
// Single async metadata() replaces the previous
|
|
|
|
|
|
// blocking .exists() + separate metadata() — one
|
|
|
|
|
|
// stat() syscall instead of two, and non-blocking.
|
|
|
|
|
|
let file_meta = match fs::metadata(&blob_path).await {
|
|
|
|
|
|
Ok(m) => m,
|
|
|
|
|
|
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
|
|
|
|
|
|
issues.push(format!("{}: file missing on disk", hash));
|
|
|
|
|
|
return issues;
|
|
|
|
|
|
}
|
|
|
|
|
|
Err(e) => {
|
|
|
|
|
|
issues.push(format!("{}: metadata error ({})", hash, e));
|
|
|
|
|
|
return issues;
|
|
|
|
|
|
}
|
|
|
|
|
|
};
|
|
|
|
|
|
|
|
|
|
|
|
// Check size
|
|
|
|
|
|
if file_meta.len() != expected_size as u64 {
|
|
|
|
|
|
issues.push(format!(
|
|
|
|
|
|
"{}: size mismatch (expected: {}, actual: {})",
|
|
|
|
|
|
hash,
|
|
|
|
|
|
expected_size,
|
|
|
|
|
|
file_meta.len(),
|
|
|
|
|
|
));
|
2026-02-24 10:45:38 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Verify hash
|
|
|
|
|
|
match Self::hash_file(&blob_path).await {
|
|
|
|
|
|
Ok(actual_hash) => {
|
|
|
|
|
|
if actual_hash != hash {
|
|
|
|
|
|
issues.push(format!(
|
|
|
|
|
|
"{}: hash mismatch (actual: {})",
|
|
|
|
|
|
hash, actual_hash,
|
|
|
|
|
|
));
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
Err(e) => {
|
|
|
|
|
|
issues.push(format!("{}: read error ({})", hash, e));
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
issues
|
2026-02-23 23:43:59 +01:00
|
|
|
|
}
|
2026-02-24 10:45:38 +01:00
|
|
|
|
})
|
|
|
|
|
|
.buffer_unordered(VERIFY_CONCURRENCY)
|
|
|
|
|
|
.flat_map(stream::iter)
|
|
|
|
|
|
.collect()
|
|
|
|
|
|
.await;
|
|
|
|
|
|
|
|
|
|
|
|
corrupted.extend(issues);
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
if is_done {
|
|
|
|
|
|
break;
|
|
|
|
|
|
}
|
|
|
|
|
|
}
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
|
|
|
|
|
if corrupted.is_empty() {
|
2026-02-23 23:43:59 +01:00
|
|
|
|
tracing::info!("Integrity check passed for {} blobs", total);
|
2026-02-14 01:29:34 +01:00
|
|
|
|
} else {
|
2026-02-14 19:30:49 +01:00
|
|
|
|
tracing::warn!("Integrity check found {} issues", corrupted.len());
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
Ok(corrupted)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 19:30:49 +01:00
|
|
|
|
/// Garbage collect orphaned blobs (ref_count = 0).
|
|
|
|
|
|
///
|
2026-02-25 23:54:09 +01:00
|
|
|
|
/// Deletes in small batches (BATCH_SIZE rows per TX) so that each
|
|
|
|
|
|
/// transaction lasts only a few milliseconds. This avoids:
|
|
|
|
|
|
/// - massive row-lock accumulation in PostgreSQL,
|
|
|
|
|
|
/// - WAL bloat from a single giant DELETE,
|
|
|
|
|
|
/// - blocking concurrent uploads that touch `storage.blobs`.
|
|
|
|
|
|
///
|
|
|
|
|
|
/// Blob files are removed **after** each batch commits, so a crash
|
|
|
|
|
|
/// mid-GC only leaves a few orphan files on disk (reclaimed next run).
|
2026-02-14 19:30:49 +01:00
|
|
|
|
pub async fn garbage_collect(&self) -> Result<(u64, u64), DomainError> {
|
2026-02-25 23:54:09 +01:00
|
|
|
|
/// Max rows deleted per mini-transaction.
|
|
|
|
|
|
const BATCH_SIZE: i64 = 500;
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
2026-02-25 23:54:09 +01:00
|
|
|
|
let mut total_deleted = 0u64;
|
|
|
|
|
|
let mut total_bytes = 0u64;
|
2026-02-14 01:29:34 +01:00
|
|
|
|
|
2026-02-25 23:54:09 +01:00
|
|
|
|
loop {
|
|
|
|
|
|
// Each DELETE is its own implicit TX — short and bounded.
|
|
|
|
|
|
// The `ctid` sub-select is the canonical way to do
|
|
|
|
|
|
// `DELETE … LIMIT` in PostgreSQL.
|
|
|
|
|
|
let batch: Vec<(String, i64)> = sqlx::query_as(
|
|
|
|
|
|
"DELETE FROM storage.blobs
|
|
|
|
|
|
WHERE ctid = ANY(
|
|
|
|
|
|
SELECT ctid FROM storage.blobs
|
|
|
|
|
|
WHERE ref_count = 0
|
|
|
|
|
|
LIMIT $1
|
|
|
|
|
|
)
|
|
|
|
|
|
RETURNING hash, size",
|
|
|
|
|
|
)
|
|
|
|
|
|
.bind(BATCH_SIZE)
|
|
|
|
|
|
.fetch_all(self.maintenance_pool.as_ref())
|
|
|
|
|
|
.await
|
2026-02-26 00:47:32 +01:00
|
|
|
|
.map_err(|e| DomainError::internal_error("Dedup", format!("GC batch failed: {e}")))?;
|
2026-02-25 23:54:09 +01:00
|
|
|
|
|
|
|
|
|
|
if batch.is_empty() {
|
|
|
|
|
|
break;
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// Delete blob files OUTSIDE the TX (already committed).
|
2026-04-12 00:50:10 +02:00
|
|
|
|
// Also clean up any thumbnail files for these blob hashes
|
|
|
|
|
|
// (thumbnails are keyed by blob_hash and live under
|
|
|
|
|
|
// storage_root/.thumbnails/{icon,preview,large}/{hash}.jpg).
|
2026-04-14 10:51:48 +02:00
|
|
|
|
let thumbnails_root = self
|
|
|
|
|
|
.blob_root
|
|
|
|
|
|
.parent()
|
|
|
|
|
|
.unwrap_or(&self.blob_root)
|
|
|
|
|
|
.join(".thumbnails");
|
2026-02-25 23:54:09 +01:00
|
|
|
|
for (hash, size) in &batch {
|
|
|
|
|
|
let blob_path = self.blob_path(hash);
|
|
|
|
|
|
if let Err(e) = fs::remove_file(&blob_path).await {
|
|
|
|
|
|
tracing::warn!("Failed to delete orphan blob file {hash}: {e}");
|
|
|
|
|
|
}
|
2026-04-12 00:50:10 +02:00
|
|
|
|
// Remove associated thumbnail files (best-effort)
|
|
|
|
|
|
for dir in &["icon", "preview", "large"] {
|
|
|
|
|
|
let thumb = thumbnails_root.join(dir).join(format!("{hash}.jpg"));
|
|
|
|
|
|
let _ = fs::remove_file(&thumb).await;
|
|
|
|
|
|
}
|
2026-02-25 23:54:09 +01:00
|
|
|
|
total_bytes += *size as u64;
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
2026-02-25 23:54:09 +01:00
|
|
|
|
total_deleted += batch.len() as u64;
|
|
|
|
|
|
|
|
|
|
|
|
// Yield so uploads / other tasks are not starved.
|
|
|
|
|
|
tokio::task::yield_now().await;
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-25 23:54:09 +01:00
|
|
|
|
if total_deleted > 0 {
|
|
|
|
|
|
tracing::info!("GC: removed {total_deleted} blobs ({total_bytes} bytes)");
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-25 23:54:09 +01:00
|
|
|
|
Ok((total_deleted, total_bytes))
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
// ─── Port implementation ─────────────────────────────────────────────────────
|
|
|
|
|
|
|
|
|
|
|
|
impl DedupPort for DedupService {
|
|
|
|
|
|
async fn store_from_file(
|
|
|
|
|
|
&self,
|
|
|
|
|
|
source_path: &Path,
|
|
|
|
|
|
content_type: Option<String>,
|
2026-02-15 17:53:25 +01:00
|
|
|
|
pre_computed_hash: Option<String>,
|
2026-02-14 01:29:34 +01:00
|
|
|
|
) -> Result<DedupResultDto, DomainError> {
|
2026-02-15 18:04:32 +01:00
|
|
|
|
self.store_from_file(source_path, content_type, pre_computed_hash)
|
|
|
|
|
|
.await
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
async fn blob_exists(&self, hash: &str) -> bool {
|
|
|
|
|
|
self.blob_exists(hash).await
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
async fn get_blob_metadata(&self, hash: &str) -> Option<BlobMetadataDto> {
|
2026-02-14 19:30:49 +01:00
|
|
|
|
self.get_blob_metadata(hash).await
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-15 17:53:25 +01:00
|
|
|
|
async fn read_blob_stream(
|
|
|
|
|
|
&self,
|
|
|
|
|
|
hash: &str,
|
|
|
|
|
|
) -> Result<Pin<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>>, DomainError>
|
|
|
|
|
|
{
|
|
|
|
|
|
self.read_blob_stream(hash).await
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
async fn read_blob_range_stream(
|
|
|
|
|
|
&self,
|
|
|
|
|
|
hash: &str,
|
|
|
|
|
|
start: u64,
|
|
|
|
|
|
end: Option<u64>,
|
|
|
|
|
|
) -> Result<Pin<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>>, DomainError>
|
|
|
|
|
|
{
|
|
|
|
|
|
self.read_blob_range_stream(hash, start, end).await
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
async fn blob_size(&self, hash: &str) -> Result<u64, DomainError> {
|
|
|
|
|
|
self.blob_size(hash).await
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 01:29:34 +01:00
|
|
|
|
async fn add_reference(&self, hash: &str) -> Result<(), DomainError> {
|
2026-02-14 19:30:49 +01:00
|
|
|
|
self.add_reference(hash).await
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
async fn remove_reference(&self, hash: &str) -> Result<bool, DomainError> {
|
2026-02-14 19:30:49 +01:00
|
|
|
|
self.remove_reference(hash).await
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
async fn hash_file(&self, path: &Path) -> Result<String, DomainError> {
|
|
|
|
|
|
DedupService::hash_file(path)
|
|
|
|
|
|
.await
|
|
|
|
|
|
.map_err(DomainError::from)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-03-04 14:02:15 +01:00
|
|
|
|
fn blob_path(&self, hash: &str) -> PathBuf {
|
|
|
|
|
|
self.blob_path(hash)
|
|
|
|
|
|
}
|
|
|
|
|
|
|
2026-02-14 01:29:34 +01:00
|
|
|
|
async fn get_stats(&self) -> DedupStatsDto {
|
2026-02-14 19:30:49 +01:00
|
|
|
|
self.get_stats().await
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
async fn flush(&self) -> Result<(), DomainError> {
|
2026-02-14 19:30:49 +01:00
|
|
|
|
// No-op: PostgreSQL handles persistence automatically via WAL/commit
|
|
|
|
|
|
Ok(())
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
async fn verify_integrity(&self) -> Result<Vec<String>, DomainError> {
|
2026-02-14 19:30:49 +01:00
|
|
|
|
self.verify_integrity().await
|
2026-02-14 01:29:34 +01:00
|
|
|
|
}
|
|
|
|
|
|
}
|