perf(blob-cache,db): release index mutex before disk I/O; skip per-acquire DB ping

CachedBlobBackend held its single tokio::Mutex<LruCache> across filesystem
syscalls, serializing every concurrent cache operation behind one lock:
- get_blob_stream / get_blob_range_stream: held across File::open()/seek()
- delete_blob: held across remove_file()
- initialize: held across the full cache-dir walk
- eviction (insert + fetch paths): held across remove_file() loops

Now the lock only guards the in-memory LRU. Presence checks bump recency
and release the guard before touching the filesystem (a vanished file falls
through to the existing fetch-and-cache path, covering the race), and
eviction selects victims under the lock then unlinks them after releasing
it. The duplicated eviction loop is extracted into
CachedRef::collect_evictions.

db: set test_before_acquire(false). With warm min_connections and a bounded
max_lifetime, the liveness ping sqlx issues on every acquire() costs more
than the rare dead connection it catches; stale sockets surface as a query
error and the pool recycles them either way.

https://claude.ai/code/session_01UtfkS3nZF1vrF5jNAps6wV
This commit is contained in:
Claude
2026-06-09 13:18:09 +00:00
parent 233531bce5
commit fc55e92299
2 changed files with 92 additions and 67 deletions
+6
View File
@@ -114,6 +114,12 @@ async fn create_pool_with_retries(
.acquire_timeout(Duration::from_secs(connect_timeout_secs)) .acquire_timeout(Duration::from_secs(connect_timeout_secs))
.idle_timeout(Duration::from_secs(idle_timeout_secs)) .idle_timeout(Duration::from_secs(idle_timeout_secs))
.max_lifetime(Duration::from_secs(max_lifetime_secs)) .max_lifetime(Duration::from_secs(max_lifetime_secs))
// Skip the liveness ping sqlx issues on every acquire() (on by
// default): with warm min_connections and a bounded max_lifetime,
// that extra round-trip per checkout costs more than the rare dead
// connection it catches. A stale socket surfaces as a query error
// and the pool recycles it either way.
.test_before_acquire(false)
.connect(connection_string) .connect(connection_string)
.await .await
{ {
@@ -96,9 +96,11 @@ impl BlobStorageBackend for CachedBlobBackend {
DomainError::internal_error("BlobCache", format!("mkdir cache_dir: {e}")) DomainError::internal_error("BlobCache", format!("mkdir cache_dir: {e}"))
})?; })?;
// Scan existing cache to rebuild index // Scan existing cache to rebuild index. Collect entries WITHOUT
// holding the index lock — a large cache directory walk must not
// serialize concurrent blob operations behind the mutex.
let mut total_bytes = 0u64; let mut total_bytes = 0u64;
let mut idx = index.lock().await; let mut entries: Vec<(String, u64)> = Vec::new();
if let Ok(mut read_dir) = fs::read_dir(&cache_dir).await { if let Ok(mut read_dir) = fs::read_dir(&cache_dir).await {
while let Ok(Some(prefix_entry)) = read_dir.next_entry().await { while let Ok(Some(prefix_entry)) = read_dir.next_entry().await {
if !prefix_entry.path().is_dir() { if !prefix_entry.path().is_dir() {
@@ -111,14 +113,20 @@ impl BlobStorageBackend for CachedBlobBackend {
&& let Some(stem) = path.file_stem().and_then(|s| s.to_str()) && let Some(stem) = path.file_stem().and_then(|s| s.to_str())
{ {
let size = fs::metadata(&path).await.map(|m| m.len()).unwrap_or(0); let size = fs::metadata(&path).await.map(|m| m.len()).unwrap_or(0);
idx.put(stem.to_string(), CacheEntry { size }); entries.push((stem.to_string(), size));
total_bytes += size; total_bytes += size;
} }
} }
} }
} }
} }
drop(idx); // 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 });
}
}
current_size.store(total_bytes, Ordering::Relaxed); current_size.store(total_bytes, Ordering::Relaxed);
tracing::info!( tracing::info!(
"Blob cache initialized: {} bytes in cache at {}", "Blob cache initialized: {} bytes in cache at {}",
@@ -196,19 +204,18 @@ impl BlobStorageBackend for CachedBlobBackend {
let max_cache_bytes = self.max_cache_bytes; let max_cache_bytes = self.max_cache_bytes;
let current_size = self.current_size.clone(); let current_size = self.current_size.clone();
Box::pin(async move { Box::pin(async move {
// Check cache // Check cache presence (and bump LRU recency) under a brief lock,
{ // then release it BEFORE touching the filesystem so concurrent
let mut idx = index.lock().await; // readers don't serialize behind a single open() syscall.
if idx.get(&hash).is_some() { if index.lock().await.get(&hash).is_some() {
if let Ok(file) = fs::File::open(&cached).await { if let Ok(file) = fs::File::open(&cached).await {
let stream: BlobStream = let stream: BlobStream =
Box::pin(ReaderStream::with_capacity(file, STREAM_CHUNK_SIZE)); Box::pin(ReaderStream::with_capacity(file, STREAM_CHUNK_SIZE));
return Ok(stream); return Ok(stream);
} }
// Cache entry stale — remove // Cache entry stale (file vanished) — drop it from the index.
if let Some(entry) = idx.pop(&hash) { if let Some(entry) = index.lock().await.pop(&hash) {
current_size.fetch_sub(entry.size, Ordering::Relaxed); current_size.fetch_sub(entry.size, Ordering::Relaxed);
}
} }
} }
@@ -243,25 +250,24 @@ impl BlobStorageBackend for CachedBlobBackend {
let max_cache_bytes = self.max_cache_bytes; let max_cache_bytes = self.max_cache_bytes;
let current_size = self.current_size.clone(); let current_size = self.current_size.clone();
Box::pin(async move { Box::pin(async move {
// Try cache first // Check cache presence (and bump LRU recency) under a brief lock,
{ // then release it BEFORE the open()/seek() syscalls so concurrent
let mut idx = index.lock().await; // range readers don't serialize behind the index mutex.
if idx.get(&hash).is_some() { if index.lock().await.get(&hash).is_some() {
if let Ok(mut file) = fs::File::open(&cached).await { if let Ok(mut file) = fs::File::open(&cached).await {
file.seek(std::io::SeekFrom::Start(start)) file.seek(std::io::SeekFrom::Start(start))
.await .await
.map_err(|e| { .map_err(|e| {
DomainError::internal_error("BlobCache", format!("seek: {e}")) DomainError::internal_error("BlobCache", format!("seek: {e}"))
})?; })?;
let take_len = end.map(|e| e - start + 1).unwrap_or(u64::MAX); let take_len = end.map(|e| e - start + 1).unwrap_or(u64::MAX);
let limited = file.take(take_len); let limited = file.take(take_len);
let stream: BlobStream = let stream: BlobStream =
Box::pin(ReaderStream::with_capacity(limited, STREAM_CHUNK_SIZE)); Box::pin(ReaderStream::with_capacity(limited, STREAM_CHUNK_SIZE));
return Ok(stream); return Ok(stream);
} }
if let Some(entry) = idx.pop(&hash) { if let Some(entry) = index.lock().await.pop(&hash) {
current_size.fetch_sub(entry.size, Ordering::Relaxed); current_size.fetch_sub(entry.size, Ordering::Relaxed);
}
} }
} }
@@ -298,9 +304,9 @@ impl BlobStorageBackend for CachedBlobBackend {
let current_size = self.current_size.clone(); let current_size = self.current_size.clone();
Box::pin(async move { Box::pin(async move {
inner.delete_blob(&hash).await?; inner.delete_blob(&hash).await?;
// Remove from cache // Remove from cache — drop the index lock before the unlink()
let mut idx = index.lock().await; // syscall so deletes don't serialize concurrent cache lookups.
if let Some(entry) = idx.pop(&hash) { if let Some(entry) = index.lock().await.pop(&hash) {
current_size.fetch_sub(entry.size, Ordering::Relaxed); current_size.fetch_sub(entry.size, Ordering::Relaxed);
} }
let _ = fs::remove_file(&cached).await; let _ = fs::remove_file(&cached).await;
@@ -402,6 +408,26 @@ impl CachedRef {
self.cache_dir.join(prefix).join(format!("{hash}.blob")) self.cache_dir.join(prefix).join(format!("{hash}.blob"))
} }
/// 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( async fn insert_into_cache_static(
&self, &self,
hash: &str, hash: &str,
@@ -423,21 +449,19 @@ impl CachedRef {
DomainError::internal_error("BlobCache", format!("cache copy failed: {e}")) DomainError::internal_error("BlobCache", format!("cache copy failed: {e}"))
})?; })?;
let mut idx = self.index.lock().await; // Update the index and pick eviction victims under a single brief
if let Some(old) = idx.put(hash.to_string(), CacheEntry { size }) { // lock, then unlink the evicted files AFTER releasing it — file
self.current_size.fetch_sub(old.size, Ordering::Relaxed); // removal must not run while the index mutex is held.
} let to_evict = {
self.current_size.fetch_add(size, Ordering::Relaxed); let mut idx = self.index.lock().await;
if let Some(old) = idx.put(hash.to_string(), CacheEntry { size }) {
while self.current_size.load(Ordering::Relaxed) > self.max_cache_bytes { self.current_size.fetch_sub(old.size, Ordering::Relaxed);
if let Some((evicted_hash, evicted_entry)) = idx.pop_lru() {
self.current_size
.fetch_sub(evicted_entry.size, Ordering::Relaxed);
let evicted_path = self.cached_path(&evicted_hash);
let _ = fs::remove_file(&evicted_path).await;
} else {
break;
} }
self.current_size.fetch_add(size, Ordering::Relaxed);
self.collect_evictions(&mut idx)
};
for path in to_evict {
let _ = fs::remove_file(&path).await;
} }
Ok(()) Ok(())
} }
@@ -482,21 +506,16 @@ impl CachedRef {
.await .await
.map_err(|e| DomainError::internal_error("BlobCache", format!("rename: {e}")))?; .map_err(|e| DomainError::internal_error("BlobCache", format!("rename: {e}")))?;
let mut idx = self.index.lock().await; let to_evict = {
if let Some(old) = idx.put(hash.to_string(), CacheEntry { size: total }) { let mut idx = self.index.lock().await;
self.current_size.fetch_sub(old.size, Ordering::Relaxed); 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);
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);
let evicted_path = self.cached_path(&evicted_hash);
let _ = fs::remove_file(&evicted_path).await;
} else {
break;
} }
self.current_size.fetch_add(total, Ordering::Relaxed);
self.collect_evictions(&mut idx)
};
for path in to_evict {
let _ = fs::remove_file(&path).await;
} }
Ok(dest) Ok(dest)