perf: round 9 — decorator PUT reactivation, session/search/dedup alloc purges, PROPFIND join!, folder-level cascade

Benchmark-gated round (benches/ROUND9.md): every change carries a
BEFORE/AFTER bench with equivalence/safety gates; verdicts below are from
the committed harnesses on 4 cores / local PG 16.

Backend:
- Blob decorators (Retry/Cached) now forward put_blob_from_bytes_unsynced
  + sync_blobs — the trait default had silently reinstated HEAD-before-PUT
  per chunk on decorated remote stacks, undoing ROUND3 §8. Full production
  stack: 500 probes -> 0, 1.9x wall at 10 ms RTT (bench_s3_put §3).
- NC PROPFIND per-page enrichment triple (favorites / oc:fileid / dead
  props) overlapped with tokio::join!: 2.07x local, 2.86x at 5 ms RTT
  (bench_nc_enrich_join, injected-latency decide-by-bench).
- Search enrichment consumes its DTOs and carries the interned Arc<str>
  display fields end-to-end (SearchFileResultDto type change, OpenAPI
  shape preserved): enrich_file 2.0x, 11.6 -> 2.2 allocs/row; the NC
  REPORT conversion stops re-running all three classifiers per row
  (bench_search_enrich).
- NC session Arc end-to-end: SharedNcSession extractor (8 -> 0 allocs),
  Arc<FolderDto> chroot cache (4 -> 0/hit), single shared Arc<CurrentUser>
  + lazy span render (11 -> 6/build) (bench_nc_session).
- Storage micro-pack: atomic create_new chunk writes (2.1x fresh),
  stream_chunks over the manifest Arc (4097 -> 0 allocs/read incl. the
  Range path), manifest single-flight (herd 64 -> 1 loads), hex_lower for
  chunk Content-MD5 (18 -> 1 allocs) (bench_storage_micro).
- OCS capabilities memoized into OnceLock<[Bytes;2]>: 237x, 102 -> 0
  allocs/poll, byte-identical (bench_capabilities_static).
- Drive::is_empty COUNT(*) sum -> EXISTS: 34.4x on a 100k-file drive
  (bench_drive_is_empty).
- favorites/recents row-map ROUND7 port: path/name/blob_hash moved,
  -2.75 allocs/row (bench_resource_row_map §2).
- Folder rows decode binary UUIDs (ROUND6 §10 port): 1.03-1.07x page
  fetch, honest verdict incl. one noise-band wash documented
  (bench_folder_uuid_decode).
- Authz: file cascade decision decomposed into memoized folder-level
  decision + direct-grant lookup (ROUND8 deferred item). Cold shared-album
  first view 592 -> 418 µs/thumb; warm path unchanged; safety gates incl.
  new direct-grant sibling isolation, revoke-flush re-verified, full
  integration authz suite green (bench_thumbnail_cascade_cache).

Frontend (vitest gates committed beside the code):
- resolveLabel/resolveRecipient O(directory) scan -> id-keyed Map: 13.9x
  (recipients.bench.test.ts).
- ResourceList selection-prune effect skips when nothing is selected
  (100 -> 0 Set builds per drain) and the photos timeline reads a
  listener-fed mobile flag instead of matchMedia per recompute
  (listDerives.bench.test.ts).

Verification: cargo fmt + clippy --all-features --all-targets -D warnings
clean; 524 unit + 554 integration (--cfg integration_tests) tests pass;
frontend npm run check clean with 293 vitest tests green.

Deferred with rationale in ROUND9.md: CalDAV authz-before-fetch reorder
(maintainer sign-off), per-page batched parent resolution, JWT-claims
Arc<str>, batch_operations signature widening.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01XDc9VtXvskJ6dnMRraSndn
This commit is contained in:
Claude
2026-07-18 16:12:04 +00:00
parent 2317d594e3
commit fdf445d2b0
40 changed files with 4279 additions and 346 deletions
@@ -187,22 +187,48 @@ impl BlobStorageBackend for CachedBlobBackend {
};
Box::pin(async move {
let size = inner.put_blob_from_bytes(&hash, data.clone()).await?;
// Also cache locally (best-effort): write bytes to cache path
let dest = self_ref.cached_path(&hash);
if let Some(parent) = dest.parent() {
let _ = fs::create_dir_all(parent).await;
}
let _ = fs::write(&dest, &data).await;
let data_len = data.len() as u64;
let mut idx = self_ref.index.lock().await;
if let Some(old) = idx.put(hash, CacheEntry { size: data_len }) {
self_ref.current_size.fetch_sub(old.size, Ordering::Relaxed);
}
self_ref.current_size.fetch_add(data_len, Ordering::Relaxed);
self_ref.cache_bytes_write_through(hash, &data).await;
Ok(size)
})
}
// Without this override the trait default would re-route the CDC chunk
// write through `put_blob_from_bytes` above, whose inner (synced) call
// pays the remote exists-probe per chunk. The local write-through cache
// population is kept identical — post-upload readers (thumbnail/EXIF/
// face hooks) hit the cache instead of re-fetching from the remote.
fn put_blob_from_bytes_unsynced(
&self,
hash: &str,
data: Bytes,
) -> Pin<Box<dyn std::future::Future<Output = Result<u64, DomainError>> + Send + '_>> {
let 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_unsynced(&hash, data.clone())
.await?;
self_ref.cache_bytes_write_through(hash, &data).await;
Ok(size)
})
}
// The durability barrier must reach the backend that buffered the
// unsynced writes; the local cache copy is disposable and needs none.
fn sync_blobs(
&self,
hashes: &[String],
) -> Pin<Box<dyn std::future::Future<Output = Result<(), DomainError>> + Send + '_>> {
self.inner.sync_blobs(hashes)
}
fn get_blob_stream(
&self,
hash: &str,
@@ -437,6 +463,24 @@ impl CachedRef {
self.cache_dir.join(prefix).join(format!("{hash}.blob"))
}
/// 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.
async fn cache_bytes_write_through(&self, hash: String, data: &Bytes) {
let dest = self.cached_path(&hash);
if let Some(parent) = dest.parent() {
let _ = fs::create_dir_all(parent).await;
}
let _ = fs::write(&dest, data).await;
let data_len = data.len() as u64;
let mut idx = self.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);
}
/// Single-flight wrapper around [`Self::fetch_and_cache_static`]: 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
@@ -834,8 +834,7 @@ impl ChunkedUploadService {
let data_clone = data.clone(); // Bytes::clone is O(1) — just an Arc increment
let actual_checksum = tokio::task::spawn_blocking(move || {
use md5::{Digest, Md5};
let hash = Md5::digest(&data_clone);
hash.iter().map(|b| format!("{b:02x}")).collect::<String>()
crate::common::fmt::hex_lower(&Md5::digest(&data_clone))
})
.await
.map_err(|e| format!("MD5 checksum task failed: {e}"))?;
+77 -40
View File
@@ -1747,18 +1747,24 @@ impl DedupService {
/// remote object stores where overlapping fetches hide per-chunk latency).
/// Shared by [`Self::read_blob_stream`] and [`Self::read_blob_bytes`] so both
/// build the chunk stream identically from a manifest's `chunk_hashes`.
/// Takes the shared manifest `Arc` and iterates its hashes by index —
/// the old `Vec<String>` signature forced every read to deep-clone the
/// whole hash list out of the cached manifest before the first byte
/// (N ~64-B String allocs per read of an N-chunk file); the per-chunk
/// `Arc` bump here is a single atomic increment.
fn stream_chunks(
&self,
chunk_hashes: Vec<String>,
manifest: Arc<ChunkManifest>,
) -> Pin<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>> {
let prefetch = self.backend.read_prefetch().max(1);
let backend = self.backend.clone();
let chunk_stream = stream::iter(chunk_hashes)
.map(move |chunk_hash| {
let chunk_stream = stream::iter(0..manifest.chunk_hashes.len())
.map(move |i| {
let backend = backend.clone();
let manifest = manifest.clone();
async move {
backend
.get_blob_stream(&chunk_hash)
.get_blob_stream(&manifest.chunk_hashes[i])
.await
.map_err(|e| std::io::Error::other(e.to_string()))
}
@@ -1771,31 +1777,56 @@ impl DedupService {
/// Cached manifest fetch for the read path (see the `manifest_cache`
/// field docs). `None` = legacy whole-file blob — never cached, so a
/// background rechunk that creates a manifest is honoured immediately.
///
/// Misses are single-flighted through `try_get_with`: K concurrent cold
/// readers of one newly-hot file (e.g. parallel Range probes on a big
/// video) coalesce onto ONE manifest SELECT instead of K. The
/// positive-only contract is preserved by routing "no manifest row" and
/// DB failures through the loader's error channel, which moka never
/// caches. The zero-alloc `get` fast path stays in front so warm reads
/// don't pay the owned-key clone `try_get_with` requires.
async fn manifest_cached(&self, hash: &str) -> Result<Option<Arc<ChunkManifest>>, DomainError> {
if let Some(m) = self.manifest_cache.get(hash).await {
return Ok(Some(m));
}
let row = sqlx::query_as::<_, (Vec<String>, Vec<i64>, i64)>(
"SELECT chunk_hashes, chunk_sizes, total_size
FROM storage.chunk_manifests WHERE file_hash = $1",
)
.bind(hash)
.fetch_optional(self.pool.as_ref())
.await
.map_err(|e| DomainError::internal_error("Dedup", format!("Manifest lookup: {}", e)))?;
match row {
Some((chunk_hashes, chunk_sizes, total_size)) => {
let m = Arc::new(ChunkManifest {
chunk_hashes,
chunk_sizes,
total_size,
});
self.manifest_cache
.insert(hash.to_string(), m.clone())
.await;
Ok(Some(m))
}
None => Ok(None),
enum MissKind {
Legacy,
Db(String),
}
let pool = self.pool.clone();
let query_hash = hash.to_string();
let result = self
.manifest_cache
.try_get_with(hash.to_string(), async move {
let row = sqlx::query_as::<_, (Vec<String>, Vec<i64>, i64)>(
"SELECT chunk_hashes, chunk_sizes, total_size
FROM storage.chunk_manifests WHERE file_hash = $1",
)
.bind(&query_hash)
.fetch_optional(pool.as_ref())
.await
.map_err(|e| MissKind::Db(e.to_string()))?;
match row {
Some((chunk_hashes, chunk_sizes, total_size)) => Ok(Arc::new(ChunkManifest {
chunk_hashes,
chunk_sizes,
total_size,
})),
None => Err(MissKind::Legacy),
}
})
.await;
match result {
Ok(m) => Ok(Some(m)),
Err(miss) => match &*miss {
MissKind::Legacy => Ok(None),
MissKind::Db(msg) => Err(DomainError::internal_error(
"Dedup",
format!("Manifest lookup: {}", msg),
)),
},
}
}
@@ -1810,7 +1841,7 @@ impl DedupService {
) -> Result<Pin<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>>, DomainError>
{
match self.manifest_cached(hash).await? {
Some(m) => Ok(self.stream_chunks(m.chunk_hashes.clone())),
Some(m) => Ok(self.stream_chunks(m)),
// Legacy whole-file blob
None => self.backend.get_blob_stream(hash).await,
}
@@ -1829,10 +1860,10 @@ impl DedupService {
/// every full-blob read (e.g. 2N queries for an N-image gallery cold load).
pub async fn read_blob_bytes(&self, hash: &str) -> Result<Bytes, DomainError> {
let (mut stream, expected_size) = match self.manifest_cached(hash).await? {
Some(m) => (
self.stream_chunks(m.chunk_hashes.clone()),
m.total_size.max(0) as usize,
),
Some(m) => {
let expected = m.total_size.max(0) as usize;
(self.stream_chunks(m), expected)
}
None => {
// Legacy whole-file blob: size + stream straight from the backend.
let size = self.backend.blob_size(hash).await? as usize;
@@ -1863,16 +1894,17 @@ impl DedupService {
) -> Result<Pin<Box<dyn Stream<Item = Result<Bytes, std::io::Error>> + Send>>, DomainError>
{
if let Some(m) = self.manifest_cached(hash).await? {
let (chunk_hashes, chunk_sizes, total_size) =
(&m.chunk_hashes, &m.chunk_sizes, m.total_size);
let end = end.unwrap_or(total_size as u64);
let end = end.unwrap_or(m.total_size as u64);
// Calculate which chunks overlap [start, end)
// Calculate which chunks overlap [start, end). Chunks are
// addressed by manifest INDEX (the hash is read through the
// shared `Arc` at fetch time) — a `bytes=0-` probe of an
// N-chunk video used to clone all N hash Strings here.
let mut offset: u64 = 0;
// (chunk_hash, range_start_within_chunk, range_end_within_chunk)
let mut selected: Vec<(String, u64, Option<u64>)> = Vec::new();
// (chunk_index, range_start_within_chunk, range_end_within_chunk)
let mut selected: Vec<(usize, u64, Option<u64>)> = Vec::new();
for (i, &chunk_size) in chunk_sizes.iter().enumerate() {
for (i, &chunk_size) in m.chunk_sizes.iter().enumerate() {
let chunk_size = chunk_size as u64;
let chunk_end = offset + chunk_size;
@@ -1883,7 +1915,7 @@ impl DedupService {
} else {
None
};
selected.push((chunk_hashes[i].clone(), range_start, range_end));
selected.push((i, range_start, range_end));
}
offset += chunk_size;
@@ -1897,11 +1929,16 @@ impl DedupService {
let prefetch = self.backend.read_prefetch().max(1);
let backend = self.backend.clone();
let chunk_stream = stream::iter(selected)
.map(move |(chunk_hash, range_start, range_end)| {
.map(move |(i, range_start, range_end)| {
let backend = backend.clone();
let manifest = m.clone();
async move {
backend
.get_blob_range_stream(&chunk_hash, range_start, range_end)
.get_blob_range_stream(
&manifest.chunk_hashes[i],
range_start,
range_end,
)
.await
.map_err(|e| std::io::Error::other(e.to_string()))
}
@@ -128,18 +128,41 @@ async fn fsync_paths_parallel(paths: Vec<PathBuf>, strict: bool) -> Result<(), D
/// (fsync now vs. deferred batch sync), or `None` when the blob already
/// existed (idempotent skip — content-addressed, so identical by definition).
async fn write_blob_bytes(blob_path: &Path, data: &Bytes) -> Result<Option<File>, DomainError> {
if fs::try_exists(blob_path).await.unwrap_or(false) {
return Ok(None);
}
let mut file = fs::File::create(blob_path).await.map_err(|e| {
DomainError::internal_error("Blob", format!("Failed to create blob file: {}", e))
})?;
// One atomic O_CREAT|O_EXCL open replaces the old stat-then-create pair:
// `AlreadyExists` IS the idempotent skip (content-addressed names mean an
// existing file has identical content), saving a syscall + a blocking-pool
// dispatch on every new chunk of every upload.
let mut file = match fs::File::options()
.write(true)
.create_new(true)
.open(blob_path)
.await
{
Ok(f) => f,
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => return Ok(None),
Err(e) => {
return Err(DomainError::internal_error(
"Blob",
format!("Failed to create blob file: {}", e),
));
}
};
file.write_all(data).await.map_err(|e| {
DomainError::internal_error("Blob", format!("Failed to write blob from bytes: {}", e))
})?;
Ok(Some(file))
}
/// Bench-only public wrapper (feature = "bench") over the private chunk
/// writer so `examples/bench_storage_micro.rs` can A/B the open strategy.
#[cfg(feature = "bench")]
pub async fn write_blob_bytes_for_bench(
blob_path: &Path,
data: &Bytes,
) -> Result<Option<File>, DomainError> {
write_blob_bytes(blob_path, data).await
}
/// Compile-time lookup table for the 256 two-digit lowercase hex prefixes ("00"…"ff").
static HEX_PREFIXES: [&str; 256] = [
"00", "01", "02", "03", "04", "05", "06", "07", "08", "09", "0a", "0b", "0c", "0d", "0e", "0f",
+106 -41
View File
@@ -116,6 +116,15 @@ const CASCADE_GRANT_CACHE_CAPACITY: u64 = 100_000;
/// invalidation tree". Short enough that any such change takes effect in <1 min.
const CASCADE_GRANT_CACHE_TTL: Duration = Duration::from_secs(30);
/// `file_parent_cache` bound/TTL: `file_id → Option<folder_id>` point rows
/// (~50 B each) resolved on the file-cascade path so an N-file album pays
/// ONE folder-cascade query instead of N (ROUND9). Parentage changes only
/// on move — an indirect path the cascade cache already self-heals via TTL,
/// so the same 30 s window applies (grant writes don't alter parentage and
/// need no flush here).
const FILE_PARENT_CACHE_CAPACITY: u64 = 100_000;
const FILE_PARENT_CACHE_TTL: Duration = Duration::from_secs(30);
pub struct PgAclEngine {
pool: Arc<PgPool>,
folder_repo: Arc<FolderDbRepository>,
@@ -202,6 +211,13 @@ pub struct PgAclEngine {
/// only positively-or-negatively for at most the TTL. A revoke via
/// `clear_role` flushes immediately; anything missed self-heals in ≤30 s.
cascade_grant_cache: Cache<(Subject, Resource, Permission), bool>,
/// `file_id → Option<parent folder_id>` memo for the file-cascade
/// decomposition (see `cascade_grant_cached`): resolving the parent lets
/// a whole folder's files share ONE folder-cascade decision, so a shared
/// album's first view runs one ltree query instead of one per file.
/// 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>>,
}
impl PgAclEngine {
@@ -243,6 +259,10 @@ impl PgAclEngine {
.max_capacity(CASCADE_GRANT_CACHE_CAPACITY)
.time_to_live(CASCADE_GRANT_CACHE_TTL)
.build(),
file_parent_cache: Cache::builder()
.max_capacity(FILE_PARENT_CACHE_CAPACITY)
.time_to_live(FILE_PARENT_CACHE_TTL)
.build(),
}
}
@@ -317,6 +337,10 @@ impl PgAclEngine {
.max_capacity(1)
.time_to_live(Duration::from_secs(1))
.build(),
file_parent_cache: Cache::builder()
.max_capacity(1)
.time_to_live(Duration::from_secs(1))
.build(),
}
}
@@ -718,11 +742,11 @@ impl PgAclEngine {
Ok(exists.is_some())
}
/// Cascading check for files: either a direct file grant OR a grant on
/// any ancestor folder of the file's containing folder. See
/// `folder_cascade_grant_exists` for the meaning of `subject_types` /
/// `subject_ids` and the D-Prep role-array migration.
async fn file_cascade_grant_exists(
/// Direct file grant only — the first branch of the historical file
/// cascade UNION, split out so `cascade_grant_cached` can amortize the
/// ancestor-folder branch per FOLDER (see the `Resource::File` arm).
/// A plain indexed `role_grants` point lookup, no ltree join.
async fn file_direct_grant_exists(
&self,
subject_types: &[&str],
subject_ids: &[Uuid],
@@ -735,30 +759,12 @@ impl PgAclEngine {
let exists: Option<i32> = sqlx::query_scalar(
r#"
SELECT 1
FROM (
-- direct file grant
SELECT 1
FROM storage.role_grants
WHERE subject_type = ANY($1)
AND subject_id = ANY($2)
AND role = ANY($3::storage.grant_role[])
AND resource_type = 'file' AND resource_id = $4
AND (expires_at IS NULL OR expires_at > NOW())
UNION ALL
-- cascading from any ancestor folder of the file's containing folder
SELECT 1
FROM storage.role_grants g
JOIN storage.folders gf ON gf.id = g.resource_id
JOIN storage.files target_f ON target_f.id = $4
WHERE g.subject_type = ANY($1)
AND g.subject_id = ANY($2)
AND g.role = ANY($3::storage.grant_role[])
AND g.resource_type = 'folder'
AND (g.expires_at IS NULL OR g.expires_at > NOW())
AND target_f.folder_id IS NOT NULL
AND gf.lpath @> (SELECT lpath FROM storage.folders
WHERE id = target_f.folder_id)
) any_match
FROM storage.role_grants
WHERE subject_type = ANY($1)
AND subject_id = ANY($2)
AND role = ANY($3::storage.grant_role[])
AND resource_type = 'file' AND resource_id = $4
AND (expires_at IS NULL OR expires_at > NOW())
LIMIT 1
"#,
)
@@ -768,11 +774,36 @@ impl PgAclEngine {
.bind(file_id)
.fetch_optional(self.pool.as_ref())
.await
.map_err(|e| DomainError::internal_error("PgAcl", format!("file cascade: {e}")))?;
.map_err(|e| DomainError::internal_error("PgAcl", format!("file direct grant: {e}")))?;
Ok(exists.is_some())
}
/// Memoised `file_id → Option<parent folder_id>` point 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).
async fn file_parent_folder_cached(
&self,
file_id: Uuid,
counters: &QueryCounters,
) -> Result<Option<Uuid>, DomainError> {
if let Some(parent) = self.file_parent_cache.get(&file_id).await {
return Ok(parent);
}
counters.sql_queries.fetch_add(1, Ordering::Relaxed);
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())
.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)
}
/// Cache-aware wrapper over the File/Folder grant cascade. Serves the
/// memoised `(subject, resource, permission)` decision when warm; on a
/// miss it expands the subject set (itself cached) and runs the matching
@@ -780,10 +811,23 @@ impl PgAclEngine {
/// precheck fails, so it never caches a decision a drive grant would have
/// satisfied — a later drive grant short-circuits above this cache.
///
/// **File decomposition (ROUND9).** The historical file query was one
/// UNION: `direct file grant ∨ grant on any ancestor of the parent
/// folder` — one ltree join per file, so a shared N-photo album's FIRST
/// view ran N near-identical ancestor queries (round 8 memoised only the
/// per-file result, covering revalidation). The arm now resolves the
/// file's parent (memoised point read) and recurses into the FOLDER arm
/// for the ancestor half — one ltree query per folder, shared by every
/// sibling — falling back to the direct-file-grant lookup only when the
/// folder half denies. The decomposition is exactly the UNION split in
/// two: no decision changes, including the parentless edge (the UNION's
/// `folder_id IS NOT NULL` guard ≡ the direct-only fallback).
///
/// The result is a pure function of the subject's group expansion + the
/// resource's grants + folder ancestry; `invalidate_cascade_grant_cache_all`
/// (on File/Folder grant writes) and the 30 s TTL (indirect changes) keep
/// it fresh. See the `cascade_grant_cache` field doc.
/// (on File/Folder grant writes — it holds file AND folder decisions in
/// the same map) and the 30 s TTL (indirect changes, incl. moves for the
/// parent memo) keep it fresh. See the `cascade_grant_cache` field doc.
async fn cascade_grant_cached(
&self,
subject: Subject,
@@ -799,9 +843,10 @@ impl PgAclEngine {
counters.cache_hit.fetch_add(1, Ordering::Relaxed);
return Ok(allowed);
}
let (subject_types, subject_ids) = self.subject_match_set(subject, counters).await?;
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,
@@ -812,14 +857,34 @@ impl PgAclEngine {
.await?
}
Resource::File(id) => {
self.file_cascade_grant_exists(
&subject_types,
&subject_ids,
permission,
id,
counters,
)
.await?
// 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),
@@ -159,6 +159,43 @@ impl BlobStorageBackend for RetryBlobBackend {
})
}
// Without this override the trait default would re-route the CDC chunk
// write through `put_blob_from_bytes` above — reinstating the remote
// backend's exists-probe (HEAD/get_properties) per chunk that the
// `_unsynced` fast path exists to skip.
fn put_blob_from_bytes_unsynced(
&self,
hash: &str,
data: Bytes,
) -> Pin<Box<dyn std::future::Future<Output = Result<u64, DomainError>> + Send + '_>> {
let inner = self.inner.clone();
let policy = self.policy.clone();
let hash = hash.to_string();
Box::pin(async move {
retry_async(
&policy,
&format!("put_blob_from_bytes_unsynced({hash})"),
|| {
let inner = inner.clone();
let hash = hash.clone();
let data = data.clone();
async move { inner.put_blob_from_bytes_unsynced(&hash, data).await }
},
)
.await
})
}
// Forwarded WITHOUT retry wrapping: a failed fsync must surface, not be
// re-issued — after an fsync error the kernel may have dropped the dirty
// pages, so a retried fsync can report success for data that was lost.
fn sync_blobs(
&self,
hashes: &[String],
) -> Pin<Box<dyn std::future::Future<Output = Result<(), DomainError>> + Send + '_>> {
self.inner.sync_blobs(hashes)
}
fn get_blob_stream(
&self,
hash: &str,