fix(cache): self-heal Redis connection + don't 500 on cache write failure
Deploy Scraper / build-and-deploy (push) Canceled after 0s
Deploy Scraper / build-and-deploy (push) Canceled after 0s
The scraper cached a single RedisCache (one multiplexed connection) in a OnceCell forever. When that one connection broke (Redis restart, idle timeout, network blip), every cache op failed with 'cache io: broken pipe', and because Cache::get_or_set propagated the post-compute write error, EVERY API (anime, anime2, komik) returned 500 until a process restart. Fix: - Build a fresh RedisCache from a freshly checked-out deadpool connection per call, so a broken connection self-heals without a process restart (deadpool recycles/drops dead connections and reconnects on checkout). - Make Cache::get_or_set treat a cache-write failure as non-fatal: return the freshly computed value (cache is best-effort), so a transient Redis outage degrades to cache-less instead of 500. Pre-existing clippy warnings (repositories/parsers) untouched — out of scope.
This commit is contained in:
+24
-25
@@ -6,39 +6,38 @@
|
|||||||
//! same `redis` crate version, a `deadpool_redis::Connection` can be converted
|
//! same `redis` crate version, a `deadpool_redis::Connection` can be converted
|
||||||
//! directly into a `redis::aio::MultiplexedConnection` via `take()`.
|
//! directly into a `redis::aio::MultiplexedConnection` via `take()`.
|
||||||
|
|
||||||
use std::sync::LazyLock;
|
|
||||||
|
|
||||||
use mytheclipse_cache::{Cache, CacheError, RedisCache};
|
use mytheclipse_cache::{Cache, CacheError, RedisCache};
|
||||||
use tokio::sync::OnceCell;
|
|
||||||
|
|
||||||
use crate::infrastructure::cache::redis_pool::redis_pool;
|
use crate::infrastructure::cache::redis_pool::redis_pool;
|
||||||
|
|
||||||
/// Lazily-initialised shared `RedisCache` built from the deadpool pool.
|
/// Build a fresh `RedisCache` from a freshly checked-out deadpool connection
|
||||||
|
/// on each call.
|
||||||
///
|
///
|
||||||
/// The multiplexed connection is cheaply cloneable (Arc-backed), so the whole
|
/// Why fresh on every call (not a cached singleton): the previous design
|
||||||
/// process shares one logical connection while deadpool manages recycling.
|
/// initialised *one* `RedisCache` (one multiplexed connection) on first use and
|
||||||
static REDIS_CACHE: LazyLock<OnceCell<RedisCache>> = LazyLock::new(OnceCell::new);
|
/// kept it forever. If that single connection broke (Redis restart, idle
|
||||||
|
/// timeout, network blip), every cache operation failed with `broken pipe`
|
||||||
/// Return a handle to the shared mytheclipse `RedisCache`, initialising it on
|
/// until the whole process restarted — and because `Cache::get_or_set`
|
||||||
/// first use from the deadpool pool.
|
/// propagated write errors, **every API returned 500**.
|
||||||
pub async fn redis_cache() -> Result<&'static RedisCache, CacheError> {
|
///
|
||||||
let cell = &*REDIS_CACHE;
|
/// Building fresh lets deadpool recycle and re-establish broken connections on
|
||||||
cell.get_or_try_init(|| async {
|
/// checkout, so caching self-heals without a process restart. The multiplexed
|
||||||
let pool = redis_pool().map_err(CacheError::Io)?;
|
/// connection is cheaply cloneable (`Arc`-backed), so a per-call pool checkout
|
||||||
let conn = pool
|
/// is negligible overhead next to the network I/O.
|
||||||
.get()
|
pub(crate) async fn fresh_redis_cache() -> Result<RedisCache, CacheError> {
|
||||||
.await
|
let pool = redis_pool().map_err(CacheError::Io)?;
|
||||||
.map_err(|e| CacheError::Io(e.to_string()))?;
|
let conn = pool
|
||||||
let mux = deadpool_redis::Connection::take(conn);
|
.get()
|
||||||
Ok(RedisCache::new(mux))
|
.await
|
||||||
})
|
.map_err(|e| CacheError::Io(e.to_string()))?;
|
||||||
.await
|
let mux = deadpool_redis::Connection::take(conn);
|
||||||
|
Ok(RedisCache::new(mux))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Convenience wrappers so callers can use the mytheclipse `Cache` methods
|
/// Convenience wrappers so callers can use the mytheclipse `Cache` methods
|
||||||
/// directly without importing the trait twice.
|
/// directly without importing the trait twice.
|
||||||
pub async fn get(key: &str) -> Result<Option<Vec<u8>>, CacheError> {
|
pub async fn get(key: &str) -> Result<Option<Vec<u8>>, CacheError> {
|
||||||
redis_cache().await?.get(key).await
|
fresh_redis_cache().await?.get(key).await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn set(
|
pub async fn set(
|
||||||
@@ -46,9 +45,9 @@ pub async fn set(
|
|||||||
value: Vec<u8>,
|
value: Vec<u8>,
|
||||||
ttl: Option<std::time::Duration>,
|
ttl: Option<std::time::Duration>,
|
||||||
) -> Result<(), CacheError> {
|
) -> Result<(), CacheError> {
|
||||||
redis_cache().await?.set(key, value, ttl).await
|
fresh_redis_cache().await?.set(key, value, ttl).await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn invalidate(key: &str) -> Result<(), CacheError> {
|
pub async fn invalidate(key: &str) -> Result<(), CacheError> {
|
||||||
redis_cache().await?.invalidate(key).await
|
fresh_redis_cache().await?.invalidate(key).await
|
||||||
}
|
}
|
||||||
|
|||||||
Vendored
+17
-6
@@ -24,13 +24,13 @@ impl<'a> Cache<'a> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn cache(&self) -> Result<&'static mytheclipse_cache::RedisCache, CacheError> {
|
async fn cache(&self) -> Result<mytheclipse_cache::RedisCache, CacheError> {
|
||||||
super::mytheclipse::redis_cache().await
|
super::mytheclipse::fresh_redis_cache().await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn get<T: DeserializeOwned>(&self, key: &str) -> Option<T> {
|
pub async fn get<T: DeserializeOwned>(&self, key: &str) -> Option<T> {
|
||||||
match self.cache().await {
|
match self.cache().await {
|
||||||
Ok(cache) => match CacheTrait::get(cache, key).await {
|
Ok(cache) => match CacheTrait::get(&cache, key).await {
|
||||||
Ok(Some(bytes)) => serde_json::from_slice(&bytes).ok(),
|
Ok(Some(bytes)) => serde_json::from_slice(&bytes).ok(),
|
||||||
Ok(None) => None,
|
Ok(None) => None,
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
@@ -69,7 +69,7 @@ impl<'a> Cache<'a> {
|
|||||||
let json = serde_json::to_vec(value).map_err(|e| e.to_string())?;
|
let json = serde_json::to_vec(value).map_err(|e| e.to_string())?;
|
||||||
let ttl = std::time::Duration::from_secs(ttl_secs);
|
let ttl = std::time::Duration::from_secs(ttl_secs);
|
||||||
let cache = self.cache().await.map_err(|e| e.to_string())?;
|
let cache = self.cache().await.map_err(|e| e.to_string())?;
|
||||||
CacheTrait::set(cache, key, json, Some(ttl))
|
CacheTrait::set(&cache, key, json, Some(ttl))
|
||||||
.await
|
.await
|
||||||
.map_err(|e| e.to_string())?;
|
.map_err(|e| e.to_string())?;
|
||||||
debug!("Cache: set key {} with TTL {}s", key, ttl_secs);
|
debug!("Cache: set key {} with TTL {}s", key, ttl_secs);
|
||||||
@@ -78,7 +78,7 @@ impl<'a> Cache<'a> {
|
|||||||
|
|
||||||
pub async fn delete(&self, key: &str) -> Result<(), String> {
|
pub async fn delete(&self, key: &str) -> Result<(), String> {
|
||||||
let cache = self.cache().await.map_err(|e| e.to_string())?;
|
let cache = self.cache().await.map_err(|e| e.to_string())?;
|
||||||
CacheTrait::invalidate(cache, key)
|
CacheTrait::invalidate(&cache, key)
|
||||||
.await
|
.await
|
||||||
.map_err(|e| e.to_string())
|
.map_err(|e| e.to_string())
|
||||||
}
|
}
|
||||||
@@ -88,6 +88,12 @@ impl<'a> Cache<'a> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Get or set: returns cached value or computes and caches new value.
|
/// Get or set: returns cached value or computes and caches new value.
|
||||||
|
///
|
||||||
|
/// The cache is best-effort: if the post-compute write fails (e.g. a
|
||||||
|
/// transient Redis outage or a broken pooled connection), the freshly
|
||||||
|
/// computed value is still returned rather than propagating a 500. Only a
|
||||||
|
/// cache *read* failure is silently tolerated; a compute failure still
|
||||||
|
/// propagates.
|
||||||
pub async fn get_or_set<T, F, Fut>(
|
pub async fn get_or_set<T, F, Fut>(
|
||||||
&self,
|
&self,
|
||||||
key: &str,
|
key: &str,
|
||||||
@@ -106,7 +112,12 @@ impl<'a> Cache<'a> {
|
|||||||
|
|
||||||
debug!("Cache miss: {}", key);
|
debug!("Cache miss: {}", key);
|
||||||
let value = compute().await?;
|
let value = compute().await?;
|
||||||
self.set_with_ttl(key, &value, ttl_secs).await?;
|
// Best-effort write: a failure here must not fail the request — the
|
||||||
|
// value is already valid. Log and continue (read path swallows errors
|
||||||
|
// too, so a broken cache degrades to cache-less, never to 500).
|
||||||
|
if let Err(e) = self.set_with_ttl(key, &value, ttl_secs).await {
|
||||||
|
debug!("Cache: failed to write {} (non-fatal): {}", key, e);
|
||||||
|
}
|
||||||
Ok(value)
|
Ok(value)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user