feat: check thumbnail cleanup on files deletion + correct ref counter

This commit is contained in:
Edouard Vanbelle
2026-05-13 11:33:45 +02:00
parent 9dd5877bb2
commit 78cb37b311
16 changed files with 444 additions and 42 deletions
+1
View File
@@ -1,4 +1,5 @@
pub mod auth_ports;
pub mod blob_lifecycle;
pub mod blob_storage_ports;
pub mod cache_ports;
pub mod calendar_ports;
+16 -3
View File
@@ -16,6 +16,7 @@ use crate::infrastructure::repositories::pg::file_blob_read_repository::FileBlob
use crate::infrastructure::repositories::pg::file_blob_write_repository::FileBlobWriteRepository;
use crate::infrastructure::repositories::pg::folder_db_repository::FolderDbRepository;
use crate::infrastructure::repositories::pg::trash_db_repository::TrashDbRepository;
use crate::infrastructure::services::dedup_service::DedupService;
use crate::infrastructure::services::file_content_cache::FileContentCache;
use crate::infrastructure::services::thumbnail_service::ThumbnailService;
@@ -45,6 +46,10 @@ pub struct TrashService {
/// Port for folder operations (get folder, trash, restore, delete)
folder_storage_port: Arc<FolderDbRepository>,
/// Dedup service — garbage-collected after bulk trash empty to clean up
/// orphaned blob files and thumbnails that the PG trigger cannot reach.
dedup_service: Arc<DedupService>,
/// Thumbnail service for cleaning up thumbnails on permanent delete
thumbnail_service: Option<Arc<ThumbnailService>>,
@@ -62,6 +67,7 @@ impl TrashService {
file_write_port: Arc<FileBlobWriteRepository>,
folder_storage_port: Arc<FolderDbRepository>,
retention_days: u32,
dedup_service: Arc<DedupService>,
thumbnail_service: Option<Arc<ThumbnailService>>,
content_cache: Option<Arc<FileContentCache>>,
) -> Self {
@@ -70,6 +76,7 @@ impl TrashService {
file_read_port,
file_write_port,
folder_storage_port,
dedup_service,
thumbnail_service,
content_cache,
retention_days,
@@ -713,18 +720,24 @@ impl TrashUseCase for TrashService {
Vec::new()
};
// clear_trash() already performs bulk SQL DELETEs in 2 queries:
// clear_trash() performs bulk SQL DELETEs in 2 queries:
// 1. DELETE FROM storage.files WHERE user_id = $1 AND is_trashed = TRUE
// 2. DELETE FROM storage.folders WHERE user_id = $1 AND is_trashed = TRUE
//
// Folder deletion cascades (FK ON DELETE CASCADE) to child folders and
// their files. The PG trigger `trg_files_decrement_blob_ref` automatically
// decrements blob ref_counts for every deleted file row — no Rust-side
// remove_reference() call is needed.
// decrements blob ref_counts for every deleted file row.
//
// Finally it clears the trash_items index for the user.
self.trash_repository.clear_trash(&user_id).await?;
// The PG trigger decremented ref_counts but cannot delete disk files or
// thumbnails. Run garbage_collect() to remove any blobs whose ref_count
// reached 0, along with their blob-keyed thumbnail files.
if let Err(e) = self.dedup_service.garbage_collect().await {
warn!("empty_trash: garbage_collect failed: {:?}", e);
}
// Invalidate content cache for all permanently deleted files.
if let Some(cc) = &self.content_cache {
for file_id in &trashed_file_ids {
+4 -2
View File
@@ -264,7 +264,7 @@ impl AppServiceFactory {
db_pool.clone(),
maintenance_pool.clone(),
)
.with_thumbnail_service(thumbnail_service.clone()),
.add_blob_hook(thumbnail_service.clone()),
);
dedup_service.initialize().await?;
@@ -451,10 +451,12 @@ impl AppServiceFactory {
repos.file_write_repository.clone(),
repos.folder_repository.clone(),
self.config.storage.trash_retention_days,
core.dedup_service.clone(),
Some(core.thumbnail_service.clone()),
Some(core.file_content_cache.clone()),
));
// Initialize cleanup service (bulk-deletes expired items in 2 SQL queries)
let cleanup_service = TrashCleanupService::new(
trash_repo.clone(),
@@ -1010,7 +1012,7 @@ impl CoreServices {
tokio::spawn(async move {
match ds.read_blob_bytes(&hash).await {
Ok(bytes) => {
ts.generate_all_sizes_background_from_bytes(file_id, hash, bytes);
ts.generate_all_sizes_background_from_bytes(file_id, hash, bytes, ds.clone());
}
Err(e) => {
tracing::warn!(
@@ -604,8 +604,27 @@ impl FileWritePort for FileBlobWriteRepository {
}
async fn delete_file_permanently(&self, file_id: &str) -> Result<(), DomainError> {
// Same as delete_file — removes from DB and decrements blob ref
self.delete_file(file_id).await
// Read blob_hash before deletion so we can clean up disk after the
// PG trigger has decremented the ref_count.
let blob_hash: Option<String> = sqlx::query_scalar(
"SELECT blob_hash FROM storage.files WHERE id = $1::uuid",
)
.bind(file_id)
.fetch_optional(self.pool.as_ref())
.await
.map_err(|e| {
DomainError::internal_error("FileBlobWrite", format!("fetch blob_hash: {e}"))
})?;
// DELETE fires trg_files_decrement_blob_ref → storage.blobs.ref_count--
self.delete_file(file_id).await?;
// If the blob is now unreferenced, remove disk file + thumbnails.
if let Some(hash) = blob_hash {
self.dedup.cleanup_if_orphaned(&hash).await;
}
Ok(())
}
async fn copy_folder_tree(
+107 -29
View File
@@ -45,11 +45,11 @@ use tokio::fs;
use tokio::io::{AsyncReadExt, AsyncSeekExt};
use crate::application::ports::blob_storage_ports::BlobStorageBackend;
use crate::application::ports::blob_lifecycle::BlobDeletionHook;
use crate::application::ports::dedup_ports::{
BlobMetadataDto, DedupPort, DedupResultDto, DedupStatsDto,
};
use crate::domain::errors::{DomainError, ErrorKind};
use crate::infrastructure::services::thumbnail_service::ThumbnailService;
// ── CDC Constants ────────────────────────────────────────────────────────────
@@ -84,9 +84,8 @@ pub struct DedupService {
/// Isolated maintenance pool for long-running operations
/// (verify_integrity, garbage_collect) that must never starve the primary.
maintenance_pool: Arc<PgPool>,
/// Optional thumbnail service — when set, blob-hash thumbnails are deleted
/// from disk whenever a blob's ref_count reaches zero.
thumbnail_service: Option<Arc<ThumbnailService>>,
/// Hooks notified when a blob's ref_count reaches zero and it is deleted.
blob_hooks: Vec<Arc<dyn BlobDeletionHook>>,
}
impl DedupService {
@@ -104,17 +103,24 @@ impl DedupService {
backend,
pool,
maintenance_pool,
thumbnail_service: None,
blob_hooks: vec![],
}
}
/// Attach a thumbnail service so that disk thumbnails are cleaned up when
/// a blob's ref_count drops to zero.
pub fn with_thumbnail_service(mut self, svc: Arc<ThumbnailService>) -> Self {
self.thumbnail_service = Some(svc);
/// Register a [`BlobDeletionHook`] to be called whenever a blob's
/// ref_count reaches zero. Hooks are called in registration order.
pub fn add_blob_hook(mut self, hook: Arc<dyn BlobDeletionHook>) -> Self {
self.blob_hooks.push(hook);
self
}
/// Fire all registered hooks for a deleted blob.
async fn fire_blob_hooks(&self, hash: &str) {
for hook in &self.blob_hooks {
hook.on_blob_deleted(hash).await;
}
}
/// Creates a stub instance for testing — never hits PG or the filesystem.
#[cfg(any(test, feature = "integration_tests"))]
pub fn new_stub() -> Self {
@@ -129,7 +135,7 @@ impl DedupService {
backend: Arc::new(LocalBlobBackend::new(Path::new("/tmp/oxicloud_stub_blobs"))),
pool: stub_pool.clone(),
maintenance_pool: stub_pool,
thumbnail_service: None,
blob_hooks: vec![],
}
}
@@ -600,7 +606,7 @@ impl DedupService {
/// references the blob identified by `hash`.
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)",
"SELECT EXISTS(SELECT 1 FROM storage.files WHERE blob_hash = $1 AND user_id = $2::uuid AND NOT is_trashed)",
)
.bind(hash)
.bind(user_id)
@@ -785,10 +791,8 @@ impl DedupService {
}
}
// Bug 4 fix: delete disk thumbnails keyed by file_hash (last reference gone)
if let Some(ts) = &self.thumbnail_service {
ts.delete_blob_thumbnails(file_hash).await;
}
// Bug 4 fix: notify hooks — e.g. thumbnail cleanup keyed by file_hash
self.fire_blob_hooks(file_hash).await;
tracing::info!(
"MANIFEST DELETED: {} ({} chunks, {} orphan chunks removed)",
@@ -867,10 +871,8 @@ impl DedupService {
tracing::warn!("Failed to delete blob file {}: {}", hash, e);
}
// Bug 3 fix: delete disk thumbnails keyed by hash (last reference gone)
if let Some(ts) = &self.thumbnail_service {
ts.delete_blob_thumbnails(hash).await;
}
// Bug 3 fix: notify hooks — e.g. thumbnail cleanup keyed by hash
self.fire_blob_hooks(hash).await;
tracing::info!("BLOB DELETED: {} (no more references)", &hash[..12]);
Ok(true)
@@ -897,6 +899,91 @@ impl DedupService {
}
}
/// Targeted cleanup for a single blob after the PG trigger has already
/// decremented its ref_count. Deletes the blob row, disk file, and
/// blob-keyed thumbnails if ref_count has reached 0.
///
/// Handles both the legacy whole-file blob path (storage.blobs) and the
/// CDC manifest path (storage.chunk_manifests). Best-effort: logs
/// warnings on failure rather than returning an error.
pub async fn cleanup_if_orphaned(&self, hash: &str) {
let short = &hash[..hash.len().min(12)];
// ── CDC manifest path (must run FIRST) ───────────────────
// For single-chunk CDC files file_hash == chunk_hash, so the PG
// trigger on storage.files already decremented storage.blobs.ref_count
// when this function is called. try_dedup_hit increments
// chunk_manifests.ref_count but NOT storage.blobs.ref_count, so
// blobs.ref_count can reach 0 while the manifest still has ref_count > 1
// (other files sharing the same blob). Checking the manifest first
// prevents premature blob + manifest deletion.
let manifest = sqlx::query_as::<_, (i32, Vec<String>)>(
"SELECT ref_count, chunk_hashes \
FROM storage.chunk_manifests WHERE file_hash = $1",
)
.bind(hash)
.fetch_optional(self.pool.as_ref())
.await
.unwrap_or(None);
if let Some((ref_count, chunk_hashes)) = manifest {
if ref_count <= 1 {
// Last reference — remove manifest and all its chunks.
if let Err(e) = self
.remove_manifest_reference(hash, ref_count, &chunk_hashes)
.await
{
tracing::warn!("cleanup_if_orphaned: manifest cleanup failed for {short}: {e}");
}
} else {
// Other files still share this blob: just decrement the manifest
// counter and undo the PG trigger's premature chunk ref_count
// decrement (blobs.ref_count is chunk-level; the manifest is the
// authoritative file-level counter).
sqlx::query(
"UPDATE storage.chunk_manifests \
SET ref_count = ref_count - 1 WHERE file_hash = $1",
)
.bind(hash)
.execute(self.pool.as_ref())
.await
.ok();
// Undo the PG trigger's decrement of storage.blobs.ref_count.
// The trigger fired with blob_hash = file_hash, so only the row
// WHERE hash = file_hash is affected. For single-chunk files
// file_hash == chunk_hash and that row exists; for multi-chunk
// files file_hash is not in storage.blobs, making this a no-op.
sqlx::query("UPDATE storage.blobs SET ref_count = ref_count + 1 WHERE hash = $1")
.bind(hash)
.execute(self.pool.as_ref())
.await
.ok();
tracing::debug!(
"cleanup_if_orphaned: manifest {short} ref_count {ref_count}→{}",
ref_count - 1
);
}
return;
}
// ── Legacy blob path (no manifest) ───────────────────────
let deleted_blob = sqlx::query_scalar::<_, String>(
"DELETE FROM storage.blobs WHERE hash = $1 AND ref_count <= 0 RETURNING hash",
)
.bind(hash)
.fetch_optional(self.pool.as_ref())
.await
.unwrap_or(None);
if deleted_blob.is_some() {
if let Err(e) = self.backend.delete_blob(hash).await {
tracing::warn!("cleanup_if_orphaned: disk delete failed for {short}: {e}");
}
self.fire_blob_hooks(hash).await;
tracing::info!("cleanup_if_orphaned: removed orphaned blob {short}");
}
}
// ── Read operations ──────────────────────────────────────────
/// Stream blob content — CDC-aware with legacy fallback.
@@ -1347,16 +1434,7 @@ impl DedupService {
if let Err(e) = self.backend.delete_blob(hash).await {
tracing::warn!("Failed to delete orphan blob {hash}: {e}");
}
// Clean up thumbnails (best-effort, only local backends)
if let Some(blob_path) = self.backend.local_blob_path(hash)
&& let Some(storage_root) = blob_path.ancestors().nth(3)
{
let thumbnails_root = storage_root.join(".thumbnails");
for dir in &["icon", "preview", "large"] {
let thumb = thumbnails_root.join(dir).join(format!("{hash}.jpg"));
let _ = fs::remove_file(&thumb).await;
}
}
self.fire_blob_hooks(hash).await;
total_bytes += *size as u64;
}
total_deleted += batch.len() as u64;
@@ -25,6 +25,7 @@ use tokio::time::timeout;
use crate::application::ports::thumbnail_ports::{
ThumbnailPort, ThumbnailSize as PortThumbnailSize, ThumbnailStatsDto,
};
use crate::infrastructure::services::dedup_service::DedupService;
use crate::domain::errors::{DomainError, ErrorKind};
/// Thumbnail sizes supported by the system
@@ -831,10 +832,22 @@ impl ThumbnailService {
file_id: String,
blob_hash: String,
original_data: Bytes,
dedup: Arc<DedupService>,
) {
tokio::spawn(async move {
tracing::info!("🖼️ Background thumbnail generation starting: {}", file_id);
// Guard: if the blob was deleted before this task ran, cleanup_if_orphaned
// already fired with no thumbnails on disk — writing them now would leak them.
// Use the DB check (manifest + blobs tables) as the authoritative source.
if !dedup.blob_exists(&blob_hash).await {
tracing::debug!(
"Blob {}… deleted before thumbnail task ran, skipping",
&blob_hash[..blob_hash.len().min(12)]
);
return;
}
let all_exist = {
let mut ok = true;
for size in ThumbnailSize::all() {
@@ -971,6 +984,17 @@ impl ThumbnailService {
}
}
// ─── BlobDeletionHook ────────────────────────────────────────────────────────
impl crate::application::ports::blob_lifecycle::BlobDeletionHook for ThumbnailService {
fn on_blob_deleted<'a>(
&'a self,
blob_hash: &'a str,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send + 'a>> {
Box::pin(async move { self.delete_blob_thumbnails(blob_hash).await })
}
}
// ─── Port implementation ─────────────────────────────────────────────────────
/// Convert port ThumbnailSize to infra ThumbnailSize.
+17 -6
View File
@@ -26,7 +26,9 @@ pub struct HashCheckResponse {
/// If exists, the size of the existing blob
#[serde(skip_serializing_if = "Option::is_none")]
pub existing_size: Option<u64>,
/// If exists, the number of references to this blob
/// Global reference count for this blob across all users.
/// Only populated when the authenticated user has the `admin` role;
/// omitted for regular users to prevent cross-user content inference.
#[serde(skip_serializing_if = "Option::is_none")]
pub ref_count: Option<u32>,
}
@@ -85,7 +87,8 @@ impl DedupHandler {
/// Check if the authenticated user already has a file with the given hash.
///
/// User-scoped: only reveals whether **this user** owns a file that
/// references the blob — never exposes global existence or ref_count.
/// references the blob — never exposes global existence to non-admins.
/// Admins additionally receive the global `ref_count` in the response.
///
/// GET /api/dedup/check/{hash}
pub(super) async fn check_hash_impl(
@@ -113,13 +116,20 @@ impl DedupHandler {
.await;
if user_has_it {
// Fetch size from metadata (safe — user owns a reference)
let size = dedup.get_blob_metadata(&hash).await.map(|m| m.size);
// Fetch size from metadata (safe — user owns a reference).
// Admins also get the global ref_count for dedup accounting tests.
let metadata = dedup.get_blob_metadata(&hash).await;
let size = metadata.as_ref().map(|m| m.size);
let ref_count = if auth_user.role == "admin" {
metadata.map(|m| m.ref_count)
} else {
None // Never expose global ref_count to regular users
};
let response = HashCheckResponse {
exists: true,
hash,
existing_size: size,
ref_count: None, // Never expose global ref_count
ref_count,
};
Response::builder()
.status(StatusCode::OK)
@@ -492,6 +502,7 @@ impl DedupHandler {
.unwrap()
.into_response()
}
}
// ── Route handlers (free functions) ──────────────────────────────────────────
@@ -511,7 +522,7 @@ impl DedupHandler {
("hash" = String, Path, description = "BLAKE3 hash (64 hex characters)"),
),
responses(
(status = 200, description = "Hash check result (user-scoped)", body = HashCheckResponse),
(status = 200, description = "Hash check result. `ref_count` is only present for admin users.", body = HashCheckResponse),
(status = 400, description = "Invalid hash format"),
),
tag = "dedup",
@@ -783,6 +783,7 @@ impl FileHandler {
file_id,
blob_hash_owned,
original_bytes,
dedup_service.clone(),
);
}
Err(err) => {