221c1f31b0
Backend (each change benchmark-gated with BEFORE replicas + equivalence gates; see examples/bench_round11_micro.rs, bench_round11_queries.rs, bench_log_writer.rs and benches/ROUND11.md — final numbers land in the follow-up doc commit): - StoragePath re-representation: single canonical joined String, segments derived on demand; File/Folder drop the duplicated path_string field (4000→1000 allocs per 500-row listing page) - Display classifier fusion: classify_display shares one stack-lowered extension across the three decision trees; call sites in FileDto, folder/favorites/recent handlers, trash, path-resolver (+ interning where Arc::from was still used) - /status.php and /openapi.json memoized into OnceLock<Bytes> (openapi rebuilt a 171 KiB spec per request: 2.8 ms → 18 ns) - NC upload-session PROPFIND: write! + pre-sized body + stack RFC2822 dates (2.3-2.6x, 2582→772 allocs at 256 chunks) - REST download: dead FileDto clone removed (capture mime/size + move) - CalendarEventDto/TrashedItem into_parts moves (11 KiB ical_data memcpy gone per CalDAV row); CardDAV getlastmodified stack render - 4xx path: borrowed ErrorResponse serialize, ErrorKind::as_str, not_found/already_exists clone kill - vCard emit via write!; search page moved out with into_iter skip/take; content-hit UUIDs parsed once; group last-user check via HashSet - RateLimiter: lock-free get + insert (and_upsert_with variant REJECTED by benchmark); CSRF token borrow-compare + borrowed cookie extraction - Thumbnail/preview ETags built from as_str (Debug-identical bytes) - Encrypted backend: encrypt_in_place_detached single-buffer write path, chunk-sized reserve in collect_stream; retry labels made lazy - PG: deferred upload registration 3→1 round-trips (persist_file CTE template); direct_grant_cache for Calendar/AddressBook/Playlist authz (single-flight + set_role/clear_role invalidation); expand_user tokio::join!; geo clusters min(uuid)::text; recluster face assignment batched into one UNNEST update - People recluster cosine: norms precomputed once (bit-identical gate) - NC capabilities poll logs demoted to debug; tracing-appender dep added for the log-writer benchmark Frontend: - ResourceList.selectedEntries O(N)-per-toggle → id-index projection O(k log k); favorites/recent consume the batchToolbar snippet param and drop their duplicate filter + dead selectedIds mirror - Recent: star state via new favoriteIds prop — a star click no longer rebuilds all N entries - admin timeAgo >30d fallback uses the cached Intl.DateTimeFormat - vitest gates in src/lib/components/round11.bench.test.ts Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01ABhTEHuGujvwoodh67Kga7
359 lines
11 KiB
Rust
359 lines
11 KiB
Rust
//! `RetryBlobBackend` — exponential-backoff retry + optional bandwidth throttling
|
||
//! decorator for remote blob backends.
|
||
//!
|
||
//! Wraps any `BlobStorageBackend` and retries transient failures with configurable
|
||
//! exponential backoff. Optionally throttles upload/download bandwidth via
|
||
//! inter-chunk sleeps.
|
||
|
||
use std::path::{Path, PathBuf};
|
||
use std::pin::Pin;
|
||
use std::sync::Arc;
|
||
use std::time::Duration;
|
||
|
||
use crate::application::ports::blob_storage_ports::{
|
||
BlobStorageBackend, BlobStream, StorageHealthStatus,
|
||
};
|
||
use crate::domain::errors::DomainError;
|
||
use bytes::Bytes;
|
||
|
||
// ── Retry policy ───────────────────────────────────────────────────
|
||
|
||
/// Exponential backoff retry configuration.
|
||
#[derive(Debug, Clone)]
|
||
pub struct RetryPolicy {
|
||
/// Maximum number of retry attempts (0 = no retries).
|
||
pub max_retries: u32,
|
||
/// Initial backoff duration before the first retry.
|
||
pub initial_backoff: Duration,
|
||
/// Maximum backoff duration (capped).
|
||
pub max_backoff: Duration,
|
||
/// Multiplier applied to backoff after each attempt.
|
||
pub backoff_multiplier: f64,
|
||
}
|
||
|
||
impl Default for RetryPolicy {
|
||
fn default() -> Self {
|
||
Self {
|
||
max_retries: 3,
|
||
initial_backoff: Duration::from_millis(100),
|
||
max_backoff: Duration::from_secs(10),
|
||
backoff_multiplier: 2.0,
|
||
}
|
||
}
|
||
}
|
||
|
||
// ── RetryBlobBackend ───────────────────────────────────────────────
|
||
|
||
/// Decorator that retries failed backend operations with exponential backoff.
|
||
pub struct RetryBlobBackend {
|
||
inner: Arc<dyn BlobStorageBackend>,
|
||
policy: RetryPolicy,
|
||
}
|
||
|
||
impl RetryBlobBackend {
|
||
pub fn new(inner: Arc<dyn BlobStorageBackend>, policy: RetryPolicy) -> Self {
|
||
Self { inner, policy }
|
||
}
|
||
}
|
||
|
||
/// Execute an async closure with exponential backoff retry.
|
||
///
|
||
/// `name` is a lazy label: the success path (the overwhelmingly common
|
||
/// case) never materializes it, so per-op `format!("op({hash})")`
|
||
/// allocations only happen on an actual retry (benches/ROUND11.md §14:
|
||
/// 64.5 → 0.7 ns, −2 allocs per blob op).
|
||
async fn retry_async<F, Fut, T, L>(
|
||
policy: &RetryPolicy,
|
||
name: L,
|
||
mut f: F,
|
||
) -> Result<T, DomainError>
|
||
where
|
||
F: FnMut() -> Fut,
|
||
Fut: std::future::Future<Output = Result<T, DomainError>>,
|
||
L: Fn() -> String,
|
||
{
|
||
let mut attempt = 0u32;
|
||
let mut backoff = policy.initial_backoff;
|
||
|
||
loop {
|
||
match f().await {
|
||
Ok(v) => return Ok(v),
|
||
Err(e) if attempt < policy.max_retries && is_retryable(&e) => {
|
||
attempt += 1;
|
||
tracing::warn!(
|
||
"Retry {}/{} for {} after error: {} (backoff {:?})",
|
||
attempt,
|
||
policy.max_retries,
|
||
name(),
|
||
e,
|
||
backoff
|
||
);
|
||
tokio::time::sleep(backoff).await;
|
||
let next =
|
||
Duration::from_secs_f64(backoff.as_secs_f64() * policy.backoff_multiplier);
|
||
backoff = next.min(policy.max_backoff);
|
||
}
|
||
Err(e) => return Err(e),
|
||
}
|
||
}
|
||
}
|
||
|
||
/// Determine if an error is likely transient (network timeout, 5xx, etc.).
|
||
fn is_retryable(err: &DomainError) -> bool {
|
||
let msg = err.to_string().to_lowercase();
|
||
msg.contains("timeout")
|
||
|| msg.contains("connection")
|
||
|| msg.contains("503")
|
||
|| msg.contains("500")
|
||
|| msg.contains("429")
|
||
|| msg.contains("temporarily")
|
||
|| msg.contains("broken pipe")
|
||
|| msg.contains("reset by peer")
|
||
}
|
||
|
||
impl BlobStorageBackend for RetryBlobBackend {
|
||
fn initialize(
|
||
&self,
|
||
) -> Pin<Box<dyn std::future::Future<Output = Result<(), DomainError>> + Send + '_>> {
|
||
let inner = self.inner.clone();
|
||
let policy = self.policy.clone();
|
||
Box::pin(async move {
|
||
retry_async(
|
||
&policy,
|
||
|| "initialize".to_string(),
|
||
|| {
|
||
let inner = inner.clone();
|
||
async move { inner.initialize().await }
|
||
},
|
||
)
|
||
.await
|
||
})
|
||
}
|
||
|
||
fn put_blob(
|
||
&self,
|
||
hash: &str,
|
||
source_path: &Path,
|
||
) -> 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();
|
||
let path = source_path.to_path_buf();
|
||
Box::pin(async move {
|
||
retry_async(
|
||
&policy,
|
||
|| format!("put_blob({hash})"),
|
||
|| {
|
||
let inner = inner.clone();
|
||
let hash = hash.clone();
|
||
let path = path.clone();
|
||
async move { inner.put_blob(&hash, &path).await }
|
||
},
|
||
)
|
||
.await
|
||
})
|
||
}
|
||
|
||
fn put_blob_from_bytes(
|
||
&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({hash})"),
|
||
|| {
|
||
let inner = inner.clone();
|
||
let hash = hash.clone();
|
||
let data = data.clone();
|
||
async move { inner.put_blob_from_bytes(&hash, data).await }
|
||
},
|
||
)
|
||
.await
|
||
})
|
||
}
|
||
|
||
// Without this override the trait default would re-route the CDC chunk
|
||
// write through `put_blob_from_bytes` above — reinstating the remote
|
||
// backend's exists-probe (HEAD/get_properties) per chunk that the
|
||
// `_unsynced` fast path exists to skip.
|
||
fn put_blob_from_bytes_unsynced(
|
||
&self,
|
||
hash: &str,
|
||
data: Bytes,
|
||
) -> Pin<Box<dyn std::future::Future<Output = Result<u64, DomainError>> + Send + '_>> {
|
||
let inner = self.inner.clone();
|
||
let policy = self.policy.clone();
|
||
let hash = hash.to_string();
|
||
Box::pin(async move {
|
||
retry_async(
|
||
&policy,
|
||
|| format!("put_blob_from_bytes_unsynced({hash})"),
|
||
|| {
|
||
let inner = inner.clone();
|
||
let hash = hash.clone();
|
||
let data = data.clone();
|
||
async move { inner.put_blob_from_bytes_unsynced(&hash, data).await }
|
||
},
|
||
)
|
||
.await
|
||
})
|
||
}
|
||
|
||
// Forwarded WITHOUT retry wrapping: a failed fsync must surface, not be
|
||
// re-issued — after an fsync error the kernel may have dropped the dirty
|
||
// pages, so a retried fsync can report success for data that was lost.
|
||
fn sync_blobs(
|
||
&self,
|
||
hashes: &[String],
|
||
) -> Pin<Box<dyn std::future::Future<Output = Result<(), DomainError>> + Send + '_>> {
|
||
self.inner.sync_blobs(hashes)
|
||
}
|
||
|
||
fn get_blob_stream(
|
||
&self,
|
||
hash: &str,
|
||
) -> Pin<Box<dyn std::future::Future<Output = Result<BlobStream, 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!("get_blob_stream({hash})"),
|
||
|| {
|
||
let inner = inner.clone();
|
||
let hash = hash.clone();
|
||
async move { inner.get_blob_stream(&hash).await }
|
||
},
|
||
)
|
||
.await
|
||
})
|
||
}
|
||
|
||
fn get_blob_range_stream(
|
||
&self,
|
||
hash: &str,
|
||
start: u64,
|
||
end: Option<u64>,
|
||
) -> Pin<Box<dyn std::future::Future<Output = Result<BlobStream, 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!("get_blob_range({hash})"),
|
||
|| {
|
||
let inner = inner.clone();
|
||
let hash = hash.clone();
|
||
async move { inner.get_blob_range_stream(&hash, start, end).await }
|
||
},
|
||
)
|
||
.await
|
||
})
|
||
}
|
||
|
||
fn delete_blob(
|
||
&self,
|
||
hash: &str,
|
||
) -> Pin<Box<dyn std::future::Future<Output = Result<(), 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!("delete_blob({hash})"),
|
||
|| {
|
||
let inner = inner.clone();
|
||
let hash = hash.clone();
|
||
async move { inner.delete_blob(&hash).await }
|
||
},
|
||
)
|
||
.await
|
||
})
|
||
}
|
||
|
||
fn blob_exists(
|
||
&self,
|
||
hash: &str,
|
||
) -> Pin<Box<dyn std::future::Future<Output = Result<bool, 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!("blob_exists({hash})"),
|
||
|| {
|
||
let inner = inner.clone();
|
||
let hash = hash.clone();
|
||
async move { inner.blob_exists(&hash).await }
|
||
},
|
||
)
|
||
.await
|
||
})
|
||
}
|
||
|
||
fn blob_size(
|
||
&self,
|
||
hash: &str,
|
||
) -> 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!("blob_size({hash})"),
|
||
|| {
|
||
let inner = inner.clone();
|
||
let hash = hash.clone();
|
||
async move { inner.blob_size(&hash).await }
|
||
},
|
||
)
|
||
.await
|
||
})
|
||
}
|
||
|
||
fn health_check(
|
||
&self,
|
||
) -> Pin<
|
||
Box<dyn std::future::Future<Output = Result<StorageHealthStatus, DomainError>> + Send + '_>,
|
||
> {
|
||
let inner = self.inner.clone();
|
||
let policy = self.policy.clone();
|
||
Box::pin(async move {
|
||
retry_async(
|
||
&policy,
|
||
|| "health_check".to_string(),
|
||
|| {
|
||
let inner = inner.clone();
|
||
async move { inner.health_check().await }
|
||
},
|
||
)
|
||
.await
|
||
})
|
||
}
|
||
|
||
fn backend_type(&self) -> &'static str {
|
||
"retry"
|
||
}
|
||
|
||
/// Transparent wrapper: the inner backend serves the bytes.
|
||
fn read_prefetch(&self) -> usize {
|
||
self.inner.read_prefetch()
|
||
}
|
||
|
||
fn local_blob_path(&self, hash: &str) -> Option<PathBuf> {
|
||
self.inner.local_blob_path(hash)
|
||
}
|
||
}
|