Merge pull request #330 from EdouardVanbelle/fix/thumbnail-on-update
bugfix/thumbnails on update
This commit is contained in:
@@ -0,0 +1,35 @@
|
||||
use std::future::Future;
|
||||
use std::pin::Pin;
|
||||
|
||||
/// Observer notified by [`DedupService`] when a genuinely new blob is stored
|
||||
/// for the first time (no dedup hit).
|
||||
///
|
||||
/// Register with [`DedupService::add_blob_creation_hook`] during DI wiring.
|
||||
pub trait BlobCreationHook: Send + Sync {
|
||||
/// Called after the new blob's chunks and manifest have been written.
|
||||
/// `blob_hash` is the BLAKE3 hex, `content_type` is the MIME type if known.
|
||||
/// Must be best-effort — must not propagate errors.
|
||||
fn on_blob_created<'a>(
|
||||
&'a self,
|
||||
blob_hash: &'a str,
|
||||
content_type: Option<&'a str>,
|
||||
) -> Pin<Box<dyn Future<Output = ()> + Send + 'a>>;
|
||||
}
|
||||
/// Observer notified by [`DedupService`] when a blob's ref_count reaches zero
|
||||
/// and it is permanently removed from storage.
|
||||
///
|
||||
/// Implement this trait on any service that needs to react to blob deletion
|
||||
/// (e.g. thumbnail cleanup, CDN invalidation, audit logging). Register with
|
||||
/// [`DedupService::add_blob_hook`] during DI wiring.
|
||||
///
|
||||
/// The boxed-future return keeps the trait dyn-compatible so multiple
|
||||
/// implementations can be stored as `Vec<Arc<dyn BlobDeletionHook>>`.
|
||||
pub trait BlobDeletionHook: Send + Sync {
|
||||
/// Called after the blob file has been removed from disk.
|
||||
/// `blob_hash` is the BLAKE3 hex string identifying the blob.
|
||||
/// Must be best-effort — must not propagate errors.
|
||||
fn on_blob_deleted<'a>(
|
||||
&'a self,
|
||||
blob_hash: &'a str,
|
||||
) -> Pin<Box<dyn Future<Output = ()> + Send + 'a>>;
|
||||
}
|
||||
@@ -0,0 +1,57 @@
|
||||
use std::future::Future;
|
||||
use std::pin::Pin;
|
||||
|
||||
/// Observer notified by [`FileUploadService`] when a new file record is created
|
||||
/// (including dedup hits where the blob already exists).
|
||||
///
|
||||
/// Register with [`FileUploadService::with_file_created_hook`] during DI wiring.
|
||||
pub trait FileCreatedHook: Send + Sync {
|
||||
/// Called after the file record has been persisted.
|
||||
/// `file_id` — opaque file UUID string.
|
||||
/// `blob_hash` — BLAKE3 hex of the blob (may already exist on disk for dedup hits).
|
||||
/// `content_type` — MIME type of the content.
|
||||
/// Must be best-effort — must not propagate errors.
|
||||
fn on_file_created<'a>(
|
||||
&'a self,
|
||||
file_id: &'a str,
|
||||
blob_hash: &'a str,
|
||||
content_type: &'a str,
|
||||
) -> Pin<Box<dyn Future<Output = ()> + Send + 'a>>;
|
||||
}
|
||||
|
||||
/// Observer notified by [`FileUploadService`] when an existing file's blob is
|
||||
/// replaced (WebDAV PUT overwrite, WOPI PutFile, Nextcloud chunked upload).
|
||||
///
|
||||
/// Implement this trait on any service that needs to react to a content swap
|
||||
/// (e.g. thumbnail invalidation + regeneration, search index update).
|
||||
/// Register with [`FileUploadService::with_file_updated_hook`] during DI wiring.
|
||||
///
|
||||
/// The boxed-future return keeps the trait dyn-compatible so multiple
|
||||
/// implementations can be stored as `Vec<Arc<dyn FileUpdatedHook>>`.
|
||||
pub trait FileUpdatedHook: Send + Sync {
|
||||
/// Called after the new blob has been stored and the file record updated.
|
||||
///
|
||||
/// `file_id` is an opaque file UUID string, `blob_hash` is the BLAKE3 hex
|
||||
/// of the new blob, `content_type` is the MIME type of the new content.
|
||||
/// Must be best-effort — must not propagate errors.
|
||||
fn on_file_updated<'a>(
|
||||
&'a self,
|
||||
file_id: &'a str,
|
||||
blob_hash: &'a str,
|
||||
content_type: &'a str,
|
||||
) -> Pin<Box<dyn Future<Output = ()> + Send + 'a>>;
|
||||
}
|
||||
|
||||
/// Observer notified by [`FileManagementService`] when a file is permanently
|
||||
/// deleted (either directly or after being emptied from trash).
|
||||
///
|
||||
/// Register with [`FileManagementService::with_file_deleted_hook`] during DI wiring.
|
||||
pub trait FileDeletedHook: Send + Sync {
|
||||
/// Called after the file record has been removed.
|
||||
/// `file_id` — opaque file UUID string.
|
||||
/// Must be best-effort — must not propagate errors.
|
||||
fn on_file_deleted<'a>(
|
||||
&'a self,
|
||||
file_id: &'a str,
|
||||
) -> Pin<Box<dyn Future<Output = ()> + Send + 'a>>;
|
||||
}
|
||||
@@ -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;
|
||||
@@ -7,6 +8,7 @@ pub mod chunked_upload_ports;
|
||||
pub mod compression_ports;
|
||||
pub mod dedup_ports;
|
||||
pub mod favorites_ports;
|
||||
pub mod file_lifecycle;
|
||||
pub mod file_ports;
|
||||
pub mod inbound;
|
||||
pub mod music_ports;
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
use std::sync::Arc;
|
||||
|
||||
use crate::application::dtos::file_dto::FileDto;
|
||||
use crate::application::ports::file_lifecycle::FileDeletedHook;
|
||||
use crate::application::ports::file_ports::FileManagementUseCase;
|
||||
use crate::application::ports::storage_ports::{CopyFolderTreeResult, FileReadPort, FileWritePort};
|
||||
use crate::application::ports::trash_ports::TrashUseCase;
|
||||
@@ -11,7 +12,6 @@ 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::services::file_content_cache::FileContentCache;
|
||||
use crate::infrastructure::services::thumbnail_service::ThumbnailService;
|
||||
use tracing::{error, info, warn};
|
||||
use uuid::Uuid;
|
||||
|
||||
@@ -26,8 +26,9 @@ pub struct FileManagementService {
|
||||
file_read: Option<Arc<FileBlobReadRepository>>,
|
||||
folder_repo: Option<Arc<FolderDbRepository>>,
|
||||
trash_service: Option<Arc<TrashService>>,
|
||||
thumbnail_service: Option<Arc<ThumbnailService>>,
|
||||
content_cache: Option<Arc<FileContentCache>>,
|
||||
/// Hooks fired after a file is permanently deleted.
|
||||
file_deleted_hooks: Vec<Arc<dyn FileDeletedHook>>,
|
||||
}
|
||||
|
||||
impl FileManagementService {
|
||||
@@ -38,8 +39,8 @@ impl FileManagementService {
|
||||
file_read: None,
|
||||
folder_repo: None,
|
||||
trash_service: None,
|
||||
thumbnail_service: None,
|
||||
content_cache: None,
|
||||
file_deleted_hooks: Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -49,7 +50,6 @@ impl FileManagementService {
|
||||
trash_service: Option<Arc<TrashService>>,
|
||||
file_read: Option<Arc<FileBlobReadRepository>>,
|
||||
folder_repo: Option<Arc<FolderDbRepository>>,
|
||||
thumbnail_service: Option<Arc<ThumbnailService>>,
|
||||
content_cache: Option<Arc<FileContentCache>>,
|
||||
) -> Self {
|
||||
Self {
|
||||
@@ -57,11 +57,17 @@ impl FileManagementService {
|
||||
file_read,
|
||||
folder_repo,
|
||||
trash_service,
|
||||
thumbnail_service,
|
||||
content_cache,
|
||||
file_deleted_hooks: Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Registers a hook to fire after a file is permanently deleted.
|
||||
pub fn with_file_deleted_hook(mut self, hook: Arc<dyn FileDeletedHook>) -> Self {
|
||||
self.file_deleted_hooks.push(hook);
|
||||
self
|
||||
}
|
||||
|
||||
/// Verifies ownership via the read repository.
|
||||
async fn verify_owner(&self, file_id: &str, caller_id: Uuid) -> Result<(), DomainError> {
|
||||
if let Some(read) = &self.file_read {
|
||||
@@ -230,15 +236,11 @@ impl FileManagementUseCase for FileManagementService {
|
||||
|
||||
async fn delete_file(&self, id: &str) -> Result<(), DomainError> {
|
||||
self.file_repository.delete_file(id).await?;
|
||||
// Invalidate content cache — file no longer exists.
|
||||
if let Some(cc) = &self.content_cache {
|
||||
cc.invalidate(id).await;
|
||||
}
|
||||
// Best-effort thumbnail cleanup
|
||||
if let Some(thumb) = &self.thumbnail_service
|
||||
&& let Err(e) = thumb.delete_thumbnails(id).await
|
||||
{
|
||||
warn!("Failed to delete thumbnails for file {}: {}", id, e);
|
||||
for hook in &self.file_deleted_hooks {
|
||||
hook.on_file_deleted(id).await;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
@@ -283,15 +285,11 @@ impl FileManagementUseCase for FileManagementService {
|
||||
// Step 2: Permanent delete — trigger handles blob ref_count
|
||||
warn!("Permanently deleting file: {}", id);
|
||||
self.file_repository.delete_file(id).await?;
|
||||
// Invalidate content cache — file permanently removed.
|
||||
if let Some(cc) = &self.content_cache {
|
||||
cc.invalidate(id).await;
|
||||
}
|
||||
// Best-effort thumbnail cleanup
|
||||
if let Some(thumb) = &self.thumbnail_service
|
||||
&& let Err(e) = thumb.delete_thumbnails(id).await
|
||||
{
|
||||
warn!("Failed to delete thumbnails for file {}: {}", id, e);
|
||||
for hook in &self.file_deleted_hooks {
|
||||
hook.on_file_deleted(id).await;
|
||||
}
|
||||
info!("File permanently deleted: {}", id);
|
||||
|
||||
|
||||
@@ -2,6 +2,7 @@ use std::path::Path;
|
||||
use std::sync::Arc;
|
||||
|
||||
use crate::application::dtos::file_dto::FileDto;
|
||||
use crate::application::ports::file_lifecycle::{FileCreatedHook, FileUpdatedHook};
|
||||
use crate::application::ports::file_ports::FileUploadUseCase;
|
||||
use crate::application::ports::storage_ports::{FileReadPort, FileWritePort};
|
||||
use crate::application::services::storage_usage_service::StorageUsageService;
|
||||
@@ -52,6 +53,10 @@ pub struct FileUploadService {
|
||||
storage_usage_service: Option<Arc<StorageUsageService>>,
|
||||
/// Content cache — invalidated on file update so stale content is never served.
|
||||
content_cache: Option<Arc<FileContentCache>>,
|
||||
/// Hooks fired after a new file record is created.
|
||||
file_created_hooks: Vec<Arc<dyn FileCreatedHook>>,
|
||||
/// Hooks fired after a file's blob is replaced (e.g. thumbnail refresh).
|
||||
file_updated_hooks: Vec<Arc<dyn FileUpdatedHook>>,
|
||||
}
|
||||
|
||||
impl FileUploadService {
|
||||
@@ -62,6 +67,8 @@ impl FileUploadService {
|
||||
file_read: None,
|
||||
storage_usage_service: None,
|
||||
content_cache: None,
|
||||
file_created_hooks: Vec::new(),
|
||||
file_updated_hooks: Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -75,6 +82,8 @@ impl FileUploadService {
|
||||
file_read: Some(file_read),
|
||||
storage_usage_service: None,
|
||||
content_cache: None,
|
||||
file_created_hooks: Vec::new(),
|
||||
file_updated_hooks: Vec::new(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -84,6 +93,18 @@ impl FileUploadService {
|
||||
self
|
||||
}
|
||||
|
||||
/// Registers a hook to fire after a new file record is created.
|
||||
pub fn with_file_created_hook(mut self, hook: Arc<dyn FileCreatedHook>) -> Self {
|
||||
self.file_created_hooks.push(hook);
|
||||
self
|
||||
}
|
||||
|
||||
/// Registers a hook to fire after a file's blob is replaced.
|
||||
pub fn with_file_updated_hook(mut self, hook: Arc<dyn FileUpdatedHook>) -> Self {
|
||||
self.file_updated_hooks.push(hook);
|
||||
self
|
||||
}
|
||||
|
||||
/// Configures the storage usage service
|
||||
pub fn with_storage_usage_service(
|
||||
mut self,
|
||||
@@ -149,6 +170,10 @@ impl FileUploadUseCase for FileUploadService {
|
||||
name, size, dto.id
|
||||
);
|
||||
self.maybe_update_storage_usage(&dto);
|
||||
for hook in &self.file_created_hooks {
|
||||
hook.on_file_created(&dto.id, &dto.etag, &dto.mime_type)
|
||||
.await;
|
||||
}
|
||||
Ok(dto)
|
||||
}
|
||||
|
||||
@@ -301,7 +326,12 @@ impl FileUploadUseCase for FileUploadService {
|
||||
}
|
||||
// Re-read to get fresh DTO with updated etag and timestamps.
|
||||
let updated = file_read.get_file(&file_id).await?;
|
||||
return Ok(FileDto::from(updated));
|
||||
let dto = FileDto::from(updated);
|
||||
for hook in &self.file_updated_hooks {
|
||||
hook.on_file_updated(&file_id, &dto.etag, content_type)
|
||||
.await;
|
||||
}
|
||||
return Ok(dto);
|
||||
}
|
||||
|
||||
// File doesn't exist — create it via streaming upload
|
||||
|
||||
@@ -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>>,
|
||||
|
||||
@@ -56,12 +61,14 @@ pub struct TrashService {
|
||||
}
|
||||
|
||||
impl TrashService {
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub fn new(
|
||||
trash_repository: Arc<TrashDbRepository>,
|
||||
file_read_port: Arc<FileBlobReadRepository>,
|
||||
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 +77,7 @@ impl TrashService {
|
||||
file_read_port,
|
||||
file_write_port,
|
||||
folder_storage_port,
|
||||
dedup_service,
|
||||
thumbnail_service,
|
||||
content_cache,
|
||||
retention_days,
|
||||
@@ -713,18 +721,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 {
|
||||
|
||||
+21
-11
@@ -62,7 +62,7 @@ use crate::infrastructure::services::image_transcode_service::ImageTranscodeServ
|
||||
use crate::infrastructure::services::jwt_service::JwtTokenService;
|
||||
use crate::infrastructure::services::password_hasher::Argon2PasswordHasher;
|
||||
use crate::infrastructure::services::path_resolver_service::PathResolverService;
|
||||
use crate::infrastructure::services::thumbnail_service::ThumbnailService;
|
||||
use crate::infrastructure::services::thumbnail_service::{ThumbnailRefreshHook, ThumbnailService};
|
||||
use crate::infrastructure::services::wopi_discovery_service::WopiDiscoveryService;
|
||||
use crate::infrastructure::services::zip_service::ZipService;
|
||||
|
||||
@@ -264,7 +264,8 @@ impl AppServiceFactory {
|
||||
blob_backend,
|
||||
db_pool.clone(),
|
||||
maintenance_pool.clone(),
|
||||
),
|
||||
)
|
||||
.add_blob_hook(thumbnail_service.clone()),
|
||||
);
|
||||
dedup_service.initialize().await?;
|
||||
|
||||
@@ -355,12 +356,18 @@ impl AppServiceFactory {
|
||||
|
||||
// Refactored services with all infrastructure ports
|
||||
// In blob model, dedup is handled by the repository — no separate write-behind needed
|
||||
let thumbnail_refresh_hook = Arc::new(ThumbnailRefreshHook::new(
|
||||
core.thumbnail_service.clone(),
|
||||
core.dedup_service.clone(),
|
||||
));
|
||||
let file_upload_service = Arc::new(
|
||||
FileUploadService::new_with_read(
|
||||
repos.file_write_repository.clone(),
|
||||
repos.file_read_repository.clone(),
|
||||
)
|
||||
.with_content_cache(core.file_content_cache.clone()),
|
||||
.with_content_cache(core.file_content_cache.clone())
|
||||
.with_file_created_hook(thumbnail_refresh_hook.clone())
|
||||
.with_file_updated_hook(thumbnail_refresh_hook),
|
||||
);
|
||||
|
||||
let file_retrieval_service = Arc::new(FileRetrievalService::new_with_cache(
|
||||
@@ -370,14 +377,16 @@ impl AppServiceFactory {
|
||||
));
|
||||
|
||||
// FileManagementService — ref_count handled by PG trigger, no dedup port needed
|
||||
let file_management_service = Arc::new(FileManagementService::with_trash(
|
||||
repos.file_write_repository.clone(),
|
||||
trash_service.clone(),
|
||||
Some(repos.file_read_repository.clone()),
|
||||
Some(repos.folder_repository.clone()),
|
||||
Some(core.thumbnail_service.clone()),
|
||||
Some(core.file_content_cache.clone()),
|
||||
));
|
||||
let file_management_service = Arc::new(
|
||||
FileManagementService::with_trash(
|
||||
repos.file_write_repository.clone(),
|
||||
trash_service.clone(),
|
||||
Some(repos.file_read_repository.clone()),
|
||||
Some(repos.folder_repository.clone()),
|
||||
Some(core.file_content_cache.clone()),
|
||||
)
|
||||
.with_file_deleted_hook(core.thumbnail_service.clone()),
|
||||
);
|
||||
|
||||
let file_use_case_factory = Arc::new(AppFileUseCaseFactory::new(
|
||||
repos.file_read_repository.clone(),
|
||||
@@ -451,6 +460,7 @@ 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()),
|
||||
));
|
||||
|
||||
@@ -604,8 +604,26 @@ 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(
|
||||
|
||||
@@ -44,6 +44,7 @@ use std::sync::Arc;
|
||||
use tokio::fs;
|
||||
use tokio::io::{AsyncReadExt, AsyncSeekExt};
|
||||
|
||||
use crate::application::ports::blob_lifecycle::{BlobCreationHook, BlobDeletionHook};
|
||||
use crate::application::ports::blob_storage_ports::BlobStorageBackend;
|
||||
use crate::application::ports::dedup_ports::{
|
||||
BlobMetadataDto, DedupPort, DedupResultDto, DedupStatsDto,
|
||||
@@ -83,6 +84,10 @@ pub struct DedupService {
|
||||
/// Isolated maintenance pool for long-running operations
|
||||
/// (verify_integrity, garbage_collect) that must never starve the primary.
|
||||
maintenance_pool: Arc<PgPool>,
|
||||
/// Hooks notified when a genuinely new blob is stored (no dedup hit).
|
||||
blob_creation_hooks: Vec<Arc<dyn BlobCreationHook>>,
|
||||
/// Hooks notified when a blob's ref_count reaches zero and it is deleted.
|
||||
blob_hooks: Vec<Arc<dyn BlobDeletionHook>>,
|
||||
}
|
||||
|
||||
impl DedupService {
|
||||
@@ -100,6 +105,36 @@ impl DedupService {
|
||||
backend,
|
||||
pool,
|
||||
maintenance_pool,
|
||||
blob_creation_hooks: vec![],
|
||||
blob_hooks: vec![],
|
||||
}
|
||||
}
|
||||
|
||||
/// Register a [`BlobCreationHook`] to be called whenever a genuinely new
|
||||
/// blob is stored. Hooks are called in registration order.
|
||||
pub fn add_blob_creation_hook(mut self, hook: Arc<dyn BlobCreationHook>) -> Self {
|
||||
self.blob_creation_hooks.push(hook);
|
||||
self
|
||||
}
|
||||
|
||||
/// 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 creation hooks for a new blob.
|
||||
async fn fire_blob_creation_hooks(&self, hash: &str, content_type: Option<&str>) {
|
||||
for hook in &self.blob_creation_hooks {
|
||||
hook.on_blob_created(hash, content_type).await;
|
||||
}
|
||||
}
|
||||
|
||||
/// 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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -117,6 +152,8 @@ impl DedupService {
|
||||
backend: Arc::new(LocalBlobBackend::new(Path::new("/tmp/oxicloud_stub_blobs"))),
|
||||
pool: stub_pool.clone(),
|
||||
maintenance_pool: stub_pool,
|
||||
blob_creation_hooks: vec![],
|
||||
blob_hooks: vec![],
|
||||
}
|
||||
}
|
||||
|
||||
@@ -334,6 +371,9 @@ impl DedupService {
|
||||
chunk_hashes.len()
|
||||
);
|
||||
|
||||
self.fire_blob_creation_hooks(&file_hash, content_type.as_deref())
|
||||
.await;
|
||||
|
||||
Ok(DedupResultDto::NewBlob {
|
||||
hash: file_hash,
|
||||
size: file_size,
|
||||
@@ -587,7 +627,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)
|
||||
@@ -772,6 +812,9 @@ impl DedupService {
|
||||
}
|
||||
}
|
||||
|
||||
// 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)",
|
||||
&file_hash[..12],
|
||||
@@ -849,6 +892,9 @@ impl DedupService {
|
||||
tracing::warn!("Failed to delete blob file {}: {}", hash, e);
|
||||
}
|
||||
|
||||
// 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)
|
||||
} else {
|
||||
@@ -874,6 +920,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.
|
||||
@@ -1324,16 +1455,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;
|
||||
|
||||
@@ -26,6 +26,7 @@ use crate::application::ports::thumbnail_ports::{
|
||||
ThumbnailPort, ThumbnailSize as PortThumbnailSize, ThumbnailStatsDto,
|
||||
};
|
||||
use crate::domain::errors::{DomainError, ErrorKind};
|
||||
use crate::infrastructure::services::dedup_service::DedupService;
|
||||
|
||||
/// Thumbnail sizes supported by the system
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
|
||||
@@ -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,122 @@ 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 })
|
||||
}
|
||||
}
|
||||
|
||||
// ─── FileUpdatedHook ─────────────────────────────────────────────────────────
|
||||
|
||||
/// Wires thumbnail invalidation + regeneration into the file-update lifecycle.
|
||||
///
|
||||
/// Registered on [`FileUploadService`] during DI. Fires whenever a file's blob
|
||||
/// is replaced (WebDAV PUT overwrite, WOPI PutFile, Nextcloud chunked upload).
|
||||
pub struct ThumbnailRefreshHook {
|
||||
thumbnail: Arc<ThumbnailService>,
|
||||
dedup: Arc<DedupService>,
|
||||
}
|
||||
|
||||
impl ThumbnailRefreshHook {
|
||||
pub fn new(thumbnail: Arc<ThumbnailService>, dedup: Arc<DedupService>) -> Self {
|
||||
Self { thumbnail, dedup }
|
||||
}
|
||||
}
|
||||
|
||||
impl crate::application::ports::file_lifecycle::FileUpdatedHook for ThumbnailRefreshHook {
|
||||
fn on_file_updated<'a>(
|
||||
&'a self,
|
||||
file_id: &'a str,
|
||||
blob_hash: &'a str,
|
||||
content_type: &'a str,
|
||||
) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send + 'a>> {
|
||||
Box::pin(async move {
|
||||
if !ThumbnailService::is_supported_image(content_type) {
|
||||
return;
|
||||
}
|
||||
if let Err(e) = self.thumbnail.delete_thumbnails(file_id).await {
|
||||
tracing::warn!(
|
||||
"Failed to invalidate thumbnail cache for {}: {}",
|
||||
file_id,
|
||||
e
|
||||
);
|
||||
}
|
||||
Self::spawn_thumbnail_generation(
|
||||
self.thumbnail.clone(),
|
||||
self.dedup.clone(),
|
||||
file_id.to_string(),
|
||||
blob_hash.to_string(),
|
||||
);
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl crate::application::ports::file_lifecycle::FileCreatedHook for ThumbnailRefreshHook {
|
||||
fn on_file_created<'a>(
|
||||
&'a self,
|
||||
file_id: &'a str,
|
||||
blob_hash: &'a str,
|
||||
content_type: &'a str,
|
||||
) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send + 'a>> {
|
||||
Box::pin(async move {
|
||||
if !ThumbnailService::is_supported_image(content_type) {
|
||||
return;
|
||||
}
|
||||
Self::spawn_thumbnail_generation(
|
||||
self.thumbnail.clone(),
|
||||
self.dedup.clone(),
|
||||
file_id.to_string(),
|
||||
blob_hash.to_string(),
|
||||
);
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl ThumbnailRefreshHook {
|
||||
fn spawn_thumbnail_generation(
|
||||
ts: Arc<ThumbnailService>,
|
||||
ds: Arc<DedupService>,
|
||||
file_id: String,
|
||||
hash: String,
|
||||
) {
|
||||
tokio::spawn(async move {
|
||||
match ds.read_blob_bytes(&hash).await {
|
||||
Ok(bytes) => {
|
||||
ts.generate_all_sizes_background_from_bytes(file_id, hash, bytes, ds.clone());
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
"Failed to read blob for thumbnail generation {}: {}",
|
||||
file_id,
|
||||
e
|
||||
);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
// ─── FileDeletedHook ─────────────────────────────────────────────────────────
|
||||
|
||||
impl crate::application::ports::file_lifecycle::FileDeletedHook for ThumbnailService {
|
||||
fn on_file_deleted<'a>(
|
||||
&'a self,
|
||||
file_id: &'a str,
|
||||
) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send + 'a>> {
|
||||
Box::pin(async move {
|
||||
if let Err(e) = self.delete_thumbnails(file_id).await {
|
||||
tracing::warn!("Failed to delete thumbnails for file {}: {}", file_id, e);
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// ─── Port implementation ─────────────────────────────────────────────────────
|
||||
|
||||
/// Convert port ThumbnailSize to infra ThumbnailSize.
|
||||
|
||||
@@ -21,12 +21,14 @@ type GlobalState = Arc<AppState>;
|
||||
pub struct HashCheckResponse {
|
||||
/// Whether a blob with this hash already exists
|
||||
pub exists: bool,
|
||||
/// The SHA-256 hash that was checked
|
||||
/// The BLAKE3 hash that was checked
|
||||
pub hash: String,
|
||||
/// 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>,
|
||||
}
|
||||
@@ -36,7 +38,7 @@ pub struct HashCheckResponse {
|
||||
pub struct DedupUploadResponse {
|
||||
/// Whether this was a new file or an existing one
|
||||
pub is_new: bool,
|
||||
/// The SHA-256 hash of the content
|
||||
/// The BLAKE3 hash of the content
|
||||
pub hash: String,
|
||||
/// The size of the content in bytes
|
||||
pub size: u64,
|
||||
@@ -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(
|
||||
@@ -95,13 +98,13 @@ impl DedupHandler {
|
||||
) -> impl IntoResponse {
|
||||
let dedup = &state.core.dedup_service;
|
||||
|
||||
// Validate hash format (SHA-256 = 64 hex chars)
|
||||
// Validate hash format (BLAKE3 = 64 hex chars)
|
||||
if hash.len() != 64 || !hash.chars().all(|c| c.is_ascii_hexdigit()) {
|
||||
return Response::builder()
|
||||
.status(StatusCode::BAD_REQUEST)
|
||||
.header(header::CONTENT_TYPE, "application/json")
|
||||
.body(Body::from(
|
||||
r#"{"error": "Invalid hash format. Expected SHA-256 (64 hex characters)"}"#,
|
||||
r#"{"error": "Invalid hash format. Expected BLAKE3 (64 hex characters)"}"#,
|
||||
))
|
||||
.unwrap()
|
||||
.into_response();
|
||||
@@ -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)
|
||||
@@ -508,10 +518,10 @@ impl DedupHandler {
|
||||
get,
|
||||
path = "/api/dedup/check/{hash}",
|
||||
params(
|
||||
("hash" = String, Path, description = "SHA-256 hash (64 hex characters)"),
|
||||
("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",
|
||||
@@ -564,7 +574,7 @@ pub async fn get_stats(state: State<GlobalState>, auth_user: AuthUser) -> impl I
|
||||
get,
|
||||
path = "/api/dedup/blob/{hash}",
|
||||
params(
|
||||
("hash" = String, Path, description = "SHA-256 hash of the blob (64 hex characters)"),
|
||||
("hash" = String, Path, description = "BLAKE3 hash of the blob (64 hex characters)"),
|
||||
),
|
||||
responses(
|
||||
(status = 200, description = "Raw blob content (user-scoped)"),
|
||||
|
||||
@@ -783,6 +783,7 @@ impl FileHandler {
|
||||
file_id,
|
||||
blob_hash_owned,
|
||||
original_bytes,
|
||||
dedup_service.clone(),
|
||||
);
|
||||
}
|
||||
Err(err) => {
|
||||
|
||||
@@ -288,7 +288,7 @@ async fn put_file(
|
||||
let _ = tokio::fs::remove_file(&temp_path).await;
|
||||
|
||||
match result {
|
||||
Ok(_) => StatusCode::OK.into_response(),
|
||||
Ok(_file_dto) => StatusCode::OK.into_response(),
|
||||
Err(e) => {
|
||||
tracing::error!("WOPI PutFile failed: {}", e);
|
||||
StatusCode::INTERNAL_SERVER_ERROR.into_response()
|
||||
|
||||
@@ -163,6 +163,7 @@ async fn handle_assemble(
|
||||
)
|
||||
.await
|
||||
.map_err(|e| AppError::internal_error(format!("Failed to update file: {}", e)))?;
|
||||
|
||||
Some(dto.etag)
|
||||
} else {
|
||||
// For new files we still need to read the temp file since create_file takes &[u8].
|
||||
|
||||
Reference in New Issue
Block a user