1use std::future::Future;
24use std::pin::Pin;
25use std::sync::Arc;
26use std::time::{Duration, SystemTime};
27
28use bytes::Bytes;
29
30use camel_api::body::Body;
31use camel_api::cache::{CacheEntry, CacheRepository, ContentType};
32use camel_api::{CamelError, Exchange, OutcomePipeline, OutcomeSegment, PipelineOutcome};
33use camel_component_api::RuntimeObservability;
34
35use crate::MessageIdSource;
36
37#[derive(Clone)]
45pub(crate) enum CoalesceTerminal {
46 Completed(Body),
48 Failed(CamelError),
50 Stopped,
52}
53
54struct InFlight {
61 terminal: std::sync::Mutex<Option<CoalesceTerminal>>,
62 notify: tokio::sync::Notify,
63}
64
65impl Default for InFlight {
66 fn default() -> Self {
67 Self {
68 terminal: std::sync::Mutex::new(None),
69 notify: tokio::sync::Notify::new(),
70 }
71 }
72}
73
74impl InFlight {
75 fn publish(&self, terminal: CoalesceTerminal) {
79 if let Ok(mut slot) = self.terminal.lock()
80 && slot.is_none()
81 {
82 *slot = Some(terminal);
83 }
84 }
85
86 fn terminal_snapshot(&self) -> Option<CoalesceTerminal> {
88 match self.terminal.lock() {
89 Ok(slot) => (*slot).clone(),
90 Err(_) => None,
91 }
92 }
93}
94
95type InFlightMap = std::sync::Mutex<std::collections::HashMap<String, std::sync::Arc<InFlight>>>;
98
99struct LeaderGuard {
108 key: String,
109 map: Arc<InFlightMap>,
110 cell: Arc<InFlight>,
111}
112
113impl LeaderGuard {
114 fn retire(map: &Arc<InFlightMap>, key: &str, cell: &Arc<InFlight>) {
120 if let Ok(mut map) = map.lock()
121 && map
122 .get(key)
123 .is_some_and(|current| Arc::ptr_eq(current, cell))
124 {
125 map.remove(key);
126 }
127 }
128}
129
130impl Drop for LeaderGuard {
131 fn drop(&mut self) {
132 self.cell
134 .publish(CoalesceTerminal::Failed(CamelError::Config(
135 "cache coalesce leader cancelled".into(),
136 )));
137 self.cell.notify.notify_waiters();
138 Self::retire(&self.map, &self.key, &self.cell);
139 }
140}
141
142pub struct CacheService {
171 repository: Arc<dyn CacheRepository>,
172 repository_name: String,
174 key_expr: MessageIdSource,
175 ttl: Option<Duration>,
176 max_entry_bytes: usize,
177 on_miss: OutcomeSegment,
178 rt: Arc<dyn RuntimeObservability>,
179 coalesce_misses: bool,
181 inflight: Arc<InFlightMap>,
184}
185
186impl CacheService {
187 pub fn new(
192 repository: Arc<dyn CacheRepository>,
193 key_expr: MessageIdSource,
194 ttl: Option<Duration>,
195 max_entry_bytes: usize,
196 on_miss: OutcomeSegment,
197 rt: Arc<dyn RuntimeObservability>,
198 ) -> Self {
199 let repository_name = repository.name().to_string();
200 Self {
201 repository,
202 repository_name,
203 key_expr,
204 ttl,
205 max_entry_bytes,
206 on_miss,
207 rt,
208 coalesce_misses: false,
209 inflight: Arc::new(InFlightMap::default()),
210 }
211 }
212
213 pub fn with_coalesce(mut self, coalesce_misses: bool) -> Self {
219 self.coalesce_misses = coalesce_misses;
220 self
221 }
222
223 pub fn repository_name(&self) -> &str {
225 &self.repository_name
226 }
227}
228
229#[allow(clippy::too_many_arguments)]
237async fn write_back(
238 repository: &Arc<dyn CacheRepository>,
239 repository_name: &str,
240 max_entry_bytes: usize,
241 ttl: Option<Duration>,
242 exchange: Exchange,
243 key: &str,
244 serialized: Vec<u8>,
245 content_type: ContentType,
246) -> PipelineOutcome {
247 if serialized.len() <= max_entry_bytes {
248 let entry = CacheEntry {
249 bytes: serialized,
250 payload_path: None,
251 content_type,
252 expires_at: None,
253 };
254 match repository.set(key, entry, ttl).await {
255 Ok(()) => {}
256 Err(e) => {
257 if matches!(&e, CamelError::Config(msg) if msg.starts_with("cache: max_entries")) {
258 tracing::debug!(
259 repository = %repository_name,
260 key = %key,
261 "cache at capacity, skipping write-back"
262 ); } else {
264 return PipelineOutcome::Failed(e);
265 }
266 }
267 }
268 } else {
269 tracing::debug!(
271 repository = %repository_name,
272 key = %key,
273 len = serialized.len(),
274 max = max_entry_bytes,
275 "cache write-back skipped: body exceeds max_entry_bytes"
276 );
277 }
278 PipelineOutcome::Completed(exchange)
279}
280
281impl Clone for CacheService {
282 fn clone(&self) -> Self {
283 Self {
284 repository: Arc::clone(&self.repository),
285 repository_name: self.repository_name.clone(),
286 key_expr: self.key_expr.clone(),
287 ttl: self.ttl,
288 max_entry_bytes: self.max_entry_bytes,
289 on_miss: self.on_miss.clone(),
290 rt: Arc::clone(&self.rt),
291 coalesce_misses: self.coalesce_misses,
292 inflight: Arc::clone(&self.inflight),
293 }
294 }
295}
296
297impl OutcomePipeline for CacheService {
298 fn clone_box(&self) -> Box<dyn OutcomePipeline> {
299 Box::new(self.clone())
300 }
301
302 fn run<'a>(
303 &'a mut self,
304 exchange: Exchange,
305 ) -> Pin<Box<dyn Future<Output = PipelineOutcome> + Send + 'a>> {
306 Box::pin(async move {
307 let key = match self.key_expr.message_id(&exchange).await {
310 Ok(Some(k)) => k,
311 Ok(None) => return self.on_miss.run(exchange).await,
312 Err(e) => return PipelineOutcome::Failed(e),
313 };
314
315 match self.repository.get(&key).await {
317 Err(e) => return PipelineOutcome::Failed(e),
318 Ok(Some(entry)) => {
319 self.rt.metrics().record_counter(
322 "camel.cache.hits",
323 1.0_f64,
324 &[("repository", &self.repository_name)],
325 );
326 match reconstruct_body(&entry) {
327 Ok(body) => {
328 let mut exchange = exchange;
329 exchange.input.body = body;
330 return PipelineOutcome::Completed(exchange);
331 }
332 Err(e) => return PipelineOutcome::Failed(e),
333 }
334 }
335 Ok(None) => {
336 self.rt.metrics().record_counter(
339 "camel.cache.misses",
340 1.0_f64,
341 &[("repository", &self.repository_name)],
342 );
343 }
344 }
345
346 if self.coalesce_misses {
348 return self.coalesced_miss(exchange, key).await;
349 }
350 self.run_miss(exchange, key).await
351 })
352 }
353}
354
355impl CacheService {
356 async fn run_miss(&mut self, exchange: Exchange, key: String) -> PipelineOutcome {
361 let mut exchange = match self.on_miss.run(exchange).await {
363 PipelineOutcome::Stopped(ex) => return PipelineOutcome::Stopped(ex),
364 PipelineOutcome::Failed(e) => return PipelineOutcome::Failed(e),
365 PipelineOutcome::Completed(ex) => ex,
366 };
367
368 let body = std::mem::replace(&mut exchange.input.body, Body::Empty);
370 match body {
371 Body::Bytes(b) => {
372 let serialized = b.to_vec();
373 exchange.input.body = Body::Bytes(b);
374 write_back(
375 &self.repository,
376 &self.repository_name,
377 self.max_entry_bytes,
378 self.ttl,
379 exchange,
380 &key,
381 serialized,
382 ContentType::Bytes,
383 )
384 .await
385 }
386 Body::Text(s) => {
387 let serialized = s.as_bytes().to_vec();
388 exchange.input.body = Body::Text(s);
389 write_back(
390 &self.repository,
391 &self.repository_name,
392 self.max_entry_bytes,
393 self.ttl,
394 exchange,
395 &key,
396 serialized,
397 ContentType::Text,
398 )
399 .await
400 }
401 Body::Json(v) => {
402 let serialized = match serde_json::to_vec(&v) {
403 Ok(b) => b,
404 Err(e) => {
405 exchange.input.body = Body::Json(v);
406 return PipelineOutcome::Failed(CamelError::TypeConversionFailed(
407 e.to_string(),
408 ));
409 }
410 };
411 exchange.input.body = Body::Json(v);
412 write_back(
413 &self.repository,
414 &self.repository_name,
415 self.max_entry_bytes,
416 self.ttl,
417 exchange,
418 &key,
419 serialized,
420 ContentType::Json,
421 )
422 .await
423 }
424 Body::Xml(s) => {
425 let serialized = s.as_bytes().to_vec();
426 exchange.input.body = Body::Xml(s);
427 write_back(
428 &self.repository,
429 &self.repository_name,
430 self.max_entry_bytes,
431 self.ttl,
432 exchange,
433 &key,
434 serialized,
435 ContentType::Xml,
436 )
437 .await
438 }
439 Body::Stream(stream_body) => {
440 let materialized = match Body::Stream(stream_body)
442 .into_bytes(self.max_entry_bytes)
443 .await
444 {
445 Ok(b) => b,
446 Err(e) => return PipelineOutcome::Failed(e),
447 };
448 let entry = CacheEntry {
450 bytes: materialized.to_vec(),
451 payload_path: None,
452 content_type: ContentType::Bytes,
453 expires_at: None,
454 };
455 if let Err(e) = self.repository.set(&key, entry, self.ttl).await {
456 if matches!(&e, CamelError::Config(msg) if msg.starts_with("cache: max_entries"))
458 {
459 tracing::debug!(
460 repository = %self.repository_name,
461 key = %key,
462 "cache at capacity, skipping write-back for stream"
463 ); exchange.input.body = Body::Bytes(materialized);
465 return PipelineOutcome::Completed(exchange);
466 }
467 exchange.input.body = Body::Bytes(materialized);
468 return PipelineOutcome::Failed(e);
469 }
470 exchange.input.body = Body::Bytes(materialized);
471 PipelineOutcome::Completed(exchange)
472 }
473 _ => {
474 exchange.input.body = body;
476 PipelineOutcome::Completed(exchange)
477 }
478 }
479 }
480
481 async fn coalesced_miss(&mut self, exchange: Exchange, key: String) -> PipelineOutcome {
503 let inflight = Arc::clone(&self.inflight);
504
505 let mut wave: Option<Arc<InFlight>> = None;
512 let mut registered = None;
513 match inflight.lock() {
514 Ok(mut map) => {
515 match map.get(&key).cloned() {
516 Some(existing) => {
517 wave = Some(existing);
520 if let Some(cell) = wave.as_ref() {
521 let mut notified = Box::pin(cell.notify.notified());
522 notified.as_mut().enable();
523 registered = Some(notified);
524 }
525 }
526 None => {
527 let cell = Arc::new(InFlight::default());
529 map.insert(key.clone(), Arc::clone(&cell));
530 wave = Some(cell);
531 }
532 }
533 }
534 Err(poisoned) => drop(poisoned),
535 }
536 let Some(cell_ref) = wave.as_ref() else {
537 return self.run_miss(exchange, key).await;
541 };
542
543 if let Some(notified) = registered {
544 let terminal = match cell_ref.terminal_snapshot() {
548 Some(t) => t,
549 None => {
550 notified.await;
551 cell_ref.terminal_snapshot().unwrap_or_else(|| {
552 CoalesceTerminal::Failed(CamelError::Config(
553 "cache coalesce waiter woke without a terminal state".into(),
554 ))
555 })
556 }
557 };
558 match terminal {
559 CoalesceTerminal::Completed(body) => {
560 let mut exchange = exchange;
561 exchange.input.body = body;
562 PipelineOutcome::Completed(exchange)
563 }
564 CoalesceTerminal::Failed(e) => PipelineOutcome::Failed(e),
565 CoalesceTerminal::Stopped => PipelineOutcome::Stopped(exchange),
566 }
567 } else {
568 let cell = Arc::clone(cell_ref);
572 let _guard = LeaderGuard {
573 key: key.clone(),
574 map: Arc::clone(&inflight),
575 cell: Arc::clone(&cell),
576 };
577 let outcome = self.run_miss(exchange, key.clone()).await;
578 let terminal = match &outcome {
579 PipelineOutcome::Completed(ex) => {
580 CoalesceTerminal::Completed(ex.input.body.clone())
581 }
582 PipelineOutcome::Failed(e) => CoalesceTerminal::Failed(e.clone()),
583 PipelineOutcome::Stopped(_) => CoalesceTerminal::Stopped,
584 };
585 cell.publish(terminal);
587 cell.notify.notify_waiters();
588 LeaderGuard::retire(&inflight, &key, &cell);
589 outcome
592 }
593 }
594}
595
596fn reconstruct_body(entry: &CacheEntry) -> Result<Body, CamelError> {
601 match entry.content_type {
602 ContentType::Bytes => Ok(Body::Bytes(Bytes::from(entry.bytes.clone()))),
603 ContentType::Text => {
604 let s = String::from_utf8(entry.bytes.clone()).map_err(|e| {
605 CamelError::TypeConversionFailed(format!("cached text is not valid UTF-8: {e}"))
606 })?;
607 Ok(Body::Text(s))
608 }
609 ContentType::Json => {
610 let v = serde_json::from_slice(&entry.bytes).map_err(|e| {
611 CamelError::TypeConversionFailed(format!("cached bytes are not valid JSON: {e}"))
612 })?;
613 Ok(Body::Json(v))
614 }
615 ContentType::Xml => {
616 let s = String::from_utf8(entry.bytes.clone()).map_err(|e| {
617 CamelError::TypeConversionFailed(format!("cached xml is not valid UTF-8: {e}"))
618 })?;
619 Ok(Body::Xml(s))
620 }
621 }
622}
623
624pub const CAMEL_CACHE_INVALIDATED_COUNT: &str = "CamelCacheInvalidatedCount";
634
635#[derive(Clone)]
638pub enum CacheInvalidateTarget {
639 Key(MessageIdSource),
641 Prefix(MessageIdSource),
643}
644
645pub struct CacheInvalidateService {
660 repository: Arc<dyn CacheRepository>,
661 target: CacheInvalidateTarget,
662 rt: Arc<dyn RuntimeObservability>,
663 repository_name: String,
664}
665
666impl CacheInvalidateService {
667 pub fn new(
668 repository: Arc<dyn CacheRepository>,
669 target: CacheInvalidateTarget,
670 rt: Arc<dyn RuntimeObservability>,
671 ) -> Self {
672 let repository_name = repository.name().to_string();
673 Self {
674 repository,
675 target,
676 rt,
677 repository_name,
678 }
679 }
680}
681
682impl Clone for CacheInvalidateService {
683 fn clone(&self) -> Self {
684 Self {
685 repository: Arc::clone(&self.repository),
686 target: self.target.clone(),
687 rt: Arc::clone(&self.rt),
688 repository_name: self.repository_name.clone(),
689 }
690 }
691}
692
693impl OutcomePipeline for CacheInvalidateService {
694 fn clone_box(&self) -> Box<dyn OutcomePipeline> {
695 Box::new(self.clone())
696 }
697
698 fn run<'a>(
699 &'a mut self,
700 exchange: Exchange,
701 ) -> Pin<Box<dyn Future<Output = PipelineOutcome> + Send + 'a>> {
702 Box::pin(async move {
703 match self.target.clone() {
704 CacheInvalidateTarget::Key(key_expr) => {
705 let key = match key_expr.message_id(&exchange).await {
706 Ok(Some(k)) => k,
707 Ok(None) => return PipelineOutcome::Completed(exchange),
708 Err(e) => return PipelineOutcome::Failed(e),
709 };
710 match self.repository.invalidate(&key).await {
711 Err(e) => PipelineOutcome::Failed(e),
712 Ok(()) => {
713 self.rt.metrics().record_counter(
715 "camel.cache.invalidations",
716 1.0_f64,
717 &[("repository", &self.repository_name)],
718 );
719 let mut exchange = exchange;
720 exchange.set_property(
721 CAMEL_CACHE_INVALIDATED_COUNT,
722 serde_json::Value::from(1u64),
723 );
724 PipelineOutcome::Completed(exchange)
725 }
726 }
727 }
728 CacheInvalidateTarget::Prefix(prefix_expr) => {
729 let prefix = match prefix_expr.message_id(&exchange).await {
730 Ok(Some(p)) => p,
731 Ok(None) => return PipelineOutcome::Completed(exchange),
732 Err(e) => return PipelineOutcome::Failed(e),
733 };
734 match self.repository.invalidate_prefix(&prefix).await {
735 Err(e) => PipelineOutcome::Failed(e),
736 Ok(count) => {
737 self.rt.metrics().record_counter(
739 "camel.cache.invalidations",
740 1.0_f64,
741 &[("repository", &self.repository_name)],
742 );
743 let mut exchange = exchange;
744 exchange.set_property(
745 CAMEL_CACHE_INVALIDATED_COUNT,
746 serde_json::Value::from(count),
747 );
748 PipelineOutcome::Completed(exchange)
749 }
750 }
751 }
752 }
753 })
754 }
755}
756
757pub const CAMEL_CACHE_PEEK_HIT: &str = "CamelCachePeekHit";
763pub const CAMEL_CACHE_PEEK_STALE: &str = "CamelCachePeekStale";
765
766#[derive(Debug, Clone, Copy, PartialEq, Eq)]
773pub enum PeekStaleMissPolicy {
774 Stop,
776 Continue,
778}
779
780impl PeekStaleMissPolicy {
781 pub fn parse_on_miss(raw: Option<&str>) -> Result<Self, CamelError> {
786 match raw {
787 None | Some("stop") => Ok(Self::Stop),
788 Some("continue") => Ok(Self::Continue),
789 Some(other) => Err(CamelError::Config(format!(
790 "cache_peek_stale: invalid on_miss '{other}'; must be \"stop\" or \"continue\""
791 ))),
792 }
793 }
794}
795
796pub struct CachePeekStaleService {
815 repository: Arc<dyn CacheRepository>,
816 key_expr: MessageIdSource,
817 miss_policy: PeekStaleMissPolicy,
818 rt: Arc<dyn RuntimeObservability>,
819 repository_name: String,
820}
821
822impl CachePeekStaleService {
823 pub fn new(
824 repository: Arc<dyn CacheRepository>,
825 key_expr: MessageIdSource,
826 miss_policy: PeekStaleMissPolicy,
827 rt: Arc<dyn RuntimeObservability>,
828 ) -> Self {
829 let repository_name = repository.name().to_string();
830 Self {
831 repository,
832 key_expr,
833 miss_policy,
834 rt,
835 repository_name,
836 }
837 }
838}
839
840impl Clone for CachePeekStaleService {
841 fn clone(&self) -> Self {
842 Self {
843 repository: Arc::clone(&self.repository),
844 key_expr: self.key_expr.clone(),
845 miss_policy: self.miss_policy,
846 rt: Arc::clone(&self.rt),
847 repository_name: self.repository_name.clone(),
848 }
849 }
850}
851
852fn set_peek_properties(exchange: &mut Exchange, hit: bool, stale: bool) {
855 exchange.set_property(CAMEL_CACHE_PEEK_HIT, serde_json::Value::Bool(hit));
856 exchange.set_property(CAMEL_CACHE_PEEK_STALE, serde_json::Value::Bool(stale));
857}
858
859impl OutcomePipeline for CachePeekStaleService {
860 fn clone_box(&self) -> Box<dyn OutcomePipeline> {
861 Box::new(self.clone())
862 }
863
864 fn run<'a>(
865 &'a mut self,
866 exchange: Exchange,
867 ) -> Pin<Box<dyn Future<Output = PipelineOutcome> + Send + 'a>> {
868 Box::pin(async move {
869 let key = match self.key_expr.message_id(&exchange).await {
870 Ok(Some(k)) => k,
871 Ok(None) => {
872 tracing::debug!(
873 step = "cache_peek_stale",
874 repository = %self.repository.name(),
875 "key expression resolved to None; stopping branch"
876 );
877 return PipelineOutcome::Stopped(exchange);
878 }
879 Err(e) => return PipelineOutcome::Failed(e),
880 };
881 match self.repository.peek_stale(&key).await {
882 Err(e) => PipelineOutcome::Failed(e),
883 Ok(Some(entry)) => {
884 self.rt.metrics().record_counter(
887 "camel.cache.peek_stale_served",
888 1.0_f64,
889 &[("repository", &self.repository_name)],
890 );
891 let stale = entry
893 .expires_at
894 .map(|t| t <= SystemTime::now())
895 .unwrap_or(false);
896 match reconstruct_body(&entry) {
897 Ok(body) => {
898 let mut exchange = exchange;
899 exchange.input.body = body;
900 set_peek_properties(&mut exchange, true, stale);
901 PipelineOutcome::Completed(exchange)
902 }
903 Err(e) => PipelineOutcome::Failed(e),
904 }
905 }
906 Ok(None) => match self.miss_policy {
907 PeekStaleMissPolicy::Stop => {
908 let mut exchange = exchange;
909 set_peek_properties(&mut exchange, false, false);
910 tracing::debug!(
911 step = "cache_peek_stale",
912 repository = %self.repository.name(),
913 "peek miss; stopping branch per on_miss=stop"
914 );
915 PipelineOutcome::Stopped(exchange)
916 }
917 PeekStaleMissPolicy::Continue => {
918 let mut exchange = exchange;
919 set_peek_properties(&mut exchange, false, false);
920 PipelineOutcome::Completed(exchange)
921 }
922 },
923 }
924 })
925 }
926}
927
928pub struct CacheClearService {
937 repository: Arc<dyn CacheRepository>,
938}
939
940impl CacheClearService {
941 pub fn new(repository: Arc<dyn CacheRepository>) -> Self {
942 Self { repository }
943 }
944}
945
946impl Clone for CacheClearService {
947 fn clone(&self) -> Self {
948 Self {
949 repository: Arc::clone(&self.repository),
950 }
951 }
952}
953
954impl OutcomePipeline for CacheClearService {
955 fn clone_box(&self) -> Box<dyn OutcomePipeline> {
956 Box::new(self.clone())
957 }
958
959 fn run<'a>(
960 &'a mut self,
961 exchange: Exchange,
962 ) -> Pin<Box<dyn Future<Output = PipelineOutcome> + Send + 'a>> {
963 Box::pin(async move {
964 match self.repository.clear().await {
965 Err(e) => PipelineOutcome::Failed(e),
966 Ok(()) => PipelineOutcome::Completed(exchange),
967 }
968 })
969 }
970}
971
972pub struct CacheStatsService {
979 repository: Arc<dyn CacheRepository>,
980 repository_name: String,
981}
982
983impl CacheStatsService {
984 pub fn new(repository: Arc<dyn CacheRepository>) -> Self {
985 let repository_name = repository.name().to_string();
986 Self {
987 repository,
988 repository_name,
989 }
990 }
991}
992
993impl Clone for CacheStatsService {
994 fn clone(&self) -> Self {
995 Self {
996 repository: Arc::clone(&self.repository),
997 repository_name: self.repository_name.clone(),
998 }
999 }
1000}
1001
1002impl OutcomePipeline for CacheStatsService {
1003 fn clone_box(&self) -> Box<dyn OutcomePipeline> {
1004 Box::new(self.clone())
1005 }
1006
1007 fn run<'a>(
1008 &'a mut self,
1009 mut exchange: Exchange,
1010 ) -> Pin<Box<dyn Future<Output = PipelineOutcome> + Send + 'a>> {
1011 Box::pin(async move {
1012 let s = self.repository.stats().await;
1013 exchange.input.body = Body::Json(serde_json::json!({
1014 "repository": self.repository_name,
1015 "hits": s.hits,
1016 "misses": s.misses,
1017 "evictions": s.evictions,
1018 "entries": s.entries,
1019 "peek_stale_served": s.peek_stale_served,
1020 "invalidations": s.invalidations,
1021 "bytes": s.bytes,
1022 }));
1023 PipelineOutcome::Completed(exchange)
1024 })
1025 }
1026}
1027
1028#[cfg(test)]
1033mod test_utils {
1034 use super::*;
1035 use async_trait::async_trait;
1036 use camel_api::cache::CacheStats;
1037 use std::collections::HashMap;
1038 use std::sync::atomic::{AtomicBool, AtomicU32, AtomicU64, Ordering};
1039 use tokio::sync::Mutex;
1040
1041 #[derive(Debug, Default)]
1045 pub struct MockCacheRepository {
1046 name: String,
1047 entries: Arc<Mutex<HashMap<String, CacheEntry>>>,
1048 get_should_fail: Arc<AtomicBool>,
1049 set_should_fail: Arc<AtomicBool>,
1050 set_call_count: Arc<AtomicU32>,
1051 last_set_ttl: Arc<Mutex<Option<Duration>>>,
1052 invalidate_call_count: Arc<AtomicU32>,
1053 last_invalidate_key: Arc<Mutex<Option<String>>>,
1054 clear_call_count: Arc<AtomicU64>,
1055 clear_should_fail: Arc<AtomicBool>,
1056 invalidate_should_fail: Arc<AtomicBool>,
1057 prefix_unsupported: Arc<AtomicBool>,
1058 stats_override: std::sync::Mutex<CacheStats>,
1059 }
1060
1061 impl MockCacheRepository {
1062 pub fn new(name: &str) -> Self {
1063 Self {
1064 name: name.to_string(),
1065 ..Default::default()
1066 }
1067 }
1068
1069 pub fn invalidate_call_count(&self) -> u32 {
1070 self.invalidate_call_count.load(Ordering::SeqCst)
1071 }
1072
1073 pub async fn last_invalidate_key(&self) -> Option<String> {
1074 self.last_invalidate_key.lock().await.clone()
1075 }
1076
1077 pub fn clear_call_count(&self) -> u64 {
1078 self.clear_call_count.load(Ordering::SeqCst)
1079 }
1080
1081 pub fn set_should_fail_clear(&self, v: bool) {
1082 self.clear_should_fail.store(v, Ordering::SeqCst);
1083 }
1084
1085 pub fn set_should_fail_invalidate(&self, v: bool) {
1086 self.invalidate_should_fail.store(v, Ordering::SeqCst);
1087 }
1088
1089 pub fn set_prefix_unsupported(&self, v: bool) {
1092 self.prefix_unsupported.store(v, Ordering::SeqCst);
1093 }
1094
1095 pub fn set_stats(&self, stats: CacheStats) {
1096 *self.stats_override.lock().unwrap() = stats; }
1098
1099 pub async fn seed(&self, key: &str, entry: CacheEntry) {
1101 self.entries.lock().await.insert(key.to_string(), entry);
1102 }
1103
1104 pub fn set_get_should_fail(&self, v: bool) {
1105 self.get_should_fail.store(v, Ordering::SeqCst);
1106 }
1107
1108 pub fn set_set_should_fail(&self, v: bool) {
1109 self.set_should_fail.store(v, Ordering::SeqCst);
1110 }
1111
1112 pub fn set_call_count(&self) -> u32 {
1113 self.set_call_count.load(Ordering::SeqCst)
1114 }
1115
1116 pub async fn last_set_ttl(&self) -> Option<Duration> {
1118 *self.last_set_ttl.lock().await
1119 }
1120
1121 pub async fn stored_entry(&self, key: &str) -> Option<CacheEntry> {
1123 self.entries.lock().await.get(key).cloned()
1124 }
1125 }
1126
1127 #[async_trait]
1128 impl CacheRepository for MockCacheRepository {
1129 fn name(&self) -> &str {
1130 &self.name
1131 }
1132
1133 async fn get(&self, key: &str) -> Result<Option<CacheEntry>, CamelError> {
1134 if self.get_should_fail.load(Ordering::SeqCst) {
1135 return Err(CamelError::ProcessorError("synthetic get failure".into()));
1136 }
1137 let found = self.entries.lock().await.get(key).cloned();
1143 tokio::task::yield_now().await;
1144 Ok(found)
1145 }
1146
1147 async fn set(
1148 &self,
1149 key: &str,
1150 value: CacheEntry,
1151 ttl: Option<Duration>,
1152 ) -> Result<(), CamelError> {
1153 self.set_call_count.fetch_add(1, Ordering::SeqCst);
1154 *self.last_set_ttl.lock().await = ttl;
1155 if self.set_should_fail.load(Ordering::SeqCst) {
1156 return Err(CamelError::ProcessorError("synthetic set failure".into()));
1157 }
1158 self.entries.lock().await.insert(key.to_string(), value);
1159 Ok(())
1160 }
1161
1162 async fn peek_stale(&self, key: &str) -> Result<Option<CacheEntry>, CamelError> {
1163 self.get(key).await
1164 }
1165
1166 async fn invalidate(&self, key: &str) -> Result<(), CamelError> {
1167 self.invalidate_call_count.fetch_add(1, Ordering::SeqCst);
1168 *self.last_invalidate_key.lock().await = Some(key.to_string());
1169 if self.invalidate_should_fail.load(Ordering::SeqCst) {
1170 return Err(CamelError::ProcessorError(
1171 "synthetic invalidate failure".into(),
1172 ));
1173 }
1174 self.entries.lock().await.remove(key);
1175 Ok(())
1176 }
1177
1178 async fn invalidate_prefix(&self, prefix: &str) -> Result<u64, CamelError> {
1179 if self.prefix_unsupported.load(Ordering::SeqCst) {
1180 return Err(CamelError::Config(format!(
1181 "cache backend '{}' does not support invalidate_prefix (no key iteration)",
1182 self.name()
1183 )));
1184 }
1185 let mut entries = self.entries.lock().await;
1186 let keys: Vec<String> = entries
1187 .keys()
1188 .filter(|k| k.starts_with(prefix))
1189 .cloned()
1190 .collect();
1191 let count = keys.len() as u64;
1192 for k in keys {
1193 entries.remove(&k);
1194 }
1195 Ok(count)
1196 }
1197
1198 async fn clear(&self) -> Result<(), CamelError> {
1199 self.clear_call_count.fetch_add(1, Ordering::SeqCst);
1200 if self.clear_should_fail.load(Ordering::SeqCst) {
1201 return Err(CamelError::ProcessorError("synthetic clear failure".into()));
1202 }
1203 self.entries.lock().await.clear();
1204 Ok(())
1205 }
1206
1207 async fn stats(&self) -> CacheStats {
1208 self.stats_override.lock().unwrap().clone() }
1210 }
1211}
1212
1213#[cfg(test)]
1218mod tests {
1219 use super::test_utils::MockCacheRepository;
1220 use super::*;
1221 use camel_api::body::{StreamBody, StreamMetadata};
1222 use camel_api::cache::CacheStats;
1223 use camel_api::metrics::NoOpMetrics;
1224 use camel_api::{Message, Value};
1225 use camel_component_api::health_registry::NoOpHealthCheckRegistry;
1226 use futures::stream;
1227 use std::sync::Mutex;
1228 use std::sync::atomic::{AtomicBool, Ordering};
1229 use std::time::SystemTime;
1230
1231 #[test]
1232 fn parse_on_miss_maps_absent_stop_and_continue() {
1233 assert_eq!(
1234 PeekStaleMissPolicy::parse_on_miss(None).unwrap(),
1235 PeekStaleMissPolicy::Stop
1236 );
1237 assert_eq!(
1238 PeekStaleMissPolicy::parse_on_miss(Some("stop")).unwrap(),
1239 PeekStaleMissPolicy::Stop
1240 );
1241 assert_eq!(
1242 PeekStaleMissPolicy::parse_on_miss(Some("continue")).unwrap(),
1243 PeekStaleMissPolicy::Continue
1244 );
1245 }
1246
1247 #[test]
1248 fn parse_on_miss_rejects_unknown_value_naming_the_step() {
1249 let err = PeekStaleMissPolicy::parse_on_miss(Some("explode")).unwrap_err();
1250 let msg = format!("{err}");
1251 assert!(msg.contains("cache_peek_stale"), "got: {msg}");
1252 assert!(msg.contains("explode"), "got: {msg}");
1253 }
1254
1255 #[derive(Clone)]
1257 struct NoopRt;
1258
1259 impl camel_component_api::HealthCheckRegistry for NoopRt {
1260 fn force_unhealthy_for_route(&self, _: &str, _: &str, _: &str) {}
1261 }
1262
1263 impl RuntimeObservability for NoopRt {
1264 fn metrics(&self) -> Arc<dyn camel_api::metrics::MetricsCollector> {
1265 Arc::new(NoOpMetrics)
1266 }
1267 fn health(&self) -> Arc<dyn camel_component_api::health_registry::HealthCheckRegistry> {
1268 Arc::new(NoOpHealthCheckRegistry)
1269 }
1270 }
1271
1272 fn noop_rt() -> Arc<dyn RuntimeObservability> {
1273 Arc::new(NoopRt)
1274 }
1275
1276 #[derive(Clone)]
1279 enum ScriptedOutcome {
1280 Complete,
1281 Stop,
1282 Fail(CamelError),
1283 }
1284
1285 struct ScriptedOnMiss {
1288 body: Option<Body>,
1289 outcome: ScriptedOutcome,
1290 invoked: Arc<AtomicBool>,
1291 }
1292
1293 impl OutcomePipeline for ScriptedOnMiss {
1294 fn clone_box(&self) -> Box<dyn OutcomePipeline> {
1295 unreachable!("clone_box not used in cache_eip tests")
1297 }
1298
1299 fn run<'a>(
1300 &'a mut self,
1301 mut exchange: Exchange,
1302 ) -> Pin<Box<dyn Future<Output = PipelineOutcome> + Send + 'a>> {
1303 self.invoked.store(true, Ordering::SeqCst);
1304 let body = self.body.take();
1305 let outcome = self.outcome.clone();
1306 Box::pin(async move {
1307 if let Some(b) = body {
1308 exchange.input.body = b;
1309 }
1310 match outcome {
1311 ScriptedOutcome::Complete => PipelineOutcome::Completed(exchange),
1312 ScriptedOutcome::Stop => PipelineOutcome::Stopped(exchange),
1313 ScriptedOutcome::Fail(e) => PipelineOutcome::Failed(e),
1314 }
1315 })
1316 }
1317 }
1318
1319 fn fixed_key() -> MessageIdSource {
1322 MessageIdSource::Sync(Arc::new(|_| Some("cache-key".to_string())))
1323 }
1324
1325 fn none_key() -> MessageIdSource {
1326 MessageIdSource::Sync(Arc::new(|_| None))
1327 }
1328
1329 fn prefix_key() -> MessageIdSource {
1330 MessageIdSource::Sync(Arc::new(|_| Some("ns:".to_string())))
1331 }
1332
1333 fn build_service(
1335 repo: Arc<MockCacheRepository>,
1336 key_expr: MessageIdSource,
1337 max_entry_bytes: usize,
1338 body: Option<Body>,
1339 outcome: ScriptedOutcome,
1340 ttl: Option<Duration>,
1341 rt: Arc<dyn RuntimeObservability>,
1342 ) -> (CacheService, Arc<AtomicBool>) {
1343 let invoked = Arc::new(AtomicBool::new(false));
1344 let on_miss = OutcomeSegment::new(Box::new(ScriptedOnMiss {
1345 body,
1346 outcome,
1347 invoked: invoked.clone(),
1348 }));
1349 let svc = CacheService::new(repo, key_expr, ttl, max_entry_bytes, on_miss, rt);
1350 (svc, invoked)
1351 }
1352
1353 fn exchange() -> Exchange {
1354 let mut ex = Exchange::new(Message::new(""));
1355 ex.input.set_header("ignored", Value::String("v".into()));
1356 ex
1357 }
1358
1359 fn stream_body(data: &'static [u8]) -> Body {
1360 let chunks: Vec<Result<Bytes, CamelError>> = vec![Ok(Bytes::from_static(data))];
1361 let s = stream::iter(chunks);
1362 Body::Stream(StreamBody {
1363 stream: Arc::new(tokio::sync::Mutex::new(Some(Box::pin(s)))),
1364 metadata: StreamMetadata::default(),
1365 })
1366 }
1367
1368 fn stub_error(msg: &str) -> CamelError {
1369 CamelError::ProcessorError(msg.into())
1370 }
1371
1372 #[tokio::test]
1375 async fn cache_hit_short_circuits_on_miss() {
1376 let repo = Arc::new(MockCacheRepository::new("mock"));
1377 repo.seed(
1378 "cache-key",
1379 CacheEntry {
1380 bytes: b"cached-payload".to_vec(),
1381 payload_path: None,
1382 content_type: ContentType::Bytes,
1383 expires_at: None,
1384 },
1385 )
1386 .await;
1387 let (mut svc, on_miss_invoked) = build_service(
1388 repo,
1389 fixed_key(),
1390 1024,
1391 Some(Body::Bytes(Bytes::from_static(b"unreached"))),
1392 ScriptedOutcome::Complete,
1393 None,
1394 noop_rt(),
1395 );
1396
1397 let outcome = svc.run(exchange()).await;
1398
1399 let ex = match outcome {
1400 PipelineOutcome::Completed(ex) => ex,
1401 other => panic!("expected Completed, got {other:?}"),
1402 };
1403 assert_eq!(
1404 ex.input.body,
1405 Body::Bytes(Bytes::from_static(b"cached-payload"))
1406 );
1407 assert!(
1408 !on_miss_invoked.load(Ordering::SeqCst),
1409 "on_miss must NOT run on a cache HIT"
1410 );
1411 }
1412
1413 #[tokio::test]
1416 async fn cache_miss_runs_on_miss_sets_continues() {
1417 let ttl = Duration::from_secs(30);
1418 let repo = Arc::new(MockCacheRepository::new("mock"));
1419 let (mut svc, on_miss_invoked) = build_service(
1420 repo.clone(),
1421 fixed_key(),
1422 1024,
1423 Some(Body::Bytes(Bytes::from_static(b"x"))),
1424 ScriptedOutcome::Complete,
1425 Some(ttl),
1426 noop_rt(),
1427 );
1428
1429 let outcome = svc.run(exchange()).await;
1430
1431 let ex = match outcome {
1432 PipelineOutcome::Completed(ex) => ex,
1433 other => panic!("expected Completed, got {other:?}"),
1434 };
1435 assert!(on_miss_invoked.load(Ordering::SeqCst));
1436 assert_eq!(ex.input.body, Body::Bytes(Bytes::from_static(b"x")));
1437 assert_eq!(repo.set_call_count(), 1, "set must be called once on miss");
1438 let stored = repo
1439 .stored_entry("cache-key")
1440 .await
1441 .expect("entry must be stored");
1442 assert_eq!(stored.bytes, b"x");
1443 assert_eq!(stored.content_type, ContentType::Bytes);
1444 assert_eq!(repo.last_set_ttl().await, Some(ttl));
1445 }
1446
1447 #[tokio::test]
1450 async fn cache_miss_oversized_materialized_body_skips_writeback() {
1451 let repo = Arc::new(MockCacheRepository::new("mock"));
1452 let (mut svc, _invoked) = build_service(
1454 repo.clone(),
1455 fixed_key(),
1456 4,
1457 Some(Body::Bytes(Bytes::from_static(b"oversized"))),
1458 ScriptedOutcome::Complete,
1459 None,
1460 noop_rt(),
1461 );
1462
1463 let outcome = svc.run(exchange()).await;
1464
1465 let ex = match outcome {
1466 PipelineOutcome::Completed(ex) => ex,
1467 other => panic!("expected Completed, got {other:?}"),
1468 };
1469 assert_eq!(ex.input.body, Body::Bytes(Bytes::from_static(b"oversized")));
1471 assert_eq!(
1472 repo.set_call_count(),
1473 0,
1474 "set must NOT be called for oversized body"
1475 );
1476 assert!(repo.stored_entry("cache-key").await.is_none());
1477 }
1478
1479 #[tokio::test]
1482 async fn cache_miss_oversized_stream_propagates_err() {
1483 let repo = Arc::new(MockCacheRepository::new("mock"));
1484 let (mut svc, _invoked) = build_service(
1485 repo.clone(),
1486 fixed_key(),
1487 4,
1488 Some(stream_body(b"way-too-big-stream")),
1489 ScriptedOutcome::Complete,
1490 None,
1491 noop_rt(),
1492 );
1493
1494 let outcome = svc.run(exchange()).await;
1495
1496 match outcome {
1497 PipelineOutcome::Failed(CamelError::StreamLimitExceeded(n)) => {
1498 assert_eq!(n, 4);
1499 }
1500 other => panic!("expected Failed(StreamLimitExceeded(4)), got {other:?}"),
1501 }
1502 assert_eq!(
1503 repo.set_call_count(),
1504 0,
1505 "set must NOT be called when stream exceeds limit"
1506 );
1507 }
1508
1509 #[tokio::test]
1512 async fn cache_on_miss_stopped_propagates_without_writeback() {
1513 let repo = Arc::new(MockCacheRepository::new("mock"));
1514 let (mut svc, _invoked) = build_service(
1515 repo.clone(),
1516 fixed_key(),
1517 1024,
1518 None,
1519 ScriptedOutcome::Stop,
1520 None,
1521 noop_rt(),
1522 );
1523
1524 let outcome = svc.run(exchange()).await;
1525
1526 assert!(
1527 matches!(outcome, PipelineOutcome::Stopped(_)),
1528 "Stopped from on_miss MUST propagate as Stopped"
1529 );
1530 assert_eq!(
1531 repo.set_call_count(),
1532 0,
1533 "set must NOT be called when on_miss Stops"
1534 );
1535 }
1536
1537 #[tokio::test]
1540 async fn cache_on_miss_err_propagates_without_writeback() {
1541 let repo = Arc::new(MockCacheRepository::new("mock"));
1542 let (mut svc, _invoked) = build_service(
1543 repo.clone(),
1544 fixed_key(),
1545 1024,
1546 None,
1547 ScriptedOutcome::Fail(stub_error("on-miss blew up")),
1548 None,
1549 noop_rt(),
1550 );
1551
1552 let outcome = svc.run(exchange()).await;
1553
1554 match outcome {
1555 PipelineOutcome::Failed(e) => {
1556 assert!(e.to_string().contains("on-miss blew up"), "got: {e}");
1557 }
1558 other => panic!("expected Failed, got {other:?}"),
1559 }
1560 assert_eq!(
1561 repo.set_call_count(),
1562 0,
1563 "set must NOT be called when on_miss fails"
1564 );
1565 }
1566
1567 #[tokio::test]
1570 async fn cache_repository_get_err_propagates() {
1571 let repo = Arc::new(MockCacheRepository::new("mock"));
1572 repo.set_get_should_fail(true);
1573 let (mut svc, on_miss_invoked) = build_service(
1574 repo,
1575 fixed_key(),
1576 1024,
1577 Some(Body::Bytes(Bytes::from_static(b"x"))),
1578 ScriptedOutcome::Complete,
1579 None,
1580 noop_rt(),
1581 );
1582
1583 let outcome = svc.run(exchange()).await;
1584
1585 match outcome {
1586 PipelineOutcome::Failed(e) => {
1587 assert!(e.to_string().contains("synthetic get failure"), "got: {e}");
1588 }
1589 other => panic!("expected Failed, got {other:?}"),
1590 }
1591 assert!(
1592 !on_miss_invoked.load(Ordering::SeqCst),
1593 "on_miss must NOT run when get fails"
1594 );
1595 }
1596
1597 #[tokio::test]
1600 async fn cache_repository_set_err_propagates() {
1601 let repo = Arc::new(MockCacheRepository::new("mock"));
1602 repo.set_set_should_fail(true);
1603 let (mut svc, _invoked) = build_service(
1604 repo.clone(),
1605 fixed_key(),
1606 1024,
1607 Some(Body::Bytes(Bytes::from_static(b"x"))),
1608 ScriptedOutcome::Complete,
1609 None,
1610 noop_rt(),
1611 );
1612
1613 let outcome = svc.run(exchange()).await;
1614
1615 match outcome {
1616 PipelineOutcome::Failed(e) => {
1617 assert!(e.to_string().contains("synthetic set failure"), "got: {e}");
1618 }
1619 other => panic!("expected Failed, got {other:?}"),
1620 }
1621 assert_eq!(repo.set_call_count(), 1, "set was attempted (and failed)");
1622 }
1623
1624 #[tokio::test]
1627 async fn cache_none_key_bypasses_to_on_miss() {
1628 let repo = Arc::new(MockCacheRepository::new("mock"));
1629 let (mut svc, on_miss_invoked) = build_service(
1630 repo.clone(),
1631 none_key(),
1632 1024,
1633 Some(Body::Bytes(Bytes::from_static(b"x"))),
1634 ScriptedOutcome::Complete,
1635 None,
1636 noop_rt(),
1637 );
1638
1639 let outcome = svc.run(exchange()).await;
1640
1641 let ex = match outcome {
1642 PipelineOutcome::Completed(ex) => ex,
1643 other => panic!("expected Completed, got {other:?}"),
1644 };
1645 assert!(
1646 on_miss_invoked.load(Ordering::SeqCst),
1647 "on_miss MUST run when key is None"
1648 );
1649 assert_eq!(ex.input.body, Body::Bytes(Bytes::from_static(b"x")));
1650 assert_eq!(
1651 repo.set_call_count(),
1652 0,
1653 "set must NOT be called when key_expr returns None"
1654 );
1655 }
1656
1657 #[tokio::test]
1660 async fn cache_content_type_reconstruction() {
1661 async fn run_case(entry: CacheEntry, expected: Body) {
1662 let repo = Arc::new(MockCacheRepository::new("mock"));
1663 repo.seed("cache-key", entry).await;
1664 let (mut svc, on_miss_invoked) = build_service(
1665 repo,
1666 fixed_key(),
1667 1024,
1668 Some(Body::Bytes(Bytes::from_static(b"unreached"))),
1669 ScriptedOutcome::Complete,
1670 None,
1671 noop_rt(),
1672 );
1673 let outcome = svc.run(exchange()).await;
1674 let ex = match outcome {
1675 PipelineOutcome::Completed(ex) => ex,
1676 other => panic!("expected Completed, got {other:?}"),
1677 };
1678 assert_eq!(ex.input.body, expected);
1679 assert!(!on_miss_invoked.load(Ordering::SeqCst));
1680 }
1681
1682 run_case(
1683 CacheEntry {
1684 bytes: b"raw".to_vec(),
1685 payload_path: None,
1686 content_type: ContentType::Bytes,
1687 expires_at: None,
1688 },
1689 Body::Bytes(Bytes::from_static(b"raw")),
1690 )
1691 .await;
1692 run_case(
1693 CacheEntry {
1694 bytes: b"hi".to_vec(),
1695 payload_path: None,
1696 content_type: ContentType::Text,
1697 expires_at: None,
1698 },
1699 Body::Text("hi".into()),
1700 )
1701 .await;
1702 run_case(
1703 CacheEntry {
1704 bytes: br#"{"k":1}"#.to_vec(),
1705 payload_path: None,
1706 content_type: ContentType::Json,
1707 expires_at: None,
1708 },
1709 Body::Json(serde_json::json!({"k": 1})),
1710 )
1711 .await;
1712 run_case(
1713 CacheEntry {
1714 bytes: b"<a/>".to_vec(),
1715 payload_path: None,
1716 content_type: ContentType::Xml,
1717 expires_at: None,
1718 },
1719 Body::Xml("<a/>".into()),
1720 )
1721 .await;
1722 }
1723
1724 #[tokio::test]
1727 async fn cache_miss_stream_body_is_materialized_and_cached() {
1728 let repo = Arc::new(MockCacheRepository::new("mock"));
1729 let (mut svc, _invoked) = build_service(
1730 repo.clone(),
1731 fixed_key(),
1732 1024,
1733 Some(stream_body(b"chunky")),
1734 ScriptedOutcome::Complete,
1735 None,
1736 noop_rt(),
1737 );
1738
1739 let outcome = svc.run(exchange()).await;
1740
1741 let ex = match outcome {
1742 PipelineOutcome::Completed(ex) => ex,
1743 other => panic!("expected Completed, got {other:?}"),
1744 };
1745 assert_eq!(ex.input.body, Body::Bytes(Bytes::from_static(b"chunky")));
1747 assert_eq!(repo.set_call_count(), 1);
1748 let stored = repo.stored_entry("cache-key").await.expect("stored");
1749 assert_eq!(stored.bytes, b"chunky");
1750 assert_eq!(stored.content_type, ContentType::Bytes);
1751 }
1752
1753 #[tokio::test]
1756 async fn cache_peek_stale_serves_post_expiry_entry() {
1757 let repo = Arc::new(MockCacheRepository::new("mock"));
1758 repo.seed(
1759 "cache-key",
1760 CacheEntry {
1761 bytes: b"stale-payload".to_vec(),
1762 payload_path: None,
1763 content_type: ContentType::Text,
1764 expires_at: None,
1765 },
1766 )
1767 .await;
1768 let mut svc =
1769 CachePeekStaleService::new(repo, fixed_key(), PeekStaleMissPolicy::Stop, noop_rt());
1770
1771 let outcome = svc.run(exchange()).await;
1772
1773 let ex = match outcome {
1774 PipelineOutcome::Completed(ex) => ex,
1775 other => panic!("expected Completed, got {other:?}"),
1776 };
1777 assert_eq!(ex.input.body, Body::Text("stale-payload".into()));
1778 }
1779
1780 #[tokio::test]
1781 async fn cache_peek_stale_on_absence_stops_branch() {
1782 let repo = Arc::new(MockCacheRepository::new("mock"));
1783 let mut svc =
1784 CachePeekStaleService::new(repo, fixed_key(), PeekStaleMissPolicy::Stop, noop_rt());
1785
1786 let outcome = svc.run(exchange()).await;
1787
1788 assert!(
1789 matches!(outcome, PipelineOutcome::Stopped(_)),
1790 "expected Stopped when no stale entry, got {outcome:?}"
1791 );
1792 }
1793
1794 #[tokio::test]
1795 async fn cache_peek_stale_none_key_stops() {
1796 let repo = Arc::new(MockCacheRepository::new("mock"));
1797 let mut svc =
1798 CachePeekStaleService::new(repo, none_key(), PeekStaleMissPolicy::Stop, noop_rt());
1799
1800 let outcome = svc.run(exchange()).await;
1801
1802 assert!(
1803 matches!(outcome, PipelineOutcome::Stopped(_)),
1804 "expected Stopped when key_expr returns None, got {outcome:?}"
1805 );
1806 }
1807
1808 #[tokio::test]
1809 async fn peek_stale_miss_stop_sets_properties_and_stops() {
1810 let repo = Arc::new(MockCacheRepository::new("mock"));
1811 let mut svc =
1812 CachePeekStaleService::new(repo, fixed_key(), PeekStaleMissPolicy::Stop, noop_rt());
1813
1814 let outcome = svc.run(exchange()).await;
1815
1816 let ex = match outcome {
1817 PipelineOutcome::Stopped(ex) => ex,
1818 other => panic!("expected Stopped, got {other:?}"),
1819 };
1820 assert_eq!(ex.property(CAMEL_CACHE_PEEK_HIT), Some(&Value::Bool(false)));
1821 assert_eq!(
1822 ex.property(CAMEL_CACHE_PEEK_STALE),
1823 Some(&Value::Bool(false))
1824 );
1825 }
1826
1827 #[tokio::test]
1828 async fn peek_stale_miss_continue_completes_with_body_untouched() {
1829 let repo = Arc::new(MockCacheRepository::new("mock"));
1830 let mut svc =
1831 CachePeekStaleService::new(repo, fixed_key(), PeekStaleMissPolicy::Continue, noop_rt());
1832
1833 let mut ex = exchange();
1834 ex.input.body = Body::Text("orig".into());
1835
1836 let outcome = svc.run(ex).await;
1837
1838 let ex = match outcome {
1839 PipelineOutcome::Completed(ex) => ex,
1840 other => panic!("expected Completed, got {other:?}"),
1841 };
1842 assert_eq!(ex.input.body, Body::Text("orig".into()));
1843 assert_eq!(ex.property(CAMEL_CACHE_PEEK_HIT), Some(&Value::Bool(false)));
1844 assert_eq!(
1845 ex.property(CAMEL_CACHE_PEEK_STALE),
1846 Some(&Value::Bool(false))
1847 );
1848 }
1849
1850 #[tokio::test]
1851 async fn peek_stale_hit_sets_hit_and_stale_properties() {
1852 let repo = Arc::new(MockCacheRepository::new("mock"));
1853 repo.seed(
1854 "cache-key",
1855 CacheEntry {
1856 bytes: b"stale-payload".to_vec(),
1857 payload_path: None,
1858 content_type: ContentType::Bytes,
1859 expires_at: Some(SystemTime::now() - Duration::from_millis(1)),
1860 },
1861 )
1862 .await;
1863 let mut svc =
1864 CachePeekStaleService::new(repo, fixed_key(), PeekStaleMissPolicy::Stop, noop_rt());
1865
1866 let outcome = svc.run(exchange()).await;
1867
1868 let ex = match outcome {
1869 PipelineOutcome::Completed(ex) => ex,
1870 other => panic!("expected Completed, got {other:?}"),
1871 };
1872 assert_eq!(
1873 ex.input.body,
1874 Body::Bytes(Bytes::from_static(b"stale-payload"))
1875 );
1876 assert_eq!(ex.property(CAMEL_CACHE_PEEK_HIT), Some(&Value::Bool(true)));
1877 assert_eq!(
1878 ex.property(CAMEL_CACHE_PEEK_STALE),
1879 Some(&Value::Bool(true))
1880 );
1881 }
1882
1883 #[tokio::test]
1884 async fn peek_stale_hit_fresh_sets_stale_false() {
1885 let repo = Arc::new(MockCacheRepository::new("mock"));
1886 repo.seed(
1887 "cache-key",
1888 CacheEntry {
1889 bytes: b"fresh-payload".to_vec(),
1890 payload_path: None,
1891 content_type: ContentType::Bytes,
1892 expires_at: Some(SystemTime::now() + Duration::from_secs(3600)),
1893 },
1894 )
1895 .await;
1896 let mut svc =
1897 CachePeekStaleService::new(repo, fixed_key(), PeekStaleMissPolicy::Stop, noop_rt());
1898
1899 let outcome = svc.run(exchange()).await;
1900
1901 let ex = match outcome {
1902 PipelineOutcome::Completed(ex) => ex,
1903 other => panic!("expected Completed, got {other:?}"),
1904 };
1905 assert_eq!(ex.property(CAMEL_CACHE_PEEK_HIT), Some(&Value::Bool(true)));
1906 assert_eq!(
1907 ex.property(CAMEL_CACHE_PEEK_STALE),
1908 Some(&Value::Bool(false))
1909 );
1910 }
1911
1912 #[derive(Default)]
1922 struct EventRecorder {
1923 records: Arc<Mutex<Vec<String>>>,
1924 }
1925
1926 impl EventRecorder {
1927 fn install(self) -> (Arc<Mutex<Vec<String>>>, tracing::subscriber::DefaultGuard) {
1930 use tracing_subscriber::prelude::*;
1931 static INIT: std::sync::OnceLock<()> = std::sync::OnceLock::new();
1937 if INIT.set(()).is_ok() {
1938 let _ = tracing::subscriber::set_global_default(tracing_subscriber::registry());
1939 }
1940 let records = Arc::clone(&self.records);
1941 let guard = tracing_subscriber::registry().with(self).set_default();
1942 (records, guard)
1943 }
1944 }
1945
1946 struct FieldFmt<'a>(&'a mut String);
1949
1950 impl tracing::field::Visit for FieldFmt<'_> {
1951 fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) {
1952 use std::fmt::Write as _;
1953 let _ = write!(self.0, " {}={:?}", field.name(), value);
1954 }
1955 fn record_str(&mut self, field: &tracing::field::Field, value: &str) {
1956 self.record_debug(field, &value);
1957 }
1958 }
1959
1960 impl<C> tracing_subscriber::Layer<C> for EventRecorder
1961 where
1962 C: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>,
1963 {
1964 fn on_event(
1965 &self,
1966 event: &tracing::Event<'_>,
1967 _ctx: tracing_subscriber::layer::Context<'_, C>,
1968 ) {
1969 let meta = event.metadata();
1970 if meta.target() != "camel_processor::cache_eip" {
1971 return;
1972 }
1973 let mut line = format!("{} ", meta.level());
1974 event.record(&mut FieldFmt(&mut line));
1975 self.records.lock().unwrap().push(line); }
1977 }
1978
1979 #[tokio::test]
1980 async fn peek_stale_miss_stop_emits_debug_log() {
1981 let repo = Arc::new(MockCacheRepository::new("mock"));
1982 let mut svc =
1983 CachePeekStaleService::new(repo, fixed_key(), PeekStaleMissPolicy::Stop, noop_rt());
1984
1985 let (records, _guard) = EventRecorder::default().install();
1986 let outcome = svc.run(exchange()).await;
1987 drop(_guard);
1988
1989 assert!(matches!(outcome, PipelineOutcome::Stopped(_)));
1990
1991 let captured = records.lock().unwrap().join("\n"); let miss_records: Vec<&str> = captured
1993 .lines()
1994 .filter(|l| l.contains("peek miss"))
1995 .collect();
1996 assert_eq!(
1997 miss_records.len(),
1998 1,
1999 "expected exactly one DEBUG record containing \"peek miss\"; got: {captured}"
2000 );
2001 assert!(
2002 miss_records[0].contains("DEBUG"),
2003 "expected DEBUG level record; got: {captured}"
2004 );
2005 assert!(
2006 miss_records[0].contains("repository=mock"),
2007 "expected repository field in record; got: {captured}"
2008 );
2009 assert!(
2010 miss_records[0].contains("step=\"cache_peek_stale\""),
2011 "expected step field in record; got: {captured}"
2012 );
2013 }
2014
2015 #[tokio::test]
2016 async fn peek_stale_key_none_stops_with_debug_log() {
2017 let repo = Arc::new(MockCacheRepository::new("mock"));
2018 let mut svc =
2019 CachePeekStaleService::new(repo, none_key(), PeekStaleMissPolicy::Stop, noop_rt());
2020
2021 let (records, _guard) = EventRecorder::default().install();
2022 let outcome = svc.run(exchange()).await;
2023 drop(_guard);
2024
2025 assert!(matches!(outcome, PipelineOutcome::Stopped(_)));
2026
2027 let captured = records.lock().unwrap().join("\n"); let none_records: Vec<&str> = captured
2029 .lines()
2030 .filter(|l| l.contains("resolved to None"))
2031 .collect();
2032 assert_eq!(
2033 none_records.len(),
2034 1,
2035 "expected exactly one DEBUG record containing \"resolved to None\"; got: {captured}"
2036 );
2037 assert!(
2038 none_records[0].contains("DEBUG"),
2039 "expected DEBUG level record; got: {captured}"
2040 );
2041 assert!(
2042 none_records[0].contains("repository=mock"),
2043 "expected repository field in record; got: {captured}"
2044 );
2045 assert!(
2046 none_records[0].contains("step=\"cache_peek_stale\""),
2047 "expected step field in record; got: {captured}"
2048 );
2049 }
2050
2051 #[tokio::test]
2054 async fn cache_invalidate_calls_repository_invalidate() {
2055 let repo = Arc::new(MockCacheRepository::new("mock"));
2056 repo.seed(
2057 "cache-key",
2058 CacheEntry {
2059 bytes: b"to-go".to_vec(),
2060 payload_path: None,
2061 content_type: ContentType::Bytes,
2062 expires_at: None,
2063 },
2064 )
2065 .await;
2066 let mut svc = CacheInvalidateService::new(
2067 repo.clone(),
2068 CacheInvalidateTarget::Key(fixed_key()),
2069 noop_rt(),
2070 );
2071
2072 let outcome = svc.run(exchange()).await;
2073
2074 let ex = match outcome {
2075 PipelineOutcome::Completed(ex) => ex,
2076 other => panic!("expected Completed, got {other:?}"),
2077 };
2078 assert_eq!(
2079 ex.property(CAMEL_CACHE_INVALIDATED_COUNT),
2080 Some(&serde_json::Value::from(1u64)),
2081 "exact-key success must set CamelCacheInvalidatedCount = 1"
2082 );
2083 assert_eq!(
2084 repo.invalidate_call_count(),
2085 1,
2086 "invalidate must be called once"
2087 );
2088 assert_eq!(
2089 repo.last_invalidate_key().await,
2090 Some("cache-key".to_string()),
2091 "invalidate must be called with the correct key"
2092 );
2093 assert!(
2094 repo.stored_entry("cache-key").await.is_none(),
2095 "entry must be removed after invalidation"
2096 );
2097 }
2098
2099 #[tokio::test]
2100 async fn cache_invalidate_none_key_completes() {
2101 let repo = Arc::new(MockCacheRepository::new("mock"));
2102 let mut svc = CacheInvalidateService::new(
2103 repo.clone(),
2104 CacheInvalidateTarget::Key(none_key()),
2105 noop_rt(),
2106 );
2107
2108 let outcome = svc.run(exchange()).await;
2109
2110 let _ex = match outcome {
2111 PipelineOutcome::Completed(ex) => ex,
2112 other => panic!("expected Completed, got {other:?}"),
2113 };
2114 assert_eq!(
2115 repo.invalidate_call_count(),
2116 0,
2117 "invalidate must NOT be called when key_expr returns None"
2118 );
2119 }
2120
2121 type CounterRecording = Vec<(String, f64, Vec<(String, String)>)>;
2125
2126 #[derive(Clone)]
2127 struct RecordingMetricsCollector {
2128 counters: Arc<Mutex<CounterRecording>>,
2129 }
2130
2131 impl RecordingMetricsCollector {
2132 fn new() -> Self {
2133 Self {
2134 counters: Arc::new(Mutex::new(Vec::new())),
2135 }
2136 }
2137 }
2138
2139 impl camel_api::metrics::MetricsCollector for RecordingMetricsCollector {
2140 fn record_exchange_duration(&self, _route_id: &str, _duration: Duration) {}
2141 fn increment_errors(&self, _route_id: &str, _error_type: &str) {}
2142 fn increment_exchanges(&self, _route_id: &str) {}
2143 fn set_queue_depth(&self, _route_id: &str, _depth: usize) {}
2144 fn record_circuit_breaker_change(&self, _route_id: &str, _from: &str, _to: &str) {}
2145 fn record_counter(&self, name: &str, value: f64, labels: &[(&str, &str)]) {
2146 self.counters.lock().unwrap().push((
2147 name.to_string(),
2148 value,
2149 labels
2150 .iter()
2151 .map(|(k, v)| (k.to_string(), v.to_string()))
2152 .collect(),
2153 ));
2154 }
2155 }
2156
2157 #[derive(Clone)]
2158 struct TestOtelmRt {
2159 collector: Arc<RecordingMetricsCollector>,
2160 }
2161
2162 impl camel_component_api::health_registry::HealthCheckRegistry for TestOtelmRt {
2163 fn force_unhealthy_for_route(&self, _: &str, _: &str, _: &str) {}
2164 }
2165
2166 impl RuntimeObservability for TestOtelmRt {
2167 fn metrics(&self) -> Arc<dyn camel_api::metrics::MetricsCollector> {
2168 self.collector.clone()
2169 }
2170 fn health(&self) -> Arc<dyn camel_component_api::health_registry::HealthCheckRegistry> {
2171 Arc::new(NoOpHealthCheckRegistry)
2172 }
2173 }
2174
2175 #[tokio::test]
2176 async fn cache_step_hit_increments_otel_counter() {
2177 let repo = Arc::new(MockCacheRepository::new("mock"));
2178 repo.seed(
2179 "cache-key",
2180 CacheEntry {
2181 bytes: b"cached".to_vec(),
2182 payload_path: None,
2183 content_type: ContentType::Bytes,
2184 expires_at: None,
2185 },
2186 )
2187 .await;
2188 let collector = RecordingMetricsCollector::new();
2189 let counters = collector.counters.clone();
2190 let rt = Arc::new(TestOtelmRt {
2191 collector: Arc::new(collector),
2192 });
2193 let (mut svc, _invoked) = build_service(
2194 repo,
2195 fixed_key(),
2196 1024,
2197 None,
2198 ScriptedOutcome::Complete,
2199 None,
2200 rt,
2201 );
2202
2203 let outcome = svc.run(exchange()).await;
2204 assert!(matches!(outcome, PipelineOutcome::Completed(_)));
2205
2206 let recorded = counters.lock().unwrap().clone();
2207 assert!(
2208 recorded.contains(&(
2209 "camel.cache.hits".to_string(),
2210 1.0,
2211 vec![("repository".to_string(), "mock".to_string())]
2212 )),
2213 "expected camel.cache.hits counter, got: {recorded:?}"
2214 );
2215 }
2216
2217 #[tokio::test]
2218 async fn cache_step_miss_increments_otel_counter() {
2219 let repo = Arc::new(MockCacheRepository::new("mock"));
2220 let collector = RecordingMetricsCollector::new();
2221 let counters = collector.counters.clone();
2222 let rt = Arc::new(TestOtelmRt {
2223 collector: Arc::new(collector),
2224 });
2225 let (mut svc, _invoked) = build_service(
2226 repo.clone(),
2227 fixed_key(),
2228 1024,
2229 Some(Body::Bytes(Bytes::from_static(b"x"))),
2230 ScriptedOutcome::Complete,
2231 None,
2232 rt,
2233 );
2234
2235 let outcome = svc.run(exchange()).await;
2236 assert!(matches!(outcome, PipelineOutcome::Completed(_)));
2237
2238 let recorded = counters.lock().unwrap().clone();
2239 assert!(
2240 recorded.contains(&(
2241 "camel.cache.misses".to_string(),
2242 1.0,
2243 vec![("repository".to_string(), "mock".to_string())]
2244 )),
2245 "expected camel.cache.misses counter, got: {recorded:?}"
2246 );
2247 }
2248
2249 #[tokio::test]
2252 async fn cache_clear_calls_repository_clear() {
2253 let repo = Arc::new(MockCacheRepository::new("mock"));
2254 repo.seed(
2255 "k",
2256 CacheEntry {
2257 bytes: b"v".to_vec(),
2258 payload_path: None,
2259 content_type: ContentType::Bytes,
2260 expires_at: None,
2261 },
2262 )
2263 .await;
2264 let mut svc = CacheClearService::new(repo.clone());
2265
2266 let outcome = svc.run(exchange()).await;
2267
2268 match outcome {
2269 PipelineOutcome::Completed(_) => {}
2270 other => panic!("expected Completed, got {other:?}"),
2271 }
2272 assert_eq!(repo.clear_call_count(), 1, "clear must be called once");
2273 assert!(
2274 repo.stored_entry("k").await.is_none(),
2275 "entry must be removed after clear"
2276 );
2277 }
2278
2279 #[tokio::test]
2280 async fn cache_clear_err_propagates_failed() {
2281 let repo = Arc::new(MockCacheRepository::new("mock"));
2282 repo.set_should_fail_clear(true);
2283 let mut svc = CacheClearService::new(repo);
2284
2285 let outcome = svc.run(exchange()).await;
2286
2287 match outcome {
2288 PipelineOutcome::Failed(e) => {
2289 assert!(
2290 e.to_string().contains("synthetic clear failure"),
2291 "got: {e}"
2292 );
2293 }
2294 other => panic!("expected Failed, got {other:?}"),
2295 }
2296 }
2297
2298 #[tokio::test]
2301 async fn cache_stats_sets_json_body() {
2302 let repo = Arc::new(MockCacheRepository::new("mock"));
2303 repo.set_stats(CacheStats {
2304 hits: 2,
2305 misses: 1,
2306 evictions: 0,
2307 entries: 3,
2308 peek_stale_served: 4,
2309 invalidations: 1,
2310 bytes: None,
2311 });
2312 let mut svc = CacheStatsService::new(repo);
2313
2314 let outcome = svc.run(exchange()).await;
2315
2316 let ex = match outcome {
2317 PipelineOutcome::Completed(ex) => ex,
2318 other => panic!("expected Completed, got {other:?}"),
2319 };
2320 let expected = serde_json::json!({
2321 "repository": "mock",
2322 "hits": 2,
2323 "misses": 1,
2324 "evictions": 0,
2325 "entries": 3,
2326 "peek_stale_served": 4,
2327 "invalidations": 1,
2328 "bytes": null
2329 });
2330 assert_eq!(ex.input.body, Body::Json(expected));
2331
2332 let Body::Json(v) = &ex.input.body else {
2335 panic!("expected Json body");
2336 };
2337 let keys: std::collections::BTreeSet<&str> = v
2338 .as_object()
2339 .expect("stats body must be a JSON object")
2340 .keys()
2341 .map(String::as_str)
2342 .collect();
2343 let expected_keys: std::collections::BTreeSet<&str> = [
2344 "repository",
2345 "hits",
2346 "misses",
2347 "evictions",
2348 "entries",
2349 "peek_stale_served",
2350 "invalidations",
2351 "bytes",
2352 ]
2353 .into_iter()
2354 .collect();
2355 assert_eq!(
2356 keys, expected_keys,
2357 "stats body must have exactly the eight canonical keys"
2358 );
2359 }
2360
2361 #[tokio::test]
2364 async fn peek_stale_hit_emits_peek_served_counter() {
2365 let repo = Arc::new(MockCacheRepository::new("mock"));
2366 repo.seed(
2367 "cache-key",
2368 CacheEntry {
2369 bytes: b"cached".to_vec(),
2370 payload_path: None,
2371 content_type: ContentType::Bytes,
2372 expires_at: None,
2373 },
2374 )
2375 .await;
2376 let collector = RecordingMetricsCollector::new();
2377 let counters = collector.counters.clone();
2378 let rt = Arc::new(TestOtelmRt {
2379 collector: Arc::new(collector),
2380 });
2381 let mut svc = CachePeekStaleService::new(repo, fixed_key(), PeekStaleMissPolicy::Stop, rt);
2382
2383 let outcome = svc.run(exchange()).await;
2384 assert!(matches!(outcome, PipelineOutcome::Completed(_)));
2385
2386 let recorded = counters.lock().unwrap().clone();
2387 assert!(
2388 recorded.contains(&(
2389 "camel.cache.peek_stale_served".to_string(),
2390 1.0,
2391 vec![("repository".to_string(), "mock".to_string())]
2392 )),
2393 "expected camel.cache.peek_stale_served counter, got: {recorded:?}"
2394 );
2395 }
2396
2397 #[tokio::test]
2398 async fn invalidate_emits_invalidations_counter() {
2399 let repo = Arc::new(MockCacheRepository::new("mock"));
2400 repo.seed(
2401 "cache-key",
2402 CacheEntry {
2403 bytes: b"to-go".to_vec(),
2404 payload_path: None,
2405 content_type: ContentType::Bytes,
2406 expires_at: None,
2407 },
2408 )
2409 .await;
2410 let collector = RecordingMetricsCollector::new();
2411 let counters = collector.counters.clone();
2412 let rt = Arc::new(TestOtelmRt {
2413 collector: Arc::new(collector),
2414 });
2415 let mut svc =
2416 CacheInvalidateService::new(repo, CacheInvalidateTarget::Key(fixed_key()), rt);
2417
2418 let outcome = svc.run(exchange()).await;
2419 assert!(matches!(outcome, PipelineOutcome::Completed(_)));
2420
2421 let recorded = counters.lock().unwrap().clone();
2422 assert!(
2423 recorded.contains(&(
2424 "camel.cache.invalidations".to_string(),
2425 1.0,
2426 vec![("repository".to_string(), "mock".to_string())]
2427 )),
2428 "expected camel.cache.invalidations counter, got: {recorded:?}"
2429 );
2430 }
2431
2432 #[tokio::test]
2433 async fn peek_stale_miss_emits_no_peek_served_counter() {
2434 let repo = Arc::new(MockCacheRepository::new("mock"));
2436 let collector = RecordingMetricsCollector::new();
2437 let counters = collector.counters.clone();
2438 let rt = Arc::new(TestOtelmRt {
2439 collector: Arc::new(collector),
2440 });
2441 let mut svc = CachePeekStaleService::new(repo, fixed_key(), PeekStaleMissPolicy::Stop, rt);
2442
2443 let outcome = svc.run(exchange()).await;
2444 assert!(matches!(outcome, PipelineOutcome::Stopped(_)));
2445
2446 let recorded = counters.lock().unwrap().clone();
2447 assert!(
2448 !recorded
2449 .iter()
2450 .any(|(name, _, _)| name == "camel.cache.peek_stale_served"),
2451 "expected zero camel.cache.peek_stale_served counters, got: {recorded:?}"
2452 );
2453 }
2454
2455 #[tokio::test]
2456 async fn invalidate_err_emits_no_invalidations_counter() {
2457 let repo = Arc::new(MockCacheRepository::new("mock"));
2459 repo.seed(
2460 "cache-key",
2461 CacheEntry {
2462 bytes: b"to-go".to_vec(),
2463 payload_path: None,
2464 content_type: ContentType::Bytes,
2465 expires_at: None,
2466 },
2467 )
2468 .await;
2469 repo.set_should_fail_invalidate(true);
2470 let collector = RecordingMetricsCollector::new();
2471 let counters = collector.counters.clone();
2472 let rt = Arc::new(TestOtelmRt {
2473 collector: Arc::new(collector),
2474 });
2475 let mut svc =
2476 CacheInvalidateService::new(repo, CacheInvalidateTarget::Key(fixed_key()), rt);
2477
2478 let outcome = svc.run(exchange()).await;
2479 assert!(matches!(outcome, PipelineOutcome::Failed(_)));
2480
2481 let recorded = counters.lock().unwrap().clone();
2482 assert!(
2483 !recorded
2484 .iter()
2485 .any(|(name, _, _)| name == "camel.cache.invalidations"),
2486 "expected zero camel.cache.invalidations counters, got: {recorded:?}"
2487 );
2488 }
2489
2490 #[tokio::test]
2493 async fn cache_invalidate_prefix_removes_namespace_sets_count() {
2494 let repo = Arc::new(MockCacheRepository::new("mock"));
2495 for key in ["ns:one", "ns:two", "other:x"] {
2496 repo.seed(
2497 key,
2498 CacheEntry {
2499 bytes: key.as_bytes().to_vec(),
2500 payload_path: None,
2501 content_type: ContentType::Bytes,
2502 expires_at: None,
2503 },
2504 )
2505 .await;
2506 }
2507
2508 let collector = RecordingMetricsCollector::new();
2509 let counters = collector.counters.clone();
2510 let rt = Arc::new(TestOtelmRt {
2511 collector: Arc::new(collector),
2512 });
2513 let mut svc = CacheInvalidateService::new(
2514 repo.clone(),
2515 CacheInvalidateTarget::Prefix(prefix_key()),
2516 rt,
2517 );
2518
2519 let outcome = svc.run(exchange()).await;
2520
2521 let ex = match outcome {
2522 PipelineOutcome::Completed(ex) => ex,
2523 other => panic!("expected Completed, got {other:?}"),
2524 };
2525 assert_eq!(
2526 ex.property(CAMEL_CACHE_INVALIDATED_COUNT),
2527 Some(&serde_json::Value::from(2u64)),
2528 "prefix purge must report the removed count"
2529 );
2530 assert!(
2531 repo.stored_entry("ns:one").await.is_none(),
2532 "ns:one must be removed"
2533 );
2534 assert!(
2535 repo.stored_entry("ns:two").await.is_none(),
2536 "ns:two must be removed"
2537 );
2538 assert!(
2539 repo.stored_entry("other:x").await.is_some(),
2540 "other:x must be preserved"
2541 );
2542
2543 let recorded = counters.lock().unwrap().clone();
2544 assert!(
2545 recorded.contains(&(
2546 "camel.cache.invalidations".to_string(),
2547 1.0,
2548 vec![("repository".to_string(), "mock".to_string())]
2549 )),
2550 "expected one camel.cache.invalidations counter, got: {recorded:?}"
2551 );
2552 }
2553
2554 #[tokio::test]
2555 async fn cache_invalidate_prefix_none_expr_completes() {
2556 let repo = Arc::new(MockCacheRepository::new("mock"));
2557 let mut svc = CacheInvalidateService::new(
2558 repo.clone(),
2559 CacheInvalidateTarget::Prefix(none_key()),
2560 noop_rt(),
2561 );
2562
2563 let outcome = svc.run(exchange()).await;
2564
2565 let ex = match outcome {
2566 PipelineOutcome::Completed(ex) => ex,
2567 other => panic!("expected Completed, got {other:?}"),
2568 };
2569 assert_eq!(
2570 repo.invalidate_call_count(),
2571 0,
2572 "no invalidate calls when prefix expr resolves to None"
2573 );
2574 assert!(
2575 ex.property(CAMEL_CACHE_INVALIDATED_COUNT).is_none(),
2576 "no count property when prefix expr resolves to None"
2577 );
2578 }
2579
2580 #[tokio::test]
2581 async fn cache_invalidate_prefix_unsupported_fails_closed() {
2582 let repo = Arc::new(MockCacheRepository::new("mock"));
2583 repo.set_prefix_unsupported(true);
2584 let mut svc = CacheInvalidateService::new(
2585 repo,
2586 CacheInvalidateTarget::Prefix(prefix_key()),
2587 noop_rt(),
2588 );
2589
2590 let outcome = svc.run(exchange()).await;
2591
2592 match outcome {
2593 PipelineOutcome::Failed(e) => {
2594 let msg = format!("{e}");
2595 assert!(
2596 msg.contains("mock"),
2597 "error must name the backend, got: {msg}"
2598 );
2599 }
2600 other => panic!("expected Failed, got {other:?}"),
2601 }
2602 }
2603
2604 use std::sync::atomic::AtomicUsize;
2607 use tokio::sync::Notify;
2608
2609 #[derive(Clone)]
2611 enum GatedOutcome {
2612 Complete(Body),
2613 Fail(CamelError),
2614 Stop,
2615 }
2616
2617 #[derive(Clone)]
2622 struct GatedOnMiss {
2623 leader_entered: Option<Arc<Notify>>,
2624 release: Option<Arc<Notify>>,
2625 invocations: Arc<AtomicUsize>,
2626 outcome: GatedOutcome,
2627 }
2628
2629 impl OutcomePipeline for GatedOnMiss {
2630 fn clone_box(&self) -> Box<dyn OutcomePipeline> {
2631 Box::new(self.clone())
2632 }
2633
2634 fn run<'a>(
2635 &'a mut self,
2636 mut exchange: Exchange,
2637 ) -> Pin<Box<dyn Future<Output = PipelineOutcome> + Send + 'a>> {
2638 let leader_entered = self.leader_entered.clone();
2639 let release = self.release.clone();
2640 let invocations = Arc::clone(&self.invocations);
2641 let outcome = self.outcome.clone();
2642 Box::pin(async move {
2643 if let Some(entered) = leader_entered.as_ref() {
2645 entered.notify_one();
2646 }
2647 if let Some(release) = release.as_ref() {
2649 release.notified().await;
2650 }
2651 invocations.fetch_add(1, Ordering::SeqCst);
2652 match outcome {
2653 GatedOutcome::Complete(body) => {
2654 exchange.input.body = body;
2655 PipelineOutcome::Completed(exchange)
2656 }
2657 GatedOutcome::Fail(e) => PipelineOutcome::Failed(e),
2658 GatedOutcome::Stop => PipelineOutcome::Stopped(exchange),
2659 }
2660 })
2661 }
2662 }
2663
2664 fn build_gated_service(
2666 repo: Arc<MockCacheRepository>,
2667 outcome: GatedOutcome,
2668 leader_entered: Option<Arc<Notify>>,
2669 release: Option<Arc<Notify>>,
2670 ) -> (CacheService, Arc<AtomicUsize>) {
2671 let invocations = Arc::new(AtomicUsize::new(0));
2672 let on_miss = OutcomeSegment::new(Box::new(GatedOnMiss {
2673 leader_entered,
2674 release,
2675 invocations: Arc::clone(&invocations),
2676 outcome,
2677 }));
2678 let svc = CacheService::new(repo, fixed_key(), None, 1024, on_miss, noop_rt())
2679 .with_coalesce(true);
2680 (svc, invocations)
2681 }
2682
2683 fn exchange_with_body(text: &str) -> Exchange {
2684 Exchange::new(Message::new(text))
2685 }
2686
2687 #[tokio::test]
2688 async fn coalesce_three_concurrent_misses_fetch_once() {
2689 let repo = Arc::new(MockCacheRepository::new("mock"));
2690 let leader_entered = Arc::new(Notify::new());
2691 let release = Arc::new(Notify::new());
2692 let (svc, invocations) = build_gated_service(
2693 repo.clone(),
2694 GatedOutcome::Complete(Body::Text("fetched".into())),
2695 Some(Arc::clone(&leader_entered)),
2696 Some(Arc::clone(&release)),
2697 );
2698
2699 let entered = leader_entered.notified();
2702 tokio::pin!(entered);
2703 entered.as_mut().enable();
2704
2705 let mut leader_svc = svc.clone();
2706 let leader = tokio::spawn(async move { leader_svc.run(exchange()).await });
2707 entered.await; let mut w1_svc = svc.clone();
2710 let mut w1 = tokio::spawn(async move { w1_svc.run(exchange()).await });
2711 let mut w2_svc = svc.clone();
2712 let mut w2 = tokio::spawn(async move { w2_svc.run(exchange()).await });
2713
2714 for waiter in [&mut w1, &mut w2] {
2716 if let Ok(done) = tokio::time::timeout(Duration::from_millis(50), waiter).await {
2717 panic!("waiter resolved before release: {done:?}")
2718 }
2719 }
2720
2721 release.notify_waiters();
2722
2723 let leader_ex = match tokio::time::timeout(Duration::from_secs(5), leader)
2724 .await
2725 .expect("leader task join timeout")
2726 .expect("leader task join")
2727 {
2728 PipelineOutcome::Completed(ex) => ex,
2729 other => panic!("expected leader Completed, got {other:?}"),
2730 };
2731 let w1_ex = match tokio::time::timeout(Duration::from_secs(5), w1)
2732 .await
2733 .expect("waiter 1 task join timeout")
2734 .expect("waiter 1 task join")
2735 {
2736 PipelineOutcome::Completed(ex) => ex,
2737 other => panic!("expected waiter 1 Completed, got {other:?}"),
2738 };
2739 let w2_ex = match tokio::time::timeout(Duration::from_secs(5), w2)
2740 .await
2741 .expect("waiter 2 task join timeout")
2742 .expect("waiter 2 task join")
2743 {
2744 PipelineOutcome::Completed(ex) => ex,
2745 other => panic!("expected waiter 2 Completed, got {other:?}"),
2746 };
2747 assert_eq!(leader_ex.input.body, Body::Text("fetched".into()));
2748 assert_eq!(w1_ex.input.body, Body::Text("fetched".into()));
2749 assert_eq!(w2_ex.input.body, Body::Text("fetched".into()));
2750 assert_eq!(invocations.load(Ordering::SeqCst), 1, "on_miss ran once");
2751 assert_eq!(repo.set_call_count(), 1, "single write-back set");
2752 }
2753
2754 #[tokio::test]
2755 async fn coalesce_leader_failure_fails_waiters_once() {
2756 let repo = Arc::new(MockCacheRepository::new("mock"));
2757 let leader_entered = Arc::new(Notify::new());
2758 let release = Arc::new(Notify::new());
2759 let (svc, invocations) = build_gated_service(
2760 repo.clone(),
2761 GatedOutcome::Fail(stub_error("coalesce-boom")),
2762 Some(Arc::clone(&leader_entered)),
2763 Some(Arc::clone(&release)),
2764 );
2765
2766 let entered = leader_entered.notified();
2767 tokio::pin!(entered);
2768 entered.as_mut().enable();
2769
2770 let mut leader_svc = svc.clone();
2771 let leader = tokio::spawn(async move { leader_svc.run(exchange()).await });
2772 entered.await;
2773
2774 let mut w1_svc = svc.clone();
2775 let mut w1 = tokio::spawn(async move { w1_svc.run(exchange()).await });
2776 let mut w2_svc = svc.clone();
2777 let mut w2 = tokio::spawn(async move { w2_svc.run(exchange()).await });
2778
2779 for waiter in [&mut w1, &mut w2] {
2780 if let Ok(done) = tokio::time::timeout(Duration::from_millis(50), waiter).await {
2781 panic!("waiter resolved before release: {done:?}")
2782 }
2783 }
2784
2785 release.notify_waiters();
2786
2787 let leader_err = match leader.await.expect("leader task join") {
2788 PipelineOutcome::Failed(e) => e,
2789 other => panic!("expected leader Failed, got {other:?}"),
2790 };
2791 let w1_err = match w1.await.expect("waiter 1 task join") {
2792 PipelineOutcome::Failed(e) => e,
2793 other => panic!("expected waiter 1 Failed, got {other:?}"),
2794 };
2795 let w2_err = match w2.await.expect("waiter 2 task join") {
2796 PipelineOutcome::Failed(e) => e,
2797 other => panic!("expected waiter 2 Failed, got {other:?}"),
2798 };
2799 assert_eq!(format!("{leader_err}"), format!("{w1_err}"));
2800 assert_eq!(format!("{leader_err}"), format!("{w2_err}"));
2801 assert_eq!(invocations.load(Ordering::SeqCst), 1, "on_miss ran once");
2802 assert_eq!(repo.set_call_count(), 0, "no write-back on failure");
2803 }
2804
2805 #[tokio::test]
2806 async fn coalesce_leader_stopped_stops_waiters() {
2807 let repo = Arc::new(MockCacheRepository::new("mock"));
2808 let leader_entered = Arc::new(Notify::new());
2809 let release = Arc::new(Notify::new());
2810 let (svc, invocations) = build_gated_service(
2811 repo.clone(),
2812 GatedOutcome::Stop,
2813 Some(Arc::clone(&leader_entered)),
2814 Some(Arc::clone(&release)),
2815 );
2816
2817 let entered = leader_entered.notified();
2818 tokio::pin!(entered);
2819 entered.as_mut().enable();
2820
2821 let mut leader_svc = svc.clone();
2822 let leader =
2823 tokio::spawn(async move { leader_svc.run(exchange_with_body("leader-orig")).await });
2824 entered.await;
2825
2826 let mut w1_svc = svc.clone();
2827 let mut w1 =
2828 tokio::spawn(async move { w1_svc.run(exchange_with_body("waiter-orig")).await });
2829
2830 if let Ok(done) = tokio::time::timeout(Duration::from_millis(50), &mut w1).await {
2831 panic!("waiter resolved before release: {done:?}")
2832 }
2833
2834 release.notify_waiters();
2835
2836 let leader_ex = match leader.await.expect("leader task join") {
2837 PipelineOutcome::Stopped(ex) => ex,
2838 other => panic!("expected leader Stopped, got {other:?}"),
2839 };
2840 let waiter_ex = match w1.await.expect("waiter task join") {
2841 PipelineOutcome::Stopped(ex) => ex,
2842 other => panic!("expected waiter Stopped, got {other:?}"),
2843 };
2844 assert_eq!(
2845 leader_ex.input.body,
2846 Body::Text("leader-orig".into()),
2847 "leader stopped with its own exchange"
2848 );
2849 assert_eq!(
2850 waiter_ex.input.body,
2851 Body::Text("waiter-orig".into()),
2852 "waiter stopped with its own exchange, body untouched"
2853 );
2854 assert_eq!(invocations.load(Ordering::SeqCst), 1, "on_miss ran once");
2855 assert_eq!(repo.set_call_count(), 0, "no write-back on stop");
2856 }
2857
2858 #[tokio::test]
2859 async fn coalesce_leader_dropped_does_not_strand_waiters() {
2860 let repo = Arc::new(MockCacheRepository::new("mock"));
2861 let leader_entered = Arc::new(Notify::new());
2862 let release = Arc::new(Notify::new());
2863 let (svc, _invocations) = build_gated_service(
2864 repo,
2865 GatedOutcome::Complete(Body::Text("fetched".into())),
2866 Some(Arc::clone(&leader_entered)),
2867 Some(Arc::clone(&release)),
2868 );
2869
2870 let entered = leader_entered.notified();
2871 tokio::pin!(entered);
2872 entered.as_mut().enable();
2873
2874 let mut leader_svc = svc.clone();
2875 let leader = tokio::spawn(async move { leader_svc.run(exchange()).await });
2876 entered.await;
2877
2878 let mut w1_svc = svc.clone();
2879 let mut w1 = tokio::spawn(async move { w1_svc.run(exchange()).await });
2880
2881 if let Ok(done) = tokio::time::timeout(Duration::from_millis(50), &mut w1).await {
2882 panic!("waiter resolved before leader abort: {done:?}")
2883 }
2884
2885 leader.abort();
2888
2889 let joined = tokio::time::timeout(Duration::from_secs(1), w1)
2890 .await
2891 .expect("waiter completes within 1s after leader drop")
2892 .expect("waiter task join");
2893 match joined {
2894 PipelineOutcome::Failed(e) => {
2895 let msg = format!("{e}");
2896 assert!(
2897 msg.contains("cancelled"),
2898 "expected cancellation terminal, got: {msg}"
2899 );
2900 }
2901 other => panic!("expected waiter Failed, got {other:?}"),
2902 }
2903 assert!(
2904 svc.inflight.lock().unwrap().is_empty(),
2905 "in-flight map must not retain the aborted wave's entry"
2906 );
2907 }
2908
2909 #[tokio::test]
2910 async fn no_coalesce_runs_per_exchange() {
2911 let repo = Arc::new(MockCacheRepository::new("mock"));
2912 let invocations = Arc::new(AtomicUsize::new(0));
2913 let on_miss = OutcomeSegment::new(Box::new(GatedOnMiss {
2914 leader_entered: None,
2915 release: None,
2916 invocations: Arc::clone(&invocations),
2917 outcome: GatedOutcome::Complete(Body::Text("per-exchange".into())),
2918 }));
2919 let svc = CacheService::new(repo.clone(), fixed_key(), None, 1024, on_miss, noop_rt());
2921
2922 let mut a = svc.clone();
2923 let mut b = svc.clone();
2924 let mut c = svc.clone();
2925 let (ra, rb, rc) = tokio::join!(a.run(exchange()), b.run(exchange()), c.run(exchange()));
2926
2927 for outcome in [ra, rb, rc] {
2928 match outcome {
2929 PipelineOutcome::Completed(ex) => {
2930 assert_eq!(ex.input.body, Body::Text("per-exchange".into()));
2931 }
2932 other => panic!("expected Completed, got {other:?}"),
2933 }
2934 }
2935 assert_eq!(
2936 invocations.load(Ordering::SeqCst),
2937 3,
2938 "on_miss ran per exchange"
2939 );
2940 assert_eq!(repo.set_call_count(), 3, "set called per exchange");
2941 }
2942}