Skip to main content

camel_processor/
recipient_list.rs

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            // H13 Batch 1: cap the resolved-URI list BEFORE any endpoint
58            // resolution. A malicious expression yielding millions of URIs
59            // would otherwise allocate a Vec of millions of &str references
60            // and resolve each one (multicast) or call each one (sequential).
61            // The default cap is 1_000 (camel-api::recipient_list).
62            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                            // stop_on_exception=false: track the representative
106                            // error — the last failing task to complete via
107                            // join_next order (ADR-0058). Pending tasks continue.
108                            last_parallel_error = Some(e);
109                        }
110                        Err(join_err) if join_err.is_panic() => {
111                            // A recipient task panicked. ADR-0058: a panic is
112                            // zero-success attempted work and MUST NOT launder to
113                            // Ok(original); convert to a representative error so
114                            // the zero-success guard fires. (Cancellation is
115                            // handled separately below — it is often self-induced
116                            // by stop_on_exception's abort_all.)
117                            last_parallel_error = Some(CamelError::ProcessorError(format!(
118                                "recipient task panicked: {join_err}"
119                            )));
120                        }
121                        Err(_) => {} // Cancellation (JoinSet abort); ignore.
122                    }
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                // ADR-0058: zero-success operational failure. At least one
132                // recipient was called and zero returned Ok — report the
133                // representative error instead of laundering to Ok(original).
134                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                            // stop_on_exception=false: track the iteration-last
164                            // error (ADR-0058) and continue to remaining recipients.
165                            last_error = Some(e);
166                            continue;
167                        }
168                    }
169                }
170                // ADR-0058: zero-success operational failure. At least one
171                // recipient was called and zero returned Ok — report the
172                // iteration-last error instead of laundering to Ok(original),
173                // which would poison an outer cache write-back with the inbound body.
174                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                    // Empty and future variants contribute no extractable value.
209                    _ => 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 and any future variant return the original exchange.
220        _ => 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    // ── H13 Batch 1: cap resolved-URI count ──────────────────────────
629
630    /// H13: an expression yielding millions of URIs is truncated to
631    /// `max_recipients` before endpoint resolution. The test uses a
632    /// cap of 4 to keep the test fast; the principle (cap the list) is
633    /// what Batch 1 enforces. The default cap is 1_000 in camel-api.
634    #[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        // Build the untrusted payload: 1_000_000 URIs as one string.
652        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        // Disposition-5 pattern: the untrusted data flows FROM the exchange
661        // (a header on the inbound message), NOT from a captured variable.
662        // The expression reads it off the passed `&Exchange` — this is what
663        // makes the cap an untrusted-data-validation control, not a local
664        // limit.
665        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    // ── ADR-0058: zero-success operational failure must not launder to Ok(original) ─
724
725    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        // ADR-0058: zero-success sequential. One recipient errors; zero Ok.
743        // MUST return Err, not Ok(original) (which would poison an outer cache).
744        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        // ADR-0058: zero-success parallel. Two recipients error; zero Ok.
771        // MUST return a representative Err, not Ok(original).
772        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        // ADR-0058 last-error determinism: the representative error is the one
795        // from the task returned by the last `JoinSet::join_next` that completed
796        // with an error. mock:a errors immediately; mock:b awaits a oneshot
797        // signal then errors. The test sends the signal after a brief yield so
798        // mock:a completes first → join_next order yields mock:b's error last.
799        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                        // Wait for the test's signal before completing.
812                        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        // Drive the call concurrently; release mock:b after mock:a has had a
833        // chance to error first.
834        let join = tokio::spawn(async move { svc.ready().await.unwrap().call(ex).await });
835        // Yield the runtime so mock:a (synchronous Err) completes before mock:b.
836        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        // ADR-0058: partial success (>=1 Ok) MUST aggregate over successes and
854        // return Ok. The invariant fires only on ZERO successes.
855        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        // ADR-0058 (e_gpt review gap): a parallel recipient_list where every
893        // spawned task PANICS produces only JoinError(panic) results. These
894        // MUST NOT launder to Ok(original); convert to a representative error
895        // so the zero-success guard fires. Cancels (self-induced abort) stay
896        // ignored.
897        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}