1use crate::error::CamelError;
2use crate::exchange::Exchange;
3use crate::value::Value;
4use std::sync::Arc;
5
6pub type BoxValueFuture =
12 std::pin::Pin<Box<dyn std::future::Future<Output = Result<Value, CamelError>> + Send>>;
13
14pub type BoxBoolFuture =
18 std::pin::Pin<Box<dyn std::future::Future<Output = Result<bool, CamelError>> + Send>>;
19
20pub struct FilterPredicate(pub Arc<dyn Fn(&Exchange) -> bool + Send + Sync>);
31
32impl FilterPredicate {
33 pub fn new<F>(f: F) -> Self
35 where
36 F: Fn(&Exchange) -> bool + Send + Sync + 'static,
37 {
38 FilterPredicate(Arc::new(f))
39 }
40}
41
42impl Clone for FilterPredicate {
43 fn clone(&self) -> Self {
44 FilterPredicate(Arc::clone(&self.0))
45 }
46}
47
48impl std::fmt::Debug for FilterPredicate {
49 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
50 f.write_str("FilterPredicate(..)")
51 }
52}
53
54impl std::ops::Deref for FilterPredicate {
55 type Target = dyn Fn(&Exchange) -> bool + Send + Sync;
56
57 fn deref(&self) -> &Self::Target {
58 &*self.0
59 }
60}
61
62#[derive(Clone)]
69#[non_exhaustive]
70pub enum PredicateSource {
71 Sync(FilterPredicate),
73 Async(Arc<dyn Fn(&Exchange) -> BoxBoolFuture + Send + Sync>),
75}
76
77impl PredicateSource {
78 pub async fn matches(&self, exchange: &Exchange) -> Result<bool, CamelError> {
80 match self {
81 Self::Sync(pred) => Ok(pred(exchange)),
82 Self::Async(f) => f(exchange).await,
83 }
84 }
85}
86
87impl From<FilterPredicate> for PredicateSource {
88 fn from(pred: FilterPredicate) -> Self {
89 Self::Sync(pred)
90 }
91}
92
93#[derive(Clone)]
95#[non_exhaustive]
96pub enum ValueSource {
97 Sync(Arc<dyn Fn(&Exchange) -> Value + Send + Sync>),
99 Async(Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>),
101}
102
103impl ValueSource {
104 pub async fn evaluate(&self, exchange: &Exchange) -> Result<Value, CamelError> {
106 match self {
107 Self::Sync(f) => Ok(f(exchange)),
108 Self::Async(f) => f(exchange).await,
109 }
110 }
111}
112
113const NON_SCALAR_TARGET: &str =
115 "router target expression returned a non-scalar value (array/object); expected a string target";
116
117fn coerce_target_value(value: Value) -> Result<Option<String>, CamelError> {
122 match value {
123 Value::Null => Ok(None),
124 Value::String(s) => Ok(Some(s)),
125 Value::Array(_) | Value::Object(_) => {
126 Err(CamelError::ProcessorError(NON_SCALAR_TARGET.to_string()))
127 }
128 other => Ok(Some(other.to_string())),
129 }
130}
131
132#[allow(clippy::type_complexity)]
136#[derive(Clone)]
137#[non_exhaustive]
138pub enum TargetSource {
139 Sync(Arc<dyn Fn(&Exchange) -> Option<String> + Send + Sync>),
141 Async(Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>),
143}
144
145impl TargetSource {
146 pub async fn resolve(&self, exchange: &Exchange) -> Result<Option<String>, CamelError> {
148 match self {
149 Self::Sync(f) => Ok(f(exchange)),
150 Self::Async(f) => coerce_target_value(f(exchange).await?),
151 }
152 }
153}
154
155#[derive(Clone)]
157#[non_exhaustive]
158pub enum RecipientSource {
159 Sync(Arc<dyn Fn(&Exchange) -> String + Send + Sync>),
161 Async(Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>),
163}
164
165impl RecipientSource {
166 pub async fn resolve(&self, exchange: &Exchange) -> Result<String, CamelError> {
170 match self {
171 Self::Sync(f) => Ok(f(exchange)),
172 Self::Async(f) => Ok(coerce_target_value(f(exchange).await?)?.unwrap_or_default()),
173 }
174 }
175}
176
177impl From<Arc<dyn Fn(&Exchange) -> Value + Send + Sync>> for ValueSource {
178 fn from(f: Arc<dyn Fn(&Exchange) -> Value + Send + Sync>) -> Self {
179 Self::Sync(f)
180 }
181}
182
183impl From<Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>> for ValueSource {
184 fn from(f: Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>) -> Self {
185 Self::Async(f)
186 }
187}
188
189impl From<Arc<dyn Fn(&Exchange) -> Option<String> + Send + Sync>> for TargetSource {
190 fn from(f: Arc<dyn Fn(&Exchange) -> Option<String> + Send + Sync>) -> Self {
191 Self::Sync(f)
192 }
193}
194
195impl From<Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>> for TargetSource {
196 fn from(f: Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>) -> Self {
197 Self::Async(f)
198 }
199}
200
201impl From<Arc<dyn Fn(&Exchange) -> String + Send + Sync>> for RecipientSource {
202 fn from(f: Arc<dyn Fn(&Exchange) -> String + Send + Sync>) -> Self {
203 Self::Sync(f)
204 }
205}
206
207impl From<Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>> for RecipientSource {
208 fn from(f: Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>) -> Self {
209 Self::Async(f)
210 }
211}
212
213#[cfg(test)]
214mod tests {
215 use super::*;
216 use crate::{Exchange, ExpressionErrorClass, Message};
217
218 fn expression_failed() -> CamelError {
219 CamelError::ExpressionFailed {
220 language: "rhai".to_string(),
221 route_id: "r1".to_string(),
222 step_id: "step#0".to_string(),
223 verb: "filter".to_string(),
224 class: ExpressionErrorClass::Runtime,
225 position: None,
226 conversion: None,
227 cause: None,
228 }
229 }
230
231 #[test]
232 fn test_filter_predicate_is_callable() {
233 let pred = FilterPredicate::new(|ex: &Exchange| ex.input.body.as_text().is_some());
234 let ex = Exchange::new(Message::new("hello"));
235 assert!(pred(&ex));
236 }
237
238 #[test]
239 fn test_filter_predicate_debug_is_redacted() {
240 let pred = FilterPredicate::new(|_: &Exchange| true);
241 assert_eq!(format!("{pred:?}"), "FilterPredicate(..)");
242 }
243
244 #[test]
245 fn test_filter_predicate_clone_shares_arc() {
246 let pred = FilterPredicate::new(|_: &Exchange| true);
247 let cloned = pred.clone();
248 assert!(matches!(cloned, FilterPredicate(_)));
249 }
250
251 #[tokio::test]
252 async fn predicate_source_sync_returns_bool() {
253 let source = PredicateSource::Sync(FilterPredicate::new(|_: &Exchange| true));
254 let ex = Exchange::new(Message::new("hello"));
255 assert!(source.matches(&ex).await.unwrap());
256 }
257
258 #[tokio::test]
259 async fn predicate_source_async_propagates_error() {
260 let source = PredicateSource::Async(Arc::new(|_: &Exchange| {
261 Box::pin(async { Err(expression_failed()) }) as BoxBoolFuture
262 }));
263 let ex = Exchange::new(Message::new("hello"));
264 let err = source.matches(&ex).await.unwrap_err();
265 assert!(matches!(err, CamelError::ExpressionFailed { .. }));
266 }
267
268 #[tokio::test]
269 async fn value_source_async_propagates_error() {
270 let source = ValueSource::Async(Arc::new(|_: &Exchange| {
271 Box::pin(async { Err(expression_failed()) }) as BoxValueFuture
272 }));
273 let ex = Exchange::new(Message::new("hello"));
274 let err = source.evaluate(&ex).await.unwrap_err();
275 assert!(matches!(err, CamelError::ExpressionFailed { .. }));
276 }
277
278 #[tokio::test]
279 async fn target_source_null_maps_to_none_and_error_propagates() {
280 let ex = Exchange::new(Message::new("hello"));
281
282 let null_source = TargetSource::Async(Arc::new(|_: &Exchange| {
283 Box::pin(async { Ok(Value::Null) }) as BoxValueFuture
284 }));
285 assert_eq!(null_source.resolve(&ex).await.unwrap(), None);
286
287 let err_source = TargetSource::Async(Arc::new(|_: &Exchange| {
288 Box::pin(async { Err(expression_failed()) }) as BoxValueFuture
289 }));
290 let err = err_source.resolve(&ex).await.unwrap_err();
291 assert!(matches!(err, CamelError::ExpressionFailed { .. }));
292 }
293
294 #[tokio::test]
295 async fn target_source_array_is_error() {
296 let source = TargetSource::Async(Arc::new(|_: &Exchange| {
297 Box::pin(async { Ok(Value::Array(vec![Value::String("a".into())])) }) as BoxValueFuture
298 }));
299 let ex = Exchange::new(Message::new("hello"));
300 match source.resolve(&ex).await.unwrap_err() {
301 CamelError::ProcessorError(msg) => {
302 assert!(msg.contains("non-scalar"), "missing non-scalar: {msg}");
303 assert!(msg.contains("string target"), "missing target: {msg}");
304 }
305 other => panic!("expected ProcessorError, got {other:?}"),
306 }
307 }
308
309 #[tokio::test]
310 async fn recipient_source_null_maps_to_empty_and_non_scalar_is_error() {
311 let ex = Exchange::new(Message::new("hello"));
312
313 let null_source = RecipientSource::Async(Arc::new(|_: &Exchange| {
314 Box::pin(async { Ok(Value::Null) }) as BoxValueFuture
315 }));
316 assert_eq!(null_source.resolve(&ex).await.unwrap(), "");
317
318 let array_source = RecipientSource::Async(Arc::new(|_: &Exchange| {
319 Box::pin(async { Ok(Value::Array(vec![])) }) as BoxValueFuture
320 }));
321 match array_source.resolve(&ex).await.unwrap_err() {
322 CamelError::ProcessorError(msg) => {
323 assert!(msg.contains("non-scalar"), "missing non-scalar: {msg}");
324 }
325 other => panic!("expected ProcessorError, got {other:?}"),
326 }
327 }
328}