Skip to main content

dataflow_rs/engine/functions/
validation.rs

1//! # Validation Function Module
2//!
3//! This module provides rule-based validation capabilities using JSONLogic expressions.
4//! The validation function evaluates a set of rules against message data and collects
5//! any validation errors that occur.
6//!
7//! ## Features
8//!
9//! - Define validation rules using JSONLogic expressions
10//! - Custom error messages for each rule
11//! - Non-destructive: validation is read-only and doesn't modify message data
12//! - Errors are collected in the message's error list
13//!
14//! ## Example Usage
15//!
16//! ```json
17//! {
18//!     "name": "validation",
19//!     "input": {
20//!         "rules": [
21//!             {
22//!                 "logic": {"!!": [{"var": "data.email"}]},
23//!                 "message": "Email is required"
24//!             },
25//!             {
26//!                 "logic": {">": [{"var": "data.age"}, 0]},
27//!                 "message": "Age must be positive"
28//!             }
29//!         ]
30//!     }
31//! }
32//! ```
33
34use crate::engine::error::{DataflowError, ErrorInfo, Result};
35use crate::engine::executor::{ArenaContext, with_arena};
36use crate::engine::message::{Change, Message};
37use crate::engine::task_outcome::TaskOutcome;
38use datalogic_rs::{Engine, Logic};
39use datavalue::DataValue;
40use log::{debug, error};
41use serde::Deserialize;
42use serde_json::Value;
43use std::sync::Arc;
44
45/// Configuration for the validation function containing a list of rules.
46///
47/// Each rule specifies a JSONLogic condition that must evaluate to `true`
48/// for the validation to pass. If a rule evaluates to anything other than
49/// `true`, its error message is added to the message's error list.
50#[derive(Debug, Clone, Deserialize)]
51pub struct ValidationConfig {
52    /// List of validation rules to evaluate.
53    pub rules: Vec<ValidationRule>,
54}
55
56/// A single validation rule with a condition and error message.
57///
58/// The rule's logic is evaluated against the message context. If it does not
59/// return exactly `true`, the validation fails and the error message is recorded.
60#[derive(Debug, Clone, Deserialize)]
61pub struct ValidationRule {
62    /// JSONLogic expression that must evaluate to `true` for validation to pass.
63    /// Any other result (false, null, etc.) is considered a validation failure.
64    pub logic: Value,
65
66    /// Error message to display if validation fails.
67    ///
68    /// Required when a rule is deserialized as part of a workflow definition —
69    /// which is the path `Engine::build` takes, so a rule without it is
70    /// rejected at build time. [`ValidationConfig::from_json`], the standalone
71    /// parser, is the one path that substitutes `"Validation failed"`.
72    pub message: String,
73
74    /// Pre-compiled JSONLogic, populated by `LogicCompiler`. `None` is
75    /// recorded as a `COMPILATION_ERROR` at execute time.
76    #[serde(skip)]
77    pub compiled_logic: Option<Arc<Logic>>,
78}
79
80impl ValidationConfig {
81    /// Parses a `ValidationConfig` from a JSON value.
82    ///
83    /// # Arguments
84    /// * `input` - JSON object containing a "rules" array
85    ///
86    /// # Errors
87    /// Returns `DataflowError::Validation` if:
88    /// - The "rules" field is missing
89    /// - The "rules" field is not an array
90    /// - Any rule is missing the "logic" field
91    pub fn from_json(input: &Value) -> Result<Self> {
92        let rules = input.get("rules").ok_or_else(|| {
93            DataflowError::Validation("Missing 'rules' array in input".to_string())
94        })?;
95
96        let rules_arr = rules
97            .as_array()
98            .ok_or_else(|| DataflowError::Validation("'rules' must be an array".to_string()))?;
99
100        let mut parsed_rules = Vec::new();
101
102        for rule in rules_arr {
103            let logic = rule
104                .get("logic")
105                .ok_or_else(|| DataflowError::Validation("Missing 'logic' in rule".to_string()))?
106                .clone();
107
108            let message = rule
109                .get("message")
110                .and_then(Value::as_str)
111                .unwrap_or("Validation failed")
112                .to_string();
113
114            parsed_rules.push(ValidationRule {
115                logic,
116                message,
117                compiled_logic: None,
118            });
119        }
120
121        Ok(ValidationConfig {
122            rules: parsed_rules,
123        })
124    }
125
126    /// Executes all validation rules using pre-compiled logic.
127    ///
128    /// Evaluates each rule sequentially against the message context.
129    /// This is a read-only operation that does not modify message data.
130    ///
131    /// # Arguments
132    /// * `message` - The message to validate (errors are added to its error list)
133    /// * `engine` - Datalogic v5 engine for evaluation
134    ///
135    /// # Returns
136    /// * `Ok((TaskOutcome::Success, []))` — all rules passed
137    /// * `Ok((TaskOutcome::Status(400), []))` — one or more rules failed,
138    ///   `ErrorInfo` entries pushed onto `message.errors`
139    pub fn execute(
140        &self,
141        message: &mut Message,
142        engine: &Arc<Engine>,
143    ) -> Result<(TaskOutcome, Vec<Change>)> {
144        // Default path: open the arena and convert context once for this
145        // task call. When called from the workflow-level sync-stretch
146        // executor (`execute_in_arena`), the conversion is reused across
147        // multiple tasks in the same stretch.
148        with_arena(|arena| {
149            let ctx_av: DataValue<'_> = message.context.to_arena(arena);
150            self.run_rules(message, ctx_av, arena, engine)
151        })
152    }
153
154    /// Run validation rules against an externally-provided `ArenaContext`.
155    /// Reuses the cached arena form built by an earlier task in the same
156    /// workflow sync stretch — the heavy `data.input` subtree stays cached
157    /// across the parse_json → map → validation pipeline.
158    pub(crate) fn execute_in_arena(
159        &self,
160        message: &mut Message,
161        arena_ctx: &mut ArenaContext<'_>,
162        engine: &Arc<Engine>,
163    ) -> Result<(TaskOutcome, Vec<Change>)> {
164        let arena = arena_ctx.arena();
165        let ctx_av = arena_ctx.as_data_value();
166        self.run_rules(message, ctx_av, arena, engine)
167    }
168
169    /// Shared inner loop: evaluate each rule against `ctx_av` and record
170    /// `ErrorInfo` entries for any failures.
171    fn run_rules(
172        &self,
173        message: &mut Message,
174        ctx_av: DataValue<'_>,
175        arena: &bumpalo::Bump,
176        engine: &Arc<Engine>,
177    ) -> Result<(TaskOutcome, Vec<Change>)> {
178        let changes = Vec::new();
179        let mut validation_errors = Vec::new();
180
181        for (idx, rule) in self.rules.iter().enumerate() {
182            debug!("Processing validation rule {}: {}", idx, rule.message);
183
184            let compiled_logic = match &rule.compiled_logic {
185                Some(logic) => logic,
186                None => {
187                    error!("Validation: Logic not compiled for rule at index {}", idx);
188                    validation_errors.push(ErrorInfo::simple_ref(
189                        "COMPILATION_ERROR",
190                        &format!("Logic not compiled for rule at index: {}", idx),
191                        None,
192                    ));
193                    continue;
194                }
195            };
196
197            // Reuse the pre-converted `ctx_av` (DataValue is Copy). The
198            // result is `&DataValue<'_>` borrowed from the arena — we
199            // only need to peek at the discriminant so we skip the
200            // `to_owned()` deep-clone too.
201            match engine.evaluate(compiled_logic, ctx_av, arena) {
202                Ok(value) => {
203                    if !matches!(value, DataValue::Bool(true)) {
204                        debug!("Validation failed for rule {}: {}", idx, rule.message);
205                        validation_errors.push(ErrorInfo::simple_ref(
206                            "VALIDATION_ERROR",
207                            &rule.message,
208                            None,
209                        ));
210                    } else {
211                        debug!("Validation passed for rule {}", idx);
212                    }
213                }
214                Err(e) => {
215                    error!("Validation: Error evaluating rule {}: {:?}", idx, e);
216                    validation_errors.push(ErrorInfo::simple_ref(
217                        "EVALUATION_ERROR",
218                        &format!("Failed to evaluate rule {}: {}", idx, e),
219                        None,
220                    ));
221                }
222            }
223        }
224
225        if !validation_errors.is_empty() {
226            message.errors.extend(validation_errors);
227            Ok((TaskOutcome::Status(400), changes))
228        } else {
229            Ok((TaskOutcome::Success, changes))
230        }
231    }
232}
233
234#[cfg(test)]
235mod tests {
236    use super::*;
237    use datavalue::OwnedDataValue;
238    use serde_json::json;
239
240    #[test]
241    fn test_validation_config_from_json() {
242        let input = json!({
243            "rules": [
244                {
245                    "logic": {"!!": [{"var": "data.required_field"}]},
246                    "path": "data",
247                    "message": "Required field is missing"
248                },
249                {
250                    "logic": {">": [{"var": "data.age"}, 18]},
251                    "message": "Must be over 18"
252                }
253            ]
254        });
255
256        let config = ValidationConfig::from_json(&input).unwrap();
257        assert_eq!(config.rules.len(), 2);
258        assert_eq!(config.rules[0].message, "Required field is missing");
259        assert_eq!(config.rules[1].message, "Must be over 18");
260    }
261
262    #[test]
263    fn test_validation_config_missing_rules() {
264        let input = json!({});
265        let result = ValidationConfig::from_json(&input);
266        assert!(result.is_err());
267    }
268
269    #[test]
270    fn test_validation_config_invalid_rules() {
271        let input = json!({
272            "rules": "not_an_array"
273        });
274        let result = ValidationConfig::from_json(&input);
275        assert!(result.is_err());
276    }
277
278    #[test]
279    fn test_validation_config_missing_logic() {
280        let input = json!({
281            "rules": [
282                {
283                    "path": "data",
284                    "message": "Some error"
285                }
286            ]
287        });
288        let result = ValidationConfig::from_json(&input);
289        assert!(result.is_err());
290    }
291
292    #[test]
293    fn test_validation_config_defaults() {
294        let input = json!({
295            "rules": [
296                {
297                    "logic": {"var": "data.field"}
298                }
299            ]
300        });
301
302        let config = ValidationConfig::from_json(&input).unwrap();
303        assert_eq!(config.rules[0].message, "Validation failed");
304    }
305
306    fn dv(v: serde_json::Value) -> OwnedDataValue {
307        OwnedDataValue::from(&v)
308    }
309
310    fn message_with_data(initial: serde_json::Value) -> crate::engine::message::Message {
311        use crate::engine::message::Message;
312        Message::builder().data(dv(initial)).build()
313    }
314
315    /// Compile each rule's `logic` and stamp the resulting `Arc<Logic>` into
316    /// the `compiled_logic` slot — mirroring `LogicCompiler`.
317    fn compile_rules(engine: &Arc<Engine>, config: &mut ValidationConfig) {
318        for rule in &mut config.rules {
319            rule.compiled_logic = Some(engine.compile_arc(&rule.logic).unwrap());
320        }
321    }
322
323    #[test]
324    fn test_validation_execute_passes() {
325        let engine = Arc::new(Engine::builder().with_templating(true).build());
326
327        let mut message = message_with_data(json!({
328            "email": "test@example.com",
329            "age": 25
330        }));
331
332        let mut config = ValidationConfig {
333            rules: vec![
334                ValidationRule {
335                    logic: json!({"!!": [{"var": "data.email"}]}),
336                    message: "Email is required".to_string(),
337                    compiled_logic: None,
338                },
339                ValidationRule {
340                    logic: json!({">": [{"var": "data.age"}, 18]}),
341                    message: "Must be over 18".to_string(),
342                    compiled_logic: None,
343                },
344            ],
345        };
346        compile_rules(&engine, &mut config);
347
348        let result = config.execute(&mut message, &engine);
349        assert!(result.is_ok());
350
351        let (outcome, changes) = result.unwrap();
352        assert_eq!(outcome, TaskOutcome::Success);
353        assert!(changes.is_empty());
354        assert!(message.errors.is_empty());
355    }
356
357    #[test]
358    fn test_validation_execute_fails() {
359        let engine = Arc::new(Engine::builder().with_templating(true).build());
360
361        let mut message = message_with_data(json!({ "age": 15 }));
362
363        let mut config = ValidationConfig {
364            rules: vec![
365                ValidationRule {
366                    logic: json!({"!!": [{"var": "data.email"}]}),
367                    message: "Email is required".to_string(),
368                    compiled_logic: None,
369                },
370                ValidationRule {
371                    logic: json!({">": [{"var": "data.age"}, 18]}),
372                    message: "Must be over 18".to_string(),
373                    compiled_logic: None,
374                },
375            ],
376        };
377        compile_rules(&engine, &mut config);
378
379        let result = config.execute(&mut message, &engine);
380        assert!(result.is_ok());
381
382        let (outcome, _changes) = result.unwrap();
383        assert_eq!(outcome, TaskOutcome::Status(400));
384        assert_eq!(message.errors.len(), 2);
385
386        let error_messages: Vec<&str> = message.errors.iter().map(|e| e.message.as_str()).collect();
387        assert!(error_messages.contains(&"Email is required"));
388        assert!(error_messages.contains(&"Must be over 18"));
389    }
390
391    #[test]
392    fn test_validation_uncompiled_logic() {
393        use crate::engine::message::Message;
394
395        let engine = Arc::new(Engine::builder().with_templating(true).build());
396
397        let mut message = Message::new(Arc::new(dv(json!({}))));
398
399        let config = ValidationConfig {
400            rules: vec![ValidationRule {
401                logic: json!(true),
402                message: "Test".to_string(),
403                compiled_logic: None,
404            }],
405        };
406
407        let result = config.execute(&mut message, &engine);
408        assert!(result.is_ok());
409
410        let (outcome, _) = result.unwrap();
411        assert_eq!(outcome, TaskOutcome::Status(400));
412        assert!(!message.errors.is_empty());
413        assert!(message.errors[0].code == "COMPILATION_ERROR");
414    }
415}