1use std::future::Future;
2use std::pin::Pin;
3use std::task::{Context, Poll};
4
5use tokio::task::JoinSet;
6use tower::Service;
7use tower::ServiceExt;
8
9use camel_api::endpoint_pipeline::{CAMEL_SLIP_ENDPOINT, EndpointPipelineConfig};
10use camel_api::recipient_list::RecipientListConfig;
11use camel_api::{Body, CamelError, Exchange, Value};
12
13use crate::endpoint_pipeline::EndpointPipelineService;
14
15#[derive(Clone)]
16pub struct RecipientListService {
17 config: RecipientListConfig,
18 pipeline: EndpointPipelineService,
19}
20
21impl RecipientListService {
22 pub fn new(
23 config: RecipientListConfig,
24 endpoint_resolver: camel_api::EndpointResolver,
25 ) -> Result<Self, CamelError> {
26 config.validate()?;
27 let pipeline_config = EndpointPipelineConfig {
28 cache_size: EndpointPipelineConfig::from_signed(1000),
29 ignore_invalid_endpoints: false,
30 };
31 Ok(Self {
32 config,
33 pipeline: EndpointPipelineService::new(endpoint_resolver, pipeline_config),
34 })
35 }
36}
37
38impl Service<Exchange> for RecipientListService {
39 type Response = Exchange;
40 type Error = CamelError;
41 type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
42
43 fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
44 Poll::Ready(Ok(()))
45 }
46
47 fn call(&mut self, mut exchange: Exchange) -> Self::Future {
48 let config = self.config.clone();
49 let pipeline = self.pipeline.clone();
50
51 Box::pin(async move {
52 let uris_raw = config.expression.resolve(&exchange).await?;
53 if uris_raw.is_empty() {
54 return Ok(exchange);
55 }
56
57 let cap = config.max_recipients;
63 let uris: Vec<&str> = uris_raw
64 .split(&config.delimiter)
65 .map(|s| s.trim())
66 .filter(|s| !s.is_empty())
67 .take(cap)
68 .collect();
69 if uris.is_empty() {
70 return Ok(exchange);
71 }
72
73 if config.parallel {
74 let original_for_aggregate = exchange.clone();
75 let mut endpoints_to_call = Vec::with_capacity(uris.len());
76 for uri in &uris {
77 if let Some(endpoint) = pipeline.resolve(uri)? {
78 endpoints_to_call.push((uri.to_string(), endpoint));
79 }
80 }
81
82 let mut results: Vec<Exchange> = Vec::with_capacity(endpoints_to_call.len());
83 let mut join_set = JoinSet::new();
84 let mut iter = endpoints_to_call.into_iter();
85 let raw_limit = config.parallel_limit.unwrap_or(results.capacity());
86 let limit = raw_limit.max(1).min(results.capacity().max(1));
87 let mut last_parallel_error: Option<CamelError> = None;
88
89 for _ in 0..limit {
90 if let Some((uri, mut endpoint)) = iter.next() {
91 let mut cloned = original_for_aggregate.clone();
92 cloned.set_property(CAMEL_SLIP_ENDPOINT, Value::String(uri));
93 join_set.spawn(async move { endpoint.ready().await?.call(cloned).await });
94 }
95 }
96
97 while let Some(result) = join_set.join_next().await {
98 match result {
99 Ok(Ok(ex)) => results.push(ex),
100 Ok(Err(e)) if config.stop_on_exception => {
101 join_set.abort_all();
102 return Err(e);
103 }
104 Ok(Err(e)) => {
105 last_parallel_error = Some(e);
109 }
110 Err(join_err) if join_err.is_panic() => {
111 last_parallel_error = Some(CamelError::ProcessorError(format!(
118 "recipient task panicked: {join_err}"
119 )));
120 }
121 Err(_) => {} }
123
124 if let Some((uri, mut endpoint)) = iter.next() {
125 let mut cloned = original_for_aggregate.clone();
126 cloned.set_property(CAMEL_SLIP_ENDPOINT, Value::String(uri));
127 join_set.spawn(async move { endpoint.ready().await?.call(cloned).await });
128 }
129 }
130
131 let zero_success_error = if results.is_empty() {
135 last_parallel_error
136 } else {
137 None
138 };
139 if let Some(err) = zero_success_error {
140 return Err(err);
141 }
142
143 exchange = aggregate_results(config.strategy, original_for_aggregate, results);
144 } else {
145 let mut results: Vec<Exchange> = Vec::new();
146 let mut last_error: Option<CamelError> = None;
147 let original_for_aggregate = exchange.clone();
148 for uri in &uris {
149 let endpoint = match pipeline.resolve(uri)? {
150 Some(e) => e,
151 None => continue,
152 };
153 exchange.set_property(CAMEL_SLIP_ENDPOINT, Value::String(uri.to_string()));
154 let mut endpoint = endpoint;
155 let result = endpoint.ready().await?.call(exchange.clone()).await;
156 match result {
157 Ok(ex) => {
158 results.push(ex.clone());
159 exchange = ex;
160 }
161 Err(e) if config.stop_on_exception => return Err(e),
162 Err(e) => {
163 last_error = Some(e);
166 continue;
167 }
168 }
169 }
170 let zero_success_error = if results.is_empty() { last_error } else { None };
175 if let Some(err) = zero_success_error {
176 return Err(err);
177 }
178 exchange = aggregate_results(config.strategy, original_for_aggregate, results);
179 }
180
181 Ok(exchange)
182 })
183 }
184}
185
186fn aggregate_results(
187 strategy: camel_api::MulticastStrategy,
188 original: Exchange,
189 results: Vec<Exchange>,
190) -> Exchange {
191 match strategy {
192 camel_api::MulticastStrategy::LastWins => results.into_iter().last().unwrap_or(original),
193 camel_api::MulticastStrategy::CollectAll => {
194 let bodies: Vec<Value> = results
195 .iter()
196 .map(|ex| match &ex.input.body {
197 Body::Text(s) => Value::String(s.clone()),
198 Body::Json(v) => v.clone(),
199 Body::Xml(s) => Value::String(s.clone()),
200 Body::Bytes(b) => Value::String(String::from_utf8_lossy(b).into_owned()),
201 Body::Stream(s) => serde_json::json!({
202 "_stream": {
203 "origin": s.metadata.origin,
204 "placeholder": true,
205 "hint": "Materialize exchange body with .into_bytes() before recipient-list aggregation"
206 }
207 }),
208 _ => Value::Null,
210 })
211 .collect();
212 let mut result = results.into_iter().last().unwrap_or(original);
213 result.input.body = camel_api::Body::from(Value::Array(bodies));
214 result
215 }
216 camel_api::MulticastStrategy::Custom(fn_) => {
217 results.into_iter().fold(original, |acc, ex| fn_(acc, ex))
218 }
219 _ => original,
221 }
222}
223
224#[cfg(test)]
225mod tests {
226 use super::*;
227 use camel_api::MulticastStrategy;
228 use camel_api::{BoxProcessor, BoxProcessorExt, CamelError, Message};
229 use std::collections::HashMap;
230 use std::sync::Arc;
231 use std::sync::atomic::{AtomicUsize, Ordering};
232 use std::time::{Duration, Instant};
233 use tokio::sync::Mutex;
234 use tokio::time::sleep;
235
236 fn mock_resolver() -> camel_api::EndpointResolver {
237 Arc::new(|uri: &str| {
238 if uri.starts_with("mock:") {
239 Some(BoxProcessor::from_fn(|ex| Box::pin(async move { Ok(ex) })))
240 } else {
241 None
242 }
243 })
244 }
245
246 #[tokio::test]
247 async fn recipient_list_single_destination() {
248 let call_count = Arc::new(AtomicUsize::new(0));
249 let count_clone = call_count.clone();
250
251 let resolver = Arc::new(move |uri: &str| {
252 if uri == "mock:a" {
253 let count = count_clone.clone();
254 Some(BoxProcessor::from_fn(move |ex| {
255 count.fetch_add(1, Ordering::SeqCst);
256 Box::pin(async move { Ok(ex) })
257 }))
258 } else {
259 None
260 }
261 });
262
263 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
264 |_ex: &Exchange| "mock:a".to_string(),
265 )));
266
267 let mut svc = RecipientListService::new(config, resolver).unwrap();
268 let ex = Exchange::new(Message::new("test"));
269 let result = svc.ready().await.unwrap().call(ex).await;
270
271 assert!(result.is_ok());
272 assert_eq!(call_count.load(Ordering::SeqCst), 1);
273 }
274
275 #[tokio::test]
276 async fn recipient_list_multiple_destinations() {
277 let call_count = Arc::new(AtomicUsize::new(0));
278 let count_clone = call_count.clone();
279
280 let resolver = Arc::new(move |uri: &str| {
281 if uri.starts_with("mock:") {
282 let count = count_clone.clone();
283 Some(BoxProcessor::from_fn(move |ex| {
284 count.fetch_add(1, Ordering::SeqCst);
285 Box::pin(async move { Ok(ex) })
286 }))
287 } else {
288 None
289 }
290 });
291
292 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
293 |_ex: &Exchange| "mock:a,mock:b,mock:c".to_string(),
294 )));
295
296 let mut svc = RecipientListService::new(config, resolver).unwrap();
297 let ex = Exchange::new(Message::new("test"));
298 let result = svc.ready().await.unwrap().call(ex).await;
299
300 assert!(result.is_ok());
301 assert_eq!(call_count.load(Ordering::SeqCst), 3);
302 }
303
304 #[tokio::test]
305 async fn recipient_list_empty_expression() {
306 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
307 |_ex: &Exchange| String::new(),
308 )));
309
310 let mut svc = RecipientListService::new(config, mock_resolver()).unwrap();
311 let ex = Exchange::new(Message::new("test"));
312 let result = svc.ready().await.unwrap().call(ex).await;
313
314 assert!(result.is_ok());
315 }
316
317 #[tokio::test]
318 async fn recipient_list_invalid_endpoint_error() {
319 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
320 |_ex: &Exchange| "invalid:endpoint".to_string(),
321 )));
322
323 let mut svc = RecipientListService::new(config, mock_resolver()).unwrap();
324 let ex = Exchange::new(Message::new("test"));
325 let result = svc.ready().await.unwrap().call(ex).await;
326
327 assert!(result.is_err());
328 assert!(result.unwrap_err().to_string().contains("Invalid endpoint"));
329 }
330
331 #[tokio::test]
332 async fn recipient_list_custom_delimiter() {
333 use std::sync::Mutex;
334
335 let order: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(Vec::new()));
336
337 let resolver = {
338 let order = order.clone();
339 Arc::new(move |uri: &str| {
340 let order = order.clone();
341 let uri = uri.to_string();
342 Some(BoxProcessor::from_fn(move |ex| {
343 order.lock().unwrap().push(uri.clone());
344 Box::pin(async move { Ok(ex) })
345 }))
346 })
347 };
348
349 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
350 |_ex: &Exchange| "mock:x|mock:y|mock:z".to_string(),
351 )))
352 .delimiter("|");
353
354 let mut svc = RecipientListService::new(config, resolver).unwrap();
355 let ex = Exchange::new(Message::new("test"));
356 svc.ready().await.unwrap().call(ex).await.unwrap();
357
358 let order = order.lock().unwrap();
359 assert_eq!(*order, vec!["mock:x", "mock:y", "mock:z"]);
360 }
361
362 #[tokio::test]
363 async fn recipient_list_expression_evaluated_once() {
364 let expr_count = Arc::new(AtomicUsize::new(0));
365 let expr_count_clone = expr_count.clone();
366
367 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
368 move |_ex: &Exchange| {
369 expr_count_clone.fetch_add(1, Ordering::SeqCst);
370 "mock:a,mock:b".to_string()
371 },
372 )));
373
374 let mut svc = RecipientListService::new(config, mock_resolver()).unwrap();
375 let ex = Exchange::new(Message::new("test"));
376 svc.ready().await.unwrap().call(ex).await.unwrap();
377
378 assert_eq!(
379 expr_count.load(Ordering::SeqCst),
380 1,
381 "Expression must be evaluated exactly once"
382 );
383 }
384
385 #[tokio::test]
386 async fn recipient_list_ignores_empty_uri_tokens() {
387 let call_count = Arc::new(AtomicUsize::new(0));
388 let call_count_clone = call_count.clone();
389
390 let resolver = Arc::new(move |uri: &str| {
391 if uri.starts_with("mock:") {
392 let count = call_count_clone.clone();
393 Some(BoxProcessor::from_fn(move |ex| {
394 count.fetch_add(1, Ordering::SeqCst);
395 Box::pin(async move { Ok(ex) })
396 }))
397 } else {
398 None
399 }
400 });
401
402 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
403 |_ex: &Exchange| " ,mock:a, ,mock:b,, ".to_string(),
404 )));
405
406 let mut svc = RecipientListService::new(config, resolver).unwrap();
407 let ex = Exchange::new(Message::new("test"));
408 let result = svc.ready().await.unwrap().call(ex).await;
409 assert!(result.is_ok());
410 assert_eq!(call_count.load(Ordering::SeqCst), 2);
411 }
412
413 #[tokio::test]
414 async fn recipient_list_mutation_between_steps() {
415 let resolver = Arc::new(|uri: &str| {
416 if uri == "mock:mutate" {
417 Some(BoxProcessor::from_fn(|mut ex| {
418 ex.input.body = camel_api::Body::Text("mutated".to_string());
419 Box::pin(async move { Ok(ex) })
420 }))
421 } else if uri == "mock:verify" {
422 Some(BoxProcessor::from_fn(|ex| {
423 let body = ex.input.body.as_text().unwrap_or("").to_string();
424 assert_eq!(body, "mutated");
425 Box::pin(async move { Ok(ex) })
426 }))
427 } else {
428 None
429 }
430 });
431
432 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
433 |_ex: &Exchange| "mock:mutate,mock:verify".to_string(),
434 )));
435
436 let mut svc = RecipientListService::new(config, resolver).unwrap();
437 let ex = Exchange::new(Message::new("original"));
438 let result = svc.ready().await.unwrap().call(ex).await;
439
440 assert!(result.is_ok());
441 }
442
443 #[tokio::test]
444 async fn recipient_list_parallel_executes_concurrently() {
445 let records: Arc<Mutex<Vec<(String, Instant, Instant)>>> = Arc::new(Mutex::new(Vec::new()));
446
447 let resolver = {
448 let records = records.clone();
449 Arc::new(move |uri: &str| {
450 if uri.starts_with("mock:") {
451 let records = records.clone();
452 let uri = uri.to_string();
453 Some(BoxProcessor::from_fn(move |ex| {
454 let records = records.clone();
455 let uri = uri.clone();
456 Box::pin(async move {
457 let start = Instant::now();
458 sleep(Duration::from_millis(100)).await;
459 let end = Instant::now();
460 records.lock().await.push((uri, start, end));
461 Ok(ex)
462 })
463 }))
464 } else {
465 None
466 }
467 })
468 };
469
470 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
471 |_ex: &Exchange| "mock:a,mock:b,mock:c".to_string(),
472 )))
473 .parallel(true);
474
475 let mut svc = RecipientListService::new(config, resolver).unwrap();
476 let ex = Exchange::new(Message::new("test"));
477 svc.ready().await.unwrap().call(ex).await.unwrap();
478
479 let records = tokio::time::timeout(Duration::from_secs(5), records.lock())
480 .await
481 .expect("records lock timeout");
482 assert_eq!(records.len(), 3);
483
484 let mut overlap_found = false;
485 for i in 0..records.len() {
486 for j in (i + 1)..records.len() {
487 let (_, a_start, a_end) = records[i];
488 let (_, b_start, b_end) = records[j];
489 if a_start < b_end && b_start < a_end {
490 overlap_found = true;
491 break;
492 }
493 }
494 if overlap_found {
495 break;
496 }
497 }
498
499 assert!(overlap_found);
500 }
501
502 #[tokio::test]
503 async fn recipient_list_parallel_stop_on_exception_returns_error() {
504 let resolver = Arc::new(|uri: &str| {
505 if uri == "mock:err" {
506 Some(BoxProcessor::from_fn(|_ex| {
507 Box::pin(async { Err(CamelError::ProcessorError("boom".to_string())) })
508 }))
509 } else if uri.starts_with("mock:") {
510 Some(BoxProcessor::from_fn(|ex| Box::pin(async move { Ok(ex) })))
511 } else {
512 None
513 }
514 });
515
516 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
517 |_ex: &Exchange| "mock:a,mock:err,mock:c".to_string(),
518 )))
519 .parallel(true)
520 .stop_on_exception(true);
521
522 let mut svc = RecipientListService::new(config, resolver).unwrap();
523 let ex = Exchange::new(Message::new("test"));
524 let result = svc.ready().await.unwrap().call(ex).await;
525 assert!(matches!(result, Err(CamelError::ProcessorError(msg)) if msg == "boom"));
526 }
527
528 #[tokio::test]
529 async fn recipient_list_parallel_limit_respects_limit() {
530 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
531 |_ex: &Exchange| "mock:a,mock:b,mock:c,mock:d".to_string(),
532 )))
533 .parallel(true)
534 .parallel_limit(2);
535
536 let resolver = Arc::new(|uri: &str| {
537 if uri.starts_with("mock:") {
538 Some(BoxProcessor::from_fn(|ex| {
539 Box::pin(async move {
540 sleep(Duration::from_millis(100)).await;
541 Ok(ex)
542 })
543 }))
544 } else {
545 None
546 }
547 });
548
549 let mut svc = RecipientListService::new(config, resolver).unwrap();
550 let ex = Exchange::new(Message::new("test"));
551 let start = Instant::now();
552 svc.ready().await.unwrap().call(ex).await.unwrap();
553 let elapsed = start.elapsed();
554
555 assert!(elapsed >= Duration::from_millis(180));
556 assert!(elapsed < Duration::from_millis(350));
557 }
558
559 #[tokio::test]
560 async fn recipient_list_collect_all_strategy() {
561 let resolver = Arc::new(|uri: &str| {
562 if uri == "mock:a" {
563 Some(BoxProcessor::from_fn(|mut ex| {
564 ex.input.body = Body::Text("a".to_string());
565 Box::pin(async move { Ok(ex) })
566 }))
567 } else if uri == "mock:b" {
568 Some(BoxProcessor::from_fn(|mut ex| {
569 ex.input.body = Body::Text("b".to_string());
570 Box::pin(async move { Ok(ex) })
571 }))
572 } else if uri == "mock:c" {
573 Some(BoxProcessor::from_fn(|mut ex| {
574 ex.input.body = Body::Text("c".to_string());
575 Box::pin(async move { Ok(ex) })
576 }))
577 } else {
578 None
579 }
580 });
581
582 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
583 |_ex: &Exchange| "mock:a,mock:b,mock:c".to_string(),
584 )))
585 .strategy(MulticastStrategy::CollectAll);
586
587 let mut svc = RecipientListService::new(config, resolver).unwrap();
588 let ex = Exchange::new(Message::new("seed"));
589 let result = svc.ready().await.unwrap().call(ex).await.unwrap();
590
591 assert_eq!(
592 result.input.body,
593 Body::from(Value::Array(vec![
594 Value::String("a".to_string()),
595 Value::String("b".to_string()),
596 Value::String("c".to_string()),
597 ]))
598 );
599 }
600
601 #[tokio::test]
602 async fn recipient_list_original_strategy() {
603 let resolver = Arc::new(|uri: &str| {
604 if uri.starts_with("mock:") {
605 let label = uri.to_string();
606 Some(BoxProcessor::from_fn(move |mut ex| {
607 let label = label.clone();
608 ex.input.body = Body::Text(format!("mutated-{label}"));
609 Box::pin(async move { Ok(ex) })
610 }))
611 } else {
612 None
613 }
614 });
615
616 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
617 |_ex: &Exchange| "mock:a,mock:b,mock:c".to_string(),
618 )))
619 .strategy(MulticastStrategy::Original);
620
621 let mut svc = RecipientListService::new(config, resolver).unwrap();
622 let ex = Exchange::new(Message::new("original"));
623 let result = svc.ready().await.unwrap().call(ex).await.unwrap();
624
625 assert_eq!(result.input.body.as_text(), Some("original"));
626 }
627
628 #[tokio::test]
635 async fn test_huge_recipient_list_is_capped() {
636 let call_count = Arc::new(AtomicUsize::new(0));
637 let count_clone = call_count.clone();
638
639 let resolver = Arc::new(move |uri: &str| {
640 if uri.starts_with("mock:") {
641 let count = count_clone.clone();
642 Some(BoxProcessor::from_fn(move |ex| {
643 count.fetch_add(1, Ordering::SeqCst);
644 Box::pin(async move { Ok(ex) })
645 }))
646 } else {
647 None
648 }
649 });
650
651 let mut many = String::with_capacity(8 * 1_000_000);
653 for i in 0..1_000_000 {
654 if i > 0 {
655 many.push(',');
656 }
657 many.push_str(&format!("mock:k{i}"));
658 }
659
660 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
666 |ex: &Exchange| {
667 ex.input
668 .header("CamelRecipients")
669 .and_then(|v| v.as_str().map(|s| s.to_string()))
670 .unwrap_or_default()
671 },
672 )))
673 .max_recipients(4);
674
675 let mut svc = RecipientListService::new(config, resolver).unwrap();
676 let mut ex = Exchange::new(Message::new("test"));
677 ex.input.set_header("CamelRecipients", Value::String(many));
678 let result = svc.ready().await.unwrap().call(ex).await;
679 assert!(result.is_ok(), "capped execution should still succeed");
680 assert_eq!(
681 call_count.load(Ordering::SeqCst),
682 4,
683 "must resolve at most max_recipients (4) endpoints"
684 );
685 }
686
687 #[tokio::test]
688 async fn recipient_list_last_wins_strategy() {
689 let payloads: Arc<HashMap<String, String>> = Arc::new(HashMap::from([
690 ("mock:a".to_string(), "first".to_string()),
691 ("mock:b".to_string(), "second".to_string()),
692 ("mock:c".to_string(), "third".to_string()),
693 ]));
694
695 let resolver = {
696 let payloads = payloads.clone();
697 Arc::new(move |uri: &str| {
698 if let Some(payload) = payloads.get(uri) {
699 let payload = payload.clone();
700 Some(BoxProcessor::from_fn(move |mut ex| {
701 let payload = payload.clone();
702 ex.input.body = Body::Text(payload);
703 Box::pin(async move { Ok(ex) })
704 }))
705 } else {
706 None
707 }
708 })
709 };
710
711 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
712 |_ex: &Exchange| "mock:a,mock:b,mock:c".to_string(),
713 )))
714 .strategy(MulticastStrategy::LastWins);
715
716 let mut svc = RecipientListService::new(config, resolver).unwrap();
717 let ex = Exchange::new(Message::new("seed"));
718 let result = svc.ready().await.unwrap().call(ex).await.unwrap();
719
720 assert_eq!(result.input.body.as_text(), Some("third"));
721 }
722
723 fn err_resolver(uri_to_err: Vec<(&'static str, CamelError)>) -> camel_api::EndpointResolver {
726 Arc::new(move |uri: &str| {
727 for (pattern, err) in &uri_to_err {
728 if uri == *pattern {
729 let err = err.clone();
730 return Some(BoxProcessor::from_fn(move |_ex| {
731 let err = err.clone();
732 Box::pin(async move { Err(err) })
733 }));
734 }
735 }
736 None
737 })
738 }
739
740 #[tokio::test]
741 async fn recipient_list_sequential_all_failed_returns_err() {
742 let resolver = err_resolver(vec![(
745 "mock:a",
746 CamelError::Config(String::from("seq-all-failed")),
747 )]);
748 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
749 |_ex: &Exchange| "mock:a".to_string(),
750 )))
751 .strategy(MulticastStrategy::LastWins);
752
753 let mut svc = RecipientListService::new(config, resolver).unwrap();
754 let mut ex = Exchange::new(Message::new("timer:t tick #1"));
755 ex.input.body = Body::Text(String::from("timer:t tick #1"));
756 let result = svc.ready().await.unwrap().call(ex).await;
757
758 assert!(
759 result.is_err(),
760 "zero-success recipient_list must return Err, not Ok(original)"
761 );
762 assert!(
763 matches!(result, Err(CamelError::Config(m)) if m == "seq-all-failed"),
764 "returned error must carry the iteration-last error"
765 );
766 }
767
768 #[tokio::test]
769 async fn recipient_list_parallel_all_failed_returns_err() {
770 let resolver = err_resolver(vec![
773 ("mock:a", CamelError::Config(String::from("par-err-a"))),
774 ("mock:b", CamelError::Config(String::from("par-err-b"))),
775 ]);
776 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
777 |_ex: &Exchange| "mock:a,mock:b".to_string(),
778 )))
779 .strategy(MulticastStrategy::LastWins)
780 .parallel(true);
781
782 let mut svc = RecipientListService::new(config, resolver).unwrap();
783 let ex = Exchange::new(Message::new("inbound"));
784 let result = svc.ready().await.unwrap().call(ex).await;
785
786 assert!(
787 result.is_err(),
788 "zero-success parallel recipient_list must return Err, not Ok(original)"
789 );
790 }
791
792 #[tokio::test]
793 async fn recipient_list_parallel_last_error_is_join_next_order() {
794 let (tx, rx) = tokio::sync::oneshot::channel::<()>();
800 let rx = Arc::new(tokio::sync::Mutex::new(Some(rx)));
801 let resolver: camel_api::EndpointResolver = Arc::new(move |uri: &str| {
802 if uri == "mock:a" {
803 Some(BoxProcessor::from_fn(|_ex| {
804 Box::pin(async move { Err(CamelError::Config(String::from("par-err-a"))) })
805 }))
806 } else if uri == "mock:b" {
807 let rx = rx.clone();
808 Some(BoxProcessor::from_fn(move |_ex| {
809 let rx = rx.clone();
810 Box::pin(async move {
811 let mut lock = rx.lock().await;
813 if let Some(rx) = lock.take() {
814 let _ = rx.await;
815 }
816 Err(CamelError::Config(String::from("par-err-b")))
817 })
818 }))
819 } else {
820 None
821 }
822 });
823 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
824 |_ex: &Exchange| "mock:a,mock:b".to_string(),
825 )))
826 .strategy(MulticastStrategy::LastWins)
827 .parallel(true);
828
829 let mut svc = RecipientListService::new(config, resolver).unwrap();
830 let ex = Exchange::new(Message::new("inbound"));
831
832 let join = tokio::spawn(async move { svc.ready().await.unwrap().call(ex).await });
835 for _ in 0..10 {
837 tokio::task::yield_now().await;
838 }
839 let _ = tx.send(());
840 let result = tokio::time::timeout(Duration::from_secs(5), join)
841 .await
842 .expect("join timeout")
843 .unwrap();
844
845 assert!(
846 matches!(result, Err(CamelError::Config(ref m)) if m == "par-err-b"),
847 "representative error must be the last failing task to complete (mock:b), got: {result:?}"
848 );
849 }
850
851 #[tokio::test]
852 async fn recipient_list_partial_success_aggregates_and_returns_ok() {
853 let call_count = Arc::new(AtomicUsize::new(0));
856 let ok_count = call_count.clone();
857 let resolver: camel_api::EndpointResolver = Arc::new(move |uri: &str| {
858 if uri == "mock:ok" {
859 let c = ok_count.clone();
860 Some(BoxProcessor::from_fn(move |mut ex| {
861 c.fetch_add(1, Ordering::SeqCst);
862 ex.input.body = Body::Text(String::from("ok-body"));
863 Box::pin(async move { Ok(ex) })
864 }))
865 } else if uri == "mock:fail" {
866 Some(BoxProcessor::from_fn(|_ex| {
867 Box::pin(async move { Err(CamelError::Config(String::from("partial-fail"))) })
868 }))
869 } else {
870 None
871 }
872 });
873 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
874 |_ex: &Exchange| "mock:fail,mock:ok".to_string(),
875 )))
876 .strategy(MulticastStrategy::LastWins);
877
878 let mut svc = RecipientListService::new(config, resolver).unwrap();
879 let ex = Exchange::new(Message::new("inbound"));
880 let result = svc.ready().await.unwrap().call(ex).await;
881
882 assert!(
883 result.is_ok(),
884 "partial success must return Ok, got: {result:?}"
885 );
886 assert_eq!(call_count.load(Ordering::SeqCst), 1);
887 assert_eq!(result.unwrap().input.body.as_text(), Some("ok-body"));
888 }
889
890 #[tokio::test]
891 async fn recipient_list_parallel_all_panic_returns_err() {
892 let resolver: camel_api::EndpointResolver = Arc::new(|uri: &str| {
898 if uri.starts_with("mock:panic") {
899 Some(BoxProcessor::from_fn(|_ex| {
900 Box::pin(async move {
901 panic!("recipient panicked");
902 })
903 }))
904 } else {
905 None
906 }
907 });
908 let config = RecipientListConfig::new(camel_api::RecipientSource::Sync(Arc::new(
909 |_ex: &Exchange| "mock:panic1,mock:panic2".to_string(),
910 )))
911 .strategy(MulticastStrategy::LastWins)
912 .parallel(true);
913
914 let mut svc = RecipientListService::new(config, resolver).unwrap();
915 let ex = Exchange::new(Message::new("inbound"));
916 let result = svc.ready().await.unwrap().call(ex).await;
917
918 assert!(
919 result.is_err(),
920 "all-panic parallel recipient_list must return Err, not Ok(original); got: {result:?}"
921 );
922 assert!(
923 matches!(result, Err(CamelError::ProcessorError(_))),
924 "panic must surface as a ProcessorError representative"
925 );
926 }
927
928 #[tokio::test]
929 async fn recipient_list_error_fails_step() {
930 use camel_api::{BoxValueFuture, RecipientSource};
931
932 let resolver_calls = Arc::new(AtomicUsize::new(0));
933 let calls_clone = resolver_calls.clone();
934 let resolver: camel_api::EndpointResolver = Arc::new(move |uri: &str| {
935 calls_clone.fetch_add(1, Ordering::SeqCst);
936 if uri.starts_with("mock:") {
937 Some(BoxProcessor::from_fn(|ex| Box::pin(async move { Ok(ex) })))
938 } else {
939 None
940 }
941 });
942
943 let config = RecipientListConfig::new(RecipientSource::Async(Arc::new(|_: &Exchange| {
944 Box::pin(async { Err(CamelError::ProcessorError("recipient boom".into())) })
945 as BoxValueFuture
946 })));
947
948 let mut svc = RecipientListService::new(config, resolver).unwrap();
949 let result = svc
950 .ready()
951 .await
952 .unwrap()
953 .call(Exchange::new(Message::new("test")))
954 .await;
955
956 assert!(
957 result.is_err(),
958 "a failed recipient expression must fail the step"
959 );
960 assert_eq!(
961 resolver_calls.load(Ordering::SeqCst),
962 0,
963 "no endpoint may be resolved when the recipient expression fails"
964 );
965 }
966}