Skip to main content

camel_core/cache/
offload.rs

1//! Payload-offload decorator for [`CacheRepository`] backends, with a
2//! pluggable [`PayloadStore`].
3//!
4//! [`OffloadRepository`] wraps any backend (the "index") and moves entry
5//! payloads into a [`PayloadStore`] (local disk, redis, …), storing only
6//! the store-side blob name in the index row. Index rows stay small; the
7//! payload is re-injected on `get`/`peek_stale`.
8//!
9//! # Blob lifecycle
10//!
11//! Blob names are `{blake3-128hex(key)}.{death_epoch_secs}.{blake3-128hex(
12//! bytes || content_type-discriminant)}.blob`. The death epoch —
13//! `expires_at + stale_retention + sweep_interval` — is encoded in the
14//! name so store-side reclaimers can reclaim dead payloads by name alone,
15//! without consulting the index.
16//!
17//! # Failure policy
18//!
19//! - A payload put that fails falls back to storing the entry inline in
20//!   the index (WARN + `inner.set` with the original entry): the
21//!   decorator never converts a store failure into a cache-write `Err`.
22//! - A vanished or corrupt blob row degrades to a miss (`Ok(None)` + WARN).
23//! - A payload that exists but cannot be read (e.g. `PermissionDenied`)
24//!   surfaces as `Err` per ADR-0023 Contract C1.
25
26use std::sync::Arc;
27use std::time::Duration;
28use std::time::SystemTime;
29use std::time::UNIX_EPOCH;
30
31use async_trait::async_trait;
32use camel_api::CamelError;
33use camel_api::cache::CacheEntry;
34use camel_api::cache::CacheRepository;
35use camel_api::cache::CacheStats;
36use camel_api::cache::ContentType;
37use tracing::warn;
38
39/// Injectable wall clock for death-epoch math and deterministic tests.
40///
41/// Mirrors `ClockFn` in `camel-redis-repo::cache_repo`.
42pub type OffloadClock = Arc<dyn Fn() -> SystemTime + Send + Sync>;
43
44/// The default production clock: [`SystemTime::now`].
45pub fn default_offload_clock() -> OffloadClock {
46    Arc::new(SystemTime::now)
47}
48
49/// Pluggable payload storage behind [`OffloadRepository`].
50///
51/// The decorator owns ALL policy — blob naming, death-epoch math, the
52/// index-row shape, and every fallback (ADR-0065) — and calls the store
53/// with finished blob names. Implementations only move opaque bytes.
54#[async_trait]
55pub trait PayloadStore: Send + Sync {
56    /// Store `bytes` under `name` until `death_epoch`.
57    ///
58    /// `name` is the decorator's content-addressed blob name; stores must
59    /// reject names that are not bare single path components. A store that
60    /// cannot honor `death_epoch` returns `Err` — the decorator's inline
61    /// fallback applies, so a payload is never stored without its
62    /// deadline.
63    async fn put(
64        &self,
65        name: &str,
66        bytes: &[u8],
67        death_epoch: SystemTime,
68    ) -> Result<(), CamelError>;
69
70    /// Read the payload stored under `name`.
71    ///
72    /// `Ok(None)` = the payload is gone (reclaimed, evicted, lost): the
73    /// decorator degrades to a MISS with WARN. An existing-but-unreadable
74    /// payload is `Err` (ADR-0023 Contract C1), never a silent `None`.
75    async fn read(&self, name: &str) -> Result<Option<Vec<u8>>, CamelError>;
76
77    /// Remove the payload stored under `name`.
78    ///
79    /// An already-absent payload counts as success (a concurrent
80    /// reclaimer won the race).
81    async fn unlink(&self, name: &str) -> Result<(), CamelError>;
82
83    /// Best-effort bulk reclamation of every stored payload.
84    ///
85    /// Default: no-op — a store whose payloads carry their own deadline
86    /// (e.g. a redis EXAT) or a background sweeper needs no eager clear.
87    /// Implementers must WARN and continue on per-payload failures: the
88    /// decorator never surfaces `clear` failures as `Err` (ADR-0065
89    /// failure policy).
90    async fn clear(&self) {}
91}
92
93/// [`CacheRepository`] decorator that offloads entry payloads to a
94/// [`PayloadStore`].
95///
96/// Wraps any index backend; see the [module docs](self) for the blob
97/// lifecycle and failure policy. `stale_retention`, `sweep_interval`, and
98/// `payload_max_ttl` must be non-zero (the payload intervals at least one
99/// second — the death epoch truncates to whole seconds) — enforced by
100/// `CacheRepoConfig` validation, not here.
101pub struct OffloadRepository {
102    /// Decorated index backend (memory, redb, redis, …).
103    inner: Arc<dyn CacheRepository>,
104    /// Payload store holding the offloaded bytes.
105    store: Arc<dyn PayloadStore>,
106    /// How long an expired entry stays peekable before reclamation.
107    stale_retention: Duration,
108    /// Background sweep cadence of the store; its length is the
109    /// death-epoch grace.
110    sweep_interval: Duration,
111    /// Fabricated TTL for entries stored without an explicit one.
112    payload_max_ttl: Duration,
113    /// Wall clock for death-epoch math.
114    clock: OffloadClock,
115}
116
117impl OffloadRepository {
118    /// Wrap `inner` with payload offload through `store` (production
119    /// clock).
120    pub fn new(
121        inner: Arc<dyn CacheRepository>,
122        store: Arc<dyn PayloadStore>,
123        stale_retention: Duration,
124        sweep_interval: Duration,
125        payload_max_ttl: Duration,
126    ) -> Self {
127        Self::with_clock(
128            inner,
129            store,
130            stale_retention,
131            sweep_interval,
132            payload_max_ttl,
133            default_offload_clock(),
134        )
135    }
136
137    /// Test seam: [`Self::new`] with an injected [`OffloadClock`].
138    ///
139    /// The injected clock drives death-epoch math only; the store's own
140    /// reclaim machinery (sweeper, TTL) always runs on the real clock.
141    pub fn with_clock(
142        inner: Arc<dyn CacheRepository>,
143        store: Arc<dyn PayloadStore>,
144        stale_retention: Duration,
145        sweep_interval: Duration,
146        payload_max_ttl: Duration,
147        clock: OffloadClock,
148    ) -> Self {
149        Self {
150            inner,
151            store,
152            stale_retention,
153            sweep_interval,
154            payload_max_ttl,
155            clock,
156        }
157    }
158
159    /// Re-inject the offloaded payload into an index row (shared by `get`
160    /// and `peek_stale`).
161    ///
162    /// Rows without `payload_path` pass through untouched (legacy/inline).
163    /// A corrupt path or a vanished payload degrades to a miss; a payload
164    /// that exists but cannot be read surfaces as `Err` (Contract C1).
165    pub(crate) async fn hydrate(
166        &self,
167        key: &str,
168        mut entry: CacheEntry,
169    ) -> Result<Option<CacheEntry>, CamelError> {
170        let Some(raw_path) = entry.payload_path.clone() else {
171            return Ok(Some(entry));
172        };
173        let Some(name) = sanitize_blob_name(&raw_path) else {
174            warn!(
175                key = key,
176                backend = self.inner.name(),
177                payload_path = %raw_path,
178                "corrupt cache row: payload_path must be a bare file name; treating as miss"
179            );
180            return Ok(None);
181        };
182        match self.store.read(name).await {
183            Ok(Some(bytes)) => {
184                entry.bytes = bytes;
185                entry.payload_path = None;
186                Ok(Some(entry))
187            }
188            Ok(None) => {
189                warn!(
190                    key = key,
191                    backend = self.inner.name(),
192                    blob = %name,
193                    "cache payload blob gone; treating as miss"
194                );
195                Ok(None)
196            }
197            Err(e) => Err(e),
198        }
199    }
200
201    /// Best-effort eager unlink of a key's predecessor payload after a
202    /// successful overwrite (ADR-0065, amendment "bd rc-uteoa").
203    ///
204    /// Row-guided, no store scan: only the name the pre-swap index row
205    /// carried is eligible, and only when it passes
206    /// [`sanitize_blob_name`], carries a parseable death epoch, and starts
207    /// with the current key's blake3-128 filename prefix — a corrupt row
208    /// naming another key's payload (or a foreign name) is never unlinked.
209    /// `keep_name` is the fresh payload's name on the successful-put path:
210    /// a same-second identical rewrite reuses the name, and only the
211    /// fresh payload owns it, so an equal name skips the reclaim. On the
212    /// inline-fallback path no fresh payload owns any name; callers pass
213    /// `None` to disable the equal-name guard. An absent payload counts
214    /// as reclaimed by someone else; any other unlink failure WARNs once
215    /// and leaves the payload to its store's reclaimer at the death
216    /// epoch. The function never returns `Err`: the reclaim adds no
217    /// failure mode to `set`.
218    async fn reclaim_predecessor(
219        &self,
220        key: &str,
221        old_name: Option<&str>,
222        keep_name: Option<&str>,
223    ) {
224        let Some(old_name) = old_name else {
225            return;
226        };
227        if keep_name == Some(old_name) {
228            return;
229        }
230        let key_prefix = format!("{}.", blake3_128hex(key.as_bytes()));
231        let eligible = sanitize_blob_name(old_name).is_some()
232            && parse_death_epoch(old_name).is_some()
233            && old_name.starts_with(&key_prefix);
234        if !eligible {
235            return;
236        }
237        if let Err(e) = self.store.unlink(old_name).await {
238            warn!(
239                key = key,
240                backend = self.inner.name(),
241                error = %e,
242                "eager reclaim of predecessor blob failed; sweeper reclaims it at its death epoch"
243            );
244        }
245    }
246}
247
248#[async_trait]
249impl CacheRepository for OffloadRepository {
250    fn name(&self) -> &str {
251        self.inner.name()
252    }
253
254    async fn get(&self, key: &str) -> Result<Option<CacheEntry>, CamelError> {
255        match self.inner.get(key).await? {
256            Some(entry) => self.hydrate(key, entry).await,
257            None => Ok(None),
258        }
259    }
260
261    async fn set(
262        &self,
263        key: &str,
264        mut entry: CacheEntry,
265        ttl: Option<Duration>,
266    ) -> Result<(), CamelError> {
267        let effective_ttl = ttl.unwrap_or(self.payload_max_ttl);
268        // Capture the predecessor's blob name before the index swap so a
269        // successful overwrite can reclaim it eagerly (ADR-0065 amendment,
270        // "bd rc-uteoa"). The SILENT maintenance read keeps the capture off
271        // every counted path — a phantom miss on first write or a phantom
272        // hit on overwrite would distort /ops/cache/stats and
273        // camel_cache_{hits,misses}_total. A failed read only skips the
274        // reclaim — the write proceeds unchanged in every case.
275        let old_name = match self.inner.peek_row_silent(key).await {
276            Ok(Some(row)) => row.payload_path,
277            Ok(None) => None,
278            Err(e) => {
279                warn!(
280                    key = key,
281                    backend = self.inner.name(),
282                    error = %e,
283                    "pre-swap row read failed; skipping eager reclaim"
284                );
285                None
286            }
287        };
288        // Death epoch = expiry + retention + sweep grace, saturating in
289        // Duration space (a pre-epoch clock clamps to the Unix epoch),
290        // truncated to whole seconds for the blob filename.
291        let death_epoch = (self.clock)()
292            .duration_since(UNIX_EPOCH)
293            .unwrap_or_default()
294            .saturating_add(effective_ttl)
295            .saturating_add(self.stale_retention)
296            .saturating_add(self.sweep_interval)
297            .as_secs();
298        // The same instant as a deadline, for stores that carry their own
299        // expiration next to the name-encoded epoch (e.g. a redis EXAT).
300        // The u64 filename seconds can exceed the platform's SystemTime
301        // range only for absurd clocks; the clamp stays a bounded
302        // deadline either way.
303        let deadline = UNIX_EPOCH
304            .checked_add(Duration::from_secs(death_epoch))
305            .unwrap_or(UNIX_EPOCH);
306        let dest_name = blob_filename(key, death_epoch, &entry);
307
308        match self.store.put(&dest_name, &entry.bytes, deadline).await {
309            Ok(()) => {
310                entry.bytes = Vec::new();
311                // Clone: `dest_name` is still needed for the equal-name
312                // guard after `entry` (carrying the same name) moves into
313                // the inner set.
314                entry.payload_path = Some(dest_name.clone());
315                // The ttl MUST be Some: every inner overwrites
316                // `expires_at` from the ttl argument, so None would wipe
317                // the fabricated expiry. The inner recomputes `expires_at`
318                // from its own clock; the sub-second skew is absorbed by
319                // the death-epoch grace.
320                let result = self.inner.set(key, entry, Some(effective_ttl)).await;
321                // Reclaim only after the inner accepted the swap: on an
322                // error the surviving row may still reference the
323                // predecessor payload.
324                if result.is_ok() {
325                    self.reclaim_predecessor(key, old_name.as_deref(), Some(&dest_name))
326                        .await;
327                }
328                result
329            }
330            Err(e) => {
331                warn!(
332                    key = key,
333                    backend = self.inner.name(),
334                    error = %e,
335                    "cache blob write failed; storing entry inline instead"
336                );
337                // Inline fallback with the original, unstripped entry: the
338                // decorator never converts a store-write failure into a
339                // cache-write error. The CAPPED ttl keeps the spec's
340                // no-TTL semantic (payload_max_ttl) even for degraded rows
341                // — an uncapped inline row would never be reclaimed. The
342                // new row no longer references the predecessor, so the
343                // reclaim runs with the equal-name guard disabled (the
344                // failed write left no fresh payload owning that name).
345                let result = self.inner.set(key, entry, Some(effective_ttl)).await;
346                if result.is_ok() {
347                    self.reclaim_predecessor(key, old_name.as_deref(), None)
348                        .await;
349                }
350                result
351            }
352        }
353    }
354
355    async fn peek_stale(&self, key: &str) -> Result<Option<CacheEntry>, CamelError> {
356        match self.inner.peek_stale(key).await? {
357            Some(entry) => self.hydrate(key, entry).await,
358            None => Ok(None),
359        }
360    }
361
362    /// Delegate-only: the index row is dropped here; the payload becomes
363    /// an orphan reclaimed asynchronously at its name-encoded death
364    /// epoch.
365    async fn invalidate(&self, key: &str) -> Result<(), CamelError> {
366        self.inner.invalidate(key).await
367    }
368
369    /// Reclaim payload space now: best-effort bulk reclamation in the
370    /// store, then delegate to the index. Store failures never turn
371    /// `clear` into `Err` — the store WARNs per payload and continues
372    /// ([`PayloadStore::clear`] contract).
373    async fn clear(&self) -> Result<(), CamelError> {
374        self.store.clear().await;
375        self.inner.clear().await
376    }
377
378    /// Delegate-only: the returned count is index-scoped; payloads are
379    /// reclaimed asynchronously at their name-encoded death epoch.
380    async fn invalidate_prefix(&self, prefix: &str) -> Result<u64, CamelError> {
381        self.inner.invalidate_prefix(prefix).await
382    }
383
384    async fn stats(&self) -> CacheStats {
385        self.inner.stats().await
386    }
387}
388
389impl std::fmt::Debug for OffloadRepository {
390    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
391        f.debug_struct("OffloadRepository")
392            .field("inner", &self.inner)
393            .field("stale_retention", &self.stale_retention)
394            .field("sweep_interval", &self.sweep_interval)
395            .field("payload_max_ttl", &self.payload_max_ttl)
396            .finish()
397    }
398}
399
400// ── Filename helpers ─────────────────────────────────────────────────────────
401
402/// One-byte discriminant of the closed [`ContentType`] enum, mixed into the
403/// content fingerprint for domain separation (identical bytes under
404/// different content types produce different fingerprints). Exhaustive
405/// match — the enum is closed by contract (ADR-0049 §Exceptions).
406fn content_type_discriminant(content_type: ContentType) -> u8 {
407    match content_type {
408        ContentType::Bytes => 0,
409        ContentType::Text => 1,
410        ContentType::Json => 2,
411        ContentType::Xml => 3,
412    }
413}
414
415/// Finalize a hasher to its first 128 bits as 32 lowercase hex chars.
416pub(crate) fn hasher_128hex(hasher: blake3::Hasher) -> String {
417    let hex = hasher.finalize().to_hex().to_string();
418    hex[..32].to_string()
419}
420
421/// blake3-128 hex of a single byte slice.
422pub(crate) fn blake3_128hex(data: &[u8]) -> String {
423    let mut hasher = blake3::Hasher::new();
424    hasher.update(data);
425    hasher_128hex(hasher)
426}
427
428/// 128-bit content fingerprint: `blake3(bytes || content_type discriminant)`.
429pub(crate) fn content_fingerprint(entry: &CacheEntry) -> String {
430    let mut hasher = blake3::Hasher::new();
431    hasher.update(&entry.bytes);
432    hasher.update(&[content_type_discriminant(entry.content_type)]);
433    hasher_128hex(hasher)
434}
435
436/// Blob name: `{key-hash}.{death_epoch}.{fingerprint}.blob`.
437fn blob_filename(key: &str, death_epoch: u64, entry: &CacheEntry) -> String {
438    format!(
439        "{}.{}.{}.blob",
440        blake3_128hex(key.as_bytes()),
441        death_epoch,
442        content_fingerprint(entry)
443    )
444}
445
446/// Death epoch (second dot-separated component) of a blob name, if it
447/// parses as `u64`.
448pub(crate) fn parse_death_epoch(file_name: &str) -> Option<u64> {
449    file_name.split('.').nth(1)?.parse().ok()
450}
451
452/// Accept only a bare file name: non-empty, no `/`, no `\`, no `..`.
453///
454/// Absolute paths necessarily contain a separator on both Unix and Windows,
455/// so the separator checks subsume the absolute-path rejection. Everything
456/// else is treated as a corrupt row.
457pub(crate) fn sanitize_blob_name(path: &str) -> Option<&str> {
458    if path.is_empty() || path.contains('/') || path.contains('\\') || path.contains("..") {
459        return None;
460    }
461    Some(path)
462}
463
464#[cfg(test)]
465mod tests {
466    use super::*;
467    use crate::cache::MemoryCacheRepository;
468
469    fn entry(bytes: Vec<u8>, content_type: ContentType) -> CacheEntry {
470        CacheEntry {
471            bytes,
472            payload_path: None,
473            content_type,
474            expires_at: None,
475        }
476    }
477
478    /// Store whose `put` always fails: models a payload store rejecting
479    /// the write (e.g. a redis EXAT overflow). The decorator must fall
480    /// back to inline storage, never surface the failure.
481    struct FailingPutStore;
482
483    #[async_trait]
484    impl PayloadStore for FailingPutStore {
485        async fn put(
486            &self,
487            _name: &str,
488            _bytes: &[u8],
489            _death_epoch: SystemTime,
490        ) -> Result<(), CamelError> {
491            Err(CamelError::Io(
492                "failing-put-store: injected put failure".into(),
493            ))
494        }
495
496        async fn read(&self, _name: &str) -> Result<Option<Vec<u8>>, CamelError> {
497            Ok(None)
498        }
499
500        async fn unlink(&self, _name: &str) -> Result<(), CamelError> {
501            Ok(())
502        }
503    }
504
505    #[tokio::test]
506    async fn offload_repository_inline_fallback_on_store_put_failure() {
507        let inner = Arc::new(MemoryCacheRepository::new("test", 100));
508        let repo = OffloadRepository::new(
509            inner,
510            Arc::new(FailingPutStore),
511            Duration::from_secs(168 * 3600),
512            Duration::from_secs(3600),
513            Duration::from_secs(24 * 3600),
514        );
515
516        let payload = vec![7; 128];
517        repo.set(
518            "k",
519            entry(payload.clone(), ContentType::Bytes),
520            Some(Duration::from_secs(60)),
521        )
522        .await
523        .expect("set must fall back inline, never Err");
524
525        let got = repo.get("k").await.expect("get").expect("present");
526        assert_eq!(
527            got.bytes, payload,
528            "bytes must be present, stored inline in the index"
529        );
530        assert_eq!(got.content_type, ContentType::Bytes);
531        assert!(got.payload_path.is_none(), "fallback row stays inline");
532    }
533
534    /// Store that accepts every put but never serves a read: models a
535    /// payload vanished from the store (reclaimed, evicted, lost).
536    struct ForgetfulStore;
537
538    #[async_trait]
539    impl PayloadStore for ForgetfulStore {
540        async fn put(
541            &self,
542            _name: &str,
543            _bytes: &[u8],
544            _death_epoch: SystemTime,
545        ) -> Result<(), CamelError> {
546            Ok(())
547        }
548
549        async fn read(&self, _name: &str) -> Result<Option<Vec<u8>>, CamelError> {
550            Ok(None)
551        }
552
553        async fn unlink(&self, _name: &str) -> Result<(), CamelError> {
554            Ok(())
555        }
556    }
557
558    #[tokio::test]
559    async fn offload_repository_missing_payload_is_miss() {
560        let inner = Arc::new(MemoryCacheRepository::new("test", 100));
561        let repo = OffloadRepository::new(
562            inner,
563            Arc::new(ForgetfulStore),
564            Duration::from_secs(168 * 3600),
565            Duration::from_secs(3600),
566            Duration::from_secs(24 * 3600),
567        );
568
569        repo.set(
570            "k",
571            entry(vec![1, 2, 3], ContentType::Bytes),
572            Some(Duration::from_secs(60)),
573        )
574        .await
575        .expect("set");
576
577        assert_eq!(
578            repo.get("k").await.expect("get"),
579            None,
580            "a vanished payload degrades to a MISS, never an error"
581        );
582    }
583}