1use 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
39const CACHE_TABLE: TableDefinition<&str, &[u8]> = TableDefinition::new("cache_entries");
44
45pub struct RedbCacheRepository {
54 name: String,
55 db: Arc<redb::Database>,
56 stale_retention: Duration,
57 max_entries: Option<usize>,
58 cache_size: usize,
60 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 entries: Arc<AtomicU64>,
72 shutdown_token: CancellationToken,
76 sweep_handle: Mutex<Option<JoinHandle<()>>>,
77}
78
79impl RedbCacheRepository {
80 #[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 let len = {
111 let wtx = db
112 .begin_write()
113 .map_err(|e| CamelError::Io(format!("redb begin_write: {e}")))?;
114 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 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 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 pub fn cache_size(&self) -> usize {
198 self.cache_size
199 }
200
201 pub fn sweep_interval(&self) -> std::time::Duration {
203 self.sweep_interval
204 }
205
206 pub fn stale_retention(&self) -> std::time::Duration {
208 self.stale_retention
209 }
210
211 #[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
229pub(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 && bytes <= 17_592_186_044_416
250 {
251 return Some(bytes);
252 }
253 None
254}
255
256pub(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
270fn 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 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 Some(threshold) => threshold < now,
299 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
320fn successor_bound(prefix: &str) -> Bound<String> {
328 match prefix.chars().last() {
329 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 None => successor_bound(rest),
342 }
343 }
344 }
345}
346
347fn increment_scalar(c: char) -> Option<char> {
351 match c {
352 '\u{D7FF}' => Some('\u{E000}'),
354 '\u{10FFFF}' => None,
356 _ => char::from_u32(c as u32 + 1),
358 }
359}
360
361#[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 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 let is_new = {
440 let prior = table
441 .get(key.as_str())
442 .map_err(|e| CamelError::Io(format!("redb get: {e}")))?;
443 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 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 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 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 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
657fn 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 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 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 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 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 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 }
865 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 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 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 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 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 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 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 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 assert_eq!(
1235 successor_bound("a\u{D7FF}"),
1236 std::ops::Bound::Excluded("a\u{E000}".to_string())
1237 );
1238 assert_eq!(
1240 successor_bound("a\u{E000}"),
1241 std::ops::Bound::Excluded("a\u{E001}".to_string())
1242 );
1243 assert_eq!(
1245 successor_bound("a\u{10FFFF}"),
1246 std::ops::Bound::Excluded("b".to_string())
1247 );
1248 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}