use std::sync::Arc;
use std::time::Duration;
use std::time::SystemTime;
use std::time::UNIX_EPOCH;
use async_trait::async_trait;
use camel_api::CamelError;
use camel_api::cache::CacheEntry;
use camel_api::cache::CacheRepository;
use camel_api::cache::CacheStats;
use camel_api::cache::ContentType;
use tracing::warn;
pub type OffloadClock = Arc<dyn Fn() -> SystemTime + Send + Sync>;
pub fn default_offload_clock() -> OffloadClock {
Arc::new(SystemTime::now)
}
#[async_trait]
pub trait PayloadStore: Send + Sync {
async fn put(
&self,
name: &str,
bytes: &[u8],
death_epoch: SystemTime,
) -> Result<(), CamelError>;
async fn read(&self, name: &str) -> Result<Option<Vec<u8>>, CamelError>;
async fn unlink(&self, name: &str) -> Result<(), CamelError>;
async fn clear(&self) {}
}
pub struct OffloadRepository {
inner: Arc<dyn CacheRepository>,
store: Arc<dyn PayloadStore>,
stale_retention: Duration,
sweep_interval: Duration,
payload_max_ttl: Duration,
clock: OffloadClock,
}
impl OffloadRepository {
pub fn new(
inner: Arc<dyn CacheRepository>,
store: Arc<dyn PayloadStore>,
stale_retention: Duration,
sweep_interval: Duration,
payload_max_ttl: Duration,
) -> Self {
Self::with_clock(
inner,
store,
stale_retention,
sweep_interval,
payload_max_ttl,
default_offload_clock(),
)
}
pub fn with_clock(
inner: Arc<dyn CacheRepository>,
store: Arc<dyn PayloadStore>,
stale_retention: Duration,
sweep_interval: Duration,
payload_max_ttl: Duration,
clock: OffloadClock,
) -> Self {
Self {
inner,
store,
stale_retention,
sweep_interval,
payload_max_ttl,
clock,
}
}
pub(crate) async fn hydrate(
&self,
key: &str,
mut entry: CacheEntry,
) -> Result<Option<CacheEntry>, CamelError> {
let Some(raw_path) = entry.payload_path.clone() else {
return Ok(Some(entry));
};
let Some(name) = sanitize_blob_name(&raw_path) else {
warn!(
key = key,
backend = self.inner.name(),
payload_path = %raw_path,
"corrupt cache row: payload_path must be a bare file name; treating as miss"
);
return Ok(None);
};
match self.store.read(name).await {
Ok(Some(bytes)) => {
entry.bytes = bytes;
entry.payload_path = None;
Ok(Some(entry))
}
Ok(None) => {
warn!(
key = key,
backend = self.inner.name(),
blob = %name,
"cache payload blob gone; treating as miss"
);
Ok(None)
}
Err(e) => Err(e),
}
}
async fn reclaim_predecessor(
&self,
key: &str,
old_name: Option<&str>,
keep_name: Option<&str>,
) {
let Some(old_name) = old_name else {
return;
};
if keep_name == Some(old_name) {
return;
}
let key_prefix = format!("{}.", blake3_128hex(key.as_bytes()));
let eligible = sanitize_blob_name(old_name).is_some()
&& parse_death_epoch(old_name).is_some()
&& old_name.starts_with(&key_prefix);
if !eligible {
return;
}
if let Err(e) = self.store.unlink(old_name).await {
warn!(
key = key,
backend = self.inner.name(),
error = %e,
"eager reclaim of predecessor blob failed; sweeper reclaims it at its death epoch"
);
}
}
}
#[async_trait]
impl CacheRepository for OffloadRepository {
fn name(&self) -> &str {
self.inner.name()
}
async fn get(&self, key: &str) -> Result<Option<CacheEntry>, CamelError> {
match self.inner.get(key).await? {
Some(entry) => self.hydrate(key, entry).await,
None => Ok(None),
}
}
async fn set(
&self,
key: &str,
mut entry: CacheEntry,
ttl: Option<Duration>,
) -> Result<(), CamelError> {
let effective_ttl = ttl.unwrap_or(self.payload_max_ttl);
let old_name = match self.inner.peek_row_silent(key).await {
Ok(Some(row)) => row.payload_path,
Ok(None) => None,
Err(e) => {
warn!(
key = key,
backend = self.inner.name(),
error = %e,
"pre-swap row read failed; skipping eager reclaim"
);
None
}
};
let death_epoch = (self.clock)()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.saturating_add(effective_ttl)
.saturating_add(self.stale_retention)
.saturating_add(self.sweep_interval)
.as_secs();
let deadline = UNIX_EPOCH
.checked_add(Duration::from_secs(death_epoch))
.unwrap_or(UNIX_EPOCH);
let dest_name = blob_filename(key, death_epoch, &entry);
match self.store.put(&dest_name, &entry.bytes, deadline).await {
Ok(()) => {
entry.bytes = Vec::new();
entry.payload_path = Some(dest_name.clone());
let result = self.inner.set(key, entry, Some(effective_ttl)).await;
if result.is_ok() {
self.reclaim_predecessor(key, old_name.as_deref(), Some(&dest_name))
.await;
}
result
}
Err(e) => {
warn!(
key = key,
backend = self.inner.name(),
error = %e,
"cache blob write failed; storing entry inline instead"
);
let result = self.inner.set(key, entry, Some(effective_ttl)).await;
if result.is_ok() {
self.reclaim_predecessor(key, old_name.as_deref(), None)
.await;
}
result
}
}
}
async fn peek_stale(&self, key: &str) -> Result<Option<CacheEntry>, CamelError> {
match self.inner.peek_stale(key).await? {
Some(entry) => self.hydrate(key, entry).await,
None => Ok(None),
}
}
async fn invalidate(&self, key: &str) -> Result<(), CamelError> {
self.inner.invalidate(key).await
}
async fn clear(&self) -> Result<(), CamelError> {
self.store.clear().await;
self.inner.clear().await
}
async fn invalidate_prefix(&self, prefix: &str) -> Result<u64, CamelError> {
self.inner.invalidate_prefix(prefix).await
}
async fn stats(&self) -> CacheStats {
self.inner.stats().await
}
}
impl std::fmt::Debug for OffloadRepository {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("OffloadRepository")
.field("inner", &self.inner)
.field("stale_retention", &self.stale_retention)
.field("sweep_interval", &self.sweep_interval)
.field("payload_max_ttl", &self.payload_max_ttl)
.finish()
}
}
fn content_type_discriminant(content_type: ContentType) -> u8 {
match content_type {
ContentType::Bytes => 0,
ContentType::Text => 1,
ContentType::Json => 2,
ContentType::Xml => 3,
}
}
pub(crate) fn hasher_128hex(hasher: blake3::Hasher) -> String {
let hex = hasher.finalize().to_hex().to_string();
hex[..32].to_string()
}
pub(crate) fn blake3_128hex(data: &[u8]) -> String {
let mut hasher = blake3::Hasher::new();
hasher.update(data);
hasher_128hex(hasher)
}
pub(crate) fn content_fingerprint(entry: &CacheEntry) -> String {
let mut hasher = blake3::Hasher::new();
hasher.update(&entry.bytes);
hasher.update(&[content_type_discriminant(entry.content_type)]);
hasher_128hex(hasher)
}
fn blob_filename(key: &str, death_epoch: u64, entry: &CacheEntry) -> String {
format!(
"{}.{}.{}.blob",
blake3_128hex(key.as_bytes()),
death_epoch,
content_fingerprint(entry)
)
}
pub(crate) fn parse_death_epoch(file_name: &str) -> Option<u64> {
file_name.split('.').nth(1)?.parse().ok()
}
pub(crate) fn sanitize_blob_name(path: &str) -> Option<&str> {
if path.is_empty() || path.contains('/') || path.contains('\\') || path.contains("..") {
return None;
}
Some(path)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::cache::MemoryCacheRepository;
fn entry(bytes: Vec<u8>, content_type: ContentType) -> CacheEntry {
CacheEntry {
bytes,
payload_path: None,
content_type,
expires_at: None,
}
}
struct FailingPutStore;
#[async_trait]
impl PayloadStore for FailingPutStore {
async fn put(
&self,
_name: &str,
_bytes: &[u8],
_death_epoch: SystemTime,
) -> Result<(), CamelError> {
Err(CamelError::Io(
"failing-put-store: injected put failure".into(),
))
}
async fn read(&self, _name: &str) -> Result<Option<Vec<u8>>, CamelError> {
Ok(None)
}
async fn unlink(&self, _name: &str) -> Result<(), CamelError> {
Ok(())
}
}
#[tokio::test]
async fn offload_repository_inline_fallback_on_store_put_failure() {
let inner = Arc::new(MemoryCacheRepository::new("test", 100));
let repo = OffloadRepository::new(
inner,
Arc::new(FailingPutStore),
Duration::from_secs(168 * 3600),
Duration::from_secs(3600),
Duration::from_secs(24 * 3600),
);
let payload = vec![7; 128];
repo.set(
"k",
entry(payload.clone(), ContentType::Bytes),
Some(Duration::from_secs(60)),
)
.await
.expect("set must fall back inline, never Err");
let got = repo.get("k").await.expect("get").expect("present");
assert_eq!(
got.bytes, payload,
"bytes must be present, stored inline in the index"
);
assert_eq!(got.content_type, ContentType::Bytes);
assert!(got.payload_path.is_none(), "fallback row stays inline");
}
struct ForgetfulStore;
#[async_trait]
impl PayloadStore for ForgetfulStore {
async fn put(
&self,
_name: &str,
_bytes: &[u8],
_death_epoch: SystemTime,
) -> Result<(), CamelError> {
Ok(())
}
async fn read(&self, _name: &str) -> Result<Option<Vec<u8>>, CamelError> {
Ok(None)
}
async fn unlink(&self, _name: &str) -> Result<(), CamelError> {
Ok(())
}
}
#[tokio::test]
async fn offload_repository_missing_payload_is_miss() {
let inner = Arc::new(MemoryCacheRepository::new("test", 100));
let repo = OffloadRepository::new(
inner,
Arc::new(ForgetfulStore),
Duration::from_secs(168 * 3600),
Duration::from_secs(3600),
Duration::from_secs(24 * 3600),
);
repo.set(
"k",
entry(vec![1, 2, 3], ContentType::Bytes),
Some(Duration::from_secs(60)),
)
.await
.expect("set");
assert_eq!(
repo.get("k").await.expect("get"),
None,
"a vanished payload degrades to a MISS, never an error"
);
}
}