diff --git a/src/application/ports/calendar_ports.rs b/src/application/ports/calendar_ports.rs index 4f8559c6..56799533 100644 --- a/src/application/ports/calendar_ports.rs +++ b/src/application/ports/calendar_ports.rs @@ -87,6 +87,14 @@ pub trait CalendarStoragePort: Send + Sync + 'static { ) -> Result; async fn delete_event(&self, event_id: &str) -> Result<(), DomainError>; async fn get_event(&self, event_id: &str) -> Result; + /// Indexed single-row lookup by iCalendar UID — the CalDAV + /// object-resource paths must use this instead of listing the whole + /// calendar (every row + its `ical_data`) and filtering client-side. + async fn find_event_by_ical_uid( + &self, + calendar_id: &str, + ical_uid: &str, + ) -> Result, DomainError>; async fn list_events_by_calendar( &self, calendar_id: &str, @@ -180,6 +188,15 @@ pub trait CalendarUseCase: Send + Sync + 'static { event_id: &str, user_id: Uuid, ) -> Result; + /// Resolve one event by its iCalendar UID (the identifier CalDAV + /// object resources are addressed by). `Ok(None)` when no event with + /// that UID exists in the calendar. + async fn get_event_by_ical_uid( + &self, + calendar_id: &str, + ical_uid: &str, + user_id: Uuid, + ) -> Result, DomainError>; async fn list_events( &self, calendar_id: &str, diff --git a/src/application/ports/carddav_ports.rs b/src/application/ports/carddav_ports.rs index 94acb1e5..ee49dcec 100644 --- a/src/application/ports/carddav_ports.rs +++ b/src/application/ports/carddav_ports.rs @@ -71,6 +71,17 @@ pub trait ContactUseCase: Send + Sync + 'static { async fn delete_contact(&self, contact_id: &str, user_id: Uuid) -> Result<(), DomainError>; async fn get_contact(&self, contact_id: &str, user_id: Uuid) -> Result; + /// Resolve one contact by its vCard UID (the identifier CardDAV + /// object resources are addressed by) with an indexed single-row + /// lookup — instead of listing the whole address book (every row + /// with its vCard + JSONB columns) and filtering client-side. + /// `Ok(None)` when no contact with that UID exists in the book. + async fn get_contact_by_uid( + &self, + address_book_id: &str, + uid: &str, + user_id: Uuid, + ) -> Result, DomainError>; async fn list_contacts( &self, address_book_id: &str, diff --git a/src/application/services/auth_application_service.rs b/src/application/services/auth_application_service.rs index 007626b9..36584601 100644 --- a/src/application/services/auth_application_service.rs +++ b/src/application/services/auth_application_service.rs @@ -11,7 +11,7 @@ use crate::common::config::OidcConfig; use crate::common::errors::{DomainError, ErrorKind}; use crate::domain::entities::magic_link_token::{MagicLinkResourceKind, MagicLinkStatus}; use crate::domain::entities::session::Session; -use crate::domain::entities::user::{User, UserRole}; +use crate::domain::entities::user::{User, UserFlags, UserRole}; use crate::domain::repositories::magic_link_token_repository::MagicLinkTokenRepository; use crate::infrastructure::repositories::pg::SessionPgRepository; use crate::infrastructure::repositories::pg::UserPgRepository; @@ -140,8 +140,20 @@ pub struct AuthApplicationService { /// Magic-link token repository — populated when the magic-link feature /// is enabled (PR 8+). `None` means redemption endpoints return 503. magic_link_repo: Option>, + /// Per-user authorization flags (`role` / `is_external` / `active`), + /// consulted by middleware guards on every WebDAV / CalDAV / CardDAV + /// request. The short TTL keeps the "role changes apply without token + /// rotation" property within seconds while removing one DB round-trip + /// per request; the known mutation paths (`change_user_role`, + /// `set_user_active`) also invalidate eagerly. + user_flags_cache: Cache, } +/// TTL for [`AuthApplicationService::user_flags_cache`]. Upper bound on how +/// long a role / external / active change can take to be observed by the +/// per-request guards when it bypasses the eager invalidation paths. +const USER_FLAGS_CACHE_TTL: Duration = Duration::from_secs(30); + impl AuthApplicationService { pub fn new( user_storage: Arc, @@ -170,6 +182,10 @@ impl AuthApplicationService { .time_to_live(Duration::from_secs(60)) .build(), magic_link_repo: None, + user_flags_cache: Cache::builder() + .max_capacity(10_000) + .time_to_live(USER_FLAGS_CACHE_TTL) + .build(), } } @@ -1141,6 +1157,24 @@ impl AuthApplicationService { Ok(UserDto::from(user)) } + /// Cached, image-free lookup of the caller's authorization flags + /// (`role` / `is_external` / `active`). This is the per-request fast + /// path for middleware guards: the full `get_user` row fetch drags the + /// `image` column (a data URI of up to 512 KiB) across the wire, which + /// a sync client issuing hundreds of DAV requests per minute paid on + /// every single one just to read a boolean. + /// + /// Staleness is bounded by [`USER_FLAGS_CACHE_TTL`]; role and active + /// changes made through this service invalidate the entry eagerly. + pub async fn get_user_flags(&self, user_id: Uuid) -> Result { + if let Some(flags) = self.user_flags_cache.get(&user_id) { + return Ok(flags); + } + let flags = self.user_storage.get_user_flags(user_id).await?; + self.user_flags_cache.insert(user_id, flags); + Ok(flags) + } + /// Apply a profile update on behalf of the calling user (PR 24). /// /// Hard rules: @@ -1797,7 +1831,9 @@ impl AuthApplicationService { pub async fn set_user_active(&self, user_id: Uuid, active: bool) -> Result<(), DomainError> { self.user_storage .set_user_active_status(user_id, active) - .await + .await?; + self.user_flags_cache.invalidate(&user_id); + Ok(()) } /// Change user role (admin only) @@ -1809,7 +1845,9 @@ impl AuthApplicationService { format!("Invalid role: {}. Must be 'admin' or 'user'", role), )); } - self.user_storage.change_role(user_id, role).await + self.user_storage.change_role(user_id, role).await?; + self.user_flags_cache.invalidate(&user_id); + Ok(()) } /// Update user's storage quota (admin only) diff --git a/src/application/services/calendar_service.rs b/src/application/services/calendar_service.rs index 7ce70475..6b51cb8a 100644 --- a/src/application/services/calendar_service.rs +++ b/src/application/services/calendar_service.rs @@ -277,6 +277,29 @@ impl CalendarUseCase for CalendarService { Ok(event) } + async fn get_event_by_ical_uid( + &self, + calendar_id: &str, + ical_uid: &str, + user_id: Uuid, + ) -> Result, DomainError> { + let has_access = self + .calendar_storage + .check_calendar_access(calendar_id, user_id) + .await?; + let calendar = self.calendar_storage.get_calendar(calendar_id).await?; + if !has_access && !calendar.is_public { + return Err(DomainError::new( + ErrorKind::AccessDenied, + "Calendar", + "You don't have permission to view events in this calendar", + )); + } + self.calendar_storage + .find_event_by_ical_uid(calendar_id, ical_uid) + .await + } + async fn list_events( &self, calendar_id: &str, diff --git a/src/application/services/contact_service.rs b/src/application/services/contact_service.rs index f95e9406..9d44b306 100644 --- a/src/application/services/contact_service.rs +++ b/src/application/services/contact_service.rs @@ -780,6 +780,22 @@ impl ContactUseCase for ContactService { Ok(ContactDto::from(contact)) } + async fn get_contact_by_uid( + &self, + address_book_id: &str, + uid: &str, + user_id: Uuid, + ) -> Result, DomainError> { + let id = Uuid::parse_str(address_book_id) + .map_err(|_| DomainError::validation_error("Invalid address book ID format"))?; + + // Check if user has access to the address book + self.check_address_book_access(&id, &user_id).await?; + + let contact = self.contact_repository.get_contact_by_uid(&id, uid).await?; + Ok(contact.map(ContactDto::from)) + } + async fn list_contacts( &self, address_book_id: &str, diff --git a/src/application/services/nextcloud_file_id_service.rs b/src/application/services/nextcloud_file_id_service.rs index 30f35614..2190e7f9 100644 --- a/src/application/services/nextcloud_file_id_service.rs +++ b/src/application/services/nextcloud_file_id_service.rs @@ -1,12 +1,25 @@ +use std::collections::HashMap; use std::sync::Arc; +use moka::future::Cache; +use uuid::Uuid; + use crate::common::errors::{DomainError, ErrorKind, Result}; use crate::infrastructure::repositories::pg::NextcloudObjectIdRepository; +/// Capacity of the in-memory UUID→numeric-id cache. The mapping is immutable +/// once created, so a warm entry never goes stale and eviction only costs a +/// re-query; ~100k entries is a few MB. +const ID_CACHE_CAPACITY: u64 = 100_000; + #[derive(Clone)] pub struct NextcloudFileIdService { repo: Option>, instance_id: String, + /// Object-UUID → stable numeric id. `moka` caches are `Arc`-backed, so all + /// clones of the service share one cache and the per-child resolution in a + /// listing costs zero queries once warm. + cache: Cache, } impl NextcloudFileIdService { @@ -14,6 +27,7 @@ impl NextcloudFileIdService { Self { repo: Some(repo), instance_id, + cache: Cache::new(ID_CACHE_CAPACITY), } } @@ -21,29 +35,76 @@ impl NextcloudFileIdService { Self { repo: None, instance_id: "ocnca".to_string(), + cache: Cache::new(ID_CACHE_CAPACITY), } } - pub async fn get_or_create_file_id(&self, file_id: &str) -> Result { - let repo = self.repo.as_ref().ok_or_else(|| { - DomainError::internal_error("NextcloudFileId", "Repository not initialized") - })?; - repo.get_or_create("file", file_id).await + /// Resolve — creating when absent — stable numeric file IDs for many + /// UUIDs at once. Cache hits cost nothing; the misses are resolved with a + /// single backing query. The returned map is keyed by the caller's + /// original id strings; unresolvable inputs are simply absent (mirroring + /// the `.ok()` behaviour the callers relied on). + pub async fn get_or_create_file_ids( + &self, + file_ids: &[String], + ) -> Result> { + self.get_or_create_many("file", file_ids).await } - pub async fn get_or_create_folder_id(&self, folder_id: &str) -> Result { - let repo = self.repo.as_ref().ok_or_else(|| { + /// Folder counterpart of [`Self::get_or_create_file_ids`]. + pub async fn get_or_create_folder_ids( + &self, + folder_ids: &[String], + ) -> Result> { + self.get_or_create_many("folder", folder_ids).await + } + + async fn get_or_create_many( + &self, + object_type: &str, + raw_ids: &[String], + ) -> Result> { + let mut result = HashMap::with_capacity(raw_ids.len()); + // Parsed-UUID → caller's original string; also dedupes the miss list. + let mut pending: HashMap = HashMap::new(); + + for raw in raw_ids { + let Ok(uuid) = Uuid::parse_str(raw) else { + continue; // Unparseable ids never had a mapping — skip silently. + }; + if let Some(id) = self.cache.get(&uuid).await { + result.insert(raw.clone(), id); + } else { + pending.entry(uuid).or_insert_with(|| raw.clone()); + } + } + + if !pending.is_empty() { + let misses: Vec = pending.keys().copied().collect(); + let resolved = self + .repo()? + .get_or_create_many(object_type, &misses) + .await?; + for (uuid, id) in resolved { + self.cache.insert(uuid, id).await; + if let Some(original) = pending.get(&uuid) { + result.insert(original.clone(), id); + } + } + } + + Ok(result) + } + + fn repo(&self) -> Result<&Arc> { + self.repo.as_ref().ok_or_else(|| { DomainError::internal_error("NextcloudFileId", "Repository not initialized") - })?; - repo.get_or_create("folder", folder_id).await + }) } /// Get the OxiCloud file UUID from a Nextcloud numeric ID. pub async fn get_oxicloud_id(&self, nc_file_id: i64) -> Result { - let repo = self.repo.as_ref().ok_or_else(|| { - DomainError::internal_error("NextcloudFileId", "Repository not initialized") - })?; - repo.get_object_id(nc_file_id, "file").await + self.repo()?.get_object_id(nc_file_id, "file").await } pub fn format_oc_id(&self, id: i64) -> String { @@ -59,6 +120,7 @@ impl NextcloudFileIdService { Self { repo: None, instance_id: instance_id.to_string(), + cache: Cache::new(ID_CACHE_CAPACITY), } } @@ -107,4 +169,25 @@ mod tests { let svc = NextcloudFileIdService::new_stub(); assert!(svc.ensure_ready().is_err()); } + + // Empty input resolves to an empty map without ever touching the repo, so + // it succeeds even on the repo-less stub. + #[tokio::test] + async fn test_get_or_create_file_ids_empty_is_noop() { + let svc = NextcloudFileIdService::new_stub(); + let map = svc.get_or_create_file_ids(&[]).await.unwrap(); + assert!(map.is_empty()); + } + + // Unparseable ids never had a mapping, so they are skipped before any repo + // call — the stub (no repo) must not error on them. + #[tokio::test] + async fn test_get_or_create_file_ids_skips_unparseable() { + let svc = NextcloudFileIdService::new_stub(); + let map = svc + .get_or_create_file_ids(&["not-a-uuid".to_string()]) + .await + .unwrap(); + assert!(map.is_empty()); + } } diff --git a/src/application/services/recipient_notification_service.rs b/src/application/services/recipient_notification_service.rs index fce62b27..ca832d2e 100644 --- a/src/application/services/recipient_notification_service.rs +++ b/src/application/services/recipient_notification_service.rs @@ -70,6 +70,11 @@ use crate::domain::services::authorization::{Resource, Subject}; use crate::infrastructure::repositories::pg::UserPgRepository; use crate::interfaces::middleware::rate_limit::RateLimiter; +/// Concurrent per-recipient dispatches in flight during a group fan-out. +/// High enough to collapse a 30-member group's serial SMTP latency, +/// low enough not to flood the relay (most reject >10 parallel sessions). +const NOTIFY_DISPATCH_CONCURRENCY: usize = 6; + /// What triggered the notification — purely an audit discriminator. /// `GrantCreated` → fired implicitly when a grant lands; `ManualResend` /// → granter explicitly clicked "Notify by email" in My Shares. @@ -268,13 +273,22 @@ impl RecipientNotificationService { ); } - let mut outcomes = Vec::with_capacity(members.len()); - for member in &members { - let outcome = self - .dispatch_to_one_user(granter, member, resource, trigger) - .await; - outcomes.push(outcome); - } + // SMTP dispatch dominates each iteration (hundreds of ms per + // recipient) and the iterations are independent — coalescing and + // rate-limiting key on (granter, recipient), which is distinct per + // member. Bounded concurrency keeps a 30-member group grant from + // holding the HTTP response for 15+ s of serial sends while still + // capping the pressure on the SMTP relay. `buffered` (not + // `buffer_unordered`) preserves the member order of the outcomes. + use futures::stream::{self, StreamExt}; + let outcomes: Vec = stream::iter(members) + .map(|member| async move { + self.dispatch_to_one_user(granter, &member, resource, trigger) + .await + }) + .buffered(NOTIFY_DISPATCH_CONCURRENCY) + .collect() + .await; Ok(NotifyOutcomeSet { outcomes }) } diff --git a/src/domain/entities/user.rs b/src/domain/entities/user.rs index 22417cd8..4431831c 100644 --- a/src/domain/entities/user.rs +++ b/src/domain/entities/user.rs @@ -20,6 +20,18 @@ impl std::fmt::Display for UserRole { } } +/// Authorization-relevant account flags, fetched without the heavyweight +/// profile columns. The full user row drags `image` along — a data URI of +/// up to 512 KiB — which per-request guards (`require_internal_user`, +/// `require_admin_user`, the NC Basic Auth external check) must never pay +/// for just to read a boolean or a role. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct UserFlags { + pub role: UserRole, + pub is_external: bool, + pub active: bool, +} + #[derive(Debug, Clone)] pub struct User { id: Uuid, diff --git a/src/infrastructure/adapters/calendar_storage_adapter.rs b/src/infrastructure/adapters/calendar_storage_adapter.rs index 05b84c08..e37c3ac5 100644 --- a/src/infrastructure/adapters/calendar_storage_adapter.rs +++ b/src/infrastructure/adapters/calendar_storage_adapter.rs @@ -402,6 +402,26 @@ impl CalendarStoragePort for CalendarStorageAdapter { Ok(CalendarEventDto::from(event)) } + async fn find_event_by_ical_uid( + &self, + calendar_id: &str, + ical_uid: &str, + ) -> Result, DomainError> { + let uuid = Uuid::parse_str(calendar_id).map_err(|_| { + DomainError::new( + ErrorKind::InvalidInput, + "Calendar", + "Invalid calendar ID format", + ) + })?; + + let event = self + .event_repository + .find_event_by_ical_uid(&uuid, ical_uid) + .await?; + Ok(event.map(CalendarEventDto::from)) + } + async fn list_events_by_calendar( &self, calendar_id: &str, diff --git a/src/infrastructure/adapters/contact_storage_adapter.rs b/src/infrastructure/adapters/contact_storage_adapter.rs index d629be2c..1f8a06d6 100644 --- a/src/infrastructure/adapters/contact_storage_adapter.rs +++ b/src/infrastructure/adapters/contact_storage_adapter.rs @@ -712,6 +712,24 @@ impl ContactUseCase for ContactStorageAdapter { Ok(ContactDto::from(contact)) } + async fn get_contact_by_uid( + &self, + address_book_id: &str, + uid: &str, + user_id: Uuid, + ) -> Result, DomainError> { + let uuid = Self::parse_uuid(address_book_id, "AddressBook")?; + + // Check read access + self.check_address_book_access(&uuid, user_id).await?; + + let contact = self + .contact_repository + .get_contact_by_uid(&uuid, uid) + .await?; + Ok(contact.map(ContactDto::from)) + } + async fn list_contacts( &self, address_book_id: &str, diff --git a/src/infrastructure/repositories/pg/nextcloud_object_id_repository.rs b/src/infrastructure/repositories/pg/nextcloud_object_id_repository.rs index 3ffddcbc..5e83ed10 100644 --- a/src/infrastructure/repositories/pg/nextcloud_object_id_repository.rs +++ b/src/infrastructure/repositories/pg/nextcloud_object_id_repository.rs @@ -1,5 +1,7 @@ use sqlx::{PgPool, Row}; +use std::collections::HashMap; use std::sync::Arc; +use uuid::Uuid; use crate::common::errors::{DomainError, ErrorKind, Result}; @@ -12,29 +14,73 @@ impl NextcloudObjectIdRepository { Self { pool } } - pub async fn get_or_create(&self, object_type: &str, object_id: &str) -> Result { - let row = sqlx::query( + /// Resolve — creating when absent — stable numeric IDs for a batch of + /// object UUIDs sharing one `object_type`. + /// + /// Two statements instead of one per id: an idempotent bulk insert that + /// leaves existing rows untouched (`ON CONFLICT DO NOTHING` — no row + /// rewrite, no WAL churn, no dead tuples, unlike the former `DO UPDATE`), + /// followed by a single read of every requested mapping. The insert + /// auto-commits before the read, so the read observes both our own rows + /// and any created concurrently. Returns a map keyed by object UUID; + /// unresolvable inputs are simply absent. + pub async fn get_or_create_many( + &self, + object_type: &str, + object_ids: &[Uuid], + ) -> Result> { + if object_ids.is_empty() { + return Ok(HashMap::new()); + } + + // 1. Create missing mappings only. `DO NOTHING` skips the write for + // UUIDs that already map, eliminating the per-listing row rewrite. + sqlx::query( r#" INSERT INTO storage.nextcloud_object_ids (object_type, object_id) - VALUES ($1, $2::uuid) - ON CONFLICT (object_type, object_id) - DO UPDATE SET object_id = EXCLUDED.object_id - RETURNING id + SELECT $1, u FROM unnest($2::uuid[]) AS u + ON CONFLICT (object_type, object_id) DO NOTHING "#, ) .bind(object_type) - .bind(object_id) - .fetch_one(&*self.pool) + .bind(object_ids) + .execute(&*self.pool) .await .map_err(|e| { DomainError::new( ErrorKind::DatabaseError, "NextcloudFileId", - format!("Failed to get/create Nextcloud ID: {}", e), + format!("Failed to create Nextcloud IDs: {}", e), ) })?; - Ok(row.get::("id")) + // 2. Read every requested mapping back in a single round-trip. + let rows = sqlx::query( + r#" + SELECT id, object_id + FROM storage.nextcloud_object_ids + WHERE object_type = $1 AND object_id = ANY($2::uuid[]) + "#, + ) + .bind(object_type) + .bind(object_ids) + .fetch_all(&*self.pool) + .await + .map_err(|e| { + DomainError::new( + ErrorKind::DatabaseError, + "NextcloudFileId", + format!("Failed to load Nextcloud IDs: {}", e), + ) + })?; + + let mut map = HashMap::with_capacity(rows.len()); + for row in rows { + let object_id: Uuid = row.get("object_id"); + let id: i64 = row.get("id"); + map.insert(object_id, id); + } + Ok(map) } /// Get the OxiCloud object ID from a Nextcloud numeric ID. diff --git a/src/infrastructure/repositories/pg/user_pg_repository.rs b/src/infrastructure/repositories/pg/user_pg_repository.rs index 82132fad..9d896a5d 100644 --- a/src/infrastructure/repositories/pg/user_pg_repository.rs +++ b/src/infrastructure/repositories/pg/user_pg_repository.rs @@ -5,7 +5,7 @@ use uuid::Uuid; use crate::application::ports::auth_ports::UserStoragePort; use crate::common::errors::DomainError; -use crate::domain::entities::user::{User, UserRole}; +use crate::domain::entities::user::{User, UserFlags, UserRole}; use crate::domain::repositories::user_repository::{ StorageStats, UserRepository, UserRepositoryError, UserRepositoryResult, }; @@ -51,6 +51,40 @@ impl UserPgRepository { } } + /// Fetch only the authorization-relevant flags of a user. Not part of + /// the `UserRepository` trait — called directly from + /// `AuthApplicationService::get_user_flags`. + /// + /// Deliberately selects three tiny columns instead of the full row: + /// the full-row SELECT includes `image` (a data URI of up to 512 KiB), + /// which per-request middleware guards were paying on every WebDAV / + /// CalDAV / CardDAV request just to read `is_external` or `role`. + pub async fn get_user_flags(&self, id: Uuid) -> UserRepositoryResult { + let row = sqlx::query( + r#" + SELECT role::text as role_text, is_external, active + FROM auth.users + WHERE id = $1 + "#, + ) + .bind(id) + .fetch_one(&*self.pool) + .await + .map_err(Self::map_sqlx_error)?; + + let role_str: Option = row.try_get("role_text").unwrap_or(None); + let role = match role_str.as_deref() { + Some("admin") => UserRole::Admin, + _ => UserRole::User, + }; + + Ok(UserFlags { + role, + is_external: row.get("is_external"), + active: row.get("active"), + }) + } + /// Updates a user's profile image (URL or data URI). Not part of the /// `UserRepository` trait — called directly from `AuthApplicationService`. pub async fn update_image( diff --git a/src/infrastructure/services/image_transcode_service.rs b/src/infrastructure/services/image_transcode_service.rs index e38058e7..a30f49d0 100644 --- a/src/infrastructure/services/image_transcode_service.rs +++ b/src/infrastructure/services/image_transcode_service.rs @@ -7,8 +7,11 @@ //! - **Dedicated `rayon` thread pool** for CPU-bound transcoding (never blocks Tokio) //! - **`moka` lock-free cache** for hot transcoded images (no write-lock on reads) //! - Disk cache for persistence across restarts -//! - Supports JPEG, PNG, GIF → WebP conversion -//! - Falls back to original if conversion fails or result is larger +//! - Supports PNG, GIF → WebP conversion (JPEG excluded — the encoder is +//! lossless-only, so photos would come out larger; see `can_transcode`) +//! - Falls back to original if conversion fails or result is larger, and +//! remembers that negative verdict (memory sentinel + disk marker) so the +//! decode + encode is never repeated for the same file use bytes::Bytes; use image::ImageFormat; @@ -188,12 +191,16 @@ impl ImageTranscodeService { Ok(()) } - /// Check if a mime type can be transcoded + /// Check if a mime type can be transcoded. + /// + /// JPEG is deliberately excluded: the `image` crate's WebP encoder is + /// lossless-only, and losslessly re-encoding an already-lossy photo + /// almost always produces a LARGER file — so every JPEG download paid + /// a full decode + encode (hundreds of ms of CPU) only to discard the + /// result. PNG/GIF → lossless WebP genuinely shrinks. Re-add JPEG only + /// together with a lossy WebP encoder. pub fn can_transcode(mime_type: &str) -> bool { - matches!( - mime_type, - "image/jpeg" | "image/jpg" | "image/png" | "image/gif" - ) + matches!(mime_type, "image/png" | "image/gif") } /// Check if transcoding should be attempted based on file size and type @@ -216,8 +223,16 @@ impl ImageTranscodeService { let cache_key = format!("{}:{}", file_id, target_format.extension()); // ── 1. Check moka memory cache (lock-free read) ── + // An empty-Bytes entry is the negative sentinel: "transcoding this + // file is not beneficial — serve the original". Without it, every + // GET of such an image repeated the full decode + encode just to + // discard the result again. if let Some(cached) = self.memory_cache.get(&cache_key).await { self.stats.cache_hits.fetch_add(1, Ordering::Relaxed); + if cached.is_empty() { + tracing::debug!("🔥 Transcode negative cache HIT: {}", file_id); + return Ok((original_content, original_mime.to_string(), false)); + } tracing::debug!("🔥 Transcode memory cache HIT: {}", file_id); return Ok((cached, target_format.mime_type().to_string(), true)); } @@ -241,6 +256,15 @@ impl ImageTranscodeService { } } + // ── 2b. Negative verdict persisted on disk (survives restarts) ── + let skip_marker = self.get_skip_marker_path(file_id, target_format); + if tokio::fs::try_exists(&skip_marker).await.unwrap_or(false) { + self.memory_cache.insert(cache_key, Bytes::new()).await; + self.stats.disk_hits.fetch_add(1, Ordering::Relaxed); + tracing::debug!("💾 Transcode negative disk marker HIT: {}", file_id); + return Ok((original_content, original_mime.to_string(), false)); + } + // ── 3. Transcode on dedicated rayon pool (never blocks Tokio) ── let content_for_rayon = original_content.clone(); // O(1) ref-count bump let mime_owned = original_mime.to_string(); @@ -269,6 +293,20 @@ impl ImageTranscodeService { original_size, transcoded_size ); + // Remember the negative verdict so the next GET doesn't repeat + // the decode + encode: empty-Bytes sentinel in memory (expires + // with the cache TTL) + zero-byte marker on disk (survives + // restarts; removed by `invalidate` when the file changes). + self.memory_cache.insert(cache_key, Bytes::new()).await; + let marker = self.get_skip_marker_path(file_id, target_format); + tokio::spawn(async move { + if let Some(parent) = marker.parent() { + let _ = fs::create_dir_all(parent).await; + } + if let Err(e) = fs::write(&marker, b"").await { + tracing::warn!("Failed to persist transcode skip marker: {}", e); + } + }); return Ok((original_content, original_mime.to_string(), false)); } @@ -319,6 +357,16 @@ impl ImageTranscodeService { .join(format!("{}.{}", file_id, format.extension())) } + /// Path of the zero-byte marker recording a negative transcode verdict + /// ("result was not smaller — serve the original"). + fn get_skip_marker_path(&self, file_id: &str, format: OutputFormat) -> PathBuf { + self.cache_dir.join(format.extension()).join(format!( + "{}.{}.skip", + file_id, + format.extension() + )) + } + /// Invalidate cached transcodes for a file pub async fn invalidate(&self, file_id: &str) { let cache_key = format!("{}:{}", file_id, OutputFormat::WebP.extension()); @@ -326,6 +374,8 @@ impl ImageTranscodeService { let cache_path = self.get_cache_path(file_id, OutputFormat::WebP); let _ = fs::remove_file(&cache_path).await; + let skip_marker = self.get_skip_marker_path(file_id, OutputFormat::WebP); + let _ = fs::remove_file(&skip_marker).await; } /// Get transcoding statistics @@ -463,9 +513,10 @@ mod tests { #[test] fn test_can_transcode() { - assert!(ImageTranscodeService::can_transcode("image/jpeg")); assert!(ImageTranscodeService::can_transcode("image/png")); assert!(ImageTranscodeService::can_transcode("image/gif")); + // JPEG excluded: the lossless-only WebP encoder makes photos LARGER + assert!(!ImageTranscodeService::can_transcode("image/jpeg")); assert!(!ImageTranscodeService::can_transcode("image/webp")); assert!(!ImageTranscodeService::can_transcode("image/svg+xml")); assert!(!ImageTranscodeService::can_transcode("application/pdf")); @@ -473,16 +524,22 @@ mod tests { #[test] fn test_should_transcode() { - // Small JPEG - yes + // Small PNG - yes assert!(ImageTranscodeService::should_transcode( - "image/jpeg", + "image/png", 1024 * 1024 )); - // Large JPEG - no (too big) + // Large PNG - no (too big) + assert!(!ImageTranscodeService::should_transcode( + "image/png", + 10 * 1024 * 1024 + )); + + // JPEG - no (lossless-only encoder, result would be larger) assert!(!ImageTranscodeService::should_transcode( "image/jpeg", - 10 * 1024 * 1024 + 1024 * 1024 )); // WebP - no (already optimal) diff --git a/src/interfaces/api/handlers/auth_handler.rs b/src/interfaces/api/handlers/auth_handler.rs index 1261558e..64556952 100644 --- a/src/interfaces/api/handlers/auth_handler.rs +++ b/src/interfaces/api/handlers/auth_handler.rs @@ -290,10 +290,7 @@ pub async fn login( // same account from a different address (issue #323). The check runs // BEFORE Argon2 to save CPU under brute-force attacks. let client_ip = client_ip_from_parts(&headers, Some(peer), false); - if let Err(lockout_secs) = auth_service - .login_lockout - .check(&dto.username, &client_ip) - { + if let Err(lockout_secs) = auth_service.login_lockout.check(&dto.username, &client_ip) { tracing::warn!( target: "audit", event = "auth.login", diff --git a/src/interfaces/api/handlers/caldav_handler.rs b/src/interfaces/api/handlers/caldav_handler.rs index 99179144..98fded62 100644 --- a/src/interfaces/api/handlers/caldav_handler.rs +++ b/src/interfaces/api/handlers/caldav_handler.rs @@ -608,12 +608,13 @@ async fn handle_put( let ical_uid = extract_uid_from_ical(&ical_data); + // Indexed single-row lookup — listing the whole calendar (every row + // with its ical_data) to find one UID made imports O(N²). let existing = if let Some(ref uid) = ical_uid { - let events = calendar_service - .list_events(calendar_id, None, None, user.id) + calendar_service + .get_event_by_ical_uid(calendar_id, uid, user.id) .await - .unwrap_or_default(); - events.into_iter().find(|e| e.ical_uid == *uid) + .unwrap_or_default() } else { None }; @@ -704,21 +705,17 @@ async fn handle_get( .body(Body::from(ical)) .unwrap()) } else { - // GET on individual event + // GET on individual event — indexed lookup by iCalendar UID. let event_file = parts[1]; let ical_uid = event_file.trim_end_matches(".ics"); - let events = calendar_service - .list_events(calendar_id, None, None, user.id) + let event = calendar_service + .get_event_by_ical_uid(calendar_id, ical_uid, user.id) .await - .map_err(|e| AppError::internal_error(format!("Failed to list events: {}", e)))?; - - let event = events - .iter() - .find(|e| e.ical_uid == ical_uid) + .map_err(|e| AppError::internal_error(format!("Failed to look up event: {}", e)))? .ok_or_else(|| AppError::not_found(format!("Event not found: {}", ical_uid)))?; - let ical = generate_event_ical(event); + let ical = generate_event_ical(&event); Ok(Response::builder() .status(StatusCode::OK) @@ -813,14 +810,11 @@ async fn handle_delete( let event_file = parts[1]; let ical_uid = event_file.trim_end_matches(".ics"); - let events = calendar_service - .list_events(calendar_id, None, None, user.id) + // Indexed lookup by iCalendar UID instead of listing the calendar. + let event = calendar_service + .get_event_by_ical_uid(calendar_id, ical_uid, user.id) .await - .map_err(|e| AppError::internal_error(format!("Failed to list events: {}", e)))?; - - let event = events - .iter() - .find(|e| e.ical_uid == ical_uid) + .map_err(|e| AppError::internal_error(format!("Failed to look up event: {}", e)))? .ok_or_else(|| AppError::not_found(format!("Event not found: {}", ical_uid)))?; calendar_service diff --git a/src/interfaces/api/handlers/carddav_handler.rs b/src/interfaces/api/handlers/carddav_handler.rs index 3a6d273e..5b58c1fd 100644 --- a/src/interfaces/api/handlers/carddav_handler.rs +++ b/src/interfaces/api/handlers/carddav_handler.rs @@ -295,19 +295,14 @@ async fn handle_propfind( .body(Body::from(response_body)) .unwrap()) } else { - // Individual contact .vcf + // Individual contact .vcf — indexed lookup by vCard UID. let contact_file = parts[1]; let contact_uid = contact_file.trim_end_matches(".vcf"); - // Look up by UID across all contacts in this address book - let contacts = contact_svc - .list_contacts(address_book_id, user.id) + let contact = contact_svc + .get_contact_by_uid(address_book_id, contact_uid, user.id) .await - .map_err(|e| AppError::internal_error(format!("Failed to list contacts: {}", e)))?; - - let contact = contacts - .iter() - .find(|c| c.uid == contact_uid) + .map_err(|e| AppError::internal_error(format!("Failed to look up contact: {}", e)))? .ok_or_else(|| { AppError::not_found(format!("Contact not found: {}", contact_uid)) })?; @@ -322,8 +317,8 @@ async fn handle_propfind( let mut response_body = Vec::new(); CardDavAdapter::generate_contacts_response( &mut response_body, - std::slice::from_ref(contact), - &[(contact.uid.clone(), contact_to_vcard(contact))], + std::slice::from_ref(&contact), + &[(contact.uid.clone(), contact_to_vcard(&contact))], &report, base_href, ) @@ -483,13 +478,13 @@ async fn handle_put( // Extract UID from vCard let vcard_uid = extract_uid_from_vcard(&vcard_data); - // Check if contact already exists + // Check if contact already exists — indexed single-row lookup + // (listing the whole address book made imports O(N²)). let existing = if let Some(ref uid) = vcard_uid { - let contacts = contact_svc - .list_contacts(address_book_id, user.id) + contact_svc + .get_contact_by_uid(address_book_id, uid, user.id) .await - .unwrap_or_default(); - contacts.into_iter().find(|c| c.uid == *uid) + .unwrap_or_default() } else { None }; @@ -579,21 +574,17 @@ async fn handle_get( .body(Body::from(vcf_data)) .unwrap()) } else { - // GET on individual contact + // GET on individual contact — indexed lookup by vCard UID. let contact_file = parts[1]; let contact_uid = contact_file.trim_end_matches(".vcf"); - let contacts = contact_svc - .list_contacts(address_book_id, user.id) + let contact = contact_svc + .get_contact_by_uid(address_book_id, contact_uid, user.id) .await - .map_err(|e| AppError::internal_error(format!("Failed to list contacts: {}", e)))?; - - let contact = contacts - .iter() - .find(|c| c.uid == contact_uid) + .map_err(|e| AppError::internal_error(format!("Failed to look up contact: {}", e)))? .ok_or_else(|| AppError::not_found(format!("Contact not found: {}", contact_uid)))?; - let vcard = contact_to_vcard(contact); + let vcard = contact_to_vcard(&contact); Ok(Response::builder() .status(StatusCode::OK) @@ -632,18 +623,14 @@ async fn handle_delete( AppError::internal_error(format!("Failed to delete address book: {}", e)) })?; } else { - // Delete contact + // Delete contact — indexed lookup by vCard UID. let contact_file = parts[1]; let contact_uid = contact_file.trim_end_matches(".vcf"); - let contacts = contact_svc - .list_contacts(address_book_id, user.id) + let contact = contact_svc + .get_contact_by_uid(address_book_id, contact_uid, user.id) .await - .map_err(|e| AppError::internal_error(format!("Failed to list contacts: {}", e)))?; - - let contact = contacts - .iter() - .find(|c| c.uid == contact_uid) + .map_err(|e| AppError::internal_error(format!("Failed to look up contact: {}", e)))? .ok_or_else(|| AppError::not_found(format!("Contact not found: {}", contact_uid)))?; contact_svc diff --git a/src/interfaces/api/handlers/file_handler.rs b/src/interfaces/api/handlers/file_handler.rs index aface7df..1162dd2b 100644 --- a/src/interfaces/api/handlers/file_handler.rs +++ b/src/interfaces/api/handlers/file_handler.rs @@ -20,6 +20,7 @@ use crate::application::ports::{file_ports::OptimizedFileContent, folder_ports:: use crate::common::di::AppState; use crate::interfaces::errors::AppError; use crate::interfaces::middleware::auth::AuthUser; +use crate::interfaces::range_requests::not_modified_response; use crate::{application::dtos::file_dto::FileDto, domain::services::authorization::Permission}; use std::sync::Arc; @@ -600,16 +601,8 @@ impl FileHandler { let etag = format!("\"{}\"", file_dto.etag); // ── ETag (304 Not Modified) ────────────────────────────────── - if let Some(inm) = headers.get(header::IF_NONE_MATCH) - && let Ok(client_etag) = inm.to_str() - && (client_etag == etag || client_etag == "*") - { - return Response::builder() - .status(StatusCode::NOT_MODIFIED) - .header(header::ETAG, &etag) - .body(Body::empty()) - .unwrap() - .into_response(); + if let Some(resp) = not_modified_response(&headers, &etag) { + return resp.into_response(); } // ── Range Requests ─────────────────────────────────────────── diff --git a/src/interfaces/api/handlers/webdav_handler.rs b/src/interfaces/api/handlers/webdav_handler.rs index e9d895ea..c0097eb3 100644 --- a/src/interfaces/api/handlers/webdav_handler.rs +++ b/src/interfaces/api/handlers/webdav_handler.rs @@ -30,6 +30,7 @@ use crate::common::di::AppState; use crate::infrastructure::services::path_resolver_service::ResolvedResource; use crate::interfaces::errors::AppError; use crate::interfaces::middleware::auth::{AuthUser, CurrentUser}; +use crate::interfaces::range_requests::{not_modified_response, range_response}; use percent_encoding::{AsciiSet, NON_ALPHANUMERIC, percent_decode_str, utf8_percent_encode}; use std::sync::Arc; @@ -127,7 +128,9 @@ const MAX_MKCOL_BODY: usize = 4096; /// Batch size for streaming PROPFIND — files and folders are fetched in pages /// of this size to keep memory constant regardless of folder contents. -const PROPFIND_BATCH_SIZE: i64 = 500; +/// `pub(crate)` so the NextCloud PROPFIND handler streams with the same +/// page size. +pub(crate) const PROPFIND_BATCH_SIZE: i64 = 500; // ──────────────────────────────────────────────────────────────────────── // Security helpers (Sol.1 — handler-level user extraction & ownership guard) @@ -767,6 +770,21 @@ async fn handle_get( f }; + let etag = format!("\"{}\"", file.etag); + + // Conditional GET — clients revalidating a cached copy get a 304 + // instead of the full body. + if let Some(resp) = not_modified_response(req.headers(), &etag) { + return Ok(resp); + } + + // Range Requests — mount-style clients (rclone, davfs2, Finder) read + // by ranges; serve 206/416 instead of re-sending the whole file on + // every seek or resume. + if let Some(resp) = range_response(req.headers(), &file, &etag, file_retrieval_service).await { + return Ok(resp); + } + // Stream file content — constant ~64 KB memory regardless of file size let stream = file_retrieval_service .get_file_stream(&file.id) @@ -778,7 +796,8 @@ async fn handle_get( .status(StatusCode::OK) .header(header::CONTENT_TYPE, &*file.mime_type) .header(header::CONTENT_LENGTH, file.size) - .header(header::ETAG, format!("\"{}\"", file.etag)) + .header(header::ETAG, etag) + .header(header::ACCEPT_RANGES, "bytes") .header( header::LAST_MODIFIED, chrono::DateTime::::from_timestamp(file.created_at as i64, 0) diff --git a/src/interfaces/middleware/trusted_proxy.rs b/src/interfaces/middleware/trusted_proxy.rs index 0330e10d..d6bb6310 100644 --- a/src/interfaces/middleware/trusted_proxy.rs +++ b/src/interfaces/middleware/trusted_proxy.rs @@ -158,9 +158,7 @@ pub fn client_ip_from_parts( if let Some(peer_addr) = peer { if is_trusted_proxy(peer_addr.ip()) { // Try X-Forwarded-For first (leftmost = original client) - if let Some(xff) = headers - .get("x-forwarded-for") - .and_then(|v| v.to_str().ok()) + if let Some(xff) = headers.get("x-forwarded-for").and_then(|v| v.to_str().ok()) && let Some(ip) = xff .split(',') .next() diff --git a/src/interfaces/middleware/user.rs b/src/interfaces/middleware/user.rs index ddce6387..4aa4f6c7 100644 --- a/src/interfaces/middleware/user.rs +++ b/src/interfaces/middleware/user.rs @@ -4,9 +4,11 @@ //! so handlers compose them uniformly as one-liners. They assume the //! caller has already been authenticated by the //! [`AuthUser`](super::auth::AuthUser) extractor, and pull the current -//! user state from the database via `AuthApplicationService` so role / -//! external-flag changes take effect on the next request without -//! waiting for token rotation. +//! user flags via `AuthApplicationService::get_user_flags` — a +//! lightweight, short-TTL-cached lookup (no `image` column) — so role / +//! external-flag changes take effect within seconds without waiting +//! for token rotation, while the hot DAV paths stop paying one full-row +//! DB fetch per request. //! //! ```ignore //! let caller_id = auth_user.id; @@ -30,6 +32,7 @@ use uuid::Uuid; use crate::application::services::auth_application_service::AuthApplicationService; use crate::common::di::AppState; +use crate::domain::entities::user::UserRole; use crate::interfaces::errors::AppError; use crate::interfaces::middleware::auth::CurrentUser; @@ -56,8 +59,8 @@ pub async fn require_internal_user( auth: &AuthApplicationService, caller_id: Uuid, ) -> Result<(), AppError> { - match auth.get_user_by_id(caller_id).await { - Ok(dto) if dto.is_external => Err(AppError::new( + match auth.get_user_flags(caller_id).await { + Ok(flags) if flags.is_external => Err(AppError::new( StatusCode::FORBIDDEN, "External users cannot access this endpoint", "Forbidden", @@ -70,9 +73,11 @@ pub async fn require_internal_user( /// admins, `Err(403)` otherwise. /// /// The check pulls the role from the user record (not from JWT -/// claims) so a role change takes effect on the next request without -/// waiting for token rotation. Mirrors [`require_internal_user`]'s -/// shape so handlers compose either of them as a one-liner via `?`. +/// claims) so a role change takes effect within the flags-cache TTL — +/// or immediately when changed through `change_user_role`, which +/// invalidates the entry — without waiting for token rotation. Mirrors +/// [`require_internal_user`]'s shape so handlers compose either of +/// them as a one-liner via `?`. /// /// Use this in handlers that already have an /// [`AuthUser`](super::auth::AuthUser) extractor (and thus a validated @@ -82,12 +87,12 @@ pub async fn require_admin_user( auth: &AuthApplicationService, caller_id: Uuid, ) -> Result<(), AppError> { - let user = auth - .get_user_by_id(caller_id) + let flags = auth + .get_user_flags(caller_id) .await .map_err(AppError::from)?; - if user.role != "admin" { + if flags.role != UserRole::Admin { return Err(AppError::new( StatusCode::FORBIDDEN, "Admin access required", diff --git a/src/interfaces/mod.rs b/src/interfaces/mod.rs index 2aab264d..fef37883 100644 --- a/src/interfaces/mod.rs +++ b/src/interfaces/mod.rs @@ -2,6 +2,7 @@ pub mod api; pub mod errors; pub mod middleware; pub mod nextcloud; +pub mod range_requests; pub mod upload_spool; pub mod web; diff --git a/src/interfaces/nextcloud/basic_auth_middleware.rs b/src/interfaces/nextcloud/basic_auth_middleware.rs index 70b3c7e1..8486a4a1 100644 --- a/src/interfaces/nextcloud/basic_auth_middleware.rs +++ b/src/interfaces/nextcloud/basic_auth_middleware.rs @@ -64,8 +64,7 @@ pub async fn basic_auth_middleware( // Check account lockout before attempting password verification (saves CPU). // The lockout is per (account, IP), see #323 for rationale. - let client_ip = - crate::interfaces::middleware::rate_limit::extract_client_ip(&request); + let client_ip = crate::interfaces::middleware::rate_limit::extract_client_ip(&request); if let Some(auth_svc) = state.auth_service.as_ref() && let Err(secs) = auth_svc.login_lockout.check(&username, &client_ip) { @@ -103,11 +102,11 @@ pub async fn basic_auth_middleware( // this is the belt-and-braces check in case one slipped // through (e.g. user later flipped to is_external). if let Some(auth_svc) = state.auth_service.as_ref() - && let Ok(user) = auth_svc + && let Ok(flags) = auth_svc .auth_application_service - .get_user_by_id(user_id) + .get_user_flags(user_id) .await - && user.is_external + && flags.is_external { tracing::info!( target: "audit", diff --git a/src/interfaces/nextcloud/ocs_handler.rs b/src/interfaces/nextcloud/ocs_handler.rs index 5c255107..6c86fb8d 100644 --- a/src/interfaces/nextcloud/ocs_handler.rs +++ b/src/interfaces/nextcloud/ocs_handler.rs @@ -5,6 +5,7 @@ use axum::{ response::{IntoResponse, Response}, }; use serde_json::json; +use std::collections::HashMap; use std::sync::Arc; use crate::application::dtos::search_dto::SearchCriteriaDto; @@ -384,6 +385,17 @@ pub async fn handle_search( let file_id_svc = state.nextcloud.as_ref().map(|n| &n.file_ids); + // Pre-resolve numeric ids for every file result in a single batch query + // (was one INSERT round-trip per result). + let file_uuids: Vec = results.files.iter().map(|f| f.id.clone()).collect(); + let file_id_map: HashMap = match file_id_svc { + Some(svc) => svc + .get_or_create_file_ids(&file_uuids) + .await + .unwrap_or_default(), + None => HashMap::new(), + }; + let mut entries: Vec = Vec::new(); // Map file results @@ -394,11 +406,7 @@ pub async fn handle_search( .unwrap_or(&file.path); let display_path = format!("/{}", display_path); - let numeric_id = if let Some(svc) = file_id_svc { - svc.get_or_create_file_id(&file.id).await.ok() - } else { - None - }; + let numeric_id = file_id_map.get(&file.id).copied(); let thumbnail_url = match numeric_id { Some(nid) => format!("/index.php/core/preview?fileId={}&x=32&y=32", nid), diff --git a/src/interfaces/nextcloud/report_handler.rs b/src/interfaces/nextcloud/report_handler.rs index c59a07f2..d5e81b38 100644 --- a/src/interfaces/nextcloud/report_handler.rs +++ b/src/interfaces/nextcloud/report_handler.rs @@ -25,8 +25,7 @@ use crate::domain::entities::file::File; use crate::interfaces::errors::AppError; use crate::interfaces::middleware::auth::CurrentUser; use crate::interfaces::nextcloud::webdav_handler::{ - format_oc_id, nc_href, resolve_file_id, resolve_folder_id, write_file_response, - write_folder_response, + batch_resolve_ids, format_oc_id, nc_href, write_file_response, write_folder_response, }; /// Handle WebDAV REPORT and SEARCH methods for Nextcloud compatibility. @@ -87,56 +86,71 @@ async fn handle_filter_files( let home_prefix = format!("My Folder - {}/", user.username); + // Pass 1: fetch the favorited DTOs (the per-item fetch is a separate + // concern from the oc:fileid resolution batched below). + let mut files: Vec = Vec::new(); + let mut folders: Vec = Vec::new(); + for fav in &favorites { + match fav.item_type.as_str() { + "file" => { + if let Ok(f) = file_service.get_file(&fav.item_id).await { + files.push(f); + } + } + "folder" => { + if let Ok(f) = folder_service.get_folder(&fav.item_id).await { + folders.push(f); + } + } + _ => {} + } + } + + // Pass 2: resolve every oc:fileid in two batch queries (was one per item). + let file_uuids: Vec = files.iter().map(|f| f.id.clone()).collect(); + let folder_uuids: Vec = folders.iter().map(|f| f.id.clone()).collect(); + let (file_id_map, folder_id_map) = + batch_resolve_ids(file_id_svc, &file_uuids, &folder_uuids).await; + + // Pass 3: write the multistatus XML (pure synchronous map lookups). let mut buf = Vec::new(); { let mut xml = Writer::new(&mut buf); write_multistatus_start(&mut xml)?; - for fav in &favorites { - match fav.item_type.as_str() { - "file" => { - let file = match file_service.get_file(&fav.item_id).await { - Ok(f) => f, - Err(_) => continue, // Deleted or inaccessible -- skip. - }; - let subpath = strip_home_prefix(&file.path, &home_prefix); - let href = nc_href(&user.username, subpath); - let fid = resolve_file_id(file_id_svc, &file.id).await; - let oc_id = fid.map(|id| format_oc_id(id, file_id_svc)); - write_file_response( - &mut xml, - &file, - &href, - fid, - oc_id.as_deref(), - &user.username, - &favorite_ids, - ) - .map_err(|e| AppError::internal_error(format!("XML write error: {}", e)))?; - } - "folder" => { - let folder = match folder_service.get_folder(&fav.item_id).await { - Ok(f) => f, - Err(_) => continue, - }; - let subpath = strip_home_prefix(&folder.path, &home_prefix); - let href = format!("{}/", nc_href(&user.username, subpath)); - let fid = resolve_folder_id(file_id_svc, &folder.id).await; - let oc_id = fid.map(|id| format_oc_id(id, file_id_svc)); - write_folder_response( - &mut xml, - &folder, - &href, - fid, - oc_id.as_deref(), - &user.username, - &favorite_ids, - ) - .map_err(|e| AppError::internal_error(format!("XML write error: {}", e)))?; - } - _ => continue, - } + for file in &files { + let subpath = strip_home_prefix(&file.path, &home_prefix); + let href = nc_href(&user.username, subpath); + let fid = file_id_map.get(&file.id).copied(); + let oc_id = fid.map(|id| format_oc_id(id, file_id_svc)); + write_file_response( + &mut xml, + file, + &href, + fid, + oc_id.as_deref(), + &user.username, + &favorite_ids, + ) + .map_err(|e| AppError::internal_error(format!("XML write error: {}", e)))?; + } + + for folder in &folders { + let subpath = strip_home_prefix(&folder.path, &home_prefix); + let href = format!("{}/", nc_href(&user.username, subpath)); + let fid = folder_id_map.get(&folder.id).copied(); + let oc_id = fid.map(|id| format_oc_id(id, file_id_svc)); + write_folder_response( + &mut xml, + folder, + &href, + fid, + oc_id.as_deref(), + &user.username, + &favorite_ids, + ) + .map_err(|e| AppError::internal_error(format!("XML write error: {}", e)))?; } xml.write_event(Event::End(BytesEnd::new("d:multistatus"))) @@ -192,6 +206,15 @@ async fn handle_search( // No favorite checking for search results -- pass an empty set. let favorite_ids: HashSet = HashSet::new(); + // Materialize DTOs, then resolve every oc:fileid in two batch queries + // (was one INSERT round-trip per result). + let files: Vec = results.files.iter().map(file_dto_from_search).collect(); + let folders: Vec = results.folders.iter().map(folder_dto_from_search).collect(); + let file_uuids: Vec = files.iter().map(|f| f.id.clone()).collect(); + let folder_uuids: Vec = folders.iter().map(|f| f.id.clone()).collect(); + let (file_id_map, folder_id_map) = + batch_resolve_ids(file_id_svc, &file_uuids, &folder_uuids).await; + let mut buf = Vec::new(); { let mut xml = Writer::new(&mut buf); @@ -199,15 +222,14 @@ async fn handle_search( write_multistatus_start(&mut xml)?; // Files. - for fr in &results.files { - let file = file_dto_from_search(fr); + for file in &files { let subpath = strip_home_prefix(&file.path, &home_prefix); let href = nc_href(&user.username, subpath); - let fid = resolve_file_id(file_id_svc, &file.id).await; + let fid = file_id_map.get(&file.id).copied(); let oc_id = fid.map(|id| format_oc_id(id, file_id_svc)); write_file_response( &mut xml, - &file, + file, &href, fid, oc_id.as_deref(), @@ -218,15 +240,14 @@ async fn handle_search( } // Folders. - for sr in &results.folders { - let folder = folder_dto_from_search(sr); + for folder in &folders { let subpath = strip_home_prefix(&folder.path, &home_prefix); let href = format!("{}/", nc_href(&user.username, subpath)); - let fid = resolve_folder_id(file_id_svc, &folder.id).await; + let fid = folder_id_map.get(&folder.id).copied(); let oc_id = fid.map(|id| format_oc_id(id, file_id_svc)); write_folder_response( &mut xml, - &folder, + folder, &href, fid, oc_id.as_deref(), diff --git a/src/interfaces/nextcloud/trashbin_handler.rs b/src/interfaces/nextcloud/trashbin_handler.rs index 7c8cb286..60464ef2 100644 --- a/src/interfaces/nextcloud/trashbin_handler.rs +++ b/src/interfaces/nextcloud/trashbin_handler.rs @@ -14,7 +14,7 @@ use crate::common::di::AppState; use crate::interfaces::errors::AppError; use crate::interfaces::middleware::auth::{AuthUser, CurrentUser}; use crate::interfaces::nextcloud::webdav_handler::{ - format_oc_id, resolve_file_id, resolve_folder_id, write_text_element, + batch_resolve_ids, format_oc_id, write_text_element, }; const HEADER_DAV: HeaderName = HeaderName::from_static("dav"); @@ -198,6 +198,7 @@ fn strip_home_prefix<'a>(original_path: &'a str, username: &str) -> &'a str { use crate::application::dtos::trash_dto::TrashedItemDto; use crate::application::services::nextcloud_file_id_service::NextcloudFileIdService; +use std::collections::HashMap; /// Generate a complete Nextcloud-compatible multistatus XML response for the trashbin. async fn write_trashbin_multistatus( @@ -219,9 +220,25 @@ async fn write_trashbin_multistatus( // Root container entry for the trash collection itself. write_trash_root_response(&mut xml, username)?; + // Pre-resolve every oc:fileid in two batch queries by object type (was one + // INSERT round-trip per item). File and folder UUIDs are disjoint, so the + // two maps merge cleanly into one keyed by original_id. + let mut file_uuids: Vec = Vec::new(); + let mut folder_uuids: Vec = Vec::new(); + for item in items { + if item.item_type == "folder" { + folder_uuids.push(item.original_id.clone()); + } else { + file_uuids.push(item.original_id.clone()); + } + } + let (mut id_map, folder_id_map) = + batch_resolve_ids(file_id_svc, &file_uuids, &folder_uuids).await; + id_map.extend(folder_id_map); + // Individual trashed items. for item in items { - write_trash_item_response(&mut xml, item, username, file_id_svc).await?; + write_trash_item_response(&mut xml, item, username, file_id_svc, &id_map)?; } xml.write_event(Event::End(BytesEnd::new("d:multistatus"))) @@ -267,11 +284,12 @@ fn write_trash_root_response( } /// Write a single trashed item as a `` element. -async fn write_trash_item_response( +fn write_trash_item_response( xml: &mut Writer, item: &TrashedItemDto, username: &str, file_id_svc: Option<&Arc>, + id_map: &HashMap, ) -> Result<(), String> { xml.write_event(Event::Start(BytesStart::new("d:response"))) .map_err(|e| e.to_string())?; @@ -318,12 +336,8 @@ async fn write_trash_item_response( // d:getcontentlength write_text_element(xml, "d:getcontentlength", "0")?; - // oc:fileid and oc:id — resolve numeric ID via file_id service - let file_id = if item.item_type == "folder" { - resolve_folder_id(file_id_svc, &item.original_id).await - } else { - resolve_file_id(file_id_svc, &item.original_id).await - }; + // oc:fileid and oc:id — resolved up front in a batch query. + let file_id = id_map.get(&item.original_id).copied(); if let Some(id) = file_id { write_text_element(xml, "oc:fileid", &id.to_string())?; let oc_id = format_oc_id(id, file_id_svc); diff --git a/src/interfaces/nextcloud/webdav_handler.rs b/src/interfaces/nextcloud/webdav_handler.rs index 5a82baa7..47e4f4d7 100644 --- a/src/interfaces/nextcloud/webdav_handler.rs +++ b/src/interfaces/nextcloud/webdav_handler.rs @@ -3,16 +3,18 @@ use axum::{ http::{HeaderName, Request, StatusCode, header}, response::Response, }; -use bytes::Buf; +use bytes::{Buf, Bytes}; use chrono::Utc; use quick_xml::{ Writer, events::{BytesEnd, BytesStart, BytesText, Event}, }; -use std::collections::HashSet; +use std::collections::{HashMap, HashSet}; use std::sync::Arc; +use uuid::Uuid; use crate::application::adapters::webdav_adapter::{PropFindRequest, WebDavAdapter}; +use crate::application::dtos::pagination::PaginationRequestDto; use crate::application::ports::favorites_ports::FavoritesUseCase; use crate::application::ports::file_ports::{ FileManagementUseCase, FileRetrievalUseCase, FileUploadUseCase, @@ -21,8 +23,10 @@ use crate::application::ports::folder_ports::FolderUseCase; use crate::application::ports::trash_ports::TrashUseCase; use crate::common::di::AppState; use crate::common::mime_detect::{filename_from_path, refine_content_type_from_file}; +use crate::interfaces::api::handlers::webdav_handler::PROPFIND_BATCH_SIZE; use crate::interfaces::errors::AppError; use crate::interfaces::middleware::auth::{AuthUser, CurrentUser}; +use crate::interfaces::range_requests::{not_modified_response, range_response}; use crate::interfaces::upload_spool::spool_body_to_temp; /// Extension trait to map XML write errors to `String` concisely. @@ -118,8 +122,8 @@ pub async fn handle_nc_webdav( let method = req.method().clone(); match method.as_str() { "OPTIONS" => handle_options(), + "GET" => handle_get(state, &user, &subpath, req.headers()).await, "PROPFIND" => handle_propfind(state, req, &user, &subpath).await, - "GET" => handle_get(state, &user, &subpath).await, "PUT" => handle_put(state, req, &user, &subpath).await, "MKCOL" => handle_mkcol(state, &user, &subpath).await, "DELETE" => handle_delete(state, &user, &subpath).await, @@ -182,7 +186,10 @@ async fn handle_propfind( .await .map_err(|e| AppError::bad_request(format!("Failed to read body: {}", e)))?; - let propfind = if body_bytes.is_empty() { + // Parse (and thereby validate) the PROPFIND body. The NC response + // always emits the full property set, so the parsed request is not + // consulted further — but malformed XML must still fail with 400. + let _propfind = if body_bytes.is_empty() { PropFindRequest { prop_find_type: crate::application::adapters::webdav_adapter::PropFindType::AllProp, } @@ -199,64 +206,17 @@ async fn handle_propfind( let folder_result = folder_service.get_folder_by_path(&internal_path).await; if let Ok(folder) = folder_result { - // It's a folder. - let (files, subfolders) = if depth != "0" { - let files = file_service - .list_files(Some(&folder.id)) - .await - .unwrap_or_default(); - let subfolders = folder_service - .list_folders(Some(&folder.id)) - .await - .unwrap_or_default(); - (files, subfolders) - } else { - (vec![], vec![]) - }; - - // Batch-check favorites for all items in this listing. - let favorite_ids = if let Some(fav_svc) = state.favorites_service.as_ref() { - let mut items: Vec<(&str, &str)> = Vec::new(); - items.push((&folder.id, "folder")); - for f in &files { - items.push((&f.id, "file")); - } - for sf in &subfolders { - items.push((&sf.id, "folder")); - } - fav_svc - .batch_check_favorites(user.id, &items) - .await - .unwrap_or_default() - } else { - HashSet::new() - }; - - // Generate Nextcloud-aware XML. - let nc = state.nextcloud.as_ref(); - let file_id_svc = nc.map(|n| &n.file_ids); - - let mut buf = Vec::new(); - write_nc_multistatus( - &mut buf, - Some(&folder), - &files, - &subfolders, - &propfind, - &depth, - &user.username, - subpath, - file_id_svc, - &favorite_ids, - ) - .await - .map_err(|e| AppError::internal_error(format!("XML generation failed: {}", e)))?; - - return Ok(Response::builder() - .status(StatusCode::MULTI_STATUS) - .header(header::CONTENT_TYPE, "application/xml; charset=utf-8") - .body(Body::from(buf)) - .unwrap()); + // It's a folder — stream the multistatus: children are fetched in + // pages and serialized chunk by chunk, so memory stays O(batch) + // regardless of how many entries the folder holds. + return Ok(build_nc_streaming_propfind( + state.clone(), + folder, + depth, + user.id, + user.username.clone(), + subpath.to_string(), + )); } // Not a folder — try as a file. @@ -277,13 +237,9 @@ async fn handle_propfind( let file_id_svc = nc.map(|n| &n.file_ids); let mut buf = Vec::new(); - write_nc_multistatus( + write_nc_file_multistatus( &mut buf, - None, - &[file], - &[], - &propfind, - "0", + &file, &user.username, subpath, file_id_svc, @@ -308,6 +264,7 @@ async fn handle_get( state: Arc, user: &CurrentUser, subpath: &str, + headers: &axum::http::HeaderMap, ) -> Result, AppError> { // GET on root folder — NC clients use this as an existence check if subpath.is_empty() || subpath == "/" { @@ -340,6 +297,25 @@ async fn handle_get( .await .map_err(|_| AppError::not_found("File not found"))?; + // ETag comes from `FileDto::etag` (populated from `File::etag()` + // in the `From` impl) — single source of truth, so GET, + // HEAD, PUT-response, MOVE, and PROPFIND all emit byte-identical + // values for the same file. NC's sync engine compares cached + // PROPFIND ETags against GET/HEAD responses; using `file.id` here + // (a UUID) while PROPFIND emitted the blob hash made NC see + // every file as "remotely changed" on first descent. + let etag = format!("\"{}\"", file.etag); + + // Conditional GET — sync clients revalidating get a 304, not the body. + if let Some(resp) = not_modified_response(headers, &etag) { + return Ok(resp); + } + + // Range Requests — serve 206/416 instead of the whole file on seeks. + if let Some(resp) = range_response(headers, &file, &etag, file_service).await { + return Ok(resp); + } + let stream = file_service .get_file_stream(&file.id) .await @@ -349,18 +325,12 @@ async fn handle_get( chrono::DateTime::::from_timestamp(timestamp_to_i64(file.modified_at), 0) .unwrap_or_else(Utc::now); - // ETag comes from `FileDto::etag` (populated from `File::etag()` - // in the `From` impl) — single source of truth, so GET, - // HEAD, PUT-response, MOVE, and PROPFIND all emit byte-identical - // values for the same file. NC's sync engine compares cached - // PROPFIND ETags against GET/HEAD responses; using `file.id` here - // (a UUID) while PROPFIND emitted the blob hash made NC see - // every file as "remotely changed" on first descent. Ok(Response::builder() .status(StatusCode::OK) .header(header::CONTENT_TYPE, file.mime_type.as_ref()) .header(header::CONTENT_LENGTH, file.size) - .header(header::ETAG, format!("\"{}\"", file.etag)) + .header(header::ETAG, etag) + .header(header::ACCEPT_RANGES, "bytes") .header(header::LAST_MODIFIED, modified_at.to_rfc2822()) .body(Body::from_stream(std::pin::Pin::from(stream))) .unwrap()) @@ -987,98 +957,47 @@ use crate::application::dtos::file_dto::FileDto; use crate::application::dtos::folder_dto::FolderDto; use crate::application::services::nextcloud_file_id_service::NextcloudFileIdService; -/// Generate a complete Nextcloud-compatible multistatus XML response. -#[allow(clippy::too_many_arguments)] -async fn write_nc_multistatus( - writer: W, - folder: Option<&FolderDto>, - files: &[FileDto], - subfolders: &[FolderDto], - _request: &PropFindRequest, - depth: &str, - username: &str, - subpath: &str, - file_id_svc: Option<&Arc>, - favorite_ids: &HashSet, -) -> Result<(), String> { - let mut xml = Writer::new(writer); - - // Root element with all required namespaces. +/// Write the `` opening tag with the full NC namespace set. +/// Shared by the streaming folder PROPFIND and the single-file variant so +/// the namespace list can never diverge between the two. +fn write_nc_multistatus_open(xml: &mut Writer) -> Result<(), String> { let mut ms = BytesStart::new("d:multistatus"); ms.push_attribute(("xmlns:d", "DAV:")); ms.push_attribute(("xmlns:oc", "http://owncloud.org/ns")); ms.push_attribute(("xmlns:nc", "http://nextcloud.org/ns")); ms.push_attribute(("xmlns:ocs", "http://open-collaboration-services.org/ns")); - xml.write_event(Event::Start(ms)).xml_err()?; + xml.write_event(Event::Start(ms)).xml_err() +} - // Current folder entry. Collection hrefs MUST end in `/` (RFC 4918 - // §5.2 + strict NC-client enforcement — see `nc_collection_href`). - if let Some(f) = folder { - let href = nc_collection_href(username, subpath); - let file_id = resolve_folder_id(file_id_svc, &f.id).await; - let oc_id = file_id.map(|id| format_oc_id(id, file_id_svc)); - write_folder_response( - &mut xml, - f, - &href, - file_id, - oc_id.as_deref(), - username, - favorite_ids, - )?; - } +/// Generate the multistatus XML for a single-file PROPFIND. The folder +/// case streams via [`build_nc_streaming_propfind`] instead. +async fn write_nc_file_multistatus( + writer: W, + file: &FileDto, + username: &str, + subpath: &str, + file_id_svc: Option<&Arc>, + favorite_ids: &HashSet, +) -> Result<(), String> { + let (file_id_map, _) = + batch_resolve_ids(file_id_svc, std::slice::from_ref(&file.id), &[]).await; - // When folder is None, files are the target resource itself (single-file - // PROPFIND) and must always be emitted. When folder is Some, files/subfolders - // are children and should only be listed when depth > 0. - let emit_children = folder.is_none() || depth != "0"; + let mut xml = Writer::new(writer); + write_nc_multistatus_open(&mut xml)?; - if emit_children { - // Files. - for file in files { - let child_sub = if folder.is_none() { - // Single-file PROPFIND — subpath already points to the file. - subpath.to_string() - } else if subpath.is_empty() { - file.name.clone() - } else { - format!("{}/{}", subpath.trim_end_matches('/'), file.name) - }; - let href = nc_href(username, &child_sub); - let file_id = resolve_file_id(file_id_svc, &file.id).await; - let oc_id = file_id.map(|id| format_oc_id(id, file_id_svc)); - write_file_response( - &mut xml, - file, - &href, - file_id, - oc_id.as_deref(), - username, - favorite_ids, - )?; - } - - // Subfolders — also collections, same trailing-slash rule. - for sf in subfolders { - let child_sub = if subpath.is_empty() { - sf.name.clone() - } else { - format!("{}/{}", subpath.trim_end_matches('/'), sf.name) - }; - let href = nc_collection_href(username, &child_sub); - let file_id = resolve_folder_id(file_id_svc, &sf.id).await; - let oc_id = file_id.map(|id| format_oc_id(id, file_id_svc)); - write_folder_response( - &mut xml, - sf, - &href, - file_id, - oc_id.as_deref(), - username, - favorite_ids, - )?; - } - } + // Single-file PROPFIND — subpath already points to the file. + let href = nc_href(username, subpath); + let file_id = file_id_map.get(&file.id).copied(); + let oc_id = file_id.map(|id| format_oc_id(id, file_id_svc)); + write_file_response( + &mut xml, + file, + &href, + file_id, + oc_id.as_deref(), + username, + favorite_ids, + )?; xml.write_event(Event::End(BytesEnd::new("d:multistatus"))) .xml_err()?; @@ -1086,6 +1005,174 @@ async fn write_nc_multistatus( Ok(()) } +/// Build a streaming 207 Multi-Status response for a folder PROPFIND. +/// +/// Mirrors the native WebDAV handler's `build_streaming_propfind_response`: +/// children are fetched in pages of [`PROPFIND_BATCH_SIZE`], each page's +/// favorites and `oc:fileid`s are resolved with two batch queries, and the +/// XML is yielded chunk by chunk — memory stays O(batch) and the response +/// starts flowing immediately, instead of materializing the full listing +/// plus its entire multistatus (~2 KB/entry) in RAM before the first byte. +fn build_nc_streaming_propfind( + state: Arc, + folder: FolderDto, + depth: String, + user_id: Uuid, + username: String, + subpath: String, +) -> Response { + let stream = async_stream::try_stream! { + let file_id_svc = state.nextcloud.as_ref().map(|n| &n.file_ids); + let fav_svc = state.favorites_service.as_ref(); + let folder_service = &state.applications.folder_service; + let file_service = &state.applications.file_retrieval_service; + + // ── + the folder's own entry ───────────────── + // Collection hrefs MUST end in `/` (RFC 4918 §5.2 + strict + // NC-client enforcement — see `nc_collection_href`). + let folder_favs = if let Some(fav) = fav_svc { + fav.batch_check_favorites(user_id, &[(folder.id.as_str(), "folder")]) + .await + .unwrap_or_default() + } else { + HashSet::new() + }; + let (_, folder_id_map) = + batch_resolve_ids(file_id_svc, &[], std::slice::from_ref(&folder.id)).await; + + let mut buf = Vec::with_capacity(4096); + { + let mut xml = Writer::new(&mut buf); + write_nc_multistatus_open(&mut xml).map_err(std::io::Error::other)?; + let href = nc_collection_href(&username, &subpath); + let fid = folder_id_map.get(&folder.id).copied(); + let oc_id = fid.map(|id| format_oc_id(id, file_id_svc)); + write_folder_response(&mut xml, &folder, &href, fid, oc_id.as_deref(), &username, &folder_favs) + .map_err(std::io::Error::other)?; + } + yield Bytes::from(buf); + + // ── Children (only if Depth != 0) ──────────────────────────── + if depth != "0" { + // Files in pages. + let mut offset: i64 = 0; + loop { + let batch = file_service + .list_files_batch_with_perms(Some(&folder.id), user_id, offset, PROPFIND_BATCH_SIZE) + .await + .map_err(|e| std::io::Error::other(e.to_string()))?; + if batch.is_empty() { + break; + } + let batch_len = batch.len(); + + // Per-page enrichment: favorites + oc:fileids, two batch queries. + let favs = if let Some(fav) = fav_svc { + let items: Vec<(&str, &str)> = + batch.iter().map(|f| (f.id.as_str(), "file")).collect(); + fav.batch_check_favorites(user_id, &items).await.unwrap_or_default() + } else { + HashSet::new() + }; + let file_uuids: Vec = batch.iter().map(|f| f.id.clone()).collect(); + let (file_id_map, _) = batch_resolve_ids(file_id_svc, &file_uuids, &[]).await; + + let mut chunk = Vec::with_capacity(batch_len * 1024); + { + let mut xml = Writer::new(&mut chunk); + for file in &batch { + let child_sub = if subpath.is_empty() { + file.name.clone() + } else { + format!("{}/{}", subpath.trim_end_matches('/'), file.name) + }; + let href = nc_href(&username, &child_sub); + let fid = file_id_map.get(&file.id).copied(); + let oc_id = fid.map(|id| format_oc_id(id, file_id_svc)); + write_file_response(&mut xml, file, &href, fid, oc_id.as_deref(), &username, &favs) + .map_err(std::io::Error::other)?; + } + } + yield Bytes::from(chunk); + + if (batch_len as i64) < PROPFIND_BATCH_SIZE { + break; + } + offset += batch_len as i64; + } + + // Subfolders in pages — also collections, same trailing-slash rule. + let mut page = 0usize; + loop { + let pag = PaginationRequestDto { + page, + page_size: PROPFIND_BATCH_SIZE as usize, + }; + let result = folder_service + .list_folders_paginated_with_perms(Some(&folder.id), user_id, &pag) + .await + .map_err(|e| std::io::Error::other(e.to_string()))?; + if result.items.is_empty() { + break; + } + + let favs = if let Some(fav) = fav_svc { + let items: Vec<(&str, &str)> = + result.items.iter().map(|sf| (sf.id.as_str(), "folder")).collect(); + fav.batch_check_favorites(user_id, &items).await.unwrap_or_default() + } else { + HashSet::new() + }; + let folder_uuids: Vec = result.items.iter().map(|sf| sf.id.clone()).collect(); + let (_, sub_id_map) = batch_resolve_ids(file_id_svc, &[], &folder_uuids).await; + + let mut chunk = Vec::with_capacity(result.items.len() * 1024); + { + let mut xml = Writer::new(&mut chunk); + for sf in &result.items { + let child_sub = if subpath.is_empty() { + sf.name.clone() + } else { + format!("{}/{}", subpath.trim_end_matches('/'), sf.name) + }; + let href = nc_collection_href(&username, &child_sub); + let fid = sub_id_map.get(&sf.id).copied(); + let oc_id = fid.map(|id| format_oc_id(id, file_id_svc)); + write_folder_response(&mut xml, sf, &href, fid, oc_id.as_deref(), &username, &favs) + .map_err(std::io::Error::other)?; + } + } + let has_more = result.pagination.has_next; + yield Bytes::from(chunk); + + if !has_more { + break; + } + page += 1; + } + } + + // ── ───────────────────────────────────────── + let mut buf = Vec::with_capacity(32); + { + let mut xml = Writer::new(&mut buf); + xml.write_event(Event::End(BytesEnd::new("d:multistatus"))) + .map_err(|e| std::io::Error::other(e.to_string()))?; + } + yield Bytes::from(buf); + }; + + use futures::TryStreamExt; + let stream = stream + .map_err(|e: std::io::Error| -> Box { Box::new(e) }); + + Response::builder() + .status(StatusCode::MULTI_STATUS) + .header(header::CONTENT_TYPE, "application/xml; charset=utf-8") + .body(Body::from_stream(stream)) + .unwrap() +} + pub fn write_folder_response( xml: &mut Writer, folder: &FolderDto, @@ -1273,20 +1360,24 @@ pub fn write_text_element( Ok(()) } -pub async fn resolve_file_id( +/// Resolve every `oc:fileid` for a listing in two batch queries (one per +/// object type) instead of one INSERT round-trip per child. Returns +/// `(file_map, folder_map)` keyed by object UUID; entries are absent when the +/// service is disabled or an id can't be resolved, mirroring the previous +/// per-call `Option` behaviour. The two batches run concurrently. +pub async fn batch_resolve_ids( svc: Option<&Arc>, - file_uuid: &str, -) -> Option { - let svc = svc?; - svc.get_or_create_file_id(file_uuid).await.ok() -} - -pub async fn resolve_folder_id( - svc: Option<&Arc>, - folder_uuid: &str, -) -> Option { - let svc = svc?; - svc.get_or_create_folder_id(folder_uuid).await.ok() + file_uuids: &[String], + folder_uuids: &[String], +) -> (HashMap, HashMap) { + let Some(svc) = svc else { + return (HashMap::new(), HashMap::new()); + }; + let (files, folders) = tokio::join!( + svc.get_or_create_file_ids(file_uuids), + svc.get_or_create_folder_ids(folder_uuids), + ); + (files.unwrap_or_default(), folders.unwrap_or_default()) } pub fn format_oc_id(id: i64, svc: Option<&Arc>) -> String { diff --git a/src/interfaces/range_requests.rs b/src/interfaces/range_requests.rs new file mode 100644 index 00000000..b29296f8 --- /dev/null +++ b/src/interfaces/range_requests.rs @@ -0,0 +1,97 @@ +//! Conditional-request (`If-None-Match` → 304) and `Range` (→ 206/416) +//! helpers shared by the native WebDAV and NextCloud GET handlers. +//! +//! Mount-style WebDAV clients (rclone, davfs2, Finder, video players) +//! read by ranges; without 206 support every seek or resume transfers +//! the whole file. The REST download handler keeps its own richer +//! variant in `file_handler.rs` (Content-Disposition, compression, +//! permission-scoped streams) but follows the same header semantics. + +use axum::body::Body; +use axum::http::{HeaderMap, Response, StatusCode, header}; +use http_range_header::parse_range_header; +use std::sync::Arc; + +use crate::application::dtos::file_dto::FileDto; +use crate::application::ports::file_ports::FileRetrievalUseCase; +use crate::application::services::file_retrieval_service::FileRetrievalService; + +/// `If-None-Match` short-circuit: returns a `304 Not Modified` response +/// when the client already holds the current representation. `etag` must +/// be the quoted form the GET would emit (same comparison as the REST +/// download endpoint: exact match or `*`). +pub fn not_modified_response(headers: &HeaderMap, etag: &str) -> Option> { + let client_etag = headers.get(header::IF_NONE_MATCH)?.to_str().ok()?; + if client_etag == etag || client_etag == "*" { + return Some( + Response::builder() + .status(StatusCode::NOT_MODIFIED) + .header(header::ETAG, etag) + .body(Body::empty()) + .unwrap(), + ); + } + None +} + +/// `Range` short-circuit for a streaming download. +/// +/// Returns `Some(206)` with the requested byte range, `Some(416)` when +/// the range cannot be satisfied, or `None` when no (parseable) range +/// was requested or the range stream could not be created — callers +/// fall through to the full-body response, mirroring the REST handler. +/// +/// The caller has already resolved access to the file (path-resolver / +/// ownership), so this uses the unscoped range stream — same contract +/// as the `get_file_stream` call in the surrounding handlers. +pub async fn range_response( + headers: &HeaderMap, + file: &FileDto, + etag: &str, + retrieval: &Arc, +) -> Option> { + let range_str = headers.get(header::RANGE)?.to_str().ok()?; + let ranges = parse_range_header(range_str).ok()?; + + let valid_ranges = match ranges.validate(file.size) { + Ok(v) => v, + Err(_) => { + return Some( + Response::builder() + .status(StatusCode::RANGE_NOT_SATISFIABLE) + .header(header::CONTENT_RANGE, format!("bytes */{}", file.size)) + .body(Body::empty()) + .unwrap(), + ); + } + }; + + let range = valid_ranges.first()?; + let start = *range.start(); + let end = *range.end(); + let range_length = end - start + 1; + + match retrieval + .get_file_range_stream(&file.id, start, Some(end + 1)) + .await + { + Ok(stream) => Some( + Response::builder() + .status(StatusCode::PARTIAL_CONTENT) + .header(header::CONTENT_TYPE, &*file.mime_type) + .header(header::CONTENT_LENGTH, range_length) + .header( + header::CONTENT_RANGE, + format!("bytes {}-{}/{}", start, end, file.size), + ) + .header(header::ACCEPT_RANGES, "bytes") + .header(header::ETAG, etag) + .body(Body::from_stream(Box::into_pin(stream))) + .unwrap(), + ), + Err(err) => { + tracing::error!("Error creating range stream: {}", err); + None // fall through to the full download + } + } +} diff --git a/static/js/app/filesView.js b/static/js/app/filesView.js index 5f7893be..8324c36d 100644 --- a/static/js/app/filesView.js +++ b/static/js/app/filesView.js @@ -438,9 +438,17 @@ async function loadFiles(options = { insertHistory: true }) { } } - await rebuildBreadCrumb(); - ui.updateBreadcrumb(); - updateHistory(options.insertHistory ?? true); + // Rebuild the breadcrumb concurrently with the first page fetch — + // the listing does not depend on it, and awaiting the ancestor + // chain first added one round-trip per depth level before the + // content even started loading. Crumbs, history and title update + // when it resolves (skipped when superseded by a newer navigation). + const requestedPath = app.currentPath; + const breadcrumbReady = rebuildBreadCrumb().then((committed) => { + if (!committed) return; + ui.updateBreadcrumb(); + updateHistory(options.insertHistory ?? true); + }); clearTimeout(spinnerTimeout); @@ -457,6 +465,17 @@ async function loadFiles(options = { insertHistory: true }) { // Hand off to _loadPage (re-use cursor/groupBy state just reset above). _loading = false; // _loadPage sets its own guard await _loadPage({ isFirstPage: true }); + await breadcrumbReady; + + // rebuildBreadCrumb falls back to the home folder when the + // requested folder is inaccessible (deleted / permission revoked). + // The old sequential flow got the home listing for free; reload to + // match it. + if (app.currentPath !== requestedPath) { + ui.resetFilesList(); + _nextCursor = null; + await _loadPage({ isFirstPage: true }); + } // Deep-link: open a specific file if requested via app.viewFile. // We don't have a flat file list anymore (cursor pages), so only try @@ -496,6 +515,21 @@ async function loadFiles(options = { insertHistory: true }) { } } +/** + * Replace the list contents with a whole result set in one batched render + * (single DocumentFragment pass — no per-item DOM scans, smooth-scrolls or + * highlight pulses). Used by search to display results; optimistic + * single-item inserts should keep using `addItem`. + * + * @param {Array} items + */ +function renderItems(items) { + const component = _ensureComponent(); + if (!component) return; + document.getElementById('files-container-error')?.classList.add('hidden'); + component.render(items); +} + /** * Re-evaluate the shared badge for every item currently rendered in the Files list. * Call this after the outgoing grants cache has been refreshed. @@ -504,4 +538,4 @@ function refreshSharedBadges() { _component?.refreshSharedBadges(); } -export { addItem, filesView, loadFiles, refreshSharedBadges }; +export { addItem, filesView, loadFiles, refreshSharedBadges, renderItems }; diff --git a/static/js/app/ui.js b/static/js/app/ui.js index 68828898..673574e1 100644 --- a/static/js/app/ui.js +++ b/static/js/app/ui.js @@ -14,6 +14,7 @@ import { fileOps } from '../features/files/fileOperations.js'; import { inlineViewer } from '../features/files/inlineViewer.js'; import { wopiEditor } from '../features/files/wopiEditor.js'; import { recent } from '../features/library/recent.js'; +import { buildBatchDownloadUrl } from '../utils/download.js'; import { positionMenu } from '../utils/menuPosition.js'; import { loadFiles } from './filesView.js'; import { updateHistory } from './main.js'; @@ -774,7 +775,7 @@ const ui = { if (item?.type === 'file') fileIds.push(item.id); else if (item) folderIds.push(item.id); }); - downloadUrl = `${window.location.origin}/api/batch/download?file_ids=${fileIds.join(',')}&folder_ids=${folderIds.join(',')}`; + downloadUrl = `${window.location.origin}${buildBatchDownloadUrl(fileIds, folderIds)}`; } e.dataTransfer.setData('DownloadURL', `application/octet-stream:${nameEncoded}:${downloadUrl}`); @@ -1099,6 +1100,69 @@ function initRubberBandSelection() { let active = false; let startX = 0, startY = 0; + let curX = 0, + curY = 0; + let rafId = 0; + + /** + * Card geometry snapshot taken once per drag (and rebuilt on scroll). + * Comparing the lasso against these cached rects means the per-frame + * pass performs zero DOM reads — no forced reflow per card. + * @type {Array<{el: HTMLElement, left: number, top: number, right: number, + * bottom: number, info: ReturnType, + * selected: boolean}> | null} + */ + let cardCache = null; + + const buildCardCache = () => { + cardCache = []; + document.querySelectorAll('#files-list .file-item').forEach((card) => { + const el = /** @type {HTMLElement} */ (card); + const r = el.getBoundingClientRect(); + cardCache.push({ + el, + left: r.left, + top: r.top, + right: r.right, + bottom: r.bottom, + info: batchToolbar ? batchToolbar._extractInfo(/** @type {HTMLDivElement} */ (el)) : null, + selected: el.classList.contains('selected') + }); + }); + }; + + // Scrolling mid-drag shifts every viewport rect — drop the snapshot so + // the next frame rebuilds it. + const invalidateCardCache = () => { + cardCache = null; + }; + + /** One classification pass per animation frame (cached rects only). */ + const classifyCards = () => { + rafId = 0; + if (!cardCache) buildCardCache(); + + const left = Math.min(startX, curX); + const top = Math.min(startY, curY); + const right = Math.max(startX, curX); + const bottom = Math.max(startY, curY); + + for (const entry of cardCache) { + const intersects = entry.left < right && entry.right > left && entry.top < bottom && entry.bottom > top; + if (intersects === entry.selected) continue; + entry.selected = intersects; + entry.el.classList.toggle('selected', intersects); + + // Sync with batchToolbar module (only on state change) + if (batchToolbar && entry.info) { + if (intersects) { + batchToolbar.select(entry.info.id, entry.info.name, entry.info.type, entry.info.parentId); + } else { + batchToolbar.deselect(entry.info.id); + } + } + } + }; // We listen on the whole files-container (covers grid + empty space) const container = document.querySelector('.files-container') || document.getElementById('files-list'); @@ -1123,6 +1187,10 @@ function initRubberBandSelection() { active = true; startX = e.clientX; startY = e.clientY; + curX = startX; + curY = startY; + cardCache = null; // built lazily on the first classification frame + document.addEventListener('scroll', invalidateCardCache, { capture: true, passive: true }); selRect.style.left = `${startX}px`; selRect.style.top = `${startY}px`; @@ -1136,8 +1204,8 @@ function initRubberBandSelection() { document.addEventListener('mousemove', (e) => { if (!active) return; - const curX = e.clientX; - const curY = e.clientY; + curX = e.clientX; + curY = e.clientY; const left = Math.min(startX, curX); const top = Math.min(startY, curY); @@ -1149,41 +1217,27 @@ function initRubberBandSelection() { selRect.style.display = 'block'; } + // Style writes only — no layout reads here. The card highlighting + // runs at most once per frame against the cached geometry. selRect.style.left = `${left}px`; selRect.style.top = `${top}px`; selRect.style.width = `${width}px`; selRect.style.height = `${height}px`; - // Highlight cards that intersect with the rectangle - const rectBounds = { left, top, right: left + width, bottom: top + height }; - - document.querySelectorAll('#files-list .file-item').forEach((card) => { - const cardRect = card.getBoundingClientRect(); - const intersects = - cardRect.left < rectBounds.right && cardRect.right > rectBounds.left && cardRect.top < rectBounds.bottom && cardRect.bottom > rectBounds.top; - - if (intersects) { - card.classList.add('selected'); - - // Sync with batchToolbar module - if (batchToolbar) { - const info = batchToolbar._extractInfo(/** @type {HTMLDivElement} */ (card)); - if (info) batchToolbar.select(info.id, info.name, info.type, info.parentId); - } - } else { - card.classList.remove('selected'); - // Deselect from batchToolbar module - if (batchToolbar) { - const info = batchToolbar._extractInfo(/** @type {HTMLDivElement} */ (card)); - if (info) batchToolbar.deselect(info.id); - } - } - }); + if (!rafId) rafId = requestAnimationFrame(classifyCards); }); document.addEventListener('mouseup', () => { if (!active) return; active = false; + document.removeEventListener('scroll', invalidateCardCache, { capture: true }); + // Apply the still-pending classification so the final lasso + // position is what determines the selection. + if (rafId) { + cancelAnimationFrame(rafId); + classifyCards(); + } + cardCache = null; const hadSelection = selRect.style.display === 'block'; selRect.style.display = 'none'; // Update the batch bar after rubber band selection completes diff --git a/static/js/components/resourceList.js b/static/js/components/resourceList.js index 0578be27..73b633fc 100644 --- a/static/js/components/resourceList.js +++ b/static/js/components/resourceList.js @@ -643,7 +643,6 @@ export class ResourceListComponent { `; el.querySelector('.resource-icon-slot')?.replaceWith(buildResourceIcon(folder, 'folder')); - this._bindItemEvents(el, folder); return el; } @@ -693,7 +692,6 @@ export class ResourceListComponent { `; el.querySelector('.resource-icon-slot')?.replaceWith(buildResourceIcon(file, 'file')); - this._bindItemEvents(el, file); return el; } @@ -716,57 +714,6 @@ export class ResourceListComponent { .join(''); } - /** - * Attach direct event listeners to interactive elements inside a .file-item. - * This covers buttons that must stop propagation before the delegated listener runs. - * @param {HTMLElement} el - * @param {FileItem|FolderItem} item - */ - _bindItemEvents(el, item) { - const cfg = this._cfg; - - // Favorite-star — direct click, stopPropagation so the card open doesn't fire - if (cfg.showFavorite && cfg.onFavoriteToggle) { - const star = el.querySelector('.favorite-star'); - star?.addEventListener('click', (e) => { - e.stopPropagation(); - e.stopImmediatePropagation(); - e.preventDefault(); - cfg.onFavoriteToggle?.(item); - }); - } - - // Custom inline actions (e.g. restore / delete-permanently on trash) — - // bound directly so they stop propagation before the card-open handler. - if (cfg.customActions?.length) { - el.querySelectorAll('button[data-custom-action]').forEach((btn) => { - btn.addEventListener('click', (e) => { - e.stopPropagation(); - e.stopImmediatePropagation(); - e.preventDefault(); - const idx = Number(/** @type {HTMLElement} */ (btn).dataset.customAction); - const action = cfg.customActions?.[idx]; - if (action) action.onClick(item); - }); - }); - } - - // Shared-badge click → open share modal (or fall back to context menu) - if (cfg.showShareBadge && (cfg.onShareBadgeClick || cfg.onContextMenu)) { - const badge = el.querySelector('.file-badge-shared'); - badge?.addEventListener('click', (e) => { - e.stopPropagation(); - e.stopImmediatePropagation(); - e.preventDefault(); - if (cfg.onShareBadgeClick) { - cfg.onShareBadgeClick(item); - } else { - cfg.onContextMenu?.(item, /** @type {MouseEvent} */ (e)); - } - }); - } - } - /** Wire one delegated listener for all pointer events in this container. */ _initDelegation() { const container = this._container; @@ -791,6 +738,39 @@ export class ResourceListComponent { return; } + // Favorite-star button — never opens the card + if (target.closest('.favorite-star')) { + e.preventDefault(); + const item = this._itemFromCard(card); + if (item) cfg.onFavoriteToggle?.(item); + return; + } + + // Custom inline actions (e.g. restore / delete-permanently on trash) + const customBtn = /** @type {HTMLElement | null} */ (target.closest('button[data-custom-action]')); + if (customBtn) { + e.preventDefault(); + const idx = Number(customBtn.dataset.customAction); + const action = cfg.customActions?.[idx]; + const item = this._itemFromCard(card); + if (action && item) action.onClick(item); + return; + } + + // Shared-badge click → open share modal (or fall back to context menu) + if (target.closest('.file-badge-shared')) { + e.preventDefault(); + const item = this._itemFromCard(card); + if (item) { + if (cfg.onShareBadgeClick) { + cfg.onShareBadgeClick(item); + } else { + cfg.onContextMenu?.(item, /** @type {MouseEvent} */ (e)); + } + } + return; + } + // Checkbox cell → selection (shift extends range) if (cfg.selectable && target.closest('.checkbox-cell')) { if (e.shiftKey) { @@ -801,9 +781,6 @@ export class ResourceListComponent { return; } - // Favorite star is handled by the direct listener in _bindItemEvents - if (target.closest('.favorite-star')) return; - // Modifier-key click → selection toggle if (e.metaKey || e.altKey || e.ctrlKey) { if (cfg.selectable) this._toggleSelection(card); diff --git a/static/js/core/icons.js b/static/js/core/icons.js index 66420925..b5c857cd 100644 --- a/static/js/core/icons.js +++ b/static/js/core/icons.js @@ -618,27 +618,53 @@ function replaceIconsInElement(container) { function oxiIconsInit() { let raf = 0; + /** + * Subtree roots added since the last animation frame. Scanning only + * these (instead of the whole document) keeps the cost proportional + * to what was inserted, not to the total DOM size. + * @type {Set} + */ + let pendingRoots = new Set(); + const scan = () => { raf = 0; - replaceIconsInElement(); + const roots = pendingRoots; + pendingRoots = new Set(); + for (const root of roots) { + if (!root.isConnected) continue; // removed (or replaced) meanwhile + if (root.matches('i[class*="fa-"]')) { + // The inserted node IS the icon — scan via its parent so the + // descendant selector pass picks it up. + replaceIconsInElement(root.parentElement || document.body); + } else { + replaceIconsInElement(root); + } + } }; // Initial sweep once the DOM is ready + const fullSweep = () => replaceIconsInElement(); if (document.readyState === 'loading') { - document.addEventListener('DOMContentLoaded', scan); + document.addEventListener('DOMContentLoaded', fullSweep); } else { - scan(); + fullSweep(); } - // Observe future mutations (dynamic renders, modals, etc.) + // Observe future mutations (dynamic renders, modals, etc.). Only element + // insertions can carry icons — text-node churn (progress counters, + // notification text) no longer triggers any scan at all. new MutationObserver((mutations) => { - if (raf) return; - for (let i = 0; i < mutations.length; i++) { - if (mutations[i].addedNodes.length) { - raf = requestAnimationFrame(scan); - return; + for (const mutation of mutations) { + for (let i = 0; i < mutation.addedNodes.length; i++) { + const node = mutation.addedNodes[i]; + if (node.nodeType === Node.ELEMENT_NODE) { + pendingRoots.add(/** @type {Element} */ (node)); + } } } + if (!raf && pendingRoots.size) { + raf = requestAnimationFrame(scan); + } }).observe(document.documentElement, { childList: true, subtree: true }); } diff --git a/static/js/features/files/batchToolbar.js b/static/js/features/files/batchToolbar.js index 99aadc1a..07c123cc 100644 --- a/static/js/features/files/batchToolbar.js +++ b/static/js/features/files/batchToolbar.js @@ -15,6 +15,7 @@ import { loadFiles } from '../../app/filesView.js'; import { app } from '../../app/state.js'; import { showConfirmDialog, ui } from '../../app/ui.js'; import { i18n } from '../../core/i18n.js'; +import { buildBatchDownloadUrl, triggerBrowserDownload } from '../../utils/download.js'; import { favorites } from '../library/favorites.js'; import { contextMenus } from './contextMenus.js'; import { getAuthHeaders } from './fileOperations.js'; @@ -126,10 +127,14 @@ const batchToolbar = { clear() { this._selected.clear(); this._lastClickedIndex = -1; - document.querySelectorAll('.file-item.selected').forEach((el) => { + // Scope the DOM sweep to the files list — the only container this + // toolbar manages (see `selectAll`) — instead of the whole document, + // and only touch checkboxes that are actually checked. + const list = document.getElementById('files-list'); + list?.querySelectorAll('.file-item.selected').forEach((el) => { el.classList.remove('selected'); }); - document.querySelectorAll('.item-checkbox').forEach((cb) => { + list?.querySelectorAll('.item-checkbox:checked').forEach((cb) => { /** @type {HTMLInputElement} */ (cb).checked = false; }); // Reset the active component's internal selection state without going @@ -173,15 +178,16 @@ const batchToolbar = { /** @type {Array} */ const folderIds = []; - // TODO optimize & check if _selected is a better use - /** @type {NodeListOf} */ (document.querySelectorAll(`div.file-item.selected`)).forEach((item) => { - if (item.dataset.fileId) { - fileIds.push(item.dataset.fileId); - } else { - // ignore selectedItem if this is the target - if (targtFolderId && targtFolderId !== item.dataset.folderId) folderIds.push(item.dataset.folderId); + // `_selected` is the source of truth (every selection path keeps it + // in sync) — no need to re-derive the selection from a DOM scan. + for (const sel of this._selected.values()) { + if (sel.type === 'file') { + fileIds.push(sel.id); + } else if (targtFolderId && targtFolderId !== sel.id) { + // ignore the selected folder if it is the drop target itself + folderIds.push(sel.id); } - }); + } return { fileIds: fileIds, @@ -454,10 +460,22 @@ const batchToolbar = { ui.showNotification('Preparing download', 'Creating ZIP archive...'); - try { - const fileIds = items.filter((i) => i.type === 'file').map((i) => i.id); - const folderIds = items.filter((i) => i.type === 'folder').map((i) => i.id); + const fileIds = items.filter((i) => i.type === 'file').map((i) => i.id); + const folderIds = items.filter((i) => i.type === 'folder').map((i) => i.id); + const zipName = `oxicloud-download-${Date.now()}.zip`; + // Browser-native download via the GET variant of the endpoint: the + // ZIP streams to disk instead of being buffered whole in the tab's + // memory (a multi-GB selection used to risk crashing the tab). + const url = buildBatchDownloadUrl(fileIds, folderIds); + if (url.length <= 4000) { + triggerBrowserDownload(url, zipName); + return; + } + + // Selections too large for a URL (~100+ items) keep the buffered + // POST path — the id list only fits in a request body. + try { const response = await fetch('/api/batch/download', { method: 'POST', headers: { ...getAuthHeaders(), 'Content-Type': 'application/json' }, @@ -467,14 +485,14 @@ const batchToolbar = { if (!response.ok) throw new Error(`Server returned ${response.status}`); const blob = await response.blob(); - const url = URL.createObjectURL(blob); + const blobUrl = URL.createObjectURL(blob); const link = document.createElement('a'); - link.href = url; - link.download = `oxicloud-download-${Date.now()}.zip`; + link.href = blobUrl; + link.download = zipName; document.body.appendChild(link); link.click(); document.body.removeChild(link); - URL.revokeObjectURL(url); + URL.revokeObjectURL(blobUrl); } catch (e) { console.error('Batch download error:', e); ui.showNotification('Error', 'Could not download selected items'); diff --git a/static/js/features/files/fileOperations.js b/static/js/features/files/fileOperations.js index 8cb6a3a7..87f7578d 100644 --- a/static/js/features/files/fileOperations.js +++ b/static/js/features/files/fileOperations.js @@ -10,6 +10,8 @@ import { showConfirmDialog, ui } from '../../app/ui.js'; import { getCsrfHeaders, getCsrfToken } from '../../core/csrf.js'; import { i18n } from '../../core/i18n.js'; import { notifications } from '../../core/notifications.js'; +import { invalidateFolderMeta } from '../../model/filesModel.js'; +import { triggerBrowserDownload } from '../../utils/download.js'; /** * @typedef {Object} BatchResult @@ -349,7 +351,9 @@ const fileOps = { } // Filter out unreadable entries (typically dropped folders/placeholders) + /** @type {File[]} */ const readableFiles = []; + /** @type {string[]} */ const skippedEntries = []; for (const f of originalFiles) { // eslint-disable-next-line no-await-in-loop @@ -386,13 +390,21 @@ const fileOps = { let uploadedCount = 0; let successCount = 0; + let quotaStop = false; - for (let i = 0; i < totalFiles; i++) { - const file = readableFiles[i]; + const targetFolderId = app.currentPath || app.userHomeFolderId; + + /** + * Upload a single readable file by index. Shared counters are + * mutated here; safe because JS runs the workers cooperatively + * (no true parallelism between awaits). + * @param {number} idx + */ + const uploadOneFile = async (idx) => { + if (quotaStop) return; + const file = readableFiles[idx]; const formData = new FormData(); - - const targetFolderId = app.currentPath || app.userHomeFolderId; if (targetFolderId) formData.append('folder_id', targetFolderId); formData.append('file', file); @@ -436,6 +448,8 @@ const fileOps = { }); } if (result.isQuotaError) { + // Stop pulling new files; in-flight uploads still finish. + quotaStop = true; const msg = result.errorMsg || i18n.t('storage_quota_exceeded'); if (notifications) { notifications.addNotification({ @@ -445,10 +459,27 @@ const fileOps = { text: msg }); } - break; } } + }; + + // Pool-based concurrency: keep up to CONCURRENCY uploads in flight + // instead of one at a time (mirrors uploadFolderEntries). Files are + // independent, so this is ~CONCURRENCY× faster for many small files. + const CONCURRENCY = 10; + let nextIdx = 0; + const runNext = async () => { + while (nextIdx < totalFiles && !quotaStop) { + const idx = nextIdx++; + await uploadOneFile(idx); + } + }; + + const workers = []; + for (let w = 0; w < Math.min(CONCURRENCY, totalFiles); w++) { + workers.push(runNext()); } + await Promise.all(workers); // All done this._finishUploadToast(successCount, totalFiles); @@ -862,6 +893,8 @@ const fileOps = { }); if (response.ok) { + // Parent changed — drop the cached breadcrumb metadata + invalidateFolderMeta(folderId); // Reload files after moving await loadFiles(); ui.showNotification('Folder moved', 'Folder moved successfully'); @@ -925,6 +958,8 @@ const fileOps = { const data = await res.json(); success += data.stats?.successful || 0; errors += data.stats?.failed || 0; + // Parents changed — drop the cached breadcrumb metadata + for (const id of folderIds) invalidateFolderMeta(id); } } catch (err) { console.error('Batch move error:', err); @@ -1117,6 +1152,8 @@ const fileOps = { console.log('Response status:', response.status); if (response.ok) { + // Name changed — drop the cached breadcrumb metadata + invalidateFolderMeta(folderId); ui.showNotification('Folder renamed', `Folder renamed to "${newName}"`); } else { const errorText = await response.text(); @@ -1340,32 +1377,13 @@ const fileOps = { }, /** - * Download a file + * Download a file — handed to the browser so it streams to disk with + * its native download UI instead of buffering the file in memory. * @param {string} fileId - File ID * @param {string} fileName - File name */ async downloadFile(fileId, fileName) { - try { - const response = await fetch(`/api/files/${fileId}`, { - headers: getAuthHeaders() - }); - if (response.ok) { - const blob = await response.blob(); - const url = URL.createObjectURL(blob); - const link = document.createElement('a'); - link.href = url; - link.download = fileName; - document.body.appendChild(link); - link.click(); - document.body.removeChild(link); - URL.revokeObjectURL(url); - } else { - ui.showNotification('Error', 'Error downloading the file'); - } - } catch (error) { - console.error('Error downloading file:', error); - ui.showNotification('Error', 'Error downloading the file'); - } + triggerBrowserDownload(`/api/files/${fileId}`, fileName); }, /** @@ -1374,30 +1392,10 @@ const fileOps = { * @param {string} folderName - Folder name */ async downloadFolder(folderId, folderName) { - try { - // Show notification to user - ui.showNotification('Preparing download', 'Preparing the folder for download...'); - - const response = await fetch(`/api/folders/${folderId}/download?format=zip`, { - headers: getAuthHeaders() - }); - if (response.ok) { - const blob = await response.blob(); - const url = URL.createObjectURL(blob); - const link = document.createElement('a'); - link.href = url; - link.download = `${folderName}.zip`; - document.body.appendChild(link); - link.click(); - document.body.removeChild(link); - URL.revokeObjectURL(url); - } else { - ui.showNotification('Error', 'Error downloading the folder'); - } - } catch (error) { - console.error('Error downloading folder:', error); - ui.showNotification('Error', 'Error downloading the folder'); - } + // Show notification to user (the server still has to assemble the + // ZIP before the browser's own download UI takes over). + ui.showNotification('Preparing download', 'Preparing the folder for download...'); + triggerBrowserDownload(`/api/folders/${folderId}/download?format=zip`, `${folderName}.zip`); } }; diff --git a/static/js/features/files/inlineViewer.js b/static/js/features/files/inlineViewer.js index c4b14c74..e97133c3 100644 --- a/static/js/features/files/inlineViewer.js +++ b/static/js/features/files/inlineViewer.js @@ -6,6 +6,7 @@ import { updateHistory } from '../../app/main.js'; import { app } from '../../app/state.js'; import { isTextViewable } from '../../core/formatters.js'; +import { triggerBrowserDownload } from '../../utils/download.js'; import { wopiEditor } from './wopiEditor.js'; /** @import {FileItem} from '../../core/types.js' */ @@ -397,111 +398,92 @@ class InlineViewer { } /** - * Creates an audio or video player using blob URL (authenticated fetch) + * Creates an audio or video player that streams straight from the API. + * The element's `src` points at the same-origin endpoint (cookies are + * sent automatically), so the browser issues Range requests and starts + * playback progressively — the file is never materialized in memory, + * and seeking works without downloading everything first. * @param {FileItem} file * @param {string} mediaType * @param {HTMLDivElement} container * @param {HTMLDivElement} loader */ - async createMediaViewer(file, mediaType, container, loader) { - try { - console.log(`Creating ${mediaType} player for:`, file.name); + createMediaViewer(file, mediaType, container, loader) { + console.log(`Creating ${mediaType} player for:`, file.name); - // Fetch file (cookie auto-sent) - const response = await fetch(`/api/files/${file.id}?inline=true`, { - credentials: 'same-origin' - }); + const streamUrl = `/api/files/${file.id}?inline=true`; - if (!response.ok) { - throw new Error(`Error fetching file: ${response.status} ${response.statusText}`); - } + // The native player has its own buffering UI — drop our spinner now. + if (loader?.parentNode) { + loader.parentNode.removeChild(loader); + } - const blob = await response.blob(); - const blobUrl = URL.createObjectURL(blob); + if (mediaType === 'audio') { + // Wrapper with icon + player + const wrapper = document.createElement('div'); + wrapper.className = 'inline-viewer-audio-wrapper'; - // Remove loader - if (loader?.parentNode) { - loader.parentNode.removeChild(loader); - } + const icon = document.createElement('div'); + icon.className = 'inline-viewer-audio-icon'; + icon.innerHTML = ''; + wrapper.appendChild(icon); - if (mediaType === 'audio') { - // Wrapper with icon + player - const wrapper = document.createElement('div'); - wrapper.className = 'inline-viewer-audio-wrapper'; + const nameEl = document.createElement('div'); + nameEl.className = 'inline-viewer-audio-name'; + nameEl.textContent = file.name; + wrapper.appendChild(nameEl); - const icon = document.createElement('div'); - icon.className = 'inline-viewer-audio-icon'; - icon.innerHTML = ''; - wrapper.appendChild(icon); + const audio = document.createElement('audio'); + audio.className = 'inline-viewer-audio'; + audio.controls = true; + audio.preload = 'metadata'; + audio.src = streamUrl; + wrapper.appendChild(audio); - const nameEl = document.createElement('div'); - nameEl.className = 'inline-viewer-audio-name'; - nameEl.textContent = file.name; - wrapper.appendChild(nameEl); - - const audio = document.createElement('audio'); - audio.className = 'inline-viewer-audio'; - audio.controls = true; - audio.preload = 'metadata'; - audio.src = blobUrl; - wrapper.appendChild(audio); - - // Fallback message for unsupported codecs - audio.addEventListener('error', () => { - console.warn('Audio playback error — codec may not be supported'); - wrapper.innerHTML = ''; - const msg = document.createElement('div'); - msg.className = 'inline-viewer-message'; - msg.innerHTML = ` + // Fallback message for unsupported codecs / failed loads + audio.addEventListener('error', () => { + console.warn('Audio playback error — codec may not be supported'); + wrapper.innerHTML = ''; + const msg = document.createElement('div'); + msg.className = 'inline-viewer-message'; + msg.innerHTML = `

Your browser cannot play this audio format.

Click "Download" to save the file.

`; - wrapper.appendChild(msg); - }); + wrapper.appendChild(msg); + }); - container.appendChild(wrapper); - } else { - const video = document.createElement('video'); - video.className = 'inline-viewer-video'; - video.controls = true; - video.preload = 'metadata'; - video.src = blobUrl; - video.setAttribute('playsinline', 'true'); + container.appendChild(wrapper); + } else { + const video = document.createElement('video'); + video.className = 'inline-viewer-video'; + video.controls = true; + video.preload = 'metadata'; + video.src = streamUrl; + video.setAttribute('playsinline', 'true'); - // Fallback message for unsupported codecs - video.addEventListener('error', () => { - console.warn('Video playback error — codec may not be supported'); - if (video.parentNode) { - video.parentNode.removeChild(video); - } - const msg = document.createElement('div'); - msg.className = 'inline-viewer-message'; - msg.innerHTML = ` + // Fallback message for unsupported codecs / failed loads + video.addEventListener('error', () => { + console.warn('Video playback error — codec may not be supported'); + if (video.parentNode) { + video.parentNode.removeChild(video); + } + const msg = document.createElement('div'); + msg.className = 'inline-viewer-message'; + msg.innerHTML = `

Your browser cannot play this video format.

Click "Download" to save the file.

`; - container.appendChild(msg); - }); + container.appendChild(msg); + }); - container.appendChild(video); - } - - // Store blob URL for cleanup on close - this.currentBlobUrl = blobUrl; - } catch (error) { - console.error(`Error creating ${mediaType} viewer:`, error); - - if (loader?.parentNode) { - loader.parentNode.removeChild(loader); - } - - this.showErrorMessage(container); + container.appendChild(video); } } @@ -549,26 +531,12 @@ class InlineViewer { } /** - * + * Download the file via a browser-native download (streams to disk, + * nothing is buffered in page memory). * @param {FileItem} file */ downloadFile(file) { - fetch(`/api/files/${file.id}`, { credentials: 'same-origin' }) - .then((res) => { - if (!res.ok) throw new Error(`HTTP ${res.status}`); - return res.blob(); - }) - .then((blob) => { - const url = URL.createObjectURL(blob); - const link = document.createElement('a'); - link.href = url; - link.download = file.name; - document.body.appendChild(link); - link.click(); - document.body.removeChild(link); - URL.revokeObjectURL(url); - }) - .catch((err) => console.error('Download error:', err)); + triggerBrowserDownload(`/api/files/${file.id}`, file.name); } /** diff --git a/static/js/features/files/search.js b/static/js/features/files/search.js index 8db9fa5a..73d5bbac 100644 --- a/static/js/features/files/search.js +++ b/static/js/features/files/search.js @@ -7,7 +7,7 @@ * displays the enriched results returned by the server. */ -import { addItem as filesViewAddItem, loadFiles } from '../../app/filesView.js'; +import { loadFiles, renderItems } from '../../app/filesView.js'; import { app } from '../../app/state.js'; import { ui } from '../../app/ui.js'; import { getAuthHeaders } from './fileOperations.js'; @@ -198,15 +198,10 @@ const search = { return; } - // Render folders (server-provided enriched data) - results.folders.forEach((folder) => { - filesViewAddItem(folder); - }); - - // Render files (server-provided enriched data) - results.files.forEach((file) => { - filesViewAddItem(file); - }); + // Render the whole result set (server-provided enriched data) in one + // batched pass — per-item inserts would trigger a DOM scan, a smooth + // scroll and a highlight pulse for each of up to 100 rows. + renderItems([...results.folders, ...results.files]); }, /** diff --git a/static/js/features/library/photos.js b/static/js/features/library/photos.js index eff922ba..d7f6f6ee 100644 --- a/static/js/features/library/photos.js +++ b/static/js/features/library/photos.js @@ -472,22 +472,43 @@ const photosView = { if (bar_delete) { bar_delete.onclick = async () => { if (!confirm('Delete selected items?')) return; - for (const fid of this.selected) { - try { - await fetch(`/api/files/${fid}`, { - method: 'DELETE', + + // One batch request per chunk instead of one DELETE per photo. + // The photos view is files-only, so every id is a file id. + const ids = [...this.selected]; + const CHUNK_SIZE = 1000; // backend MAX_BATCH_SIZE + const trashed = new Set(); + + try { + for (let i = 0; i < ids.length; i += CHUNK_SIZE) { + const chunk = ids.slice(i, i + CHUNK_SIZE); + const response = await fetch('/api/batch/trash', { + method: 'POST', credentials: 'include', - headers: this._headers() + headers: this._headers(true), + body: JSON.stringify({ file_ids: chunk, folder_ids: [] }) }); - } catch (err) { - console.error('Delete failed:', fid, err); + // 200 = all trashed, 206 = partial; both carry `successful`. + if (!response.ok && response.status !== 206) { + console.error('Batch trash failed:', response.status); + continue; + } + const data = await response.json(); + const ok = Array.isArray(data?.successful) ? data.successful : chunk; + for (const id of ok) trashed.add(id); } + } catch (err) { + console.error('Batch trash error:', err); } - this.items = this.items.filter((f) => !this.selected.has(f.id)); - this.selected.clear(); - this._hideSelectionBar(); - this._renderedCount = 0; - this._renderFull(); + + if (trashed.size > 0) { + this.items = this.items.filter((f) => !trashed.has(f.id)); + for (const id of trashed) this.selected.delete(id); + this._renderedCount = 0; + this._renderFull(); + } + // Refresh (or hide) the bar to reflect any items left selected. + this._updateSelectionBar(); }; } diff --git a/static/js/model/filesModel.js b/static/js/model/filesModel.js index f483b5e9..9b6abf96 100644 --- a/static/js/model/filesModel.js +++ b/static/js/model/filesModel.js @@ -33,6 +33,48 @@ async function getFolder(id) { return Promise.reject(null); } +/** + * Session cache of folder metadata for breadcrumb resolution (`id → + * FolderItem`). Ancestors of the current folder have almost always been + * visited already, so a warm navigation rebuilds the whole crumb trail + * with zero fetches instead of one round-trip per depth level. + * Invalidated on rename/move via {@link invalidateFolderMeta}. + * @type {Map} + */ +const _folderMetaCache = new Map(); + +/** + * Resolve breadcrumb metadata for one folder, consulting the session + * cache first. + * @param {string} id + * @returns {Promise} + */ +async function _getFolderMeta(id) { + const cached = _folderMetaCache.get(id); + if (cached) return cached; + const info = await getFolder(id); + _folderMetaCache.set(id, info); + return info; +} + +/** + * Drop the cached breadcrumb metadata for a folder. Call after any + * operation that changes its name or parent (rename, move) so the next + * breadcrumb rebuild re-fetches the fresh row. + * @param {string} folderId + */ +function invalidateFolderMeta(folderId) { + _folderMetaCache.delete(folderId); +} + +/** + * Monotonic token identifying the most recent {@link rebuildBreadCrumb} + * call. Rebuilds now run concurrently with the listing fetch, so a rapid + * second navigation can supersede one still in flight — the superseded + * run must not commit its (stale) trail over the newer one. + */ +let _breadcrumbGeneration = 0; + /** * Walk up the folder hierarchy to rebuild `app.breadcrumbPath`. * @@ -41,29 +83,39 @@ async function getFolder(id) { * handles shared folders the user cannot traverse beyond. * * An error on the target folder itself is treated as a real error and falls - * back to the home folder. + * back to the home folder (resetting `app.currentPath`). * - * @returns {Promise} + * The trail is built locally and committed to `app` atomically at the end, + * and only when this call is still the most recent one — callers run this + * concurrently with the listing fetch. + * + * @returns {Promise} `true` when the trail was committed; `false` + * when this rebuild was superseded by a newer navigation. */ async function rebuildBreadCrumb() { + const generation = ++_breadcrumbGeneration; + /** @type {FolderItem|null} */ let currentFolderInfo = null; - app.breadcrumbPath = []; + /** @type {Array<{id: string, name: string}>} */ + const crumbs = []; /** @type {string|null} */ let id = app.currentPath; while (id !== null) { try { - const folderInfo = await getFolder(id); + const folderInfo = await _getFolderMeta(id); + if (generation !== _breadcrumbGeneration) return false; if (currentFolderInfo === null) currentFolderInfo = folderInfo; - app.breadcrumbPath.unshift({ id: folderInfo.id, name: folderInfo.name }); + crumbs.unshift({ id: folderInfo.id, name: folderInfo.name }); id = folderInfo.parent_id; } catch (_e) { + if (generation !== _breadcrumbGeneration) return false; if (currentFolderInfo === null) { console.warn(`Cannot access target folder ${app.currentPath}, falling back to home`); uiNotifications.show('error: folder not found or permission denied', 'the given folder is not available or you do not have sufficient rights'); - app.breadcrumbPath = []; + crumbs.length = 0; id = app.userHomeFolderId; if (id) app.currentPath = id; } else { @@ -73,7 +125,10 @@ async function rebuildBreadCrumb() { } } + if (generation !== _breadcrumbGeneration) return false; + app.breadcrumbPath = crumbs; app.currentFolderInfo = currentFolderInfo; + return true; } /** @@ -175,4 +230,4 @@ async function fetchResourcesPage(folderId, { cursor = null, orderBy = 'name', l return { items, nextCursor: data.next_cursor ?? null }; } -export { fetchListing, fetchResourcesPage, getFolder, rebuildBreadCrumb }; +export { fetchListing, fetchResourcesPage, getFolder, invalidateFolderMeta, rebuildBreadCrumb }; diff --git a/static/js/utils/download.js b/static/js/utils/download.js new file mode 100644 index 00000000..5fa7ad17 --- /dev/null +++ b/static/js/utils/download.js @@ -0,0 +1,44 @@ +// @ts-check + +/** + * Browser-native download helpers. + * + * Downloads are handed to the browser as same-origin navigations: the + * response streams straight to disk with the browser's own progress UI, + * and auth cookies travel automatically. Nothing is buffered in page + * memory — unlike the old `fetch → blob → objectURL` pattern, which + * materialized the entire payload in the tab's heap before the save + * dialog could even appear. + */ + +/** + * Trigger a browser-native download for a same-origin URL. + * + * @param {string} url - Same-origin URL of the resource to download + * @param {string} [filename] - Suggested file name. The server's + * `Content-Disposition` filename wins when present; an empty string + * keeps whatever the server (or URL) provides. + */ +export function triggerBrowserDownload(url, filename = '') { + const link = document.createElement('a'); + link.href = url; + link.download = filename; + document.body.appendChild(link); + link.click(); + document.body.removeChild(link); +} + +/** + * Build the GET URL for the batch ZIP download endpoint + * (`GET /api/batch/download` accepts comma-separated id lists). + * Shared by the batch toolbar download and the drag-out `DownloadURL` + * builder so both stay in sync with the endpoint's query contract. + * + * @param {string[]} fileIds + * @param {string[]} folderIds + * @returns {string} Root-relative URL (prepend `window.location.origin` + * when an absolute URL is required, e.g. for `DataTransfer.setData`). + */ +export function buildBatchDownloadUrl(fileIds, folderIds) { + return `/api/batch/download?file_ids=${fileIds.join(',')}&folder_ids=${folderIds.join(',')}`; +}