Skip to main content

camel_processor/
validate.rs

1use std::future::Future;
2use std::pin::Pin;
3use std::task::{Context, Poll};
4
5use camel_api::{CamelError, Exchange, FilterPredicate, PredicateSource};
6use tower::Service;
7
8/// Tower Service implementing the Validate EIP.
9///
10/// If the predicate evaluates to `true`, the exchange continues (returned as `Ok`).
11/// If `false`, a `CamelError::ValidationError` is returned. A failed (async)
12/// predicate evaluation surfaces that typed error directly — it is never
13/// rewritten into a `ValidationError`.
14#[derive(Clone)]
15pub struct ValidateService {
16    predicate: PredicateSource,
17    expression_source: String,
18}
19
20impl ValidateService {
21    /// Create from a closure predicate and an expression source string (for error messages).
22    pub fn new(
23        predicate: impl Fn(&Exchange) -> bool + Send + Sync + 'static,
24        expression_source: impl Into<String>,
25    ) -> Self {
26        Self {
27            predicate: PredicateSource::Sync(FilterPredicate::new(predicate)),
28            expression_source: expression_source.into(),
29        }
30    }
31
32    /// Create from a fallible `PredicateSource` (used by `resolve_steps`).
33    pub fn from_predicate(
34        predicate: PredicateSource,
35        expression_source: impl Into<String>,
36    ) -> Self {
37        Self {
38            predicate,
39            expression_source: expression_source.into(),
40        }
41    }
42}
43
44impl Service<Exchange> for ValidateService {
45    type Response = Exchange;
46    type Error = CamelError;
47    type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
48
49    fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
50        Poll::Ready(Ok(()))
51    }
52
53    fn call(&mut self, exchange: Exchange) -> Self::Future {
54        // Clone-and-replace: the future owns its state so the predicate can be
55        // awaited before deciding the outcome.
56        let predicate = self.predicate.clone();
57        let source = self.expression_source.clone();
58        Box::pin(async move {
59            match predicate.matches(&exchange).await {
60                Ok(true) => Ok(exchange),
61                Ok(false) => Err(CamelError::ValidationError(format!(
62                    "validate('{source}'): predicate returned false",
63                ))),
64                // Typed predicate failure: surface as-is, not ValidationError.
65                Err(err) => Err(err),
66            }
67        })
68    }
69}
70
71#[cfg(test)]
72mod tests {
73    use super::*;
74    use camel_api::Message;
75
76    // ── ValidateService tests ──
77
78    // 1. Passing predicate returns Ok(exchange)
79    #[tokio::test]
80    async fn test_validate_passing_predicate_returns_ok() {
81        let mut svc = ValidateService::new(|_ex: &Exchange| true, "true predicate");
82        let ex = Exchange::new(Message::new("hello"));
83        let result = svc.call(ex).await;
84        assert!(result.is_ok());
85        assert_eq!(result.unwrap().input.body.as_text(), Some("hello"));
86    }
87
88    // 2. Failing predicate returns Err(ValidationError)
89    #[tokio::test]
90    async fn test_validate_failing_predicate_returns_err() {
91        let mut svc = ValidateService::new(|_ex: &Exchange| false, "false predicate");
92        let ex = Exchange::new(Message::new("hello"));
93        let result = svc.call(ex).await;
94        assert!(result.is_err());
95        match result.unwrap_err() {
96            CamelError::ValidationError(msg) => {
97                assert!(
98                    msg.contains("false predicate"),
99                    "error message should contain expression source, got: {msg}"
100                );
101            }
102            other => panic!("expected ValidationError, got: {other:?}"),
103        }
104    }
105
106    // 3. Predicate evaluates the exchange (body-based validation)
107    #[tokio::test]
108    async fn test_validate_predicate_evaluates_body() {
109        let mut svc = ValidateService::new(
110            |ex: &Exchange| ex.input.body.as_text().is_some_and(|s| s.len() > 3),
111            "body length > 3",
112        );
113        let short = Exchange::new(Message::new("ab"));
114        let long = Exchange::new(Message::new("abcdef"));
115
116        assert!(svc.call(short).await.is_err());
117        assert!(svc.call(long).await.is_ok());
118    }
119
120    // 4. ValidateService is Clone
121    #[tokio::test]
122    async fn test_validate_clone_is_independent() {
123        let svc = ValidateService::new(|_ex: &Exchange| true, "true predicate");
124        let mut cloned = svc.clone();
125        let ex = Exchange::new(Message::new("hi"));
126        let result = cloned.call(ex).await;
127        assert!(result.is_ok());
128    }
129
130    // 5. poll_ready is always Ready(Ok(()))
131    #[tokio::test]
132    async fn test_validate_poll_ready() {
133        let mut svc = ValidateService::new(|_ex: &Exchange| true, "true predicate");
134        let poll = svc.poll_ready(&mut Context::from_waker(futures::task::noop_waker_ref()));
135        assert!(poll.is_ready());
136        // unwrap the Poll<Result<...>>
137        match poll {
138            std::task::Poll::Ready(Ok(())) => {}
139            other => panic!("expected Ready(Ok(())), got: {other:?}"),
140        }
141    }
142
143    // ── Fallible predicate path (language-value-boundary task 1.4) ──
144
145    use std::sync::Arc;
146
147    use camel_api::{ExpressionErrorClass, PredicateSource};
148
149    fn expression_failed() -> CamelError {
150        CamelError::ExpressionFailed {
151            language: "rhai".to_string(),
152            route_id: "r1".to_string(),
153            step_id: "step#0".to_string(),
154            verb: "validate".to_string(),
155            class: ExpressionErrorClass::Runtime,
156            position: None,
157            conversion: None,
158            cause: None,
159        }
160    }
161
162    fn async_err_predicate(err: CamelError) -> PredicateSource {
163        PredicateSource::Async(Arc::new(move |_: &Exchange| {
164            let err = err.clone();
165            Box::pin(async move { Err(err) }) as camel_api::BoxBoolFuture
166        }))
167    }
168
169    #[tokio::test]
170    async fn validate_predicate_error_is_typed_not_validation() {
171        let mut svc =
172            ValidateService::from_predicate(async_err_predicate(expression_failed()), "expr");
173        let result = svc.call(Exchange::new(Message::new("x"))).await;
174        match result {
175            Err(CamelError::ExpressionFailed { .. }) => {}
176            Err(other) => panic!("expected ExpressionFailed, got {other:?}"),
177            Ok(_) => panic!("expected Err, got Ok"),
178        }
179    }
180}