Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 24 additions & 16 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

20 changes: 10 additions & 10 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -8,18 +8,18 @@ license = "MIT OR Apache-2.0"
debug = true # Keeps debug symbols without sacrificing release optimizations

[build-dependencies]
mangle-ast = { git = "https://codeberg.org/ajwdev/mangle-rs.git", branch = "indexed-store" }
mangle-driver = { git = "https://codeberg.org/ajwdev/mangle-rs.git", branch = "indexed-store" }
mangle-ast = "0.9.1"
mangle-driver = "0.9.1"
glob = "0.3"

[dependencies]
mangle-ast = { git = "https://codeberg.org/ajwdev/mangle-rs.git", branch = "indexed-store" }
mangle-parse = { git = "https://codeberg.org/ajwdev/mangle-rs.git", branch = "indexed-store" }
mangle-common = { git = "https://codeberg.org/ajwdev/mangle-rs.git", branch = "indexed-store" }
mangle-driver = { git = "https://codeberg.org/ajwdev/mangle-rs.git", branch = "indexed-store" }
mangle-interpreter = { git = "https://codeberg.org/ajwdev/mangle-rs.git", branch = "indexed-store" }
mangle-analysis = { git = "https://codeberg.org/ajwdev/mangle-rs.git", branch = "indexed-store" }
mangle-ir = { git = "https://codeberg.org/ajwdev/mangle-rs.git", branch = "indexed-store" }
mangle-ast = "0.9.1"
mangle-parse = "0.9.1"
mangle-common = "0.9.1"
mangle-driver = "0.9.1"
mangle-interpreter = "0.9.1"
mangle-analysis = "0.9.1"
mangle-ir = "0.9.1"

serde = { version = "1", features = ["derive", "rc"] }
serde_json = { version = "1", features = ["preserve_order"] }
Expand All @@ -42,7 +42,7 @@ z3 = "0.12"
[dev-dependencies]
criterion = { version = "0.5", features = ["html_reports"] }
# Used by integration tests (tests/) to construct `Value` fixtures.
mangle-common = { git = "https://codeberg.org/ajwdev/mangle-rs.git", branch = "indexed-store" }
mangle-common = "0.9.1"

[[bench]]
name = "bulk_load"
Expand Down
100 changes: 80 additions & 20 deletions src/dd/build.rs
Original file line number Diff line number Diff line change
Expand Up @@ -50,20 +50,23 @@ fn eval_cmp(op: CmpOp, left: &Val, right: &Val) -> bool {
// Expr evaluation
// ---------------------------------------------------------------------------

fn eval_expr(expr: &OwnedExpr, row: &Row) -> Val {
/// Evaluate a `let` expression against one row.
///
/// Calls the interpreter's own `eval_function`, so errors (e.g. `fn:plus` on
/// a string) carry the same messages and fail the same evaluations as upstream
/// mangle. See `Step::Let` for how they are surfaced.
fn eval_expr(expr: &OwnedExpr, row: &Row) -> Result<Val> {
match expr {
OwnedExpr::Value(slot) => slot_val(slot, row),
OwnedExpr::Value(slot) => Ok(slot_val(slot, row)),
OwnedExpr::Call { func, args } => {
let vals: Vec<Value> = args.iter().map(|s| slot_val(s, row).into()).collect();
let result =
eval_function(func, &vals).unwrap_or_else(|e| panic!("Let fn:{func} failed: {e}"));
Val::from(&result)
Ok(Val::from(&eval_function(func, &vals)?))
}
}
}

// ---------------------------------------------------------------------------
// String builtin filters (Condition::Call)
// Built-in predicate filters (Condition::Call)
// ---------------------------------------------------------------------------

/// The set of `CallFilter` builtins the DD backend supports.
Expand All @@ -76,40 +79,72 @@ const SUPPORTED_CALL_FILTERS: &[&str] = &[
":string:ends_with",
":string:contains",
":match_prefix",
// Check modes, only reached via negation (`!:list:member`, `!:match_field`).
// The positive forms bind variables and lower to IterateList / MatchField.
":list:member",
":match_field",
];

/// True if `func` is a `CallFilter` builtin the DD backend can evaluate.
pub(crate) fn is_supported_call_filter(func: &str) -> bool {
SUPPORTED_CALL_FILTERS.contains(&func)
}

/// Evaluate a built-in predicate against one row.
///
/// Type errors follow upstream mangle's interpreter (`eval_builtin_predicate`):
/// `:string:*` and `:match_prefix` return `Err` on wrong argument types, while
/// the `:list:member` / `:match_field` check modes return `false`. The error
/// messages match the interpreter's too. This is deliberate parity, not a
/// considered semantics; if upstream (or we) decide a type mismatch should
/// just be `false`, change it here and drop the error collection in
/// `Step::CallFilter`.
fn eval_call_filter(func: &str, args: &[Slot], row: &Row) -> Result<bool> {
let vals: Vec<Val> = args.iter().map(|s| slot_val(s, row)).collect();
match func {
":string:starts_with" => {
let (Val::String(s), Val::String(prefix)) = (&vals[0], &vals[1]) else {
return Ok(false);
bail!(":string:starts_with: expected string arguments");
};
Ok(s.starts_with(&**prefix))
}
":string:ends_with" => {
let (Val::String(s), Val::String(suffix)) = (&vals[0], &vals[1]) else {
return Ok(false);
bail!(":string:ends_with: expected string arguments");
};
Ok(s.ends_with(&**suffix))
}
":string:contains" => {
let (Val::String(s), Val::String(needle)) = (&vals[0], &vals[1]) else {
return Ok(false);
bail!(":string:contains: expected string arguments");
};
Ok(s.contains(&**needle))
}
":match_prefix" => {
let (Val::Name(name), Val::Name(prefix)) = (&vals[0], &vals[1]) else {
return Ok(false);
bail!(":match_prefix: expected name arguments");
};
Ok(name.starts_with(&**prefix))
// Strictly longer, matching the interpreter: `/a` is not a prefix match of `/a`.
Ok(name.starts_with(&**prefix) && name.len() > prefix.len())
}
// :list:member(Elem, List) check mode. A non-list yields false,
// matching the interpreter (so the negation keeps the row).
":list:member" => match &vals[1] {
Val::Compound(CompoundKindMirror::List, elems) => Ok(elems.contains(&vals[0])),
_ => Ok(false),
},
// :match_field(Struct, Field, Value) check mode: the struct has the
// field with that value. A non-struct or missing field yields false,
// matching the interpreter.
":match_field" => match (&vals[0], &vals[1]) {
(Val::Compound(CompoundKindMirror::Struct, kvs), Val::Name(_)) => {
// Struct layout: [k1, v1, k2, v2, ...]
Ok(kvs
.chunks_exact(2)
.any(|kv| kv[0] == vals[1] && kv[1] == vals[2]))
}
_ => Ok(false),
},
other => bail!("unsupported CallFilter function: {other}"),
}
}
Expand Down Expand Up @@ -212,13 +247,17 @@ fn eval_aggregate(agg: &LoweredAggregate, input: &[(&Row, isize)]) -> Val {
/// `T = u64` for the top-level batch scope; `T = Product<u64, u32>` for the
/// recursive inner scope in Phase 3.
///
/// Runtime evaluation errors (see [`eval_call_filter`]) are pushed onto
/// `errors` as single-column rows holding the message.
///
/// Returns the output `VecCollection<Row>` — its rows are the tuples to be
/// inserted into `rule.head_rel`. Call `.distinct()` after concatenating all
/// rules for the same head relation.
pub fn build_rule<'scope, T>(
rule: &LoweredRule,
rels: &HashMap<String, VecCollection<'scope, T, Row>>,
unit_coll: &VecCollection<'scope, T, Row>,
errors: &mut Vec<VecCollection<'scope, T, Row>>,
) -> Result<VecCollection<'scope, T, Row>>
where
T: Timestamp + differential_dataflow::lattice::Lattice + Ord + 'static,
Expand Down Expand Up @@ -352,7 +391,7 @@ where
// ---------------------------------------------------------------
// CallFilter — Phase 2 string builtins.
// ---------------------------------------------------------------
Step::CallFilter { func, args } => {
Step::CallFilter { func, args, negate } => {
let pipeline = curr
.take()
.ok_or_else(|| anyhow::anyhow!("CallFilter before Scan"))?;
Expand All @@ -364,10 +403,22 @@ where
}
let func = func.clone();
let args = args.clone();
curr = Some(
pipeline
.filter(move |row| eval_call_filter(&func, &args, row).unwrap_or(false)),
);
let negate = *negate;
// A dataflow can't abort mid-evaluation the way the interpreter
// does, so an evaluation error drops the row (in both polarities)
// and is emitted into `errors` instead. The session refuses to
// answer reads while any error rows are live, which matches the
// interpreter failing the whole evaluation. Retracting the
// offending fact retracts its error row too.
let (err_func, err_args) = (func.clone(), args.clone());
errors.push(pipeline.clone().flat_map(move |row| {
eval_call_filter(&err_func, &err_args, &row)
.err()
.map(|e| Row(vec![Val::String(e.to_string().into())].into()))
}));
curr = Some(pipeline.filter(move |row| {
eval_call_filter(&func, &args, row).is_ok_and(|b| b != negate)
}));
}

// ---------------------------------------------------------------
Expand All @@ -378,10 +429,19 @@ where
.take()
.ok_or_else(|| anyhow::anyhow!("Let before Scan"))?;
let expr = expr.clone();
curr = Some(pipeline.map(move |row| {
let v = eval_expr(&expr, &row);
row.appended(std::iter::once(v))
}));
// Same scheme as `CallFilter`: a failing row is dropped and its
// error emitted into `errors`. Evaluate once and tag each row
// with its outcome, since functions can be costlier than checks.
let evaluated = pipeline.map(move |row| match eval_expr(&expr, &row) {
Ok(v) => (None, row.appended(std::iter::once(v))),
Err(e) => (Some(e.to_string()), row),
});
errors.push(
evaluated.clone().flat_map(|(err, _row)| {
err.map(|m| Row(vec![Val::String(m.into())].into()))
}),
);
curr = Some(evaluated.flat_map(|(err, row)| err.is_none().then_some(row)));
}

// ---------------------------------------------------------------
Expand Down
Loading
Loading