pub mod dataflow;
pub mod keys;
pub mod logic;
pub mod operators;
use std::collections::BTreeMap;
use serde_json::Value;
use crate::config::AppConfig;
use crate::definitions::json::Document;
use crate::definitions::{DefinitionSet, Entity, SharedDefinitions};
pub use dataflow::Reads;
pub use logic::{Evaluator, Expr};
pub fn selection_context() -> Value {
serde_json::json!({ "data": {}, "temp_data": {}, "metadata": {} })
}
pub struct Analysis<'a> {
pub source: &'a DefinitionSet,
pub compiled: &'a DefinitionSet,
pub shared: &'a SharedDefinitions,
pub config: Option<&'a AppConfig>,
pub evaluator: Evaluator,
pub workflows: Vec<WorkflowFacts>,
pub channels: BTreeMap<String, String>,
documents: BTreeMap<String, Document>,
}
pub struct WorkflowFacts {
pub origin: String,
pub name: String,
pub workflow_id: Option<String>,
pub doc: Value,
pub condition: Expr,
pub has_loop: bool,
pub loop_counter: Option<String>,
pub steps: Vec<StepFacts>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StepKind {
Task,
Group,
}
pub struct StepFacts {
pub path: String,
pub list: String,
pub index: usize,
pub parent: Option<usize>,
pub kind: StepKind,
pub id: String,
pub node: Value,
pub condition: Option<Expr>,
pub terminal: bool,
pub certain: bool,
pub function: Option<String>,
pub expressions: Vec<(String, Expr)>,
pub writes: Vec<String>,
}
impl StepFacts {
pub fn reads(&self) -> Reads {
let mut out = Reads::default();
for expr in self
.condition
.iter()
.chain(self.expressions.iter().map(|(_, e)| e))
{
out.paths.extend(expr.reads.paths.iter().cloned());
out.computed |= expr.reads.computed;
out.scoped |= expr.reads.scoped;
}
out
}
pub fn is_unconditional(&self) -> bool {
self.condition.as_ref().is_none_or(Expr::is_constant_true)
}
}
impl<'a> Analysis<'a> {
pub fn new(
source: &'a DefinitionSet,
compiled: &'a DefinitionSet,
shared: &'a SharedDefinitions,
config: Option<&'a AppConfig>,
) -> Self {
let evaluator = Evaluator::new();
let workflows = compiled
.iter(Entity::Workflow)
.map(|def| workflow_facts(&def.origin, &def.doc, &evaluator))
.collect();
let channels = compiled
.iter(Entity::Channel)
.filter_map(|def| {
Some((
def.doc.get("name")?.as_str()?.to_string(),
def.doc.get("workflow_id")?.as_str()?.to_string(),
))
})
.collect();
let documents = source
.definitions
.iter()
.filter(|def| crate::definitions::compile::residue(&def.doc, "").is_empty())
.filter_map(|def| Some((def.origin.clone(), def.spans.clone()?)))
.collect();
Self {
source,
compiled,
shared,
config,
evaluator,
workflows,
channels,
documents,
}
}
pub fn locate(&self, origin: &str, path: &str) -> Option<(usize, usize)> {
let doc = self.documents.get(origin)?;
let span = doc.locate(path)?;
Some(doc.line_col(span.start))
}
pub fn workflow_for_channel(&self, channel: &str) -> Option<&WorkflowFacts> {
let id = self.channels.get(channel)?;
self.workflows
.iter()
.find(|w| w.workflow_id.as_deref() == Some(id))
}
}
fn workflow_facts(origin: &str, doc: &Value, evaluator: &Evaluator) -> WorkflowFacts {
let condition = evaluator.expr(doc.get("condition").unwrap_or(&Value::Bool(true)));
let loop_config = doc.get("loop").filter(|l| !l.is_null());
let loop_counter = loop_config
.and_then(|l| l.get("counter"))
.and_then(Value::as_str)
.map(|c| format!("temp_data.{c}"));
let mut steps = Vec::new();
if let Some(tasks) = doc.get("tasks") {
walk(tasks, "tasks", None, true, evaluator, &mut steps);
}
WorkflowFacts {
origin: origin.to_string(),
name: doc
.get("name")
.and_then(Value::as_str)
.unwrap_or("")
.to_string(),
workflow_id: doc
.get("workflow_id")
.and_then(Value::as_str)
.map(str::to_string),
doc: doc.clone(),
condition,
has_loop: loop_config.is_some(),
loop_counter,
steps,
}
}
fn walk(
tasks: &Value,
list: &str,
parent: Option<usize>,
parent_certain: bool,
evaluator: &Evaluator,
out: &mut Vec<StepFacts>,
) {
let Some(items) = tasks.as_array() else {
return;
};
for (index, node) in items.iter().enumerate() {
let path = format!("{list}[{index}]");
let condition = node.get("condition").map(|c| evaluator.expr(c));
let certain = parent_certain && condition.as_ref().is_none_or(Expr::is_constant_true);
let terminal = node
.get("terminal")
.and_then(Value::as_bool)
.unwrap_or(false);
let id = node
.get("id")
.and_then(Value::as_str)
.unwrap_or("")
.to_string();
let kind = if crate::engine::is_group(node) {
StepKind::Group
} else {
StepKind::Task
};
let function = node
.get("function")
.and_then(|f| f.get("name"))
.and_then(Value::as_str)
.map(str::to_string);
let expressions = match (&function, node.get("function").and_then(|f| f.get("input"))) {
(Some(name), Some(input)) => operators::input_expressions(name, input)
.into_iter()
.map(|(p, e)| (p, evaluator.expr(e)))
.collect(),
_ => Vec::new(),
};
let writes = match kind {
StepKind::Task => dataflow::task_writes(node),
StepKind::Group => Vec::new(), };
let me = out.len();
out.push(StepFacts {
path: path.clone(),
list: list.to_string(),
index,
parent,
kind,
id,
node: node.clone(),
condition,
terminal,
certain,
function,
expressions,
writes,
});
if kind == StepKind::Group
&& let Some(members) = node.get("tasks")
{
walk(
members,
&format!("{path}.tasks"),
Some(me),
certain,
evaluator,
out,
);
let member_writes: Vec<String> = out[me + 1..]
.iter()
.flat_map(|s| s.writes.iter().cloned())
.collect();
out[me].writes = member_writes;
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn analysis_of(doc: Value) -> Vec<StepFacts> {
workflow_facts("wf.json", &doc, &Evaluator::new()).steps
}
#[test]
fn steps_are_flattened_in_document_order_with_certainty() {
let steps = analysis_of(json!({
"name": "w",
"tasks": [
{"id": "a", "name": "A", "function": {"name": "parse_json", "input": {"source": "payload", "target": "req"}}},
{"id": "g", "condition": {"var": "data.flag"}, "terminal": true, "tasks": [
{"id": "b", "name": "B", "function": {"name": "map", "input": {"mappings": [{"path": "data.x", "logic": {"var": "data.req.x"}}]}}},
{"id": "c", "name": "C", "condition": true, "function": {"name": "log", "input": {"message": "hi"}}}
]},
{"id": "d", "name": "D", "condition": {"==": [1, 1]}, "function": {"name": "log", "input": {"message": "hi"}}}
]
}));
let paths: Vec<&str> = steps.iter().map(|s| s.path.as_str()).collect();
assert_eq!(
paths,
[
"tasks[0]",
"tasks[1]",
"tasks[1].tasks[0]",
"tasks[1].tasks[1]",
"tasks[2]"
]
);
assert!(steps[0].certain);
assert!(!steps[1].certain, "a conditional group");
assert!(!steps[2].certain, "inside a conditional group");
assert!(steps[2].is_unconditional(), "but unconditional itself");
assert!(
steps[4].certain,
"a condition folded to true is no condition"
);
assert_eq!(steps[1].kind, StepKind::Group);
assert!(steps[1].terminal);
assert_eq!(
steps[1].writes,
["data.x"],
"a group writes what its members write"
);
assert_eq!(steps[2].parent, Some(1));
assert_eq!(steps[0].writes, ["data.req"]);
assert_eq!(steps[2].reads().paths, ["data.req.x"]);
assert_eq!(steps[1].reads().paths, ["data.flag"]);
}
#[test]
fn a_loop_counter_is_named_as_a_temp_data_path() {
let facts = workflow_facts(
"wf.json",
&json!({"name": "w", "loop": {"counter": "i", "max": 3}, "tasks": []}),
&Evaluator::new(),
);
assert!(facts.has_loop);
assert_eq!(facts.loop_counter.as_deref(), Some("temp_data.i"));
assert!(
facts.condition.is_constant_true(),
"absent condition is true"
);
}
#[test]
fn every_operator_is_classified() {
for op in crate::engine::operators::operator_names() {
let scoping = operators::SCOPING.contains(&op.as_str());
let plain = operators::NON_SCOPING.contains(&op.as_str());
assert!(
scoping ^ plain,
"operator `{op}` must be in exactly one of SCOPING / NON_SCOPING"
);
}
for op in operators::SCOPING.iter().chain(operators::NON_SCOPING) {
assert!(
crate::engine::operators::is_operator(op),
"`{op}` is classified but is not an operator this build registers"
);
}
}
#[test]
fn scoping_operators_rebind_var_to_the_element() {
let ev = Evaluator::new();
let ctx = json!({"data": {"items": [{"payload": 1}, {"payload": 2}]}, "payload": 99});
assert_eq!(
ev.evaluate(
&json!({"map": [{"var": "data.items"}, {"var": "payload"}]}),
&ctx
),
Some(json!([1, 2]))
);
assert_eq!(
ev.evaluate(
&json!({"filter": [{"var": "data.items"}, {"==": [{"var": "payload"}, 2]}]}),
&ctx
),
Some(json!([{"payload": 2}]))
);
assert_eq!(
ev.evaluate(
&json!({"some": [{"var": "data.items"}, {"==": [{"var": "payload"}, 2]}]}),
&ctx
),
Some(json!(true))
);
}
}