Skip to main content

camel_processor/
loop_eip.rs

1use std::future::Future;
2use std::pin::Pin;
3use std::task::{Context, Poll};
4
5use tower::{Service, ServiceExt};
6
7use camel_api::loop_eip::{LoopConfig, LoopMode};
8use camel_api::{BoxProcessor, CamelError, Exchange, Value};
9
10pub const CAMEL_LOOP_INDEX: &str = "CamelLoopIndex";
11pub const CAMEL_LOOP_SIZE: &str = "CamelLoopSize";
12
13#[derive(Clone)]
14pub struct LoopService {
15    config: LoopConfig,
16    sub_pipeline: BoxProcessor,
17}
18
19impl LoopService {
20    pub fn new(config: LoopConfig, sub_pipeline: BoxProcessor) -> Self {
21        Self {
22            config,
23            sub_pipeline,
24        }
25    }
26}
27
28impl Service<Exchange> for LoopService {
29    type Response = Exchange;
30    type Error = CamelError;
31    type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
32
33    fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
34        Poll::Ready(Ok(()))
35    }
36
37    fn call(&mut self, mut exchange: Exchange) -> Self::Future {
38        let config = self.config.clone();
39        let mut pipeline = self.sub_pipeline.clone();
40
41        Box::pin(async move {
42            match config.mode {
43                LoopMode::Count(n) => {
44                    // clamp to config.max_iterations. The audit finding was that
45                    // Count(n) had no cap (only While did); a `count: u32::MAX`
46                    // would exhaust CPU. The clamp is applied uniformly to Count
47                    // and While. The CAMEL_LOOP_SIZE property is set to the
48                    // *clamped* count so downstream steps observe the effective
49                    // iteration count.
50                    let n_clamped = n.min(config.max_iterations);
51                    if n > config.max_iterations {
52                        tracing::warn!(
53                            requested = n,
54                            clamped_to = config.max_iterations,
55                            "LoopMode::Count exceeded max_iterations; clamping"
56                        );
57                    }
58                    exchange.set_property(CAMEL_LOOP_SIZE, Value::from(n_clamped as u64));
59                    for i in 0..n_clamped {
60                        exchange.set_property(CAMEL_LOOP_INDEX, Value::from(i as u64));
61                        exchange = pipeline.ready().await?.call(exchange).await?;
62                    }
63                }
64                LoopMode::While(ref predicate) => {
65                    exchange.set_property(CAMEL_LOOP_SIZE, Value::from(0u64));
66                    // Every predicate evaluation error propagates (never
67                    // loop-end success), including the post-exhaustion
68                    // safety-guard probe. `exhausted` tracks whether that
69                    // re-check applies so a side-effectful predicate is not
70                    // re-evaluated after a normal false exit.
71                    let mut exhausted = true;
72                    for i in 0..config.max_iterations {
73                        match predicate.matches(&exchange).await {
74                            Ok(true) => {}
75                            Ok(false) => {
76                                exhausted = false;
77                                break;
78                            }
79                            Err(err) => return Err(err),
80                        }
81                        exchange.set_property(CAMEL_LOOP_INDEX, Value::from(i as u64));
82                        exchange = pipeline.ready().await?.call(exchange).await?;
83                    }
84                    if exhausted {
85                        match predicate.matches(&exchange).await {
86                            Ok(true) => {
87                                tracing::warn!(
88                                    "Loop while-mode hit max_iterations ({}) safety guard. Predicate still true.",
89                                    config.max_iterations
90                                );
91                            }
92                            Ok(false) => {}
93                            Err(err) => return Err(err),
94                        }
95                    }
96                }
97                // Future loop modes: no iteration, pass exchange through unchanged.
98                _ => {}
99            }
100            Ok(exchange)
101        })
102    }
103}
104
105// ── LoopSegment (ADR-0025 OutcomePipeline) ─────────────────────────────
106
107/// Outcome-aware structural EIP segment for the Loop pattern.
108///
109/// Operates at the `PipelineOutcome` layer so that `Stopped(ex)` from a
110/// sub-step (e.g. Stop EIP) is preserved with the exchange including all
111/// mutations. Supports both Count and While modes, mirroring `LoopService`
112/// semantics exactly.
113///
114/// Unlike `LoopService` (which operates at the Tower layer), `LoopSegment`
115/// correctly short-circuits on `PipelineOutcome::Stopped` or `Failed`.
116pub struct LoopSegment {
117    pub config: camel_api::loop_eip::LoopConfig,
118    pub body: camel_api::OutcomeSegment,
119}
120
121impl Clone for LoopSegment {
122    fn clone(&self) -> Self {
123        Self {
124            config: self.config.clone(),
125            body: self.body.clone(),
126        }
127    }
128}
129
130impl camel_api::OutcomePipeline for LoopSegment {
131    fn clone_box(&self) -> Box<dyn camel_api::OutcomePipeline> {
132        Box::new(self.clone())
133    }
134
135    fn run<'a>(
136        &'a mut self,
137        exchange: camel_api::Exchange,
138    ) -> Pin<Box<dyn Future<Output = camel_api::PipelineOutcome> + Send + 'a>> {
139        use camel_api::{PipelineOutcome, Value};
140
141        let config = self.config.clone();
142        let body = &mut self.body;
143
144        Box::pin(async move {
145            match config.mode {
146                camel_api::loop_eip::LoopMode::Count(n) => {
147                    let n_clamped = n.min(config.max_iterations);
148                    if n > config.max_iterations {
149                        tracing::warn!(
150                            requested = n,
151                            clamped_to = config.max_iterations,
152                            "LoopMode::Count exceeded max_iterations; clamping"
153                        );
154                    }
155                    let mut ex = exchange;
156                    ex.set_property(CAMEL_LOOP_SIZE, Value::from(n_clamped as u64));
157                    for i in 0..n_clamped {
158                        ex.set_property(CAMEL_LOOP_INDEX, Value::from(i as u64));
159                        match body.run(ex).await {
160                            PipelineOutcome::Completed(next) => {
161                                ex = next;
162                            }
163                            other => return other,
164                        }
165                    }
166                    PipelineOutcome::Completed(ex)
167                }
168                camel_api::loop_eip::LoopMode::While(ref predicate) => {
169                    let mut ex = exchange;
170                    ex.set_property(CAMEL_LOOP_SIZE, Value::from(0u64));
171                    // Every predicate evaluation error propagates as
172                    // `Failed(err)` (never as loop-end Completed), including
173                    // the post-exhaustion safety-guard probe.
174                    let mut i = 0u64;
175                    let mut exhausted = true;
176                    while i < config.max_iterations as u64 {
177                        match predicate.matches(&ex).await {
178                            Ok(true) => {}
179                            Ok(false) => {
180                                exhausted = false;
181                                break;
182                            }
183                            Err(err) => return PipelineOutcome::Failed(err),
184                        }
185                        ex.set_property(CAMEL_LOOP_INDEX, Value::from(i));
186                        match body.run(ex).await {
187                            PipelineOutcome::Completed(next) => {
188                                ex = next;
189                            }
190                            other => return other,
191                        }
192                        i += 1;
193                    }
194                    if exhausted {
195                        match predicate.matches(&ex).await {
196                            Ok(true) => {
197                                tracing::warn!(
198                                    "Loop while-mode hit max_iterations ({}) safety guard. Predicate still true.",
199                                    config.max_iterations
200                                );
201                            }
202                            Ok(false) => {}
203                            Err(err) => return PipelineOutcome::Failed(err),
204                        }
205                    }
206                    PipelineOutcome::Completed(ex)
207                }
208                // Future loop modes: no iteration, complete with the exchange unchanged.
209                _ => PipelineOutcome::Completed(exchange),
210            }
211        })
212    }
213}
214
215#[cfg(test)]
216mod tests {
217    use std::sync::atomic::{AtomicUsize, Ordering};
218    use std::sync::{Arc, Mutex};
219
220    use camel_api::loop_eip::{LoopConfig, LoopMode, MAX_LOOP_ITERATIONS};
221    use camel_api::{
222        Body, BoxProcessor, BoxProcessorExt, CamelError, Exchange, FilterPredicate,
223        IdentityProcessor, Message,
224    };
225    use tower::{Service, ServiceExt};
226
227    use super::{CAMEL_LOOP_INDEX, CAMEL_LOOP_SIZE, LoopSegment, LoopService};
228
229    fn identity_pipeline() -> BoxProcessor {
230        BoxProcessor::new(IdentityProcessor)
231    }
232
233    fn counter_pipeline(counter: Arc<AtomicUsize>) -> BoxProcessor {
234        BoxProcessor::from_fn(move |exchange: Exchange| {
235            let counter = Arc::clone(&counter);
236            Box::pin(async move {
237                counter.fetch_add(1, Ordering::SeqCst);
238                Ok(exchange)
239            })
240        })
241    }
242
243    #[tokio::test]
244    async fn test_loop_count_iterates_n_times() {
245        let counter = Arc::new(AtomicUsize::new(0));
246        let config = LoopConfig::new(LoopMode::Count(3));
247        let mut service = LoopService::new(config, counter_pipeline(Arc::clone(&counter)));
248
249        let exchange = Exchange::new(Message::new("test"));
250        let result = service.ready().await.unwrap().call(exchange).await;
251
252        assert!(result.is_ok());
253        assert_eq!(counter.load(Ordering::SeqCst), 3);
254    }
255
256    #[tokio::test]
257    async fn test_loop_count_sets_properties() {
258        let seen_indices = Arc::new(Mutex::new(Vec::<u64>::new()));
259        let seen_indices_for_pipeline = Arc::clone(&seen_indices);
260
261        let pipeline = BoxProcessor::from_fn(move |exchange: Exchange| {
262            let seen_indices = Arc::clone(&seen_indices_for_pipeline);
263            Box::pin(async move {
264                if let Some(index) = exchange.property(CAMEL_LOOP_INDEX).and_then(|v| v.as_u64()) {
265                    seen_indices.lock().unwrap().push(index);
266                }
267                Ok(exchange)
268            })
269        });
270
271        let config = LoopConfig::new(LoopMode::Count(3));
272        let mut service = LoopService::new(config, pipeline);
273
274        let exchange = Exchange::new(Message::new("test"));
275        let result = service.ready().await.unwrap().call(exchange).await.unwrap();
276
277        assert_eq!(*seen_indices.lock().unwrap(), vec![0, 1, 2]);
278        assert_eq!(
279            result.property(CAMEL_LOOP_SIZE).and_then(|v| v.as_u64()),
280            Some(3)
281        );
282    }
283
284    #[tokio::test]
285    async fn test_loop_count_zero_is_noop() {
286        let config = LoopConfig::new(LoopMode::Count(0));
287        let mut service = LoopService::new(config, identity_pipeline());
288
289        let exchange = Exchange::new(Message::new("test"));
290        let result = service.ready().await.unwrap().call(exchange).await.unwrap();
291
292        assert_eq!(result.input.body.as_text(), Some("test"));
293        assert_eq!(
294            result.property(CAMEL_LOOP_SIZE).and_then(|v| v.as_u64()),
295            Some(0)
296        );
297        assert!(result.property(CAMEL_LOOP_INDEX).is_none());
298    }
299
300    #[tokio::test]
301    async fn test_loop_while_stops_when_predicate_false() {
302        let counter = Arc::new(AtomicUsize::new(0));
303
304        let predicate = FilterPredicate::new(|exchange: &Exchange| {
305            exchange
306                .property("iterations")
307                .and_then(|v| v.as_u64())
308                .unwrap_or(0)
309                < 2
310        });
311
312        let counter_for_pipeline = Arc::clone(&counter);
313        let pipeline = BoxProcessor::from_fn(move |mut exchange: Exchange| {
314            let counter = Arc::clone(&counter_for_pipeline);
315            Box::pin(async move {
316                let current = exchange
317                    .property("iterations")
318                    .and_then(|v| v.as_u64())
319                    .unwrap_or(0);
320                exchange.set_property("iterations", current + 1);
321                counter.fetch_add(1, Ordering::SeqCst);
322                Ok(exchange)
323            })
324        });
325
326        let config = LoopConfig::new(LoopMode::While(PredicateSource::Sync(predicate)));
327        let mut service = LoopService::new(config, pipeline);
328
329        let exchange = Exchange::new(Message::new("test"));
330        let result = service.ready().await.unwrap().call(exchange).await.unwrap();
331
332        assert_eq!(counter.load(Ordering::SeqCst), 2);
333        assert_eq!(
334            result.property("iterations").and_then(|v| v.as_u64()),
335            Some(2)
336        );
337        assert_eq!(
338            result.property(CAMEL_LOOP_INDEX).and_then(|v| v.as_u64()),
339            Some(1)
340        );
341        assert_eq!(
342            result.property(CAMEL_LOOP_SIZE).and_then(|v| v.as_u64()),
343            Some(0)
344        );
345    }
346
347    #[tokio::test]
348    async fn test_loop_while_respects_max_iterations() {
349        let counter = Arc::new(AtomicUsize::new(0));
350        let predicate = FilterPredicate::new(|_exchange: &Exchange| true);
351        let config = LoopConfig::new(LoopMode::While(PredicateSource::Sync(predicate)));
352        let mut service = LoopService::new(config, counter_pipeline(Arc::clone(&counter)));
353
354        let exchange = Exchange::new(Message::new("test"));
355        let result = service.ready().await.unwrap().call(exchange).await;
356
357        assert!(result.is_ok());
358        assert_eq!(counter.load(Ordering::SeqCst), MAX_LOOP_ITERATIONS);
359    }
360
361    #[tokio::test]
362    async fn test_loop_error_propagation() {
363        let pipeline = BoxProcessor::from_fn(|_exchange: Exchange| {
364            Box::pin(async { Err(CamelError::ProcessorError("boom".into())) })
365        });
366
367        let config = LoopConfig::new(LoopMode::Count(3));
368        let mut service = LoopService::new(config, pipeline);
369
370        let exchange = Exchange::new(Message::new("test"));
371        let result = service.ready().await.unwrap().call(exchange).await;
372
373        assert!(matches!(result, Err(CamelError::ProcessorError(msg)) if msg == "boom"));
374    }
375
376    // ── D-M9 Batch 1: LoopMode::Count(n) clamped to MAX_LOOP_ITERATIONS ──
377
378    /// D-M9: a `Count(u32::MAX as usize)` (or any value above
379    /// `MAX_LOOP_ITERATIONS`) is clamped to `MAX_LOOP_ITERATIONS` and
380    /// runs at most that many iterations. Without the clamp, a malicious
381    /// or typo'd `count: 4294967295` would exhaust CPU. The pre-existing
382    /// `Count(3)` test continues to pass — small values are unaffected.
383    #[tokio::test]
384    async fn test_loop_count_clamped_to_max_iterations() {
385        let counter = Arc::new(AtomicUsize::new(0));
386        let config = LoopConfig::new(LoopMode::Count(usize::MAX));
387        let mut service = LoopService::new(config, counter_pipeline(Arc::clone(&counter)));
388
389        let exchange = Exchange::new(Message::new("test"));
390        let result = service.ready().await.unwrap().call(exchange).await;
391
392        assert!(result.is_ok());
393        // Must run exactly MAX_LOOP_ITERATIONS, not u32::MAX iterations.
394        assert_eq!(counter.load(Ordering::SeqCst), MAX_LOOP_ITERATIONS);
395        // The CAMEL_LOOP_SIZE property is set to the *clamped* count.
396        assert_eq!(
397            result
398                .unwrap()
399                .property(CAMEL_LOOP_SIZE)
400                .and_then(|v| v.as_u64()),
401            Some(MAX_LOOP_ITERATIONS as u64)
402        );
403    }
404
405    #[tokio::test]
406    async fn test_loop_count_with_custom_max_iterations() {
407        let counter = Arc::new(AtomicUsize::new(0));
408        let config = LoopConfig::new(LoopMode::Count(15_000)).with_max_iterations(15_000);
409        let mut service = LoopService::new(config, counter_pipeline(Arc::clone(&counter)));
410        let exchange = Exchange::new(Message::new("test"));
411        let _ = service.ready().await.unwrap().call(exchange).await.unwrap();
412        assert_eq!(counter.load(Ordering::SeqCst), 15_000);
413    }
414
415    #[tokio::test]
416    async fn test_loop_while_with_custom_max_iterations() {
417        let counter = Arc::new(AtomicUsize::new(0));
418        let predicate = FilterPredicate::new(|_| true);
419        let config = LoopConfig::new(LoopMode::While(PredicateSource::Sync(predicate)))
420            .with_max_iterations(50);
421        let mut service = LoopService::new(config, counter_pipeline(Arc::clone(&counter)));
422        let exchange = Exchange::new(Message::new("test"));
423        let _ = service.ready().await.unwrap().call(exchange).await.unwrap();
424        assert_eq!(counter.load(Ordering::SeqCst), 50);
425    }
426
427    #[tokio::test]
428    async fn test_loop_pipeline_chaining() {
429        let pipeline = BoxProcessor::from_fn(|mut exchange: Exchange| {
430            Box::pin(async move {
431                if let Body::Text(s) = &exchange.input.body {
432                    exchange.input.body = Body::Text(format!("{s}x"));
433                }
434                Ok(exchange)
435            })
436        });
437
438        let config = LoopConfig::new(LoopMode::Count(3));
439        let mut service = LoopService::new(config, pipeline);
440
441        let exchange = Exchange::new(Message::new("start"));
442        let result = service.ready().await.unwrap().call(exchange).await.unwrap();
443
444        assert_eq!(result.input.body.as_text(), Some("startxxx"));
445    }
446
447    // ── Fallible predicate path (language-value-boundary task 1.4) ──
448
449    use camel_api::{ExpressionErrorClass, PredicateSource};
450
451    fn expression_failed() -> CamelError {
452        CamelError::ExpressionFailed {
453            language: "rhai".to_string(),
454            route_id: "r1".to_string(),
455            step_id: "step#0".to_string(),
456            verb: "loop".to_string(),
457            class: ExpressionErrorClass::Runtime,
458            position: None,
459            conversion: None,
460            cause: None,
461        }
462    }
463
464    fn async_err_predicate(err: CamelError) -> PredicateSource {
465        PredicateSource::Async(Arc::new(move |_: &Exchange| {
466            let err = err.clone();
467            Box::pin(async move { Err(err) }) as camel_api::BoxBoolFuture
468        }))
469    }
470
471    /// Async predicate that returns `Ok(true)` on its first call and
472    /// `Err(expression_failed())` afterwards. Paired with
473    /// `max_iterations = 1` it drives the loop into the exhaustion
474    /// safety-guard probe on the second evaluation.
475    fn exhausted_then_err_predicate(calls: Arc<AtomicUsize>) -> PredicateSource {
476        PredicateSource::Async(Arc::new(move |_: &Exchange| {
477            let call = calls.fetch_add(1, Ordering::SeqCst);
478            let err = expression_failed();
479            Box::pin(async move { if call == 0 { Ok(true) } else { Err(err) } })
480                as camel_api::BoxBoolFuture
481        }))
482    }
483
484    /// OutcomePipeline body that always completes with the exchange unchanged.
485    #[derive(Clone)]
486    struct CompletedBody;
487
488    impl camel_api::OutcomePipeline for CompletedBody {
489        fn clone_box(&self) -> Box<dyn camel_api::OutcomePipeline> {
490            Box::new(self.clone())
491        }
492
493        fn run<'a>(
494            &'a mut self,
495            exchange: Exchange,
496        ) -> std::pin::Pin<
497            Box<dyn std::future::Future<Output = camel_api::PipelineOutcome> + Send + 'a>,
498        > {
499            Box::pin(async move { camel_api::PipelineOutcome::Completed(exchange) })
500        }
501    }
502
503    #[tokio::test]
504    async fn loop_while_predicate_error_fails_loop() {
505        let counter = Arc::new(AtomicUsize::new(0));
506        let config = LoopConfig::new(LoopMode::While(async_err_predicate(expression_failed())));
507        let mut service = LoopService::new(config, counter_pipeline(Arc::clone(&counter)));
508
509        let exchange = Exchange::new(Message::new("test"));
510        let result = service.ready().await.unwrap().call(exchange).await;
511
512        match result {
513            Err(CamelError::ExpressionFailed { .. }) => {}
514            other => panic!("expected Err(ExpressionFailed), got {other:?}"),
515        }
516        assert_eq!(
517            counter.load(Ordering::SeqCst),
518            0,
519            "loop body must not run when the while predicate errors"
520        );
521    }
522
523    #[tokio::test]
524    async fn loop_guard_predicate_error_propagates() {
525        let calls = Arc::new(AtomicUsize::new(0));
526        let counter = Arc::new(AtomicUsize::new(0));
527        let config = LoopConfig::new(LoopMode::While(exhausted_then_err_predicate(Arc::clone(
528            &calls,
529        ))))
530        .with_max_iterations(1);
531        let mut service = LoopService::new(config, counter_pipeline(Arc::clone(&counter)));
532
533        let exchange = Exchange::new(Message::new("test"));
534        let result = service.ready().await.unwrap().call(exchange).await;
535
536        match result {
537            Err(CamelError::ExpressionFailed { .. }) => {}
538            other => panic!("expected Err(ExpressionFailed), got {other:?}"),
539        }
540        assert_eq!(
541            counter.load(Ordering::SeqCst),
542            1,
543            "body runs once before the exhaustion guard probe errors"
544        );
545        assert_eq!(
546            calls.load(Ordering::SeqCst),
547            2,
548            "predicate evaluated once in-loop and once by the safety guard"
549        );
550    }
551
552    #[tokio::test]
553    async fn loop_guard_predicate_error_propagates_segment() {
554        let calls = Arc::new(AtomicUsize::new(0));
555        let config = LoopConfig::new(LoopMode::While(exhausted_then_err_predicate(Arc::clone(
556            &calls,
557        ))))
558        .with_max_iterations(1);
559        let mut segment = LoopSegment {
560            config,
561            body: camel_api::OutcomeSegment::new(Box::new(CompletedBody)),
562        };
563
564        let exchange = Exchange::new(Message::new("test"));
565        let result = camel_api::OutcomePipeline::run(&mut segment, exchange).await;
566
567        match result {
568            camel_api::PipelineOutcome::Failed(CamelError::ExpressionFailed { .. }) => {}
569            other => panic!("expected Failed(ExpressionFailed), got {other:?}"),
570        }
571        assert_eq!(
572            calls.load(Ordering::SeqCst),
573            2,
574            "predicate evaluated once in-loop and once by the safety guard"
575        );
576    }
577}