dataflow_rs/engine/functions/
validation.rs1use 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#[derive(Debug, Clone, Deserialize)]
51pub struct ValidationConfig {
52 pub rules: Vec<ValidationRule>,
54}
55
56#[derive(Debug, Clone, Deserialize)]
61pub struct ValidationRule {
62 pub logic: Value,
65
66 pub message: String,
73
74 #[serde(skip)]
77 pub compiled_logic: Option<Arc<Logic>>,
78}
79
80impl ValidationConfig {
81 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 pub fn execute(
140 &self,
141 message: &mut Message,
142 engine: &Arc<Engine>,
143 ) -> Result<(TaskOutcome, Vec<Change>)> {
144 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 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 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 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 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}