perf: round 4 — one-pass row paths, drive-selector cache, CalDAV single-parse, streamed Azure, batched hydration
Nine benchmark-gated changes (benches/ROUND4.md; every one ships with a BEFORE/AFTER bench + equivalence gate, rollback rule as ROUND2/3): - Row→entity path build: one-pass StoragePath::from_folder_and_name / from_joined + normalize_storage_name_owned + alloc-free Display — 743→417 ns/file-row (1.78x), −5 allocs/row on every listing surface. - WebDAV drive-selector: per-user readable_cache (single-flight, 30 s TTL, explicit invalidation incl. membership + group changes) replaces the grants join per request — 441 µs → 0.8 µs (~550x), 0 queries warm. - CalDAV from_ical/update_ical_data: 8 full IcalParser runs per VEVENT → 1 (7.1x per PUT, 4.4x on 50-event imports); alloc-free split_vevents, chunk scan without the whole-body uppercase copy (1.4x), borrowed-key UID grouping (1.3x), REPORT props no longer cloned. - PROPFIND emit: partition Vecs dropped (single-pass 404 list) + stack rendered RFC 3339/2822 dates, sizes, quoted etags (common::fmt, chrono-byte-identical, sweep-tested) on both DAV surfaces — 1.22x per page, 17.9→12.0 allocs/row. - Grant-listing hydration: calendars/address books/playlists batch hydrate via = ANY($1) — 15 serial queries → 1 (~13x per sync poll). - user-flags cache: get→insert → try_get_with single-flight (32→1 queries per cold herd). - Azure downloads: whole-blob Vec buffering → streamed SDK pages — TTFB 349→4 ms (87x), peak heap 480→1.9 MiB (254x) on 256 MiB blobs; new OXICLOUD_AZURE_ENDPOINT_URL override (Azurite/bench hook). - Face indexing: unbounded per-image tokio::spawn → core-count semaphore, permit before blob read — peak heap 1175→176 MiB (6.7x). Checks: cargo fmt, clippy --all-features --all-targets -D warnings, cargo test --workspace (523 passed) + --features test_utils. hurl API suite and dockerized integration DB not runnable in this environment — left to CI. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_017aJu9ghvuT8WqC31ZEGTBA
This commit is contained in:
@@ -9,7 +9,7 @@ use std::pin::Pin;
|
||||
use azure_storage::StorageCredentials;
|
||||
use azure_storage_blobs::prelude::*;
|
||||
use bytes::Bytes;
|
||||
use futures::StreamExt;
|
||||
use futures::{StreamExt, TryStreamExt};
|
||||
use tokio::fs;
|
||||
|
||||
use crate::application::ports::blob_storage_ports::{
|
||||
@@ -33,8 +33,21 @@ impl AzureBlobBackend {
|
||||
StorageCredentials::access_key(&config.account_name, config.account_key.clone())
|
||||
};
|
||||
|
||||
let container_client = ClientBuilder::new(&config.account_name, credentials)
|
||||
.container_client(&config.container);
|
||||
// Custom endpoint (Azurite emulator / private deployment /
|
||||
// benches) mirrors S3's `endpoint_url`; default is the public
|
||||
// cloud URL derived from the account name.
|
||||
let container_client = match &config.endpoint_url {
|
||||
Some(uri) => ClientBuilder::with_location(
|
||||
azure_storage::CloudLocation::Custom {
|
||||
account: config.account_name.clone(),
|
||||
uri: uri.trim_end_matches('/').to_string(),
|
||||
},
|
||||
credentials,
|
||||
)
|
||||
.container_client(&config.container),
|
||||
None => ClientBuilder::new(&config.account_name, credentials)
|
||||
.container_client(&config.container),
|
||||
};
|
||||
|
||||
Self {
|
||||
container_client,
|
||||
@@ -169,29 +182,46 @@ impl BlobStorageBackend for AzureBlobBackend {
|
||||
Box::pin(async move {
|
||||
let client = self.blob_client(&hash);
|
||||
|
||||
let mut result_data: Vec<u8> = Vec::new();
|
||||
let mut stream = client.get().into_stream();
|
||||
|
||||
while let Some(response) = stream.next().await {
|
||||
let response = response.map_err(|e| {
|
||||
DomainError::new(
|
||||
// The old implementation drained the ENTIRE blob into one
|
||||
// `Vec<u8>` before yielding a single mega-chunk — whole-blob
|
||||
// RAM residency per reader, and with `read_prefetch() = 8`
|
||||
// up to 8 entire chunk-blobs resident at once during CDC
|
||||
// reassembly. Now the SDK's page/body streams forward
|
||||
// directly. The FIRST page is still awaited eagerly so a
|
||||
// missing blob surfaces as the same up-front NotFound the
|
||||
// old code produced; later pages/chunks map to io::Error
|
||||
// items like every other backend's stream.
|
||||
let mut pages = client.get().into_stream();
|
||||
let first = match pages.next().await {
|
||||
Some(Ok(response)) => response,
|
||||
Some(Err(e)) => {
|
||||
return Err(DomainError::new(
|
||||
ErrorKind::NotFound,
|
||||
"Azure",
|
||||
format!("Failed to get blob {hash}: {e}"),
|
||||
)
|
||||
})?;
|
||||
let mut body = response.data;
|
||||
while let Some(chunk) = body.next().await {
|
||||
let chunk = chunk.map_err(|e| {
|
||||
DomainError::internal_error("Azure", format!("Stream read error: {e}"))
|
||||
})?;
|
||||
result_data.extend_from_slice(&chunk);
|
||||
));
|
||||
}
|
||||
}
|
||||
None => {
|
||||
let empty: BlobStream =
|
||||
Box::pin(futures::stream::once(async move { Ok(Bytes::new()) }));
|
||||
return Ok(empty);
|
||||
}
|
||||
};
|
||||
|
||||
let stream: BlobStream = Box::pin(futures::stream::once(async move {
|
||||
Ok(Bytes::from(result_data))
|
||||
}));
|
||||
let first_body = first.data.map(|chunk| {
|
||||
chunk.map_err(|e| std::io::Error::other(format!("Stream read error: {e}")))
|
||||
});
|
||||
let tail = pages
|
||||
.map(|page| match page {
|
||||
Ok(response) => Ok(response.data.map(|chunk| {
|
||||
chunk.map_err(|e| std::io::Error::other(format!("Stream read error: {e}")))
|
||||
})),
|
||||
Err(e) => Err(std::io::Error::other(format!(
|
||||
"Failed to get blob page: {e}"
|
||||
))),
|
||||
})
|
||||
.try_flatten();
|
||||
let stream: BlobStream = Box::pin(first_body.chain(tail));
|
||||
Ok(stream)
|
||||
})
|
||||
}
|
||||
@@ -212,32 +242,42 @@ impl BlobStorageBackend for AzureBlobBackend {
|
||||
None => azure_core::request_options::Range::new(start, u64::MAX),
|
||||
};
|
||||
|
||||
let mut result_data: Vec<u8> = Vec::new();
|
||||
let mut stream = client.get().range(range).into_stream();
|
||||
|
||||
while let Some(response) = stream.next().await {
|
||||
let response = response.map_err(|e| {
|
||||
DomainError::new(
|
||||
// Same forwarding shape as `get_blob_stream` — a ranged read
|
||||
// doubly so: the caller explicitly asked NOT to pay for the
|
||||
// whole blob, yet the old code buffered the full range.
|
||||
let mut pages = client.get().range(range).into_stream();
|
||||
let first = match pages.next().await {
|
||||
Some(Ok(response)) => response,
|
||||
Some(Err(e)) => {
|
||||
return Err(DomainError::new(
|
||||
ErrorKind::NotFound,
|
||||
"Azure",
|
||||
format!("Failed to get blob range {hash}: {e}"),
|
||||
)
|
||||
})?;
|
||||
let mut body = response.data;
|
||||
while let Some(chunk) = body.next().await {
|
||||
let chunk = chunk.map_err(|e| {
|
||||
DomainError::internal_error(
|
||||
"Azure",
|
||||
format!("Stream range read error: {e}"),
|
||||
)
|
||||
})?;
|
||||
result_data.extend_from_slice(&chunk);
|
||||
));
|
||||
}
|
||||
}
|
||||
None => {
|
||||
let empty: BlobStream =
|
||||
Box::pin(futures::stream::once(async move { Ok(Bytes::new()) }));
|
||||
return Ok(empty);
|
||||
}
|
||||
};
|
||||
|
||||
let stream: BlobStream = Box::pin(futures::stream::once(async move {
|
||||
Ok(Bytes::from(result_data))
|
||||
}));
|
||||
let first_body = first.data.map(|chunk| {
|
||||
chunk.map_err(|e| std::io::Error::other(format!("Stream range read error: {e}")))
|
||||
});
|
||||
let tail = pages
|
||||
.map(|page| match page {
|
||||
Ok(response) => Ok(response.data.map(|chunk| {
|
||||
chunk.map_err(|e| {
|
||||
std::io::Error::other(format!("Stream range read error: {e}"))
|
||||
})
|
||||
})),
|
||||
Err(e) => Err(std::io::Error::other(format!(
|
||||
"Failed to get blob range page: {e}"
|
||||
))),
|
||||
})
|
||||
.try_flatten();
|
||||
let stream: BlobStream = Box::pin(first_body.chain(tail));
|
||||
Ok(stream)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -28,11 +28,35 @@ fn is_image(content_type: &str) -> bool {
|
||||
content_type.starts_with("image/")
|
||||
}
|
||||
|
||||
/// Concurrent index-task budget. Env override
|
||||
/// `OXICLOUD_FACES_INDEX_CONCURRENCY`, else the effective core count —
|
||||
/// each task is a full-image read + decode + ONNX inference, so more
|
||||
/// permits than cores only adds RAM pressure, not throughput.
|
||||
fn max_concurrent_index() -> usize {
|
||||
std::env::var("OXICLOUD_FACES_INDEX_CONCURRENCY")
|
||||
.ok()
|
||||
.and_then(|v| v.parse().ok())
|
||||
.filter(|&n: &usize| n > 0)
|
||||
.unwrap_or_else(|| {
|
||||
std::thread::available_parallelism()
|
||||
.map(|n| n.get())
|
||||
.unwrap_or(2)
|
||||
})
|
||||
}
|
||||
|
||||
pub struct FaceIndexingService {
|
||||
pool: Arc<PgPool>,
|
||||
repo: Arc<FacePgRepository>,
|
||||
analyzer: Arc<dyn FaceAnalyzerPort>,
|
||||
blob_root: PathBuf,
|
||||
/// Bounds concurrent indexing tasks. The lifecycle hooks spawn one
|
||||
/// task per uploaded/copied image with no ceiling, so a bulk upload
|
||||
/// used to fan out N simultaneous full-image reads + decodes +
|
||||
/// inferences — peak RSS N × image size plus CPU thrash. Same
|
||||
/// invariant as `ThumbnailService::decode_semaphore`: the permit is
|
||||
/// acquired BEFORE the blob read, so peak memory is
|
||||
/// `permits × image size` regardless of upload concurrency.
|
||||
index_semaphore: Arc<tokio::sync::Semaphore>,
|
||||
}
|
||||
|
||||
impl FaceIndexingService {
|
||||
@@ -43,6 +67,7 @@ impl FaceIndexingService {
|
||||
repo,
|
||||
analyzer,
|
||||
blob_root,
|
||||
index_semaphore: Arc::new(tokio::sync::Semaphore::new(max_concurrent_index())),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -60,7 +85,15 @@ impl FaceIndexingService {
|
||||
let repo = self.repo.clone();
|
||||
let analyzer = self.analyzer.clone();
|
||||
let blob_path = self.blob_path(&blob_hash);
|
||||
let semaphore = self.index_semaphore.clone();
|
||||
tokio::spawn(async move {
|
||||
// Queue behind the concurrency budget BEFORE touching the
|
||||
// blob — excess tasks wait holding only this tiny future,
|
||||
// not a decoded image.
|
||||
let _permit = semaphore
|
||||
.acquire_owned()
|
||||
.await
|
||||
.expect("face index semaphore never closes");
|
||||
if delete_first {
|
||||
let _ = repo.delete_faces_for_file(file_id).await;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user