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 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 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 _ => {}
99 }
100 Ok(exchange)
101 })
102 }
103}
104
105pub 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 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 _ => 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 #[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 assert_eq!(counter.load(Ordering::SeqCst), MAX_LOOP_ITERATIONS);
395 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 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 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 #[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}