camel_processor/
dynamic_router.rs1use 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 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 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 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}