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