feat(recent): update recent list server side
initially the recent was done client side
recent files are now directly updated on serverside when accessing a file
note: nextcloud and webdav voluntary not included
This commit is contained in:
@@ -4,6 +4,7 @@ use crate::application::dtos::file_dto::FileDto;
|
||||
use crate::application::ports::authorization_ports::AuthorizationEngine;
|
||||
use crate::application::ports::file_lifecycle::FileLifecycleHook;
|
||||
use crate::application::ports::file_ports::FileManagementUseCase;
|
||||
use crate::application::ports::resource_access_hook::ResourceAccessHook;
|
||||
use crate::application::ports::storage_ports::{CopyFolderTreeResult, FileWritePort};
|
||||
use crate::application::ports::trash_ports::TrashUseCase;
|
||||
use crate::application::services::trash_service::TrashService;
|
||||
@@ -31,6 +32,11 @@ pub struct FileManagementService {
|
||||
authz: Arc<PgAclEngine>,
|
||||
/// Lifecycle hook dispatcher — fired on file created (copy) and deleted.
|
||||
file_lifecycle_hook: Option<Arc<dyn FileLifecycleHook>>,
|
||||
/// Read/write access hook — fired so Recent reflects "this is the file
|
||||
/// I just copied / renamed / moved", same way the read paths surface
|
||||
/// downloads. Distinct from the lifecycle hook because lifecycle hooks
|
||||
/// don't carry the `caller_id` the recording side needs.
|
||||
resource_access_hook: Option<Arc<dyn ResourceAccessHook>>,
|
||||
}
|
||||
|
||||
impl FileManagementService {
|
||||
@@ -53,6 +59,7 @@ impl FileManagementService {
|
||||
content_cache,
|
||||
authz,
|
||||
file_lifecycle_hook: None,
|
||||
resource_access_hook: None,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -62,6 +69,19 @@ impl FileManagementService {
|
||||
self
|
||||
}
|
||||
|
||||
/// Registers the read/write access hook (Recent list recorder).
|
||||
pub fn with_resource_access_hook(mut self, hook: Arc<dyn ResourceAccessHook>) -> Self {
|
||||
self.resource_access_hook = Some(hook);
|
||||
self
|
||||
}
|
||||
|
||||
/// Internal helper: fire the access hook if registered.
|
||||
fn notify_file_accessed(&self, caller_id: Uuid, file_id: &str) {
|
||||
if let Some(hook) = &self.resource_access_hook {
|
||||
hook.on_file_accessed(caller_id, file_id);
|
||||
}
|
||||
}
|
||||
|
||||
/// Engine check for a file resource. Parses the id into a `Uuid` and
|
||||
/// requires the specified permission.
|
||||
async fn require_file_perm(
|
||||
@@ -156,6 +176,9 @@ impl FileManagementService {
|
||||
if let Some(hook) = &self.file_lifecycle_hook {
|
||||
hook.on_file_copied(&dto.id, &dto.content_hash, &dto.mime_type, file_id);
|
||||
}
|
||||
// The caller just spawned a fresh file — show it in their Recent
|
||||
// list. The source file isn't recorded; only the visible target.
|
||||
self.notify_file_accessed(caller_id, &dto.id);
|
||||
Ok(dto)
|
||||
}
|
||||
|
||||
|
||||
@@ -6,6 +6,7 @@ use std::sync::Arc;
|
||||
use crate::application::dtos::file_dto::FileDto;
|
||||
use crate::application::ports::authorization_ports::AuthorizationEngine;
|
||||
use crate::application::ports::file_ports::{FileRetrievalUseCase, OptimizedFileContent};
|
||||
use crate::application::ports::resource_access_hook::ResourceAccessHook;
|
||||
use crate::application::ports::storage_ports::FileReadPort;
|
||||
use crate::common::errors::DomainError;
|
||||
use crate::domain::services::authorization::{Permission, Resource, Subject};
|
||||
@@ -33,6 +34,11 @@ pub struct FileRetrievalService {
|
||||
content_cache: Option<Arc<FileContentCache>>,
|
||||
transcode: Option<Arc<ImageTranscodeService>>,
|
||||
authz: Option<Arc<PgAclEngine>>,
|
||||
/// Optional read-event observer. Currently fans out to the Recent-list
|
||||
/// recorder; future observers (audit trail, "last seen by", …) attach
|
||||
/// to the same hook so service code only knows the trait, not the impl.
|
||||
/// `None` for the test/stub path that constructs via [`Self::new`].
|
||||
resource_access_hook: Option<Arc<dyn ResourceAccessHook>>,
|
||||
}
|
||||
|
||||
impl FileRetrievalService {
|
||||
@@ -45,6 +51,7 @@ impl FileRetrievalService {
|
||||
content_cache: None,
|
||||
transcode: None,
|
||||
authz: None,
|
||||
resource_access_hook: None,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -61,6 +68,32 @@ impl FileRetrievalService {
|
||||
content_cache: Some(content_cache),
|
||||
transcode: Some(transcode),
|
||||
authz: Some(authz),
|
||||
resource_access_hook: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Builder: attach a [`ResourceAccessHook`] that fires after every
|
||||
/// authorised `_with_perms` read. Without it the service is silent —
|
||||
/// existing behaviour for stub / test paths.
|
||||
pub fn with_resource_access_hook(mut self, hook: Arc<dyn ResourceAccessHook>) -> Self {
|
||||
self.resource_access_hook = Some(hook);
|
||||
self
|
||||
}
|
||||
|
||||
/// Fire the access hook if registered. Called from every `_with_perms`
|
||||
/// read after the authZ + lookup has succeeded (never on failure
|
||||
/// paths — denied reads must not surface in Recent).
|
||||
///
|
||||
/// `pub` because the WebDAV / NextCloud DAV handlers resolve files
|
||||
/// by path and authorise via that resolver, not via the
|
||||
/// `*_with_perms` service methods — they then serve content through
|
||||
/// the no-perms `get_file_stream` / `get_file_range_stream`. Those
|
||||
/// handlers must call this directly after their own authZ has
|
||||
/// passed so cross-protocol downloads (NC desktop, davx5, native
|
||||
/// `/webdav/`) also surface in Recent.
|
||||
pub fn notify_file_accessed(&self, caller_id: Uuid, file_id: &str) {
|
||||
if let Some(hook) = &self.resource_access_hook {
|
||||
hook.on_file_accessed(caller_id, file_id);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -262,6 +295,11 @@ impl FileRetrievalUseCase for FileRetrievalService {
|
||||
async fn get_file_with_perms(&self, id: &str, caller_id: Uuid) -> Result<FileDto, DomainError> {
|
||||
self.require_file(id, Permission::Read, caller_id).await?;
|
||||
let file = self.file_read.get_file(id).await?;
|
||||
// After authZ + lookup succeed: this caller has just inspected the
|
||||
// file. Recent listing observes via the hook. The throttle in the
|
||||
// recording impl coalesces repeat metadata fetches against the same
|
||||
// file (file viewer poll, browse-then-download pattern).
|
||||
self.notify_file_accessed(caller_id, id);
|
||||
Ok(FileDto::from(file))
|
||||
}
|
||||
|
||||
@@ -333,6 +371,7 @@ impl FileRetrievalUseCase for FileRetrievalService {
|
||||
caller_id: Uuid,
|
||||
) -> Result<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>, DomainError> {
|
||||
self.require_file(id, Permission::Read, caller_id).await?;
|
||||
self.notify_file_accessed(caller_id, id);
|
||||
self.file_read.get_file_stream(id).await
|
||||
}
|
||||
|
||||
@@ -359,6 +398,7 @@ impl FileRetrievalUseCase for FileRetrievalService {
|
||||
self.require_file(id, Permission::Read, caller_id).await?;
|
||||
let file = self.file_read.get_file(id).await?;
|
||||
let dto = FileDto::from(file);
|
||||
self.notify_file_accessed(caller_id, id);
|
||||
self.optimized_inner(id, dto, accept_webp, prefer_original)
|
||||
.await
|
||||
}
|
||||
@@ -393,6 +433,10 @@ impl FileRetrievalUseCase for FileRetrievalService {
|
||||
end: Option<u64>,
|
||||
) -> Result<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>, DomainError> {
|
||||
self.require_file(id, Permission::Read, caller_id).await?;
|
||||
// Range requests are bursty (video seeks, NC chunked downloads) —
|
||||
// the recording hook's per-(caller, file) throttle absorbs the
|
||||
// storm so one watched video lands as one Recent row, not 1000.
|
||||
self.notify_file_accessed(caller_id, id);
|
||||
self.file_read.get_file_range_stream(id, start, end).await
|
||||
}
|
||||
|
||||
|
||||
@@ -5,6 +5,7 @@ use crate::application::dtos::file_dto::FileDto;
|
||||
use crate::application::ports::authorization_ports::AuthorizationEngine;
|
||||
use crate::application::ports::file_lifecycle::FileLifecycleHook;
|
||||
use crate::application::ports::file_ports::{FileUploadUseCase, StoredBlob};
|
||||
use crate::application::ports::resource_access_hook::ResourceAccessHook;
|
||||
use crate::application::ports::storage_ports::{FileReadPort, FileWritePort, StorageUsagePort};
|
||||
use crate::application::services::storage_usage_service::StorageUsageService;
|
||||
use crate::common::errors::DomainError;
|
||||
@@ -35,6 +36,12 @@ pub struct FileUploadService {
|
||||
content_cache: Option<Arc<FileContentCache>>,
|
||||
/// Single lifecycle dispatcher — fires on_file_created / on_file_updated.
|
||||
file_lifecycle_hook: Option<Arc<dyn FileLifecycleHook>>,
|
||||
/// Read-event hook — fires "caller just touched this file" so Recent
|
||||
/// records uploads / overwrites alongside reads. Distinct from
|
||||
/// `file_lifecycle_hook` because the lifecycle dispatcher only knows
|
||||
/// `(file_id, blob_hash, content_type)`; the recording side needs the
|
||||
/// `caller_id` the service already has in hand.
|
||||
resource_access_hook: Option<Arc<dyn ResourceAccessHook>>,
|
||||
/// Dependencies of the instant-upload path
|
||||
/// (`create_file_from_owned_blob_with_perms`); `None` in minimal test
|
||||
/// wiring.
|
||||
@@ -58,6 +65,7 @@ impl FileUploadService {
|
||||
storage_usage_service: None,
|
||||
content_cache: None,
|
||||
file_lifecycle_hook: None,
|
||||
resource_access_hook: None,
|
||||
instant_upload: None,
|
||||
}
|
||||
}
|
||||
@@ -73,6 +81,7 @@ impl FileUploadService {
|
||||
storage_usage_service: None,
|
||||
content_cache: None,
|
||||
file_lifecycle_hook: None,
|
||||
resource_access_hook: None,
|
||||
instant_upload: None,
|
||||
}
|
||||
}
|
||||
@@ -105,6 +114,19 @@ impl FileUploadService {
|
||||
self
|
||||
}
|
||||
|
||||
/// Registers the read/write access hook (Recent list recorder).
|
||||
pub fn with_resource_access_hook(mut self, hook: Arc<dyn ResourceAccessHook>) -> Self {
|
||||
self.resource_access_hook = Some(hook);
|
||||
self
|
||||
}
|
||||
|
||||
/// Internal helper: fire the access hook if registered.
|
||||
fn notify_file_accessed(&self, caller_id: Uuid, file_id: &str) {
|
||||
if let Some(hook) = &self.resource_access_hook {
|
||||
hook.on_file_accessed(caller_id, file_id);
|
||||
}
|
||||
}
|
||||
|
||||
/// Configures the storage usage service
|
||||
pub fn with_storage_usage_service(
|
||||
mut self,
|
||||
@@ -282,6 +304,9 @@ impl FileUploadService {
|
||||
if let Some(hook) = &self.file_lifecycle_hook {
|
||||
hook.on_file_updated(file_id, &dto.content_hash, &dto.mime_type);
|
||||
}
|
||||
// Delta-upload commit path — record the swap so Recent reflects
|
||||
// "this is the file I just delta-updated".
|
||||
self.notify_file_accessed(caller_id, file_id);
|
||||
Ok(dto)
|
||||
}
|
||||
|
||||
@@ -372,6 +397,9 @@ impl FileUploadUseCase for FileUploadService {
|
||||
if let Some(hook) = &self.file_lifecycle_hook {
|
||||
hook.on_file_created(&dto.id, &dto.content_hash, &dto.mime_type, blob.is_new_blob);
|
||||
}
|
||||
// The caller just created this file — surface it in Recent so the
|
||||
// "I just uploaded X" UX matches the pre-SvelteKit behaviour.
|
||||
self.notify_file_accessed(caller_id, &dto.id);
|
||||
Ok(dto)
|
||||
}
|
||||
|
||||
@@ -429,6 +457,7 @@ impl FileUploadUseCase for FileUploadService {
|
||||
if let Some(hook) = &self.file_lifecycle_hook {
|
||||
hook.on_file_updated(&file_id, &dto.content_hash, content_type);
|
||||
}
|
||||
self.notify_file_accessed(caller_id, &file_id);
|
||||
return Ok(dto);
|
||||
}
|
||||
|
||||
@@ -474,6 +503,7 @@ impl FileUploadUseCase for FileUploadService {
|
||||
if let Some(hook) = &self.file_lifecycle_hook {
|
||||
hook.on_file_created(&dto.id, &dto.content_hash, content_type, is_new_blob);
|
||||
}
|
||||
self.notify_file_accessed(caller_id, &dto.id);
|
||||
Ok(dto)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,10 +1,11 @@
|
||||
use crate::application::dtos::cursor::PageCursor;
|
||||
use crate::application::dtos::recent_dto::{RecentCursor, RecentItemDto, RecentResourceRow};
|
||||
use crate::application::ports::recent_ports::{RecentItemsRepositoryPort, RecentItemsUseCase};
|
||||
use crate::application::ports::resource_access_hook::ResourceAccessHook;
|
||||
use crate::common::errors::{DomainError, ErrorKind, Result};
|
||||
use crate::domain::services::authorization::ResourceKind;
|
||||
use crate::infrastructure::repositories::pg::RecentItemsPgRepository;
|
||||
use std::sync::Arc;
|
||||
use std::sync::{Arc, OnceLock};
|
||||
use tracing::info;
|
||||
use uuid::Uuid;
|
||||
|
||||
@@ -15,6 +16,14 @@ use uuid::Uuid;
|
||||
pub struct RecentService {
|
||||
repo: Arc<RecentItemsPgRepository>,
|
||||
max_recent_items: i32,
|
||||
/// Set after construction via [`Self::set_resource_access_hook`].
|
||||
/// The hook is built FROM this service (it wraps an `Arc<Self>`), so
|
||||
/// we can't take it as a constructor arg without circular ownership;
|
||||
/// the OnceLock holds the back-edge so this service can notify the
|
||||
/// hook when the user clears or removes Recent rows. The notification
|
||||
/// lets the hook drop its in-memory throttle entries — otherwise a
|
||||
/// freshly-cleared Recent refuses to re-record until the TTL expires.
|
||||
resource_access_hook: OnceLock<Arc<dyn ResourceAccessHook>>,
|
||||
}
|
||||
|
||||
impl RecentService {
|
||||
@@ -23,6 +32,25 @@ impl RecentService {
|
||||
Self {
|
||||
repo,
|
||||
max_recent_items: max_recent_items.clamp(1, 100),
|
||||
resource_access_hook: OnceLock::new(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Wire the access hook in after construction. Idempotent: a second
|
||||
/// `set` is a no-op (returns the existing value as `Err`). Called
|
||||
/// from DI once `RecentRecordingHook::new(Arc<Self>)` has produced
|
||||
/// the back-edge that closes the loop.
|
||||
pub fn set_resource_access_hook(&self, hook: Arc<dyn ResourceAccessHook>) {
|
||||
let _ = self.resource_access_hook.set(hook);
|
||||
}
|
||||
|
||||
/// Internal helper: notify the hook (if registered) that `user_id`
|
||||
/// has emptied their Recent list — wholly or by removing a single
|
||||
/// row. The hook drops its in-memory throttle entries so the very
|
||||
/// next access re-records into the freshly-empty table.
|
||||
fn notify_recents_cleared(&self, user_id: Uuid) {
|
||||
if let Some(hook) = self.resource_access_hook.get() {
|
||||
hook.on_recents_cleared(user_id);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -100,6 +128,12 @@ impl RecentItemsUseCase for RecentService {
|
||||
item_id,
|
||||
user_id
|
||||
);
|
||||
// Drop the throttle entries so the next access re-records. We
|
||||
// notify on every call (even when `removed == false`) so the
|
||||
// semantics are "the user expressed intent to forget this" —
|
||||
// the hook owns the per-(user, item) cache anyway, dropping a
|
||||
// miss is a no-op.
|
||||
self.notify_recents_cleared(user_id);
|
||||
Ok(removed)
|
||||
}
|
||||
|
||||
@@ -108,6 +142,7 @@ impl RecentItemsUseCase for RecentService {
|
||||
info!("Clearing all recent items for user {}", user_id);
|
||||
self.repo.clear_all(user_id).await?;
|
||||
info!("Cleared all recent items for user {}", user_id);
|
||||
self.notify_recents_cleared(user_id);
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user