camel_processor/
validate.rs1use 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#[derive(Clone)]
15pub struct ValidateService {
16 predicate: PredicateSource,
17 expression_source: String,
18}
19
20impl ValidateService {
21 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 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 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 Err(err) => Err(err),
66 }
67 })
68 }
69}
70
71#[cfg(test)]
72mod tests {
73 use super::*;
74 use camel_api::Message;
75
76 #[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 #[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 #[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 #[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 #[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 match poll {
138 std::task::Poll::Ready(Ok(())) => {}
139 other => panic!("expected Ready(Ok(())), got: {other:?}"),
140 }
141 }
142
143 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}