use std::sync::Mutex;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
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};
const NO_EXPIRY: u64 = 0;
const PREFIX_LEN: usize = 24;
const EXPIRY_LEN: usize = 8;
const DIGEST_LEN: usize = 16;
pub struct L4SledBackend<V> {
tree: sled::Tree,
health: Mutex<TierHealth>,
_marker: std::marker::PhantomData<V>,
}
impl<V> std::fmt::Debug for L4SledBackend<V> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("L4SledBackend").finish_non_exhaustive()
}
}
fn now_millis() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(NO_EXPIRY, |d| d.as_millis() as u64)
}
impl<V: ByteValue> L4SledBackend<V> {
pub fn open<P: AsRef<std::path::Path>>(path: P) -> Result<Self, CacheError> {
let db = sled::open(path).map_err(|_| CacheError::ConfigurationError)?;
let tree = db
.open_tree(b"thesix")
.map_err(|_| CacheError::ConfigurationError)?;
Ok(L4SledBackend {
tree,
health: Mutex::new(TierHealth::default()),
_marker: std::marker::PhantomData,
})
}
fn succeed(&self) {
if let Ok(mut h) = self.health.lock() {
h.consecutive_failures = 0;
h.health_score = 1.0;
}
}
fn finish_mutation(&self, result: &sled::Result<Option<sled::IVec>>) -> Result<(), CacheError> {
if result.is_ok() {
self.succeed();
Ok(())
} else {
self.fail();
Err(CacheError::TierUnavailable)
}
}
fn fail(&self) {
if let Ok(mut h) = self.health.lock() {
h.consecutive_failures += 1;
h.last_failure_timestamp = Some(SystemTime::now());
h.health_score = (h.health_score - 0.1).max(0.0);
}
}
fn is_expired(deadline_ms: u64) -> bool {
deadline_ms != NO_EXPIRY && now_millis() >= deadline_ms
}
fn decode_record(raw: &[u8]) -> Result<Option<(u64, V)>, CacheError>
where
V: ByteValue,
{
if raw.len() < PREFIX_LEN {
return Err(CacheError::SerializationFailed);
}
let mut expiry = [0u8; EXPIRY_LEN];
expiry.copy_from_slice(&raw[..EXPIRY_LEN]);
let mut digest = [0u8; DIGEST_LEN];
digest.copy_from_slice(&raw[EXPIRY_LEN..PREFIX_LEN]);
let deadline = u64::from_le_bytes(expiry);
if Self::is_expired(deadline) {
return Ok(None);
}
let payload = &raw[PREFIX_LEN..];
let expected = u128::from_le_bytes(digest);
if crate::integrity::ContentDigest::of_bytes(payload).value() != expected {
return Err(CacheError::Corrupted);
}
V::decode_bytes(payload).map(|v| Some((deadline, v)))
}
}
#[async_trait::async_trait]
impl<V: ByteValue> CacheTier<V> for L4SledBackend<V> {
fn name(&self) -> String {
"L4-sled-persistent".to_string()
}
fn backend(&self) -> BackendKind {
BackendKind::Sled
}
fn capability(&self) -> crate::capability::TierCapability {
crate::capability::TierCapability::new(
BackendKind::Sled,
crate::capability::CapabilityFlags::PERSISTENT
| crate::capability::CapabilityFlags::BLOCKING_IO
| crate::capability::CapabilityFlags::ATOMIC_WRITE_OR_ERROR,
crate::capability::OperationalState::Healthy,
crate::capability::DurabilityClass::Delegated,
)
}
async fn get(&self, key: &KeyRef<'_>) -> Result<Option<V>, CacheError> {
match self.tree.get(key.0) {
Ok(Some(raw)) => match Self::decode_record(&raw) {
Ok(None) => {
let _ = self.tree.remove(key.0);
self.succeed();
Ok(None)
}
Ok(Some((_, v))) => {
self.succeed();
Ok(Some(v))
}
Err(e) => {
self.fail();
Err(e)
}
},
Ok(None) => {
self.succeed();
Ok(None)
}
Err(_) => {
self.fail();
Err(CacheError::TierUnavailable)
}
}
}
async fn set(
&self,
key: &KeyRef<'_>,
value: V,
ttl: Option<Duration>,
) -> Result<(), CacheError> {
let payload = value.encode_bytes()?;
let deadline = ttl.map_or(NO_EXPIRY, |t| {
now_millis().saturating_add(t.as_millis() as u64)
});
let mut record = Vec::with_capacity(PREFIX_LEN + payload.len());
record.extend_from_slice(&deadline.to_le_bytes());
record.extend_from_slice(
&crate::integrity::ContentDigest::of_bytes(&payload)
.value()
.to_le_bytes(),
);
record.extend_from_slice(&payload);
self.finish_mutation(&self.tree.insert(key.0, record))
}
async fn remove(&self, key: &KeyRef<'_>) -> Result<(), CacheError> {
self.finish_mutation(&self.tree.remove(key.0))
}
async fn contains(&self, key: &KeyRef<'_>) -> Result<bool, CacheError> {
match self.tree.get(key.0) {
Ok(Some(raw)) => {
let present = Self::decode_record(&raw)?.is_some();
self.succeed();
Ok(present)
}
Ok(None) => {
self.succeed();
Ok(false)
}
Err(_) => {
self.fail();
Err(CacheError::TierUnavailable)
}
}
}
fn health(&self) -> TierHealth {
self.health
.lock()
.map_or_else(|_| TierHealth::default(), |h| h.clone())
}
fn tier_id(&self) -> TierId {
TierId::L4
}
}
impl<V: ByteValue> L4SledBackend<V> {
pub async fn truncate_record_for_test(&self, key: &[u8], len: usize) -> bool {
match self.tree.get(key) {
Ok(Some(raw)) if raw.len() > len => self.tree.insert(key, &raw[..len]).is_ok(),
_ => false,
}
}
pub async fn corrupt_payload_for_test(&self, key: &[u8], offset: usize) -> bool {
match self.tree.get(key) {
Ok(Some(raw)) if raw.len() > PREFIX_LEN + offset => {
let mut damaged = raw.to_vec();
let idx = PREFIX_LEN + offset;
damaged[idx] ^= 0xFF;
self.tree.insert(key, damaged).is_ok()
}
_ => false,
}
}
pub async fn flush_for_test(&self) {
let _ = self.tree.flush();
}
}