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 payload_path: None,
707 content_type: camel_api::cache::ContentType::Bytes,
708 expires_at: None,
709 }
710 }
711
712 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 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 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 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 }
866 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 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 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 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 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 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 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 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 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 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 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 assert_eq!(
1278 successor_bound("a\u{D7FF}"),
1279 std::ops::Bound::Excluded("a\u{E000}".to_string())
1280 );
1281 assert_eq!(
1283 successor_bound("a\u{E000}"),
1284 std::ops::Bound::Excluded("a\u{E001}".to_string())
1285 );
1286 assert_eq!(
1288 successor_bound("a\u{10FFFF}"),
1289 std::ops::Bound::Excluded("b".to_string())
1290 );
1291 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}