Skip to main content

camel_processor/
dynamic_router.rs

1use std::future::Future;
2use std::pin::Pin;
3use std::task::{Context, Poll};
4use std::time::Instant;
5
6use tower::Service;
7use tower::ServiceExt;
8
9use camel_api::endpoint_pipeline::{CAMEL_SLIP_ENDPOINT, EndpointPipelineConfig, EndpointResolver};
10use camel_api::{CamelError, DynamicRouterConfig, Exchange, Value};
11
12use crate::endpoint_pipeline::EndpointPipelineService;
13
14#[derive(Clone)]
15pub struct DynamicRouterService {
16    config: DynamicRouterConfig,
17    pipeline: EndpointPipelineService,
18}
19
20impl DynamicRouterService {
21    pub fn new(config: DynamicRouterConfig, endpoint_resolver: EndpointResolver) -> Self {
22        let pipeline_config = EndpointPipelineConfig {
23            cache_size: EndpointPipelineConfig::from_signed(config.cache_size),
24            ignore_invalid_endpoints: config.ignore_invalid_endpoints,
25        };
26
27        Self {
28            config,
29            pipeline: EndpointPipelineService::new(endpoint_resolver, pipeline_config),
30        }
31    }
32}
33
34impl Service<Exchange> for DynamicRouterService {
35    type Response = Exchange;
36    type Error = CamelError;
37    type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
38
39    fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
40        Poll::Ready(Ok(()))
41    }
42
43    fn call(&mut self, mut exchange: Exchange) -> Self::Future {
44        let config = self.config.clone();
45        let pipeline = self.pipeline.clone();
46
47        Box::pin(async move {
48            let start = Instant::now();
49            let mut iterations = 0;
50            let mut last_destinations: Option<String> = None;
51
52            loop {
53                iterations += 1;
54
55                if iterations > config.max_iterations {
56                    return Err(CamelError::ProcessorError(format!(
57                        "Dynamic router exceeded max iterations ({})",
58                        config.max_iterations
59                    )));
60                }
61
62                if let Some(timeout) = config.timeout
63                    && start.elapsed() > timeout
64                {
65                    return Err(CamelError::ProcessorError(format!(
66                        "Dynamic router timed out after {:?}",
67                        timeout
68                    )));
69                }
70
71                let destinations = match config.expression.resolve(&exchange).await {
72                    Ok(None) => break,
73                    Ok(Some(uris)) => uris,
74                    Err(e) => return Err(e),
75                };
76
77                if last_destinations.as_deref() == Some(destinations.as_str()) {
78                    return Err(CamelError::ProcessorError(format!(
79                        "Dynamic router detected infinite loop: expression returned the same destination '{}' on consecutive iterations. The destination endpoint must clear or update the routing header to signal completion.",
80                        destinations
81                    )));
82                }
83
84                last_destinations = Some(destinations.clone());
85
86                for uri in destinations.split(&config.uri_delimiter) {
87                    let uri = uri.trim();
88                    if uri.is_empty() {
89                        continue;
90                    }
91
92                    let endpoint = match pipeline.resolve(uri)? {
93                        Some(e) => e,
94                        None => {
95                            continue;
96                        }
97                    };
98
99                    exchange.set_property(CAMEL_SLIP_ENDPOINT, Value::String(uri.to_string()));
100
101                    let mut endpoint = endpoint;
102                    exchange = endpoint.ready().await?.call(exchange).await?;
103                }
104            }
105
106            Ok(exchange)
107        })
108    }
109}
110
111#[cfg(test)]
112mod tests {
113    use super::*;
114    use camel_api::{BoxProcessor, BoxProcessorExt, Message};
115    use std::sync::Arc;
116    use std::sync::atomic::{AtomicUsize, Ordering};
117    use tower::ServiceExt;
118
119    fn make_config<F>(f: F) -> DynamicRouterConfig
120    where
121        F: Fn(&Exchange) -> Option<String> + Send + Sync + 'static,
122    {
123        DynamicRouterConfig::new(camel_api::TargetSource::Sync(Arc::new(f)))
124    }
125
126    fn mock_resolver() -> EndpointResolver {
127        Arc::new(|uri: &str| {
128            if uri.starts_with("mock:") {
129                Some(BoxProcessor::from_fn(|ex| Box::pin(async move { Ok(ex) })))
130            } else {
131                None
132            }
133        })
134    }
135
136    #[tokio::test]
137    async fn test_dynamic_router_single_destination() {
138        let call_count = Arc::new(AtomicUsize::new(0));
139        let count_clone = call_count.clone();
140        let expr_count = Arc::new(AtomicUsize::new(0));
141        let expr_count_clone = expr_count.clone();
142
143        let resolver = Arc::new(move |uri: &str| {
144            if uri == "mock:a" {
145                let count = count_clone.clone();
146                Some(BoxProcessor::from_fn(move |ex| {
147                    count.fetch_add(1, Ordering::SeqCst);
148                    Box::pin(async move { Ok(ex) })
149                }))
150            } else {
151                None
152            }
153        });
154
155        let config = DynamicRouterConfig::new(camel_api::TargetSource::Sync(Arc::new(
156            move |ex: &Exchange| {
157                let count = expr_count_clone.fetch_add(1, Ordering::SeqCst);
158                if count == 0 {
159                    ex.input
160                        .header("dest")
161                        .and_then(|v| v.as_str().map(|s| s.to_string()))
162                } else {
163                    None
164                }
165            },
166        )));
167
168        let mut svc = DynamicRouterService::new(config, resolver);
169
170        let mut ex = Exchange::new(Message::new("test"));
171        ex.input.set_header("dest", Value::String("mock:a".into()));
172
173        let _result = svc.ready().await.unwrap().call(ex).await.unwrap();
174        assert_eq!(call_count.load(Ordering::SeqCst), 1);
175    }
176
177    #[tokio::test]
178    async fn test_dynamic_router_loop_terminates_on_none() {
179        let iterations = Arc::new(AtomicUsize::new(0));
180        let iterations_clone = iterations.clone();
181
182        let config = DynamicRouterConfig::new(camel_api::TargetSource::Sync(Arc::new(
183            move |_ex: &Exchange| {
184                let count = iterations_clone.fetch_add(1, Ordering::SeqCst);
185                match count {
186                    0 => Some("mock:a".to_string()),
187                    1 => Some("mock:b".to_string()),
188                    _ => None,
189                }
190            },
191        )));
192
193        let mut svc = DynamicRouterService::new(config, mock_resolver());
194
195        let ex = Exchange::new(Message::new("test"));
196        let result = svc.ready().await.unwrap().call(ex).await;
197
198        assert!(result.is_ok());
199        assert_eq!(iterations.load(Ordering::SeqCst), 3);
200    }
201
202    #[tokio::test]
203    async fn test_dynamic_router_max_iterations() {
204        let config = make_config(|_| Some("mock:a".to_string())).max_iterations(5);
205
206        let mut svc = DynamicRouterService::new(config, mock_resolver());
207
208        let ex = Exchange::new(Message::new("test"));
209        let result = svc.ready().await.unwrap().call(ex).await;
210
211        assert!(result.is_err());
212        let err = result.unwrap_err().to_string();
213        assert!(err.contains("infinite loop"));
214        assert!(err.contains("same destination"));
215    }
216
217    #[tokio::test]
218    async fn test_dynamic_router_detects_same_destination_loop() {
219        let config = make_config(|_| Some("mock:loop".to_string())).max_iterations(100);
220
221        let mut svc = DynamicRouterService::new(config, mock_resolver());
222
223        let ex = Exchange::new(Message::new("test"));
224        let result = svc.ready().await.unwrap().call(ex).await;
225
226        assert!(result.is_err());
227        let err = result.unwrap_err().to_string();
228        assert!(err.contains("infinite loop"));
229        assert!(err.contains("mock:loop"));
230    }
231
232    #[tokio::test]
233    async fn test_dynamic_router_invalid_endpoint_error() {
234        let config =
235            make_config(|_| Some("invalid:endpoint".to_string())).ignore_invalid_endpoints(false);
236
237        let mut svc = DynamicRouterService::new(config, mock_resolver());
238
239        let ex = Exchange::new(Message::new("test"));
240        let result = svc.ready().await.unwrap().call(ex).await;
241
242        assert!(result.is_err());
243        let err = result.unwrap_err().to_string();
244        assert!(err.contains("Invalid endpoint"));
245    }
246
247    #[tokio::test]
248    async fn test_dynamic_router_ignore_invalid_endpoint() {
249        let call_count = Arc::new(AtomicUsize::new(0));
250        let count_clone = call_count.clone();
251        let expr_count = Arc::new(AtomicUsize::new(0));
252        let expr_count_clone = expr_count.clone();
253
254        let resolver = Arc::new(move |uri: &str| {
255            if uri == "mock:valid" {
256                let count = count_clone.clone();
257                Some(BoxProcessor::from_fn(move |ex| {
258                    count.fetch_add(1, Ordering::SeqCst);
259                    Box::pin(async move { Ok(ex) })
260                }))
261            } else {
262                None
263            }
264        });
265
266        let config = DynamicRouterConfig::new(camel_api::TargetSource::Sync(Arc::new(
267            move |_ex: &Exchange| {
268                let count = expr_count_clone.fetch_add(1, Ordering::SeqCst);
269                if count == 0 {
270                    Some("invalid:endpoint,mock:valid".to_string())
271                } else {
272                    None
273                }
274            },
275        )))
276        .ignore_invalid_endpoints(true);
277
278        let mut svc = DynamicRouterService::new(config, resolver);
279
280        let ex = Exchange::new(Message::new("test"));
281        let result = svc.ready().await.unwrap().call(ex).await;
282
283        assert!(result.is_ok());
284        assert_eq!(call_count.load(Ordering::SeqCst), 1);
285    }
286
287    #[tokio::test]
288    async fn test_dynamic_router_cache_size_enforced() {
289        // Verify that the cache never exceeds the configured capacity.
290        let resolver_call_count = Arc::new(AtomicUsize::new(0));
291        let count_clone = resolver_call_count.clone();
292
293        let resolver: EndpointResolver = Arc::new(move |uri: &str| {
294            if uri.starts_with("mock:") {
295                count_clone.fetch_add(1, Ordering::SeqCst);
296                Some(BoxProcessor::from_fn(|ex| Box::pin(async move { Ok(ex) })))
297            } else {
298                None
299            }
300        });
301
302        // Cache capacity of 2; we will route through 3 distinct URIs.
303        let expr_count = Arc::new(AtomicUsize::new(0));
304        let expr_count_clone = expr_count.clone();
305        let config = DynamicRouterConfig::new(camel_api::TargetSource::Sync(Arc::new(
306            move |_ex: &Exchange| {
307                let n = expr_count_clone.fetch_add(1, Ordering::SeqCst);
308                match n {
309                    0 => Some("mock:a".to_string()),
310                    1 => Some("mock:b".to_string()),
311                    2 => Some("mock:c".to_string()),
312                    _ => None,
313                }
314            },
315        )))
316        .cache_size(2);
317
318        let mut svc = DynamicRouterService::new(config, resolver);
319
320        let ex = Exchange::new(Message::new("test"));
321        svc.ready().await.unwrap().call(ex).await.unwrap();
322
323        // The cache capacity is 2 so mock:a, mock:b, mock:c each required a resolver call
324        // (mock:a is evicted before mock:c is inserted). All three resolver calls must have happened.
325        assert_eq!(resolver_call_count.load(Ordering::SeqCst), 3);
326    }
327
328    #[tokio::test]
329    async fn dynamic_router_target_error_fails_step() {
330        use camel_api::{BoxValueFuture, TargetSource};
331
332        let resolver_calls = Arc::new(AtomicUsize::new(0));
333        let calls_clone = resolver_calls.clone();
334        let resolver: EndpointResolver = Arc::new(move |uri: &str| {
335            calls_clone.fetch_add(1, Ordering::SeqCst);
336            if uri.starts_with("mock:") {
337                Some(BoxProcessor::from_fn(|ex| Box::pin(async move { Ok(ex) })))
338            } else {
339                None
340            }
341        });
342
343        let config = DynamicRouterConfig::new(TargetSource::Async(Arc::new(|_: &Exchange| {
344            Box::pin(async { Err(CamelError::ProcessorError("router boom".into())) })
345                as BoxValueFuture
346        })));
347
348        let mut svc = DynamicRouterService::new(config, resolver);
349        let result = svc
350            .ready()
351            .await
352            .unwrap()
353            .call(Exchange::new(Message::new("test")))
354            .await;
355
356        assert!(
357            result.is_err(),
358            "a failed target expression must fail the step"
359        );
360        assert_eq!(
361            resolver_calls.load(Ordering::SeqCst),
362            0,
363            "no endpoint may be resolved when the target expression fails"
364        );
365    }
366}