Skip to main content

camel_processor/
routing_slip.rs

1use std::future::Future;
2use std::pin::Pin;
3use std::task::{Context, Poll};
4
5use tower::Service;
6use tower::ServiceExt;
7
8use camel_api::endpoint_pipeline::CAMEL_SLIP_ENDPOINT;
9use camel_api::{CamelError, EndpointPipelineConfig, Exchange, RoutingSlipConfig, Value};
10
11use crate::endpoint_pipeline::EndpointPipelineService;
12
13/// Routing Slip EIP implementation.
14///
15/// Evaluates an expression once to get a list of endpoint URIs, then routes
16/// the message through them sequentially in a pipeline fashion.
17#[derive(Clone)]
18pub struct RoutingSlipService {
19    config: RoutingSlipConfig,
20    pipeline: EndpointPipelineService,
21}
22
23impl RoutingSlipService {
24    pub fn new(config: RoutingSlipConfig, endpoint_resolver: camel_api::EndpointResolver) -> Self {
25        let pipeline_config = EndpointPipelineConfig {
26            cache_size: EndpointPipelineConfig::from_signed(config.cache_size),
27            ignore_invalid_endpoints: config.ignore_invalid_endpoints,
28        };
29        Self {
30            config,
31            pipeline: EndpointPipelineService::new(endpoint_resolver, pipeline_config),
32        }
33    }
34}
35
36impl Service<Exchange> for RoutingSlipService {
37    type Response = Exchange;
38    type Error = CamelError;
39    type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
40
41    fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
42        Poll::Ready(Ok(()))
43    }
44
45    fn call(&mut self, mut exchange: Exchange) -> Self::Future {
46        let config = self.config.clone();
47        let pipeline = self.pipeline.clone();
48
49        Box::pin(async move {
50            let slip = match config.expression.resolve(&exchange).await {
51                Ok(None) => return Ok(exchange),
52                Ok(Some(s)) => s,
53                Err(e) => return Err(e),
54            };
55
56            for uri in slip.split(&config.uri_delimiter) {
57                let uri = uri.trim();
58                if uri.is_empty() {
59                    continue;
60                }
61
62                let endpoint = match pipeline.resolve(uri)? {
63                    Some(e) => e,
64                    None => continue,
65                };
66
67                exchange.set_property(CAMEL_SLIP_ENDPOINT, Value::String(uri.to_string()));
68
69                let mut endpoint = endpoint;
70                exchange = endpoint.ready().await?.call(exchange).await?;
71            }
72
73            Ok(exchange)
74        })
75    }
76}
77
78#[cfg(test)]
79mod tests {
80    use super::*;
81    use camel_api::{BoxProcessor, BoxProcessorExt, Message};
82    use std::sync::Arc;
83    use std::sync::atomic::{AtomicUsize, Ordering};
84
85    fn mock_resolver() -> camel_api::EndpointResolver {
86        Arc::new(|uri: &str| {
87            if uri.starts_with("mock:") {
88                Some(BoxProcessor::from_fn(|ex| Box::pin(async move { Ok(ex) })))
89            } else {
90                None
91            }
92        })
93    }
94
95    #[tokio::test]
96    async fn routing_slip_single_destination() {
97        let call_count = Arc::new(AtomicUsize::new(0));
98        let count_clone = call_count.clone();
99
100        let resolver = Arc::new(move |uri: &str| {
101            if uri == "mock:a" {
102                let count = count_clone.clone();
103                Some(BoxProcessor::from_fn(move |ex| {
104                    count.fetch_add(1, Ordering::SeqCst);
105                    Box::pin(async move { Ok(ex) })
106                }))
107            } else {
108                None
109            }
110        });
111
112        let config =
113            RoutingSlipConfig::new(camel_api::TargetSource::Sync(Arc::new(|_ex: &Exchange| {
114                Some("mock:a".to_string())
115            })));
116
117        let mut svc = RoutingSlipService::new(config, resolver);
118        let ex = Exchange::new(Message::new("test"));
119        let result = svc.ready().await.unwrap().call(ex).await;
120
121        assert!(result.is_ok());
122        assert_eq!(call_count.load(Ordering::SeqCst), 1);
123    }
124
125    #[tokio::test]
126    async fn routing_slip_multiple_destinations() {
127        let call_count = Arc::new(AtomicUsize::new(0));
128        let count_clone = call_count.clone();
129
130        let resolver = Arc::new(move |uri: &str| {
131            if uri.starts_with("mock:") {
132                let count = count_clone.clone();
133                Some(BoxProcessor::from_fn(move |ex| {
134                    count.fetch_add(1, Ordering::SeqCst);
135                    Box::pin(async move { Ok(ex) })
136                }))
137            } else {
138                None
139            }
140        });
141
142        let config =
143            RoutingSlipConfig::new(camel_api::TargetSource::Sync(Arc::new(|_ex: &Exchange| {
144                Some("mock:a,mock:b,mock:c".to_string())
145            })));
146
147        let mut svc = RoutingSlipService::new(config, resolver);
148        let ex = Exchange::new(Message::new("test"));
149        let result = svc.ready().await.unwrap().call(ex).await;
150
151        assert!(result.is_ok());
152        assert_eq!(call_count.load(Ordering::SeqCst), 3);
153    }
154
155    #[tokio::test]
156    async fn routing_slip_empty_expression() {
157        let config =
158            RoutingSlipConfig::new(camel_api::TargetSource::Sync(Arc::new(|_ex: &Exchange| {
159                None
160            })));
161
162        let mut svc = RoutingSlipService::new(config, mock_resolver());
163        let ex = Exchange::new(Message::new("test"));
164        let result = svc.ready().await.unwrap().call(ex).await;
165
166        assert!(result.is_ok());
167    }
168
169    #[tokio::test]
170    async fn routing_slip_empty_string() {
171        let config =
172            RoutingSlipConfig::new(camel_api::TargetSource::Sync(Arc::new(|_ex: &Exchange| {
173                Some(String::new())
174            })));
175
176        let mut svc = RoutingSlipService::new(config, mock_resolver());
177        let ex = Exchange::new(Message::new("test"));
178        let result = svc.ready().await.unwrap().call(ex).await;
179
180        assert!(result.is_ok());
181    }
182
183    #[tokio::test]
184    async fn routing_slip_invalid_endpoint_error() {
185        let config =
186            RoutingSlipConfig::new(camel_api::TargetSource::Sync(Arc::new(|_ex: &Exchange| {
187                Some("invalid:endpoint".to_string())
188            })))
189            .ignore_invalid_endpoints(false);
190
191        let mut svc = RoutingSlipService::new(config, mock_resolver());
192        let ex = Exchange::new(Message::new("test"));
193        let result = svc.ready().await.unwrap().call(ex).await;
194
195        assert!(result.is_err());
196        assert!(result.unwrap_err().to_string().contains("Invalid endpoint"));
197    }
198
199    #[tokio::test]
200    async fn routing_slip_ignore_invalid_endpoint() {
201        let call_count = Arc::new(AtomicUsize::new(0));
202        let count_clone = call_count.clone();
203
204        let resolver = Arc::new(move |uri: &str| {
205            if uri == "mock:valid" {
206                let count = count_clone.clone();
207                Some(BoxProcessor::from_fn(move |ex| {
208                    count.fetch_add(1, Ordering::SeqCst);
209                    Box::pin(async move { Ok(ex) })
210                }))
211            } else {
212                None
213            }
214        });
215
216        let config =
217            RoutingSlipConfig::new(camel_api::TargetSource::Sync(Arc::new(|_ex: &Exchange| {
218                Some("invalid:endpoint,mock:valid".to_string())
219            })))
220            .ignore_invalid_endpoints(true);
221
222        let mut svc = RoutingSlipService::new(config, resolver);
223        let ex = Exchange::new(Message::new("test"));
224        let result = svc.ready().await.unwrap().call(ex).await;
225
226        assert!(result.is_ok());
227        assert_eq!(call_count.load(Ordering::SeqCst), 1);
228    }
229
230    #[tokio::test]
231    async fn routing_slip_order_preserved() {
232        use std::sync::Mutex;
233
234        let order: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(Vec::new()));
235        let order_clone = order.clone();
236
237        let resolver = Arc::new(move |uri: &str| {
238            let order = order_clone.clone();
239            let uri = uri.to_string();
240            Some(BoxProcessor::from_fn(move |ex| {
241                order.lock().unwrap().push(uri.clone());
242                Box::pin(async move { Ok(ex) })
243            }))
244        });
245
246        let config =
247            RoutingSlipConfig::new(camel_api::TargetSource::Sync(Arc::new(|_ex: &Exchange| {
248                Some("mock:first,mock:second,mock:third".to_string())
249            })));
250
251        let mut svc = RoutingSlipService::new(config, resolver);
252        let ex = Exchange::new(Message::new("test"));
253        svc.ready().await.unwrap().call(ex).await.unwrap();
254
255        let order = order.lock().unwrap();
256        assert_eq!(*order, vec!["mock:first", "mock:second", "mock:third"]);
257    }
258
259    #[tokio::test]
260    async fn routing_slip_endpoint_property_set() {
261        let last_uri: Arc<std::sync::Mutex<Option<String>>> = Arc::new(std::sync::Mutex::new(None));
262        let last_uri_clone = last_uri.clone();
263
264        let resolver = Arc::new(move |uri: &str| {
265            let last = last_uri_clone.clone();
266            let _uri = uri.to_string();
267            Some(BoxProcessor::from_fn(move |ex| {
268                let prop = ex.property(CAMEL_SLIP_ENDPOINT).cloned();
269                *last.lock().unwrap() = prop.and_then(|v| v.as_str().map(String::from));
270                Box::pin(async move { Ok(ex) })
271            }))
272        });
273
274        let config =
275            RoutingSlipConfig::new(camel_api::TargetSource::Sync(Arc::new(|_ex: &Exchange| {
276                Some("mock:a,mock:b".to_string())
277            })));
278
279        let mut svc = RoutingSlipService::new(config, resolver);
280        let ex = Exchange::new(Message::new("test"));
281        svc.ready().await.unwrap().call(ex).await.unwrap();
282
283        let last = last_uri.lock().unwrap();
284        assert_eq!(last.as_deref(), Some("mock:b"));
285    }
286
287    #[tokio::test]
288    async fn routing_slip_mutation_between_steps() {
289        let resolver = Arc::new(|uri: &str| {
290            if uri == "mock:mutate" {
291                Some(BoxProcessor::from_fn(|mut ex| {
292                    ex.input.body = camel_api::Body::Text("mutated".to_string());
293                    Box::pin(async move { Ok(ex) })
294                }))
295            } else if uri == "mock:verify" {
296                Some(BoxProcessor::from_fn(|ex| {
297                    let body = ex.input.body.as_text().unwrap_or("").to_string();
298                    assert_eq!(body, "mutated");
299                    Box::pin(async move { Ok(ex) })
300                }))
301            } else {
302                None
303            }
304        });
305
306        let config =
307            RoutingSlipConfig::new(camel_api::TargetSource::Sync(Arc::new(|_ex: &Exchange| {
308                Some("mock:mutate,mock:verify".to_string())
309            })));
310
311        let mut svc = RoutingSlipService::new(config, resolver);
312        let ex = Exchange::new(Message::new("original"));
313        let result = svc.ready().await.unwrap().call(ex).await;
314
315        assert!(result.is_ok());
316    }
317
318    #[tokio::test]
319    async fn routing_slip_cache_hit() {
320        let resolve_count = Arc::new(AtomicUsize::new(0));
321        let resolve_clone = resolve_count.clone();
322
323        let resolver = Arc::new(move |uri: &str| {
324            if uri.starts_with("mock:") {
325                resolve_clone.fetch_add(1, Ordering::SeqCst);
326                Some(BoxProcessor::from_fn(|ex| Box::pin(async move { Ok(ex) })))
327            } else {
328                None
329            }
330        });
331
332        let call_count = Arc::new(AtomicUsize::new(0));
333        let call_clone = call_count.clone();
334
335        let config = RoutingSlipConfig::new(camel_api::TargetSource::Sync(Arc::new(
336            move |_ex: &Exchange| {
337                let n = call_clone.fetch_add(1, Ordering::SeqCst);
338                if n < 2 {
339                    Some("mock:a,mock:b".to_string())
340                } else {
341                    None
342                }
343            },
344        )));
345
346        let mut svc = RoutingSlipService::new(config, resolver);
347        let ex1 = Exchange::new(Message::new("test1"));
348        svc.ready().await.unwrap().call(ex1).await.unwrap();
349        let ex2 = Exchange::new(Message::new("test2"));
350        svc.ready().await.unwrap().call(ex2).await.unwrap();
351
352        assert_eq!(resolve_count.load(Ordering::SeqCst), 2);
353    }
354
355    #[tokio::test]
356    async fn routing_slip_custom_delimiter() {
357        let order: Arc<std::sync::Mutex<Vec<String>>> = Arc::new(std::sync::Mutex::new(Vec::new()));
358        let order_clone = order.clone();
359
360        let resolver = Arc::new(move |uri: &str| {
361            let order = order_clone.clone();
362            let uri = uri.to_string();
363            Some(BoxProcessor::from_fn(move |ex| {
364                order.lock().unwrap().push(uri.clone());
365                Box::pin(async move { Ok(ex) })
366            }))
367        });
368
369        let config =
370            RoutingSlipConfig::new(camel_api::TargetSource::Sync(Arc::new(|_ex: &Exchange| {
371                Some("mock:x|mock:y|mock:z".to_string())
372            })))
373            .uri_delimiter("|");
374
375        let mut svc = RoutingSlipService::new(config, resolver);
376        let ex = Exchange::new(Message::new("test"));
377        svc.ready().await.unwrap().call(ex).await.unwrap();
378
379        let order = order.lock().unwrap();
380        assert_eq!(*order, vec!["mock:x", "mock:y", "mock:z"]);
381    }
382
383    #[tokio::test]
384    async fn routing_slip_expression_evaluated_once() {
385        let expr_count = Arc::new(AtomicUsize::new(0));
386        let expr_count_clone = expr_count.clone();
387
388        let resolver = Arc::new(|uri: &str| {
389            if uri.starts_with("mock:") {
390                Some(BoxProcessor::from_fn(|ex| Box::pin(async move { Ok(ex) })))
391            } else {
392                None
393            }
394        });
395
396        let config = RoutingSlipConfig::new(camel_api::TargetSource::Sync(Arc::new(
397            move |_ex: &Exchange| {
398                expr_count_clone.fetch_add(1, Ordering::SeqCst);
399                Some("mock:a,mock:b".to_string())
400            },
401        )));
402
403        let mut svc = RoutingSlipService::new(config, resolver);
404        let ex = Exchange::new(Message::new("test"));
405        svc.ready().await.unwrap().call(ex).await.unwrap();
406
407        assert_eq!(
408            expr_count.load(Ordering::SeqCst),
409            1,
410            "Expression must be evaluated exactly once"
411        );
412    }
413
414    #[tokio::test]
415    async fn routing_slip_target_error_fails_step() {
416        use camel_api::{BoxValueFuture, TargetSource};
417
418        let config = RoutingSlipConfig::new(TargetSource::Async(Arc::new(|_: &Exchange| {
419            Box::pin(async { Err(CamelError::ProcessorError("slip boom".into())) })
420                as BoxValueFuture
421        })));
422
423        let mut svc = RoutingSlipService::new(config, mock_resolver());
424        let result = svc
425            .ready()
426            .await
427            .unwrap()
428            .call(Exchange::new(Message::new("test")))
429            .await;
430
431        assert!(
432            result.is_err(),
433            "a failed slip expression must fail the step"
434        );
435    }
436}