perf: round 10 — auth alloc purge, parent-herd batching, query-shape pack, NC 304s
Benchmark-gated (benches/ROUND10.md; every change carries a BEFORE/AFTER harness with equivalence/safety gates — two designs were rejected or rewritten by their own benches before adoption): - Auth hot path: TokenClaims/CurrentUser display fields to Arc<str>, role to inline SmolStr end-to-end (Bearer, cookie, Basic-auth cache) — 4→1 allocs per authenticated request, 3→0 per warm DAV request; JWT Encoding/Decoding/Validation built once. - Cold shared-album herd: leader-inline parent batching in PgAclEngine (+ cascade try_get_with single-flight) — 100→2 parent queries per 100-thumb cold herd, herd wall 1.9x, sequential + warm paths unchanged, all ROUND8/9 safety gates plus new herd-equivalence gates. - Query-shape pack: share download double-fetch 2→1 (2.18x), contact-group COUNT(*) 14.9x, save_faces UNNEST 3.9x, playlist reorder UNNEST 63.7x (now atomic), search files∥folders join! 1.45x, move drive-lookup join! 2.14x, trash partial (drive_id, trashed_at) indexes, CalDAV event-gate narrow read, favorites/recents binary-decode port, dead count_files removed. - NC surface: preview + avatar honour If-None-Match (e2e: 5 KB and 197 KB → 0 bytes per revalidation), avatar WebP→PNG transcode memoised, PROPFIND/trashbin integer+date emits on stack formatters, folder-header enrichment join!, chunk-PUT retry stat folded into create_new open. - common::fmt integer rendering rewritten on the std 2-digit LUT after the round's own bench caught the div-loop losing to to_string (16.1 ns vs 22.5; speeds every prior-round call site). - Micro-pack: WebDAV scope probe borrow-only, ShareService base_url snapshot, cookie_secure OnceLock, Arc'd AES-GCM cipher, stack request-id, tantivy analyzer clone dropped. - SPA: search stale-guard + AbortController (10→1 completed round-trips, stale-clobber gone), getFolder in-flight dedup, gridColumns matchMedia hoist (10k→0 style reads). Backend: cargo fmt + clippy -D warnings clean, 524 tests green. Frontend: npm run check clean, 301 vitest green. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_018DdM7V7M3QPW7HEHg3gLov
This commit is contained in:
@@ -391,6 +391,18 @@ impl CalendarStoragePort for CalendarStorageAdapter {
|
||||
Ok(CalendarEventDto::from(event))
|
||||
}
|
||||
|
||||
async fn calendar_id_for_event(&self, event_id: &str) -> Result<String, DomainError> {
|
||||
let uuid = Uuid::parse_str(event_id).map_err(|_| {
|
||||
DomainError::new(ErrorKind::InvalidInput, "Event", "Invalid event ID format")
|
||||
})?;
|
||||
|
||||
let calendar_id = self
|
||||
.event_repository
|
||||
.find_calendar_id_by_event_id(&uuid)
|
||||
.await?;
|
||||
Ok(calendar_id.to_string())
|
||||
}
|
||||
|
||||
async fn find_event_by_ical_uid(
|
||||
&self,
|
||||
calendar_id: &str,
|
||||
|
||||
@@ -230,6 +230,12 @@ impl ContactStoragePort for ContactStorageAdapter {
|
||||
.await
|
||||
}
|
||||
|
||||
async fn count_contacts_in_group(&self, group_id: &Uuid) -> Result<i64, DomainError> {
|
||||
self.contact_group_repository
|
||||
.count_contacts_in_group(group_id)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn get_groups_for_contact(
|
||||
&self,
|
||||
contact_id: &Uuid,
|
||||
|
||||
@@ -213,6 +213,17 @@ impl CalendarEventRepository for CalendarEventPgRepository {
|
||||
Ok(events)
|
||||
}
|
||||
|
||||
async fn find_calendar_id_by_event_id(&self, id: &Uuid) -> CalendarEventRepositoryResult<Uuid> {
|
||||
sqlx::query_scalar("SELECT calendar_id FROM caldav.calendar_events WHERE id = $1")
|
||||
.bind(id)
|
||||
.fetch_optional(&*self.pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
DomainError::database_error(format!("Failed to get event calendar id: {}", e))
|
||||
})?
|
||||
.ok_or_else(|| DomainError::not_found("Calendar Event", id.to_string()))
|
||||
}
|
||||
|
||||
async fn find_event_by_id(&self, id: &Uuid) -> CalendarEventRepositoryResult<CalendarEvent> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
|
||||
@@ -177,6 +177,20 @@ impl ContactGroupRepository for ContactGroupPgRepository {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn count_contacts_in_group(&self, group_id: &Uuid) -> ContactRepositoryResult<i64> {
|
||||
sqlx::query_scalar("SELECT COUNT(*) FROM carddav.group_memberships WHERE group_id = $1")
|
||||
.bind(group_id)
|
||||
.fetch_one(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
DomainError::new(
|
||||
ErrorKind::InternalError,
|
||||
"ContactGroup",
|
||||
format!("Failed to count contacts in group: {}", e),
|
||||
)
|
||||
})
|
||||
}
|
||||
|
||||
async fn get_contacts_in_group(
|
||||
&self,
|
||||
group_id: &Uuid,
|
||||
|
||||
@@ -115,29 +115,70 @@ impl FaceRepository for FacePgRepository {
|
||||
if faces.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
let mut tx = self.pool.begin().await.map_err(|e| db_err("begin", e))?;
|
||||
// One multi-row INSERT over parallel UNNEST arrays instead of one
|
||||
// round-trip per face — a group photo yields many faces per indexed
|
||||
// image. The `bbox` float4[] can't ride an array-of-arrays through
|
||||
// unnest (PG flattens), so its 4 components travel as 4 parallel
|
||||
// arrays and are reassembled server-side. A single statement is
|
||||
// atomic on its own; the per-row transaction wrapper is gone.
|
||||
let n = faces.len();
|
||||
let mut ids = Vec::with_capacity(n);
|
||||
let mut file_ids = Vec::with_capacity(n);
|
||||
let mut user_ids = Vec::with_capacity(n);
|
||||
let mut person_ids: Vec<Option<Uuid>> = Vec::with_capacity(n);
|
||||
let (mut bx, mut by, mut bw, mut bh) = (
|
||||
Vec::with_capacity(n),
|
||||
Vec::with_capacity(n),
|
||||
Vec::with_capacity(n),
|
||||
Vec::with_capacity(n),
|
||||
);
|
||||
let mut det_scores = Vec::with_capacity(n);
|
||||
let mut qualities: Vec<Option<f32>> = Vec::with_capacity(n);
|
||||
let mut embeddings = Vec::with_capacity(n);
|
||||
let mut blob_hashes: Vec<Option<&str>> = Vec::with_capacity(n);
|
||||
for f in faces {
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO faces.faces
|
||||
(id, file_id, user_id, person_id, bbox, det_score, quality, embedding, blob_hash)
|
||||
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
|
||||
"#,
|
||||
)
|
||||
.bind(f.id)
|
||||
.bind(f.file_id)
|
||||
.bind(f.user_id)
|
||||
.bind(f.person_id)
|
||||
.bind(f.bbox.to_array())
|
||||
.bind(f.det_score)
|
||||
.bind(f.quality)
|
||||
.bind(embedding_to_bytes(&f.embedding))
|
||||
.bind(f.blob_hash.as_deref())
|
||||
.execute(&mut *tx)
|
||||
.await
|
||||
.map_err(|e| db_err("save_faces", e))?;
|
||||
ids.push(f.id);
|
||||
file_ids.push(f.file_id);
|
||||
user_ids.push(f.user_id);
|
||||
person_ids.push(f.person_id);
|
||||
bx.push(f.bbox.x);
|
||||
by.push(f.bbox.y);
|
||||
bw.push(f.bbox.w);
|
||||
bh.push(f.bbox.h);
|
||||
det_scores.push(f.det_score);
|
||||
qualities.push(f.quality);
|
||||
embeddings.push(embedding_to_bytes(&f.embedding));
|
||||
blob_hashes.push(f.blob_hash.as_deref());
|
||||
}
|
||||
tx.commit().await.map_err(|e| db_err("commit", e))?;
|
||||
sqlx::query(
|
||||
r#"
|
||||
INSERT INTO faces.faces
|
||||
(id, file_id, user_id, person_id, bbox, det_score, quality, embedding, blob_hash)
|
||||
SELECT t.id, t.file_id, t.user_id, t.person_id,
|
||||
ARRAY[t.bx, t.by, t.bw, t.bh]::real[],
|
||||
t.det_score, t.quality, t.embedding, t.blob_hash
|
||||
FROM unnest($1::uuid[], $2::uuid[], $3::uuid[], $4::uuid[],
|
||||
$5::real[], $6::real[], $7::real[], $8::real[],
|
||||
$9::real[], $10::real[], $11::bytea[], $12::text[])
|
||||
AS t(id, file_id, user_id, person_id,
|
||||
bx, by, bw, bh, det_score, quality, embedding, blob_hash)
|
||||
"#,
|
||||
)
|
||||
.bind(&ids)
|
||||
.bind(&file_ids)
|
||||
.bind(&user_ids)
|
||||
.bind(&person_ids)
|
||||
.bind(&bx)
|
||||
.bind(&by)
|
||||
.bind(&bw)
|
||||
.bind(&bh)
|
||||
.bind(&det_scores)
|
||||
.bind(&qualities)
|
||||
.bind(&embeddings)
|
||||
.bind(&blob_hashes)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| db_err("save_faces", e))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
@@ -24,18 +24,21 @@ impl FavoritesPgRepository {
|
||||
|
||||
impl FavoritesRepositoryPort for FavoritesPgRepository {
|
||||
async fn get_favorites(&self, user_id: Uuid) -> Result<Vec<FavoriteItemDto>> {
|
||||
// `id`/`user_id`/`parent_id` decode as binary UUIDs (16 B on the wire,
|
||||
// no server-side `::TEXT` cast) and render app-side — the ROUND6 §10
|
||||
// pattern the two legacy listing methods here never picked up.
|
||||
let rows = sqlx::query(
|
||||
r#"
|
||||
SELECT
|
||||
uf.id::TEXT AS "id",
|
||||
uf.user_id::TEXT AS "user_id",
|
||||
uf.id AS "id",
|
||||
uf.user_id AS "user_id",
|
||||
uf.item_id AS "item_id",
|
||||
uf.item_type AS "item_type",
|
||||
uf.created_at AS "created_at",
|
||||
COALESCE(f.name, fld.name) AS "item_name",
|
||||
f.size AS "item_size",
|
||||
f.mime_type AS "item_mime_type",
|
||||
COALESCE(f.folder_id::TEXT, fld.parent_id::TEXT) AS "parent_id",
|
||||
COALESCE(f.folder_id, fld.parent_id) AS "parent_id",
|
||||
COALESCE(f.updated_at, fld.updated_at) AS "modified_at",
|
||||
CASE
|
||||
WHEN uf.item_type = 'folder' THEN fld.path
|
||||
@@ -70,15 +73,19 @@ impl FavoritesRepositoryPort for FavoritesPgRepository {
|
||||
.iter()
|
||||
.map(|row| {
|
||||
FavoriteItemDto {
|
||||
id: row.get("id"),
|
||||
user_id: row.get("user_id"),
|
||||
id: row.get::<i32, _>("id").to_string(),
|
||||
user_id: row.get::<Uuid, _>("user_id").to_string(),
|
||||
item_id: row.get("item_id"),
|
||||
item_type: row.get("item_type"),
|
||||
created_at: row.get("created_at"),
|
||||
item_name: row.try_get("item_name").ok(),
|
||||
item_size: row.try_get("item_size").ok(),
|
||||
item_mime_type: row.try_get("item_mime_type").ok(),
|
||||
parent_id: row.try_get("parent_id").ok(),
|
||||
parent_id: row
|
||||
.try_get::<Option<Uuid>, _>("parent_id")
|
||||
.ok()
|
||||
.flatten()
|
||||
.map(|u| u.to_string()),
|
||||
modified_at: row.try_get("modified_at").ok(),
|
||||
item_path: row.try_get("item_path").ok(),
|
||||
// Temporary defaults; with_display_fields() computes the real values
|
||||
|
||||
@@ -1394,19 +1394,6 @@ impl FileReadPort for FileBlobReadRepository {
|
||||
Ok((files, total_count))
|
||||
}
|
||||
|
||||
/// Count files matching the search criteria (without loading them).
|
||||
async fn count_files(
|
||||
&self,
|
||||
folder_id: Option<&str>,
|
||||
criteria: &SearchCriteriaDto,
|
||||
caller_id: Uuid,
|
||||
) -> Result<usize, DomainError> {
|
||||
let (_, count) = self
|
||||
.search_files_paginated(folder_id, criteria, caller_id)
|
||||
.await?;
|
||||
Ok(count)
|
||||
}
|
||||
|
||||
#[allow(clippy::type_complexity)]
|
||||
async fn suggest_files_by_name(
|
||||
&self,
|
||||
|
||||
@@ -544,17 +544,27 @@ impl PlaylistItemRepository for PlaylistItemPgRepository {
|
||||
playlist_id: &Uuid,
|
||||
item_ids: &[Uuid],
|
||||
) -> PlaylistItemRepositoryResult<()> {
|
||||
for (index, item_id) in item_ids.iter().enumerate() {
|
||||
sqlx::query(
|
||||
"UPDATE audio.playlist_items SET position = $2 WHERE id = $1 AND playlist_id = $3",
|
||||
)
|
||||
.bind(item_id)
|
||||
.bind(index as i32)
|
||||
.bind(playlist_id)
|
||||
.execute(&*self.pool)
|
||||
.await
|
||||
.map_err(|e| DomainError::database_error(format!("Failed to reorder: {}", e)))?;
|
||||
if item_ids.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
// One UNNEST-driven UPDATE instead of one autocommit round-trip per
|
||||
// track — a full drag-reorder of an N-track playlist was N statements
|
||||
// (and non-atomic: a mid-loop failure left a half-applied order).
|
||||
// `WITH ORDINALITY` numbers the ids in array order, 1-based, so
|
||||
// `ord - 1` reproduces the historical 0-based positions.
|
||||
sqlx::query(
|
||||
r#"
|
||||
UPDATE audio.playlist_items AS pi
|
||||
SET position = (t.ord - 1)::int
|
||||
FROM unnest($1::uuid[]) WITH ORDINALITY AS t(id, ord)
|
||||
WHERE pi.id = t.id AND pi.playlist_id = $2
|
||||
"#,
|
||||
)
|
||||
.bind(item_ids)
|
||||
.bind(playlist_id)
|
||||
.execute(&*self.pool)
|
||||
.await
|
||||
.map_err(|e| DomainError::database_error(format!("Failed to reorder: {}", e)))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
@@ -21,18 +21,20 @@ impl RecentItemsPgRepository {
|
||||
|
||||
impl RecentItemsRepositoryPort for RecentItemsPgRepository {
|
||||
async fn get_recent_items(&self, user_id: Uuid, limit: i32) -> Result<Vec<RecentItemDto>> {
|
||||
// Binary UUID decode + app-side render (ROUND6 §10 pattern) — no
|
||||
// server-side `::TEXT` casts, 16 B per id on the wire instead of 36.
|
||||
let rows = sqlx::query(
|
||||
r#"
|
||||
SELECT
|
||||
ur.id::TEXT AS "id",
|
||||
ur.user_id::TEXT AS "user_id",
|
||||
ur.id AS "id",
|
||||
ur.user_id AS "user_id",
|
||||
ur.item_id AS "item_id",
|
||||
ur.item_type AS "item_type",
|
||||
ur.accessed_at AS "accessed_at",
|
||||
COALESCE(f.name, fld.name) AS "item_name",
|
||||
f.size AS "item_size",
|
||||
f.mime_type AS "item_mime_type",
|
||||
COALESCE(f.folder_id::TEXT, fld.parent_id::TEXT) AS "parent_id",
|
||||
COALESCE(f.folder_id, fld.parent_id) AS "parent_id",
|
||||
CASE
|
||||
WHEN ur.item_type = 'folder' THEN fld.path
|
||||
WHEN ur.item_type = 'file' THEN COALESCE(pfld.path || '/' || f.name, f.name)
|
||||
@@ -67,15 +69,19 @@ impl RecentItemsRepositoryPort for RecentItemsPgRepository {
|
||||
.iter()
|
||||
.map(|row| {
|
||||
RecentItemDto {
|
||||
id: row.get("id"),
|
||||
user_id: row.get("user_id"),
|
||||
id: row.get::<i32, _>("id").to_string(),
|
||||
user_id: row.get::<Uuid, _>("user_id").to_string(),
|
||||
item_id: row.get("item_id"),
|
||||
item_type: row.get("item_type"),
|
||||
accessed_at: row.get("accessed_at"),
|
||||
item_name: row.try_get("item_name").ok(),
|
||||
item_size: row.try_get("item_size").ok(),
|
||||
item_mime_type: row.try_get("item_mime_type").ok(),
|
||||
parent_id: row.try_get("parent_id").ok(),
|
||||
parent_id: row
|
||||
.try_get::<Option<Uuid>, _>("parent_id")
|
||||
.ok()
|
||||
.flatten()
|
||||
.map(|u| u.to_string()),
|
||||
item_path: row.try_get("item_path").ok(),
|
||||
// Temporary defaults; with_display_fields() computes the real values
|
||||
icon_class: String::new(),
|
||||
|
||||
@@ -56,7 +56,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 +67,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 }
|
||||
}
|
||||
|
||||
|
||||
@@ -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,9 +40,9 @@ 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,
|
||||
}
|
||||
@@ -80,8 +85,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 +138,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 +184,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 +197,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 +227,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 +306,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]
|
||||
|
||||
@@ -31,12 +31,13 @@
|
||||
|
||||
use std::collections::HashSet;
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicU32, Ordering};
|
||||
use std::sync::atomic::{AtomicU32, AtomicU64, Ordering};
|
||||
use std::time::Duration;
|
||||
use uuid::Uuid;
|
||||
|
||||
use moka::future::Cache;
|
||||
use sqlx::PgPool;
|
||||
use tokio::sync::oneshot;
|
||||
|
||||
use crate::application::ports::authorization_ports::AuthorizationEngine;
|
||||
use crate::common::errors::DomainError;
|
||||
@@ -218,6 +219,61 @@ pub struct PgAclEngine {
|
||||
/// Grant writes don't affect parentage — only the TTL applies (moves are
|
||||
/// an indirect path, same self-heal contract as `cascade_grant_cache`).
|
||||
file_parent_cache: Cache<Uuid, Option<Uuid>>,
|
||||
/// Natural-batching collector for cold `file_parent_cache` misses — the
|
||||
/// ROUND9 §10 deferred item. A shared N-photo album's cold first view
|
||||
/// arrives as N near-simultaneous thumbnail requests, each missing the
|
||||
/// parent memo and each paying a point `SELECT folder_id`.
|
||||
///
|
||||
/// Leader-runs-inline shape: an idle miss marks itself leader (one
|
||||
/// mutex op) and runs its point query exactly as before — the
|
||||
/// SEQUENTIAL path gains no hop, no task, no extra latency (a
|
||||
/// channel-task variant benchmarked at ~66 µs/miss of pure overhead and
|
||||
/// was rejected). Misses arriving while a leader is in flight park a
|
||||
/// oneshot in this queue; the leader drains them into ONE `= ANY($1)`
|
||||
/// batch after its own query, so a K-wide herd collapses to ~2 queries.
|
||||
/// If a leader future is dropped mid-flight, its guard wakes every
|
||||
/// parked waiter to retry (and re-elect); waiters that exhaust retries
|
||||
/// fall back to the inline point query — strictly additive.
|
||||
parent_batch: Arc<std::sync::Mutex<Option<Vec<ParentWaiter>>>>,
|
||||
/// Total parent-resolution queries actually issued (point + batches) —
|
||||
/// exposed via [`Self::parent_query_count`] for benches/operators.
|
||||
parent_queries: Arc<AtomicU64>,
|
||||
}
|
||||
|
||||
/// One parked parent-resolution request: file id + reply slot. A dropped
|
||||
/// sender (leader cancelled) is the retry signal. Errors are shared behind
|
||||
/// `Arc` because `DomainError` carries a non-clonable source chain (same
|
||||
/// convention as the Basic-auth single-flight).
|
||||
type ParentWaiter = (
|
||||
Uuid,
|
||||
oneshot::Sender<Result<Option<Uuid>, Arc<DomainError>>>,
|
||||
);
|
||||
|
||||
/// Upper bound on ids drained into one `= ANY` parent batch. A browser herd
|
||||
/// is O(100); this only guards pathological queue growth.
|
||||
const PARENT_BATCH_MAX: usize = 256;
|
||||
|
||||
/// How many times a parked waiter re-runs the elect-or-park protocol after
|
||||
/// a leader vanished before giving up and querying inline itself.
|
||||
const PARENT_WAIT_RETRIES: usize = 3;
|
||||
|
||||
/// RAII release of parent-resolution leadership. If the leader future is
|
||||
/// dropped at an await point (client disconnect cancels the request), this
|
||||
/// clears the in-flight marker and drops every parked waiter's sender —
|
||||
/// their `oneshot` recv errors and they re-run the election, so a vanished
|
||||
/// leader can never strand the queue.
|
||||
struct ParentLeaderGuard<'a> {
|
||||
engine: &'a PgAclEngine,
|
||||
}
|
||||
|
||||
impl Drop for ParentLeaderGuard<'_> {
|
||||
fn drop(&mut self) {
|
||||
let mut slot = match self.engine.parent_batch.lock() {
|
||||
Ok(s) => s,
|
||||
Err(poisoned) => poisoned.into_inner(),
|
||||
};
|
||||
*slot = None;
|
||||
}
|
||||
}
|
||||
|
||||
impl PgAclEngine {
|
||||
@@ -263,6 +319,8 @@ impl PgAclEngine {
|
||||
.max_capacity(FILE_PARENT_CACHE_CAPACITY)
|
||||
.time_to_live(FILE_PARENT_CACHE_TTL)
|
||||
.build(),
|
||||
parent_batch: Arc::new(std::sync::Mutex::new(None)),
|
||||
parent_queries: Arc::new(AtomicU64::new(0)),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -341,6 +399,8 @@ impl PgAclEngine {
|
||||
.max_capacity(1)
|
||||
.time_to_live(Duration::from_secs(1))
|
||||
.build(),
|
||||
parent_batch: Arc::new(std::sync::Mutex::new(None)),
|
||||
parent_queries: Arc::new(AtomicU64::new(0)),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -779,11 +839,16 @@ impl PgAclEngine {
|
||||
Ok(exists.is_some())
|
||||
}
|
||||
|
||||
/// Memoised `file_id → Option<parent folder_id>` point read backing the
|
||||
/// Memoised `file_id → Option<parent folder_id>` read backing the
|
||||
/// file-cascade decomposition. `None` covers both a missing row and a
|
||||
/// NULL `folder_id` — in either case only the direct-file-grant branch
|
||||
/// can match (mirroring the historical UNION's `folder_id IS NOT NULL`
|
||||
/// guard).
|
||||
///
|
||||
/// Cold misses run the leader-inline batching protocol (see
|
||||
/// `parent_batch`): an idle miss queries inline exactly as before;
|
||||
/// misses concurrent with an in-flight leader park and are answered by
|
||||
/// the leader's single `= ANY` charity batch.
|
||||
async fn file_parent_folder_cached(
|
||||
&self,
|
||||
file_id: Uuid,
|
||||
@@ -793,15 +858,221 @@ impl PgAclEngine {
|
||||
return Ok(parent);
|
||||
}
|
||||
counters.sql_queries.fetch_add(1, Ordering::Relaxed);
|
||||
|
||||
enum Elect {
|
||||
Lead,
|
||||
Park(oneshot::Receiver<Result<Option<Uuid>, Arc<DomainError>>>),
|
||||
Overflow,
|
||||
}
|
||||
for _ in 0..=PARENT_WAIT_RETRIES {
|
||||
// Elect-or-park. The guard lives only inside this block — the
|
||||
// decision is acted on AFTER it drops, so no lock is ever held
|
||||
// across an await (and the handler futures stay `Send`).
|
||||
let outcome = {
|
||||
let mut slot = self.parent_batch.lock().expect("parent_batch poisoned");
|
||||
match slot.as_mut() {
|
||||
// A leader is in flight — park a oneshot in its queue.
|
||||
Some(queue) if queue.len() < PARENT_BATCH_MAX => {
|
||||
let (tx, rx) = oneshot::channel();
|
||||
queue.push((file_id, tx));
|
||||
Elect::Park(rx)
|
||||
}
|
||||
// Queue full — behave as if idle contention: inline below.
|
||||
Some(_) => Elect::Overflow,
|
||||
// Idle — become the leader.
|
||||
None => {
|
||||
*slot = Some(Vec::new());
|
||||
Elect::Lead
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
match outcome {
|
||||
Elect::Lead => return self.parent_leader_resolve(file_id).await,
|
||||
Elect::Park(rx) => match rx.await {
|
||||
Ok(Ok(parent)) => return Ok(parent),
|
||||
Ok(Err(shared)) => {
|
||||
return Err(DomainError::new(
|
||||
shared.kind,
|
||||
shared.entity_type,
|
||||
shared.message.clone(),
|
||||
));
|
||||
}
|
||||
// Leader vanished (cancelled mid-flight) — retry the
|
||||
// election; a fresh leader (possibly us) takes over.
|
||||
Err(_) => continue,
|
||||
},
|
||||
// Queue overflow: don't wait — resolve inline.
|
||||
Elect::Overflow => break,
|
||||
}
|
||||
}
|
||||
|
||||
// Retries exhausted or queue overflow: the historical inline read.
|
||||
self.parent_queries.fetch_add(1, Ordering::Relaxed);
|
||||
let parent = Self::query_parent_point(&self.pool, file_id).await?;
|
||||
self.file_parent_cache.insert(file_id, parent).await;
|
||||
Ok(parent)
|
||||
}
|
||||
|
||||
/// Leader half of the parent-resolution protocol: run own point query
|
||||
/// inline (the exact pre-round-10 cost), then serve everything that
|
||||
/// parked during it with ONE `= ANY` batch. A second wave arriving
|
||||
/// during the charity batch is handed to a detached drainer task so the
|
||||
/// leader's own response is never delayed by more than one batch.
|
||||
///
|
||||
/// Cancellation-safe: `ParentLeaderGuard` releases leadership on drop
|
||||
/// and wakes parked waiters (their `oneshot` senders drop → they retry
|
||||
/// and re-elect).
|
||||
async fn parent_leader_resolve(&self, file_id: Uuid) -> Result<Option<Uuid>, DomainError> {
|
||||
let guard = ParentLeaderGuard { engine: self };
|
||||
|
||||
self.parent_queries.fetch_add(1, Ordering::Relaxed);
|
||||
let own = Self::query_parent_point(&self.pool, file_id).await;
|
||||
if let Ok(parent) = &own {
|
||||
self.file_parent_cache.insert(file_id, *parent).await;
|
||||
}
|
||||
|
||||
// Take the first charity wave (leave `Some(vec![])` so later
|
||||
// arrivals keep parking while the batch runs).
|
||||
let wave = {
|
||||
let mut slot = self.parent_batch.lock().expect("parent_batch poisoned");
|
||||
match slot.as_mut() {
|
||||
Some(queue) if !queue.is_empty() => std::mem::take(queue),
|
||||
_ => Vec::new(),
|
||||
}
|
||||
};
|
||||
if !wave.is_empty() {
|
||||
self.parent_queries.fetch_add(1, Ordering::Relaxed);
|
||||
Self::serve_parent_wave(&self.pool, &self.file_parent_cache, wave).await;
|
||||
}
|
||||
|
||||
// Release leadership — or, if a second wave parked during the
|
||||
// charity batch, hand leadership to a detached drainer so the
|
||||
// leader's own response isn't delayed further. The drainer loops:
|
||||
// it keeps the slot marked in-flight (newer misses keep parking)
|
||||
// and only clears it when the queue drains empty.
|
||||
let second_wave = {
|
||||
let mut slot = self.parent_batch.lock().expect("parent_batch poisoned");
|
||||
match slot.as_mut() {
|
||||
Some(queue) if !queue.is_empty() => Some(std::mem::take(queue)),
|
||||
_ => {
|
||||
*slot = None; // idle again
|
||||
None
|
||||
}
|
||||
}
|
||||
};
|
||||
std::mem::forget(guard); // leadership released or handed to the drainer
|
||||
if let Some(first) = second_wave {
|
||||
let engine = self.clone_batch_handles();
|
||||
tokio::spawn(async move {
|
||||
let mut wave = first;
|
||||
loop {
|
||||
engine.2.fetch_add(1, Ordering::Relaxed);
|
||||
Self::serve_parent_wave(&engine.0, &engine.1, wave).await;
|
||||
let mut slot = match engine.3.lock() {
|
||||
Ok(s) => s,
|
||||
Err(p) => p.into_inner(),
|
||||
};
|
||||
match slot.as_mut() {
|
||||
Some(queue) if !queue.is_empty() => {
|
||||
wave = std::mem::take(queue);
|
||||
}
|
||||
_ => {
|
||||
*slot = None;
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
own
|
||||
}
|
||||
|
||||
/// The `Arc`'d handles the detached drainer needs (pool, memo cache,
|
||||
/// query counter, queue slot). Cloned individually because the drainer
|
||||
/// outlives this call and the engine isn't guaranteed to sit behind an
|
||||
/// `Arc` here.
|
||||
#[allow(clippy::type_complexity)]
|
||||
fn clone_batch_handles(
|
||||
&self,
|
||||
) -> (
|
||||
Arc<PgPool>,
|
||||
Cache<Uuid, Option<Uuid>>,
|
||||
Arc<AtomicU64>,
|
||||
Arc<std::sync::Mutex<Option<Vec<ParentWaiter>>>>,
|
||||
) {
|
||||
(
|
||||
Arc::clone(&self.pool),
|
||||
self.file_parent_cache.clone(),
|
||||
Arc::clone(&self.parent_queries),
|
||||
Arc::clone(&self.parent_batch),
|
||||
)
|
||||
}
|
||||
|
||||
/// Resolve one file's parent with the point query (shared by the
|
||||
/// leader's own read and the no-batching fallback).
|
||||
async fn query_parent_point(pool: &PgPool, file_id: Uuid) -> Result<Option<Uuid>, DomainError> {
|
||||
let parent: Option<Option<Uuid>> =
|
||||
sqlx::query_scalar("SELECT folder_id FROM storage.files WHERE id = $1")
|
||||
.bind(file_id)
|
||||
.fetch_optional(self.pool.as_ref())
|
||||
.fetch_optional(pool)
|
||||
.await
|
||||
.map_err(|e| DomainError::internal_error("PgAcl", format!("file parent: {e}")))?;
|
||||
let parent = parent.flatten();
|
||||
self.file_parent_cache.insert(file_id, parent).await;
|
||||
Ok(parent)
|
||||
Ok(parent.flatten())
|
||||
}
|
||||
|
||||
/// Serve a parked wave with one `= ANY` query: memoise every id
|
||||
/// (requested-but-absent rows memoise as `None`, matching the point
|
||||
/// read) and answer every oneshot. On error the shared failure is
|
||||
/// fanned out instead.
|
||||
async fn serve_parent_wave(
|
||||
pool: &PgPool,
|
||||
cache: &Cache<Uuid, Option<Uuid>>,
|
||||
wave: Vec<ParentWaiter>,
|
||||
) {
|
||||
let mut ids: Vec<Uuid> = Vec::with_capacity(wave.len());
|
||||
for (id, _) in &wave {
|
||||
if !ids.contains(id) {
|
||||
ids.push(*id);
|
||||
}
|
||||
}
|
||||
let fetched: Result<Vec<(Uuid, Option<Uuid>)>, sqlx::Error> =
|
||||
sqlx::query_as("SELECT id, folder_id FROM storage.files WHERE id = ANY($1)")
|
||||
.bind(&ids)
|
||||
.fetch_all(pool)
|
||||
.await;
|
||||
match fetched {
|
||||
Ok(rows) => {
|
||||
let mut by_id: std::collections::HashMap<Uuid, Option<Uuid>> =
|
||||
rows.into_iter().collect();
|
||||
for id in &ids {
|
||||
by_id.entry(*id).or_insert(None);
|
||||
}
|
||||
for (id, parent) in &by_id {
|
||||
cache.insert(*id, *parent).await;
|
||||
}
|
||||
for (id, reply) in wave {
|
||||
let parent = by_id.get(&id).copied().unwrap_or(None);
|
||||
let _ = reply.send(Ok(parent));
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
let shared = Arc::new(DomainError::internal_error(
|
||||
"PgAcl",
|
||||
format!("file parent batch: {e}"),
|
||||
));
|
||||
for (_, reply) in wave {
|
||||
let _ = reply.send(Err(Arc::clone(&shared)));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Total parent-resolution queries actually issued (point + `= ANY`
|
||||
/// batches). With batching, this is ≤ the number of cold misses — the
|
||||
/// gap is the herd-collapse win. Exposed for benches and operators.
|
||||
pub fn parent_query_count(&self) -> u64 {
|
||||
self.parent_queries.load(Ordering::Relaxed)
|
||||
}
|
||||
|
||||
/// Cache-aware wrapper over the File/Folder grant cascade. Serves the
|
||||
@@ -843,56 +1114,71 @@ impl PgAclEngine {
|
||||
counters.cache_hit.fetch_add(1, Ordering::Relaxed);
|
||||
return Ok(allowed);
|
||||
}
|
||||
let allowed = match resource {
|
||||
Resource::Folder(id) => {
|
||||
let (subject_types, subject_ids) =
|
||||
self.subject_match_set(subject, counters).await?;
|
||||
self.folder_cascade_grant_exists(
|
||||
&subject_types,
|
||||
&subject_ids,
|
||||
permission,
|
||||
id,
|
||||
counters,
|
||||
)
|
||||
.await?
|
||||
}
|
||||
Resource::File(id) => {
|
||||
// Ancestor half first — amortized to one query per FOLDER
|
||||
// via the recursive Folder arm (its own cache entry).
|
||||
let folder_allowed = match self.file_parent_folder_cached(id, counters).await? {
|
||||
Some(parent) => {
|
||||
Box::pin(self.cascade_grant_cached(
|
||||
subject,
|
||||
Resource::Folder(parent),
|
||||
permission,
|
||||
counters,
|
||||
))
|
||||
.await?
|
||||
}
|
||||
None => false,
|
||||
};
|
||||
if folder_allowed {
|
||||
true
|
||||
} else {
|
||||
let (subject_types, subject_ids) =
|
||||
self.subject_match_set(subject, counters).await?;
|
||||
self.file_direct_grant_exists(
|
||||
&subject_types,
|
||||
&subject_ids,
|
||||
permission,
|
||||
id,
|
||||
counters,
|
||||
)
|
||||
.await?
|
||||
}
|
||||
}
|
||||
// Only File/Folder reach this helper (see `check_inner`).
|
||||
_ => return Ok(false),
|
||||
};
|
||||
// Only File/Folder reach this helper (see `check_inner`); keep the
|
||||
// defensive arm OUTSIDE the loader so it stays uncached, as before.
|
||||
if !matches!(resource, Resource::Folder(_) | Resource::File(_)) {
|
||||
return Ok(false);
|
||||
}
|
||||
// `try_get_with`: a cold herd on the same key — every photo of an
|
||||
// album recursing into the SAME folder decision at once — coalesces
|
||||
// into ONE loader run. The old get→compute→insert let K concurrent
|
||||
// misses each run the ltree query (ROUND10; the ROUND3 auth-herd
|
||||
// pattern). moka never caches loader errors, preserving the
|
||||
// historical error semantics.
|
||||
self.cascade_grant_cache
|
||||
.insert((subject, resource, permission), allowed)
|
||||
.await;
|
||||
Ok(allowed)
|
||||
.try_get_with((subject, resource, permission), async {
|
||||
match resource {
|
||||
Resource::Folder(id) => {
|
||||
let (subject_types, subject_ids) =
|
||||
self.subject_match_set(subject, counters).await?;
|
||||
self.folder_cascade_grant_exists(
|
||||
&subject_types,
|
||||
&subject_ids,
|
||||
permission,
|
||||
id,
|
||||
counters,
|
||||
)
|
||||
.await
|
||||
}
|
||||
Resource::File(id) => {
|
||||
// Ancestor half first — amortized to one query per
|
||||
// FOLDER via the recursive Folder arm (its own cache
|
||||
// entry + its own single-flight).
|
||||
let folder_allowed =
|
||||
match self.file_parent_folder_cached(id, counters).await? {
|
||||
Some(parent) => {
|
||||
Box::pin(self.cascade_grant_cached(
|
||||
subject,
|
||||
Resource::Folder(parent),
|
||||
permission,
|
||||
counters,
|
||||
))
|
||||
.await?
|
||||
}
|
||||
None => false,
|
||||
};
|
||||
if folder_allowed {
|
||||
Ok(true)
|
||||
} else {
|
||||
let (subject_types, subject_ids) =
|
||||
self.subject_match_set(subject, counters).await?;
|
||||
self.file_direct_grant_exists(
|
||||
&subject_types,
|
||||
&subject_ids,
|
||||
permission,
|
||||
id,
|
||||
counters,
|
||||
)
|
||||
.await
|
||||
}
|
||||
}
|
||||
_ => unreachable!("guarded above"),
|
||||
}
|
||||
})
|
||||
.await
|
||||
.map_err(|e: Arc<DomainError>| {
|
||||
DomainError::new(e.kind, e.entity_type, e.message.clone())
|
||||
})
|
||||
}
|
||||
|
||||
/// Cached resolution of `(subject, drive_id) → Option<Role>` — the
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user