use std::sync::Mutex;
use std::time::Duration;
use crate::error::CacheError;
use crate::key::KeyRef;
use crate::tier::TierId;
use crate::tier::backends::ByteValue;
use crate::tier::tier_trait::{BackendKind, CacheTier, TierHealth};
pub type OriginFetcher = dyn Fn(&[u8]) -> Result<Option<Vec<u8>>, CacheError> + Send + Sync;
pub type OriginWriter =
dyn Fn(&[u8], &[u8], Option<Duration>) -> Result<(), CacheError> + Send + Sync;
pub struct L5OriginBackend<V> {
fetcher: Box<OriginFetcher>,
writer: Option<Box<OriginWriter>>,
health: Mutex<TierHealth>,
_marker: std::marker::PhantomData<V>,
}
impl<V> std::fmt::Debug for L5OriginBackend<V> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("L5OriginBackend").finish_non_exhaustive()
}
}
impl<V: ByteValue> L5OriginBackend<V> {
pub fn new(fetcher: Box<OriginFetcher>) -> Self {
L5OriginBackend {
fetcher,
writer: None,
health: Mutex::new(TierHealth::default()),
_marker: std::marker::PhantomData,
}
}
pub fn with_writer(mut self, writer: Box<OriginWriter>) -> Self {
self.writer = Some(writer);
self
}
fn succeed(&self) {
if let Ok(mut h) = self.health.lock() {
h.consecutive_failures = 0;
h.health_score = 1.0;
}
}
fn fail(&self) {
if let Ok(mut h) = self.health.lock() {
h.consecutive_failures += 1;
h.last_failure_timestamp = Some(std::time::SystemTime::now());
h.health_score = (h.health_score - 0.1).max(0.0);
}
}
}
#[async_trait::async_trait]
impl<V: ByteValue> CacheTier<V> for L5OriginBackend<V> {
fn name(&self) -> String {
"L5-origin-fallback".to_string()
}
fn backend(&self) -> BackendKind {
BackendKind::Origin
}
fn capability(&self) -> crate::capability::TierCapability {
crate::capability::TierCapability::new(
BackendKind::Origin,
crate::capability::CapabilityFlags::SHARED
| crate::capability::CapabilityFlags::AUTHORITATIVE,
crate::capability::OperationalState::Healthy,
crate::capability::DurabilityClass::Delegated,
)
}
async fn get(&self, key: &KeyRef<'_>) -> Result<Option<V>, CacheError> {
match (self.fetcher)(key.0) {
Ok(Some(bytes)) => {
self.succeed();
V::decode_bytes(&bytes).map(Some)
}
Ok(None) => {
self.succeed();
Ok(None)
}
Err(e) => {
self.fail();
Err(e)
}
}
}
async fn set(
&self,
key: &KeyRef<'_>,
value: V,
ttl: Option<Duration>,
) -> Result<(), CacheError> {
let Some(writer) = &self.writer else {
return Ok(());
};
let bytes = value.encode_bytes()?;
match writer(key.0, &bytes, ttl) {
Ok(()) => {
self.succeed();
Ok(())
}
Err(e) => {
self.fail();
Err(e)
}
}
}
async fn remove(&self, key: &KeyRef<'_>) -> Result<(), CacheError> {
let Some(writer) = &self.writer else {
return Ok(());
};
match writer(key.0, &[], None) {
Ok(()) => {
self.succeed();
Ok(())
}
Err(e) => {
self.fail();
Err(e)
}
}
}
async fn contains(&self, key: &KeyRef<'_>) -> Result<bool, CacheError> {
match (self.fetcher)(key.0) {
Ok(opt) => {
self.succeed();
Ok(opt.is_some())
}
Err(e) => {
self.fail();
Err(e)
}
}
}
fn health(&self) -> TierHealth {
self.health
.lock()
.map_or_else(|_| TierHealth::default(), |h| h.clone())
}
fn tier_id(&self) -> TierId {
TierId::L5
}
}