use crate::engine::error::Result;
use crate::engine::executor::{ArenaContext, with_arena};
use crate::engine::message::{Change, Message};
use crate::engine::task_outcome::TaskOutcome;
use datalogic_rs::{Engine, Logic};
use datavalue::DataValue;
use log::{debug, info};
use serde::Deserialize;
use serde_json::Value;
use std::sync::Arc;
#[derive(Debug, Clone, Deserialize, Default)]
#[serde(rename_all = "lowercase")]
pub enum RejectAction {
#[default]
Halt,
Skip,
}
#[derive(Debug, Clone, Deserialize)]
pub struct FilterConfig {
pub condition: Value,
#[serde(default)]
pub on_reject: RejectAction,
#[serde(skip)]
pub compiled_condition: Option<Arc<Logic>>,
}
impl FilterConfig {
pub fn execute(
&self,
message: &mut Message,
engine: &Arc<Engine>,
) -> Result<(TaskOutcome, Vec<Change>)> {
with_arena(|arena| {
let mut arena_ctx = ArenaContext::from_owned(&message.context, arena);
self.execute_in_arena(message, &mut arena_ctx, engine)
})
}
pub(crate) fn execute_in_arena(
&self,
_message: &mut Message,
arena_ctx: &mut ArenaContext<'_>,
engine: &Arc<Engine>,
) -> Result<(TaskOutcome, Vec<Change>)> {
let condition_met = match &self.compiled_condition {
Some(compiled) => {
let ctx_av = arena_ctx.as_data_value();
match engine.evaluate(compiled, ctx_av, arena_ctx.arena()) {
Ok(DataValue::Bool(true)) => true,
Ok(_) => false,
Err(e) => {
debug!("Filter: condition evaluation error: {:?}", e);
false
}
}
}
None => {
debug!("Filter: condition not compiled, treating as not met");
false
}
};
if condition_met {
debug!("Filter: condition passed");
Ok((TaskOutcome::Success, vec![]))
} else {
match self.on_reject {
RejectAction::Halt => {
info!("Filter: condition not met, halting workflow");
Ok((TaskOutcome::Halt, vec![]))
}
RejectAction::Skip => {
debug!("Filter: condition not met, skipping");
Ok((TaskOutcome::Skip, vec![]))
}
}
}
}
}