use jaq_core::load::{self, parse::BinaryOp, parse::Term};
use jaq_core::path::Part;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Plan {
Slurp,
PerFileStream,
Decompose { reduce: String },
}
impl Plan {
pub fn describe(&self) -> String {
match self {
Plan::Slurp => "slurp (whole array; not decomposable)".to_string(),
Plan::PerFileStream => {
"stream per file (element-wise; bounded to one file)".to_string()
}
Plan::Decompose { reduce } => {
format!("decompose per file, then reduce with `{reduce}` (bounded)")
}
}
}
}
pub fn analyze(code: &str) -> Plan {
let Some(term) = load::parse(code, |p| p.term()) else {
return Plan::Slurp;
};
let Some(segs) = pipe_segments(&term) else {
return Plan::Slurp;
};
let plan = classify(code, &segs);
if let Plan::Decompose { reduce } = &plan
&& !reparses(reduce)
{
return Plan::Slurp;
}
plan
}
fn reparses(code: &str) -> bool {
load::parse(code, |p| p.term()).is_some()
}
fn pipe_segments<'a, 's>(t: &'a Term<&'s str>) -> Option<Vec<&'a Term<&'s str>>> {
let mut segs = Vec::new();
let mut cur = t;
loop {
match cur {
Term::BinOp(l, BinaryOp::Pipe(None), r) => {
segs.push(l.as_ref());
cur = r.as_ref();
}
Term::BinOp(_, BinaryOp::Pipe(Some(_)), _) => return None,
other => {
segs.push(other);
break;
}
}
}
Some(segs)
}
fn classify(code: &str, segs: &[&Term<&str>]) -> Plan {
let last = *segs.last().expect("at least one segment");
if is_iterate(segs[0]) {
return Plan::PerFileStream;
}
if segs.len() == 1 {
if is_map(last) || is_collected_iteration(last) {
return Plan::Decompose {
reduce: "add".to_string(),
};
}
if let Some(reduce) = scalar_reduce(last) {
return Plan::Decompose {
reduce: reduce.to_string(),
};
}
return Plan::Slurp;
}
if is_slice_upto(last) {
let before = &segs[..segs.len() - 1];
if let Some(sort_seg) = before.last().copied().filter(|s| is_sort_by(s))
&& prefix_ok(&before[..before.len() - 1])
{
let tail = tail_source(code, sort_seg);
return Plan::Decompose {
reduce: format!("add | {tail}"),
};
}
return Plan::Slurp;
}
if let Some(reduce) = scalar_reduce(last)
&& prefix_ok(&segs[..segs.len() - 1])
{
return Plan::Decompose {
reduce: reduce.to_string(),
};
}
Plan::Slurp
}
fn tail_source<'a>(code: &'a str, call: &Term<&str>) -> &'a str {
match call {
Term::Call(name, _) => &code[load::span(code, name).start..],
_ => code,
}
}
fn prefix_ok(prefix: &[&Term<&str>]) -> bool {
match prefix {
[] => true,
[seg] => is_map(seg) || is_collected_iteration(seg),
_ => false,
}
}
fn is_iterate(t: &Term<&str>) -> bool {
matches!(t, Term::Path(inner, path)
if matches!(inner.as_ref(), Term::Id)
&& path.0.len() == 1
&& matches!(path.0[0].0, Part::Range(None, None)))
}
fn is_slice_upto(t: &Term<&str>) -> bool {
if let Term::Path(inner, path) = t
&& matches!(inner.as_ref(), Term::Id)
&& path.0.len() == 1
&& let Part::Range(None, Some(bound)) = &path.0[0].0
{
return matches!(bound, Term::Num(s) if s.parse::<u64>().is_ok());
}
false
}
fn is_map(t: &Term<&str>) -> bool {
matches!(t, Term::Call(name, args) if *name == "map" && args.len() == 1)
}
fn is_sort_by(t: &Term<&str>) -> bool {
matches!(t, Term::Call(name, args) if *name == "sort_by" && args.len() == 1)
}
fn is_collected_iteration(t: &Term<&str>) -> bool {
match t {
Term::Arr(Some(inner)) => match pipe_segments(inner) {
Some(segs) => is_iterate(segs[0]),
None => false,
},
_ => false,
}
}
fn scalar_reduce(t: &Term<&str>) -> Option<&'static str> {
match t {
Term::Call(name, args) if args.is_empty() && *name == "length" => Some("add"),
_ => None,
}
}
#[cfg(test)]
mod tests {
use super::*;
fn decompose(reduce: &str) -> Plan {
Plan::Decompose {
reduce: reduce.to_string(),
}
}
#[test]
fn stream_for_iteration() {
assert_eq!(analyze(".[]"), Plan::PerFileStream);
assert_eq!(analyze(".[] | select(.dead_end)"), Plan::PerFileStream);
assert_eq!(
analyze(r#".[] | select(.step.actor | startswith("agent:")) | .step.id"#),
Plan::PerFileStream
);
}
#[test]
fn map_decomposes_to_add() {
assert_eq!(analyze("map(select(.dead_end))"), decompose("add"));
assert_eq!(analyze("[.[] | select(.dead_end)]"), decompose("add"));
}
#[test]
fn top_n_decomposes_with_sort_tail() {
assert_eq!(
analyze("sort_by(-.tokens) | .[:10]"),
decompose("add | sort_by(-.tokens) | .[:10]")
);
assert_eq!(
analyze("map({t: .step.id}) | sort_by(-.t) | .[:5]"),
decompose("add | sort_by(-.t) | .[:5]")
);
}
#[test]
fn scalar_reductions() {
assert_eq!(analyze("length"), decompose("add"));
assert_eq!(analyze("map(select(.dead_end)) | length"), decompose("add"));
}
#[test]
fn scalar_add_slurps_because_float_sums_reassociate() {
assert_eq!(analyze("add"), Plan::Slurp);
assert_eq!(analyze("map(.tokens) | add"), Plan::Slurp);
assert_eq!(analyze("[.[].tokens] | add"), Plan::Slurp);
}
#[test]
fn non_distributive_prefix_slurps() {
assert_eq!(analyze("unique | sort_by(.x) | .[:10]"), Plan::Slurp);
assert_eq!(
analyze("group_by(.path.meta.source) | map({n: length})"),
Plan::Slurp
);
}
#[test]
fn bare_slice_without_sort_slurps() {
assert_eq!(analyze("map(.step) | .[:10]"), Plan::Slurp);
}
#[test]
fn as_binding_and_reduce_slurp() {
assert_eq!(analyze(".tokens as $t | $t"), Plan::Slurp);
assert_eq!(analyze("reduce .[] as $x (0; . + $x.n)"), Plan::Slurp);
}
#[test]
fn bare_identity_slurps() {
assert_eq!(analyze("."), Plan::Slurp);
}
}