Skip to main content

camel_core/cache/
redb.rs

1//! Redb-backed persistent cache repository.
2//!
3//! Mirrors `idempotent/redb_repository.rs`: all redb I/O is offloaded to
4//! `tokio::task::spawn_blocking` because `redb::Database` is blocking. Values
5//! are `serde_json`-serialized [`CacheEntry`] blobs.
6//!
7//! # Sweep
8//!
9//! A background task wakes every `sweep_interval` and reclaims entries whose
10//! `expires_at + stale_retention < now`. It is bound to the context's
11//! [`CancellationToken`]: when the context shuts down, the sweep exits. The
12//! token is context-owned — dropping this repository aborts only the sweep
13//! task handle, never the token (cancelling it would tear down the whole
14//! context).
15
16use std::fmt;
17use std::ops::Bound;
18use std::path::Path;
19use std::path::PathBuf;
20use std::sync::Arc;
21use std::sync::atomic::AtomicU64;
22use std::sync::atomic::Ordering;
23use std::time::Duration;
24use std::time::SystemTime;
25
26use async_trait::async_trait;
27use camel_api::CamelError;
28use camel_api::cache::CacheEntry;
29use camel_api::cache::CacheRepository;
30use camel_api::cache::CacheStats;
31use parking_lot::Mutex;
32use redb::ReadableDatabase;
33use redb::ReadableTable;
34use redb::ReadableTableMetadata;
35use redb::TableDefinition;
36use tokio::task::JoinHandle;
37use tokio_util::sync::CancellationToken;
38
39// ── Table definition ──────────────────────────────────────────────────────────
40
41/// `key → serde_json(CacheEntry)`. Mirrors the table-definition style of
42/// `idempotent/redb_repository.rs`.
43const CACHE_TABLE: TableDefinition<&str, &[u8]> = TableDefinition::new("cache_entries");
44
45// ── Repository ────────────────────────────────────────────────────────────────
46
47/// Redb-backed implementation of [`CacheRepository`].
48///
49/// All counters are `Arc<AtomicU64>` so the spawned sweep task can update
50/// them via cloned references. `cache_size`, `sweep_interval`, and
51/// `stale_retention` are retained so the propagation seam required by the
52/// eip-cache spec can expose them via accessors.
53pub struct RedbCacheRepository {
54    name: String,
55    db: Arc<redb::Database>,
56    stale_retention: Duration,
57    max_entries: Option<usize>,
58    /// redb page-cache size in bytes, passed to `redb::Builder::set_cache_size`.
59    cache_size: usize,
60    /// Recorded background-sweep interval, consumed by the spawned sweep task.
61    sweep_interval: Duration,
62    hits: Arc<AtomicU64>,
63    misses: Arc<AtomicU64>,
64    evictions: Arc<AtomicU64>,
65    peek_stale_served: Arc<AtomicU64>,
66    invalidations: Arc<AtomicU64>,
67    /// Best-effort approximation of `table.len()` for stats display.
68    /// The authoritative count is `table.len()` inside write transactions,
69    /// used for capacity enforcement. Sweep and invalidate decrement via
70    /// saturating-fetch-sub to prevent underflow.
71    entries: Arc<AtomicU64>,
72    /// Context-owned shutdown token. Cloned into the sweep task; never
73    /// cancelled by this repository (doing so would shut down the entire
74    /// context when a single repo is dropped).
75    shutdown_token: CancellationToken,
76    sweep_handle: Mutex<Option<JoinHandle<()>>>,
77}
78
79impl RedbCacheRepository {
80    /// Open (or create) the redb database at `path`, seed the `entries`
81    /// counter from `table.len()`, and spawn the background sweep task.
82    ///
83    /// `shutdown_token` is the **context's** token — binding the sweep to
84    /// context shutdown. The whole open sequence runs in
85    /// `spawn_blocking` because `redb::Database::create` is blocking.
86    #[allow(clippy::too_many_arguments)]
87    pub async fn new(
88        name: impl Into<String>,
89        path: impl Into<PathBuf>,
90        stale_retention: Duration,
91        max_entries: Option<usize>,
92        cache_size: usize,
93        sweep_interval: Duration,
94        shutdown_token: CancellationToken,
95    ) -> Result<Self, CamelError> {
96        let name = name.into();
97        let path: PathBuf = path.into();
98        let path_for_db = path.clone();
99        let (db, initial_len) = tokio::task::spawn_blocking(move || {
100            if let Some(parent) = path_for_db.parent() {
101                std::fs::create_dir_all(parent)
102                    .map_err(|e| CamelError::Io(format!("redb create_dir_all: {e}")))?;
103            }
104            let db = redb::Builder::new()
105                .set_cache_size(cache_size)
106                .create(&path_for_db)
107                .map_err(|e| CamelError::Io(format!("redb open: {e}")))?;
108            // Create the table on first open AND read the persisted entry
109            // count in a single write txn so the counter survives reopen.
110            let len = {
111                let wtx = db
112                    .begin_write()
113                    .map_err(|e| CamelError::Io(format!("redb begin_write: {e}")))?;
114                // `table` borrows `wtx` mutably, so it must drop before
115                // `wtx.commit()` (which moves `wtx`). Read `len` then drop.
116                let len = {
117                    let table = wtx
118                        .open_table(CACHE_TABLE)
119                        .map_err(|e| CamelError::Io(format!("redb open_table: {e}")))?;
120                    table
121                        .len()
122                        .map_err(|e| CamelError::Io(format!("redb len: {e}")))?
123                };
124                wtx.commit()
125                    .map_err(|e| CamelError::Io(format!("redb commit: {e}")))?;
126                len
127            };
128            Ok::<_, CamelError>((Arc::new(db), len))
129        })
130        .await
131        .map_err(|e| CamelError::Io(format!("spawn_blocking join: {e}")))??;
132
133        let hits = Arc::new(AtomicU64::new(0));
134        let misses = Arc::new(AtomicU64::new(0));
135        let evictions = Arc::new(AtomicU64::new(0));
136        let peek_stale_served = Arc::new(AtomicU64::new(0));
137        let invalidations = Arc::new(AtomicU64::new(0));
138        let entries = Arc::new(AtomicU64::new(initial_len));
139
140        // Container memory guardrail — diagnostic only, never fails
141        // construction. Warns once when the redb page cache cannot fit within
142        // the container's cgroup memory limit.
143        emit_memory_guardrail(
144            cache_size,
145            Path::new("/sys/fs/cgroup/memory.max"),
146            Path::new("/sys/fs/cgroup/memory/memory.limit_in_bytes"),
147        );
148
149        // Spawn the sweep loop. All shared state is captured as cloned Arcs;
150        // the context token clone drives termination.
151        let db_clone = Arc::clone(&db);
152        let evictions_clone = Arc::clone(&evictions);
153        let entries_clone = Arc::clone(&entries);
154        let token_clone = shutdown_token.clone();
155        let retention = stale_retention;
156        let handle = tokio::spawn(async move {
157            let mut ticker = tokio::time::interval(sweep_interval);
158            loop {
159                tokio::select! {
160                    _ = ticker.tick() => {
161                        let db = Arc::clone(&db_clone);
162                        let reclaimed = tokio::task::spawn_blocking(move || {
163                            sweep_reclaim(&db, retention).unwrap_or(0)
164                        })
165                        .await
166                        .unwrap_or(0);
167                        evictions_clone.fetch_add(reclaimed, Ordering::Relaxed);
168                        let current = entries_clone.load(Ordering::Relaxed);
169                        let sub = std::cmp::min(current, reclaimed);
170                        entries_clone.fetch_sub(sub, Ordering::Relaxed);
171                    }
172                    _ = token_clone.cancelled() => break,
173                }
174            }
175        });
176
177        Ok(Self {
178            name,
179            db,
180            stale_retention,
181            max_entries,
182            cache_size,
183            sweep_interval,
184            hits,
185            misses,
186            evictions,
187            peek_stale_served,
188            invalidations,
189            entries,
190            shutdown_token,
191            sweep_handle: Mutex::new(Some(handle)),
192        })
193    }
194
195    /// Recorded redb page-cache size in bytes (propagation seam for the
196    /// eip-cache spec).
197    pub fn cache_size(&self) -> usize {
198        self.cache_size
199    }
200
201    /// Recorded background-sweep interval.
202    pub fn sweep_interval(&self) -> std::time::Duration {
203        self.sweep_interval
204    }
205
206    /// Recorded stale-retention window.
207    pub fn stale_retention(&self) -> std::time::Duration {
208        self.stale_retention
209    }
210
211    /// Run a single reclamation pass and return the number of entries
212    /// reclaimed. Tests call this directly for deterministic sweep coverage
213    /// without waiting on the background ticker.
214    #[cfg(test)]
215    pub(crate) async fn sweep_once(&self) -> Result<u64, CamelError> {
216        let db = Arc::clone(&self.db);
217        let retention = self.stale_retention;
218        let reclaimed = tokio::task::spawn_blocking(move || sweep_reclaim(&db, retention))
219            .await
220            .map_err(|e| CamelError::Io(format!("spawn_blocking join: {e}")))??;
221        self.evictions.fetch_add(reclaimed, Ordering::Relaxed);
222        let current = self.entries.load(Ordering::Relaxed);
223        let sub = std::cmp::min(current, reclaimed);
224        self.entries.fetch_sub(sub, Ordering::Relaxed);
225        Ok(reclaimed)
226    }
227}
228
229// ── Memory guardrail ──────────────────────────────────────────────────────────
230
231/// Read the container memory limit (bytes) from the cgroup filesystem,
232/// preferring cgroup v2 (`memory.max`) with cgroup v1
233/// (`memory.limit_in_bytes`) as fallback.
234///
235/// v2 `"max"` (unlimited) or unparseable content falls through to v1; a v1
236/// value above 16 TiB is the v1 "unlimited" sentinel and is reported as no
237/// limit. Missing/unreadable files at either path fall through to `None`.
238/// All reads are best-effort (`std::fs::read_to_string(...).ok()`) — this is a
239/// diagnostic seam, never a failure path.
240pub(crate) fn memory_limit_from_paths(v2: &Path, v1: &Path) -> Option<u64> {
241    if let Ok(content) = std::fs::read_to_string(v2)
242        && let Ok(bytes) = content.trim().parse::<u64>()
243    {
244        return Some(bytes);
245    }
246    if let Ok(content) = std::fs::read_to_string(v1)
247        && let Ok(bytes) = content.trim().parse::<u64>()
248        // cgroup v1 "unlimited" sentinel: anything above 16 TiB.
249        && bytes <= 17_592_186_044_416
250    {
251        return Some(bytes);
252    }
253    None
254}
255
256/// Diagnostic-only guardrail: when the configured redb cache size exceeds the
257/// container's cgroup memory limit, emit a single warning naming both values.
258/// Never fails — a missing limit or unreadable files simply skip the warning.
259pub(crate) fn emit_memory_guardrail(cache_size: usize, v2: &Path, v1: &Path) {
260    let cache_size = cache_size as u64;
261    if let Some(limit) = memory_limit_from_paths(v2, v1)
262        && cache_size > limit
263    {
264        tracing::warn!(
265            "redb cache_size ({cache_size} bytes) exceeds container memory limit ({limit} bytes)"
266        );
267    }
268}
269
270/// Reclaim entries whose `expires_at + stale_retention < now`.
271///
272/// Entries with `expires_at = None` (no expiry) are never reclaimed. Returns
273/// the reclaimed count. Used by both the background sweep loop (errors mapped
274/// to `0`) and [`RedbCacheRepository::sweep_once`] (errors propagated).
275fn sweep_reclaim(db: &redb::Database, stale_retention: Duration) -> Result<u64, CamelError> {
276    let txn = db
277        .begin_write()
278        .map_err(|e| CamelError::Io(format!("redb begin_write: {e}")))?;
279    let reclaimed = {
280        let mut table = txn
281            .open_table(CACHE_TABLE)
282            .map_err(|e| CamelError::Io(format!("redb open_table: {e}")))?;
283        let now = SystemTime::now();
284        // Collect keys first — `table.iter()` borrows the table immutably and
285        // would conflict with `remove` (mirrors `idempotent/redb_repository`
286        // clear pattern).
287        let mut to_delete: Vec<String> = Vec::new();
288        for row in table
289            .iter()
290            .map_err(|e| CamelError::Io(format!("redb iter: {e}")))?
291        {
292            let (k, v) = row.map_err(|e| CamelError::Io(format!("redb iter item: {e}")))?;
293            let entry: CacheEntry = serde_json::from_slice(v.value())
294                .map_err(|e| CamelError::Io(format!("cache deserialization: {e}")))?;
295            let should_delete = match entry.expires_at {
296                Some(exp) => match exp.checked_add(stale_retention) {
297                    // Threshold in the past → past stale-retention window.
298                    Some(threshold) => threshold < now,
299                    // SystemTime addition overflowed → treat as reclaimable.
300                    None => true,
301                },
302                None => false,
303            };
304            if should_delete {
305                to_delete.push(k.value().to_string());
306            }
307        }
308        for k in &to_delete {
309            let _ = table
310                .remove(k.as_str())
311                .map_err(|e| CamelError::Io(format!("redb remove: {e}")))?;
312        }
313        to_delete.len() as u64
314    };
315    txn.commit()
316        .map_err(|e| CamelError::Io(format!("redb commit: {e}")))?;
317    Ok(reclaimed)
318}
319
320/// Compute the smallest string that sorts after every string beginning with
321/// `prefix`, as the exclusive upper [`Bound`] of a range scan.
322///
323/// The successor is `prefix` with its last Unicode scalar value incremented by
324/// one, skipping the UTF-16 surrogate gap (U+D7FF → U+E000). A trailing
325/// U+10FFFF carries into the preceding scalar; an empty or all-U+10FFFF
326/// prefix has no successor, so the bound is [`Bound::Unbounded`].
327fn successor_bound(prefix: &str) -> Bound<String> {
328    match prefix.chars().last() {
329        // Empty string has no scalar to increment — match every key.
330        None => Bound::Unbounded,
331        Some(last) => {
332            let rest = &prefix[..prefix.len() - last.len_utf8()];
333            match increment_scalar(last) {
334                Some(next) => {
335                    let mut s = String::with_capacity(rest.len() + next.len_utf8());
336                    s.push_str(rest);
337                    s.push(next);
338                    Bound::Excluded(s)
339                }
340                // U+10FFFF has no successor scalar — carry into the rest.
341                None => successor_bound(rest),
342            }
343        }
344    }
345}
346
347/// Increment a Unicode scalar value by one, skipping the surrogate range.
348///
349/// Returns `None` for U+10FFFF (the maximum scalar has no successor).
350fn increment_scalar(c: char) -> Option<char> {
351    match c {
352        // U+D7FF jumps over the surrogate range to U+E000.
353        '\u{D7FF}' => Some('\u{E000}'),
354        // U+10FFFF is the maximum scalar value.
355        '\u{10FFFF}' => None,
356        // Surrogates are not valid `char`, so every remaining scalar +1 is valid.
357        _ => char::from_u32(c as u32 + 1),
358    }
359}
360
361// ── CacheRepository impl ──────────────────────────────────────────────────────
362
363#[async_trait]
364impl CacheRepository for RedbCacheRepository {
365    fn name(&self) -> &str {
366        &self.name
367    }
368
369    async fn get(&self, key: &str) -> Result<Option<CacheEntry>, CamelError> {
370        let db = Arc::clone(&self.db);
371        let key = key.to_string();
372        let result =
373            tokio::task::spawn_blocking(move || -> Result<Option<CacheEntry>, CamelError> {
374                let rtx = db
375                    .begin_read()
376                    .map_err(|e| CamelError::Io(format!("redb begin_read: {e}")))?;
377                let table = rtx
378                    .open_table(CACHE_TABLE)
379                    .map_err(|e| CamelError::Io(format!("redb open_table: {e}")))?;
380                match table
381                    .get(key.as_str())
382                    .map_err(|e| CamelError::Io(format!("redb get: {e}")))?
383                {
384                    Some(guard) => {
385                        let entry: CacheEntry = serde_json::from_slice(guard.value())
386                            .map_err(|e| CamelError::Io(format!("cache deserialization: {e}")))?;
387                        Ok(Some(entry))
388                    }
389                    None => Ok(None),
390                }
391            })
392            .await
393            .map_err(|e| CamelError::Io(format!("spawn_blocking join: {e}")))??;
394        // Expiry check happens outside the blocking closure, mirroring
395        // `MemoryCacheRepository::get` semantics.
396        match result {
397            Some(entry) => {
398                let expired = entry
399                    .expires_at
400                    .map(|e| e <= SystemTime::now())
401                    .unwrap_or(false);
402                if expired {
403                    self.misses.fetch_add(1, Ordering::Relaxed);
404                    Ok(None)
405                } else {
406                    self.hits.fetch_add(1, Ordering::Relaxed);
407                    Ok(Some(entry))
408                }
409            }
410            None => {
411                self.misses.fetch_add(1, Ordering::Relaxed);
412                Ok(None)
413            }
414        }
415    }
416
417    async fn set(
418        &self,
419        key: &str,
420        mut value: CacheEntry,
421        ttl: Option<Duration>,
422    ) -> Result<(), CamelError> {
423        value.expires_at = ttl.map(|d| SystemTime::now() + d);
424        let serialized = serde_json::to_vec(&value)
425            .map_err(|e| CamelError::Io(format!("cache serialization: {e}")))?;
426        let db = Arc::clone(&self.db);
427        let key = key.to_string();
428        let max_entries = self.max_entries;
429        let was_new = tokio::task::spawn_blocking(move || {
430            let txn = db
431                .begin_write()
432                .map_err(|e| CamelError::Io(format!("redb begin_write: {e}")))?;
433            let was_new = {
434                let mut table = txn
435                    .open_table(CACHE_TABLE)
436                    .map_err(|e| CamelError::Io(format!("redb open_table: {e}")))?;
437                // Scope `prior` so its immutable borrow of `table` ends before
438                // the mutable `insert` below (redb guards borrow the table).
439                let is_new = {
440                    let prior = table
441                        .get(key.as_str())
442                        .map_err(|e| CamelError::Io(format!("redb get: {e}")))?;
443                    // Capacity check only applies to genuinely new keys —
444                    // overwrites don't grow the table.
445                    if prior.is_none()
446                        && let Some(max) = max_entries
447                    {
448                        let count = table
449                            .len()
450                            .map_err(|e| CamelError::Io(format!("redb len: {e}")))?
451                            as usize;
452                        if count >= max {
453                            return Err(CamelError::Config(format!(
454                                "cache: max_entries ({max}) exceeded"
455                            )));
456                        }
457                    }
458                    prior.is_none()
459                };
460                table
461                    .insert(key.as_str(), serialized.as_slice())
462                    .map_err(|e| CamelError::Io(format!("redb insert: {e}")))?;
463                is_new
464            };
465            txn.commit()
466                .map_err(|e| CamelError::Io(format!("redb commit: {e}")))?;
467            Ok::<bool, CamelError>(was_new)
468        })
469        .await
470        .map_err(|e| CamelError::Io(format!("spawn_blocking join: {e}")))??;
471        if was_new {
472            self.entries.fetch_add(1, Ordering::Relaxed);
473        }
474        Ok(())
475    }
476
477    async fn peek_row_silent(&self, key: &str) -> Result<Option<CacheEntry>, CamelError> {
478        // Same raw read as peek_stale: redb counts nothing here, and the
479        // trait contract forbids the counted get path.
480        self.peek_stale(key).await
481    }
482
483    async fn peek_stale(&self, key: &str) -> Result<Option<CacheEntry>, CamelError> {
484        let db = Arc::clone(&self.db);
485        let key = key.to_string();
486        let result =
487            tokio::task::spawn_blocking(move || -> Result<Option<CacheEntry>, CamelError> {
488                let rtx = db
489                    .begin_read()
490                    .map_err(|e| CamelError::Io(format!("redb begin_read: {e}")))?;
491                let table = rtx
492                    .open_table(CACHE_TABLE)
493                    .map_err(|e| CamelError::Io(format!("redb open_table: {e}")))?;
494                match table
495                    .get(key.as_str())
496                    .map_err(|e| CamelError::Io(format!("redb get: {e}")))?
497                {
498                    Some(guard) => {
499                        let entry: CacheEntry = serde_json::from_slice(guard.value())
500                            .map_err(|e| CamelError::Io(format!("cache deserialization: {e}")))?;
501                        Ok(Some(entry))
502                    }
503                    None => Ok(None),
504                }
505            })
506            .await
507            .map_err(|e| CamelError::Io(format!("spawn_blocking join: {e}")))??;
508        if result.is_some() {
509            self.peek_stale_served.fetch_add(1, Ordering::Relaxed);
510        }
511        Ok(result)
512    }
513
514    async fn invalidate(&self, key: &str) -> Result<(), CamelError> {
515        let db = Arc::clone(&self.db);
516        let key = key.to_string();
517        let was_present = tokio::task::spawn_blocking(move || {
518            let txn = db
519                .begin_write()
520                .map_err(|e| CamelError::Io(format!("redb begin_write: {e}")))?;
521            let was_present = {
522                let mut table = txn
523                    .open_table(CACHE_TABLE)
524                    .map_err(|e| CamelError::Io(format!("redb open_table: {e}")))?;
525                // `remove` is idempotent per the trait contract; the returned
526                // Option tells us whether a value was actually present. Read
527                // `.is_some()` now so the guard (borrowing `table`) drops
528                // before `commit` moves `txn`.
529                table
530                    .remove(key.as_str())
531                    .map_err(|e| CamelError::Io(format!("redb remove: {e}")))?
532                    .is_some()
533            };
534            txn.commit()
535                .map_err(|e| CamelError::Io(format!("redb commit: {e}")))?;
536            Ok::<bool, CamelError>(was_present)
537        })
538        .await
539        .map_err(|e| CamelError::Io(format!("spawn_blocking join: {e}")))??;
540        if was_present {
541            let current = self.entries.load(Ordering::Relaxed);
542            let sub = std::cmp::min(current, 1);
543            self.entries.fetch_sub(sub, Ordering::Relaxed);
544        }
545        self.invalidations.fetch_add(1, Ordering::Relaxed);
546        Ok(())
547    }
548
549    async fn invalidate_prefix(&self, prefix: &str) -> Result<u64, CamelError> {
550        let db = Arc::clone(&self.db);
551        let prefix = prefix.to_string();
552        let deleted = tokio::task::spawn_blocking(move || -> Result<u64, CamelError> {
553            // Collect matching keys in a read txn, then delete in one write
554            // txn (mirrors the collect-then-remove pattern of `clear`).
555            let keys: Vec<String> = {
556                let rtx = db
557                    .begin_read()
558                    .map_err(|e| CamelError::Io(format!("redb begin_read: {e}")))?;
559                let table = rtx
560                    .open_table(CACHE_TABLE)
561                    .map_err(|e| CamelError::Io(format!("redb open_table: {e}")))?;
562                // `successor_bound` yields an owned bound; redb's range needs
563                // `&str` bounds for `&str` keys (`String` is not `Borrow<&str>`).
564                let upper: Bound<String> = successor_bound(&prefix);
565                let upper_ref: Bound<&str> = match &upper {
566                    Bound::Included(s) => Bound::Included(s.as_str()),
567                    Bound::Excluded(s) => Bound::Excluded(s.as_str()),
568                    Bound::Unbounded => Bound::Unbounded,
569                };
570                let mut keys = Vec::new();
571                for row in table
572                    .range::<&str>((Bound::Included(prefix.as_str()), upper_ref))
573                    .map_err(|e| CamelError::Io(format!("redb range: {e}")))?
574                {
575                    let (k, _v) =
576                        row.map_err(|e| CamelError::Io(format!("redb range item: {e}")))?;
577                    keys.push(k.value().to_string());
578                }
579                keys
580            };
581            let wtx = db
582                .begin_write()
583                .map_err(|e| CamelError::Io(format!("redb begin_write: {e}")))?;
584            {
585                let mut table = wtx
586                    .open_table(CACHE_TABLE)
587                    .map_err(|e| CamelError::Io(format!("redb open_table: {e}")))?;
588                for k in &keys {
589                    let _ = table
590                        .remove(k.as_str())
591                        .map_err(|e| CamelError::Io(format!("redb remove: {e}")))?;
592                }
593            }
594            wtx.commit()
595                .map_err(|e| CamelError::Io(format!("redb commit: {e}")))?;
596            Ok(keys.len() as u64)
597        })
598        .await
599        .map_err(|e| CamelError::Io(format!("spawn_blocking join: {e}")))??;
600        self.invalidations.fetch_add(1, Ordering::Relaxed);
601        if deleted > 0 {
602            let current = self.entries.load(Ordering::Relaxed);
603            let sub = std::cmp::min(current, deleted);
604            self.entries.fetch_sub(sub, Ordering::Relaxed);
605        }
606        Ok(deleted)
607    }
608
609    async fn clear(&self) -> Result<(), CamelError> {
610        let db = Arc::clone(&self.db);
611        tokio::task::spawn_blocking(move || {
612            let txn = db
613                .begin_write()
614                .map_err(|e| CamelError::Io(format!("redb begin_write: {e}")))?;
615            {
616                let mut table = txn
617                    .open_table(CACHE_TABLE)
618                    .map_err(|e| CamelError::Io(format!("redb open_table: {e}")))?;
619                // Collect-then-remove: `iter()` borrows immutably and cannot
620                // coexist with `remove`. Mirrors `idempotent/redb_repository`
621                // clear pattern.
622                let keys: Vec<String> = table
623                    .iter()
624                    .map_err(|e| CamelError::Io(format!("redb iter: {e}")))?
625                    .map(|r| {
626                        r.map(|(k, _v)| k.value().to_string())
627                            .map_err(|e| CamelError::Io(format!("redb iter item: {e}")))
628                    })
629                    .collect::<Result<_, _>>()?;
630                for k in &keys {
631                    let _ = table
632                        .remove(k.as_str())
633                        .map_err(|e| CamelError::Io(format!("redb remove: {e}")))?;
634                }
635            }
636            txn.commit()
637                .map_err(|e| CamelError::Io(format!("redb commit: {e}")))?;
638            Ok::<_, CamelError>(())
639        })
640        .await
641        .map_err(|e| CamelError::Io(format!("spawn_blocking join: {e}")))??;
642        self.entries.store(0, Ordering::Relaxed);
643        Ok(())
644    }
645
646    async fn stats(&self) -> CacheStats {
647        let db = Arc::clone(&self.db);
648        let bytes = tokio::task::spawn_blocking(move || total_bytes(&db))
649            .await
650            .unwrap_or_default();
651        CacheStats {
652            hits: self.hits.load(Ordering::Relaxed),
653            misses: self.misses.load(Ordering::Relaxed),
654            evictions: self.evictions.load(Ordering::Relaxed),
655            entries: self.entries.load(Ordering::Relaxed),
656            peek_stale_served: self.peek_stale_served.load(Ordering::Relaxed),
657            invalidations: self.invalidations.load(Ordering::Relaxed),
658            bytes,
659        }
660    }
661}
662
663/// Sum every entry's `bytes.len()` over the full table range.
664///
665/// Called only inside `spawn_blocking` from `stats()`; `None` when the table
666/// cannot be read or any entry fails to deserialize — the `bytes` field is a
667/// best-effort report, never an error.
668fn total_bytes(db: &redb::Database) -> Option<u64> {
669    let rtx = db.begin_read().ok()?;
670    let table = rtx.open_table(CACHE_TABLE).ok()?;
671    let mut total: u64 = 0;
672    for row in table.iter().ok()? {
673        let (_key, value) = row.ok()?;
674        let entry: CacheEntry = serde_json::from_slice(value.value()).ok()?;
675        total = total.saturating_add(entry.bytes.len() as u64);
676    }
677    Some(total)
678}
679
680impl fmt::Debug for RedbCacheRepository {
681    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
682        f.debug_struct("RedbCacheRepository")
683            .field("name", &self.name)
684            .field("stale_retention", &self.stale_retention)
685            .field("max_entries", &self.max_entries)
686            .field("cache_size", &self.cache_size)
687            .field("sweep_interval", &self.sweep_interval)
688            .field("shutdown_cancelled", &self.shutdown_token.is_cancelled())
689            .finish()
690    }
691}
692
693impl Drop for RedbCacheRepository {
694    fn drop(&mut self) {
695        // Abort ONLY the sweep task. Never cancel the context-owned token —
696        // that would shut down the entire context when one repo drops.
697        if let Some(handle) = self.sweep_handle.lock().take() {
698            handle.abort();
699        }
700    }
701}
702
703#[cfg(test)]
704mod tests {
705    use super::*;
706    use tempfile::TempDir;
707    use tempfile::tempdir;
708
709    fn entry() -> CacheEntry {
710        CacheEntry {
711            bytes: vec![1, 2, 3],
712            payload_path: None,
713            content_type: camel_api::cache::ContentType::Bytes,
714            expires_at: None,
715        }
716    }
717
718    /// Open a repo at `<tmp>/cache.redb` with a 60s stale-retention, no cap,
719    /// a 256 MiB cache size, and a 1h sweep interval (so the background loop
720    /// stays dormant during sequential tests).
721    async fn new_repo(tmp: &TempDir, shutdown_token: CancellationToken) -> RedbCacheRepository {
722        new_repo_with(
723            tmp,
724            shutdown_token,
725            Duration::from_secs(60),
726            None,
727            256 * 1024 * 1024,
728            Duration::from_secs(3600),
729        )
730        .await
731    }
732
733    /// Full-parameter variant of [`new_repo`] for tests that need custom
734    /// stale-retention, cap, cache size, or sweep interval values.
735    async fn new_repo_with(
736        tmp: &TempDir,
737        shutdown_token: CancellationToken,
738        stale_retention: Duration,
739        max_entries: Option<usize>,
740        cache_size: usize,
741        sweep_interval: Duration,
742    ) -> RedbCacheRepository {
743        let path = tmp.path().join("cache.redb");
744        RedbCacheRepository::new(
745            "redb",
746            path,
747            stale_retention,
748            max_entries,
749            cache_size,
750            sweep_interval,
751            shutdown_token,
752        )
753        .await
754        .expect("open redb cache repo")
755    }
756
757    #[tokio::test]
758    async fn cache_size_recorded_and_accessible() {
759        let dir = tempdir().expect("tempdir");
760        let token = CancellationToken::new();
761        let repo = new_repo_with(
762            &dir,
763            token,
764            Duration::from_secs(60),
765            None,
766            536_870_912,
767            Duration::from_secs(3600),
768        )
769        .await;
770        assert_eq!(repo.cache_size(), 536_870_912);
771    }
772
773    #[tokio::test]
774    async fn sweep_interval_recorded_and_accessible() {
775        let dir = tempdir().expect("tempdir");
776        let token = CancellationToken::new();
777        let repo = new_repo_with(
778            &dir,
779            token,
780            Duration::from_secs(60),
781            None,
782            256 * 1024 * 1024,
783            Duration::from_secs(1800),
784        )
785        .await;
786        assert_eq!(repo.sweep_interval(), Duration::from_secs(1800));
787    }
788
789    #[tokio::test]
790    async fn stale_retention_recorded_and_accessible() {
791        let dir = tempdir().expect("tempdir");
792        let token = CancellationToken::new();
793        let repo = new_repo_with(
794            &dir,
795            token,
796            Duration::from_secs(3600),
797            None,
798            256 * 1024 * 1024,
799            Duration::from_secs(3600),
800        )
801        .await;
802        assert_eq!(repo.stale_retention(), Duration::from_secs(3600));
803    }
804
805    #[tokio::test]
806    async fn explicit_cache_size_round_trip() {
807        let dir = tempdir().expect("tempdir");
808        let token = CancellationToken::new();
809        let repo = new_repo_with(
810            &dir,
811            token,
812            Duration::from_secs(60),
813            None,
814            512 * 1024 * 1024,
815            Duration::from_secs(3600),
816        )
817        .await;
818        repo.set("k", entry(), Some(Duration::from_secs(3600)))
819            .await
820            .expect("set");
821        let found = repo.get("k").await.expect("get");
822        assert!(
823            found.is_some(),
824            "entry must round-trip through the builder-opened database"
825        );
826    }
827
828    #[tokio::test]
829    async fn entries_survive_handle_drop_and_reopen() {
830        let dir = tempdir().expect("tempdir");
831        let path = dir.path().join("cache.redb");
832        let token = CancellationToken::new();
833        // Keep a clone so the sweep can be stopped via the real shutdown path
834        // (token cancellation) before reopening the same file. redb rejects
835        // reopening a file whose in-process handle is still alive, and the
836        // sweep task holds a cloned Arc<Database> until it exits.
837        let token_for_shutdown = token.clone();
838        {
839            let repo = RedbCacheRepository::new(
840                "redb",
841                path.clone(),
842                Duration::from_secs(60),
843                None,
844                256 * 1024 * 1024,
845                Duration::from_secs(3600),
846                token,
847            )
848            .await
849            .expect("open repo");
850            repo.set("k", entry(), Some(Duration::from_secs(3600)))
851                .await
852                .expect("set");
853            assert_eq!(
854                repo.stats().await.entries,
855                1,
856                "entries counter must be 1 after first insert"
857            );
858            // Stop the sweep via its bound token and await the task so its
859            // cloned Arc<Database> releases before the repo drops. Bind the
860            // handle separately so the MutexGuard drops before the `.await`
861            // (clippy::await-holding-lock).
862            token_for_shutdown.cancel();
863            let sweep_handle = repo.sweep_handle.lock().take();
864            if let Some(handle) = sweep_handle {
865                handle
866                    .await
867                    .expect("sweep task must exit cleanly on token cancel");
868            }
869            // repo dropped here — handle already taken; self.db Arc → 0 → DB
870            // closes → redb frees its in-process lock.
871        }
872        // Reopen with a fresh token; entries must be reloaded from table.len().
873        let token2 = CancellationToken::new();
874        let repo = RedbCacheRepository::new(
875            "redb",
876            path,
877            Duration::from_secs(60),
878            None,
879            256 * 1024 * 1024,
880            Duration::from_secs(3600),
881            token2,
882        )
883        .await
884        .expect("reopen repo");
885        let found = repo.get("k").await.expect("get after reopen");
886        assert!(
887            found.is_some(),
888            "persisted entry must survive drop + reopen"
889        );
890        assert_eq!(
891            repo.stats().await.entries,
892            1,
893            "entries counter must be restored from table.len() on reopen"
894        );
895    }
896
897    #[tokio::test]
898    async fn peek_stale_returns_post_expiry_entry_on_redb() {
899        let dir = tempdir().expect("tempdir");
900        let token = CancellationToken::new();
901        let repo = new_repo(&dir, token).await;
902        repo.set("k", entry(), Some(Duration::from_millis(1)))
903            .await
904            .expect("set");
905        tokio::time::sleep(Duration::from_millis(10)).await;
906        let stale = repo.peek_stale("k").await.expect("peek_stale");
907        assert!(
908            stale.is_some(),
909            "peek_stale must return the expired-but-present entry"
910        );
911    }
912
913    #[tokio::test]
914    async fn sweep_once_removes_entries_past_stale_retention() {
915        let dir = tempdir().expect("tempdir");
916        let token = CancellationToken::new();
917        let path = dir.path().join("cache.redb");
918        // stale_retention=10ms; entry expires at +1ms; sleep 50ms so
919        // (expires_at + 10ms) is well in the past → reclaimable.
920        let repo = RedbCacheRepository::new(
921            "redb",
922            path,
923            Duration::from_millis(10),
924            None,
925            256 * 1024 * 1024,
926            Duration::from_secs(3600),
927            token,
928        )
929        .await
930        .expect("open repo");
931        // Neutralize the background sweep so its immediate first tick (which
932        // fires whenever the runtime next polls it) cannot steal the reclaim
933        // from sweep_once. This makes the test deterministic: sweep_once is
934        // the sole reclaimer.
935        if let Some(handle) = repo.sweep_handle.lock().take() {
936            handle.abort();
937        }
938        repo.set("k", entry(), Some(Duration::from_millis(1)))
939            .await
940            .expect("set");
941        tokio::time::sleep(Duration::from_millis(50)).await;
942        let reclaimed = repo.sweep_once().await.expect("sweep_once");
943        assert!(
944            reclaimed >= 1,
945            "sweep_once must reclaim at least 1 entry, got {reclaimed}"
946        );
947        let stale = repo.peek_stale("k").await.expect("peek_stale after sweep");
948        assert!(
949            stale.is_none(),
950            "entry must be gone after sweep_once reclaimed it"
951        );
952    }
953
954    #[tokio::test]
955    async fn sweep_stops_on_context_shutdown() {
956        let dir = tempdir().expect("tempdir");
957        let token = CancellationToken::new();
958        let path = dir.path().join("cache.redb");
959        // Short sweep interval so the loop is definitely armed.
960        let repo = RedbCacheRepository::new(
961            "redb",
962            path,
963            Duration::from_secs(60),
964            None,
965            256 * 1024 * 1024,
966            Duration::from_millis(10),
967            token.clone(),
968        )
969        .await
970        .expect("open repo");
971        token.cancel();
972        // Take the handle out so we can await it directly; Drop will see None.
973        let handle = repo
974            .sweep_handle
975            .lock()
976            .take()
977            .expect("sweep handle must be present after construct");
978        let completed = tokio::time::timeout(Duration::from_secs(5), handle).await;
979        assert!(
980            completed.is_ok(),
981            "sweep task must complete within 5s of context shutdown"
982        );
983    }
984
985    #[tokio::test]
986    async fn redb_errors_surface_as_err() {
987        let dir = tempdir().expect("tempdir");
988        // A regular file where a directory is required — `create_dir_all`
989        // must fail, surfacing as Err(CamelError::Io(_)).
990        let blocker = dir.path().join("blocker");
991        std::fs::write(&blocker, b"not a dir").expect("write blocker");
992        let path = blocker.join("cache.redb");
993        let result = RedbCacheRepository::new(
994            "redb",
995            path,
996            Duration::from_secs(60),
997            None,
998            256 * 1024 * 1024,
999            Duration::from_secs(3600),
1000            CancellationToken::new(),
1001        )
1002        .await;
1003        assert!(
1004            matches!(result, Err(CamelError::Io(_))),
1005            "expected Err(CamelError::Io(_)), got {result:?}"
1006        );
1007    }
1008
1009    #[tokio::test]
1010    async fn overwrite_does_not_inflate_entries() {
1011        let dir = tempdir().expect("tempdir");
1012        let token = CancellationToken::new();
1013        let repo = new_repo(&dir, token).await;
1014        repo.set("k", entry(), None).await.expect("first set");
1015        repo.set("k", entry(), None).await.expect("second set");
1016        assert_eq!(
1017            repo.stats().await.entries,
1018            1,
1019            "overwriting an existing key must not inflate the entries counter"
1020        );
1021    }
1022
1023    #[tokio::test]
1024    async fn stats_reports_bytes_sum() {
1025        let dir = tempdir().expect("tempdir");
1026        let token = CancellationToken::new();
1027        let repo = new_repo(&dir, token).await;
1028        let a = CacheEntry {
1029            bytes: vec![1, 2, 3],
1030            payload_path: None,
1031            content_type: camel_api::cache::ContentType::Bytes,
1032            expires_at: None,
1033        };
1034        let b = CacheEntry {
1035            bytes: vec![1, 2, 3, 4, 5],
1036            payload_path: None,
1037            content_type: camel_api::cache::ContentType::Bytes,
1038            expires_at: None,
1039        };
1040        repo.set("a", a, None).await.expect("set a");
1041        repo.set("b", b, None).await.expect("set b");
1042        assert_eq!(repo.stats().await.bytes, Some(8));
1043    }
1044
1045    #[tokio::test]
1046    async fn stats_counters_reported_alongside_bytes() {
1047        let dir = tempdir().expect("tempdir");
1048        let token = CancellationToken::new();
1049        let repo = new_repo(&dir, token).await;
1050        let a = CacheEntry {
1051            bytes: vec![1, 2, 3],
1052            payload_path: None,
1053            content_type: camel_api::cache::ContentType::Bytes,
1054            expires_at: None,
1055        };
1056        repo.set("a", a, None).await.expect("set a");
1057        let s = repo.stats().await;
1058        assert_eq!(s.entries, 1);
1059        assert_eq!(s.bytes, Some(3));
1060    }
1061
1062    #[tokio::test]
1063    async fn stats_degrades_bytes_none_when_entry_corrupt() {
1064        let dir = tempdir().expect("tempdir");
1065        let token = CancellationToken::new();
1066        let repo = new_repo(&dir, token).await;
1067
1068        // Exercise every operation counter with a deterministic value.
1069        repo.set("good", entry(), None).await.expect("set good");
1070        repo.get("good").await.expect("get good");
1071        repo.get("absent").await.expect("get absent");
1072        repo.peek_stale("good").await.expect("peek_stale good");
1073        repo.invalidate("absent").await.expect("invalidate absent");
1074
1075        // Corrupt the table out-of-band: a valid key whose value blob is not
1076        // deserializable `CacheEntry` JSON. This bypasses `set`, so the
1077        // operation-derived `entries` counter is untouched (still 1) while the
1078        // physical table now holds a second row.
1079        let garbage: &[u8] = b"not-json";
1080        let txn = repo.db.begin_write().expect("begin_write");
1081        {
1082            let mut table = txn.open_table(CACHE_TABLE).expect("open_table");
1083            table
1084                .insert("corrupt", garbage)
1085                .expect("insert corrupt blob");
1086        }
1087        txn.commit().expect("commit");
1088
1089        // The byte-sum scan hits the corrupt blob and degrades to None; the
1090        // operation counters are unaffected.
1091        let s = repo.stats().await;
1092        assert_eq!(s.hits, 1);
1093        assert_eq!(s.misses, 1);
1094        assert_eq!(s.evictions, 0);
1095        assert_eq!(s.entries, 1);
1096        assert_eq!(s.peek_stale_served, 1);
1097        assert_eq!(s.invalidations, 1);
1098        assert_eq!(s.bytes, None);
1099    }
1100
1101    #[tokio::test]
1102    async fn max_entries_rejects_new_key_allows_overwrite() {
1103        let dir = tempdir().expect("tempdir");
1104        let token = CancellationToken::new();
1105        let path = dir.path().join("cache.redb");
1106        let repo = RedbCacheRepository::new(
1107            "redb",
1108            path,
1109            Duration::from_secs(60),
1110            Some(2),
1111            256 * 1024 * 1024,
1112            Duration::from_secs(3600),
1113            token,
1114        )
1115        .await
1116        .expect("open repo");
1117        repo.set("a", entry(), None).await.expect("set a");
1118        repo.set("b", entry(), None).await.expect("set b");
1119        let over = repo.set("c", entry(), None).await;
1120        assert!(
1121            over.is_err(),
1122            "third distinct key must be rejected at max_entries, got {over:?}"
1123        );
1124        let overw = repo.set("a", entry(), None).await;
1125        assert!(
1126            overw.is_ok(),
1127            "overwrite of an existing key must succeed at max_entries, got {overw:?}"
1128        );
1129    }
1130
1131    #[tokio::test]
1132    async fn invalidate_prefix_removes_namespace_only() {
1133        let dir = tempdir().expect("tempdir");
1134        let token = CancellationToken::new();
1135        let repo = new_repo(&dir, token).await;
1136        repo.set("rainviewer:a", entry(), None)
1137            .await
1138            .expect("set rainviewer:a");
1139        repo.set("rainviewer:b", entry(), None)
1140            .await
1141            .expect("set rainviewer:b");
1142        repo.set("gibs:a", entry(), None).await.expect("set gibs:a");
1143        let deleted = repo
1144            .invalidate_prefix("rainviewer:")
1145            .await
1146            .expect("invalidate_prefix");
1147        assert_eq!(deleted, 2, "only the rainviewer namespace must be removed");
1148        assert!(
1149            repo.get("rainviewer:a").await.expect("get").is_none(),
1150            "rainviewer:a must be gone"
1151        );
1152        assert!(
1153            repo.get("rainviewer:b").await.expect("get").is_none(),
1154            "rainviewer:b must be gone"
1155        );
1156        assert!(
1157            repo.get("gibs:a").await.expect("get").is_some(),
1158            "gibs:a must survive"
1159        );
1160    }
1161
1162    #[tokio::test]
1163    async fn invalidate_prefix_does_not_delete_successor_key() {
1164        let dir = tempdir().expect("tempdir");
1165        let token = CancellationToken::new();
1166        let repo = new_repo(&dir, token).await;
1167        repo.set("ns:", entry(), None).await.expect("set ns:");
1168        repo.set("ns;", entry(), None).await.expect("set ns;");
1169        let deleted = repo
1170            .invalidate_prefix("ns:")
1171            .await
1172            .expect("invalidate_prefix");
1173        assert_eq!(deleted, 1, "only the ns: key must be removed");
1174        assert!(
1175            repo.get("ns:").await.expect("get").is_none(),
1176            "ns: must be gone"
1177        );
1178        assert!(
1179            repo.get("ns;").await.expect("get").is_some(),
1180            "successor key ns; must survive"
1181        );
1182    }
1183
1184    // ── cgroup memory-limit guardrail ─────────────────────────────────────────
1185
1186    /// Runs `f` under a thread-local default `fmt` subscriber that appends into
1187    /// a shared buffer, then returns the captured text.
1188    fn capture_guardrail(f: impl FnOnce()) -> String {
1189        let buf = Arc::new(Mutex::new(Vec::new()));
1190        let subscriber = tracing_subscriber::fmt::Subscriber::builder()
1191            .with_writer(TestWriter {
1192                buf: Arc::clone(&buf),
1193            })
1194            .with_ansi(false)
1195            .finish();
1196        tracing::subscriber::with_default(subscriber, f);
1197        let captured = buf.lock().clone();
1198        String::from_utf8(captured).expect("captured output must be UTF-8")
1199    }
1200
1201    /// `fmt` writer that appends into a shared `Arc<Mutex<Vec<u8>>>`.
1202    struct TestWriter {
1203        buf: Arc<Mutex<Vec<u8>>>,
1204    }
1205
1206    impl std::io::Write for TestWriter {
1207        fn write(&mut self, data: &[u8]) -> std::io::Result<usize> {
1208            self.buf.lock().extend_from_slice(data);
1209            Ok(data.len())
1210        }
1211
1212        fn flush(&mut self) -> std::io::Result<()> {
1213            Ok(())
1214        }
1215    }
1216
1217    impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for TestWriter {
1218        type Writer = TestWriter;
1219
1220        fn make_writer(&'a self) -> Self::Writer {
1221            TestWriter {
1222                buf: Arc::clone(&self.buf),
1223            }
1224        }
1225    }
1226
1227    #[test]
1228    fn cgroup_v2_limit_parsed() {
1229        let dir = tempdir().expect("tempdir");
1230        let v2 = dir.path().join("memory.max");
1231        std::fs::write(&v2, "805306368\n").expect("write v2");
1232        let missing_v1 = dir.path().join("missing-v1");
1233        assert_eq!(memory_limit_from_paths(&v2, &missing_v1), Some(805_306_368));
1234    }
1235
1236    #[test]
1237    fn cgroup_v2_max_means_unlimited() {
1238        let dir = tempdir().expect("tempdir");
1239        let v2 = dir.path().join("memory.max");
1240        std::fs::write(&v2, "max").expect("write v2");
1241        let missing_v1 = dir.path().join("missing-v1");
1242        assert_eq!(memory_limit_from_paths(&v2, &missing_v1), None);
1243    }
1244
1245    #[test]
1246    fn cgroup_v2_malformed_falls_through() {
1247        let dir = tempdir().expect("tempdir");
1248        let v2 = dir.path().join("memory.max");
1249        std::fs::write(&v2, "not-a-number").expect("write v2");
1250        let v1 = dir.path().join("memory.limit_in_bytes");
1251        std::fs::write(&v1, "1073741824").expect("write v1");
1252        assert_eq!(memory_limit_from_paths(&v2, &v1), Some(1_073_741_824));
1253    }
1254
1255    #[test]
1256    fn cgroup_v1_sentinel_unlimited() {
1257        let dir = tempdir().expect("tempdir");
1258        let missing_v2 = dir.path().join("missing-v2");
1259        let v1 = dir.path().join("memory.limit_in_bytes");
1260        std::fs::write(&v1, "9223372036854771712").expect("write v1");
1261        assert_eq!(memory_limit_from_paths(&missing_v2, &v1), None);
1262    }
1263
1264    #[test]
1265    fn cgroup_v1_exactly_16tib_is_a_limit() {
1266        let dir = tempdir().expect("tempdir");
1267        let missing_v2 = dir.path().join("missing-v2");
1268        let v1 = dir.path().join("memory.limit_in_bytes");
1269        std::fs::write(&v1, "17592186044416\n").expect("write v1");
1270        assert_eq!(
1271            memory_limit_from_paths(&missing_v2, &v1),
1272            Some(17_592_186_044_416)
1273        );
1274    }
1275
1276    #[test]
1277    fn successor_bound_unit_tests() {
1278        assert_eq!(
1279            successor_bound("ns:"),
1280            std::ops::Bound::Excluded("ns;".to_string())
1281        );
1282        // Prefix ending U+D7FF must jump the surrogate gap to U+E000.
1283        assert_eq!(
1284            successor_bound("a\u{D7FF}"),
1285            std::ops::Bound::Excluded("a\u{E000}".to_string())
1286        );
1287        // Prefix ending U+E000 increments to U+E001.
1288        assert_eq!(
1289            successor_bound("a\u{E000}"),
1290            std::ops::Bound::Excluded("a\u{E001}".to_string())
1291        );
1292        // Trailing U+10FFFF carries into the preceding scalar: "a…" → "b".
1293        assert_eq!(
1294            successor_bound("a\u{10FFFF}"),
1295            std::ops::Bound::Excluded("b".to_string())
1296        );
1297        // A prefix of only U+10FFFF has no successor.
1298        assert_eq!(
1299            successor_bound("\u{10FFFF}\u{10FFFF}"),
1300            std::ops::Bound::Unbounded
1301        );
1302    }
1303
1304    #[tokio::test]
1305    async fn invalidate_prefix_empty_prefix_removes_all_seeded() {
1306        let dir = tempdir().expect("tempdir");
1307        let token = CancellationToken::new();
1308        let repo = new_repo(&dir, token).await;
1309        repo.set("ns:a", entry(), None).await.expect("set ns:a");
1310        repo.set("ns:b", entry(), None).await.expect("set ns:b");
1311        repo.set("other:c", entry(), None)
1312            .await
1313            .expect("set other:c");
1314        let deleted = repo.invalidate_prefix("").await.expect("invalidate_prefix");
1315        assert_eq!(deleted, 3, "empty prefix must remove every entry");
1316    }
1317
1318    #[tokio::test]
1319    async fn invalidate_prefix_empty_namespace_returns_zero() {
1320        let dir = tempdir().expect("tempdir");
1321        let token = CancellationToken::new();
1322        let repo = new_repo(&dir, token).await;
1323        let deleted = repo
1324            .invalidate_prefix("ns:")
1325            .await
1326            .expect("invalidate_prefix");
1327        assert_eq!(deleted, 0, "absent namespace must report zero removals");
1328    }
1329
1330    #[test]
1331    fn cgroup_files_missing() {
1332        let dir = tempdir().expect("tempdir");
1333        let missing_v2 = dir.path().join("missing-v2");
1334        let missing_v1 = dir.path().join("missing-v1");
1335        assert_eq!(memory_limit_from_paths(&missing_v2, &missing_v1), None);
1336    }
1337
1338    #[test]
1339    fn guardrail_warns_when_exceeds() {
1340        let dir = tempdir().expect("tempdir");
1341        let v2 = dir.path().join("memory.max");
1342        std::fs::write(&v2, "805306368\n").expect("write v2");
1343        let missing_v1 = dir.path().join("missing-v1");
1344        let output = capture_guardrail(|| {
1345            emit_memory_guardrail(1_073_741_824, &v2, &missing_v1);
1346        });
1347        assert!(output.contains("1073741824"), "output: {output}");
1348        assert!(output.contains("805306368"), "output: {output}");
1349        assert_eq!(
1350            output.matches("exceeds container memory limit").count(),
1351            1,
1352            "warn line must appear exactly once: {output}"
1353        );
1354    }
1355
1356    #[test]
1357    fn guardrail_silent_when_fits() {
1358        let dir = tempdir().expect("tempdir");
1359        let v2 = dir.path().join("memory.max");
1360        std::fs::write(&v2, "805306368\n").expect("write v2");
1361        let missing_v1 = dir.path().join("missing-v1");
1362        let output = capture_guardrail(|| {
1363            emit_memory_guardrail(268_435_456, &v2, &missing_v1);
1364        });
1365        assert!(output.is_empty(), "expected no output, got: {output}");
1366    }
1367
1368    #[test]
1369    fn guardrail_silent_when_files_missing() {
1370        let dir = tempdir().expect("tempdir");
1371        let missing_v2 = dir.path().join("missing-v2");
1372        let missing_v1 = dir.path().join("missing-v1");
1373        let output = capture_guardrail(|| {
1374            emit_memory_guardrail(1_073_741_824, &missing_v2, &missing_v1);
1375        });
1376        assert!(output.is_empty(), "expected no output, got: {output}");
1377    }
1378}