use crate::engine::functions::schema::WriteShape;
use serde_json::Value;
use super::operators;
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct Reads {
pub paths: Vec<String>,
pub computed: bool,
pub scoped: bool,
}
impl Reads {
pub fn uncertain(&self) -> bool {
self.computed || self.scoped
}
pub fn touches(&self, path: &str) -> bool {
self.paths.iter().any(|read| overlaps(read, path))
}
}
pub fn reads(value: &Value) -> Reads {
let mut out = Reads::default();
collect_reads(value, false, &mut out);
out
}
fn collect_reads(value: &Value, scoped: bool, out: &mut Reads) {
match value {
Value::Array(items) => items.iter().for_each(|v| collect_reads(v, scoped, out)),
Value::Object(map) => {
if map.len() == 1 {
let (key, arg) = map.iter().next().expect("one member");
match key.as_str() {
"var" => {
read_var(arg, scoped, out);
if let Some(rest) = arg.as_array().and_then(|a| a.get(1..)) {
rest.iter().for_each(|v| collect_reads(v, scoped, out));
}
return;
}
"val" => {
read_val(arg, scoped, out);
return;
}
op if operators::is_scoping(op) => {
match arg {
Value::Array(args) => {
if let Some(first) = args.first() {
collect_reads(first, scoped, out);
}
for later in args.iter().skip(1) {
collect_reads(later, true, out);
}
}
other => collect_reads(other, scoped, out),
}
return;
}
_ => {}
}
}
map.values().for_each(|v| collect_reads(v, scoped, out));
}
_ => {}
}
}
fn read_var(arg: &Value, scoped: bool, out: &mut Reads) {
let path = match arg {
Value::String(s) => Some(s.as_str()),
Value::Array(items) => match items.first() {
Some(Value::String(s)) => Some(s.as_str()),
None => Some(""),
Some(_) => None,
},
Value::Null => Some(""),
_ => None,
};
match (path, scoped) {
(Some(_), true) => out.scoped = true,
(Some(p), false) => out.paths.push(p.to_string()),
(None, _) => out.computed = true,
}
}
fn read_val(arg: &Value, scoped: bool, out: &mut Reads) {
let segments: Option<Vec<String>> = match arg {
Value::String(s) => Some(vec![s.clone()]),
Value::Array(items) => items
.iter()
.map(|seg| match seg {
Value::String(s) => Some(s.clone()),
Value::Number(n) => Some(n.to_string()),
_ => None,
})
.collect(),
_ => None,
};
match (segments, scoped) {
(Some(_), true) => out.scoped = true,
(Some(segs), false) => out.paths.push(segs.join(".")),
(None, _) => out.computed = true,
}
}
pub fn data_reads(value: &Value) -> Vec<String> {
reads(value)
.paths
.into_iter()
.filter(|p| p == "data" || p.starts_with("data."))
.collect()
}
pub fn task_writes(task: &Value) -> Vec<String> {
let Some(function) = task.get("function") else {
return Vec::new();
};
let name = function.get("name").and_then(Value::as_str).unwrap_or("");
let Some(input) = function.get("input") else {
return Vec::new();
};
let Some(shape) = crate::engine::functions::schema::write_shape(name) else {
return Vec::new();
};
let mut out = Vec::new();
match shape {
WriteShape::Target => {
if let Some(target) = input.get("target").and_then(Value::as_str) {
out.push(format!("data.{target}"));
}
}
WriteShape::Mappings => {
if let Some(mappings) = input.get("mappings").and_then(Value::as_array) {
for mapping in mappings {
if let Some(path) = mapping.get("path").and_then(Value::as_str) {
out.push(path.to_string());
}
}
}
}
WriteShape::OutputPath { default_root } => {
match input
.get("output")
.or_else(|| input.get("response_path"))
.and_then(Value::as_str)
{
Some(path) => out.push(path.to_string()),
None => out.extend(default_root.map(str::to_string)),
}
}
WriteShape::Nothing => {}
}
out
}
pub fn is_written(path: &str, written: &[String]) -> bool {
written.iter().any(|w| w == "data" || overlaps(w, path))
}
pub fn overlaps(a: &str, b: &str) -> bool {
a == b
|| a.is_empty()
|| b.is_empty()
|| a.starts_with(&format!("{b}."))
|| b.starts_with(&format!("{a}."))
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn literal_reads_are_collected_with_their_full_paths() {
let r = reads(
&json!({"and": [{">": [{"var": "data.a"}, 1]}, {"val": ["metadata", "vars", "x"]}]}),
);
assert_eq!(r.paths, ["data.a", "metadata.vars.x"]);
assert!(!r.uncertain());
}
#[test]
fn a_default_argument_is_read_too() {
let r = reads(&json!({"var": ["data.nick", {"var": "data.name"}]}));
assert_eq!(r.paths, ["data.nick", "data.name"]);
}
#[test]
fn a_read_under_a_scoped_argument_is_not_a_context_read() {
let r = reads(&json!({"some": [{"var": "data.items"}, {">": [{"var": "qty"}, 0]}]}));
assert_eq!(
r.paths,
["data.items"],
"only the array is read from the context"
);
assert!(r.scoped && r.uncertain());
}
#[test]
fn a_computed_val_is_uncertain() {
let r = reads(&json!({"val": ["data", "items", {"val": ["temp_data", "i"]}]}));
assert!(r.computed && r.uncertain());
assert!(r.paths.is_empty());
}
#[test]
fn writes_by_function() {
let parse = json!({"function": {"name": "parse_json", "input": {"source": "payload", "target": "order"}}});
assert_eq!(task_writes(&parse), ["data.order"]);
let map = json!({"function": {"name": "map", "input": {"mappings": [{"path": "data.a", "logic": 1}, {"path": "temp_data.b", "logic": 2}]}}});
assert_eq!(task_writes(&map), ["data.a", "temp_data.b"]);
let call = json!({"function": {"name": "http_call", "input": {"connector": "c", "output": "data.resp"}}});
assert_eq!(task_writes(&call), ["data.resp"]);
let query = json!({"function": {"name": "data_query", "input": {"connector": "c"}}});
assert_eq!(task_writes(&query), ["data"]);
}
#[test]
fn overlap_is_prefix_in_either_direction() {
assert!(overlaps("data.order", "data.order.total"));
assert!(overlaps("data.order.total", "data.order"));
assert!(!overlaps("data.order", "data.orders"));
assert!(is_written("data.x.y", &["data".to_string()]));
}
}