camel_processor/
routing_slip.rs1use 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#[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}