Skip to main content

camel_processor/
cache_eip.rs

1//! Cache EIP — outcome-aware Segment implementation.
2//!
3//! Implements the Caching pattern (lookup → on-miss sub-pipeline → write-back)
4//! at the `OutcomePipeline` layer (one layer above Tower), mirroring
5//! [`IdempotentConsumerSegment`](crate::IdempotentConsumerSegment). On a cache HIT the body is reconstructed from
6//! the stored [`CacheEntry`] and the on-miss sub-pipeline is skipped entirely;
7//! on a MISS the sub-pipeline runs and its result body is written back into the
8//! repository (subject to `max_entry_bytes`).
9//!
10//! # Why Segment-mode (NOT Process-mode)
11//!
12//! Same rationale as the idempotent consumer: a Tower `Service<Exchange>` cannot
13//! propagate `PipelineOutcome::Stopped` distinctly from `Ok(ex)`. By implementing
14//! [`OutcomePipeline`] directly, a `Stopped` from the on-miss sub-pipeline flows
15//! out with the Exchange intact and NO write-back occurs (ADR-0024, ADR-0025).
16//!
17//! # Contract C1 (ADR-0023)
18//!
19//! [`CacheRepository::get`] / [`CacheRepository::set`] surface backend failures
20//! as `Err(CamelError)`. The segment propagates those as `PipelineOutcome::Failed`
21//! — it NEVER treats a failed read as a miss.
22
23use 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// ── Singleflight miss coalescing (cache-admin task 2.4) ──
38
39/// Terminal state a coalescing leader publishes for its waiters.
40///
41/// Mirrors the leader's own `PipelineOutcome` in a clonable shape so
42/// every waiter of the wave receives it (waiters clone-read; nobody
43/// consumes the slot).
44#[derive(Clone)]
45pub(crate) enum CoalesceTerminal {
46    /// Leader completed: waiters adopt the leader's resulting body.
47    Completed(Body),
48    /// Leader failed: waiters fail with the same error (anti-burst).
49    Failed(CamelError),
50    /// Leader stopped: waiters stop their own exchanges (branch-filter).
51    Stopped,
52}
53
54/// One in-flight coalescing wave for a resolved cache key.
55///
56/// `terminal` is a write-once slot filled BEFORE `notify_waiters()` is
57/// called, so a woken waiter always re-reads a filled slot (no lost
58/// wakeup: `notify_waiters` alone wakes only currently-registered
59/// waiters).
60struct 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    /// Publish the terminal state. Write-once: an already-filled slot is
76    /// never cleared or overwritten (a late `LeaderGuard::drop` after a
77    /// normal completion is a no-op here).
78    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    /// Clone-read the terminal state, if published.
87    fn terminal_snapshot(&self) -> Option<CoalesceTerminal> {
88        match self.terminal.lock() {
89            Ok(slot) => (*slot).clone(),
90            Err(_) => None,
91        }
92    }
93}
94
95/// In-flight coalescing waves, keyed by resolved cache key, scoped per
96/// compiled route-step instance (shared across `CacheService` clones).
97type InFlightMap = std::sync::Mutex<std::collections::HashMap<String, std::sync::Arc<InFlight>>>;
98
99/// Cancellation guard for the coalescing leader.
100///
101/// If the leader future is dropped before it publishes its terminal
102/// state (route shutdown, task abort), `Drop` publishes a cancellation
103/// terminal (`Failed`) into the write-once slot, wakes waiters, and
104/// removes the map entry — so waiters are never stranded. On normal
105/// completion the slot is already filled (publish is skipped) and the
106/// entry is already retired; both `Drop` actions become no-ops.
107struct LeaderGuard {
108    key: String,
109    map: Arc<InFlightMap>,
110    cell: Arc<InFlight>,
111}
112
113impl LeaderGuard {
114    /// Remove the map entry iff it still identifies `cell`.
115    ///
116    /// The `Arc::ptr_eq` identity check keeps a late guard (or a
117    /// completed leader) from evicting a NEWER wave's entry for the
118    /// same key.
119    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        // Write-once: no-op when the leader already published a terminal.
133        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
142/// Outcome-aware Cache segment (Caching EIP).
143///
144/// Wraps a named [`CacheRepository`] and an on-miss sub-pipeline
145/// ([`OutcomeSegment`]). On each exchange:
146///
147/// 1. Evaluate `key_expr`. `None` → not cacheable; forward directly to the
148///    on-miss sub-pipeline (no lookup, no write-back).
149/// 2. `repository.get(&key)`:
150///    - `Err(e)` → `Failed(e)` (contract C1).
151///    - `Ok(Some(entry))` → HIT: reconstruct `Body` from the entry, set it on
152///      the exchange, return `Completed` (skip on-miss).
153///    - `Ok(None)` → MISS: proceed to step 3.
154/// 3. Run the on-miss sub-pipeline.
155///    - `Stopped(ex)` / `Failed(e)` → propagate as-is (NO write-back).
156///    - `Completed(ex)` → proceed to write-back.
157/// 4. Write-back the resulting body (when it fits `max_entry_bytes`):
158///    - materialized variants (`Bytes`/`Text`/`Json`/`Xml`) → serialize, store.
159///    - `Stream` → materialize via [`Body::into_bytes`] (consumes the body,
160///      replaces it with `Body::Bytes`); `StreamLimitExceeded` propagates.
161///    - `Empty` / oversized body → pass through uncached, return `Completed`.
162///
163/// With `coalesce_misses` enabled ([`CacheService::with_coalesce`]),
164/// concurrent misses on the same resolved key are coalesced
165/// (singleflight): the first exchange (leader) runs `on_miss` and the
166/// single write-back `set`; concurrent exchanges (waiters) await the
167/// leader's terminal state instead of running `on_miss`. HIT, key-`None`,
168/// and `coalesce_misses == false` paths bypass the in-flight map
169/// entirely.
170pub struct CacheService {
171    repository: Arc<dyn CacheRepository>,
172    /// Cached `repository.name()` for OTel span tagging (Task 3.3).
173    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    /// Singleflight miss coalescing toggle (default `false`).
180    coalesce_misses: bool,
181    /// In-flight coalescing waves. `Clone` clones the `Arc`, so every
182    /// service clone of one compiled route-step shares the same map.
183    inflight: Arc<InFlightMap>,
184}
185
186impl CacheService {
187    /// Build a new cache segment.
188    ///
189    /// `repository_name` is derived from `repository.name()` so OTel tags stay
190    /// in sync with the resolved backend.
191    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    /// Enable (or explicitly disable) singleflight miss coalescing.
214    ///
215    /// With coalescing on, concurrent misses on the same resolved key
216    /// run the `on_miss` sub-pipeline exactly once per wave (leader
217    /// runs + writes back; waiters receive the leader's terminal state).
218    pub fn with_coalesce(mut self, coalesce_misses: bool) -> Self {
219        self.coalesce_misses = coalesce_misses;
220        self
221    }
222
223    /// The configured repository name (for OTel tagging).
224    pub fn repository_name(&self) -> &str {
225        &self.repository_name
226    }
227}
228
229/// Shared write-back tail for materialized bodies.
230///
231/// Checks `max_entry_bytes`, builds a [`CacheEntry`], stores via
232/// the repository, and returns `Completed(exchange)`. The exchange body
233/// is not modified — it passes through as-is. On oversized body, logs a
234/// debug! skip message and returns `Completed(exchange)` without storing.
235/// On repository error, returns `Failed(e)`.
236#[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                    ); // log-policy: g:cache:capacity-full-skip
263                } else {
264                    return PipelineOutcome::Failed(e);
265                }
266            }
267        }
268    } else {
269        // log-policy: g:cache:oversized-skip
270        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            // 1. Evaluate key. None → not cacheable, bypass straight to
308            //    on_miss; a failed evaluation fails the step (contract C1).
309            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            // 2. Lookup (contract C1: propagate Err, never treat as miss).
316            match self.repository.get(&key).await {
317                Err(e) => return PipelineOutcome::Failed(e),
318                Ok(Some(entry)) => {
319                    // HIT: record metric, reconstruct body, skip on-miss sub-pipeline.
320                    // allow-open-label rc-ycts (repository: user-declared cache repository name)
321                    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                    // MISS: record metric, fall through to on-miss sub-pipeline.
337                    // allow-open-label rc-ycts (repository: user-declared cache repository name)
338                    self.rt.metrics().record_counter(
339                        "camel.cache.misses",
340                        1.0_f64,
341                        &[("repository", &self.repository_name)],
342                    );
343                }
344            }
345
346            // 3./4. MISS flow, singleflight-coalesced when enabled.
347            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    /// The un-coalesced MISS flow (spec steps 3-4): run the on-miss
357    /// sub-pipeline, then write the resulting body back (subject to
358    /// `max_entry_bytes` and the materialization policy). Also the
359    /// leader's flow under coalescing.
360    async fn run_miss(&mut self, exchange: Exchange, key: String) -> PipelineOutcome {
361        // 3. Run the on-miss sub-pipeline.
362        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        // 4. Write-back. Take the body out so the Stream arm can consume it.
369        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                // Materialize (consumes the stream). StreamLimitExceeded propagates.
441                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                // into_bytes already enforced max_entry_bytes, so it fits by construction.
449                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                    // Degrade capacity-exceeded to uncached — same policy as write_back.
457                    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                        ); // log-policy: g:cache:capacity-full-skip
464                        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                // Empty (or any future variant): pass through uncached.
475                exchange.input.body = body;
476                PipelineOutcome::Completed(exchange)
477            }
478        }
479    }
480
481    /// The coalesced MISS flow (singleflight, cache-admin task 2.4).
482    ///
483    /// The first exchange on a key (leader) inserts the in-flight cell
484    /// and runs [`CacheService::run_miss`] (on_miss + the single
485    /// write-back `set`) under a [`LeaderGuard`]; concurrent misses on
486    /// the same key (waiters) do NOT run on_miss — they await the
487    /// leader's terminal state and clone-read it.
488    ///
489    /// Protocol (cancellation-safe, race-free):
490    /// - Waiter registration is atomic with the map lookup: the
491    ///   pinned-`Notified` `enable()` happens while STILL holding the
492    ///   map lock, so the leader's publish+notify cannot slip between
493    ///   the lookup and the registration.
494    /// - The terminal slot is filled BEFORE `notify_waiters()`, and a
495    ///   woken waiter re-reads the slot (no lost wakeup —
496    ///   `notify_waiters` alone wakes only currently-registered
497    ///   waiters).
498    /// - Map removal happens only on `Arc::ptr_eq` identity with this
499    ///   leader's cell (a late guard cannot evict a newer wave's entry).
500    /// - The slot is write-once: once filled it is never cleared or
501    ///   overwritten.
502    async fn coalesced_miss(&mut self, exchange: Exchange, key: String) -> PipelineOutcome {
503        let inflight = Arc::clone(&self.inflight);
504
505        // Resolve the role under ONE short lock scope; the guard (and
506        // the lock Result) never crosses an await. A waiter registers
507        // its Notified (Box::pin + enable) while STILL holding the map
508        // lock — registration atomic with the lookup. The registration
509        // borrows the outer `wave` binding (which outlives the guard),
510        // so it can be awaited after the guard is gone.
511        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                        // WAITER: claim the existing wave, then register
518                        // (pin + enable) while STILL holding the lock.
519                        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                        // LEADER: claim the key.
528                        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            // Unreachable on the Ok path (both arms set `wave`); a
538            // poisoned in-flight map degrades to un-coalesced
539            // execution rather than stranding exchanges.
540            return self.run_miss(exchange, key).await;
541        };
542
543        if let Some(notified) = registered {
544            // WAITER. The slot may already be filled (leader finished
545            // between registration and this read): clone-read without
546            // parking. Otherwise await the leader's notify and re-read.
547            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            // LEADER. Run the miss flow under a cancellation guard,
569            // publish the terminal state, wake waiters, retire the
570            // map entry.
571            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            // Slot BEFORE notify: woken waiters re-read a filled slot.
586            cell.publish(terminal);
587            cell.notify.notify_waiters();
588            LeaderGuard::retire(&inflight, &key, &cell);
589            // Guard drop is a no-op now: the write-once slot is
590            // filled and the entry is retired.
591            outcome
592        }
593    }
594}
595
596/// Reconstruct a [`Body`] from a stored [`CacheEntry`].
597///
598/// Maps each [`ContentType`] back to the matching `Body` variant, decoding
599/// UTF-8 / JSON failures into `CamelError::TypeConversionFailed`.
600fn 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
624// ===========================================================================
625// CacheInvalidateService — invalidate a single cache entry or a namespace
626// ===========================================================================
627
628/// Exchange property set to the number of entries removed by a successful
629/// `cache_invalidate` step — always `1` for exact-key (removal is not
630/// observable by the backend), the returned count for a namespace purge. Not
631/// set when the key/prefix expression resolves to `None` or the backend
632/// reports an error.
633pub const CAMEL_CACHE_INVALIDATED_COUNT: &str = "CamelCacheInvalidatedCount";
634
635/// The invalidation target of a [`CacheInvalidateService`]: an exact key or a
636/// namespace prefix.
637#[derive(Clone)]
638pub enum CacheInvalidateTarget {
639    /// Invalidate the single entry under the resolved key.
640    Key(MessageIdSource),
641    /// Invalidate every entry whose key starts with the resolved prefix.
642    Prefix(MessageIdSource),
643}
644
645/// Outcome-aware segment that invalidates a single cache entry or a namespace.
646///
647/// Evaluates the configured [`CacheInvalidateTarget`]:
648/// - [`Key`](CacheInvalidateTarget::Key): expression `None` → `Completed`
649///   (nothing to invalidate); `Some(key)` → `repository.invalidate(&key).await`.
650///   - `Err(e)` → `Failed(e)`.
651///   - `Ok(())` → sets `CAMEL_CACHE_INVALIDATED_COUNT = 1`, emits
652///     `camel.cache.invalidations` +1, `Completed(exchange)`.
653/// - [`Prefix`](CacheInvalidateTarget::Prefix): expression `None` → `Completed`;
654///   `Some(prefix)` → `repository.invalidate_prefix(&prefix).await`.
655///   - `Err(e)` → `Failed(e)` (an unsupported backend surfaces as failure —
656///     fail-closed).
657///   - `Ok(count)` → sets `CAMEL_CACHE_INVALIDATED_COUNT = count`, emits
658///     `camel.cache.invalidations` +1, `Completed(exchange)`.
659pub 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                            // allow-open-label rc-ycts (repository: user-declared cache repository name)
714                            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                            // allow-open-label rc-ycts (repository: user-declared cache repository name)
738                            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
757// ===========================================================================
758// CachePeekStaleService — serve a stale entry after expiry
759// ===========================================================================
760
761/// Exchange property set to `true` when a `cache_peek_stale` HIT occurred.
762pub const CAMEL_CACHE_PEEK_HIT: &str = "CamelCachePeekHit";
763/// Exchange property set to `true` when the served entry was stale (post-expiry).
764pub const CAMEL_CACHE_PEEK_STALE: &str = "CamelCachePeekStale";
765
766/// On-miss policy for [`CachePeekStaleService`].
767///
768/// - [`Stop`](PeekStaleMissPolicy::Stop) (default) preserves the
769///   `CircuitBreaker.fallback` absence-Stops contract.
770/// - [`Continue`](PeekStaleMissPolicy::Continue) leaves the body untouched on
771///   MISS so `choice` can branch on [`CAMEL_CACHE_PEEK_HIT`].
772#[derive(Debug, Clone, Copy, PartialEq, Eq)]
773pub enum PeekStaleMissPolicy {
774    /// MISS Stops the branch (no stale available — `CircuitBreaker.fallback`).
775    Stop,
776    /// MISS continues with the body unchanged.
777    Continue,
778}
779
780impl PeekStaleMissPolicy {
781    /// Parses the canonical/DSL `cache_peek_stale.on_miss` knob:
782    /// absent or `"stop"` → [`Stop`](Self::Stop), `"continue"` →
783    /// [`Continue`](Self::Continue). Any other value fails closed naming
784    /// the step.
785    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
796/// Outcome-aware segment that serves a stale (post-expiry) cache entry.
797///
798/// Evaluates `key_expr`:
799/// - `None` → `Stopped(exchange)` with a `debug` log (an anomalous key
800///   resolution is fail-closed, not a miss).
801/// - `Some(key)` → `repository.peek_stale(&key).await`.
802///   - `Err(e)` → `Failed(e)`.
803///   - `Ok(Some(entry))` → reconstruct body from entry, set
804///     `CamelCachePeekHit=true` and `CamelCachePeekStale` (true when the
805///     entry's `expires_at` has elapsed at evaluation time; false when absent
806///     or not elapsed), return `Completed(exchange)`.
807///   - `Ok(None)` → MISS (absence), governed by [`PeekStaleMissPolicy`]:
808///     - `Stop` (default): set `CamelCachePeekHit=false` and
809///       `CamelCachePeekStale=false`, log at `debug`, return `Stopped(exchange)`
810///       (absence in `CircuitBreaker.fallback` means "no stale available").
811///     - `Continue`: set `CamelCachePeekHit=false` and
812///       `CamelCachePeekStale=false`, leave the body unchanged, return
813///       `Completed(exchange)` so `choice` can branch on `CamelCachePeekHit`.
814pub 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
852/// Write the peek-result exchange properties under [`CAMEL_CACHE_PEEK_HIT`] and
853/// [`CAMEL_CACHE_PEEK_STALE`] as `serde_json::Value::Bool` values.
854fn 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                    // Emit peek_stale_served (fresh or stale — both are serves).
885                    // allow-open-label rc-ycts (repository: user-declared cache repository name)
886                    self.rt.metrics().record_counter(
887                        "camel.cache.peek_stale_served",
888                        1.0_f64,
889                        &[("repository", &self.repository_name)],
890                    );
891                    // Staleness read before body reconstruction; reconstruct_body borrows the entry, so ordering is stylistic.
892                    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
928// ===========================================================================
929// CacheClearService — remove all entries from the repository
930// ===========================================================================
931
932/// Outcome-aware segment that clears the entire cache repository.
933///
934/// - `Err(e)` → `Failed(e)`.
935/// - `Ok(())` → `Completed(exchange)` with the body unchanged.
936pub 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
972// ===========================================================================
973// CacheStatsService — emit the repository stats as a JSON body
974// ===========================================================================
975
976/// Outcome-aware segment that replaces the exchange body with a JSON snapshot
977/// of the repository's [`CacheStats`](camel_api::CacheStats).
978pub 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// ===========================================================================
1029// Test utilities
1030// ===========================================================================
1031
1032#[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    /// In-memory mock [`CacheRepository`] for cache segment tests. Allows tests
1042    /// to pre-seed entries, force `get`/`set` failures, and inspect the last
1043    /// `set` call (entry + TTL).
1044    #[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        /// Force `invalidate_prefix` to report the backend-naming unsupported
1090        /// error (mirrors the default `CacheRepository::invalidate_prefix`).
1091        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; // allow-unwrap: test-only
1097        }
1098
1099        /// Pre-seed a key so `get` returns a HIT.
1100        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        /// The TTL passed to the most recent `set` call.
1117        pub async fn last_set_ttl(&self) -> Option<Duration> {
1118            *self.last_set_ttl.lock().await
1119        }
1120
1121        /// Inspect the entry currently stored for `key` (if any).
1122        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            // Read first, then yield: concurrent same-key lookups that
1138            // reach `get` before a write-back all observe the miss
1139            // (without the yield the uncontended tokio Mutex fast path
1140            // completes without rescheduling, serializing "concurrent"
1141            // exchanges and hiding the per-exchange behavior under test).
1142            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() // allow-unwrap: test-only
1209        }
1210    }
1211}
1212
1213// ===========================================================================
1214// Tests
1215// ===========================================================================
1216
1217#[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    /// Minimal no-op RuntimeObservability for tests that don't need OTel.
1256    #[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    // ── Scripted on-miss sub-pipeline ──
1277
1278    #[derive(Clone)]
1279    enum ScriptedOutcome {
1280        Complete,
1281        Stop,
1282        Fail(CamelError),
1283    }
1284
1285    /// Test sub-pipeline: optionally replaces the body, then returns a
1286    /// scripted outcome. Records whether it was invoked.
1287    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            // clone_box is required by the trait but unused by these tests.
1296            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    // ── Builders ──
1320
1321    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    /// Build a CacheService whose on-miss sets `body` and returns `outcome`.
1334    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    // ── Test 1: cache HIT short-circuits, on_miss NOT executed ──
1373
1374    #[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    // ── Test 2: cache MISS runs on_miss, writes back, continues ──
1414
1415    #[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    // ── Test 3: oversized materialized body skips write-back ──
1448
1449    #[tokio::test]
1450    async fn cache_miss_oversized_materialized_body_skips_writeback() {
1451        let repo = Arc::new(MockCacheRepository::new("mock"));
1452        // max_entry_bytes = 4; on_miss produces 9 bytes.
1453        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        // Body passes through unchanged.
1470        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    // ── Test 4: oversized Stream propagates StreamLimitExceeded ──
1480
1481    #[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    // ── Test 5: on_miss Stopped propagates without write-back ──
1510
1511    #[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    // ── Test 6: on_miss Err propagates without write-back ──
1538
1539    #[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    // ── Test 7: repository get Err propagates ──
1568
1569    #[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    // ── Test 8: repository set Err propagates ──
1598
1599    #[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    // ── Test 9: None key bypasses to on_miss, no set ──
1625
1626    #[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    // ── Extra: HIT reconstruction for each ContentType ──
1658
1659    #[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    // ── Extra: Stream body write-back materializes into Body::Bytes ──
1725
1726    #[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        // Stream is replaced by materialized Bytes.
1746        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    // ── CachePeekStaleService tests ──
1754
1755    #[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    // --- Tracing capture helper for debug-log assertions ---
1913
1914    // ── log-capture harness ──────────────────────────────────────────────────
1915
1916    /// Records dispatched events from this module as `LEVEL field=value…`
1917    /// lines. `on_event` fires only for events actually dispatched to this
1918    /// layer: unlike fmt-writer capture it has no registration-time side
1919    /// effects and does not depend on the process-wide callsite interest
1920    /// cache, which parallel tests rebuild concurrently (bd rc-pna5).
1921    #[derive(Default)]
1922    struct EventRecorder {
1923        records: Arc<Mutex<Vec<String>>>,
1924    }
1925
1926    impl EventRecorder {
1927        /// Installs the recorder as this thread's default subscriber and
1928        /// returns the shared record list plus the dispatcher guard.
1929        fn install(self) -> (Arc<Mutex<Vec<String>>>, tracing::subscriber::DefaultGuard) {
1930            use tracing_subscriber::prelude::*;
1931            // OnceLock-gated global registry: heals/prevents callsite-
1932            // interest poisoning of the shared `cache_peek_stale` debug
1933            // callsites (`cache_eip.rs:258/270/457`), which subscriber-less
1934            // sibling cache tests in this binary hit first
1935            // (fix pattern: c3853198; bd rc-img5).
1936            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    /// Formats each visited field as `name=debug-value`, preserving str
1947    /// quoting so `step="cache_peek_stale"`-style assertions keep working.
1948    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); // allow-unwrap: test-only
1976        }
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"); // allow-unwrap: test-only
1992        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"); // allow-unwrap: test-only
2028        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    // ── CacheInvalidateService tests ──
2052
2053    #[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    // ── OTel metrics tests ──
2122
2123    /// Records every `record_counter` call for test assertions.
2124    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    // ── CacheClearService tests ──
2250
2251    #[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    // ── CacheStatsService tests ──
2299
2300    #[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        // Exact key-set assertion: the stats JSON snapshot contract is frozen
2333        // to the eight canonical keys (bd rc-22wj) — no extras, none missing.
2334        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    // ── Peek/invalidate OTel counter tests ──
2362
2363    #[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        // Absent key -> MISS path must NOT emit camel.cache.peek_stale_served.
2435        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        // Failing invalidate -> Failed outcome must NOT emit camel.cache.invalidations.
2458        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    // ── CacheInvalidateService prefix tests ──
2491
2492    #[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    // ── Task 2.4: singleflight miss coalescing ──
2605
2606    use std::sync::atomic::AtomicUsize;
2607    use tokio::sync::Notify;
2608
2609    /// Terminal behavior of the gated on-miss sub-pipeline.
2610    #[derive(Clone)]
2611    enum GatedOutcome {
2612        Complete(Body),
2613        Fail(CamelError),
2614        Stop,
2615    }
2616
2617    /// Test on-miss sub-pipeline with two Notify gates: signals
2618    /// `leader_entered` as its FIRST action, parks on `release.notified()`,
2619    /// then bumps the invocation counter and returns the scripted outcome.
2620    /// Both gates are optional so the same struct serves ungated tests.
2621    #[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                // FIRST action: the leader is inside on_miss.
2644                if let Some(entered) = leader_entered.as_ref() {
2645                    entered.notify_one();
2646                }
2647                // Park until the test releases the wave.
2648                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    /// Build a coalescing CacheService around a `GatedOnMiss`.
2665    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        // Register on leader_entered BEFORE spawning the leader so its
2700        // notify_one cannot be missed (enable-before-check).
2701        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; // deterministic proof: the leader is inside on_miss.
2708
2709        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        // Both waiters park (registered on the wave, not running on_miss).
2715        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        // Drop the leader future mid-flight: the cancellation guard must
2886        // publish a Failed terminal and retire the wave's map entry.
2887        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        // Default construction (no with_coalesce): per-exchange execution.
2920        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}