Skip to main content

camel_api/
filter.rs

1use crate::error::CamelError;
2use crate::exchange::Exchange;
3use crate::value::Value;
4use std::sync::Arc;
5
6/// Boxed future yielding a dynamic [`Value`] or a [`CamelError`].
7///
8/// Used by the async arms of [`ValueSource`], [`TargetSource`] and
9/// [`RecipientSource`] so language expression engines can be evaluated
10/// lazily and fallibly from sync call sites.
11pub type BoxValueFuture =
12    std::pin::Pin<Box<dyn std::future::Future<Output = Result<Value, CamelError>> + Send>>;
13
14/// Boxed future yielding a `bool` predicate result or a [`CamelError`].
15///
16/// Used by the async arm of [`PredicateSource`].
17pub type BoxBoolFuture =
18    std::pin::Pin<Box<dyn std::future::Future<Output = Result<bool, CamelError>> + Send>>;
19
20/// Predicate that determines whether an exchange passes the filter.
21/// Returns `true` to forward the exchange into the filter body; `false` to skip it.
22///
23/// This is a newtype around `Arc<dyn Fn(&Exchange) -> bool + Send + Sync>` so that
24/// it can implement `Debug` (used by `#[derive(Debug)]` on `BuilderStep` and friends).
25/// The `Deref` impl keeps the call-site ergonomic: `predicate(&exchange)` still works
26/// because `FilterPredicate` derefs to the inner `dyn Fn`.
27///
28/// Pre-v1.0: this used to be a type alias. Converted to a newtype (H2) so the
29/// containing structs (`WhenStep`, etc.) can be `#[derive(Debug)]`.
30pub struct FilterPredicate(pub Arc<dyn Fn(&Exchange) -> bool + Send + Sync>);
31
32impl FilterPredicate {
33    /// Create a new `FilterPredicate` from a closure or function pointer.
34    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/// Fallible predicate source: either a synchronous [`FilterPredicate`] or an
63/// asynchronous, fallible language predicate.
64///
65/// The sync arm always succeeds (`Ok(bool)`); the async arm may fail with a
66/// [`CamelError`] that callers must propagate instead of silently treating the
67/// exchange as skipped.
68#[derive(Clone)]
69#[non_exhaustive]
70pub enum PredicateSource {
71    /// Programmatic synchronous predicate.
72    Sync(FilterPredicate),
73    /// Language-backed asynchronous predicate.
74    Async(Arc<dyn Fn(&Exchange) -> BoxBoolFuture + Send + Sync>),
75}
76
77impl PredicateSource {
78    /// Evaluate the predicate against `exchange`, propagating async failures.
79    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/// Fallible value source for dynamic setters and log messages.
94#[derive(Clone)]
95#[non_exhaustive]
96pub enum ValueSource {
97    /// Programmatic synchronous closure producing a [`Value`].
98    Sync(Arc<dyn Fn(&Exchange) -> Value + Send + Sync>),
99    /// Language-backed asynchronous expression.
100    Async(Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>),
101}
102
103impl ValueSource {
104    /// Evaluate the value against `exchange`, propagating async failures.
105    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
113/// Error returned when a target/recipient expression yields a non-scalar.
114const NON_SCALAR_TARGET: &str =
115    "router target expression returned a non-scalar value (array/object); expected a string target";
116
117/// Coerce a dynamic [`Value`] into an optional target string.
118///
119/// `Null` maps to `None`; strings pass through; other scalars are stringified;
120/// arrays and objects are rejected.
121fn 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/// Fallible target source for dynamic routers and routing slips.
133// The closure shape is part of the published contract; factoring it into a
134// private alias would not change the type, so keep the signature literal.
135#[allow(clippy::type_complexity)]
136#[derive(Clone)]
137#[non_exhaustive]
138pub enum TargetSource {
139    /// Programmatic synchronous closure producing an optional target.
140    Sync(Arc<dyn Fn(&Exchange) -> Option<String> + Send + Sync>),
141    /// Language-backed asynchronous expression.
142    Async(Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>),
143}
144
145impl TargetSource {
146    /// Resolve the target, propagating async failures and non-scalar values.
147    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/// Fallible recipient source for the recipient list EIP.
156#[derive(Clone)]
157#[non_exhaustive]
158pub enum RecipientSource {
159    /// Programmatic synchronous closure producing a recipient.
160    Sync(Arc<dyn Fn(&Exchange) -> String + Send + Sync>),
161    /// Language-backed asynchronous expression.
162    Async(Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>),
163}
164
165impl RecipientSource {
166    /// Resolve the recipient, propagating async failures and non-scalar values.
167    ///
168    /// `Null` coerces to the empty string; the empty string is allowed.
169    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}