1use crate::engine::error::{DataflowError, Result};
2use crate::engine::functions::{ConnectorName, FunctionConfig};
3use crate::engine::task::Task;
4use chrono::{DateTime, Utc};
5use datalogic_rs::Logic;
6use serde::{Deserialize, Serialize};
7use serde_json::Value;
8use std::fs;
9use std::path::Path;
10use std::sync::Arc;
11
12pub use crate::engine::rollout::{Rollout, RolloutError};
13
14#[derive(Clone, Debug, Deserialize, PartialEq, Eq)]
41pub struct LoopConfig {
42 #[serde(default)]
53 pub counter: Option<String>,
54
55 #[serde(default)]
57 pub init: i64,
58
59 #[serde(default = "default_increment")]
62 pub increment: i64,
63
64 pub max: i64,
69
70 #[doc(hidden)]
74 #[serde(skip)]
75 pub counter_parts: Arc<[Arc<str>]>,
76}
77
78fn default_increment() -> i64 {
79 1
80}
81
82impl LoopConfig {
83 fn validate(&self, workflow_id: &str) -> Result<()> {
87 if self.increment < 1 {
88 return Err(DataflowError::Workflow(format!(
89 "Workflow {workflow_id}: loop increment must be >= 1, got {} \
90 (a non-advancing counter would never reach max)",
91 self.increment
92 )));
93 }
94 if self.max <= self.init {
95 return Err(DataflowError::Workflow(format!(
96 "Workflow {workflow_id}: loop max ({}) must be greater than init ({}) — \
97 the bound is half-open, so this could never run a sweep",
98 self.max, self.init
99 )));
100 }
101 if let Some(counter) = &self.counter {
102 if counter.is_empty() || counter.split('.').any(str::is_empty) {
103 return Err(DataflowError::Workflow(format!(
104 "Workflow {workflow_id}: loop counter must be a non-empty \
105 temp_data field path, got {counter:?}"
106 )));
107 }
108 }
109 Ok(())
110 }
111
112 #[doc(hidden)]
117 pub fn precompute_counter_path(&mut self) {
118 self.counter_parts = match &self.counter {
119 Some(counter) => crate::engine::utils::compute_path_parts("temp_data", counter),
120 None => Arc::from([] as [Arc<str>; 0]),
121 };
122 }
123}
124
125#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
127#[serde(rename_all = "lowercase")]
128pub enum WorkflowStatus {
129 #[default]
130 Active,
131 Paused,
132 Archived,
133}
134
135#[derive(Clone, Debug, Deserialize)]
148#[non_exhaustive]
149pub struct Workflow {
150 pub id: String,
151 #[doc(hidden)]
156 #[serde(skip)]
157 pub id_arc: Arc<str>,
158 pub name: String,
159 #[serde(default)]
160 pub priority: u32,
161 pub description: Option<String>,
162 #[serde(default = "crate::engine::utils::default_condition")]
163 pub condition: Value,
164 #[doc(hidden)]
168 #[serde(skip)]
169 pub compiled_condition: Option<Arc<Logic>>,
170 #[doc(hidden)]
176 #[serde(skip, default)]
177 pub fully_sync: bool,
178 #[serde(deserialize_with = "crate::engine::steps::flatten")]
187 pub tasks: Vec<Task>,
188 #[serde(default)]
189 pub continue_on_error: bool,
190 #[serde(default = "default_channel")]
192 pub channel: String,
193 #[serde(default = "default_version")]
195 pub version: u32,
196 #[serde(default)]
198 pub status: WorkflowStatus,
199 #[serde(default)]
206 pub rollout: Option<Rollout>,
207 #[serde(default, rename = "loop")]
213 pub loop_config: Option<LoopConfig>,
214 #[serde(default)]
216 pub tags: Vec<String>,
217 #[serde(default)]
219 pub created_at: Option<DateTime<Utc>>,
220 #[serde(default)]
222 pub updated_at: Option<DateTime<Utc>>,
223}
224
225fn default_channel() -> String {
226 "default".to_string()
227}
228
229fn default_version() -> u32 {
230 1
231}
232
233impl Default for Workflow {
234 fn default() -> Self {
235 Self::new()
236 }
237}
238
239impl Workflow {
240 pub fn new() -> Self {
241 Self {
242 id: String::new(),
243 id_arc: Arc::from(""),
244 name: String::new(),
245 priority: 0,
246 description: None,
247 condition: Value::Bool(true),
248 compiled_condition: None,
249 fully_sync: false,
250 tasks: Vec::new(),
251 continue_on_error: false,
252 channel: default_channel(),
253 version: 1,
254 status: WorkflowStatus::Active,
255 rollout: None,
256 loop_config: None,
257 tags: Vec::new(),
258 created_at: None,
259 updated_at: None,
260 }
261 }
262
263 pub fn rule(id: &str, name: &str, condition: Value, tasks: Vec<Task>) -> Self {
274 Self {
275 id: id.to_string(),
276 id_arc: Arc::from(id),
277 name: name.to_string(),
278 priority: 0,
279 description: None,
280 condition,
281 compiled_condition: None,
282 fully_sync: false,
283 tasks,
284 continue_on_error: false,
285 channel: default_channel(),
286 version: 1,
287 status: WorkflowStatus::Active,
288 rollout: None,
289 loop_config: None,
290 tags: Vec::new(),
291 created_at: None,
292 updated_at: None,
293 }
294 }
295
296 pub fn from_json(json_str: &str) -> Result<Self> {
298 serde_json::from_str(json_str).map_err(DataflowError::from_serde)
299 }
300
301 pub fn from_file<P: AsRef<Path>>(path: P) -> Result<Self> {
303 let json_str = fs::read_to_string(path).map_err(DataflowError::from_io)?;
304
305 Self::from_json(&json_str)
306 }
307
308 pub fn validate(&self) -> Result<()> {
310 if self.id.is_empty() {
312 return Err(DataflowError::Workflow(
313 "Workflow id cannot be empty".to_string(),
314 ));
315 }
316
317 if self.name.is_empty() {
318 return Err(DataflowError::Workflow(
319 "Workflow name cannot be empty".to_string(),
320 ));
321 }
322
323 if self.tasks.is_empty() {
325 return Err(DataflowError::Workflow(
326 "Workflow must have at least one task".to_string(),
327 ));
328 }
329
330 let mut step_ids = std::collections::HashSet::new();
334 for task in &self.tasks {
335 for group in &task.group_starts {
336 if !step_ids.insert(group.id.as_str()) {
337 return Err(DataflowError::Workflow(format!(
338 "Duplicate step ID '{}' in workflow — task group IDs share the task ID namespace",
339 group.id
340 )));
341 }
342 }
343 if !step_ids.insert(task.id.as_str()) {
344 return Err(DataflowError::Workflow(format!(
345 "Duplicate task ID '{}' in workflow",
346 task.id
347 )));
348 }
349 }
350
351 if let Some(loop_config) = &self.loop_config {
355 loop_config.validate(&self.id)?;
356 }
357
358 Ok(())
359 }
360}
361
362#[derive(Debug, Clone, Copy)]
371pub struct ConnectorRef<'a> {
372 pub workflow_id: &'a str,
374 pub task_id: &'a str,
376 pub function: &'a str,
378 pub connector: ConnectorName<'a>,
381 pub config: &'a FunctionConfig,
383}
384
385impl Workflow {
386 pub fn connector_refs(&self) -> impl Iterator<Item = ConnectorRef<'_>> {
399 self.tasks.iter().filter_map(move |task| {
402 task.function.connector().map(|connector| ConnectorRef {
403 workflow_id: &self.id,
404 task_id: &task.id,
405 function: task.function.function_name(),
406 connector,
407 config: &task.function,
408 })
409 })
410 }
411}
412
413#[cfg(test)]
414mod tests {
415 use super::*;
416
417 fn wf(tasks_json: &str) -> Workflow {
418 Workflow::from_json(&format!(
419 r#"{{ "id": "w", "name": "w", "priority": 0, "condition": true,
420 "tasks": [{tasks_json}] }}"#
421 ))
422 .expect("workflow should parse")
423 }
424
425 const HTTP: &str = r#"{ "id": "call", "name": "call", "function": {
426 "name": "http_call", "input": { "connector": "user_service" } } }"#;
427 const KAFKA: &str = r#"{ "id": "pub", "name": "pub", "function": {
428 "name": "publish_kafka",
429 "input": { "connector": "events", "topic": "t" } } }"#;
430 const MAP: &str = r#"{ "id": "m", "name": "m", "function": {
431 "name": "map", "input": { "mappings": [] } } }"#;
432 const LOG: &str = r#"{ "id": "l", "name": "l", "function": {
433 "name": "log", "input": { "message": "hi" } } }"#;
434
435 #[test]
436 fn connector_refs_yields_only_connector_tasks_in_task_order() {
437 let workflow = wf(&format!("{MAP},{HTTP},{LOG},{KAFKA}"));
438 let refs: Vec<_> = workflow.connector_refs().collect();
439
440 assert_eq!(refs.len(), 2);
441 assert_eq!(refs[0].task_id, "call");
442 assert_eq!(refs[0].function, "http_call");
443 assert_eq!(refs[0].connector.as_static(), Some("user_service"));
444 assert_eq!(refs[1].task_id, "pub");
445 assert_eq!(refs[1].function, "publish_kafka");
446 assert_eq!(refs[1].connector.as_static(), Some("events"));
447 }
448
449 #[test]
450 fn connector_refs_carries_the_owning_workflow_id() {
451 let workflow = wf(HTTP);
452 assert!(workflow.connector_refs().all(|r| r.workflow_id == "w"));
453
454 let empty = Workflow::new();
456 assert_eq!(empty.id, "");
457 assert_eq!(empty.connector_refs().count(), 0);
458 }
459
460 #[test]
461 fn connector_refs_is_empty_for_no_tasks() {
462 assert_eq!(Workflow::new().connector_refs().count(), 0);
465 }
466
467 #[test]
468 fn connector_refs_does_not_deduplicate() {
469 let a = r#"{ "id": "a", "name": "a", "function": {
470 "name": "http_call", "input": { "connector": "same" } } }"#;
471 let b = r#"{ "id": "b", "name": "b", "function": {
472 "name": "enrich",
473 "input": { "connector": "same", "merge_path": "data.out" } } }"#;
474 let workflow = wf(&format!("{a},{b}"));
475
476 let refs: Vec<_> = workflow.connector_refs().collect();
477 assert_eq!(refs.len(), 2, "one item per task, not a distinct set");
478 assert!(refs.iter().all(|r| r.connector.as_static() == Some("same")));
479 }
480
481 #[test]
482 fn connector_refs_works_on_an_uncompiled_workflow() {
483 let workflow = wf(HTTP);
486 assert!(workflow.compiled_condition.is_none());
487 assert_eq!(workflow.connector_refs().count(), 1);
488 }
489
490 #[test]
491 fn connector_ref_is_copy() {
492 let workflow = wf(HTTP);
493 let r = workflow.connector_refs().next().unwrap();
494 let copied = r;
495 assert_eq!(r.connector, copied.connector);
497 assert_eq!(r.task_id, copied.task_id);
498 }
499
500 #[test]
501 fn connector_ref_config_supports_a_cross_field_rule() {
502 let custom = r#"{ "id": "db", "name": "db", "function": {
505 "name": "pg_query",
506 "input": { "connector": "pg_main", "database": "orders" } } }"#;
507 let workflow = wf(custom);
508
509 let r = workflow.connector_refs().next().expect("custom connector");
510 assert_eq!(r.connector.as_static(), Some("pg_main"));
511 match r.config {
512 FunctionConfig::Custom { input, .. } => {
513 assert_eq!(
514 input.get("database").and_then(|v| v.as_str()),
515 Some("orders")
516 );
517 }
518 other => panic!("expected Custom, got {other:?}"),
519 }
520 }
521
522 #[test]
523 fn rollout_defaults_to_none_on_every_construction_path() {
524 assert_eq!(Workflow::new().rollout, None);
525 assert_eq!(Workflow::default().rollout, None);
526 assert_eq!(
527 Workflow::rule("r", "r", Value::Bool(true), Vec::new()).rollout,
528 None
529 );
530 assert_eq!(wf(MAP).rollout, None, "absent JSON key gives None");
531 }
532
533 fn loop_wf(loop_json: &str) -> Result<Workflow> {
539 let workflow = Workflow::from_json(&format!(
540 r#"{{ "id": "w", "name": "w", "loop": {loop_json}, "tasks": [{MAP}] }}"#
541 ))?;
542 workflow.validate()?;
543 Ok(workflow)
544 }
545
546 #[test]
547 fn loop_config_defaults_init_zero_increment_one() {
548 let cfg = loop_wf(r#"{"max": 5}"#)
549 .expect("valid loop")
550 .loop_config
551 .expect("loop config present");
552 assert_eq!(cfg.init, 0);
553 assert_eq!(cfg.increment, 1);
554 assert_eq!(cfg.max, 5);
555 assert_eq!(cfg.counter, None);
556 }
557
558 #[test]
559 fn loop_config_is_absent_on_every_construction_path() {
560 assert!(wf(MAP).loop_config.is_none(), "absent JSON key gives None");
561 assert!(Workflow::new().loop_config.is_none());
562 assert!(Workflow::default().loop_config.is_none());
563 assert!(
564 Workflow::rule("r", "r", Value::Bool(true), Vec::new())
565 .loop_config
566 .is_none()
567 );
568 }
569
570 #[test]
571 fn loop_config_rejects_a_bound_that_could_never_run_a_sweep() {
572 assert!(loop_wf(r#"{"max": 0}"#).is_err());
575 assert!(loop_wf(r#"{"init": 5, "max": 5}"#).is_err());
576 assert!(loop_wf(r#"{"init": 5, "max": 2}"#).is_err());
577 }
578
579 #[test]
580 fn loop_config_rejects_a_non_advancing_increment() {
581 assert!(loop_wf(r#"{"max": 5, "increment": 0}"#).is_err());
582 assert!(loop_wf(r#"{"max": 5, "increment": -1}"#).is_err());
583 }
584
585 #[test]
586 fn loop_config_rejects_an_empty_counter_path() {
587 assert!(loop_wf(r#"{"max": 5, "counter": ""}"#).is_err());
588 assert!(loop_wf(r#"{"max": 5, "counter": "a..b"}"#).is_err());
589 assert!(loop_wf(r#"{"max": 5, "counter": "a."}"#).is_err());
590 }
591
592 #[test]
593 fn loop_config_requires_max() {
594 assert!(
597 Workflow::from_json(r#"{ "id": "w", "name": "w", "loop": {}, "tasks": [] }"#).is_err()
598 );
599 }
600
601 #[test]
602 fn loop_config_deserializes_every_combination_of_optional_fields() {
603 for (json, counter, init, increment) in [
607 (r#"{"max": 9}"#, None, 0, 1),
608 (r#"{"max": 9, "counter": "i"}"#, Some("i"), 0, 1),
609 (r#"{"max": 9, "init": 4}"#, None, 4, 1),
610 (r#"{"max": 9, "increment": 3}"#, None, 0, 3),
611 (r#"{"max": 9, "counter": "i", "init": 4}"#, Some("i"), 4, 1),
612 (
613 r#"{"max": 9, "counter": "i", "increment": 3}"#,
614 Some("i"),
615 0,
616 3,
617 ),
618 (r#"{"max": 9, "init": 4, "increment": 3}"#, None, 4, 3),
619 (
620 r#"{"max": 9, "counter": "i", "init": 4, "increment": 3}"#,
621 Some("i"),
622 4,
623 3,
624 ),
625 ] {
626 let cfg = loop_wf(json)
627 .unwrap_or_else(|e| panic!("{json} should be valid: {e}"))
628 .loop_config
629 .expect("loop config present");
630 assert_eq!(cfg.counter.as_deref(), counter, "counter for {json}");
631 assert_eq!(cfg.init, init, "init for {json}");
632 assert_eq!(cfg.increment, increment, "increment for {json}");
633 assert_eq!(cfg.max, 9, "max for {json}");
634 }
635 }
636
637 #[test]
638 fn loop_config_validation_matrix_over_init_increment_and_max() {
639 for (init, increment, max, valid) in [
643 (0_i64, 1_i64, 1_i64, true),
645 (0, 1, 100, true),
646 (0, 7, 3, true), (5, 1, 6, true),
648 (-5, 1, 0, true),
650 (-5, 2, -4, true),
651 (-1, 1, 1, true),
652 (0, 1, 0, false),
654 (5, 1, 5, false),
655 (5, 1, 4, false),
656 (0, 1, -1, false),
657 (-5, 1, -5, false),
658 (0, 0, 10, false),
660 (0, -1, 10, false),
661 (0, -100, 10, false),
662 ] {
663 let json = format!(r#"{{"init": {init}, "increment": {increment}, "max": {max}}}"#);
664 assert_eq!(
665 loop_wf(&json).is_ok(),
666 valid,
667 "init={init} increment={increment} max={max} should be {}",
668 if valid { "accepted" } else { "rejected" }
669 );
670 }
671 }
672
673 #[test]
674 fn loop_config_counter_path_matrix() {
675 for (counter, valid) in [
678 ("i", true),
679 ("index", true),
680 ("cursor.index", true),
681 ("a.b.c.d", true),
682 ("#7", true), ("", false),
684 (".", false),
685 ("a.", false),
686 (".a", false),
687 ("a..b", false),
688 ] {
689 let json = format!(r#"{{"max": 5, "counter": "{counter}"}}"#);
690 assert_eq!(
691 loop_wf(&json).is_ok(),
692 valid,
693 "counter {counter:?} should be {}",
694 if valid { "accepted" } else { "rejected" }
695 );
696 }
697 }
698
699 #[test]
700 fn precompute_counter_path_matrix() {
701 for (counter, expected) in [
702 ("i", vec!["temp_data", "i"]),
703 ("cursor.index", vec!["temp_data", "cursor", "index"]),
704 ("a.b.c", vec!["temp_data", "a", "b", "c"]),
705 ("#7", vec!["temp_data", "#7"]),
708 ] {
709 let mut cfg = loop_wf(&format!(r#"{{"max": 5, "counter": "{counter}"}}"#))
710 .expect("valid loop")
711 .loop_config
712 .expect("loop config present");
713 cfg.precompute_counter_path();
714 let parts: Vec<&str> = cfg.counter_parts.iter().map(Arc::as_ref).collect();
715 assert_eq!(parts, expected, "for counter {counter:?}");
716 }
717 }
718
719 #[test]
720 fn precompute_counter_path_is_idempotent() {
721 let mut cfg = loop_wf(r#"{"max": 5, "counter": "cursor.index"}"#)
724 .expect("valid loop")
725 .loop_config
726 .expect("loop config present");
727 cfg.precompute_counter_path();
728 let first: Vec<Arc<str>> = cfg.counter_parts.to_vec();
729 cfg.precompute_counter_path();
730 assert_eq!(cfg.counter_parts.to_vec(), first);
731 }
732
733 #[test]
734 fn loop_config_rejects_a_non_object_and_a_non_numeric_max() {
735 for json in [r#""five""#, "5", "[]", r#"{"max": "5"}"#, "true"] {
736 assert!(loop_wf(json).is_err(), "{json} is not a valid loop config");
737 }
738 }
739
740 #[test]
741 fn an_explicit_null_loop_means_no_loop() {
742 let workflow = loop_wf("null").expect("explicit null should be accepted");
746 assert!(workflow.loop_config.is_none());
747 }
748
749 #[test]
750 fn a_workflow_with_a_loop_still_validates_its_other_rules() {
751 let duplicate_tasks = Workflow::from_json(
754 r#"{ "id": "w", "name": "w", "loop": {"max": 5}, "tasks": [
755 {"id": "t", "name": "t", "function": {"name": "map", "input": {"mappings": []}}},
756 {"id": "t", "name": "t", "function": {"name": "map", "input": {"mappings": []}}}] }"#,
757 )
758 .expect("should parse");
759 assert!(duplicate_tasks.validate().is_err(), "duplicate task ids");
760
761 let no_tasks =
762 Workflow::from_json(r#"{ "id": "w", "name": "w", "loop": {"max": 5}, "tasks": [] }"#)
763 .expect("should parse");
764 assert!(no_tasks.validate().is_err(), "empty task list");
765 }
766
767 #[test]
768 fn loop_config_coexists_with_every_other_workflow_field() {
769 let workflow = Workflow::from_json(&format!(
772 r#"{{ "id": "w", "name": "w", "priority": 7, "description": "d",
773 "condition": {{"==": [1, 1]}},
774 "loop": {{"counter": "i", "max": 5}},
775 "continue_on_error": true, "channel": "c", "version": 3,
776 "status": "paused",
777 "rollout": {{"bucket_start": 0, "bucket_end": 50}},
778 "tags": ["x"], "tasks": [{MAP}] }}"#
779 ))
780 .expect("should parse");
781 workflow.validate().expect("should validate");
782
783 assert_eq!(workflow.priority, 7);
784 assert_eq!(workflow.channel, "c");
785 assert_eq!(workflow.version, 3);
786 assert_eq!(workflow.status, WorkflowStatus::Paused);
787 assert!(workflow.continue_on_error);
788 assert_eq!(
789 workflow.rollout,
790 Some(Rollout {
791 bucket_start: 0,
792 bucket_end: 50
793 })
794 );
795 assert_eq!(workflow.tags, ["x"]);
796 assert_eq!(
797 workflow
798 .loop_config
799 .expect("loop present")
800 .counter
801 .as_deref(),
802 Some("i")
803 );
804 }
805
806 #[test]
807 fn loop_config_accepts_a_valid_counter() {
808 let cfg = loop_wf(r#"{"max": 5, "counter": "cursor.index"}"#)
809 .expect("valid loop")
810 .loop_config
811 .expect("loop config present");
812 assert_eq!(cfg.counter.as_deref(), Some("cursor.index"));
813 }
814
815 #[test]
816 fn precompute_counter_path_prefixes_temp_data() {
817 let mut cfg = loop_wf(r#"{"max": 5, "counter": "cursor.index"}"#)
818 .expect("valid loop")
819 .loop_config
820 .expect("loop config present");
821 assert!(
822 cfg.counter_parts.is_empty(),
823 "uncompiled workflows start with no pre-split path"
824 );
825
826 cfg.precompute_counter_path();
827
828 let parts: Vec<&str> = cfg.counter_parts.iter().map(Arc::as_ref).collect();
829 assert_eq!(parts, ["temp_data", "cursor", "index"]);
830 }
831
832 #[test]
833 fn precompute_counter_path_is_empty_without_a_counter_name() {
834 let mut cfg = loop_wf(r#"{"max": 5}"#)
835 .expect("valid loop")
836 .loop_config
837 .expect("loop config present");
838 cfg.precompute_counter_path();
839 assert!(cfg.counter_parts.is_empty());
840 }
841
842 #[test]
843 fn rollout_deserializes_from_json() {
844 let workflow = Workflow::from_json(
845 r#"{ "id": "w", "name": "w", "condition": true,
846 "rollout": { "bucket_start": 0, "bucket_end": 50 },
847 "tasks": [] }"#,
848 )
849 .unwrap();
850 assert_eq!(
851 workflow.rollout,
852 Some(Rollout {
853 bucket_start: 0,
854 bucket_end: 50
855 })
856 );
857 }
858}