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}