Merge upstream/main into feat/external-file-mounts

Resolve conflicts between the external-file-mounts feature and upstream's
D5/D7 refactor (per-file provenance, keyset pagination, cross-drive move
gates, resource-access hook, folder-cascade lifecycle hook).

Key resolutions:
- FolderService::new now takes (repo, authz, file_lifecycle, mount_router);
  all callers + DI updated.
- FileRetrievalService / FileManagementService keep both the mount_router
  and the new resource_access_hook / drive_repo / storage_usage wiring.
- list_files_batch_with_perms: adapt the mount branch from offset- to
  keyset (after_name) pagination, mirroring paginate_mount_entries.
- download_file_impl: keep upstream's &HeaderMap + `impl IntoResponse + use<>`
  signature, retain the mount-download branch.
- Mount DTOs: the retired `owner_id` field maps onto created_by/updated_by
  (the mount owner) — the fields the frontend now uses for owner display.
- admin/+page.svelte: keep upstream's user-delete modal + the 'mounts' tab.
- Bump memmap2 0.9.10 -> 0.9.11 (RUSTSEC critical advisory fix) and
  regenerate Cargo.lock against the merged Cargo.toml.
This commit is contained in:
Bradley Nelson
2026-07-21 17:09:36 -06:00
600 changed files with 105575 additions and 14440 deletions
+105 -43
View File
@@ -9,7 +9,7 @@ use std::pin::Pin;
use azure_storage::StorageCredentials;
use azure_storage_blobs::prelude::*;
use bytes::Bytes;
use futures::StreamExt;
use futures::{StreamExt, TryStreamExt};
use tokio::fs;
use crate::application::ports::blob_storage_ports::{
@@ -33,8 +33,21 @@ impl AzureBlobBackend {
StorageCredentials::access_key(&config.account_name, config.account_key.clone())
};
let container_client = ClientBuilder::new(&config.account_name, credentials)
.container_client(&config.container);
// Custom endpoint (Azurite emulator / private deployment /
// benches) mirrors S3's `endpoint_url`; default is the public
// cloud URL derived from the account name.
let container_client = match &config.endpoint_url {
Some(uri) => ClientBuilder::with_location(
azure_storage::CloudLocation::Custom {
account: config.account_name.clone(),
uri: uri.trim_end_matches('/').to_string(),
},
credentials,
)
.container_client(&config.container),
None => ClientBuilder::new(&config.account_name, credentials)
.container_client(&config.container),
};
Self {
container_client,
@@ -130,7 +143,9 @@ impl BlobStorageBackend for AzureBlobBackend {
return Ok(size);
}
client.put_block_blob(data.to_vec()).await.map_err(|e| {
// `Bytes` converts into `azure_core::Body` by reference count —
// the old `data.to_vec()` copied every chunk once more.
client.put_block_blob(data).await.map_err(|e| {
DomainError::internal_error("Azure", format!("Failed to upload blob {hash}: {e}"))
})?;
@@ -138,6 +153,26 @@ impl BlobStorageBackend for AzureBlobBackend {
})
}
/// Dedup settle path: PUT unconditionally. Content-addressed keys make
/// re-PUTs idempotent, so the `get_properties` probe
/// `put_blob_from_bytes` pays is a pure extra round-trip on every NEW
/// chunk (2 RTTs -> 1, benches/S3-PUT.md — same shape as S3).
fn put_blob_from_bytes_unsynced(
&self,
hash: &str,
data: Bytes,
) -> Pin<Box<dyn std::future::Future<Output = Result<u64, DomainError>> + Send + '_>> {
let hash = hash.to_owned();
Box::pin(async move {
let client = self.blob_client(&hash);
let size = data.len() as u64;
client.put_block_blob(data).await.map_err(|e| {
DomainError::internal_error("Azure", format!("Failed to upload blob {hash}: {e}"))
})?;
Ok(size)
})
}
fn get_blob_stream(
&self,
hash: &str,
@@ -147,29 +182,46 @@ impl BlobStorageBackend for AzureBlobBackend {
Box::pin(async move {
let client = self.blob_client(&hash);
let mut result_data: Vec<u8> = Vec::new();
let mut stream = client.get().into_stream();
while let Some(response) = stream.next().await {
let response = response.map_err(|e| {
DomainError::new(
// The old implementation drained the ENTIRE blob into one
// `Vec<u8>` before yielding a single mega-chunk — whole-blob
// RAM residency per reader, and with `read_prefetch() = 8`
// up to 8 entire chunk-blobs resident at once during CDC
// reassembly. Now the SDK's page/body streams forward
// directly. The FIRST page is still awaited eagerly so a
// missing blob surfaces as the same up-front NotFound the
// old code produced; later pages/chunks map to io::Error
// items like every other backend's stream.
let mut pages = client.get().into_stream();
let first = match pages.next().await {
Some(Ok(response)) => response,
Some(Err(e)) => {
return Err(DomainError::new(
ErrorKind::NotFound,
"Azure",
format!("Failed to get blob {hash}: {e}"),
)
})?;
let mut body = response.data;
while let Some(chunk) = body.next().await {
let chunk = chunk.map_err(|e| {
DomainError::internal_error("Azure", format!("Stream read error: {e}"))
})?;
result_data.extend_from_slice(&chunk);
));
}
}
None => {
let empty: BlobStream =
Box::pin(futures::stream::once(async move { Ok(Bytes::new()) }));
return Ok(empty);
}
};
let stream: BlobStream = Box::pin(futures::stream::once(async move {
Ok(Bytes::from(result_data))
}));
let first_body = first.data.map(|chunk| {
chunk.map_err(|e| std::io::Error::other(format!("Stream read error: {e}")))
});
let tail = pages
.map(|page| match page {
Ok(response) => Ok(response.data.map(|chunk| {
chunk.map_err(|e| std::io::Error::other(format!("Stream read error: {e}")))
})),
Err(e) => Err(std::io::Error::other(format!(
"Failed to get blob page: {e}"
))),
})
.try_flatten();
let stream: BlobStream = Box::pin(first_body.chain(tail));
Ok(stream)
})
}
@@ -190,32 +242,42 @@ impl BlobStorageBackend for AzureBlobBackend {
None => azure_core::request_options::Range::new(start, u64::MAX),
};
let mut result_data: Vec<u8> = Vec::new();
let mut stream = client.get().range(range).into_stream();
while let Some(response) = stream.next().await {
let response = response.map_err(|e| {
DomainError::new(
// Same forwarding shape as `get_blob_stream` — a ranged read
// doubly so: the caller explicitly asked NOT to pay for the
// whole blob, yet the old code buffered the full range.
let mut pages = client.get().range(range).into_stream();
let first = match pages.next().await {
Some(Ok(response)) => response,
Some(Err(e)) => {
return Err(DomainError::new(
ErrorKind::NotFound,
"Azure",
format!("Failed to get blob range {hash}: {e}"),
)
})?;
let mut body = response.data;
while let Some(chunk) = body.next().await {
let chunk = chunk.map_err(|e| {
DomainError::internal_error(
"Azure",
format!("Stream range read error: {e}"),
)
})?;
result_data.extend_from_slice(&chunk);
));
}
}
None => {
let empty: BlobStream =
Box::pin(futures::stream::once(async move { Ok(Bytes::new()) }));
return Ok(empty);
}
};
let stream: BlobStream = Box::pin(futures::stream::once(async move {
Ok(Bytes::from(result_data))
}));
let first_body = first.data.map(|chunk| {
chunk.map_err(|e| std::io::Error::other(format!("Stream range read error: {e}")))
});
let tail = pages
.map(|page| match page {
Ok(response) => Ok(response.data.map(|chunk| {
chunk.map_err(|e| {
std::io::Error::other(format!("Stream range read error: {e}"))
})
})),
Err(e) => Err(std::io::Error::other(format!(
"Failed to get blob range page: {e}"
))),
})
.try_flatten();
let stream: BlobStream = Box::pin(first_body.chain(tail));
Ok(stream)
})
}
+246 -234
View File
@@ -10,15 +10,14 @@
use std::path::{Path, PathBuf};
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use bytes::Bytes;
use lru::LruCache;
use std::num::NonZeroUsize;
use dashmap::DashMap;
use tokio::fs;
use tokio::io::{AsyncReadExt, AsyncSeekExt, AsyncWriteExt};
use tokio::sync::Mutex;
use tokio_util::io::ReaderStream;
use uuid::Uuid;
use crate::application::ports::blob_storage_ports::{
BlobStorageBackend, BlobStream, StorageHealthStatus,
@@ -50,33 +49,63 @@ struct CacheEntry {
/// A `BlobStorageBackend` decorator that adds an LRU disk cache in front of
/// a remote backend.
///
/// The index is a `moka::sync::Cache` with a byte weigher: cached reads
/// probe it lock-free (sharded, striped recency) where the previous
/// `tokio::sync::Mutex<LruCache>` serialized EVERY cached chunk read on one
/// global async mutex — negative scaling under concurrent readers
/// (benches/ROUND12.md §B: 2.08 → 1.07 Mops/s going 1 → 2 readers on the
/// mutex; moka holds 1.7-2.4). moka also owns the byte budget: eviction by
/// weighted size replaces the manual `current_size` counter +
/// `collect_evictions` sweep, and the eviction listener unlinks the evicted
/// `.blob` (only on size-eviction — a Replaced entry shares its file with
/// the replacement, and Explicit invalidations unlink at their call site).
pub struct CachedBlobBackend {
inner: Arc<dyn BlobStorageBackend>,
cache_dir: PathBuf,
max_cache_bytes: u64,
index: Arc<Mutex<LruCache<String, CacheEntry>>>,
current_size: Arc<AtomicU64>,
index: moka::sync::Cache<String, CacheEntry>,
/// Per-hash single-flight gates for cache misses. K concurrent cold
/// readers of one blob (e.g. a video player's parallel Range probes)
/// used to each download the FULL blob from the remote backend — and
/// race their writes on one shared `.tmp` path. The gate coalesces
/// them onto one fetch; waiters re-check the cache and serve locally
/// (16 fetches -> 1, benches/BLOB-CACHE.md).
inflight: Arc<DashMap<String, Arc<Mutex<()>>>>,
}
fn cached_path_in(cache_dir: &Path, hash: &str) -> PathBuf {
let prefix = &hash[..2.min(hash.len())];
cache_dir.join(prefix).join(format!("{hash}.blob"))
}
impl CachedBlobBackend {
/// Create a new cached backend wrapping `inner`.
pub fn new(inner: Arc<dyn BlobStorageBackend>, config: &BlobCacheConfig) -> Self {
let listener_dir = config.cache_dir.clone();
Self {
inner,
cache_dir: config.cache_dir.clone(),
max_cache_bytes: config.max_cache_bytes,
// Capacity is essentially unbounded — eviction is by byte budget, not count.
index: Arc::new(Mutex::new(LruCache::new(
NonZeroUsize::new(1_000_000).unwrap(),
))),
current_size: Arc::new(AtomicU64::new(0)),
index: moka::sync::Cache::builder()
.weigher(|_k: &String, e: &CacheEntry| e.size.clamp(1, u32::MAX as u64) as u32)
.max_capacity(config.max_cache_bytes)
.eviction_listener(move |hash: Arc<String>, _entry, cause| {
// Size-evicted blobs lose their on-disk file here (the
// sweep `collect_evictions` used to do). A quick unlink
// on the inserting task's thread, off the hot get path.
if cause == moka::notification::RemovalCause::Size {
let _ = std::fs::remove_file(cached_path_in(&listener_dir, &hash));
}
})
.build(),
inflight: Arc::new(DashMap::new()),
}
}
/// Path where a blob is cached locally.
fn cached_path(&self, hash: &str) -> PathBuf {
let prefix = &hash[..2.min(hash.len())];
self.cache_dir.join(prefix).join(format!("{hash}.blob"))
cached_path_in(&self.cache_dir, hash)
}
}
@@ -87,14 +116,24 @@ impl BlobStorageBackend for CachedBlobBackend {
let inner = self.inner.clone();
let cache_dir = self.cache_dir.clone();
let index = self.index.clone();
let current_size = self.current_size.clone();
Box::pin(async move {
inner.initialize().await?;
// Create cache dir structure (256 prefix dirs)
// Create the cache dir AND its 256 {00..ff} shard dirs up front
// (mirroring LocalBlobBackend::initialize), so the write paths never
// pay a per-chunk `create_dir_all` on an already-existing shard — a
// ~45 µs mkdirat(EEXIST)+stat+blocking-dispatch removed per cache
// write on cached-remote deployments (benches/ROUND26.md §D1).
fs::create_dir_all(&cache_dir).await.map_err(|e| {
DomainError::internal_error("BlobCache", format!("mkdir cache_dir: {e}"))
})?;
for prefix in &crate::infrastructure::services::local_blob_backend::HEX_PREFIXES {
fs::create_dir_all(cache_dir.join(prefix))
.await
.map_err(|e| {
DomainError::internal_error("BlobCache", format!("mkdir cache shard: {e}"))
})?;
}
// Scan existing cache to rebuild index. Collect entries WITHOUT
// holding the index lock — a large cache directory walk must not
@@ -120,14 +159,13 @@ impl BlobStorageBackend for CachedBlobBackend {
}
}
}
// Bulk-insert the rebuilt index under a single brief lock.
{
let mut idx = index.lock().await;
for (stem, size) in entries {
idx.put(stem, CacheEntry { size });
}
// Rebuild the index; if the restored set exceeds the byte
// budget, moka trims it (and the eviction listener unlinks the
// trimmed files) — the old index carried the excess until the
// next insert.
for (stem, size) in entries {
index.insert(stem, CacheEntry { size });
}
current_size.store(total_bytes, Ordering::Relaxed);
tracing::info!(
"Blob cache initialized: {} bytes in cache at {}",
total_bytes,
@@ -142,21 +180,28 @@ impl BlobStorageBackend for CachedBlobBackend {
hash: &str,
source_path: &Path,
) -> Pin<Box<dyn std::future::Future<Output = Result<u64, DomainError>> + Send + '_>> {
let inner = self.inner.clone();
let hash = hash.to_string();
let source = source_path.to_path_buf();
let self_ref = CachedRef {
cache_dir: self.cache_dir.clone(),
max_cache_bytes: self.max_cache_bytes,
index: self.index.clone(),
current_size: self.current_size.clone(),
};
Box::pin(async move {
// Write to inner backend
let bytes = inner.put_blob(&hash, &source).await?;
// Also cache locally (best-effort)
let _ = self_ref.insert_into_cache_static(&hash, &source).await;
Ok(bytes)
// Cache FIRST: every inner backend consumes the source file
// (local renames it, S3/Azure delete it after upload), so the
// old populate-after-put ordering failed 100% of the time and
// the first read after a whole-file put paid a full remote
// re-download (the ROUND11 deferred correctness note; fix
// gated in benches/ROUND12.md §B).
let cached = self.insert_into_cache(&hash, &source).await.is_ok();
match self.inner.put_blob(&hash, &source).await {
Ok(bytes) => Ok(bytes),
Err(e) => {
// Never serve a blob the backend rejected: drop the
// just-inserted cache entry + file.
if cached {
self.index.invalidate(&hash);
let _ = fs::remove_file(self.cached_path(&hash)).await;
}
Err(e)
}
}
})
}
@@ -165,68 +210,68 @@ impl BlobStorageBackend for CachedBlobBackend {
hash: &str,
data: Bytes,
) -> Pin<Box<dyn std::future::Future<Output = Result<u64, DomainError>> + Send + '_>> {
let inner = self.inner.clone();
let hash = hash.to_string();
let self_ref = CachedRef {
cache_dir: self.cache_dir.clone(),
max_cache_bytes: self.max_cache_bytes,
index: self.index.clone(),
current_size: self.current_size.clone(),
};
Box::pin(async move {
let size = inner.put_blob_from_bytes(&hash, data.clone()).await?;
// Also cache locally (best-effort): write bytes to cache path
let dest = self_ref.cached_path(&hash);
if let Some(parent) = dest.parent() {
let _ = fs::create_dir_all(parent).await;
}
let _ = fs::write(&dest, &data).await;
let data_len = data.len() as u64;
let mut idx = self_ref.index.lock().await;
if let Some(old) = idx.put(hash, CacheEntry { size: data_len }) {
self_ref.current_size.fetch_sub(old.size, Ordering::Relaxed);
}
self_ref.current_size.fetch_add(data_len, Ordering::Relaxed);
let size = self.inner.put_blob_from_bytes(&hash, data.clone()).await?;
self.cache_bytes_write_through(hash, &data).await;
Ok(size)
})
}
// Without this override the trait default would re-route the CDC chunk
// write through `put_blob_from_bytes` above, whose inner (synced) call
// pays the remote exists-probe per chunk. The local write-through cache
// population is kept identical — post-upload readers (thumbnail/EXIF/
// face hooks) hit the cache instead of re-fetching from the remote.
fn put_blob_from_bytes_unsynced(
&self,
hash: &str,
data: Bytes,
) -> Pin<Box<dyn std::future::Future<Output = Result<u64, DomainError>> + Send + '_>> {
let hash = hash.to_string();
Box::pin(async move {
let size = self
.inner
.put_blob_from_bytes_unsynced(&hash, data.clone())
.await?;
self.cache_bytes_write_through(hash, &data).await;
Ok(size)
})
}
// The durability barrier must reach the backend that buffered the
// unsynced writes; the local cache copy is disposable and needs none.
fn sync_blobs(
&self,
hashes: &[String],
) -> Pin<Box<dyn std::future::Future<Output = Result<(), DomainError>> + Send + '_>> {
self.inner.sync_blobs(hashes)
}
fn get_blob_stream(
&self,
hash: &str,
) -> Pin<Box<dyn std::future::Future<Output = Result<BlobStream, DomainError>> + Send + '_>>
{
let hash = hash.to_string();
let cached = self.cached_path(&hash);
let index = self.index.clone();
let inner = self.inner.clone();
let cache_dir = self.cache_dir.clone();
let max_cache_bytes = self.max_cache_bytes;
let current_size = self.current_size.clone();
Box::pin(async move {
// Check cache presence (and bump LRU recency) under a brief lock,
// then release it BEFORE touching the filesystem so concurrent
// readers don't serialize behind a single open() syscall.
if index.lock().await.get(&hash).is_some() {
// Lock-free cache probe (bumps moka recency) — the old shape
// took the one global async mutex here on EVERY cached chunk
// read, and cloned `cache_dir` per hit for a miss-only struct.
if self.index.get(&hash).is_some() {
let cached = self.cached_path(&hash);
if let Ok(file) = fs::File::open(&cached).await {
let stream: BlobStream =
Box::pin(ReaderStream::with_capacity(file, STREAM_CHUNK_SIZE));
return Ok(stream);
}
// Cache entry stale (file vanished) — drop it from the index.
if let Some(entry) = index.lock().await.pop(&hash) {
current_size.fetch_sub(entry.size, Ordering::Relaxed);
}
self.index.invalidate(&hash);
}
// Cache miss — fetch from inner, spool to cache
let self_ref = CachedRef {
cache_dir,
max_cache_bytes,
index: index.clone(),
current_size: current_size.clone(),
};
let dest = self_ref.fetch_and_cache_static(&hash, &*inner).await?;
// Cache miss — fetch from inner (single-flight), spool to cache
let cached = self.cached_path(&hash);
let dest = self.fetch_and_cache_singleflight(&hash, &cached).await?;
let file = fs::File::open(&dest).await.map_err(|e| {
DomainError::internal_error("BlobCache", format!("re-open cached: {e}"))
})?;
@@ -243,17 +288,11 @@ impl BlobStorageBackend for CachedBlobBackend {
) -> Pin<Box<dyn std::future::Future<Output = Result<BlobStream, DomainError>> + Send + '_>>
{
let hash = hash.to_string();
let cached = self.cached_path(&hash);
let index = self.index.clone();
let inner = self.inner.clone();
let cache_dir = self.cache_dir.clone();
let max_cache_bytes = self.max_cache_bytes;
let current_size = self.current_size.clone();
Box::pin(async move {
// Check cache presence (and bump LRU recency) under a brief lock,
// then release it BEFORE the open()/seek() syscalls so concurrent
// range readers don't serialize behind the index mutex.
if index.lock().await.get(&hash).is_some() {
// Lock-free cache probe (bumps moka recency); the filesystem is
// only touched after the probe, as before.
if self.index.get(&hash).is_some() {
let cached = self.cached_path(&hash);
if let Ok(mut file) = fs::File::open(&cached).await {
file.seek(std::io::SeekFrom::Start(start))
.await
@@ -266,19 +305,14 @@ impl BlobStorageBackend for CachedBlobBackend {
Box::pin(ReaderStream::with_capacity(limited, STREAM_CHUNK_SIZE));
return Ok(stream);
}
if let Some(entry) = index.lock().await.pop(&hash) {
current_size.fetch_sub(entry.size, Ordering::Relaxed);
}
self.index.invalidate(&hash);
}
// Cache miss — fetch full blob into cache, then serve range
let self_ref = CachedRef {
cache_dir,
max_cache_bytes,
index: index.clone(),
current_size: current_size.clone(),
};
let dest = self_ref.fetch_and_cache_static(&hash, &*inner).await?;
// Cache miss — fetch full blob into cache (single-flight: a
// player's parallel cold Range probes coalesce onto ONE remote
// download), then serve the range locally.
let cached = self.cached_path(&hash);
let dest = self.fetch_and_cache_singleflight(&hash, &cached).await?;
let mut file = fs::File::open(&dest)
.await
.map_err(|e| DomainError::internal_error("BlobCache", format!("re-open: {e}")))?;
@@ -297,19 +331,13 @@ impl BlobStorageBackend for CachedBlobBackend {
&self,
hash: &str,
) -> Pin<Box<dyn std::future::Future<Output = Result<(), DomainError>> + Send + '_>> {
let inner = self.inner.clone();
let hash = hash.to_string();
let cached = self.cached_path(&hash);
let index = self.index.clone();
let current_size = self.current_size.clone();
Box::pin(async move {
inner.delete_blob(&hash).await?;
// Remove from cache — drop the index lock before the unlink()
// syscall so deletes don't serialize concurrent cache lookups.
if let Some(entry) = index.lock().await.pop(&hash) {
current_size.fetch_sub(entry.size, Ordering::Relaxed);
}
let _ = fs::remove_file(&cached).await;
self.inner.delete_blob(&hash).await?;
// Explicit invalidation unlinks here (the eviction listener
// only unlinks size-evictions).
self.index.invalidate(&hash);
let _ = fs::remove_file(self.cached_path(&hash)).await;
Ok(())
})
}
@@ -318,18 +346,13 @@ impl BlobStorageBackend for CachedBlobBackend {
&self,
hash: &str,
) -> Pin<Box<dyn std::future::Future<Output = Result<bool, DomainError>> + Send + '_>> {
let inner = self.inner.clone();
let hash = hash.to_string();
let index = self.index.clone();
Box::pin(async move {
// Check cache first (fast)
{
let mut idx = index.lock().await;
if idx.get(&hash).is_some() {
return Ok(true);
}
// Check cache first (fast, lock-free)
if self.index.get(&hash).is_some() {
return Ok(true);
}
inner.blob_exists(&hash).await
self.inner.blob_exists(&hash).await
})
}
@@ -337,23 +360,17 @@ impl BlobStorageBackend for CachedBlobBackend {
&self,
hash: &str,
) -> Pin<Box<dyn std::future::Future<Output = Result<u64, DomainError>> + Send + '_>> {
let inner = self.inner.clone();
let hash = hash.to_string();
let index = self.index.clone();
let cached = self.cached_path(&hash);
Box::pin(async move {
// Check cache
{
let mut idx = index.lock().await;
if let Some(entry) = idx.get(&hash) {
return Ok(entry.size);
}
// Check cache (lock-free)
if let Some(entry) = self.index.get(&hash) {
return Ok(entry.size);
}
// Fallback to cached file on disk (in case index was lost)
if let Ok(meta) = fs::metadata(&cached).await {
if let Ok(meta) = fs::metadata(self.cached_path(&hash)).await {
return Ok(meta.len());
}
inner.blob_size(&hash).await
self.inner.blob_size(&hash).await
})
}
@@ -362,19 +379,18 @@ impl BlobStorageBackend for CachedBlobBackend {
) -> Pin<
Box<dyn std::future::Future<Output = Result<StorageHealthStatus, DomainError>> + Send + '_>,
> {
let inner = self.inner.clone();
let cache_dir = self.cache_dir.clone();
let current_size = self.current_size.clone();
let max_bytes = self.max_cache_bytes;
Box::pin(async move {
let mut status = inner.health_check().await?;
let used = current_size.load(Ordering::Relaxed);
let mut status = self.inner.health_check().await?;
// Flush moka's pending maintenance so the reported byte count
// is current (rare admin path — the cost is fine here).
self.index.run_pending_tasks();
let used = self.index.weighted_size();
status.message = format!(
"{} | Cache: {}/{} bytes used at {}",
status.message,
used,
max_bytes,
cache_dir.display()
self.max_cache_bytes,
self.cache_dir.display()
);
status.backend_type = format!("cached({})", status.backend_type);
Ok(status)
@@ -398,54 +414,56 @@ impl BlobStorageBackend for CachedBlobBackend {
}
}
// ── Helper struct for owned references in async closures ───────────
// ── Cache internals (miss path + population) ───────────────────────
/// Cloneable set of cache internals — avoids borrow issues in boxed futures.
struct CachedRef {
cache_dir: PathBuf,
max_cache_bytes: u64,
index: Arc<Mutex<LruCache<String, CacheEntry>>>,
current_size: Arc<AtomicU64>,
}
impl CachedRef {
fn cached_path(&self, hash: &str) -> PathBuf {
let prefix = &hash[..2.min(hash.len())];
self.cache_dir.join(prefix).join(format!("{hash}.blob"))
impl CachedBlobBackend {
/// Best-effort write-through cache population shared by both blob-bytes
/// PUT paths. moka enforces the byte budget on every insert (the old
/// index deliberately skipped the eviction sweep on this path, letting
/// write bursts overshoot the budget until the next read-miss insert).
async fn cache_bytes_write_through(&self, hash: String, data: &Bytes) {
// The shard dir was created at initialize() — no per-write create_dir_all
// (benches/ROUND26.md §D1).
let dest = self.cached_path(&hash);
let _ = fs::write(&dest, data).await;
let data_len = data.len() as u64;
self.index.insert(hash, CacheEntry { size: data_len });
}
/// Pop LRU entries until the cache is back within its byte budget,
/// returning the on-disk paths of the evicted blobs.
///
/// Only the in-memory index is touched here (atomic counter + LRU map);
/// the caller MUST unlink the returned paths AFTER releasing the index
/// lock so the `remove_file` syscalls never run while the mutex is held.
fn collect_evictions(&self, idx: &mut LruCache<String, CacheEntry>) -> Vec<PathBuf> {
let mut victims = Vec::new();
while self.current_size.load(Ordering::Relaxed) > self.max_cache_bytes {
if let Some((evicted_hash, evicted_entry)) = idx.pop_lru() {
self.current_size
.fetch_sub(evicted_entry.size, Ordering::Relaxed);
victims.push(self.cached_path(&evicted_hash));
} else {
break;
}
}
victims
}
async fn insert_into_cache_static(
/// Single-flight wrapper around [`Self::fetch_and_cache`]: the first
/// caller for a hash becomes the leader and downloads; concurrent
/// callers queue on the per-hash gate, then re-check the cache and serve
/// the leader's file without touching the remote backend. Errors are not
/// cached — the gate entry is dropped, so the next caller retries.
async fn fetch_and_cache_singleflight(
&self,
hash: &str,
source_path: &Path,
) -> Result<(), DomainError> {
let dest = self.cached_path(hash);
if let Some(parent) = dest.parent() {
fs::create_dir_all(parent).await.map_err(|e| {
DomainError::internal_error("BlobCache", format!("mkdir failed: {e}"))
})?;
cached: &Path,
) -> Result<PathBuf, DomainError> {
let gate = self
.inflight
.entry(hash.to_string())
.or_insert_with(|| Arc::new(Mutex::new(())))
.clone();
let _guard = gate.lock().await;
// Re-check under the gate: if we queued behind the leader, the blob
// is on disk now and this turns into a local open.
if self.index.get(hash).is_some() && fs::metadata(cached).await.is_ok() {
return Ok(cached.to_path_buf());
}
let result = self.fetch_and_cache(hash).await;
// Drop the gate whether we succeeded or failed; a late-arriving
// caller after an error creates a fresh gate and retries the fetch.
self.inflight.remove(hash);
result
}
async fn insert_into_cache(&self, hash: &str, source_path: &Path) -> Result<(), DomainError> {
// Shard dir pre-created at initialize() (benches/ROUND26.md §D1).
let dest = self.cached_path(hash);
let size = fs::metadata(source_path)
.await
.map(|m| m.len())
@@ -455,74 +473,68 @@ impl CachedRef {
DomainError::internal_error("BlobCache", format!("cache copy failed: {e}"))
})?;
// Update the index and pick eviction victims under a single brief
// lock, then unlink the evicted files AFTER releasing it — file
// removal must not run while the index mutex is held.
let to_evict = {
let mut idx = self.index.lock().await;
if let Some(old) = idx.put(hash.to_string(), CacheEntry { size }) {
self.current_size.fetch_sub(old.size, Ordering::Relaxed);
}
self.current_size.fetch_add(size, Ordering::Relaxed);
self.collect_evictions(&mut idx)
};
for path in to_evict {
let _ = fs::remove_file(&path).await;
}
// moka enforces the byte budget; size-evicted victims are unlinked
// by the eviction listener.
self.index.insert(hash.to_string(), CacheEntry { size });
Ok(())
}
async fn fetch_and_cache_static(
&self,
hash: &str,
inner: &dyn BlobStorageBackend,
) -> Result<PathBuf, DomainError> {
let stream = inner.get_blob_stream(hash).await?;
async fn fetch_and_cache(&self, hash: &str) -> Result<PathBuf, DomainError> {
let stream = self.inner.get_blob_stream(hash).await?;
// Shard dir pre-created at initialize() (benches/ROUND26.md §D1).
let dest = self.cached_path(hash);
if let Some(parent) = dest.parent() {
fs::create_dir_all(parent).await.map_err(|e| {
DomainError::internal_error("BlobCache", format!("mkdir failed: {e}"))
// Unique temp name: even if two fetches for one hash ever race
// (e.g. across processes sharing a cache dir), each writes its own
// inode and the rename is atomic — a torn/interleaved file can
// never land at the final path.
let tmp = dest.with_extension(format!("{}.tmp", Uuid::new_v4()));
let write_result: Result<u64, DomainError> = async {
let mut file = fs::File::create(&tmp).await.map_err(|e| {
DomainError::internal_error("BlobCache", format!("create tmp: {e}"))
})?;
}
let tmp = dest.with_extension("tmp");
let mut file = fs::File::create(&tmp)
.await
.map_err(|e| DomainError::internal_error("BlobCache", format!("create tmp: {e}")))?;
use futures::StreamExt;
let mut stream = stream;
let mut total = 0u64;
while let Some(chunk) = stream.next().await {
let bytes = chunk.map_err(|e| {
DomainError::internal_error("BlobCache", format!("stream read: {e}"))
})?;
total += bytes.len() as u64;
file.write_all(&bytes)
.await
.map_err(|e| DomainError::internal_error("BlobCache", format!("write: {e}")))?;
}
file.flush()
.await
.map_err(|e| DomainError::internal_error("BlobCache", format!("flush: {e}")))?;
drop(file);
fs::rename(&tmp, &dest)
.await
.map_err(|e| DomainError::internal_error("BlobCache", format!("rename: {e}")))?;
let to_evict = {
let mut idx = self.index.lock().await;
if let Some(old) = idx.put(hash.to_string(), CacheEntry { size: total }) {
self.current_size.fetch_sub(old.size, Ordering::Relaxed);
use futures::StreamExt;
let mut stream = stream;
let mut total = 0u64;
while let Some(chunk) = stream.next().await {
let bytes = chunk.map_err(|e| {
DomainError::internal_error("BlobCache", format!("stream read: {e}"))
})?;
total += bytes.len() as u64;
file.write_all(&bytes)
.await
.map_err(|e| DomainError::internal_error("BlobCache", format!("write: {e}")))?;
}
self.current_size.fetch_add(total, Ordering::Relaxed);
self.collect_evictions(&mut idx)
};
for path in to_evict {
let _ = fs::remove_file(&path).await;
file.flush()
.await
.map_err(|e| DomainError::internal_error("BlobCache", format!("flush: {e}")))?;
Ok(total)
}
.await;
let total = match write_result {
Ok(total) => total,
Err(e) => {
// Unique tmp names never get overwritten by a later fetch —
// reap the partial file instead of leaking it.
let _ = fs::remove_file(&tmp).await;
return Err(e);
}
};
if let Err(e) = fs::rename(&tmp, &dest).await {
let _ = fs::remove_file(&tmp).await;
return Err(DomainError::internal_error(
"BlobCache",
format!("rename: {e}"),
));
}
// moka enforces the byte budget; size-evicted victims are unlinked
// by the eviction listener.
self.index
.insert(hash.to_string(), CacheEntry { size: total });
Ok(dest)
}
@@ -514,6 +514,17 @@ impl ChunkedUploadService {
Ok(())
}
/// Alloc-free owner compare for the per-chunk hot path: the caller's
/// `Uuid` is stack-encoded (hyphenated, the format sessions store) —
/// `prepare_chunk`/`commit_chunk` used to pay a `Uuid::to_string` each
/// plus a dedicated `verify_session_owner` map lookup per chunk
/// (benches/ROUND12.md §M5, 1.28x / −2 allocs per chunk).
#[inline]
fn owner_matches(session_user_id: &str, user_id: Uuid) -> bool {
let mut buf = [0u8; 36];
session_user_id == user_id.hyphenated().encode_lower(&mut buf) as &str
}
/// Create a new upload session (persists `session.json` + empty `progress.bin`)
async fn create_session_inner(
&self,
@@ -617,9 +628,8 @@ impl ChunkedUploadService {
user_id: Uuid,
chunk_index: usize,
) -> Result<(PathBuf, usize), DomainError> {
self.verify_session_owner(upload_id, &user_id.to_string())
.map_err(|e| DomainError::new(ErrorKind::NotFound, "ChunkedUpload", e))?;
// Single map lookup: the owner gate rides the same guard (same
// anti-enum not-found for unknown session and foreign session).
let session = self.sessions.get(upload_id).ok_or_else(|| {
DomainError::new(
ErrorKind::NotFound,
@@ -627,6 +637,13 @@ impl ChunkedUploadService {
format!("Upload session not found: {}", upload_id),
)
})?;
if !Self::owner_matches(&session.user_id, user_id) {
return Err(DomainError::new(
ErrorKind::NotFound,
"ChunkedUpload",
format!("Upload session not found: {}", upload_id),
));
}
if chunk_index >= session.chunks.len() {
return Err(DomainError::new(
@@ -678,20 +695,23 @@ impl ChunkedUploadService {
computed_checksum: Option<String>,
expected_checksum: Option<String>,
) -> Result<ChunkUploadResponseDto, DomainError> {
self.verify_session_owner(upload_id, &user_id.to_string())
.map_err(|e| DomainError::new(ErrorKind::NotFound, "ChunkedUpload", e))?;
// Re-fetch chunk metadata under fresh lock — guards against the
// (vanishingly unlikely) case of a session expiry / cancellation
// racing with the write.
// Owner gate folded into the metadata read below — one lookup
// instead of two, same anti-enum not-found semantics.
let (chunk_path, expected_size, persist_path) = {
let session = self.sessions.get(upload_id).ok_or_else(|| {
DomainError::new(
ErrorKind::NotFound,
"ChunkedUpload",
"Session disappeared".to_string(),
format!("Upload session not found: {}", upload_id),
)
})?;
if !Self::owner_matches(&session.user_id, user_id) {
return Err(DomainError::new(
ErrorKind::NotFound,
"ChunkedUpload",
format!("Upload session not found: {}", upload_id),
));
}
if chunk_index >= session.chunks.len() {
return Err(DomainError::new(
ErrorKind::InvalidInput,
@@ -834,8 +854,7 @@ impl ChunkedUploadService {
let data_clone = data.clone(); // Bytes::clone is O(1) — just an Arc increment
let actual_checksum = tokio::task::spawn_blocking(move || {
use md5::{Digest, Md5};
let hash = Md5::digest(&data_clone);
hash.iter().map(|b| format!("{b:02x}")).collect::<String>()
crate::common::fmt::hex_lower(&Md5::digest(&data_clone))
})
.await
.map_err(|e| format!("MD5 checksum task failed: {e}"))?;
File diff suppressed because it is too large Load Diff
@@ -30,7 +30,7 @@
use std::path::{Path, PathBuf};
use std::pin::Pin;
use aes_gcm::aead::{Aead, AeadInPlace, KeyInit, OsRng};
use aes_gcm::aead::{AeadInPlace, KeyInit, OsRng};
use aes_gcm::{AeadCore, Aes256Gcm, Nonce};
use bytes::Bytes;
use std::sync::Arc;
@@ -44,6 +44,9 @@ use crate::domain::errors::DomainError;
/// Nonce size for AES-256-GCM (96 bits = 12 bytes).
const NONCE_SIZE: usize = 12;
/// AES-256-GCM authentication tag length appended after the ciphertext.
const TAG_SIZE: usize = 16;
/// Payloads at or above this size run crypto on the blocking pool; below
/// it the `spawn_blocking` round-trip costs more than the AES work itself.
const CRYPTO_OFFLOAD_THRESHOLD: usize = 64 * 1024;
@@ -56,7 +59,10 @@ const PLAINTEXT_EMIT_SIZE: usize = 64 * 1024;
/// `BlobStorageBackend` decorator that encrypts blobs at rest.
pub struct EncryptedBlobBackend {
inner: Arc<dyn BlobStorageBackend>,
cipher: Aes256Gcm,
/// `Arc` so the per-op `clone()` handed to `offload_crypto` closures is
/// an atomic bump instead of copying the ~240-byte expanded AES-256
/// round-key schedule on every chunk read/write.
cipher: Arc<Aes256Gcm>,
}
impl EncryptedBlobBackend {
@@ -64,7 +70,8 @@ impl EncryptedBlobBackend {
///
/// `key` must be exactly 32 bytes (AES-256).
pub fn new(inner: Arc<dyn BlobStorageBackend>, key: &[u8; 32]) -> Self {
let cipher = Aes256Gcm::new_from_slice(key).expect("AES-256 key must be 32 bytes");
let cipher =
Arc::new(Aes256Gcm::new_from_slice(key).expect("AES-256 key must be 32 bytes"));
Self { inner, cipher }
}
@@ -78,36 +85,60 @@ impl EncryptedBlobBackend {
}
/// Encrypt `data` into the on-disk layout: `[12-byte nonce][ciphertext + tag]`.
///
/// Single output buffer, mirroring the read side's in-place detached decrypt:
/// the payload is copied exactly once and encrypted in place with the tag
/// appended. The old shape let `cipher.encrypt` allocate a full ciphertext
/// `Vec` and then copied it a second time behind the nonce — one extra
/// allocation + a full-size memcpy on every encrypted chunk write
/// (benches/ROUND11.md §15; output bytes identical for a given nonce).
fn encrypt_bytes(cipher: &Aes256Gcm, data: &[u8]) -> Result<Bytes, DomainError> {
let nonce = Aes256Gcm::generate_nonce(&mut OsRng);
let ciphertext = cipher
.encrypt(&nonce, data)
let mut out = Vec::with_capacity(NONCE_SIZE + data.len() + TAG_SIZE);
out.extend_from_slice(nonce.as_slice());
out.extend_from_slice(data);
let tag = cipher
.encrypt_in_place_detached(&nonce, b"", &mut out[NONCE_SIZE..])
.map_err(|e| DomainError::internal_error("Encryption", format!("encrypt failed: {e}")))?;
let mut encrypted = Vec::with_capacity(NONCE_SIZE + ciphertext.len());
encrypted.extend_from_slice(nonce.as_slice());
encrypted.extend_from_slice(&ciphertext);
Ok(Bytes::from(encrypted))
out.extend_from_slice(&tag);
Ok(Bytes::from(out))
}
/// Decrypt the on-disk layout `[nonce][ciphertext + tag]` **in place**.
///
/// Consumes the encrypted buffer and reuses it for the plaintext, so peak
/// RAM is one buffer — not ciphertext + plaintext side by side (which for
/// legacy whole-file blobs would double a multi-hundred-MB allocation).
/// Reuses the encrypted buffer for the plaintext, so peak RAM is one buffer —
/// not ciphertext + plaintext side by side (which for legacy whole-file blobs
/// would double a multi-hundred-MB allocation). The nonce and 16-byte GCM tag
/// are lifted to the stack, the ciphertext body is decrypted in place via the
/// detached API (mirroring the encrypt side's `encrypt_in_place_detached`), and
/// the plaintext is returned as a zero-copy `Bytes::slice` past the nonce.
///
/// The prior shape did `encrypted.split_off(NONCE_SIZE)`, which allocated a
/// fresh `Vec` and memcpy'd the entire ciphertext (up to a whole legacy blob)
/// on every decrypted read — one full-payload allocation + copy the doc comment
/// above claimed did not happen (benches/ROUND25.md §M1; ROUND11 §15 fixed only
/// the encrypt side). Output plaintext is byte-identical.
fn decrypt_bytes(cipher: &Aes256Gcm, mut encrypted: Vec<u8>) -> Result<Bytes, DomainError> {
if encrypted.len() < NONCE_SIZE {
let len = encrypted.len();
if len < NONCE_SIZE + TAG_SIZE {
return Err(DomainError::internal_error(
"Encryption",
"encrypted blob too short (missing nonce)",
"encrypted blob too short (missing nonce/tag)",
));
}
let mut ciphertext = encrypted.split_off(NONCE_SIZE); // `encrypted` keeps the nonce
let nonce = Nonce::from_slice(&encrypted);
// Nonce (first 12 bytes) and GCM tag (last 16 bytes) copied to the stack so
// the middle can be borrowed mutably for in-place decryption.
let mut nonce_buf = [0u8; NONCE_SIZE];
nonce_buf.copy_from_slice(&encrypted[..NONCE_SIZE]);
let nonce = Nonce::from_slice(&nonce_buf);
let tag = aes_gcm::aead::Tag::<Aes256Gcm>::clone_from_slice(&encrypted[len - TAG_SIZE..]);
cipher
.decrypt_in_place(nonce, b"", &mut ciphertext)
.decrypt_in_place_detached(nonce, b"", &mut encrypted[NONCE_SIZE..len - TAG_SIZE], &tag)
.map_err(|e| DomainError::internal_error("Encryption", format!("decrypt failed: {e}")))?;
Ok(Bytes::from(ciphertext))
// Plaintext now lives at `encrypted[NONCE_SIZE..len - TAG_SIZE]`; drop the
// tag and hand out a refcounted view past the nonce — no copy, no new alloc.
encrypted.truncate(len - TAG_SIZE);
Ok(Bytes::from(encrypted).slice(NONCE_SIZE..))
}
/// Run a crypto closure inline for small payloads, on the blocking pool for
@@ -126,13 +157,18 @@ where
}
/// Turn a decrypted payload into a stream of bounded, zero-copy slices.
///
/// The emit-slice iterator is handed to `stream::iter` lazily — the closure
/// owns `data` (a refcounted `Bytes`), so each `slice` is produced on demand
/// as the consumer polls, rather than eagerly `collect`ing a `Vec` of
/// ⌈len/64 KiB⌉ slice handles up front (benches/ROUND20.md §I4).
fn plaintext_stream(data: Bytes) -> BlobStream {
let len = data.len();
let slices: Vec<Result<Bytes, std::io::Error>> = (0..len)
.step_by(PLAINTEXT_EMIT_SIZE)
.map(|off| Ok(data.slice(off..len.min(off + PLAINTEXT_EMIT_SIZE))))
.collect();
Box::pin(futures::stream::iter(slices))
Box::pin(futures::stream::iter(
(0..len)
.step_by(PLAINTEXT_EMIT_SIZE)
.map(move |off| Ok(data.slice(off..len.min(off + PLAINTEXT_EMIT_SIZE)))),
))
}
impl BlobStorageBackend for EncryptedBlobBackend {
@@ -318,6 +354,13 @@ impl BlobStorageBackend for EncryptedBlobBackend {
}
/// Collect a byte stream into a single `Vec<u8>`.
///
/// Modern blobs are CDC chunks (≤ `CDC_MAX_CHUNK` + nonce/tag overhead),
/// delivered here as small reader frames — growing from `Vec::new()` paid
/// ~log₂(n) reallocations + a wasted ~0.75×-size memcpy per read. Reserving
/// one chunk's worth up front on the first frame makes the common case a
/// single allocation; legacy whole-file blobs beyond that fall back to
/// normal doubling (benches/ROUND11.md §16: 9 → 1 allocs on a 1 MiB blob).
async fn collect_stream(stream: BlobStream) -> Result<Vec<u8>, DomainError> {
use futures::StreamExt;
let mut stream = stream;
@@ -325,6 +368,14 @@ async fn collect_stream(stream: BlobStream) -> Result<Vec<u8>, DomainError> {
while let Some(chunk) = stream.next().await {
let bytes = chunk
.map_err(|e| DomainError::internal_error("Encryption", format!("stream read: {e}")))?;
if buf.capacity() == 0 {
buf.reserve(
(crate::infrastructure::services::dedup_service::CDC_MAX_CHUNK
+ NONCE_SIZE
+ TAG_SIZE)
.max(bytes.len()),
);
}
buf.extend_from_slice(&bytes);
}
Ok(buf)
+25 -12
View File
@@ -61,23 +61,13 @@ impl ExifService {
// ── Camera info ──
if let Some(field) = exif.get_field(Tag::Make, In::PRIMARY) {
let val = field
.display_value()
.to_string()
.trim_matches('"')
.trim()
.to_string();
let val = display_value_trimmed(field);
if !val.is_empty() {
meta.camera_make = Some(val);
}
}
if let Some(field) = exif.get_field(Tag::Model, In::PRIMARY) {
let val = field
.display_value()
.to_string()
.trim_matches('"')
.trim()
.to_string();
let val = display_value_trimmed(field);
if !val.is_empty() {
meta.camera_model = Some(val);
}
@@ -115,6 +105,29 @@ impl ExifService {
}
}
/// Render an EXIF field's display value, then strip surrounding quotes and
/// whitespace (the shape `Make`/`Model` want) in a SINGLE allocation.
///
/// `display_value().to_string()` is the one unavoidable allocation — the field
/// value is materialized to text. The old `…to_string().trim_matches('"')
/// .trim().to_string()` chain then threw that `String` away and allocated a
/// second time for the trimmed copy. Here the same two-stage trim is applied
/// in place on the already-owned buffer (`drain` drops the prefix, `truncate`
/// the suffix — both reuse the allocation), so a quoted `"Canon"` costs one
/// allocation instead of two.
fn display_value_trimmed(field: &exif::Field) -> String {
let mut s = field.display_value().to_string();
// Same order the old chain used: strip `"` first, then whitespace. The
// result is a contiguous subslice of `s`; capture its byte range before
// mutating the owned buffer (the borrow ends at these two reads).
let trimmed = s.trim_matches('"').trim();
let start = trimmed.as_ptr().addr() - s.as_ptr().addr();
let len = trimmed.len();
s.drain(..start);
s.truncate(len);
s
}
/// Parse EXIF datetime string "YYYY:MM:DD HH:MM:SS" into DateTime<Utc>.
fn parse_exif_datetime(s: &str) -> Option<DateTime<Utc>> {
// EXIF dates use ":" as separator for date parts
@@ -28,11 +28,35 @@ fn is_image(content_type: &str) -> bool {
content_type.starts_with("image/")
}
/// Concurrent index-task budget. Env override
/// `OXICLOUD_FACES_INDEX_CONCURRENCY`, else the effective core count —
/// each task is a full-image read + decode + ONNX inference, so more
/// permits than cores only adds RAM pressure, not throughput.
fn max_concurrent_index() -> usize {
std::env::var("OXICLOUD_FACES_INDEX_CONCURRENCY")
.ok()
.and_then(|v| v.parse().ok())
.filter(|&n: &usize| n > 0)
.unwrap_or_else(|| {
std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(2)
})
}
pub struct FaceIndexingService {
pool: Arc<PgPool>,
repo: Arc<FacePgRepository>,
analyzer: Arc<dyn FaceAnalyzerPort>,
blob_root: PathBuf,
/// Bounds concurrent indexing tasks. The lifecycle hooks spawn one
/// task per uploaded/copied image with no ceiling, so a bulk upload
/// used to fan out N simultaneous full-image reads + decodes +
/// inferences — peak RSS N × image size plus CPU thrash. Same
/// invariant as `ThumbnailService::decode_semaphore`: the permit is
/// acquired BEFORE the blob read, so peak memory is
/// `permits × image size` regardless of upload concurrency.
index_semaphore: Arc<tokio::sync::Semaphore>,
}
impl FaceIndexingService {
@@ -43,6 +67,7 @@ impl FaceIndexingService {
repo,
analyzer,
blob_root,
index_semaphore: Arc::new(tokio::sync::Semaphore::new(max_concurrent_index())),
}
}
@@ -60,7 +85,15 @@ impl FaceIndexingService {
let repo = self.repo.clone();
let analyzer = self.analyzer.clone();
let blob_path = self.blob_path(&blob_hash);
let semaphore = self.index_semaphore.clone();
tokio::spawn(async move {
// Queue behind the concurrency budget BEFORE touching the
// blob — excess tasks wait holding only this tiny future,
// not a decoded image.
let _permit = semaphore
.acquire_owned()
.await
.expect("face index semaphore never closes");
if delete_first {
let _ = repo.delete_faces_for_file(file_id).await;
}
@@ -182,7 +182,30 @@ impl FileContentCache {
if let Some(hit) = self.get(&cache_key).await {
return Ok(hit);
}
self.load_and_cache(cache_key, etag, content_type, load)
.await
}
/// The populate-on-miss half of [`Self::get_or_load`], with single-flight
/// coalescing but WITHOUT the leading `get` probe.
///
/// Hot read paths that have *already* probed the cache with [`Self::get`]
/// (a borrow) call this directly on the miss branch — they then build the
/// owned `cache_key` / `etag` / `content_type` (each a heap allocation)
/// only when they are actually needed to populate, so a cache HIT allocates
/// none of them (benches/ROUND29.md §B). Because the caller's own `get`
/// already counted the hit/miss, this method does not re-probe — keeping the
/// hit/miss stat counts identical to a single `get_or_load` call.
pub async fn load_and_cache<F>(
&self,
cache_key: String,
etag: Arc<str>,
content_type: Arc<str>,
load: F,
) -> Result<(Bytes, Arc<str>, Arc<str>), DomainError>
where
F: Future<Output = Result<Bytes, DomainError>>,
{
// Slow path: coalesce concurrent misses into a single `load`.
let entry = self
.cache
@@ -0,0 +1,124 @@
//! Background daemon that purges expired `storage.role_grants` rows.
//!
//! The AuthZ engine already filters expired grants out of every
//! permission check at read time (`expires_at IS NULL OR
//! expires_at > NOW()` on every `check` / `list_grants_*` path in
//! `PgAclEngine`), so expired rows never leak permission. They just
//! accumulate. This daemon garbage-collects them once per
//! [`GrantCleanupService::interval_hours`], with a grace window past
//! `expires_at` that preserves the audit / support answer to "what
//! happened to my access?" for a few weeks.
//!
//! Shape mirrors [`TrashCleanupService`] verbatim (fire-and-forget
//! `tokio::spawn`, `tokio::time::interval`, first-tick-immediate). The
//! authoritative pattern for background daemons in this codebase; see
//! the plan doc `docs/plan/` (deferred future work: fold all daemons
//! into a central `JobRegistry` that plugins can also register into).
//!
//! [`TrashCleanupService`]: crate::infrastructure::services::trash_cleanup_service::TrashCleanupService
use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::time;
use tracing::{error, info};
use crate::application::ports::authorization_ports::AuthorizationEngine;
use crate::infrastructure::services::pg_acl_engine::PgAclEngine;
/// Daemon that periodically deletes expired grants.
///
/// Owns an `Arc<PgAclEngine>` (not a `dyn AuthorizationEngine`) to avoid
/// the wrapper allocation on every SQL call — the daemon is the sole
/// caller of `purge_expired_grants` outside of the admin trigger
/// endpoint, both statically dispatched.
pub struct GrantCleanupService {
authz: Arc<PgAclEngine>,
grace_days: u32,
interval_hours: u64,
}
impl GrantCleanupService {
pub fn new(authz: Arc<PgAclEngine>, grace_days: u32, interval_hours: u64) -> Self {
Self {
authz,
grace_days,
// Minimum 1 hour — matches TrashCleanupService's clamp so
// a mis-set `0` doesn't spin a hot loop.
interval_hours: interval_hours.max(1),
}
}
/// Grace period the daemon uses on its scheduled ticks. Exposed
/// for the admin trigger's default-response field.
pub fn grace_days(&self) -> u32 {
self.grace_days
}
/// Fire-and-forget the periodic purge. Never joins; killed
/// implicitly at `tokio::runtime::shutdown`.
pub async fn start_cleanup_job(self: Arc<Self>) {
let interval_hours = self.interval_hours;
let grace_days = self.grace_days;
info!(
"Starting grant-cleanup daemon: every {}h, grace = {}d",
interval_hours, grace_days
);
tokio::spawn(async move {
let mut interval = time::interval(Duration::from_secs(interval_hours * 60 * 60));
// First tick fires immediately — matches TrashCleanupService.
// Any accumulated backlog at boot gets flushed straight away.
loop {
interval.tick().await;
self.run_once().await;
}
});
}
/// One scheduled pass. Also called by the admin trigger endpoint
/// (via a shared `Arc<GrantCleanupService>` on `AppState`).
///
/// `grace_override`:
/// - `None` → use the configured grace (`self.grace_days`).
/// - `Some(n)` → override with `n`. The admin `?force=true` trigger
/// passes `Some(0)` so Hurl regressions can hit expired grants
/// without waiting the configured grace out.
pub async fn purge(&self, grace_override: Option<u32>) -> u64 {
let grace = grace_override.unwrap_or(self.grace_days);
let start = Instant::now();
match self.authz.purge_expired_grants(grace).await {
Ok(count) => {
// Audit-channel logging: bulk deletion of authorization
// rows is security-relevant enough to keep it in the
// audit stream even when the count is zero (proves the
// daemon is reachable).
info!(
target: "audit",
event = "grant_cleanup.purged",
count = count,
grace_days = grace,
elapsed_ms = start.elapsed().as_millis() as u64,
"👮🏻‍♂️ Purged {} expired grant(s) older than {} days",
count,
grace,
);
count
}
Err(e) => {
error!(
target: "audit",
event = "grant_cleanup.failed",
grace_days = grace,
error = %e,
"Grant cleanup failed"
);
0
}
}
}
/// Convenience for the scheduled loop.
async fn run_once(&self) {
let _ = self.purge(None).await;
}
}
+45 -33
View File
@@ -24,6 +24,11 @@ use crate::domain::entities::user::User;
/// Internal JWT claims structure for serialization.
/// This is the actual JWT payload structure used by jsonwebtoken crate.
///
/// `username` / `email` deserialize straight into `Arc<str>` (serde `rc`,
/// one allocation — same count as `String`) so the `TokenClaims` conversion
/// below is a plain move and the port-level claims can hand refcount bumps
/// to every consumer.
#[derive(Debug, Serialize, Deserialize)]
struct JwtClaims {
/// Subject identifier - contains the user ID
@@ -35,16 +40,23 @@ struct JwtClaims {
/// JWT unique ID for token tracking and revocation
pub jti: String,
/// Username for display and identification purposes
pub username: String,
pub username: Arc<str>,
/// User email for communication and identification
pub email: String,
pub email: Arc<str>,
/// User role for authorization checks
pub role: String,
}
impl From<JwtClaims> for TokenClaims {
fn from(claims: JwtClaims) -> Self {
// Pre-parse the subject once at decode time (amortized over the
// validation-cache TTL) so the auth middleware reads a `Copy` instead
// of re-parsing the 36-char string per request. A verified token we
// signed always carries a UUID `sub`; nil is a safe sentinel the
// middleware rejects. See benches/ROUND14.md §A3.
let sub_id = uuid::Uuid::parse_str(&claims.sub).unwrap_or_else(|_| uuid::Uuid::nil());
TokenClaims {
sub_id,
sub: claims.sub,
exp: claims.exp,
iat: claims.iat,
@@ -80,8 +92,16 @@ impl From<JwtClaims> for TokenClaims {
/// unique-token flooding.
/// - Expired tokens are never cached (decode itself rejects them first).
pub struct JwtTokenService {
/// Secret key used for signing JWT tokens
jwt_secret: String,
/// Pre-built signing key — `EncodingKey::from_secret` copies the secret
/// into a fresh buffer, so building it per `generate_access_token` call
/// paid an allocation per login/refresh for a process-invariant value.
encoding_key: EncodingKey,
/// Pre-built verification key (same rationale, on the validation-cache
/// miss path — every new token and every token once per TTL window).
decoding_key: DecodingKey,
/// Pre-built HS256 validation config — `Validation::new` allocates a
/// `HashSet{"exp"}` + algorithm `Vec` on every call otherwise.
validation: Validation,
/// Expiration time for access tokens in seconds
access_token_expiry: i64,
/// Expiration time for refresh tokens in seconds
@@ -125,7 +145,9 @@ impl JwtTokenService {
);
Self {
jwt_secret,
encoding_key: EncodingKey::from_secret(jwt_secret.as_bytes()),
decoding_key: DecodingKey::from_secret(jwt_secret.as_bytes()),
validation: Validation::new(Algorithm::HS256),
access_token_expiry: access_token_expiry_secs,
refresh_token_expiry: refresh_token_expiry_secs,
validation_cache,
@@ -169,9 +191,9 @@ impl TokenServicePort for JwtTokenService {
exp: now + self.access_token_expiry,
iat: now,
jti: Uuid::new_v4().to_string(),
username: user.username().unwrap_or("").to_string(),
email: user.email().to_string(),
role: format!("{}", user.role()),
username: Arc::from(user.username().unwrap_or("")),
email: Arc::from(user.email()),
role: user.role().as_str().to_string(),
};
// Log JWT claims for debugging
@@ -182,12 +204,7 @@ impl TokenServicePort for JwtTokenService {
claims.iat
);
encode(
&Header::default(),
&claims,
&EncodingKey::from_secret(self.jwt_secret.as_bytes()),
)
.map_err(|e| {
encode(&Header::default(), &claims, &self.encoding_key).map_err(|e| {
tracing::error!("Error generating token: {}", e);
DomainError::new(
ErrorKind::InternalError,
@@ -217,23 +234,18 @@ impl TokenServicePort for JwtTokenService {
// ── 2. Slow-path: full HMAC-SHA256 verification ─────────
self.cache_misses.fetch_add(1, Ordering::Relaxed);
let validation = Validation::new(Algorithm::HS256);
let token_data = decode::<JwtClaims>(
token,
&DecodingKey::from_secret(self.jwt_secret.as_bytes()),
&validation,
)
.map_err(|e| match e.kind() {
jsonwebtoken::errors::ErrorKind::ExpiredSignature => {
DomainError::new(ErrorKind::AccessDenied, "TokenService", "Token expired")
}
_ => DomainError::new(
ErrorKind::AccessDenied,
"TokenService",
format!("Invalid token: {}", e),
),
})?;
let token_data = decode::<JwtClaims>(token, &self.decoding_key, &self.validation).map_err(
|e| match e.kind() {
jsonwebtoken::errors::ErrorKind::ExpiredSignature => {
DomainError::new(ErrorKind::AccessDenied, "TokenService", "Token expired")
}
_ => DomainError::new(
ErrorKind::AccessDenied,
"TokenService",
format!("Invalid token: {}", e),
),
},
)?;
let claims = Arc::new(TokenClaims::from(token_data.claims));
@@ -301,8 +313,8 @@ mod tests {
.validate_token(&token)
.expect("Should validate token");
assert_eq!(claims.sub, user.id().to_string());
assert_eq!(Some(claims.username.as_str()), user.username());
assert_eq!(claims.email, user.email());
assert_eq!(Some(&*claims.username), user.username());
assert_eq!(&*claims.email, user.email());
}
#[test]
@@ -128,20 +128,43 @@ async fn fsync_paths_parallel(paths: Vec<PathBuf>, strict: bool) -> Result<(), D
/// (fsync now vs. deferred batch sync), or `None` when the blob already
/// existed (idempotent skip — content-addressed, so identical by definition).
async fn write_blob_bytes(blob_path: &Path, data: &Bytes) -> Result<Option<File>, DomainError> {
if fs::try_exists(blob_path).await.unwrap_or(false) {
return Ok(None);
}
let mut file = fs::File::create(blob_path).await.map_err(|e| {
DomainError::internal_error("Blob", format!("Failed to create blob file: {}", e))
})?;
// One atomic O_CREAT|O_EXCL open replaces the old stat-then-create pair:
// `AlreadyExists` IS the idempotent skip (content-addressed names mean an
// existing file has identical content), saving a syscall + a blocking-pool
// dispatch on every new chunk of every upload.
let mut file = match fs::File::options()
.write(true)
.create_new(true)
.open(blob_path)
.await
{
Ok(f) => f,
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => return Ok(None),
Err(e) => {
return Err(DomainError::internal_error(
"Blob",
format!("Failed to create blob file: {}", e),
));
}
};
file.write_all(data).await.map_err(|e| {
DomainError::internal_error("Blob", format!("Failed to write blob from bytes: {}", e))
})?;
Ok(Some(file))
}
/// Bench-only public wrapper (feature = "bench") over the private chunk
/// writer so `examples/bench_storage_micro.rs` can A/B the open strategy.
#[cfg(feature = "bench")]
pub async fn write_blob_bytes_for_bench(
blob_path: &Path,
data: &Bytes,
) -> Result<Option<File>, DomainError> {
write_blob_bytes(blob_path, data).await
}
/// Compile-time lookup table for the 256 two-digit lowercase hex prefixes ("00"…"ff").
static HEX_PREFIXES: [&str; 256] = [
pub(crate) static HEX_PREFIXES: [&str; 256] = [
"00", "01", "02", "03", "04", "05", "06", "07", "08", "09", "0a", "0b", "0c", "0d", "0e", "0f",
"10", "11", "12", "13", "14", "15", "16", "17", "18", "19", "1a", "1b", "1c", "1d", "1e", "1f",
"20", "21", "22", "23", "24", "25", "26", "27", "28", "29", "2a", "2b", "2c", "2d", "2e", "2f",
@@ -68,7 +68,24 @@ impl LoginLockoutService {
fn key(username: &str, client_ip: &str) -> String {
// `|` is not valid in either a username or an IP literal so it makes
// the username/ip boundary unambiguous.
format!("{}|{}", username.to_lowercase(), client_ip)
//
// The lowercased composite is written into ONE pre-sized buffer instead
// of the `to_lowercase()` (alloc) + `format!` (alloc) two-step. App
// passwords authenticate with an already-lowercase ASCII username in
// ~all traffic, so the fast branch covers it; the rare non-ASCII branch
// keeps `str::to_lowercase` for exact Unicode (e.g. final-sigma)
// semantics. Byte-identical key either way (benches/ROUND29.md §D).
if username.is_ascii() {
let mut k = String::with_capacity(username.len() + 1 + client_ip.len());
for &b in username.as_bytes() {
k.push(b.to_ascii_lowercase() as char);
}
k.push('|');
k.push_str(client_ip);
k
} else {
format!("{}|{}", username.to_lowercase(), client_ip)
}
}
/// Check whether the (account, IP) pair is currently locked.
@@ -96,22 +96,33 @@ impl MediaMetadataService {
}
if Self::is_image_file(mime_type) {
// ONE disk read: kamadak needs the full buffer anyway, and
// nom-exif 3.6+ parses from in-RAM bytes zero-copy
// (`MediaSource::from_memory` over the same allocation). This
// path used to re-open the file 1-2 more times — nom-exif's
// `read_exif(path)` plus a `read_track(path)` fallback for
// date-less images (2-3 opens per image, benches/ROUND12.md §M4:
// 1.44x warm geomean, 2-3x cold-cache).
let buf = std::fs::read(path).ok()?;
// Rich EXIF (GPS / camera / orientation / dimensions + naive date)
// from the proven kamadak extractor.
let kamadak = std::fs::read(path)
.ok()
.and_then(|b| ExifService::extract(&b));
let kamadak = ExifService::extract(&buf);
// nom-exif complements kamadak: a timezone-correct capture date and,
// crucially, the date + GPS for files kamadak rejects outright
// ("Unexpected next IFD"), where `kamadak` is None and the GPS would
// otherwise be lost. See `merge_image_metadata`.
merge_image_metadata(kamadak, read_nom_exif(path))
let bytes = bytes::Bytes::from(buf);
merge_image_metadata(kamadak, read_nom_exif_from_bytes(&bytes))
} else if Self::is_video_file(mime_type) {
// Videos carry no EXIF — pull the container creation time only.
read_nom_exif(path).captured_at.map(|dt| ExifMetadata {
captured_at: Some(dt),
..Default::default()
})
// Single open + header sniff; the old shape opened twice (a
// doomed `read_exif` sniff, then `read_track`).
read_nom_exif_video(path)
.captured_at
.map(|dt| ExifMetadata {
captured_at: Some(dt),
..Default::default()
})
} else {
None
}
@@ -375,34 +386,49 @@ struct NomExif {
/// carries `OffsetTimeOriginal` (or a tz-aware container time); otherwise the
/// naive wall-clock is interpreted as UTC. Either way it is converted to a true
/// UTC instant. GPS is returned as signed decimal degrees.
fn read_nom_exif(path: &Path) -> NomExif {
use nom_exif::{EntryValue, ExifTag, TrackInfoTag, read_exif, read_track};
fn nom_to_utc(ev: &nom_exif::EntryValue) -> Option<DateTime<Utc>> {
let edt = ev.as_datetime()?;
let utc0 = FixedOffset::east_opt(0)?;
Some(edt.or_offset(utc0).with_timezone(&Utc))
}
// Captures nothing → `Copy`, so it can be reused across the calls below.
let to_utc = |ev: &EntryValue| -> Option<DateTime<Utc>> {
let edt = ev.as_datetime()?;
let utc0 = FixedOffset::east_opt(0)?;
Some(edt.or_offset(utc0).with_timezone(&Utc))
};
fn nom_fill_from_exif(exif: &nom_exif::Exif, out: &mut NomExif) {
use nom_exif::ExifTag;
out.captured_at = exif
.get(ExifTag::DateTimeOriginal)
.and_then(nom_to_utc)
.or_else(|| exif.get(ExifTag::CreateDate).and_then(nom_to_utc));
if let Some(gps) = exif.gps_info() {
out.latitude = gps.latitude_decimal();
out.longitude = gps.longitude_decimal();
}
}
/// Image arm: nom-exif fed from the buffer the kamadak pass already read —
/// `MediaSource::from_memory` shares the `Bytes` refcount, so this re-parses
/// without touching the disk again (the old shape re-opened the file once,
/// plus a second time for date-less images). The track fallback stays (fed
/// from the same bytes): it covers MIME-mislabeled rows whose actual
/// container is a video — the only case where it ever produced a date.
fn read_nom_exif_from_bytes(bytes: &bytes::Bytes) -> NomExif {
use nom_exif::{MediaParser, MediaSource, TrackInfoTag};
let mut out = NomExif::default();
let mut parser = MediaParser::new();
// Images: EXIF DateTimeOriginal → DateTimeDigitized (CreateDate), plus GPS.
if let Ok(exif) = read_exif(path) {
out.captured_at = exif
.get(ExifTag::DateTimeOriginal)
.and_then(to_utc)
.or_else(|| exif.get(ExifTag::CreateDate).and_then(to_utc));
if let Some(gps) = exif.gps_info() {
out.latitude = gps.latitude_decimal();
out.longitude = gps.longitude_decimal();
}
if let Ok(ms) = MediaSource::from_memory(bytes.clone())
&& let Ok(iter) = parser.parse_exif(ms)
{
let exif: nom_exif::Exif = iter.into();
nom_fill_from_exif(&exif, &mut out);
}
// Videos / audio containers (mov/mp4/mkv): track creation time.
if out.captured_at.is_none()
&& let Ok(track) = read_track(path)
&& let Some(dt) = track.get(TrackInfoTag::CreateDate).and_then(to_utc)
&& let Ok(ms) = MediaSource::from_memory(bytes.clone())
&& let Ok(track) = parser.parse_track(ms)
&& let Some(dt) = track.get(TrackInfoTag::CreateDate).and_then(nom_to_utc)
{
out.captured_at = Some(dt);
}
@@ -410,6 +436,42 @@ fn read_nom_exif(path: &Path) -> NomExif {
out
}
/// Video arm: ONE open, dispatched on the sniffed container kind. Matches
/// the old `read_exif(path)`-then-`read_track(path)` observable behaviour
/// exactly — a Track container never parsed as EXIF (the old first open was
/// pure waste) and an Image container never parsed as a track, so the
/// two-open sequence always reduced to a single effective parse.
fn read_nom_exif_video(path: &Path) -> NomExif {
use nom_exif::{MediaKind, MediaParser, MediaSource, TrackInfoTag};
let mut out = NomExif::default();
let Ok(file) = std::fs::File::open(path) else {
return out;
};
let Ok(ms) = MediaSource::seekable(file) else {
return out;
};
let mut parser = MediaParser::new();
match ms.kind() {
MediaKind::Image => {
// MIME said video, bytes say image (mislabeled row): same EXIF
// extraction the old `read_exif(path)` performed.
if let Ok(iter) = parser.parse_exif(ms) {
let exif: nom_exif::Exif = iter.into();
nom_fill_from_exif(&exif, &mut out);
}
}
MediaKind::Track => {
if let Ok(track) = parser.parse_track(ms)
&& let Some(dt) = track.get(TrackInfoTag::CreateDate).and_then(nom_to_utc)
{
out.captured_at = Some(dt);
}
}
}
out
}
/// Combine kamadak's rich EXIF with nom-exif's date + GPS.
///
/// nom-exif's tz-correct date wins whenever present; its GPS only fills gaps
+3
View File
@@ -12,6 +12,7 @@ pub mod face_indexing_service;
pub mod ffmpeg_video_frame_service;
pub mod file_content_cache;
pub mod file_system_i18n_service;
pub mod grant_cleanup_service;
pub mod image_transcode_service;
pub mod jwt_service;
pub mod local_blob_backend;
@@ -33,6 +34,7 @@ pub mod path_service;
pub mod pg_acl_engine;
#[cfg(feature = "plugins")]
pub mod plugins;
pub mod recent_recording_hook;
pub mod retry_blob_backend;
pub mod s3_blob_backend;
pub mod search_index;
@@ -43,6 +45,7 @@ pub mod thumbnail_service;
mod thumbnail_service_test;
pub mod trash_cleanup_service;
pub mod tree_etag_flush_service;
pub mod webdav_dead_property_store;
pub mod webdav_lock_service;
pub mod wopi_discovery_service;
pub mod zip_service;
@@ -1,25 +1,84 @@
use std::path::PathBuf;
use std::time::Duration;
use tokio::fs;
use tokio::io::AsyncWriteExt;
use crate::common::errors::{DomainError, Result};
/// In-RAM running byte counter per upload session (`user/upload_id` →
/// bytes accepted so far). The per-chunk quota gate used to recompute
/// this by listing the whole session directory and stat-ing every chunk
/// on EVERY chunk PUT — O(k) stats for chunk k, O(N²/2) over an upload
/// (~500k stats for a 10 GB / 1000-chunk upload). The counter makes the
/// gate O(1); a cache miss (process restart, eviction) lazily rebuilds
/// from the directory listing, so crash-correctness is unchanged
/// (benches/NC-CHUNK-GATE.md). Sessions are forgotten on cleanup; the
/// TTL reaps counters for sessions the client abandoned.
fn build_session_bytes_cache() -> moka::sync::Cache<String, u64> {
moka::sync::Cache::builder()
.max_capacity(100_000)
.time_to_idle(Duration::from_secs(24 * 3600))
.build()
}
#[derive(Clone)]
pub struct NextcloudChunkedUploadService {
pub base_dir: PathBuf,
/// See [`build_session_bytes_cache`]. Cloning the service shares the
/// counter (moka `Cache` clones are handles to the same store).
session_bytes: moka::sync::Cache<String, u64>,
}
impl NextcloudChunkedUploadService {
pub fn new(base_dir: PathBuf) -> Self {
Self { base_dir }
Self {
base_dir,
session_bytes: build_session_bytes_cache(),
}
}
pub fn new_stub() -> Self {
Self {
base_dir: PathBuf::from("./storage/.uploads/nextcloud"),
session_bytes: build_session_bytes_cache(),
}
}
fn bytes_key(user: &str, upload_id: &str) -> String {
format!("{user}/{upload_id}")
}
/// Session bytes accepted so far, if the counter is warm.
/// `None` = rebuild from the directory listing and call
/// [`Self::set_session_bytes`].
pub fn cached_session_bytes(&self, user: &str, upload_id: &str) -> Option<u64> {
self.session_bytes.get(&Self::bytes_key(user, upload_id))
}
/// Seed / overwrite the session counter (post-rebuild or on MKCOL).
pub fn set_session_bytes(&self, user: &str, upload_id: &str, bytes: u64) {
self.session_bytes
.insert(Self::bytes_key(user, upload_id), bytes);
}
/// Add an accepted chunk's bytes to the counter (no-op when cold —
/// the next gate rebuilds from disk). Two racing PUTs on one session
/// could drop an increment; the counter is a gate hint, and the
/// MOVE-time quota check stays authoritative.
pub fn bump_session_bytes(&self, user: &str, upload_id: &str, delta: u64) {
let key = Self::bytes_key(user, upload_id);
if let Some(current) = self.session_bytes.get(&key) {
self.session_bytes
.insert(key, current.saturating_add(delta));
}
}
/// Drop the counter (session cleanup, or a chunk overwrite made the
/// running total untrustworthy — rebuilt lazily on next use).
pub fn forget_session_bytes(&self, user: &str, upload_id: &str) {
self.session_bytes
.invalidate(&Self::bytes_key(user, upload_id));
}
/// Validate that a path component contains no traversal characters.
fn validate_path_component(name: &str, label: &str) -> Result<()> {
if name.is_empty()
@@ -49,6 +108,7 @@ impl NextcloudChunkedUploadService {
fs::create_dir_all(&session_dir)
.await
.map_err(|e| DomainError::internal_error("ChunkedUpload", e.to_string()))?;
self.set_session_bytes(user, upload_id, 0);
Ok(())
}
@@ -76,6 +136,20 @@ impl NextcloudChunkedUploadService {
/// `interfaces/upload_ingest::stream_body_to_path` helper to stream the
/// HTTP body directly to disk and avoid materialising the whole chunk
/// in RAM.
///
/// Uses `tokio::fs::write` (single `spawn_blocking` around
/// `std::fs::write`) rather than manually driving
/// `create + write_all` and letting the tokio handle drop close the
/// fd. The manual shape leaked a race: `tokio::fs::File::drop`
/// dispatches `close(2)` to the blocking pool without awaiting it,
/// and until close completes the dirent update may not be visible
/// to a subsequent `read_dir` — on macOS APFS routinely, on Linux
/// under I/O contention. In practice that turned into
/// `ordered_chunk_paths` silently missing a just-uploaded chunk;
/// the NC assembly path (`handle_assemble` → `ordered_chunk_paths`)
/// would then produce a truncated file with no error to the client.
/// `std::fs::write` opens, writes, and synchronously closes before
/// returning, so the dirent is guaranteed visible on `.await`.
pub async fn store_chunk(
&self,
user: &str,
@@ -84,12 +158,16 @@ impl NextcloudChunkedUploadService {
data: &[u8],
) -> Result<()> {
let chunk_path = self.safe_chunk_path(user, upload_id, chunk_name)?;
let mut file = fs::File::create(&chunk_path)
.await
.map_err(|e| DomainError::internal_error("ChunkedUpload", e.to_string()))?;
file.write_all(data)
let overwrite = fs::metadata(&chunk_path).await.is_ok();
fs::write(&chunk_path, data)
.await
.map_err(|e| DomainError::internal_error("ChunkedUpload", e.to_string()))?;
if overwrite {
// Retried chunk — running total is stale; rebuild lazily.
self.forget_session_bytes(user, upload_id);
} else {
self.bump_session_bytes(user, upload_id, data.len() as u64);
}
Ok(())
}
@@ -137,6 +215,7 @@ impl NextcloudChunkedUploadService {
.await
.map_err(|e| DomainError::internal_error("ChunkedUpload", e.to_string()))?;
}
self.forget_session_bytes(user, upload_id);
Ok(())
}
@@ -10,11 +10,13 @@ use std::sync::Arc;
use uuid::Uuid;
use crate::application::dtos::display_helpers::{
category_for, format_file_size, icon_class_for, icon_special_class_for,
classify_display, format_file_size, intern_display, intern_mime,
};
use crate::application::dtos::file_dto::FileDto;
use crate::application::dtos::folder_dto::FolderDto;
use crate::common::errors::DomainError;
use crate::domain::entities::file::File;
use crate::domain::entities::folder::Folder;
/// Result of resolving a WebDAV path — either a folder or a file.
#[derive(Debug, Clone)]
@@ -33,14 +35,20 @@ impl PathResolverService {
Self { pool }
}
/// Resolve `path` to a folder or file **owned by `user_id`**.
/// Resolve `path` to a folder or file **within the given drive**.
///
/// Adds `AND fo.user_id = $4` / `AND fi.user_id = $4` so that one
/// user can never resolve another user's resources.
pub async fn resolve_path_for_user(
/// Filters on `fo.drive_id = $4` / `fi.drive_id = $4`. Callers
/// pre-resolve which drive they're operating in — native WebDAV
/// derives it from the caller's default drive
/// (`resolve_drive_id_for_native_webdav`); NC WebDAV takes it from
/// the URL-selected chroot (`chroot.drive_id`). Shared by both
/// surfaces so the single-query UNION ALL optimisation lands
/// consistently and no path lookup keys on the doomed
/// `storage.{files,folders}.user_id` column.
pub async fn resolve_path_in_drive(
&self,
path: &str,
user_id: Uuid,
drive_id: Uuid,
) -> Result<ResolvedResource, DomainError> {
let path = path.trim_start_matches('/').trim_end_matches('/');
if path.is_empty() {
@@ -55,6 +63,13 @@ impl PathResolverService {
String::new()
};
// Widened SELECT: also fetches `blob_hash` (for file ETag) and
// `tree_modified_at` (for folder ETag). Both share the same
// canonical formulas as the rest of the codebase — see
// [`File::compute_etag`] and [`Folder::compute_etag`]. Without
// these two extra columns the resolver used to emit empty
// ETag strings, and NC's `If-Match` round-trips broke
// (see the F6b regression on `test_nc_put_mkcol_blake3.sh`).
let row = sqlx::query_as::<
_,
(
@@ -63,34 +78,37 @@ impl PathResolverService {
String, // name
String, // path
Option<String>, // parent_id
Option<String>, // user_id
Uuid, // drive_id
i64, // created_at
i64, // modified_at
Option<i64>, // size
Option<String>, // mime_type
Option<String>, // folder_id
Option<String>, // blob_hash (files only)
Option<i64>, // tree_modified_at (folders only)
),
>(
r#"
SELECT resource_type, id, name, path, parent_id, user_id, drive_id,
created_at, modified_at, size, mime_type, folder_id
SELECT resource_type, id, name, path, parent_id, drive_id,
created_at, modified_at, size, mime_type, folder_id,
blob_hash, tree_modified_at
FROM (
SELECT 'folder'::text AS resource_type,
fo.id::text,
fo.name,
fo.path,
fo.parent_id::text,
fo.user_id::text,
fo.drive_id,
EXTRACT(EPOCH FROM fo.created_at)::bigint AS created_at,
EXTRACT(EPOCH FROM fo.updated_at)::bigint AS modified_at,
NULL::bigint AS size,
NULL::text AS mime_type,
NULL::text AS folder_id
NULL::text AS folder_id,
NULL::text AS blob_hash,
EXTRACT(EPOCH FROM fo.tree_modified_at)::bigint AS tree_modified_at
FROM storage.folders fo
WHERE fo.path = $1 AND NOT fo.is_trashed
AND fo.user_id = $4
AND fo.drive_id = $4
UNION ALL
@@ -103,13 +121,14 @@ impl PathResolverService {
ELSE fi.name
END AS path,
NULL::text AS parent_id,
fi.user_id::text,
fi.drive_id,
EXTRACT(EPOCH FROM fi.created_at)::bigint AS created_at,
EXTRACT(EPOCH FROM fi.updated_at)::bigint AS modified_at,
fi.size,
fi.mime_type,
fi.folder_id::text
fi.folder_id::text,
fi.blob_hash,
NULL::bigint AS tree_modified_at
FROM storage.files fi
LEFT JOIN storage.folders fo ON fo.id = fi.folder_id
WHERE fi.name = $2
@@ -118,7 +137,7 @@ impl PathResolverService {
OR fo.path = $3
)
AND NOT fi.is_trashed
AND fi.user_id = $4
AND fi.drive_id = $4
) sub
LIMIT 1
"#,
@@ -126,10 +145,10 @@ impl PathResolverService {
.bind(path) // $1
.bind(filename) // $2
.bind(&folder_path) // $3
.bind(user_id) // $4
.bind(drive_id) // $4
.fetch_optional(self.pool.as_ref())
.await
.map_err(|e| DomainError::internal_error("PathResolver", format!("resolve_for_user: {e}")))?
.map_err(|e| DomainError::internal_error("PathResolver", format!("resolve_in_drive: {e}")))?
.ok_or_else(|| DomainError::not_found("Resource", path))?;
let (
@@ -138,62 +157,63 @@ impl PathResolverService {
name,
res_path,
parent_id,
uid,
drive_id,
created_at,
modified_at,
size,
mime_type,
folder_id,
blob_hash,
tree_modified_at,
) = row;
match resource_type.as_str() {
"folder" => Ok(ResolvedResource::Folder(FolderDto {
etag: id.clone(),
id,
name: name.clone(),
path: res_path,
parent_id,
owner_id: uid,
drive_id,
created_at: created_at as u64,
modified_at: modified_at as u64,
is_root: false,
icon_class: Arc::from("fas fa-folder"),
icon_special_class: Arc::from("folder-icon"),
category: Arc::from("Folder"),
// §14 provenance not selected by this resolver path —
// it's used for existence/type discrimination, not
// detailed DTO emission. Callers that need provenance
// reload through the repo.
created_by: None,
updated_by: None,
})),
"folder" => {
let tree_mod = tree_modified_at.unwrap_or(modified_at) as u64;
Ok(ResolvedResource::Folder(FolderDto {
etag: Folder::compute_etag(&id, tree_mod),
id,
name: name.clone(),
path: res_path,
parent_id,
drive_id,
created_at: created_at as u64,
modified_at: modified_at as u64,
is_root: false,
icon_class: intern_display("fas fa-folder"),
icon_special_class: intern_display("folder-icon"),
category: intern_display("Folder"),
// §14 provenance not selected by this resolver path —
// it's used for existence/type discrimination, not
// detailed DTO emission. Callers that need provenance
// reload through the repo.
created_by: None,
updated_by: None,
}))
}
_ => {
let mime = mime_type.unwrap_or_else(|| "application/octet-stream".to_string());
let sz = size.unwrap_or(0) as u64;
// `content_hash`/`etag` are empty here: this resolver
// path doesn't select `blob_hash` from SQL — callers
// are doing existence/type discrimination, not ETag
// emission. If a caller ever needs an ETag from this
// codepath, widen the SELECT and populate properly.
let hash = blob_hash.unwrap_or_default();
let modified_at_u = modified_at as u64;
let etag = File::compute_etag(&hash, modified_at_u);
let classes = classify_display(&name, &mime);
Ok(ResolvedResource::File(FileDto {
id,
name: name.clone(),
path: res_path,
size: sz,
mime_type: Arc::from(&*mime),
mime_type: intern_mime(&mime),
folder_id,
created_at: created_at as u64,
modified_at: modified_at as u64,
icon_class: Arc::from(icon_class_for(&name, &mime)),
icon_special_class: Arc::from(icon_special_class_for(&name, &mime)),
category: Arc::from(category_for(&name, &mime)),
modified_at: modified_at_u,
icon_class: intern_display(classes.icon_class),
icon_special_class: intern_display(classes.icon_special_class),
category: intern_display(classes.category),
size_formatted: format_file_size(sz),
owner_id: uid,
sort_date: None,
content_hash: String::new(),
etag: String::new(),
content_hash: hash,
etag,
// §14 provenance not selected by this resolver path
created_by: None,
updated_by: None,
@@ -203,7 +223,10 @@ impl PathResolverService {
}
/// Returns `true` if the resource at `path` belongs to `user_id`.
pub async fn exists_for_user(&self, path: &str, user_id: Uuid) -> Result<bool, DomainError> {
/// Check whether `path` resolves to a folder or file within the
/// given drive. Companion to `resolve_path_in_drive` — same scope
/// filter, existence-only projection.
pub async fn exists_in_drive(&self, path: &str, drive_id: Uuid) -> Result<bool, DomainError> {
let path = path.trim_start_matches('/').trim_end_matches('/');
if path.is_empty() {
return Ok(false);
@@ -221,7 +244,7 @@ impl PathResolverService {
r#"
SELECT EXISTS(
SELECT 1 FROM storage.folders
WHERE path = $1 AND NOT is_trashed AND user_id = $4
WHERE path = $1 AND NOT is_trashed AND drive_id = $4
) OR EXISTS(
SELECT 1
FROM storage.files fi
@@ -229,18 +252,18 @@ impl PathResolverService {
WHERE fi.name = $2
AND (($3 = '' AND fi.folder_id IS NULL) OR fo.path = $3)
AND NOT fi.is_trashed
AND fi.user_id = $4
AND fi.drive_id = $4
)
"#,
)
.bind(path)
.bind(filename)
.bind(&folder_path)
.bind(user_id)
.bind(drive_id)
.fetch_one(self.pool.as_ref())
.await
.map_err(|e| {
DomainError::internal_error("PathResolver", format!("exists_for_user: {e}"))
DomainError::internal_error("PathResolver", format!("exists_in_drive: {e}"))
})?;
Ok(exists)
+1 -1
View File
@@ -95,7 +95,7 @@ impl PathService {
/// Validates a path to ensure it doesn't contain dangerous components
pub fn validate_path(&self, path: &StoragePath) -> Result<(), DomainError> {
// Check for empty segments
if path.segments().iter().any(|s| s.is_empty()) {
if path.segments().any(|s| s.is_empty()) {
return Err(DomainError::new(
ErrorKind::InvalidInput,
"Path",
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,120 @@
//! Recording side of [`ResourceAccessHook`] — turns a successful file access
//! into a row in `auth.user_recent_files` via [`RecentService`].
//!
//! Wiring lives in `common/di.rs`: this hook is registered once, every
//! `_with_perms` file method on `FileRetrievalService` / `FileManagementService`
//! fans through it, and any future read-path or write-path service can opt in
//! by holding an `Option<Arc<dyn ResourceAccessHook>>` and calling
//! `on_file_accessed` after authZ.
//!
//! Two non-obvious behaviours, with rationale:
//!
//! * **Per-(caller, file) 60-second throttle.** Range-stream downloads send
//! one GET per chunk (NC desktop, video seek, resumable transfers); without
//! throttling each chunk would trigger an upsert against the same row.
//! Moka's `time_to_live` gives us bounded memory and lock-free reads. The
//! underlying `INSERT … ON CONFLICT DO UPDATE accessed_at = now()` is
//! idempotent, so the rare TOCTOU window between `contains_key` and `insert`
//! is harmless — at worst we record twice for the same instant.
//!
//! * **Fire-and-forget via `tokio::spawn`.** The `ResourceAccessHook` method
//! is synchronous by contract (every `with_perms` caller would otherwise
//! have to `await` the side-effect). The spawn lets the user-facing
//! response return immediately; a DB hiccup in Recent recording never
//! bubbles up to the GET / PUT that triggered it. Failures log at warn.
use std::sync::Arc;
use std::time::Duration;
use moka::sync::Cache;
use uuid::Uuid;
use crate::application::ports::resource_access_hook::ResourceAccessHook;
use crate::application::services::recent_service::RecentService;
/// How long a successful recording suppresses repeat upserts for the same
/// `(caller, file)`. Sized to span a typical streamed range-GET burst while
/// still updating `accessed_at` often enough that the Recent list reflects
/// "this is the file I was just looking at".
const THROTTLE_TTL_SECONDS: u64 = 60;
/// Bound on simultaneous in-flight throttle entries. Each entry is a tuple
/// `(Uuid, String) -> ()` ≈ 80 B; 16 384 entries ≈ 1.3 MB worst case. LRU
/// eviction keeps memory bounded even if a pathological client touches a
/// million files in a minute.
const THROTTLE_MAX_ENTRIES: u64 = 16_384;
/// `ResourceAccessHook` implementation that records file accesses into
/// `auth.user_recent_files`, throttled per (caller, file).
pub struct RecentRecordingHook {
recent: Arc<RecentService>,
throttle: Cache<(Uuid, String), ()>,
}
impl RecentRecordingHook {
pub fn new(recent: Arc<RecentService>) -> Self {
// `support_invalidation_closures` is the moka opt-in needed by
// `invalidate_entries_if` (the per-user throttle reset on
// `on_recents_cleared`). Without it, the predicate-based
// invalidate call silently no-ops and a freshly-cleared Recent
// list refuses to re-record the same file until the TTL
// expires — exactly the bug surfaced by tests/api/recent.hurl
// step 8.
let throttle = Cache::builder()
.max_capacity(THROTTLE_MAX_ENTRIES)
.time_to_live(Duration::from_secs(THROTTLE_TTL_SECONDS))
.support_invalidation_closures()
.build();
Self { recent, throttle }
}
}
impl ResourceAccessHook for RecentRecordingHook {
fn on_file_accessed(&self, caller_id: Uuid, file_id: &str) {
let key = (caller_id, file_id.to_string());
if self.throttle.contains_key(&key) {
return;
}
// Insert before spawning: even if the spawned task races with another
// call for the same key, the cache entry suppresses the duplicate
// before it reaches the DB. The ON CONFLICT clause covers the
// sub-microsecond TOCTOU window between contains_key and insert.
self.throttle.insert(key.clone(), ());
let recent = Arc::clone(&self.recent);
let (caller_id, file_id) = key;
tokio::spawn(async move {
// Fast path: skip the trait's `authz.require(Read, …)`
// (upstream `_with_perms` service already gated). The
// extra SQL round-trip pushes the upsert past the client's
// immediate `GET /api/recent/resources` in
// `tests/api/recent.hurl` step 7 — the whole reason for
// the internal variant.
if let Err(e) = recent
.record_item_access_internal(caller_id, &file_id, "file")
.await
{
tracing::warn!(
target: "oxicloud::recent",
caller_id = %caller_id,
file_id = %file_id,
"recent recording failed: {e}",
);
}
});
}
fn on_recents_cleared(&self, caller_id: Uuid) {
// Drop every throttle entry that would otherwise suppress the
// next recording for this user. moka schedules the predicate to
// run during the next maintenance pass — it's not synchronous.
// The DB clear has already happened by the time we get here, so
// any racing access between the clear and the next maintenance
// pass just re-records via ON CONFLICT — the worst case is a row
// that surfaces in Recent a few ms after the clear, which is
// exactly what the user asked for.
let _ = self
.throttle
.invalidate_entries_if(move |(k_caller, _), _| *k_caller == caller_id);
}
}
+127 -48
View File
@@ -57,14 +57,20 @@ impl RetryBlobBackend {
}
/// Execute an async closure with exponential backoff retry.
async fn retry_async<F, Fut, T>(
///
/// `name` is a lazy label: the success path (the overwhelmingly common
/// case) never materializes it, so per-op `format!("op({hash})")`
/// allocations only happen on an actual retry (benches/ROUND11.md §14:
/// 64.5 → 0.7 ns, −2 allocs per blob op).
async fn retry_async<F, Fut, T, L>(
policy: &RetryPolicy,
name: &str,
name: L,
mut f: F,
) -> Result<T, DomainError>
where
F: FnMut() -> Fut,
Fut: std::future::Future<Output = Result<T, DomainError>>,
L: Fn() -> String,
{
let mut attempt = 0u32;
let mut backoff = policy.initial_backoff;
@@ -78,7 +84,7 @@ where
"Retry {}/{} for {} after error: {} (backoff {:?})",
attempt,
policy.max_retries,
name,
name(),
e,
backoff
);
@@ -112,10 +118,14 @@ impl BlobStorageBackend for RetryBlobBackend {
let inner = self.inner.clone();
let policy = self.policy.clone();
Box::pin(async move {
retry_async(&policy, "initialize", || {
let inner = inner.clone();
async move { inner.initialize().await }
})
retry_async(
&policy,
|| "initialize".to_string(),
|| {
let inner = inner.clone();
async move { inner.initialize().await }
},
)
.await
})
}
@@ -130,12 +140,16 @@ impl BlobStorageBackend for RetryBlobBackend {
let hash = hash.to_string();
let path = source_path.to_path_buf();
Box::pin(async move {
retry_async(&policy, &format!("put_blob({hash})"), || {
let inner = inner.clone();
let hash = hash.clone();
let path = path.clone();
async move { inner.put_blob(&hash, &path).await }
})
retry_async(
&policy,
|| format!("put_blob({hash})"),
|| {
let inner = inner.clone();
let hash = hash.clone();
let path = path.clone();
async move { inner.put_blob(&hash, &path).await }
},
)
.await
})
}
@@ -149,16 +163,57 @@ impl BlobStorageBackend for RetryBlobBackend {
let policy = self.policy.clone();
let hash = hash.to_string();
Box::pin(async move {
retry_async(&policy, &format!("put_blob_from_bytes({hash})"), || {
let inner = inner.clone();
let hash = hash.clone();
let data = data.clone();
async move { inner.put_blob_from_bytes(&hash, data).await }
})
retry_async(
&policy,
|| format!("put_blob_from_bytes({hash})"),
|| {
let inner = inner.clone();
let hash = hash.clone();
let data = data.clone();
async move { inner.put_blob_from_bytes(&hash, data).await }
},
)
.await
})
}
// Without this override the trait default would re-route the CDC chunk
// write through `put_blob_from_bytes` above — reinstating the remote
// backend's exists-probe (HEAD/get_properties) per chunk that the
// `_unsynced` fast path exists to skip.
fn put_blob_from_bytes_unsynced(
&self,
hash: &str,
data: Bytes,
) -> Pin<Box<dyn std::future::Future<Output = Result<u64, DomainError>> + Send + '_>> {
let inner = self.inner.clone();
let policy = self.policy.clone();
let hash = hash.to_string();
Box::pin(async move {
retry_async(
&policy,
|| format!("put_blob_from_bytes_unsynced({hash})"),
|| {
let inner = inner.clone();
let hash = hash.clone();
let data = data.clone();
async move { inner.put_blob_from_bytes_unsynced(&hash, data).await }
},
)
.await
})
}
// Forwarded WITHOUT retry wrapping: a failed fsync must surface, not be
// re-issued — after an fsync error the kernel may have dropped the dirty
// pages, so a retried fsync can report success for data that was lost.
fn sync_blobs(
&self,
hashes: &[String],
) -> Pin<Box<dyn std::future::Future<Output = Result<(), DomainError>> + Send + '_>> {
self.inner.sync_blobs(hashes)
}
fn get_blob_stream(
&self,
hash: &str,
@@ -168,11 +223,15 @@ impl BlobStorageBackend for RetryBlobBackend {
let policy = self.policy.clone();
let hash = hash.to_string();
Box::pin(async move {
retry_async(&policy, &format!("get_blob_stream({hash})"), || {
let inner = inner.clone();
let hash = hash.clone();
async move { inner.get_blob_stream(&hash).await }
})
retry_async(
&policy,
|| format!("get_blob_stream({hash})"),
|| {
let inner = inner.clone();
let hash = hash.clone();
async move { inner.get_blob_stream(&hash).await }
},
)
.await
})
}
@@ -188,11 +247,15 @@ impl BlobStorageBackend for RetryBlobBackend {
let policy = self.policy.clone();
let hash = hash.to_string();
Box::pin(async move {
retry_async(&policy, &format!("get_blob_range({hash})"), || {
let inner = inner.clone();
let hash = hash.clone();
async move { inner.get_blob_range_stream(&hash, start, end).await }
})
retry_async(
&policy,
|| format!("get_blob_range({hash})"),
|| {
let inner = inner.clone();
let hash = hash.clone();
async move { inner.get_blob_range_stream(&hash, start, end).await }
},
)
.await
})
}
@@ -205,11 +268,15 @@ impl BlobStorageBackend for RetryBlobBackend {
let policy = self.policy.clone();
let hash = hash.to_string();
Box::pin(async move {
retry_async(&policy, &format!("delete_blob({hash})"), || {
let inner = inner.clone();
let hash = hash.clone();
async move { inner.delete_blob(&hash).await }
})
retry_async(
&policy,
|| format!("delete_blob({hash})"),
|| {
let inner = inner.clone();
let hash = hash.clone();
async move { inner.delete_blob(&hash).await }
},
)
.await
})
}
@@ -222,11 +289,15 @@ impl BlobStorageBackend for RetryBlobBackend {
let policy = self.policy.clone();
let hash = hash.to_string();
Box::pin(async move {
retry_async(&policy, &format!("blob_exists({hash})"), || {
let inner = inner.clone();
let hash = hash.clone();
async move { inner.blob_exists(&hash).await }
})
retry_async(
&policy,
|| format!("blob_exists({hash})"),
|| {
let inner = inner.clone();
let hash = hash.clone();
async move { inner.blob_exists(&hash).await }
},
)
.await
})
}
@@ -239,11 +310,15 @@ impl BlobStorageBackend for RetryBlobBackend {
let policy = self.policy.clone();
let hash = hash.to_string();
Box::pin(async move {
retry_async(&policy, &format!("blob_size({hash})"), || {
let inner = inner.clone();
let hash = hash.clone();
async move { inner.blob_size(&hash).await }
})
retry_async(
&policy,
|| format!("blob_size({hash})"),
|| {
let inner = inner.clone();
let hash = hash.clone();
async move { inner.blob_size(&hash).await }
},
)
.await
})
}
@@ -256,10 +331,14 @@ impl BlobStorageBackend for RetryBlobBackend {
let inner = self.inner.clone();
let policy = self.policy.clone();
Box::pin(async move {
retry_async(&policy, "health_check", || {
let inner = inner.clone();
async move { inner.health_check().await }
})
retry_async(
&policy,
|| "health_check".to_string(),
|| {
let inner = inner.clone();
async move { inner.health_check().await }
},
)
.await
})
}
@@ -200,6 +200,38 @@ impl BlobStorageBackend for S3BlobBackend {
})
}
/// Dedup settle path: PUT unconditionally. Keys are content-addressed
/// (BLAKE3), so a re-PUT writes identical bytes — overwrite-safe
/// idempotency without the HEAD probe `put_blob_from_bytes` pays. The
/// dedup layer already filtered out chunks the database knows about,
/// so the probe was a pure extra round-trip on every NEW chunk of
/// every upload (2 RTTs -> 1, benches/S3-PUT.md).
fn put_blob_from_bytes_unsynced(
&self,
hash: &str,
data: Bytes,
) -> Pin<Box<dyn std::future::Future<Output = Result<u64, DomainError>> + Send + '_>> {
let hash = hash.to_owned();
Box::pin(async move {
let key = Self::object_key(&hash);
let size = data.len() as u64;
self.client
.put_object()
.bucket(&self.bucket)
.key(&key)
.body(ByteStream::from(data))
.send()
.await
.map_err(|e| {
DomainError::internal_error(
"S3",
format!("Failed to upload blob {}: {}", hash, e),
)
})?;
Ok(size)
})
}
fn get_blob_stream(
&self,
hash: &str,
@@ -242,30 +242,43 @@ impl ContentIndexWorker {
// Authoritative state re-read: a queued 'upsert' whose row vanished
// or got trashed in the meantime becomes a delete.
let files: Vec<(Uuid, String, String, String, String, String, i64)> =
if upsert_candidates.is_empty() {
Vec::new()
} else {
sqlx::query_as(
"SELECT fi.id, fi.user_id::text, fi.drive_id::text, fi.name,
fi.blob_hash, fi.mime_type, fi.size
FROM storage.files fi
WHERE fi.id = ANY($1) AND NOT fi.is_trashed",
)
.bind(&upsert_candidates)
.fetch_all(self.maintenance_pool.as_ref())
.await?
};
//
// Post-D7: `fi.user_id` is dropped — no longer projected. The
// Tantivy `user_id` field survives as defence-in-depth but now
// always indexes `""`. Every query is Must-scoped by `drive_id`.
// (file_id, drive_id, name, blob_hash, mime, size).
type FileIndexRow = (Uuid, String, String, String, String, i64);
let files: Vec<FileIndexRow> = if upsert_candidates.is_empty() {
Vec::new()
} else {
sqlx::query_as(
"SELECT fi.id, fi.drive_id::text, fi.name,
fi.blob_hash, fi.mime_type, fi.size
FROM storage.files fi
WHERE fi.id = ANY($1) AND NOT fi.is_trashed",
)
.bind(&upsert_candidates)
.fetch_all(self.maintenance_pool.as_ref())
.await?
};
let found: HashSet<Uuid> = files.iter().map(|f| f.0).collect();
deletes.extend(upsert_candidates.iter().filter(|id| !found.contains(id)));
// `supports` lowercases the MIME (and, on a generic MIME, the extension)
// — 1–2 allocations per call. Classify each file ONCE here and thread the
// flag through both the wanted-hashes filter and the per-file records
// loop below, where it used to be re-derived a second time per file.
let supported: Vec<bool> = files
.iter()
.map(|(_, _, name, _, mime, _)| text_extractor::supports(name, mime))
.collect();
// Per-blob text: batch-read the extraction cache, extract misses.
let wanted_hashes: Vec<String> = files
.iter()
.filter(|(_, _, _, name, _, mime, size)| {
text_extractor::supports(name, mime) && *size as u64 <= self.max_extract_file_bytes
})
.map(|f| f.4.clone())
.zip(&supported)
.filter(|&(f, sup)| *sup && f.5 as u64 <= self.max_extract_file_bytes)
.map(|(f, _)| f.3.clone())
.collect();
let mut text_by_hash: HashMap<String, Option<String>> = HashMap::new();
if !wanted_hashes.is_empty() {
@@ -282,8 +295,9 @@ impl ContentIndexWorker {
}
let mut records = Vec::with_capacity(files.len());
for (file_id, user_id, drive_id, name, blob_hash, mime, size) in files {
let supported = text_extractor::supports(&name, &mime);
for ((file_id, drive_id, name, blob_hash, mime, size), supported) in
files.into_iter().zip(supported)
{
let content = if !supported {
None
} else if let Some(cached) = text_by_hash.get(&blob_hash) {
@@ -301,7 +315,7 @@ impl ContentIndexWorker {
.map(|t| truncate_on_char(t, PREVIEW_BYTES));
records.push(IndexDocRecord {
file_id: file_id.to_string(),
user_id,
user_id: String::new(),
drive_id,
name,
content,
@@ -238,8 +238,10 @@ impl TantivyContentIndex {
}
/// Tokenize `raw` with the index analyzer (simple split + lowercase).
fn query_tokens(analyzer: &TextAnalyzer, raw: &str) -> Vec<String> {
let mut analyzer = analyzer.clone();
/// Takes the analyzer by value — the caller's per-search clone is the
/// only one needed; cloning the boxed tokenizer chain again here doubled
/// the per-query allocation for nothing.
fn query_tokens(mut analyzer: TextAnalyzer, raw: &str) -> Vec<String> {
let mut tokens = Vec::new();
let mut stream = analyzer.token_stream(raw);
while stream.advance() && tokens.len() < MAX_QUERY_TOKENS {
@@ -337,7 +339,7 @@ impl TantivyContentIndex {
raw_query: &str,
limit: usize,
) -> Result<Vec<ContentHitDto>, DomainError> {
let tokens = Self::query_tokens(&analyzer, raw_query);
let tokens = Self::query_tokens(analyzer, raw_query);
if tokens.is_empty() {
return Ok(Vec::new());
}
@@ -347,6 +349,14 @@ impl TantivyContentIndex {
.search(&query, &TopDocs::with_limit(limit.max(1)).order_by_score())
.map_err(|e| DomainError::internal_error("ContentIndex", format!("search: {e}")))?;
// No hits → no documents to highlight. `SnippetGenerator::create`
// compiles the query against the index (term lookups + weight build);
// for a query that matched nothing that is pure waste on the search
// request path, and the per-hit loop below never runs. Return early.
if top_docs.is_empty() {
return Ok(Vec::new());
}
// Snippets highlight CONTENT matches; an empty fragment means the hit
// came from the name (or a fuzzy variant) — no snippet then.
let snippet_generator = SnippetGenerator::create(&searcher, &*query, fields.content)
@@ -243,7 +243,10 @@ fn collect_xml_text<R: std::io::BufRead>(
}
match xml.read_event_into(&mut buf) {
Ok(Event::Text(t)) => {
if let Ok(decoded) = t.xml_content() {
// quick-xml 0.41+ makes XmlVersion explicit on xml_content()
// so callers pick 1.0 vs 1.1 entity-normalization rules. Text
// extraction is version-agnostic — 1.0 is the sane default.
if let Ok(decoded) = t.xml_content(quick_xml::XmlVersion::Implicit1_0) {
out.push_str(&decoded);
}
}
@@ -1479,7 +1479,7 @@ impl crate::application::ports::blob_lifecycle::BlobLifecycleHook for ThumbnailS
for format in [ThumbnailFormat::Webp, ThumbnailFormat::Jpeg] {
let path =
root.join(size.dir_name())
.join(format!("{}.{}", &blob_hash, format.ext()));
.join(format!("{}.{}", blob_hash, format.ext()));
if tokio::fs::metadata(&path).await.is_ok() {
let _ = tokio::fs::remove_file(&path).await;
}
@@ -0,0 +1,264 @@
//! PostgreSQL-backed dead property store for WebDAV PROPPATCH / PROPFIND compliance.
//!
//! RFC 4918 §4.2 defines "dead properties" as those stored verbatim by the
//! server without interpreting their value. Properties are persisted to
//! `storage.webdav_dead_properties` and survive server restarts.
//!
//! Keying contract (after migration 20260830000001): the row is keyed by
//! the underlying resource id — exactly one of `folder_id` / `file_id` is
//! set — not by the resource's current path. Three consequences:
//!
//! * Every delete code path (REST, WebDAV, NextCloud DAV, trash empty,
//! folder cascade) reaps dead-property rows for free via FK
//! `ON DELETE CASCADE`. The store has no `remove_resource()` method
//! because it isn't needed: deleting the file/folder row reaps the
//! attached dead properties as a database invariant.
//! * MOVE / RENAME never changes the resource id, so dead properties
//! follow the resource without any store-side bookkeeping. The store
//! has no `rename_resource()` method for the same reason.
//! * Dead properties are RESOURCE state (RFC 4918 §4.2), not user
//! state. Two users on a shared drive PROPFIND'ing the same resource
//! see the same dead properties. The `user_id` scope key from the
//! pre-rekey schema is gone; user-delete cleanup happens
//! transitively through `auth.users` → `storage.{folders,files}` →
//! this table.
//!
//! Queries use `sqlx::query()` (runtime-bound) rather than `sqlx::query!()`
//! to keep fresh checkouts compilable without a DB connection — the
//! codebase's standing convention.
//!
//! COPY semantics (RFC 4918 §8.8 — dead properties MUST be duplicated)
//! are NOT handled here. The COPY handler is responsible for explicitly
//! reading the source's dead properties via `get_all()` and writing them
//! against the new resource id via `set()`. Not done in this migration —
//! it was not handled by the path-based store either, so this is a
//! parity decision, not a regression.
use std::collections::HashMap;
use std::sync::Arc;
use sqlx::{PgPool, Row};
use uuid::Uuid;
use crate::application::adapters::webdav_adapter::QualifiedName;
use crate::domain::errors::DomainError;
/// Polymorphic reference to the resource a dead property hangs off.
///
/// Exactly one variant — folder or file — is ever stored in a single
/// row. The CHECK constraint
/// `(folder_id IS NULL) <> (file_id IS NULL)` enforces this at the
/// database level so the application layer cannot accidentally write a
/// row that's both or neither.
#[derive(Clone, Copy, Debug)]
pub enum ResourceRef {
Folder(Uuid),
File(Uuid),
}
pub struct DeadPropertyStore {
pool: Arc<PgPool>,
}
impl DeadPropertyStore {
pub fn new(pool: Arc<PgPool>) -> Self {
Self { pool }
}
/// Upsert a dead property. `value = None` means an empty XML element.
///
/// The two SQL branches are deliberately kept separate so each
/// ON CONFLICT clause can target the matching partial unique
/// index (`idx_webdav_dead_props_folder_unique` /
/// `idx_webdav_dead_props_file_unique`). A combined upsert would
/// require a non-partial unique index that treats NULL as
/// distinct, which doesn't match the (folder XOR file) shape.
pub async fn set(
&self,
r: ResourceRef,
name: QualifiedName,
value: Option<String>,
) -> Result<(), DomainError> {
match r {
ResourceRef::Folder(folder_id) => {
sqlx::query(
r#"
INSERT INTO storage.webdav_dead_properties
(folder_id, namespace, local_name, value)
VALUES ($1, $2, $3, $4)
ON CONFLICT (folder_id, namespace, local_name)
WHERE folder_id IS NOT NULL
DO UPDATE SET value = EXCLUDED.value, updated_at = CURRENT_TIMESTAMP
"#,
)
.bind(folder_id)
.bind(&name.namespace)
.bind(&name.name)
.bind(&value)
.execute(&*self.pool)
.await
.map_err(|e| {
DomainError::internal_error("DeadPropertyStore", format!("set folder: {e}"))
})?;
}
ResourceRef::File(file_id) => {
sqlx::query(
r#"
INSERT INTO storage.webdav_dead_properties
(file_id, namespace, local_name, value)
VALUES ($1, $2, $3, $4)
ON CONFLICT (file_id, namespace, local_name)
WHERE file_id IS NOT NULL
DO UPDATE SET value = EXCLUDED.value, updated_at = CURRENT_TIMESTAMP
"#,
)
.bind(file_id)
.bind(&name.namespace)
.bind(&name.name)
.bind(&value)
.execute(&*self.pool)
.await
.map_err(|e| {
DomainError::internal_error("DeadPropertyStore", format!("set file: {e}"))
})?;
}
}
Ok(())
}
/// Delete a specific dead property. No-op if not present.
///
/// Filters on the concrete id column (`folder_id = $1` / `file_id = $1`)
/// rather than the old `IS NOT DISTINCT FROM` pair — PostgreSQL cannot
/// serve `IS NOT DISTINCT FROM` from a B-tree index, so every lookup
/// degraded to a sequential scan as the table grew. The `=` shape is
/// served by the partial unique indexes from migration 20260830000001.
/// (Same rationale for `get_all` / `get` / the batched readers below —
/// measured in `benches/DEAD-PROPS.md`.)
pub async fn remove(&self, r: ResourceRef, name: &QualifiedName) -> Result<(), DomainError> {
let (column, id) = split_ref(r);
sqlx::query(&format!(
"DELETE FROM storage.webdav_dead_properties
WHERE {column} = $1
AND namespace = $2
AND local_name = $3",
))
.bind(id)
.bind(&name.namespace)
.bind(&name.name)
.execute(&*self.pool)
.await
.map_err(|e| DomainError::internal_error("DeadPropertyStore", format!("remove: {e}")))?;
Ok(())
}
/// Return all dead properties for the given resource.
pub async fn get_all(
&self,
r: ResourceRef,
) -> Result<Vec<(QualifiedName, Option<String>)>, DomainError> {
let (column, id) = split_ref(r);
let rows = sqlx::query(&format!(
"SELECT namespace, local_name, value
FROM storage.webdav_dead_properties
WHERE {column} = $1",
))
.bind(id)
.fetch_all(&*self.pool)
.await
.map_err(|e| DomainError::internal_error("DeadPropertyStore", format!("get_all: {e}")))?;
Ok(rows.into_iter().map(row_to_prop).collect())
}
/// Batched variant of [`get_all`] for every file in a PROPFIND page:
/// ONE `file_id = ANY($1)` round-trip instead of N sequential queries.
/// Files with no dead properties are simply absent from the map.
pub async fn get_all_for_files(
&self,
file_ids: &[Uuid],
) -> Result<HashMap<Uuid, Vec<(QualifiedName, Option<String>)>>, DomainError> {
self.get_all_batched("file_id", file_ids).await
}
/// Batched variant of [`get_all`] for every subfolder in a PROPFIND page.
pub async fn get_all_for_folders(
&self,
folder_ids: &[Uuid],
) -> Result<HashMap<Uuid, Vec<(QualifiedName, Option<String>)>>, DomainError> {
self.get_all_batched("folder_id", folder_ids).await
}
async fn get_all_batched(
&self,
column: &str,
ids: &[Uuid],
) -> Result<HashMap<Uuid, Vec<(QualifiedName, Option<String>)>>, DomainError> {
if ids.is_empty() {
return Ok(HashMap::new());
}
let rows = sqlx::query(&format!(
"SELECT {column} AS resource_id, namespace, local_name, value
FROM storage.webdav_dead_properties
WHERE {column} = ANY($1)",
))
.bind(ids)
.fetch_all(&*self.pool)
.await
.map_err(|e| {
DomainError::internal_error("DeadPropertyStore", format!("get_all_batched: {e}"))
})?;
let mut map: HashMap<Uuid, Vec<(QualifiedName, Option<String>)>> = HashMap::new();
for row in rows {
let resource_id: Uuid = row.get("resource_id");
map.entry(resource_id).or_default().push(row_to_prop(row));
}
Ok(map)
}
/// Return a specific dead property, or `None` if not stored.
/// Returns `Some(None)` when the property exists with an empty value.
pub async fn get(
&self,
r: ResourceRef,
name: &QualifiedName,
) -> Result<Option<Option<String>>, DomainError> {
let (column, id) = split_ref(r);
let row = sqlx::query(&format!(
"SELECT value FROM storage.webdav_dead_properties
WHERE {column} = $1
AND namespace = $2
AND local_name = $3",
))
.bind(id)
.bind(&name.namespace)
.bind(&name.name)
.fetch_optional(&*self.pool)
.await
.map_err(|e| DomainError::internal_error("DeadPropertyStore", format!("get: {e}")))?;
Ok(row.map(|r| r.get::<Option<String>, _>("value")))
}
}
/// Maps a `ResourceRef` onto the column that stores it plus the id to bind.
/// The column name is one of two compile-time literals — never user input —
/// so interpolating it into the SQL text is safe.
fn split_ref(r: ResourceRef) -> (&'static str, Uuid) {
match r {
ResourceRef::Folder(id) => ("folder_id", id),
ResourceRef::File(id) => ("file_id", id),
}
}
fn row_to_prop(r: sqlx::postgres::PgRow) -> (QualifiedName, Option<String>) {
let namespace: String = r.get("namespace");
let local_name: String = r.get("local_name");
let value: Option<String> = r.get("value");
(QualifiedName::new(namespace, local_name), value)
}
pub fn create_dead_property_store(pool: Arc<PgPool>) -> Arc<DeadPropertyStore> {
Arc::new(DeadPropertyStore::new(pool))
}
@@ -32,6 +32,14 @@ const MAX_LOCK_TIMEOUT_SECS: u64 = 86_400; // 24 hours
pub struct LockEntry {
pub info: LockInfo,
pub path: String,
/// The user who acquired the lock. `None` for entries seeded by
/// unit tests or refresh paths that don't carry a caller (the
/// refresh flow rebuilds from the existing entry without a new
/// caller context, so we preserve whatever was there). RFC 4918
/// §9.11's "MUST be requested by the owner" rule for UNLOCK is
/// enforced by comparing this against the caller in
/// `handle_unlock`.
pub caller_user_id: Option<uuid::Uuid>,
}
/// Per-entry expiration policy for the `by_path` cache.
@@ -106,20 +114,39 @@ impl WebDavLockStore {
/// Attempt to acquire a lock on `path`.
///
/// Returns `Ok(LockEntry)` on success, or `Err(existing)` if the resource
/// is already exclusively locked by a different token.
/// Returns `Ok(LockEntry)` on success, or `Err(existing)` when:
/// - The existing lock is exclusive (blocks any new lock), or
/// - The new lock is exclusive and any lock already exists (RFC 4918 §7.8).
#[allow(clippy::result_large_err)]
pub fn acquire(&self, path: &str, info: LockInfo) -> Result<LockEntry, LockEntry> {
// Check for existing conflicting lock
if let Some(existing) = self.by_path.get(path)
&& existing.info.scope == LockScope::Exclusive
{
return Err(existing);
pub fn acquire(
&self,
path: &str,
info: LockInfo,
caller_user_id: Option<uuid::Uuid>,
) -> Result<LockEntry, LockEntry> {
if let Some(existing) = self.by_path.get(path) {
// Exclusive existing lock → blocks everything.
// New exclusive lock → blocked by any existing lock (shared or exclusive).
if existing.info.scope == LockScope::Exclusive || info.scope == LockScope::Exclusive {
return Err(existing);
}
// Both shared: keep the first holder as the enforcement sentinel in
// `by_path` so releasing a secondary holder cannot clear the lock.
// Register the new token only in the reverse index so UNLOCK works.
let entry = LockEntry {
info,
path: path.to_owned(),
caller_user_id,
};
self.by_token
.insert(entry.info.token.clone(), path.to_owned());
return Ok(entry);
}
let entry = LockEntry {
info,
path: path.to_owned(),
caller_user_id,
};
// `LockExpiry` derives the TTL from `entry.info.timeout` on insert —
@@ -242,6 +269,7 @@ mod tests {
LockEntry {
info: lock_info(token, timeout, LockScope::Exclusive),
path: "/file.txt".to_owned(),
caller_user_id: None,
}
}
@@ -293,7 +321,7 @@ mod tests {
let store = WebDavLockStore::new(16);
let info = lock_info("urn:token-1", Some("Second-600"), LockScope::Exclusive);
let acquired = store.acquire("/a.txt", info).expect("acquire");
let acquired = store.acquire("/a.txt", info, None).expect("acquire");
assert_eq!(acquired.info.token, "urn:token-1");
// Resolvable by both indexes.
@@ -320,12 +348,14 @@ mod tests {
.acquire(
"/a.txt",
lock_info("urn:token-1", Some("Second-600"), LockScope::Exclusive),
None,
)
.expect("first acquire");
let conflict = store.acquire(
"/a.txt",
lock_info("urn:token-2", Some("Second-600"), LockScope::Exclusive),
None,
);
assert!(conflict.is_err());
// The original holder is returned so the caller can report it.
@@ -339,6 +369,7 @@ mod tests {
.acquire(
"/a.txt",
lock_info("urn:token-1", Some("Infinite"), LockScope::Exclusive),
None,
)
.expect("acquire");
@@ -203,7 +203,10 @@ impl WopiDiscoveryService {
for attr in e.attributes().flatten() {
let value = attr
.decode_and_unescape_value(reader.decoder())
.decoded_and_normalized_value(
quick_xml::XmlVersion::Implicit1_0,
reader.decoder(),
)
.map(|value| value.into_owned())
.unwrap_or_else(|_| String::from_utf8_lossy(&attr.value).to_string());
+131 -62
View File
@@ -44,15 +44,22 @@ impl From<ZipError> for DomainError {
}
}
/// Type alias for the fully-async ZIP writer backed by a buffered tokio file.
type AsyncZipWriter = ZipFileWriter<Compat<BufWriter<tokio::fs::File>>>;
/// Fully-async ZIP writer over any buffered tokio sink (temp file for the
/// legacy path, one half of a `tokio::io::duplex` for the streaming path).
type AsyncZipWriter<W> = ZipFileWriter<Compat<BufWriter<W>>>;
/// One planned archive entry, in final ZIP order.
enum ZipPlanEntry {
/// Directory entry (Stored, zero-length body).
Dir(String),
/// File entry: ZIP-relative path + file id to stream from the blob store.
File { zip_path: String, file_id: String },
/// `compression` is picked from the file's MIME type at plan time —
/// `Stored` for already-compressed media (JPEG/MP4/…), `Deflate` otherwise.
File {
zip_path: String,
file_id: String,
compression: Compression,
},
}
/// Message protocol from the prefetch task to the ZIP writer. For each
@@ -74,8 +81,11 @@ const PREFETCH_BUFFER_CHUNKS: usize = 64;
///
/// Uses `async_zip` for fully-async archive creation. Every write (headers,
/// compressed chunk data, central directory) goes through
/// `tokio::io::BufWriter` → `tokio::fs::File`, so **no Tokio worker is ever
/// blocked** by disk I/O or compression.
/// `tokio::io::BufWriter` → `tokio::fs::File`, so no Tokio worker is ever
/// blocked by disk I/O. Deflate itself DOES run inline on the writing task
/// (async_zip compresses inside `poll_write`), which is why entries whose
/// MIME says the content is already compressed are `Stored` instead — that
/// turns the archive hot path from ~1 CPU core per download into CRC + memcpy.
///
/// Archive creation is a 2-stage pipeline: a prefetch task reads file
/// content from the blob store ahead of the writer, so the next file's
@@ -110,6 +120,110 @@ impl ZipService {
folder_id: &str,
folder_name: &str,
) -> Result<NamedTempFile> {
let plan = self.plan_archive(folder_id, folder_name).await?;
// ── Open the temp file + ZIP writer ──────────────────────────────
let temp = NamedTempFile::new().map_err(ZipError::IoError)?;
let tokio_file = tokio::fs::File::create(temp.path())
.await
.map_err(ZipError::IoError)?;
let (tx, mut rx) = tokio::sync::mpsc::channel::<Prefetched>(PREFETCH_BUFFER_CHUNKS);
let _prefetcher = tokio::spawn(Self::prefetch_files(
self.file_service.clone(),
Self::planned_file_ids(&plan),
tx,
));
Self::write_archive(tokio_file, &plan, &mut rx).await?;
Ok(temp)
}
/// Streaming variant: the archive bytes are produced on a spawned task
/// and yielded as they are written — the client's first byte arrives
/// after the first entry starts, not after the whole archive has been
/// built (the temp-file variant's time-to-first-byte grows with folder
/// size; benches/ZIP-STREAM.md). The plan phase still runs inline so
/// planning errors surface as proper HTTP errors; a blob-read error
/// mid-archive can only truncate the stream (no central directory →
/// clients detect the corrupt archive), which is the standard tradeoff
/// for streamed ZIPs.
pub async fn create_folder_zip_stream(
&self,
folder_id: &str,
folder_name: &str,
) -> Result<impl futures::Stream<Item = std::io::Result<bytes::Bytes>> + Send + use<>> {
let plan = self.plan_archive(folder_id, folder_name).await?;
let (writer, reader) = tokio::io::duplex(256 * 1024);
let (tx, mut rx) = tokio::sync::mpsc::channel::<Prefetched>(PREFETCH_BUFFER_CHUNKS);
let _prefetcher = tokio::spawn(Self::prefetch_files(
self.file_service.clone(),
Self::planned_file_ids(&plan),
tx,
));
tokio::spawn(async move {
if let Err(e) = Self::write_archive(writer, &plan, &mut rx).await {
// Dropping the writer EOFs the reader early — the truncated
// archive has no central directory, so clients flag it.
warn!("Streaming ZIP aborted mid-archive: {e}");
}
});
Ok(tokio_util::io::ReaderStream::new(reader))
}
/// File ids of the plan, in archive order (the prefetcher's read list).
fn planned_file_ids(plan: &[ZipPlanEntry]) -> Vec<String> {
plan.iter()
.filter_map(|entry| match entry {
ZipPlanEntry::File { file_id, .. } => Some(file_id.clone()),
ZipPlanEntry::Dir(_) => None,
})
.collect()
}
/// Write every planned entry through a buffered ZIP writer over `sink`,
/// then finalize (central directory + flush). Shared by the temp-file
/// and streaming variants.
async fn write_archive<W: tokio::io::AsyncWrite + Unpin>(
sink: W,
plan: &[ZipPlanEntry],
rx: &mut tokio::sync::mpsc::Receiver<Prefetched>,
) -> Result<()> {
let buf_writer = BufWriter::with_capacity(256 * 1024, sink);
let mut zip = ZipFileWriter::with_tokio(buf_writer);
for entry in plan {
match entry {
ZipPlanEntry::Dir(zip_dir) => {
let dir_entry =
ZipEntryBuilder::new(zip_dir.clone().into(), Compression::Stored);
match zip.write_entry_whole(dir_entry, &[]).await {
Ok(()) => debug!("Folder added to ZIP: {}", zip_dir),
Err(e) => {
warn!("Could not add folder entry (may already exist): {}", e);
}
}
}
ZipPlanEntry::File {
zip_path,
compression,
..
} => {
Self::write_prefetched_file(&mut zip, zip_path, *compression, rx).await?;
}
}
}
let mut compat_writer = zip.close().await.map_err(ZipError::AsyncZipError)?;
compat_writer.close().await.map_err(ZipError::IoError)?;
Ok(())
}
/// Resolve the folder, fetch its subtree (2 bulk queries) and lay out
/// the archive entries in final ZIP order.
async fn plan_archive(&self, folder_id: &str, folder_name: &str) -> Result<Vec<ZipPlanEntry>> {
info!(
"Creating ZIP for folder: {} (ID: {})",
folder_name, folder_id
@@ -183,62 +297,15 @@ impl ZipService {
plan.push(ZipPlanEntry::File {
zip_path: format!("{}{}", zip_dir, file.name),
file_id: file.id.to_string(),
compression: crate::common::mime_detect::zip_entry_compression(
&file.mime_type,
),
});
}
}
}
// ── 5. Open the temp file + ZIP writer ───────────────────────────
let temp = NamedTempFile::new().map_err(ZipError::IoError)?;
let tokio_file = tokio::fs::File::create(temp.path())
.await
.map_err(ZipError::IoError)?;
let buf_writer = BufWriter::with_capacity(256 * 1024, tokio_file);
let mut zip = ZipFileWriter::with_tokio(buf_writer);
// ── 6. Write entries: 2-stage pipeline ───────────────────────────
// The prefetch task reads blob streams for the planned files, in
// order, ahead of the writer — the next file's blob-store latency
// overlaps the current file's deflate. If the writer bails out,
// dropping the receiver makes the prefetcher's next send fail and
// it stops on its own.
let file_ids: Vec<String> = plan
.iter()
.filter_map(|entry| match entry {
ZipPlanEntry::File { file_id, .. } => Some(file_id.clone()),
ZipPlanEntry::Dir(_) => None,
})
.collect();
let (tx, mut rx) = tokio::sync::mpsc::channel::<Prefetched>(PREFETCH_BUFFER_CHUNKS);
let _prefetcher = tokio::spawn(Self::prefetch_files(
self.file_service.clone(),
file_ids,
tx,
));
for entry in &plan {
match entry {
ZipPlanEntry::Dir(zip_dir) => {
let dir_entry =
ZipEntryBuilder::new(zip_dir.clone().into(), Compression::Stored);
match zip.write_entry_whole(dir_entry, &[]).await {
Ok(()) => debug!("Folder added to ZIP: {}", zip_dir),
Err(e) => {
warn!("Could not add folder entry (may already exist): {}", e);
}
}
}
ZipPlanEntry::File { zip_path, .. } => {
Self::write_prefetched_file(&mut zip, zip_path, &mut rx).await?;
}
}
}
// ── 7. Finalize ──────────────────────────────────────────────────
let mut compat_writer = zip.close().await.map_err(ZipError::AsyncZipError)?;
compat_writer.close().await.map_err(ZipError::IoError)?;
Ok(temp)
Ok(plan)
}
/// Prefetch stage: streams each planned file's content from the blob
@@ -282,17 +349,19 @@ impl ZipService {
}
}
/// Writer stage: drains one file's prefetched chunks into a Deflate
/// ZIP entry. Peak memory stays bounded by the channel, independent
/// of individual file sizes.
async fn write_prefetched_file(
zip: &mut AsyncZipWriter,
/// Writer stage: drains one file's prefetched chunks into a ZIP entry
/// (`Stored` for already-compressed media, `Deflate` otherwise — see
/// `entry_compression`). Peak memory stays bounded by the channel,
/// independent of individual file sizes.
async fn write_prefetched_file<W: tokio::io::AsyncWrite + Unpin>(
zip: &mut AsyncZipWriter<W>,
zip_path: &str,
compression: Compression,
rx: &mut tokio::sync::mpsc::Receiver<Prefetched>,
) -> Result<()> {
info!("Adding file to ZIP: {}", zip_path);
let entry = ZipEntryBuilder::new(zip_path.to_string().into(), Compression::Deflate);
let entry = ZipEntryBuilder::new(zip_path.to_string().into(), compression);
let mut entry_writer = zip
.write_entry_stream(entry)
.await