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            content_type: camel_api::cache::ContentType::Bytes,
707            expires_at: None,
708        }
709    }
710
711    /// Open a repo at `<tmp>/cache.redb` with a 60s stale-retention, no cap,
712    /// a 256 MiB cache size, and a 1h sweep interval (so the background loop
713    /// stays dormant during sequential tests).
714    async fn new_repo(tmp: &TempDir, shutdown_token: CancellationToken) -> RedbCacheRepository {
715        new_repo_with(
716            tmp,
717            shutdown_token,
718            Duration::from_secs(60),
719            None,
720            256 * 1024 * 1024,
721            Duration::from_secs(3600),
722        )
723        .await
724    }
725
726    /// Full-parameter variant of [`new_repo`] for tests that need custom
727    /// stale-retention, cap, cache size, or sweep interval values.
728    async fn new_repo_with(
729        tmp: &TempDir,
730        shutdown_token: CancellationToken,
731        stale_retention: Duration,
732        max_entries: Option<usize>,
733        cache_size: usize,
734        sweep_interval: Duration,
735    ) -> RedbCacheRepository {
736        let path = tmp.path().join("cache.redb");
737        RedbCacheRepository::new(
738            "redb",
739            path,
740            stale_retention,
741            max_entries,
742            cache_size,
743            sweep_interval,
744            shutdown_token,
745        )
746        .await
747        .expect("open redb cache repo")
748    }
749
750    #[tokio::test]
751    async fn cache_size_recorded_and_accessible() {
752        let dir = tempdir().expect("tempdir");
753        let token = CancellationToken::new();
754        let repo = new_repo_with(
755            &dir,
756            token,
757            Duration::from_secs(60),
758            None,
759            536_870_912,
760            Duration::from_secs(3600),
761        )
762        .await;
763        assert_eq!(repo.cache_size(), 536_870_912);
764    }
765
766    #[tokio::test]
767    async fn sweep_interval_recorded_and_accessible() {
768        let dir = tempdir().expect("tempdir");
769        let token = CancellationToken::new();
770        let repo = new_repo_with(
771            &dir,
772            token,
773            Duration::from_secs(60),
774            None,
775            256 * 1024 * 1024,
776            Duration::from_secs(1800),
777        )
778        .await;
779        assert_eq!(repo.sweep_interval(), Duration::from_secs(1800));
780    }
781
782    #[tokio::test]
783    async fn stale_retention_recorded_and_accessible() {
784        let dir = tempdir().expect("tempdir");
785        let token = CancellationToken::new();
786        let repo = new_repo_with(
787            &dir,
788            token,
789            Duration::from_secs(3600),
790            None,
791            256 * 1024 * 1024,
792            Duration::from_secs(3600),
793        )
794        .await;
795        assert_eq!(repo.stale_retention(), Duration::from_secs(3600));
796    }
797
798    #[tokio::test]
799    async fn explicit_cache_size_round_trip() {
800        let dir = tempdir().expect("tempdir");
801        let token = CancellationToken::new();
802        let repo = new_repo_with(
803            &dir,
804            token,
805            Duration::from_secs(60),
806            None,
807            512 * 1024 * 1024,
808            Duration::from_secs(3600),
809        )
810        .await;
811        repo.set("k", entry(), Some(Duration::from_secs(3600)))
812            .await
813            .expect("set");
814        let found = repo.get("k").await.expect("get");
815        assert!(
816            found.is_some(),
817            "entry must round-trip through the builder-opened database"
818        );
819    }
820
821    #[tokio::test]
822    async fn entries_survive_handle_drop_and_reopen() {
823        let dir = tempdir().expect("tempdir");
824        let path = dir.path().join("cache.redb");
825        let token = CancellationToken::new();
826        // Keep a clone so the sweep can be stopped via the real shutdown path
827        // (token cancellation) before reopening the same file. redb rejects
828        // reopening a file whose in-process handle is still alive, and the
829        // sweep task holds a cloned Arc<Database> until it exits.
830        let token_for_shutdown = token.clone();
831        {
832            let repo = RedbCacheRepository::new(
833                "redb",
834                path.clone(),
835                Duration::from_secs(60),
836                None,
837                256 * 1024 * 1024,
838                Duration::from_secs(3600),
839                token,
840            )
841            .await
842            .expect("open repo");
843            repo.set("k", entry(), Some(Duration::from_secs(3600)))
844                .await
845                .expect("set");
846            assert_eq!(
847                repo.stats().await.entries,
848                1,
849                "entries counter must be 1 after first insert"
850            );
851            // Stop the sweep via its bound token and await the task so its
852            // cloned Arc<Database> releases before the repo drops. Bind the
853            // handle separately so the MutexGuard drops before the `.await`
854            // (clippy::await-holding-lock).
855            token_for_shutdown.cancel();
856            let sweep_handle = repo.sweep_handle.lock().take();
857            if let Some(handle) = sweep_handle {
858                handle
859                    .await
860                    .expect("sweep task must exit cleanly on token cancel");
861            }
862            // repo dropped here — handle already taken; self.db Arc → 0 → DB
863            // closes → redb frees its in-process lock.
864        }
865        // Reopen with a fresh token; entries must be reloaded from table.len().
866        let token2 = CancellationToken::new();
867        let repo = RedbCacheRepository::new(
868            "redb",
869            path,
870            Duration::from_secs(60),
871            None,
872            256 * 1024 * 1024,
873            Duration::from_secs(3600),
874            token2,
875        )
876        .await
877        .expect("reopen repo");
878        let found = repo.get("k").await.expect("get after reopen");
879        assert!(
880            found.is_some(),
881            "persisted entry must survive drop + reopen"
882        );
883        assert_eq!(
884            repo.stats().await.entries,
885            1,
886            "entries counter must be restored from table.len() on reopen"
887        );
888    }
889
890    #[tokio::test]
891    async fn peek_stale_returns_post_expiry_entry_on_redb() {
892        let dir = tempdir().expect("tempdir");
893        let token = CancellationToken::new();
894        let repo = new_repo(&dir, token).await;
895        repo.set("k", entry(), Some(Duration::from_millis(1)))
896            .await
897            .expect("set");
898        tokio::time::sleep(Duration::from_millis(10)).await;
899        let stale = repo.peek_stale("k").await.expect("peek_stale");
900        assert!(
901            stale.is_some(),
902            "peek_stale must return the expired-but-present entry"
903        );
904    }
905
906    #[tokio::test]
907    async fn sweep_once_removes_entries_past_stale_retention() {
908        let dir = tempdir().expect("tempdir");
909        let token = CancellationToken::new();
910        let path = dir.path().join("cache.redb");
911        // stale_retention=10ms; entry expires at +1ms; sleep 50ms so
912        // (expires_at + 10ms) is well in the past → reclaimable.
913        let repo = RedbCacheRepository::new(
914            "redb",
915            path,
916            Duration::from_millis(10),
917            None,
918            256 * 1024 * 1024,
919            Duration::from_secs(3600),
920            token,
921        )
922        .await
923        .expect("open repo");
924        // Neutralize the background sweep so its immediate first tick (which
925        // fires whenever the runtime next polls it) cannot steal the reclaim
926        // from sweep_once. This makes the test deterministic: sweep_once is
927        // the sole reclaimer.
928        if let Some(handle) = repo.sweep_handle.lock().take() {
929            handle.abort();
930        }
931        repo.set("k", entry(), Some(Duration::from_millis(1)))
932            .await
933            .expect("set");
934        tokio::time::sleep(Duration::from_millis(50)).await;
935        let reclaimed = repo.sweep_once().await.expect("sweep_once");
936        assert!(
937            reclaimed >= 1,
938            "sweep_once must reclaim at least 1 entry, got {reclaimed}"
939        );
940        let stale = repo.peek_stale("k").await.expect("peek_stale after sweep");
941        assert!(
942            stale.is_none(),
943            "entry must be gone after sweep_once reclaimed it"
944        );
945    }
946
947    #[tokio::test]
948    async fn sweep_stops_on_context_shutdown() {
949        let dir = tempdir().expect("tempdir");
950        let token = CancellationToken::new();
951        let path = dir.path().join("cache.redb");
952        // Short sweep interval so the loop is definitely armed.
953        let repo = RedbCacheRepository::new(
954            "redb",
955            path,
956            Duration::from_secs(60),
957            None,
958            256 * 1024 * 1024,
959            Duration::from_millis(10),
960            token.clone(),
961        )
962        .await
963        .expect("open repo");
964        token.cancel();
965        // Take the handle out so we can await it directly; Drop will see None.
966        let handle = repo
967            .sweep_handle
968            .lock()
969            .take()
970            .expect("sweep handle must be present after construct");
971        let completed = tokio::time::timeout(Duration::from_secs(5), handle).await;
972        assert!(
973            completed.is_ok(),
974            "sweep task must complete within 5s of context shutdown"
975        );
976    }
977
978    #[tokio::test]
979    async fn redb_errors_surface_as_err() {
980        let dir = tempdir().expect("tempdir");
981        // A regular file where a directory is required — `create_dir_all`
982        // must fail, surfacing as Err(CamelError::Io(_)).
983        let blocker = dir.path().join("blocker");
984        std::fs::write(&blocker, b"not a dir").expect("write blocker");
985        let path = blocker.join("cache.redb");
986        let result = RedbCacheRepository::new(
987            "redb",
988            path,
989            Duration::from_secs(60),
990            None,
991            256 * 1024 * 1024,
992            Duration::from_secs(3600),
993            CancellationToken::new(),
994        )
995        .await;
996        assert!(
997            matches!(result, Err(CamelError::Io(_))),
998            "expected Err(CamelError::Io(_)), got {result:?}"
999        );
1000    }
1001
1002    #[tokio::test]
1003    async fn overwrite_does_not_inflate_entries() {
1004        let dir = tempdir().expect("tempdir");
1005        let token = CancellationToken::new();
1006        let repo = new_repo(&dir, token).await;
1007        repo.set("k", entry(), None).await.expect("first set");
1008        repo.set("k", entry(), None).await.expect("second set");
1009        assert_eq!(
1010            repo.stats().await.entries,
1011            1,
1012            "overwriting an existing key must not inflate the entries counter"
1013        );
1014    }
1015
1016    #[tokio::test]
1017    async fn stats_reports_bytes_sum() {
1018        let dir = tempdir().expect("tempdir");
1019        let token = CancellationToken::new();
1020        let repo = new_repo(&dir, token).await;
1021        let a = CacheEntry {
1022            bytes: vec![1, 2, 3],
1023            content_type: camel_api::cache::ContentType::Bytes,
1024            expires_at: None,
1025        };
1026        let b = CacheEntry {
1027            bytes: vec![1, 2, 3, 4, 5],
1028            content_type: camel_api::cache::ContentType::Bytes,
1029            expires_at: None,
1030        };
1031        repo.set("a", a, None).await.expect("set a");
1032        repo.set("b", b, None).await.expect("set b");
1033        assert_eq!(repo.stats().await.bytes, Some(8));
1034    }
1035
1036    #[tokio::test]
1037    async fn stats_counters_reported_alongside_bytes() {
1038        let dir = tempdir().expect("tempdir");
1039        let token = CancellationToken::new();
1040        let repo = new_repo(&dir, token).await;
1041        let a = CacheEntry {
1042            bytes: vec![1, 2, 3],
1043            content_type: camel_api::cache::ContentType::Bytes,
1044            expires_at: None,
1045        };
1046        repo.set("a", a, None).await.expect("set a");
1047        let s = repo.stats().await;
1048        assert_eq!(s.entries, 1);
1049        assert_eq!(s.bytes, Some(3));
1050    }
1051
1052    #[tokio::test]
1053    async fn max_entries_rejects_new_key_allows_overwrite() {
1054        let dir = tempdir().expect("tempdir");
1055        let token = CancellationToken::new();
1056        let path = dir.path().join("cache.redb");
1057        let repo = RedbCacheRepository::new(
1058            "redb",
1059            path,
1060            Duration::from_secs(60),
1061            Some(2),
1062            256 * 1024 * 1024,
1063            Duration::from_secs(3600),
1064            token,
1065        )
1066        .await
1067        .expect("open repo");
1068        repo.set("a", entry(), None).await.expect("set a");
1069        repo.set("b", entry(), None).await.expect("set b");
1070        let over = repo.set("c", entry(), None).await;
1071        assert!(
1072            over.is_err(),
1073            "third distinct key must be rejected at max_entries, got {over:?}"
1074        );
1075        let overw = repo.set("a", entry(), None).await;
1076        assert!(
1077            overw.is_ok(),
1078            "overwrite of an existing key must succeed at max_entries, got {overw:?}"
1079        );
1080    }
1081
1082    #[tokio::test]
1083    async fn invalidate_prefix_removes_namespace_only() {
1084        let dir = tempdir().expect("tempdir");
1085        let token = CancellationToken::new();
1086        let repo = new_repo(&dir, token).await;
1087        repo.set("rainviewer:a", entry(), None)
1088            .await
1089            .expect("set rainviewer:a");
1090        repo.set("rainviewer:b", entry(), None)
1091            .await
1092            .expect("set rainviewer:b");
1093        repo.set("gibs:a", entry(), None).await.expect("set gibs:a");
1094        let deleted = repo
1095            .invalidate_prefix("rainviewer:")
1096            .await
1097            .expect("invalidate_prefix");
1098        assert_eq!(deleted, 2, "only the rainviewer namespace must be removed");
1099        assert!(
1100            repo.get("rainviewer:a").await.expect("get").is_none(),
1101            "rainviewer:a must be gone"
1102        );
1103        assert!(
1104            repo.get("rainviewer:b").await.expect("get").is_none(),
1105            "rainviewer:b must be gone"
1106        );
1107        assert!(
1108            repo.get("gibs:a").await.expect("get").is_some(),
1109            "gibs:a must survive"
1110        );
1111    }
1112
1113    #[tokio::test]
1114    async fn invalidate_prefix_does_not_delete_successor_key() {
1115        let dir = tempdir().expect("tempdir");
1116        let token = CancellationToken::new();
1117        let repo = new_repo(&dir, token).await;
1118        repo.set("ns:", entry(), None).await.expect("set ns:");
1119        repo.set("ns;", entry(), None).await.expect("set ns;");
1120        let deleted = repo
1121            .invalidate_prefix("ns:")
1122            .await
1123            .expect("invalidate_prefix");
1124        assert_eq!(deleted, 1, "only the ns: key must be removed");
1125        assert!(
1126            repo.get("ns:").await.expect("get").is_none(),
1127            "ns: must be gone"
1128        );
1129        assert!(
1130            repo.get("ns;").await.expect("get").is_some(),
1131            "successor key ns; must survive"
1132        );
1133    }
1134
1135    // ── cgroup memory-limit guardrail ─────────────────────────────────────────
1136
1137    /// Runs `f` under a thread-local default `fmt` subscriber that appends into
1138    /// a shared buffer, then returns the captured text.
1139    fn capture_guardrail(f: impl FnOnce()) -> String {
1140        let buf = Arc::new(Mutex::new(Vec::new()));
1141        let subscriber = tracing_subscriber::fmt::Subscriber::builder()
1142            .with_writer(TestWriter {
1143                buf: Arc::clone(&buf),
1144            })
1145            .with_ansi(false)
1146            .finish();
1147        tracing::subscriber::with_default(subscriber, f);
1148        let captured = buf.lock().clone();
1149        String::from_utf8(captured).expect("captured output must be UTF-8")
1150    }
1151
1152    /// `fmt` writer that appends into a shared `Arc<Mutex<Vec<u8>>>`.
1153    struct TestWriter {
1154        buf: Arc<Mutex<Vec<u8>>>,
1155    }
1156
1157    impl std::io::Write for TestWriter {
1158        fn write(&mut self, data: &[u8]) -> std::io::Result<usize> {
1159            self.buf.lock().extend_from_slice(data);
1160            Ok(data.len())
1161        }
1162
1163        fn flush(&mut self) -> std::io::Result<()> {
1164            Ok(())
1165        }
1166    }
1167
1168    impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for TestWriter {
1169        type Writer = TestWriter;
1170
1171        fn make_writer(&'a self) -> Self::Writer {
1172            TestWriter {
1173                buf: Arc::clone(&self.buf),
1174            }
1175        }
1176    }
1177
1178    #[test]
1179    fn cgroup_v2_limit_parsed() {
1180        let dir = tempdir().expect("tempdir");
1181        let v2 = dir.path().join("memory.max");
1182        std::fs::write(&v2, "805306368\n").expect("write v2");
1183        let missing_v1 = dir.path().join("missing-v1");
1184        assert_eq!(memory_limit_from_paths(&v2, &missing_v1), Some(805_306_368));
1185    }
1186
1187    #[test]
1188    fn cgroup_v2_max_means_unlimited() {
1189        let dir = tempdir().expect("tempdir");
1190        let v2 = dir.path().join("memory.max");
1191        std::fs::write(&v2, "max").expect("write v2");
1192        let missing_v1 = dir.path().join("missing-v1");
1193        assert_eq!(memory_limit_from_paths(&v2, &missing_v1), None);
1194    }
1195
1196    #[test]
1197    fn cgroup_v2_malformed_falls_through() {
1198        let dir = tempdir().expect("tempdir");
1199        let v2 = dir.path().join("memory.max");
1200        std::fs::write(&v2, "not-a-number").expect("write v2");
1201        let v1 = dir.path().join("memory.limit_in_bytes");
1202        std::fs::write(&v1, "1073741824").expect("write v1");
1203        assert_eq!(memory_limit_from_paths(&v2, &v1), Some(1_073_741_824));
1204    }
1205
1206    #[test]
1207    fn cgroup_v1_sentinel_unlimited() {
1208        let dir = tempdir().expect("tempdir");
1209        let missing_v2 = dir.path().join("missing-v2");
1210        let v1 = dir.path().join("memory.limit_in_bytes");
1211        std::fs::write(&v1, "9223372036854771712").expect("write v1");
1212        assert_eq!(memory_limit_from_paths(&missing_v2, &v1), None);
1213    }
1214
1215    #[test]
1216    fn cgroup_v1_exactly_16tib_is_a_limit() {
1217        let dir = tempdir().expect("tempdir");
1218        let missing_v2 = dir.path().join("missing-v2");
1219        let v1 = dir.path().join("memory.limit_in_bytes");
1220        std::fs::write(&v1, "17592186044416\n").expect("write v1");
1221        assert_eq!(
1222            memory_limit_from_paths(&missing_v2, &v1),
1223            Some(17_592_186_044_416)
1224        );
1225    }
1226
1227    #[test]
1228    fn successor_bound_unit_tests() {
1229        assert_eq!(
1230            successor_bound("ns:"),
1231            std::ops::Bound::Excluded("ns;".to_string())
1232        );
1233        // Prefix ending U+D7FF must jump the surrogate gap to U+E000.
1234        assert_eq!(
1235            successor_bound("a\u{D7FF}"),
1236            std::ops::Bound::Excluded("a\u{E000}".to_string())
1237        );
1238        // Prefix ending U+E000 increments to U+E001.
1239        assert_eq!(
1240            successor_bound("a\u{E000}"),
1241            std::ops::Bound::Excluded("a\u{E001}".to_string())
1242        );
1243        // Trailing U+10FFFF carries into the preceding scalar: "a…" → "b".
1244        assert_eq!(
1245            successor_bound("a\u{10FFFF}"),
1246            std::ops::Bound::Excluded("b".to_string())
1247        );
1248        // A prefix of only U+10FFFF has no successor.
1249        assert_eq!(
1250            successor_bound("\u{10FFFF}\u{10FFFF}"),
1251            std::ops::Bound::Unbounded
1252        );
1253    }
1254
1255    #[tokio::test]
1256    async fn invalidate_prefix_empty_prefix_removes_all_seeded() {
1257        let dir = tempdir().expect("tempdir");
1258        let token = CancellationToken::new();
1259        let repo = new_repo(&dir, token).await;
1260        repo.set("ns:a", entry(), None).await.expect("set ns:a");
1261        repo.set("ns:b", entry(), None).await.expect("set ns:b");
1262        repo.set("other:c", entry(), None)
1263            .await
1264            .expect("set other:c");
1265        let deleted = repo.invalidate_prefix("").await.expect("invalidate_prefix");
1266        assert_eq!(deleted, 3, "empty prefix must remove every entry");
1267    }
1268
1269    #[tokio::test]
1270    async fn invalidate_prefix_empty_namespace_returns_zero() {
1271        let dir = tempdir().expect("tempdir");
1272        let token = CancellationToken::new();
1273        let repo = new_repo(&dir, token).await;
1274        let deleted = repo
1275            .invalidate_prefix("ns:")
1276            .await
1277            .expect("invalidate_prefix");
1278        assert_eq!(deleted, 0, "absent namespace must report zero removals");
1279    }
1280
1281    #[test]
1282    fn cgroup_files_missing() {
1283        let dir = tempdir().expect("tempdir");
1284        let missing_v2 = dir.path().join("missing-v2");
1285        let missing_v1 = dir.path().join("missing-v1");
1286        assert_eq!(memory_limit_from_paths(&missing_v2, &missing_v1), None);
1287    }
1288
1289    #[test]
1290    fn guardrail_warns_when_exceeds() {
1291        let dir = tempdir().expect("tempdir");
1292        let v2 = dir.path().join("memory.max");
1293        std::fs::write(&v2, "805306368\n").expect("write v2");
1294        let missing_v1 = dir.path().join("missing-v1");
1295        let output = capture_guardrail(|| {
1296            emit_memory_guardrail(1_073_741_824, &v2, &missing_v1);
1297        });
1298        assert!(output.contains("1073741824"), "output: {output}");
1299        assert!(output.contains("805306368"), "output: {output}");
1300        assert_eq!(
1301            output.matches("exceeds container memory limit").count(),
1302            1,
1303            "warn line must appear exactly once: {output}"
1304        );
1305    }
1306
1307    #[test]
1308    fn guardrail_silent_when_fits() {
1309        let dir = tempdir().expect("tempdir");
1310        let v2 = dir.path().join("memory.max");
1311        std::fs::write(&v2, "805306368\n").expect("write v2");
1312        let missing_v1 = dir.path().join("missing-v1");
1313        let output = capture_guardrail(|| {
1314            emit_memory_guardrail(268_435_456, &v2, &missing_v1);
1315        });
1316        assert!(output.is_empty(), "expected no output, got: {output}");
1317    }
1318
1319    #[test]
1320    fn guardrail_silent_when_files_missing() {
1321        let dir = tempdir().expect("tempdir");
1322        let missing_v2 = dir.path().join("missing-v2");
1323        let missing_v1 = dir.path().join("missing-v1");
1324        let output = capture_guardrail(|| {
1325            emit_memory_guardrail(1_073_741_824, &missing_v2, &missing_v1);
1326        });
1327        assert!(output.is_empty(), "expected no output, got: {output}");
1328    }
1329}