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_stale(&self, key: &str) -> Result<Option<CacheEntry>, CamelError> {
478        let db = Arc::clone(&self.db);
479        let key = key.to_string();
480        let result =
481            tokio::task::spawn_blocking(move || -> Result<Option<CacheEntry>, CamelError> {
482                let rtx = db
483                    .begin_read()
484                    .map_err(|e| CamelError::Io(format!("redb begin_read: {e}")))?;
485                let table = rtx
486                    .open_table(CACHE_TABLE)
487                    .map_err(|e| CamelError::Io(format!("redb open_table: {e}")))?;
488                match table
489                    .get(key.as_str())
490                    .map_err(|e| CamelError::Io(format!("redb get: {e}")))?
491                {
492                    Some(guard) => {
493                        let entry: CacheEntry = serde_json::from_slice(guard.value())
494                            .map_err(|e| CamelError::Io(format!("cache deserialization: {e}")))?;
495                        Ok(Some(entry))
496                    }
497                    None => Ok(None),
498                }
499            })
500            .await
501            .map_err(|e| CamelError::Io(format!("spawn_blocking join: {e}")))??;
502        if result.is_some() {
503            self.peek_stale_served.fetch_add(1, Ordering::Relaxed);
504        }
505        Ok(result)
506    }
507
508    async fn invalidate(&self, key: &str) -> Result<(), CamelError> {
509        let db = Arc::clone(&self.db);
510        let key = key.to_string();
511        let was_present = tokio::task::spawn_blocking(move || {
512            let txn = db
513                .begin_write()
514                .map_err(|e| CamelError::Io(format!("redb begin_write: {e}")))?;
515            let was_present = {
516                let mut table = txn
517                    .open_table(CACHE_TABLE)
518                    .map_err(|e| CamelError::Io(format!("redb open_table: {e}")))?;
519                // `remove` is idempotent per the trait contract; the returned
520                // Option tells us whether a value was actually present. Read
521                // `.is_some()` now so the guard (borrowing `table`) drops
522                // before `commit` moves `txn`.
523                table
524                    .remove(key.as_str())
525                    .map_err(|e| CamelError::Io(format!("redb remove: {e}")))?
526                    .is_some()
527            };
528            txn.commit()
529                .map_err(|e| CamelError::Io(format!("redb commit: {e}")))?;
530            Ok::<bool, CamelError>(was_present)
531        })
532        .await
533        .map_err(|e| CamelError::Io(format!("spawn_blocking join: {e}")))??;
534        if was_present {
535            let current = self.entries.load(Ordering::Relaxed);
536            let sub = std::cmp::min(current, 1);
537            self.entries.fetch_sub(sub, Ordering::Relaxed);
538        }
539        self.invalidations.fetch_add(1, Ordering::Relaxed);
540        Ok(())
541    }
542
543    async fn invalidate_prefix(&self, prefix: &str) -> Result<u64, CamelError> {
544        let db = Arc::clone(&self.db);
545        let prefix = prefix.to_string();
546        let deleted = tokio::task::spawn_blocking(move || -> Result<u64, CamelError> {
547            // Collect matching keys in a read txn, then delete in one write
548            // txn (mirrors the collect-then-remove pattern of `clear`).
549            let keys: Vec<String> = {
550                let rtx = db
551                    .begin_read()
552                    .map_err(|e| CamelError::Io(format!("redb begin_read: {e}")))?;
553                let table = rtx
554                    .open_table(CACHE_TABLE)
555                    .map_err(|e| CamelError::Io(format!("redb open_table: {e}")))?;
556                // `successor_bound` yields an owned bound; redb's range needs
557                // `&str` bounds for `&str` keys (`String` is not `Borrow<&str>`).
558                let upper: Bound<String> = successor_bound(&prefix);
559                let upper_ref: Bound<&str> = match &upper {
560                    Bound::Included(s) => Bound::Included(s.as_str()),
561                    Bound::Excluded(s) => Bound::Excluded(s.as_str()),
562                    Bound::Unbounded => Bound::Unbounded,
563                };
564                let mut keys = Vec::new();
565                for row in table
566                    .range::<&str>((Bound::Included(prefix.as_str()), upper_ref))
567                    .map_err(|e| CamelError::Io(format!("redb range: {e}")))?
568                {
569                    let (k, _v) =
570                        row.map_err(|e| CamelError::Io(format!("redb range item: {e}")))?;
571                    keys.push(k.value().to_string());
572                }
573                keys
574            };
575            let wtx = db
576                .begin_write()
577                .map_err(|e| CamelError::Io(format!("redb begin_write: {e}")))?;
578            {
579                let mut table = wtx
580                    .open_table(CACHE_TABLE)
581                    .map_err(|e| CamelError::Io(format!("redb open_table: {e}")))?;
582                for k in &keys {
583                    let _ = table
584                        .remove(k.as_str())
585                        .map_err(|e| CamelError::Io(format!("redb remove: {e}")))?;
586                }
587            }
588            wtx.commit()
589                .map_err(|e| CamelError::Io(format!("redb commit: {e}")))?;
590            Ok(keys.len() as u64)
591        })
592        .await
593        .map_err(|e| CamelError::Io(format!("spawn_blocking join: {e}")))??;
594        self.invalidations.fetch_add(1, Ordering::Relaxed);
595        if deleted > 0 {
596            let current = self.entries.load(Ordering::Relaxed);
597            let sub = std::cmp::min(current, deleted);
598            self.entries.fetch_sub(sub, Ordering::Relaxed);
599        }
600        Ok(deleted)
601    }
602
603    async fn clear(&self) -> Result<(), CamelError> {
604        let db = Arc::clone(&self.db);
605        tokio::task::spawn_blocking(move || {
606            let txn = db
607                .begin_write()
608                .map_err(|e| CamelError::Io(format!("redb begin_write: {e}")))?;
609            {
610                let mut table = txn
611                    .open_table(CACHE_TABLE)
612                    .map_err(|e| CamelError::Io(format!("redb open_table: {e}")))?;
613                // Collect-then-remove: `iter()` borrows immutably and cannot
614                // coexist with `remove`. Mirrors `idempotent/redb_repository`
615                // clear pattern.
616                let keys: Vec<String> = table
617                    .iter()
618                    .map_err(|e| CamelError::Io(format!("redb iter: {e}")))?
619                    .map(|r| {
620                        r.map(|(k, _v)| k.value().to_string())
621                            .map_err(|e| CamelError::Io(format!("redb iter item: {e}")))
622                    })
623                    .collect::<Result<_, _>>()?;
624                for k in &keys {
625                    let _ = table
626                        .remove(k.as_str())
627                        .map_err(|e| CamelError::Io(format!("redb remove: {e}")))?;
628                }
629            }
630            txn.commit()
631                .map_err(|e| CamelError::Io(format!("redb commit: {e}")))?;
632            Ok::<_, CamelError>(())
633        })
634        .await
635        .map_err(|e| CamelError::Io(format!("spawn_blocking join: {e}")))??;
636        self.entries.store(0, Ordering::Relaxed);
637        Ok(())
638    }
639
640    async fn stats(&self) -> CacheStats {
641        let db = Arc::clone(&self.db);
642        let bytes = tokio::task::spawn_blocking(move || total_bytes(&db))
643            .await
644            .unwrap_or_default();
645        CacheStats {
646            hits: self.hits.load(Ordering::Relaxed),
647            misses: self.misses.load(Ordering::Relaxed),
648            evictions: self.evictions.load(Ordering::Relaxed),
649            entries: self.entries.load(Ordering::Relaxed),
650            peek_stale_served: self.peek_stale_served.load(Ordering::Relaxed),
651            invalidations: self.invalidations.load(Ordering::Relaxed),
652            bytes,
653        }
654    }
655}
656
657/// Sum every entry's `bytes.len()` over the full table range.
658///
659/// Called only inside `spawn_blocking` from `stats()`; `None` when the table
660/// cannot be read or any entry fails to deserialize — the `bytes` field is a
661/// best-effort report, never an error.
662fn total_bytes(db: &redb::Database) -> Option<u64> {
663    let rtx = db.begin_read().ok()?;
664    let table = rtx.open_table(CACHE_TABLE).ok()?;
665    let mut total: u64 = 0;
666    for row in table.iter().ok()? {
667        let (_key, value) = row.ok()?;
668        let entry: CacheEntry = serde_json::from_slice(value.value()).ok()?;
669        total = total.saturating_add(entry.bytes.len() as u64);
670    }
671    Some(total)
672}
673
674impl fmt::Debug for RedbCacheRepository {
675    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
676        f.debug_struct("RedbCacheRepository")
677            .field("name", &self.name)
678            .field("stale_retention", &self.stale_retention)
679            .field("max_entries", &self.max_entries)
680            .field("cache_size", &self.cache_size)
681            .field("sweep_interval", &self.sweep_interval)
682            .field("shutdown_cancelled", &self.shutdown_token.is_cancelled())
683            .finish()
684    }
685}
686
687impl Drop for RedbCacheRepository {
688    fn drop(&mut self) {
689        // Abort ONLY the sweep task. Never cancel the context-owned token —
690        // that would shut down the entire context when one repo drops.
691        if let Some(handle) = self.sweep_handle.lock().take() {
692            handle.abort();
693        }
694    }
695}
696
697#[cfg(test)]
698mod tests {
699    use super::*;
700    use tempfile::TempDir;
701    use tempfile::tempdir;
702
703    fn entry() -> CacheEntry {
704        CacheEntry {
705            bytes: vec![1, 2, 3],
706            payload_path: None,
707            content_type: camel_api::cache::ContentType::Bytes,
708            expires_at: None,
709        }
710    }
711
712    /// Open a repo at `<tmp>/cache.redb` with a 60s stale-retention, no cap,
713    /// a 256 MiB cache size, and a 1h sweep interval (so the background loop
714    /// stays dormant during sequential tests).
715    async fn new_repo(tmp: &TempDir, shutdown_token: CancellationToken) -> RedbCacheRepository {
716        new_repo_with(
717            tmp,
718            shutdown_token,
719            Duration::from_secs(60),
720            None,
721            256 * 1024 * 1024,
722            Duration::from_secs(3600),
723        )
724        .await
725    }
726
727    /// Full-parameter variant of [`new_repo`] for tests that need custom
728    /// stale-retention, cap, cache size, or sweep interval values.
729    async fn new_repo_with(
730        tmp: &TempDir,
731        shutdown_token: CancellationToken,
732        stale_retention: Duration,
733        max_entries: Option<usize>,
734        cache_size: usize,
735        sweep_interval: Duration,
736    ) -> RedbCacheRepository {
737        let path = tmp.path().join("cache.redb");
738        RedbCacheRepository::new(
739            "redb",
740            path,
741            stale_retention,
742            max_entries,
743            cache_size,
744            sweep_interval,
745            shutdown_token,
746        )
747        .await
748        .expect("open redb cache repo")
749    }
750
751    #[tokio::test]
752    async fn cache_size_recorded_and_accessible() {
753        let dir = tempdir().expect("tempdir");
754        let token = CancellationToken::new();
755        let repo = new_repo_with(
756            &dir,
757            token,
758            Duration::from_secs(60),
759            None,
760            536_870_912,
761            Duration::from_secs(3600),
762        )
763        .await;
764        assert_eq!(repo.cache_size(), 536_870_912);
765    }
766
767    #[tokio::test]
768    async fn sweep_interval_recorded_and_accessible() {
769        let dir = tempdir().expect("tempdir");
770        let token = CancellationToken::new();
771        let repo = new_repo_with(
772            &dir,
773            token,
774            Duration::from_secs(60),
775            None,
776            256 * 1024 * 1024,
777            Duration::from_secs(1800),
778        )
779        .await;
780        assert_eq!(repo.sweep_interval(), Duration::from_secs(1800));
781    }
782
783    #[tokio::test]
784    async fn stale_retention_recorded_and_accessible() {
785        let dir = tempdir().expect("tempdir");
786        let token = CancellationToken::new();
787        let repo = new_repo_with(
788            &dir,
789            token,
790            Duration::from_secs(3600),
791            None,
792            256 * 1024 * 1024,
793            Duration::from_secs(3600),
794        )
795        .await;
796        assert_eq!(repo.stale_retention(), Duration::from_secs(3600));
797    }
798
799    #[tokio::test]
800    async fn explicit_cache_size_round_trip() {
801        let dir = tempdir().expect("tempdir");
802        let token = CancellationToken::new();
803        let repo = new_repo_with(
804            &dir,
805            token,
806            Duration::from_secs(60),
807            None,
808            512 * 1024 * 1024,
809            Duration::from_secs(3600),
810        )
811        .await;
812        repo.set("k", entry(), Some(Duration::from_secs(3600)))
813            .await
814            .expect("set");
815        let found = repo.get("k").await.expect("get");
816        assert!(
817            found.is_some(),
818            "entry must round-trip through the builder-opened database"
819        );
820    }
821
822    #[tokio::test]
823    async fn entries_survive_handle_drop_and_reopen() {
824        let dir = tempdir().expect("tempdir");
825        let path = dir.path().join("cache.redb");
826        let token = CancellationToken::new();
827        // Keep a clone so the sweep can be stopped via the real shutdown path
828        // (token cancellation) before reopening the same file. redb rejects
829        // reopening a file whose in-process handle is still alive, and the
830        // sweep task holds a cloned Arc<Database> until it exits.
831        let token_for_shutdown = token.clone();
832        {
833            let repo = RedbCacheRepository::new(
834                "redb",
835                path.clone(),
836                Duration::from_secs(60),
837                None,
838                256 * 1024 * 1024,
839                Duration::from_secs(3600),
840                token,
841            )
842            .await
843            .expect("open repo");
844            repo.set("k", entry(), Some(Duration::from_secs(3600)))
845                .await
846                .expect("set");
847            assert_eq!(
848                repo.stats().await.entries,
849                1,
850                "entries counter must be 1 after first insert"
851            );
852            // Stop the sweep via its bound token and await the task so its
853            // cloned Arc<Database> releases before the repo drops. Bind the
854            // handle separately so the MutexGuard drops before the `.await`
855            // (clippy::await-holding-lock).
856            token_for_shutdown.cancel();
857            let sweep_handle = repo.sweep_handle.lock().take();
858            if let Some(handle) = sweep_handle {
859                handle
860                    .await
861                    .expect("sweep task must exit cleanly on token cancel");
862            }
863            // repo dropped here — handle already taken; self.db Arc → 0 → DB
864            // closes → redb frees its in-process lock.
865        }
866        // Reopen with a fresh token; entries must be reloaded from table.len().
867        let token2 = CancellationToken::new();
868        let repo = RedbCacheRepository::new(
869            "redb",
870            path,
871            Duration::from_secs(60),
872            None,
873            256 * 1024 * 1024,
874            Duration::from_secs(3600),
875            token2,
876        )
877        .await
878        .expect("reopen repo");
879        let found = repo.get("k").await.expect("get after reopen");
880        assert!(
881            found.is_some(),
882            "persisted entry must survive drop + reopen"
883        );
884        assert_eq!(
885            repo.stats().await.entries,
886            1,
887            "entries counter must be restored from table.len() on reopen"
888        );
889    }
890
891    #[tokio::test]
892    async fn peek_stale_returns_post_expiry_entry_on_redb() {
893        let dir = tempdir().expect("tempdir");
894        let token = CancellationToken::new();
895        let repo = new_repo(&dir, token).await;
896        repo.set("k", entry(), Some(Duration::from_millis(1)))
897            .await
898            .expect("set");
899        tokio::time::sleep(Duration::from_millis(10)).await;
900        let stale = repo.peek_stale("k").await.expect("peek_stale");
901        assert!(
902            stale.is_some(),
903            "peek_stale must return the expired-but-present entry"
904        );
905    }
906
907    #[tokio::test]
908    async fn sweep_once_removes_entries_past_stale_retention() {
909        let dir = tempdir().expect("tempdir");
910        let token = CancellationToken::new();
911        let path = dir.path().join("cache.redb");
912        // stale_retention=10ms; entry expires at +1ms; sleep 50ms so
913        // (expires_at + 10ms) is well in the past → reclaimable.
914        let repo = RedbCacheRepository::new(
915            "redb",
916            path,
917            Duration::from_millis(10),
918            None,
919            256 * 1024 * 1024,
920            Duration::from_secs(3600),
921            token,
922        )
923        .await
924        .expect("open repo");
925        // Neutralize the background sweep so its immediate first tick (which
926        // fires whenever the runtime next polls it) cannot steal the reclaim
927        // from sweep_once. This makes the test deterministic: sweep_once is
928        // the sole reclaimer.
929        if let Some(handle) = repo.sweep_handle.lock().take() {
930            handle.abort();
931        }
932        repo.set("k", entry(), Some(Duration::from_millis(1)))
933            .await
934            .expect("set");
935        tokio::time::sleep(Duration::from_millis(50)).await;
936        let reclaimed = repo.sweep_once().await.expect("sweep_once");
937        assert!(
938            reclaimed >= 1,
939            "sweep_once must reclaim at least 1 entry, got {reclaimed}"
940        );
941        let stale = repo.peek_stale("k").await.expect("peek_stale after sweep");
942        assert!(
943            stale.is_none(),
944            "entry must be gone after sweep_once reclaimed it"
945        );
946    }
947
948    #[tokio::test]
949    async fn sweep_stops_on_context_shutdown() {
950        let dir = tempdir().expect("tempdir");
951        let token = CancellationToken::new();
952        let path = dir.path().join("cache.redb");
953        // Short sweep interval so the loop is definitely armed.
954        let repo = RedbCacheRepository::new(
955            "redb",
956            path,
957            Duration::from_secs(60),
958            None,
959            256 * 1024 * 1024,
960            Duration::from_millis(10),
961            token.clone(),
962        )
963        .await
964        .expect("open repo");
965        token.cancel();
966        // Take the handle out so we can await it directly; Drop will see None.
967        let handle = repo
968            .sweep_handle
969            .lock()
970            .take()
971            .expect("sweep handle must be present after construct");
972        let completed = tokio::time::timeout(Duration::from_secs(5), handle).await;
973        assert!(
974            completed.is_ok(),
975            "sweep task must complete within 5s of context shutdown"
976        );
977    }
978
979    #[tokio::test]
980    async fn redb_errors_surface_as_err() {
981        let dir = tempdir().expect("tempdir");
982        // A regular file where a directory is required — `create_dir_all`
983        // must fail, surfacing as Err(CamelError::Io(_)).
984        let blocker = dir.path().join("blocker");
985        std::fs::write(&blocker, b"not a dir").expect("write blocker");
986        let path = blocker.join("cache.redb");
987        let result = RedbCacheRepository::new(
988            "redb",
989            path,
990            Duration::from_secs(60),
991            None,
992            256 * 1024 * 1024,
993            Duration::from_secs(3600),
994            CancellationToken::new(),
995        )
996        .await;
997        assert!(
998            matches!(result, Err(CamelError::Io(_))),
999            "expected Err(CamelError::Io(_)), got {result:?}"
1000        );
1001    }
1002
1003    #[tokio::test]
1004    async fn overwrite_does_not_inflate_entries() {
1005        let dir = tempdir().expect("tempdir");
1006        let token = CancellationToken::new();
1007        let repo = new_repo(&dir, token).await;
1008        repo.set("k", entry(), None).await.expect("first set");
1009        repo.set("k", entry(), None).await.expect("second set");
1010        assert_eq!(
1011            repo.stats().await.entries,
1012            1,
1013            "overwriting an existing key must not inflate the entries counter"
1014        );
1015    }
1016
1017    #[tokio::test]
1018    async fn stats_reports_bytes_sum() {
1019        let dir = tempdir().expect("tempdir");
1020        let token = CancellationToken::new();
1021        let repo = new_repo(&dir, token).await;
1022        let a = CacheEntry {
1023            bytes: vec![1, 2, 3],
1024            payload_path: None,
1025            content_type: camel_api::cache::ContentType::Bytes,
1026            expires_at: None,
1027        };
1028        let b = CacheEntry {
1029            bytes: vec![1, 2, 3, 4, 5],
1030            payload_path: None,
1031            content_type: camel_api::cache::ContentType::Bytes,
1032            expires_at: None,
1033        };
1034        repo.set("a", a, None).await.expect("set a");
1035        repo.set("b", b, None).await.expect("set b");
1036        assert_eq!(repo.stats().await.bytes, Some(8));
1037    }
1038
1039    #[tokio::test]
1040    async fn stats_counters_reported_alongside_bytes() {
1041        let dir = tempdir().expect("tempdir");
1042        let token = CancellationToken::new();
1043        let repo = new_repo(&dir, token).await;
1044        let a = CacheEntry {
1045            bytes: vec![1, 2, 3],
1046            payload_path: None,
1047            content_type: camel_api::cache::ContentType::Bytes,
1048            expires_at: None,
1049        };
1050        repo.set("a", a, None).await.expect("set a");
1051        let s = repo.stats().await;
1052        assert_eq!(s.entries, 1);
1053        assert_eq!(s.bytes, Some(3));
1054    }
1055
1056    #[tokio::test]
1057    async fn stats_degrades_bytes_none_when_entry_corrupt() {
1058        let dir = tempdir().expect("tempdir");
1059        let token = CancellationToken::new();
1060        let repo = new_repo(&dir, token).await;
1061
1062        // Exercise every operation counter with a deterministic value.
1063        repo.set("good", entry(), None).await.expect("set good");
1064        repo.get("good").await.expect("get good");
1065        repo.get("absent").await.expect("get absent");
1066        repo.peek_stale("good").await.expect("peek_stale good");
1067        repo.invalidate("absent").await.expect("invalidate absent");
1068
1069        // Corrupt the table out-of-band: a valid key whose value blob is not
1070        // deserializable `CacheEntry` JSON. This bypasses `set`, so the
1071        // operation-derived `entries` counter is untouched (still 1) while the
1072        // physical table now holds a second row.
1073        let garbage: &[u8] = b"not-json";
1074        let txn = repo.db.begin_write().expect("begin_write");
1075        {
1076            let mut table = txn.open_table(CACHE_TABLE).expect("open_table");
1077            table
1078                .insert("corrupt", garbage)
1079                .expect("insert corrupt blob");
1080        }
1081        txn.commit().expect("commit");
1082
1083        // The byte-sum scan hits the corrupt blob and degrades to None; the
1084        // operation counters are unaffected.
1085        let s = repo.stats().await;
1086        assert_eq!(s.hits, 1);
1087        assert_eq!(s.misses, 1);
1088        assert_eq!(s.evictions, 0);
1089        assert_eq!(s.entries, 1);
1090        assert_eq!(s.peek_stale_served, 1);
1091        assert_eq!(s.invalidations, 1);
1092        assert_eq!(s.bytes, None);
1093    }
1094
1095    #[tokio::test]
1096    async fn max_entries_rejects_new_key_allows_overwrite() {
1097        let dir = tempdir().expect("tempdir");
1098        let token = CancellationToken::new();
1099        let path = dir.path().join("cache.redb");
1100        let repo = RedbCacheRepository::new(
1101            "redb",
1102            path,
1103            Duration::from_secs(60),
1104            Some(2),
1105            256 * 1024 * 1024,
1106            Duration::from_secs(3600),
1107            token,
1108        )
1109        .await
1110        .expect("open repo");
1111        repo.set("a", entry(), None).await.expect("set a");
1112        repo.set("b", entry(), None).await.expect("set b");
1113        let over = repo.set("c", entry(), None).await;
1114        assert!(
1115            over.is_err(),
1116            "third distinct key must be rejected at max_entries, got {over:?}"
1117        );
1118        let overw = repo.set("a", entry(), None).await;
1119        assert!(
1120            overw.is_ok(),
1121            "overwrite of an existing key must succeed at max_entries, got {overw:?}"
1122        );
1123    }
1124
1125    #[tokio::test]
1126    async fn invalidate_prefix_removes_namespace_only() {
1127        let dir = tempdir().expect("tempdir");
1128        let token = CancellationToken::new();
1129        let repo = new_repo(&dir, token).await;
1130        repo.set("rainviewer:a", entry(), None)
1131            .await
1132            .expect("set rainviewer:a");
1133        repo.set("rainviewer:b", entry(), None)
1134            .await
1135            .expect("set rainviewer:b");
1136        repo.set("gibs:a", entry(), None).await.expect("set gibs:a");
1137        let deleted = repo
1138            .invalidate_prefix("rainviewer:")
1139            .await
1140            .expect("invalidate_prefix");
1141        assert_eq!(deleted, 2, "only the rainviewer namespace must be removed");
1142        assert!(
1143            repo.get("rainviewer:a").await.expect("get").is_none(),
1144            "rainviewer:a must be gone"
1145        );
1146        assert!(
1147            repo.get("rainviewer:b").await.expect("get").is_none(),
1148            "rainviewer:b must be gone"
1149        );
1150        assert!(
1151            repo.get("gibs:a").await.expect("get").is_some(),
1152            "gibs:a must survive"
1153        );
1154    }
1155
1156    #[tokio::test]
1157    async fn invalidate_prefix_does_not_delete_successor_key() {
1158        let dir = tempdir().expect("tempdir");
1159        let token = CancellationToken::new();
1160        let repo = new_repo(&dir, token).await;
1161        repo.set("ns:", entry(), None).await.expect("set ns:");
1162        repo.set("ns;", entry(), None).await.expect("set ns;");
1163        let deleted = repo
1164            .invalidate_prefix("ns:")
1165            .await
1166            .expect("invalidate_prefix");
1167        assert_eq!(deleted, 1, "only the ns: key must be removed");
1168        assert!(
1169            repo.get("ns:").await.expect("get").is_none(),
1170            "ns: must be gone"
1171        );
1172        assert!(
1173            repo.get("ns;").await.expect("get").is_some(),
1174            "successor key ns; must survive"
1175        );
1176    }
1177
1178    // ── cgroup memory-limit guardrail ─────────────────────────────────────────
1179
1180    /// Runs `f` under a thread-local default `fmt` subscriber that appends into
1181    /// a shared buffer, then returns the captured text.
1182    fn capture_guardrail(f: impl FnOnce()) -> String {
1183        let buf = Arc::new(Mutex::new(Vec::new()));
1184        let subscriber = tracing_subscriber::fmt::Subscriber::builder()
1185            .with_writer(TestWriter {
1186                buf: Arc::clone(&buf),
1187            })
1188            .with_ansi(false)
1189            .finish();
1190        tracing::subscriber::with_default(subscriber, f);
1191        let captured = buf.lock().clone();
1192        String::from_utf8(captured).expect("captured output must be UTF-8")
1193    }
1194
1195    /// `fmt` writer that appends into a shared `Arc<Mutex<Vec<u8>>>`.
1196    struct TestWriter {
1197        buf: Arc<Mutex<Vec<u8>>>,
1198    }
1199
1200    impl std::io::Write for TestWriter {
1201        fn write(&mut self, data: &[u8]) -> std::io::Result<usize> {
1202            self.buf.lock().extend_from_slice(data);
1203            Ok(data.len())
1204        }
1205
1206        fn flush(&mut self) -> std::io::Result<()> {
1207            Ok(())
1208        }
1209    }
1210
1211    impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for TestWriter {
1212        type Writer = TestWriter;
1213
1214        fn make_writer(&'a self) -> Self::Writer {
1215            TestWriter {
1216                buf: Arc::clone(&self.buf),
1217            }
1218        }
1219    }
1220
1221    #[test]
1222    fn cgroup_v2_limit_parsed() {
1223        let dir = tempdir().expect("tempdir");
1224        let v2 = dir.path().join("memory.max");
1225        std::fs::write(&v2, "805306368\n").expect("write v2");
1226        let missing_v1 = dir.path().join("missing-v1");
1227        assert_eq!(memory_limit_from_paths(&v2, &missing_v1), Some(805_306_368));
1228    }
1229
1230    #[test]
1231    fn cgroup_v2_max_means_unlimited() {
1232        let dir = tempdir().expect("tempdir");
1233        let v2 = dir.path().join("memory.max");
1234        std::fs::write(&v2, "max").expect("write v2");
1235        let missing_v1 = dir.path().join("missing-v1");
1236        assert_eq!(memory_limit_from_paths(&v2, &missing_v1), None);
1237    }
1238
1239    #[test]
1240    fn cgroup_v2_malformed_falls_through() {
1241        let dir = tempdir().expect("tempdir");
1242        let v2 = dir.path().join("memory.max");
1243        std::fs::write(&v2, "not-a-number").expect("write v2");
1244        let v1 = dir.path().join("memory.limit_in_bytes");
1245        std::fs::write(&v1, "1073741824").expect("write v1");
1246        assert_eq!(memory_limit_from_paths(&v2, &v1), Some(1_073_741_824));
1247    }
1248
1249    #[test]
1250    fn cgroup_v1_sentinel_unlimited() {
1251        let dir = tempdir().expect("tempdir");
1252        let missing_v2 = dir.path().join("missing-v2");
1253        let v1 = dir.path().join("memory.limit_in_bytes");
1254        std::fs::write(&v1, "9223372036854771712").expect("write v1");
1255        assert_eq!(memory_limit_from_paths(&missing_v2, &v1), None);
1256    }
1257
1258    #[test]
1259    fn cgroup_v1_exactly_16tib_is_a_limit() {
1260        let dir = tempdir().expect("tempdir");
1261        let missing_v2 = dir.path().join("missing-v2");
1262        let v1 = dir.path().join("memory.limit_in_bytes");
1263        std::fs::write(&v1, "17592186044416\n").expect("write v1");
1264        assert_eq!(
1265            memory_limit_from_paths(&missing_v2, &v1),
1266            Some(17_592_186_044_416)
1267        );
1268    }
1269
1270    #[test]
1271    fn successor_bound_unit_tests() {
1272        assert_eq!(
1273            successor_bound("ns:"),
1274            std::ops::Bound::Excluded("ns;".to_string())
1275        );
1276        // Prefix ending U+D7FF must jump the surrogate gap to U+E000.
1277        assert_eq!(
1278            successor_bound("a\u{D7FF}"),
1279            std::ops::Bound::Excluded("a\u{E000}".to_string())
1280        );
1281        // Prefix ending U+E000 increments to U+E001.
1282        assert_eq!(
1283            successor_bound("a\u{E000}"),
1284            std::ops::Bound::Excluded("a\u{E001}".to_string())
1285        );
1286        // Trailing U+10FFFF carries into the preceding scalar: "a…" → "b".
1287        assert_eq!(
1288            successor_bound("a\u{10FFFF}"),
1289            std::ops::Bound::Excluded("b".to_string())
1290        );
1291        // A prefix of only U+10FFFF has no successor.
1292        assert_eq!(
1293            successor_bound("\u{10FFFF}\u{10FFFF}"),
1294            std::ops::Bound::Unbounded
1295        );
1296    }
1297
1298    #[tokio::test]
1299    async fn invalidate_prefix_empty_prefix_removes_all_seeded() {
1300        let dir = tempdir().expect("tempdir");
1301        let token = CancellationToken::new();
1302        let repo = new_repo(&dir, token).await;
1303        repo.set("ns:a", entry(), None).await.expect("set ns:a");
1304        repo.set("ns:b", entry(), None).await.expect("set ns:b");
1305        repo.set("other:c", entry(), None)
1306            .await
1307            .expect("set other:c");
1308        let deleted = repo.invalidate_prefix("").await.expect("invalidate_prefix");
1309        assert_eq!(deleted, 3, "empty prefix must remove every entry");
1310    }
1311
1312    #[tokio::test]
1313    async fn invalidate_prefix_empty_namespace_returns_zero() {
1314        let dir = tempdir().expect("tempdir");
1315        let token = CancellationToken::new();
1316        let repo = new_repo(&dir, token).await;
1317        let deleted = repo
1318            .invalidate_prefix("ns:")
1319            .await
1320            .expect("invalidate_prefix");
1321        assert_eq!(deleted, 0, "absent namespace must report zero removals");
1322    }
1323
1324    #[test]
1325    fn cgroup_files_missing() {
1326        let dir = tempdir().expect("tempdir");
1327        let missing_v2 = dir.path().join("missing-v2");
1328        let missing_v1 = dir.path().join("missing-v1");
1329        assert_eq!(memory_limit_from_paths(&missing_v2, &missing_v1), None);
1330    }
1331
1332    #[test]
1333    fn guardrail_warns_when_exceeds() {
1334        let dir = tempdir().expect("tempdir");
1335        let v2 = dir.path().join("memory.max");
1336        std::fs::write(&v2, "805306368\n").expect("write v2");
1337        let missing_v1 = dir.path().join("missing-v1");
1338        let output = capture_guardrail(|| {
1339            emit_memory_guardrail(1_073_741_824, &v2, &missing_v1);
1340        });
1341        assert!(output.contains("1073741824"), "output: {output}");
1342        assert!(output.contains("805306368"), "output: {output}");
1343        assert_eq!(
1344            output.matches("exceeds container memory limit").count(),
1345            1,
1346            "warn line must appear exactly once: {output}"
1347        );
1348    }
1349
1350    #[test]
1351    fn guardrail_silent_when_fits() {
1352        let dir = tempdir().expect("tempdir");
1353        let v2 = dir.path().join("memory.max");
1354        std::fs::write(&v2, "805306368\n").expect("write v2");
1355        let missing_v1 = dir.path().join("missing-v1");
1356        let output = capture_guardrail(|| {
1357            emit_memory_guardrail(268_435_456, &v2, &missing_v1);
1358        });
1359        assert!(output.is_empty(), "expected no output, got: {output}");
1360    }
1361
1362    #[test]
1363    fn guardrail_silent_when_files_missing() {
1364        let dir = tempdir().expect("tempdir");
1365        let missing_v2 = dir.path().join("missing-v2");
1366        let missing_v1 = dir.path().join("missing-v1");
1367        let output = capture_guardrail(|| {
1368            emit_memory_guardrail(1_073_741_824, &missing_v2, &missing_v1);
1369        });
1370        assert!(output.is_empty(), "expected no output, got: {output}");
1371    }
1372}