perf: round 12 — auth write-path narrowing, fused quota gate, moka blob-cache index, media single-read, sized listing JSON

Benchmark-gated round (benches/ROUND12.md; every change ships with a
BEFORE/AFTER harness + equivalence gates, one candidate rejected by its
own bench):

DB / query shapes (bench_round12_queries):
- NC sharee search: username-only projection instead of the 21-column row
  (incl. the <=512 KiB avatar) per match, + gin_trgm_ops indexes on
  auth.users for the leading-wildcard ILIKE (4.98x; 54.7x with index).
- Password login: delete the redundant full-row update_user — create_session
  already stamps last_login_at in its own txn (4.45x per login).
- Email-verified stamp: narrow conditional UPDATE (8.9x); OIDC repeat login
  now compares profile state in memory and issues ZERO queries when nothing
  changed (was: full 17-column rewrite per login).
- Refresh rotation: revoke+insert+stamp fused into one transaction via new
  rotate_session port method (1.18x).
- WOPI CheckFileInfo / authorize_wopi_access: require(Read) + get_file +
  check(Update) overlapped with tokio::join!, original result precedence
  (cold 1.34x).
- Upload quota gate: user-envelope + drive-cap checks fused into ONE
  round-trip (check_upload_quotas) — the NC chunked PUT pays this per
  chunk (1.81x, 2 -> 1 queries/chunk); shared verdict evaluators keep
  error shapes byte-identical.

CPU / allocs (bench_round12_micro):
- sized_json: pre-sized listing serialization replacing axum Json's 128 B
  seed + doubling-realloc chain on files/folder-resources/photos/search
  responses (1.40x, 13 -> 2 allocs per 500-row page; byte-identical).
- Security headers: 4 SetResponseHeaderLayer folded into the CSP middleware
  pass (5 layers -> 1; 1.43x per request, -26 allocs; header set gated
  byte-identical incl. 304s).
- Media capture-metadata: single-read extraction — nom-exif now parses the
  buffer kamadak already read (zero-copy Bytes) and videos open once with a
  kind() dispatch; per-image opens 2-3 -> 1 (1.44x warm geomean, 1.6-3.2x
  cold cache; extraction outputs gated identical incl. the MIME-mislabel
  track fallback).
- Chunked-upload session ops: owner gate folded into the operation's own
  DashMap lookup + stack-encoded uuid compare (5 -> 3 lookups, -2 allocs,
  1.28x per chunk).

Blob cache (bench_blob_cache_index + round-3 regression guard):
- CachedBlobBackend index: tokio::sync::Mutex<LruCache> -> moka::sync::Cache
  with byte weigher. The mutex serialized every cached chunk read and scaled
  NEGATIVELY (2.08 -> 1.07 Mops/s from 1 -> 2 readers); moka probes are
  lock-free (2.17x at K=2). Byte budget now enforced by moka (manual
  current_size + collect_evictions machinery deleted); eviction listener
  unlinks size-evicted files only (Replaced entries keep their file —
  gated). Single-flight miss gate unchanged (16 concurrent misses -> 1
  fetch re-verified via the round-3 harness).
- put_blob now populates the cache BEFORE the inner backend consumes the
  source file (the old order failed 100% of the time — local renames,
  S3/Azure delete the source — so the first read after a whole-file put
  re-downloaded from the remote); inner-put failure invalidates the entry.

Frontend (vitest gates):
- List-view thumbnails request the 150px icon rendition instead of 400px
  preview into a 40px slot (~7.1x fewer pixels, ~4-5x fewer bytes per
  thumbnail across list views); grid keeps preview.

Rejected by its own bench (kept as evidence in bench_round12_micro §2):
- Single-pass compression predicate: the monomorphized And-chain already
  costs ~4.6 ns / 0 allocs total; the fused node measured within noise.

New migration: 20260719000000_users_search_trgm.sql (trgm indexes).
Deferred with prepared design: grouped file/grid view virtualization
(single-VirtualRows flatten, the photos pattern) — next round's headline.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01BfidAJD5AHw23jtvBUNamB
This commit is contained in:
Claude
2026-07-19 01:32:00 +00:00
parent a793cd62eb
commit 50eca0627f
33 changed files with 3989 additions and 410 deletions
+41
View File
@@ -131,6 +131,37 @@ pub trait UserStoragePort: Send + Sync + 'static {
include_external: bool,
) -> Result<Vec<User>, DomainError>;
/// Username-only projection of [`search_users`] — same WHERE / ORDER /
/// LIMIT semantics, but skips hydrating the 21-column row (incl. the
/// up-to-512 KiB avatar `image`) when the caller only needs handles.
/// Rows whose username is NULL are returned as `None` so callers can
/// keep the wide flow's post-limit filtering semantics.
async fn search_usernames(
&self,
query: &str,
limit: i64,
include_external: bool,
) -> Result<Vec<Option<String>>, DomainError>;
/// Stamps `email_verified_at = NOW()` iff it is still NULL (idempotent,
/// preserves the first timestamp — the SQL twin of
/// `User::mark_email_verified`). Narrow single-column write; avoids the
/// full-row [`update_user`] (incl. the avatar `image`) on the
/// magic-link redemption path.
async fn mark_email_verified(&self, user_id: Uuid) -> Result<(), DomainError>;
/// OIDC repeat-login profile sync: persists the IdP-provided avatar and
/// stamps `email_verified_at` (guarded, idempotent) in ONE narrow
/// statement. The `IS DISTINCT FROM` guard makes the common case (same
/// avatar, already verified) a zero-write no-op — vs the full 17-column
/// row rewrite this path used to pay per login. `last_login_at` is NOT
/// touched here: session creation stamps it, as on every login path.
async fn sync_oidc_login_profile(
&self,
user_id: Uuid,
image: Option<&str>,
) -> Result<(), DomainError>;
/// Lists users by role (e.g., "admin" or "user")
async fn list_users_by_role(&self, role: &str) -> Result<Vec<User>, DomainError>;
@@ -239,6 +270,16 @@ pub trait SessionStoragePort: Send + Sync + 'static {
/// Creates a new session
async fn create_session(&self, session: Session) -> Result<Session, DomainError>;
/// Refresh-token rotation: revokes `old_session_id` and creates
/// `new_session` in ONE transaction (the refresh path used to pay two
/// full BEGIN/COMMIT round-trip pairs per rotation). Also stamps the
/// user's `last_login_at` exactly like [`create_session`] does.
async fn rotate_session(
&self,
old_session_id: Uuid,
new_session: Session,
) -> Result<Session, DomainError>;
/// Gets a session by refresh token
async fn get_session_by_refresh_token(
&self,
@@ -786,9 +786,14 @@ impl AuthApplicationService {
lc.dispatch_login(&user).await;
}
// Update last login
// Update last login (in-memory only — the DTO below carries it).
// The full-row `update_user` this path used to issue was 100%
// redundant: `create_session` stamps `last_login_at`/`updated_at`
// in its own transaction right below, and nothing re-reads the row
// in between. Dropping it removes one transaction + a 17-column
// rewrite (incl. the up-to-512 KiB avatar) per password login
// (benches/ROUND12.md §2, 4.45x).
user.register_login();
self.user_storage.update_user(user.clone()).await?;
// Generate tokens using the injected token service
let access_token = self.token_service.generate_access_token(&user)?;
@@ -1017,9 +1022,12 @@ impl AuthApplicationService {
// PR 23: clicking the magic-link IS proof of email control —
// stamp the verification (idempotent, preserves the first
// timestamp). Applies to both invitation and login-via-email
// tokens.
// tokens. Narrow single-column write: `last_login_at` is stamped
// by `create_session` below, so the full-row `update_user` this
// path used to issue only ever contributed the verification
// timestamp (benches/ROUND12.md §3, 8.9x).
user.mark_email_verified();
self.user_storage.update_user(user.clone()).await?;
self.user_storage.mark_email_verified(user.id()).await?;
let access_token = self.token_service.generate_access_token(&user)?;
let refresh_token = self.token_service.generate_refresh_token();
@@ -1163,15 +1171,15 @@ impl AuthApplicationService {
));
}
// Revoke current session before issuing the next token in the family
self.session_storage.revoke_session(session.id()).await?;
// Generate new tokens
let access_token = self.token_service.generate_access_token(&user)?;
let new_refresh_token = self.token_service.generate_refresh_token();
// New session inherits the family_id so reuse of any ancestor triggers
// full-family revocation
// full-family revocation. Revoking the old session and inserting the
// new one happen in ONE transaction (`rotate_session`) — this path
// used to pay two BEGIN/COMMIT pairs per refresh, and DAV clients
// rotate constantly (benches/ROUND12.md §4).
let new_session = Session::new(
user.id(),
new_refresh_token.clone(),
@@ -1181,7 +1189,9 @@ impl AuthApplicationService {
session.family_id(),
);
self.session_storage.create_session(new_session).await?;
self.session_storage
.rotate_session(session.id(), new_session)
.await?;
Ok(AuthResponseDto {
user: UserDto::from(user),
@@ -2030,6 +2040,24 @@ impl AuthApplicationService {
Ok(users.into_iter().map(UserDto::from).collect())
}
/// Username-only search for the NC sharee autocomplete: identical
/// predicate / order / limit to [`search_users`], but the repository
/// projects just `username` — no 21-column hydration (incl. the
/// up-to-512 KiB avatar `image`) per matched row, per keystroke
/// (benches/ROUND12.md §1). NULL usernames (email-only signups) are
/// filtered app-side, exactly like the wide flow's post-limit filter.
pub async fn search_sharee_usernames(
&self,
query: &str,
limit: i64,
) -> Result<Vec<String>, DomainError> {
let names = self
.user_storage
.search_usernames(query, limit, false)
.await?;
Ok(names.into_iter().flatten().collect())
}
// ========================================================================
// Admin User Management Methods
// ========================================================================
@@ -2607,6 +2635,15 @@ impl AuthApplicationService {
if let Some(lc) = &self.user_lifecycle {
lc.dispatch_login(&existing_user).await;
}
// Decide BEFORE mutating: the row just fetched already
// carries the stored avatar + verification stamp, so the
// repeat-login common case (same IdP picture, already
// verified) skips the DB entirely — the old shape rewrote
// all 17 columns per login, and even a guarded UPDATE
// would ship the avatar over the wire just to compare it
// (benches/ROUND12.md §3b).
let needs_profile_sync = existing_user.email_verified_at().is_none()
|| existing_user.image() != claims.picture.as_deref();
existing_user.register_login();
existing_user.set_image(claims.picture.clone());
// PR 23: retroactive email verification for OIDC users
@@ -2615,7 +2652,16 @@ impl AuthApplicationService {
// any user reaching this branch has a verified email
// by the IdP's word; stamping is safe and idempotent.
existing_user.mark_email_verified();
self.user_storage.update_user(existing_user.clone()).await?;
// Narrow guarded sync instead of the 17-column row rewrite:
// persists the IdP avatar + the verification stamp only
// when either actually changed; `last_login_at` is stamped
// by `create_session` at the end of this flow
// (benches/ROUND12.md §3).
if needs_profile_sync {
self.user_storage
.sync_oidc_login_profile(existing_user.id(), claims.picture.as_deref())
.await?;
}
existing_user
}
Err(_) => {
+134 -30
View File
@@ -16,6 +16,10 @@ use uuid::Uuid;
* Storage usage is calculated directly from the `storage.files` table
* by summing file sizes for each user (using the `user_id` column).
*/
/// Fused quota-gate row: `(user_used, user_quota, drive_used, drive_quota,
/// drive_found)` — see [`StorageUsageService::check_upload_quotas`].
type QuotaPairRow = (i64, i64, Option<i64>, Option<i64>, bool);
pub struct StorageUsageService {
pool: Arc<PgPool>,
user_repository: Arc<UserPgRepository>,
@@ -372,6 +376,18 @@ impl StorageUsageService {
// first, so this branch fires only on a deleted-drive race.
return Err(DomainError::not_found("Drive", drive_id.to_string()));
};
Self::eval_drive_cap(used, quota, additional_bytes)
}
/// Drive-cap verdict over already-fetched counters. Shared by
/// [`Self::check_drive_quota`] and the fused
/// [`Self::check_upload_quotas`] pair so both produce byte-identical
/// errors.
fn eval_drive_cap(
used: i64,
quota: Option<i64>,
additional_bytes: u64,
) -> Result<(), DomainError> {
let Some(quota) = quota else {
return Ok(()); // unlimited
};
@@ -391,6 +407,123 @@ impl StorageUsageService {
Ok(())
}
/// User-envelope verdict over already-fetched counters. Shared by
/// `check_storage_quota` and the fused [`Self::check_upload_quotas`]
/// pair so both produce byte-identical errors.
fn eval_user_envelope(used: i64, quota: i64, additional_bytes: u64) -> Result<(), DomainError> {
// Quota of 0 means unlimited
if quota <= 0 {
return Ok(());
}
let additional = additional_bytes as i64;
// Case 1: the single file alone exceeds the entire quota
if additional > quota {
let quota_fmt = format_bytes(quota);
let file_fmt = format_bytes(additional);
return Err(DomainError::quota_exceeded(format!(
"File size ({}) exceeds your total storage quota ({})",
file_fmt, quota_fmt
)));
}
// Case 2: the upload would push usage over the quota
if used + additional > quota {
let available = (quota - used).max(0);
let avail_fmt = format_bytes(available);
let file_fmt = format_bytes(additional);
return Err(DomainError::quota_exceeded(format!(
"Not enough storage space. File size: {}, available: {}",
file_fmt, avail_fmt
)));
}
Ok(())
}
/// Fused pre-upload gate: user envelope + drive cap in ONE round-trip.
///
/// Upload entry points used to run `check_storage_quota` then
/// `check_drive_quota` as two serial point reads — and the NC chunked
/// PUT pays that pair on EVERY chunk. One `LEFT JOIN` row carries both
/// counter pairs; verdict precedence (user envelope first, then drive
/// existence, then drive cap) and every error shape are identical to
/// the two-call sequence (benches/ROUND12.md §6, 1.81x).
///
/// Row shape shared with [`Self::check_upload_quotas_by_folder`]:
/// `(user_used, user_quota, drive_used, drive_quota, drive_found)`.
pub async fn check_upload_quotas(
&self,
user_id: Uuid,
drive_id: Uuid,
additional_bytes: u64,
) -> Result<(), DomainError> {
let row: Option<QuotaPairRow> = sqlx::query_as(
r#"
SELECT u.storage_used_bytes, u.storage_quota_bytes,
d.used_bytes, d.quota_bytes, (d.id IS NOT NULL)
FROM auth.users u
LEFT JOIN storage.drives d ON d.id = $2
WHERE u.id = $1
"#,
)
.bind(user_id)
.bind(drive_id)
.fetch_optional(self.pool.as_ref())
.await
.map_err(|e| {
DomainError::internal_error("StorageUsage", format!("upload quota lookup: {e}"))
})?;
let Some((uused, uquota, dused, dquota, drive_found)) = row else {
return Err(DomainError::not_found("User", user_id.to_string()));
};
Self::eval_user_envelope(uused, uquota, additional_bytes)?;
if !drive_found {
return Err(DomainError::not_found("Drive", drive_id.to_string()));
}
Self::eval_drive_cap(dused.unwrap_or(0), dquota, additional_bytes)
}
/// [`Self::check_upload_quotas`] with the drive resolved from a parent
/// folder id — for the REST upload paths, which hold `folder_id`.
/// A missing folder (or a folder whose drive vanished mid-race) maps to
/// `not_found("Folder")`, exactly like `check_drive_quota_by_folder`.
pub async fn check_upload_quotas_by_folder(
&self,
user_id: Uuid,
folder_id: Uuid,
additional_bytes: u64,
) -> Result<(), DomainError> {
let row: Option<QuotaPairRow> = sqlx::query_as(
r#"
SELECT u.storage_used_bytes, u.storage_quota_bytes,
d.used_bytes, d.quota_bytes, (d.id IS NOT NULL)
FROM auth.users u
LEFT JOIN storage.folders f ON f.id = $2
LEFT JOIN storage.drives d ON d.id = f.drive_id
WHERE u.id = $1
"#,
)
.bind(user_id)
.bind(folder_id)
.fetch_optional(self.pool.as_ref())
.await
.map_err(|e| {
DomainError::internal_error("StorageUsage", format!("upload quota lookup: {e}"))
})?;
let Some((uused, uquota, dused, dquota, drive_found)) = row else {
return Err(DomainError::not_found("User", user_id.to_string()));
};
Self::eval_user_envelope(uused, uquota, additional_bytes)?;
if !drive_found {
return Err(DomainError::not_found("Folder", folder_id.to_string()));
}
Self::eval_drive_cap(dused.unwrap_or(0), dquota, additional_bytes)
}
/// Same as [`Self::check_drive_quota`] but resolves the drive id
/// from a parent folder id. Mirrors
/// [`Self::add_drive_storage_usage_delta_by_folder`] so the upload
@@ -563,36 +696,7 @@ impl StorageUsagePort for StorageUsageService {
// Narrow 2-column read — the full user row carries the up-to-512 KiB
// avatar `image` column, paid on every upload quota check otherwise.
let (used, quota) = self.user_repository.get_storage_usage(user_id).await?;
// Quota of 0 means unlimited
if quota <= 0 {
return Ok(());
}
let additional = additional_bytes as i64;
// Case 1: the single file alone exceeds the entire quota
if additional > quota {
let quota_fmt = format_bytes(quota);
let file_fmt = format_bytes(additional);
return Err(DomainError::quota_exceeded(format!(
"File size ({}) exceeds your total storage quota ({})",
file_fmt, quota_fmt
)));
}
// Case 2: the upload would push usage over the quota
if used + additional > quota {
let available = (quota - used).max(0);
let avail_fmt = format_bytes(available);
let file_fmt = format_bytes(additional);
return Err(DomainError::quota_exceeded(format!(
"Not enough storage space. File size: {}, available: {}",
file_fmt, avail_fmt
)));
}
Ok(())
Self::eval_user_envelope(used, quota, additional_bytes)
}
async fn get_user_storage_info(&self, user_id: Uuid) -> Result<(i64, i64), DomainError> {
@@ -327,6 +327,77 @@ impl SessionStoragePort for SessionPgRepository {
.map_err(DomainError::from)
}
/// Revoke + insert + last-login stamp in ONE transaction — the refresh
/// rotation used to pay two full BEGIN/COMMIT round-trip pairs
/// (`revoke_session` then `create_session`) per token refresh.
async fn rotate_session(
&self,
old_session_id: Uuid,
new_session: Session,
) -> Result<Session, DomainError> {
let session_clone = new_session.clone();
with_transaction(&self.pool, "rotate_session", |tx| {
Box::pin(async move {
sqlx::query("UPDATE auth.sessions SET revoked = true WHERE id = $1")
.bind(old_session_id)
.execute(&mut **tx)
.await
.map_err(Self::map_sqlx_error)?;
sqlx::query(
r#"
INSERT INTO auth.sessions (
id, user_id, refresh_token, expires_at,
ip_address, user_agent, created_at, revoked, family_id
) VALUES (
$1, $2, $3, $4, $5, $6, $7, $8, $9
)
"#,
)
.bind(session_clone.id())
.bind(session_clone.user_id())
.bind(session_clone.refresh_token())
.bind(session_clone.expires_at())
.bind(session_clone.ip_address())
.bind(session_clone.user_agent())
.bind(session_clone.created_at())
.bind(session_clone.is_revoked())
.bind(session_clone.family_id())
.execute(&mut **tx)
.await
.map_err(Self::map_sqlx_error)?;
sqlx::query(
r#"
UPDATE auth.users
SET last_login_at = NOW(), updated_at = NOW()
WHERE id = $1
"#,
)
.bind(session_clone.user_id())
.execute(&mut **tx)
.await
.map_err(|e| {
tracing::warn!(
"Could not update last_login_at for user {}: {}",
session_clone.user_id(),
e
);
SessionRepositoryError::DatabaseError(format!(
"Session rotated but could not update last_login_at: {}",
e
))
})?;
Ok(session_clone)
}) as BoxFuture<'_, SessionRepositoryResult<Session>>
})
.await
.map_err(DomainError::from)?;
Ok(new_session)
}
async fn get_session_by_refresh_token(
&self,
refresh_token: &str,
@@ -1067,6 +1067,81 @@ impl UserStoragePort for UserPgRepository {
.map_err(DomainError::from)
}
async fn search_usernames(
&self,
query: &str,
limit: i64,
include_external: bool,
) -> Result<Vec<Option<String>>, DomainError> {
// Same predicate / order / limit as `search_users`, username-only
// projection — the sharee autocomplete path reads nothing else, and
// the wide row drags the avatar `image` per matched user.
let pattern = format!("%{}%", query);
let rows = sqlx::query(
r#"
SELECT username
FROM auth.users
WHERE (username ILIKE $1 OR email ILIKE $1)
AND ($3 OR is_external = FALSE)
ORDER BY username
LIMIT $2
"#,
)
.bind(&pattern)
.bind(limit)
.bind(include_external)
.fetch_all(&*self.pool)
.await
.map_err(Self::map_sqlx_error)
.map_err(DomainError::from)?;
Ok(rows.into_iter().map(|row| row.get("username")).collect())
}
async fn mark_email_verified(&self, user_id: Uuid) -> Result<(), DomainError> {
// SQL twin of `User::mark_email_verified` — stamps once, keeps the
// first timestamp, and touches only the two columns involved.
sqlx::query(
r#"
UPDATE auth.users
SET email_verified_at = NOW(), updated_at = NOW()
WHERE id = $1 AND email_verified_at IS NULL
"#,
)
.bind(user_id)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)
.map_err(DomainError::from)?;
Ok(())
}
async fn sync_oidc_login_profile(
&self,
user_id: Uuid,
image: Option<&str>,
) -> Result<(), DomainError> {
// `IS DISTINCT FROM` guard (the `update_storage_usage` pattern): the
// common repeat-login case — same IdP avatar, already verified —
// writes nothing at all (no dead tuple, no WAL).
sqlx::query(
r#"
UPDATE auth.users
SET image = $2,
email_verified_at = COALESCE(email_verified_at, NOW()),
updated_at = NOW()
WHERE id = $1
AND (image IS DISTINCT FROM $2 OR email_verified_at IS NULL)
"#,
)
.bind(user_id)
.bind(image)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)
.map_err(DomainError::from)?;
Ok(())
}
async fn list_users_by_role(&self, role: &str) -> Result<Vec<User>, DomainError> {
UserRepository::list_users_by_role(self, role)
.await
+117 -225
View File
@@ -10,12 +10,9 @@
use std::path::{Path, PathBuf};
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use bytes::Bytes;
use dashmap::DashMap;
use lru::LruCache;
use std::num::NonZeroUsize;
use tokio::fs;
use tokio::io::{AsyncReadExt, AsyncSeekExt, AsyncWriteExt};
use tokio::sync::Mutex;
@@ -52,12 +49,22 @@ 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
@@ -67,26 +74,38 @@ pub struct CachedBlobBackend {
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)
}
}
@@ -97,7 +116,6 @@ 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?;
@@ -130,14 +148,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,
@@ -152,22 +169,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(),
inflight: self.inflight.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)
}
}
})
}
@@ -176,18 +199,10 @@ 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(),
inflight: self.inflight.clone(),
};
Box::pin(async move {
let size = inner.put_blob_from_bytes(&hash, data.clone()).await?;
self_ref.cache_bytes_write_through(hash, &data).await;
let size = self.inner.put_blob_from_bytes(&hash, data.clone()).await?;
self.cache_bytes_write_through(hash, &data).await;
Ok(size)
})
}
@@ -202,20 +217,13 @@ 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(),
inflight: self.inflight.clone(),
};
Box::pin(async move {
let size = inner
let size = self
.inner
.put_blob_from_bytes_unsynced(&hash, data.clone())
.await?;
self_ref.cache_bytes_write_through(hash, &data).await;
self.cache_bytes_write_through(hash, &data).await;
Ok(size)
})
}
@@ -235,40 +243,24 @@ 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();
let inflight = self.inflight.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 (single-flight), spool to cache
let self_ref = CachedRef {
cache_dir,
max_cache_bytes,
index: index.clone(),
current_size: current_size.clone(),
inflight,
};
let dest = self_ref
.fetch_and_cache_singleflight(&hash, &*inner, &cached)
.await?;
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}"))
})?;
@@ -285,18 +277,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();
let inflight = self.inflight.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
@@ -309,24 +294,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 (single-flight: a
// player's parallel cold Range probes coalesce onto ONE remote
// download), then serve the range locally.
let self_ref = CachedRef {
cache_dir,
max_cache_bytes,
index: index.clone(),
current_size: current_size.clone(),
inflight,
};
let dest = self_ref
.fetch_and_cache_singleflight(&hash, &*inner, &cached)
.await?;
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}")))?;
@@ -345,19 +320,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(())
})
}
@@ -366,18 +335,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
})
}
@@ -385,23 +349,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
})
}
@@ -410,19 +368,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)
@@ -446,27 +403,13 @@ impl BlobStorageBackend for CachedBlobBackend {
}
}
// ── Helper struct for owned references in async closures ───────────
/// 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>,
inflight: Arc<DashMap<String, Arc<Mutex<()>>>>,
}
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"))
}
// ── Cache internals (miss path + population) ───────────────────────
impl CachedBlobBackend {
/// Best-effort write-through cache population shared by both blob-bytes
/// PUT paths. Deliberately no eviction sweep here — the byte budget is
/// enforced on read-miss inserts (`insert_into_cache_static`), matching
/// the historical write-path behavior.
/// 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) {
let dest = self.cached_path(&hash);
if let Some(parent) = dest.parent() {
@@ -474,22 +417,17 @@ impl CachedRef {
}
let _ = fs::write(&dest, data).await;
let data_len = data.len() as u64;
let mut idx = self.index.lock().await;
if let Some(old) = idx.put(hash, CacheEntry { size: data_len }) {
self.current_size.fetch_sub(old.size, Ordering::Relaxed);
}
self.current_size.fetch_add(data_len, Ordering::Relaxed);
self.index.insert(hash, CacheEntry { size: data_len });
}
/// Single-flight wrapper around [`Self::fetch_and_cache_static`]: the
/// first caller for a hash becomes the leader and downloads; concurrent
/// 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,
inner: &dyn BlobStorageBackend,
cached: &Path,
) -> Result<PathBuf, DomainError> {
let gate = self
@@ -501,42 +439,18 @@ impl CachedRef {
// 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.lock().await.get(hash).is_some() && fs::metadata(cached).await.is_ok() {
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_static(hash, inner).await;
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
}
/// 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(
&self,
hash: &str,
source_path: &Path,
) -> Result<(), DomainError> {
async fn insert_into_cache(&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| {
@@ -553,29 +467,14 @@ 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?;
let dest = self.cached_path(hash);
if let Some(parent) = dest.parent() {
@@ -630,17 +529,10 @@ impl CachedRef {
));
}
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);
}
self.current_size.fetch_add(total, 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: 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,
@@ -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
+7 -1
View File
@@ -864,7 +864,13 @@ impl FileHandler {
}
tracing::info!("Found {} files", files.len());
let mut resp = (StatusCode::OK, Json(files)).into_response();
// Pre-sized serialization — this listing is unbounded (no
// page cap), the axum Json 128-byte seed reallocs ~11 times
// on a big folder (benches/ROUND12.md §M1).
let mut resp = crate::interfaces::api::sized_json::sized_json(
64 + files.len() * crate::interfaces::api::sized_json::EST_ROW_BYTES,
&files,
);
resp.headers_mut()
.insert(header::ETAG, header::HeaderValue::from_str(&etag).unwrap());
resp
@@ -551,11 +551,15 @@ pub async fn list_folder_resources(
})
.collect();
(
StatusCode::OK,
Json(FolderResourcesDto::with_cursor(items, next_cursor)),
)
.into_response()
{
// Pre-sized serialization (benches/ROUND12.md §M1).
let body = FolderResourcesDto::with_cursor(items, next_cursor);
crate::interfaces::api::sized_json::sized_json(
128 + body.items.len()
* crate::interfaces::api::sized_json::EST_WRAPPED_ROW_BYTES,
&body,
)
}
}
Err(e) => AppError::from(e).into_response(),
}
@@ -119,7 +119,11 @@ pub async fn list_photos(
})
.collect();
let mut response = Json(&dtos).into_response();
// Pre-sized serialization (benches/ROUND12.md §M1).
let mut response = crate::interfaces::api::sized_json::sized_json(
64 + dtos.len() * crate::interfaces::api::sized_json::EST_WRAPPED_ROW_BYTES,
&dtos,
);
{
let h = response.headers_mut();
h.insert(header::ETAG, header::HeaderValue::from_str(&etag).unwrap());
+16 -2
View File
@@ -84,7 +84,14 @@ impl SearchHandler {
results.files.len(),
results.folders.len()
);
(StatusCode::OK, Json(&*results)).into_response()
{
// Pre-sized serialization (benches/ROUND12.md §M1).
let rows = results.files.len() + results.folders.len();
crate::interfaces::api::sized_json::sized_json(
256 + rows * crate::interfaces::api::sized_json::EST_WRAPPED_ROW_BYTES,
&*results,
)
}
}
Err(err) => {
error!("Search error: {}", err);
@@ -125,7 +132,14 @@ impl SearchHandler {
results.files.len(),
results.folders.len()
);
(StatusCode::OK, Json(&*results)).into_response()
{
// Pre-sized serialization (benches/ROUND12.md §M1).
let rows = results.files.len() + results.folders.len();
crate::interfaces::api::sized_json::sized_json(
256 + rows * crate::interfaces::api::sized_json::EST_WRAPPED_ROW_BYTES,
&*results,
)
}
}
Err(err) => {
error!("Search error: {}", err);
+72 -59
View File
@@ -83,14 +83,21 @@ pub struct CheckFileInfoResponse {
/// structured `audit` line on denial internally, so ops sees the real
/// reason without the attacker being able to distinguish "gone" from
/// "revoked".
/// Shared id parsing for the WOPI authz paths: a malformed caller sub is a
/// bad token (401), a malformed file id can't exist (404, anti-enum).
fn parse_wopi_ids(caller_sub: &str, file_id: &str) -> Result<(uuid::Uuid, uuid::Uuid), StatusCode> {
let caller_uuid = uuid::Uuid::parse_str(caller_sub).map_err(|_| StatusCode::UNAUTHORIZED)?;
let file_uuid = uuid::Uuid::parse_str(file_id).map_err(|_| StatusCode::NOT_FOUND)?;
Ok((caller_uuid, file_uuid))
}
async fn require_wopi_perm(
authz: &PgAclEngine,
caller_sub: &str,
file_id: &str,
perm: Permission,
) -> Result<(uuid::Uuid, uuid::Uuid), StatusCode> {
let caller_uuid = uuid::Uuid::parse_str(caller_sub).map_err(|_| StatusCode::UNAUTHORIZED)?;
let file_uuid = uuid::Uuid::parse_str(file_id).map_err(|_| StatusCode::NOT_FOUND)?;
let (caller_uuid, file_uuid) = parse_wopi_ids(caller_sub, file_id)?;
authz
.require(Subject::User(caller_uuid), perm, Resource::File(file_uuid))
.await
@@ -118,25 +125,52 @@ async fn check_file_info(
// Redemption-time authz: even with a valid token, the caller must
// still hold Read on this file. Catches revoked-grant-mid-session.
if let Err(status) = require_wopi_perm(
state.app_state.authorization.as_ref(),
&claims.sub,
&file_id,
Permission::Read,
)
.await
{
return status.into_response();
//
// The Read gate, the metadata fetch and the Update probe are three
// independent lookups keyed only off (caller, file) — overlapped with
// `tokio::join!` (benches/ROUND12.md §5). Results are evaluated in the
// original precedence: Read gate first, then file existence.
let (caller_uuid, file_uuid) = match parse_wopi_ids(&claims.sub, &file_id) {
Ok(ids) => ids,
Err(status) => return status.into_response(),
};
let authz = state.app_state.authorization.as_ref();
let (read_gate, file, can_write_now) = tokio::join!(
authz.require(
Subject::User(caller_uuid),
Permission::Read,
Resource::File(file_uuid)
),
state
.app_state
.applications
.file_retrieval_service
.get_file(&file_id),
// `user_can_write` = actual current Update permission ∧ token's
// can_write flag. If the caller's Update was revoked since the
// token was minted (e.g. their grant was downgraded from Editor
// to Viewer), the editor sees the file as read-only and won't
// even attempt PutFile. The stricter `require_wopi_perm(Update)`
// in put_file is the actual gate; this field is a UI hint.
async {
if claims.can_write {
authz
.check(
Subject::User(caller_uuid),
Permission::Update,
Resource::File(file_uuid),
)
.await
.unwrap_or(false)
} else {
false
}
}
);
if read_gate.is_err() {
return StatusCode::NOT_FOUND.into_response();
}
// Fetch file metadata
let file = match state
.app_state
.applications
.file_retrieval_service
.get_file(&file_id)
.await
{
let file = match file {
Ok(f) => f,
Err(_) => return StatusCode::NOT_FOUND.into_response(),
};
@@ -146,24 +180,6 @@ async fn check_file_info(
.map(|dt| dt.to_rfc3339())
.unwrap_or_default();
// `user_can_write` = actual current Update permission ∧ token's
// can_write flag. If the caller's Update was revoked since the
// token was minted (e.g. their grant was downgraded from Editor
// to Viewer), the editor sees the file as read-only and won't
// even attempt PutFile. The stricter `require_wopi_perm(Update)`
// in put_file is the actual gate; this field is a UI hint.
let can_write_now = claims.can_write
&& state
.app_state
.authorization
.check(
Subject::User(uuid::Uuid::parse_str(&claims.sub).unwrap_or(uuid::Uuid::nil())),
Permission::Update,
Resource::File(uuid::Uuid::parse_str(&file_id).unwrap_or(uuid::Uuid::nil())),
)
.await
.unwrap_or(false);
let response = CheckFileInfoResponse {
base_file_name: file.name.clone(),
// WOPI's `OwnerId` field is required. Post-D7 the DTO no
@@ -550,34 +566,31 @@ async fn authorize_wopi_access<S: FileRetrievalUseCase>(
) -> Result<(crate::application::dtos::file_dto::FileDto, bool), StatusCode> {
let file_uuid = uuid::Uuid::parse_str(file_id).map_err(|_| StatusCode::NOT_FOUND)?;
// Step 1 — Read is required to even open the file.
authz
.require(
// The Read gate (step 1), the metadata fetch and the Update probe
// (step 2) are independent — overlapped with `tokio::join!`
// (benches/ROUND12.md §5); results evaluated in the original order.
//
// Step 2 rationale — can_write reflects real Update, not the client's
// action-string. `check` returns bool without throwing; failure
// just means the caller lacks Update, so we degrade the token to
// read-only. Deliberately no `require` there — a Viewer opening
// the file is legitimate; only the write claim is suppressed.
let (read_gate, file, has_update) = tokio::join!(
authz.require(
Subject::User(caller_id),
Permission::Read,
Resource::File(file_uuid),
)
.await
.map_err(|_| StatusCode::NOT_FOUND)?;
let file = file_retrieval
.get_file(file_id)
.await
.map_err(|_| StatusCode::NOT_FOUND)?;
// Step 2 — can_write reflects real Update, not the client's
// action-string. `check` returns bool without throwing; failure
// just means the caller lacks Update, so we degrade the token to
// read-only. Deliberately no `require` here — a Viewer opening
// the file is legitimate; only the write claim is suppressed.
let has_update = authz
.check(
),
file_retrieval.get_file(file_id),
authz.check(
Subject::User(caller_id),
Permission::Update,
Resource::File(file_uuid),
)
.await
.unwrap_or(false);
);
read_gate.map_err(|_| StatusCode::NOT_FOUND)?;
let file = file.map_err(|_| StatusCode::NOT_FOUND)?;
let has_update = has_update.unwrap_or(false);
// Step 3 — allow explicit view-mode downgrade for Editors.
let can_write = has_update && requested_action != "view";
+1
View File
@@ -2,6 +2,7 @@ pub mod cookie_auth;
pub mod deserializer;
pub mod handlers;
pub mod routes;
pub mod sized_json;
pub use routes::create_api_routes;
pub use routes::create_health_routes;
+54
View File
@@ -0,0 +1,54 @@
//! Pre-sized JSON responses for listing endpoints.
//!
//! `axum::Json` serializes into a `BytesMut::with_capacity(128)` — a 500-row
//! listing grows that seed through ~11 doubling reallocations, memcpy-ing
//! ~1.3× the payload on every hot listing response (files, folder
//! resources, photos timeline, search). `sized_json` serializes into one
//! right-sized `Vec` instead: 2 allocations total and no copy chain
//! (benches/ROUND12.md §M1, 1.40x / −11 allocs on a 500-row page).
//!
//! The per-row estimates are calibrated against the serialized DTOs (a
//! realistic `FileDto` row measures ~380 B). Underestimates cost one extra
//! doubling — still far better than the 128-byte seed; overestimates waste
//! transient capacity only (the buffer is freed after the response).
use axum::http::{HeaderValue, StatusCode, header};
use axum::response::{IntoResponse, Response};
use bytes::Bytes;
use serde::Serialize;
/// Serialized size estimate for one file/folder row (FileDto ≈ 380 B).
pub const EST_ROW_BYTES: usize = 384;
/// Serialized size estimate for one wrapped resource row (PhotoDto /
/// FolderResourcesDto items carry a FileDto plus wrapper fields).
pub const EST_WRAPPED_ROW_BYTES: usize = 448;
/// Serialize `value` into a single pre-sized buffer and wrap it as an
/// `application/json` response — drop-in for `Json(value).into_response()`
/// (byte-identical body, gated in `bench_round12_micro` §1), minus the
/// doubling-realloc chain.
pub fn sized_json<T: Serialize>(estimated_bytes: usize, value: &T) -> Response {
let mut buf = Vec::with_capacity(estimated_bytes.max(128));
match serde_json::to_writer(&mut buf, value) {
Ok(()) => (
StatusCode::OK,
[(
header::CONTENT_TYPE,
HeaderValue::from_static("application/json"),
)],
Bytes::from(buf),
)
.into_response(),
// Mirror axum's Json error arm: 500 + plain-text serializer error.
Err(err) => (
StatusCode::INTERNAL_SERVER_ERROR,
[(
header::CONTENT_TYPE,
HeaderValue::from_static("text/plain; charset=utf-8"),
)],
err.to_string(),
)
.into_response(),
}
}
+11 -10
View File
@@ -342,20 +342,21 @@ pub async fn handle_sharees_search(
None => return sharees_response(vec![]).into_response(),
};
// SQL-level ILIKE search with limit — avoids loading all users into memory.
let users = auth_service
.search_users(&search, 26)
// SQL-level ILIKE search with limit — avoids loading all users into
// memory. Username-only projection: the wide `search_users` row drags
// the up-to-512 KiB avatar `image` per matched user, per keystroke
// (benches/ROUND12.md §1). NULL-username (email-only signup) rows are
// already filtered by the service, preserving the old post-limit
// filtering semantics.
let usernames = auth_service
.search_sharee_usernames(&search, 26)
.await
.unwrap_or_default();
// Skip users with no claimed username — NC sharees autocomplete relies
// on a username being typeable; users still on the email-only signup
// path can't be addressed here. Also skip self (don't suggest sharing
// with yourself).
let matches: Vec<serde_json::Value> = users
// Skip self (don't suggest sharing with yourself).
let matches: Vec<serde_json::Value> = usernames
.into_iter()
.filter_map(|u| {
let handle = u.username.clone()?;
.filter_map(|handle| {
if handle.as_str() == &*user.username {
return None;
}
+4 -5
View File
@@ -7,7 +7,6 @@ use std::sync::Arc;
use uuid::Uuid;
use crate::application::ports::file_ports::FileUploadUseCase;
use crate::application::ports::storage_ports::StorageUsagePort;
use crate::common::di::AppState;
use crate::common::mime_detect::filename_from_path;
use crate::interfaces::errors::AppError;
@@ -52,10 +51,10 @@ async fn refuse_if_over_quota(
// remains authoritative.
return Ok(());
};
svc.check_storage_quota(user_id, additional)
.await
.map_err(AppError::from)?;
svc.check_drive_quota(drive_id, additional)
// Fused single round-trip (user envelope + drive cap) — this gate runs
// on EVERY chunk PUT, and the serial pair cost two point reads per
// chunk (benches/ROUND12.md §6). Verdict precedence unchanged.
svc.check_upload_quotas(user_id, drive_id, additional)
.await
.map_err(AppError::from)
}
+27 -19
View File
@@ -24,7 +24,6 @@ use oxicloud::access_log;
use oxicloud::interfaces::middleware::trace_span::{ClientIpMakeSpan, UuidRequestId};
use tower_http::limit::RequestBodyLimitLayer;
use tower_http::request_id::{PropagateRequestIdLayer, SetRequestIdLayer};
use tower_http::set_header::SetResponseHeaderLayer;
use tower_http::trace::TraceLayer;
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};
@@ -918,11 +917,37 @@ async fn run() -> Result<(), Box<dyn std::error::Error>> {
// • form-action 'https:': the WOPI office editor is launched by POSTing a
// token form to a cross-origin, admin-configured Collabora/OnlyOffice
// host. Mirrors the SPA meta policy in frontend/svelte.config.js.
// The four static security headers ride in the same response pass —
// they used to be four separate `SetResponseHeaderLayer`s stacked on
// top of this middleware (5 tower layers per response). Folding them
// here measured 1.43x per request / −26 allocs with a byte-identical
// header set, including on 304s (benches/ROUND12.md §M3). They are
// inserted BEFORE the 304 early-return below because the standalone
// layers stamped 304s too.
async fn content_security_policy(
req: axum::extract::Request,
next: axum::middleware::Next,
) -> axum::response::Response {
let mut res = next.run(req).await;
{
let h = res.headers_mut();
h.insert(
HeaderName::from_static("x-content-type-options"),
HeaderValue::from_static("nosniff"),
);
h.insert(
HeaderName::from_static("x-frame-options"),
HeaderValue::from_static("DENY"),
);
h.insert(
HeaderName::from_static("referrer-policy"),
HeaderValue::from_static("strict-origin-when-cross-origin"),
);
h.insert(
HeaderName::from_static("permissions-policy"),
HeaderValue::from_static("camera=(), microphone=(), geolocation=()"),
);
}
// A 304 Not Modified carries no entity headers (no Content-Type) since
// there's no body — `is_html` would read `None` and misclassify it as
// "not html", attaching the strict headerless CSP below. Browsers merge
@@ -982,24 +1007,7 @@ async fn run() -> Result<(), Box<dyn std::error::Error>> {
res
}
app = app
.layer(axum::middleware::from_fn(content_security_policy))
.layer(SetResponseHeaderLayer::overriding(
HeaderName::from_static("x-content-type-options"),
HeaderValue::from_static("nosniff"),
))
.layer(SetResponseHeaderLayer::overriding(
HeaderName::from_static("x-frame-options"),
HeaderValue::from_static("DENY"),
))
.layer(SetResponseHeaderLayer::overriding(
HeaderName::from_static("referrer-policy"),
HeaderValue::from_static("strict-origin-when-cross-origin"),
))
.layer(SetResponseHeaderLayer::overriding(
HeaderName::from_static("permissions-policy"),
HeaderValue::from_static("camera=(), microphone=(), geolocation=()"),
));
app = app.layer(axum::middleware::from_fn(content_security_policy));
// Warn once at startup if auth cookies are not Secure.
// HttpOnly + SameSite protection is nullified over plain HTTP because tokens