use crate::error::CamelError;
use crate::exchange::Exchange;
use crate::value::Value;
use std::sync::Arc;
pub type BoxValueFuture =
std::pin::Pin<Box<dyn std::future::Future<Output = Result<Value, CamelError>> + Send>>;
pub type BoxBoolFuture =
std::pin::Pin<Box<dyn std::future::Future<Output = Result<bool, CamelError>> + Send>>;
pub struct FilterPredicate(pub Arc<dyn Fn(&Exchange) -> bool + Send + Sync>);
impl FilterPredicate {
pub fn new<F>(f: F) -> Self
where
F: Fn(&Exchange) -> bool + Send + Sync + 'static,
{
FilterPredicate(Arc::new(f))
}
}
impl Clone for FilterPredicate {
fn clone(&self) -> Self {
FilterPredicate(Arc::clone(&self.0))
}
}
impl std::fmt::Debug for FilterPredicate {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("FilterPredicate(..)")
}
}
impl std::ops::Deref for FilterPredicate {
type Target = dyn Fn(&Exchange) -> bool + Send + Sync;
fn deref(&self) -> &Self::Target {
&*self.0
}
}
#[derive(Clone)]
#[non_exhaustive]
pub enum PredicateSource {
Sync(FilterPredicate),
Async(Arc<dyn Fn(&Exchange) -> BoxBoolFuture + Send + Sync>),
}
impl PredicateSource {
pub async fn matches(&self, exchange: &Exchange) -> Result<bool, CamelError> {
match self {
Self::Sync(pred) => Ok(pred(exchange)),
Self::Async(f) => f(exchange).await,
}
}
}
impl From<FilterPredicate> for PredicateSource {
fn from(pred: FilterPredicate) -> Self {
Self::Sync(pred)
}
}
#[derive(Clone)]
#[non_exhaustive]
pub enum ValueSource {
Sync(Arc<dyn Fn(&Exchange) -> Value + Send + Sync>),
Async(Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>),
}
impl ValueSource {
pub async fn evaluate(&self, exchange: &Exchange) -> Result<Value, CamelError> {
match self {
Self::Sync(f) => Ok(f(exchange)),
Self::Async(f) => f(exchange).await,
}
}
}
const NON_SCALAR_TARGET: &str =
"router target expression returned a non-scalar value (array/object); expected a string target";
fn coerce_target_value(value: Value) -> Result<Option<String>, CamelError> {
match value {
Value::Null => Ok(None),
Value::String(s) => Ok(Some(s)),
Value::Array(_) | Value::Object(_) => {
Err(CamelError::ProcessorError(NON_SCALAR_TARGET.to_string()))
}
other => Ok(Some(other.to_string())),
}
}
#[allow(clippy::type_complexity)]
#[derive(Clone)]
#[non_exhaustive]
pub enum TargetSource {
Sync(Arc<dyn Fn(&Exchange) -> Option<String> + Send + Sync>),
Async(Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>),
}
impl TargetSource {
pub async fn resolve(&self, exchange: &Exchange) -> Result<Option<String>, CamelError> {
match self {
Self::Sync(f) => Ok(f(exchange)),
Self::Async(f) => coerce_target_value(f(exchange).await?),
}
}
}
#[derive(Clone)]
#[non_exhaustive]
pub enum RecipientSource {
Sync(Arc<dyn Fn(&Exchange) -> String + Send + Sync>),
Async(Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>),
}
impl RecipientSource {
pub async fn resolve(&self, exchange: &Exchange) -> Result<String, CamelError> {
match self {
Self::Sync(f) => Ok(f(exchange)),
Self::Async(f) => Ok(coerce_target_value(f(exchange).await?)?.unwrap_or_default()),
}
}
}
impl From<Arc<dyn Fn(&Exchange) -> Value + Send + Sync>> for ValueSource {
fn from(f: Arc<dyn Fn(&Exchange) -> Value + Send + Sync>) -> Self {
Self::Sync(f)
}
}
impl From<Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>> for ValueSource {
fn from(f: Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>) -> Self {
Self::Async(f)
}
}
impl From<Arc<dyn Fn(&Exchange) -> Option<String> + Send + Sync>> for TargetSource {
fn from(f: Arc<dyn Fn(&Exchange) -> Option<String> + Send + Sync>) -> Self {
Self::Sync(f)
}
}
impl From<Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>> for TargetSource {
fn from(f: Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>) -> Self {
Self::Async(f)
}
}
impl From<Arc<dyn Fn(&Exchange) -> String + Send + Sync>> for RecipientSource {
fn from(f: Arc<dyn Fn(&Exchange) -> String + Send + Sync>) -> Self {
Self::Sync(f)
}
}
impl From<Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>> for RecipientSource {
fn from(f: Arc<dyn Fn(&Exchange) -> BoxValueFuture + Send + Sync>) -> Self {
Self::Async(f)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{Exchange, ExpressionErrorClass, Message};
fn expression_failed() -> CamelError {
CamelError::ExpressionFailed {
language: "rhai".to_string(),
route_id: "r1".to_string(),
step_id: "step#0".to_string(),
verb: "filter".to_string(),
class: ExpressionErrorClass::Runtime,
position: None,
conversion: None,
cause: None,
}
}
#[test]
fn test_filter_predicate_is_callable() {
let pred = FilterPredicate::new(|ex: &Exchange| ex.input.body.as_text().is_some());
let ex = Exchange::new(Message::new("hello"));
assert!(pred(&ex));
}
#[test]
fn test_filter_predicate_debug_is_redacted() {
let pred = FilterPredicate::new(|_: &Exchange| true);
assert_eq!(format!("{pred:?}"), "FilterPredicate(..)");
}
#[test]
fn test_filter_predicate_clone_shares_arc() {
let pred = FilterPredicate::new(|_: &Exchange| true);
let cloned = pred.clone();
assert!(matches!(cloned, FilterPredicate(_)));
}
#[tokio::test]
async fn predicate_source_sync_returns_bool() {
let source = PredicateSource::Sync(FilterPredicate::new(|_: &Exchange| true));
let ex = Exchange::new(Message::new("hello"));
assert!(source.matches(&ex).await.unwrap());
}
#[tokio::test]
async fn predicate_source_async_propagates_error() {
let source = PredicateSource::Async(Arc::new(|_: &Exchange| {
Box::pin(async { Err(expression_failed()) }) as BoxBoolFuture
}));
let ex = Exchange::new(Message::new("hello"));
let err = source.matches(&ex).await.unwrap_err();
assert!(matches!(err, CamelError::ExpressionFailed { .. }));
}
#[tokio::test]
async fn value_source_async_propagates_error() {
let source = ValueSource::Async(Arc::new(|_: &Exchange| {
Box::pin(async { Err(expression_failed()) }) as BoxValueFuture
}));
let ex = Exchange::new(Message::new("hello"));
let err = source.evaluate(&ex).await.unwrap_err();
assert!(matches!(err, CamelError::ExpressionFailed { .. }));
}
#[tokio::test]
async fn target_source_null_maps_to_none_and_error_propagates() {
let ex = Exchange::new(Message::new("hello"));
let null_source = TargetSource::Async(Arc::new(|_: &Exchange| {
Box::pin(async { Ok(Value::Null) }) as BoxValueFuture
}));
assert_eq!(null_source.resolve(&ex).await.unwrap(), None);
let err_source = TargetSource::Async(Arc::new(|_: &Exchange| {
Box::pin(async { Err(expression_failed()) }) as BoxValueFuture
}));
let err = err_source.resolve(&ex).await.unwrap_err();
assert!(matches!(err, CamelError::ExpressionFailed { .. }));
}
#[tokio::test]
async fn target_source_array_is_error() {
let source = TargetSource::Async(Arc::new(|_: &Exchange| {
Box::pin(async { Ok(Value::Array(vec![Value::String("a".into())])) }) as BoxValueFuture
}));
let ex = Exchange::new(Message::new("hello"));
match source.resolve(&ex).await.unwrap_err() {
CamelError::ProcessorError(msg) => {
assert!(msg.contains("non-scalar"), "missing non-scalar: {msg}");
assert!(msg.contains("string target"), "missing target: {msg}");
}
other => panic!("expected ProcessorError, got {other:?}"),
}
}
#[tokio::test]
async fn recipient_source_null_maps_to_empty_and_non_scalar_is_error() {
let ex = Exchange::new(Message::new("hello"));
let null_source = RecipientSource::Async(Arc::new(|_: &Exchange| {
Box::pin(async { Ok(Value::Null) }) as BoxValueFuture
}));
assert_eq!(null_source.resolve(&ex).await.unwrap(), "");
let array_source = RecipientSource::Async(Arc::new(|_: &Exchange| {
Box::pin(async { Ok(Value::Array(vec![])) }) as BoxValueFuture
}));
match array_source.resolve(&ex).await.unwrap_err() {
CamelError::ProcessorError(msg) => {
assert!(msg.contains("non-scalar"), "missing non-scalar: {msg}");
}
other => panic!("expected ProcessorError, got {other:?}"),
}
}
}