diff --git a/Cargo.lock b/Cargo.lock index 4a0d1b6..5712fb2 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1102,8 +1102,9 @@ checksum = "f9f8bd3e56ce4dfc153cf470fffbfa98c7620958b312ca5c3a4b8d5181fd13c6" [[package]] name = "mangle-analysis" -version = "0.9.0" -source = "git+https://codeberg.org/ajwdev/mangle-rs.git?branch=indexed-store#914de0c05f69de13d011d2691c4d5c71be4bb94e" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d125cc3d936c7032f3114869b1f0f1b9424b73cd1ac43e8f2bbbed0e7fbf5627" dependencies = [ "anyhow", "googletest", @@ -1115,8 +1116,9 @@ dependencies = [ [[package]] name = "mangle-ast" -version = "0.9.0" -source = "git+https://codeberg.org/ajwdev/mangle-rs.git?branch=indexed-store#914de0c05f69de13d011d2691c4d5c71be4bb94e" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e18df1a428f1345bfadb4ee53098ab00773d13ae7aeb77758c7e93874f25e4b" dependencies = [ "bumpalo", "googletest", @@ -1125,8 +1127,9 @@ dependencies = [ [[package]] name = "mangle-codegen" -version = "0.9.0" -source = "git+https://codeberg.org/ajwdev/mangle-rs.git?branch=indexed-store#914de0c05f69de13d011d2691c4d5c71be4bb94e" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "77b4bd4ef94859715d9208c07d977ab86a84446d4c159e4da77ba728368b590e" dependencies = [ "anyhow", "mangle-analysis", @@ -1138,8 +1141,9 @@ dependencies = [ [[package]] name = "mangle-common" -version = "0.9.0" -source = "git+https://codeberg.org/ajwdev/mangle-rs.git?branch=indexed-store#914de0c05f69de13d011d2691c4d5c71be4bb94e" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2bb54b86e00dcb774dd2aeba6606fe80ffb804c8e8cbd614f97e5ae1d9e0949d" dependencies = [ "anyhow", "mangle-ast", @@ -1148,8 +1152,9 @@ dependencies = [ [[package]] name = "mangle-driver" -version = "0.9.0" -source = "git+https://codeberg.org/ajwdev/mangle-rs.git?branch=indexed-store#914de0c05f69de13d011d2691c4d5c71be4bb94e" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3e04aeeb55bea36c5e3266d715f1e4146f19ef5d8204baba7cb16bcea1cc3a42" dependencies = [ "anyhow", "mangle-analysis", @@ -1163,8 +1168,9 @@ dependencies = [ [[package]] name = "mangle-interpreter" -version = "0.9.0" -source = "git+https://codeberg.org/ajwdev/mangle-rs.git?branch=indexed-store#914de0c05f69de13d011d2691c4d5c71be4bb94e" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7b8130ccb090fa514e75eb0142f4a8221ff3e6e120d19cc582b959173811f15e" dependencies = [ "anyhow", "mangle-ast", @@ -1174,13 +1180,15 @@ dependencies = [ [[package]] name = "mangle-ir" -version = "0.9.0" -source = "git+https://codeberg.org/ajwdev/mangle-rs.git?branch=indexed-store#914de0c05f69de13d011d2691c4d5c71be4bb94e" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dd10e9ef415e853ee61988a502da8776fcbcc55ca8818962261a007b0e3b8b73" [[package]] name = "mangle-parse" -version = "0.9.0" -source = "git+https://codeberg.org/ajwdev/mangle-rs.git?branch=indexed-store#914de0c05f69de13d011d2691c4d5c71be4bb94e" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c832749116021c7ac9db30c75b68f20ef82e6cbf2a0f9b1b1e1e7d94ca7aaa60" dependencies = [ "anyhow", "googletest", diff --git a/Cargo.toml b/Cargo.toml index 37d6813..a367830 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -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"] } @@ -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" diff --git a/src/dd/build.rs b/src/dd/build.rs index 11613d9..bd8ecbd 100644 --- a/src/dd/build.rs +++ b/src/dd/build.rs @@ -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 { match expr { - OwnedExpr::Value(slot) => slot_val(slot, row), + OwnedExpr::Value(slot) => Ok(slot_val(slot, row)), OwnedExpr::Call { func, args } => { let vals: Vec = 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. @@ -76,6 +79,10 @@ 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. @@ -83,33 +90,61 @@ 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 { let vals: Vec = 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}"), } } @@ -212,6 +247,9 @@ fn eval_aggregate(agg: &LoweredAggregate, input: &[(&Row, isize)]) -> Val { /// `T = u64` for the top-level batch scope; `T = Product` 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` — its rows are the tuples to be /// inserted into `rule.head_rel`. Call `.distinct()` after concatenating all /// rules for the same head relation. @@ -219,6 +257,7 @@ pub fn build_rule<'scope, T>( rule: &LoweredRule, rels: &HashMap>, unit_coll: &VecCollection<'scope, T, Row>, + errors: &mut Vec>, ) -> Result> where T: Timestamp + differential_dataflow::lattice::Lattice + Ord + 'static, @@ -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"))?; @@ -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) + })); } // --------------------------------------------------------------- @@ -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))); } // --------------------------------------------------------------- diff --git a/src/dd/lower.rs b/src/dd/lower.rs index 61c29a8..5580ab8 100644 --- a/src/dd/lower.rs +++ b/src/dd/lower.rs @@ -106,8 +106,13 @@ pub enum Step { const_filters: Vec<(usize, Val)>, }, - /// String-builtin filter (`:string:starts_with`, etc.). - CallFilter { func: String, args: Vec }, + /// Built-in predicate filter (`:string:starts_with`, etc.). With `negate` + /// set, keeps rows where the predicate is false. + CallFilter { + func: String, + args: Vec, + negate: bool, + }, /// Let-binding: append one computed column to each row. Let { expr: OwnedExpr }, @@ -265,6 +270,90 @@ fn lower_aggregate(agg: &Aggregate, ir: &Ir, schema: &[String]) -> Result, +) -> Result<()> { + match cond { + Condition::Cmp { op, left, right } => { + use mangle_ir::physical::CmpOp as MiCmpOp; + let my_op = match (op, negate) { + (MiCmpOp::Eq, false) | (MiCmpOp::Neq, true) => CmpOp::Eq, + (MiCmpOp::Neq, false) | (MiCmpOp::Eq, true) => CmpOp::Neq, + (MiCmpOp::Lt, false) | (MiCmpOp::Ge, true) => CmpOp::Lt, + (MiCmpOp::Le, false) | (MiCmpOp::Gt, true) => CmpOp::Le, + (MiCmpOp::Gt, false) | (MiCmpOp::Le, true) => CmpOp::Gt, + (MiCmpOp::Ge, false) | (MiCmpOp::Lt, true) => CmpOp::Ge, + }; + let left_slot = resolve_operand(left, ir, schema)?; + let right_slot = resolve_operand(right, ir, schema)?; + steps.push(Step::Cmp { + op: my_op, + left: left_slot, + right: right_slot, + }); + } + Condition::Negation { .. } if negate => { + // Double negation of a relation lookup is a semi-join; the planner + // never emits it. + bail!("negated relation negation is not supported by the DD backend: !{cond:?}") + } + Condition::Negation { relation, args } => { + let rel_name = ir.resolve_name(*relation).to_string(); + let mut left_key_slots = Vec::new(); + let mut right_key_cols = Vec::new(); + let mut const_filters = Vec::new(); + + for (right_col, arg) in args.iter().enumerate() { + match arg { + Operand::Var(name_id) => { + let name = ir.resolve_name(*name_id); + if let Some(left_pos) = schema.iter().position(|v| v == name) { + // Shared variable: join on it. + left_key_slots.push(Slot::Col(left_pos)); + right_key_cols.push(right_col); + } + // Anonymous/wildcard var (not in schema) → no constraint. + } + Operand::Const(c) => { + const_filters.push((right_col, resolve_constant(c, ir))); + } + } + } + steps.push(Step::Antijoin { + rel: rel_name, + left_key_slots, + right_key_cols, + const_filters, + }); + } + Condition::Call { function, args } => { + let func_name = ir.resolve_name(*function).to_string(); + let arg_slots: Result> = args + .iter() + .map(|a| resolve_operand(a, ir, schema)) + .collect(); + steps.push(Step::CallFilter { + func: func_name, + args: arg_slots?, + negate, + }); + } + Condition::Not(inner) => lower_cond(inner, !negate, ir, schema, steps)?, + } + Ok(()) +} + fn lower_inner( op: &Op, ir: &Ir, @@ -359,66 +448,7 @@ fn lower_inner( } Op::Filter { cond, body } => { - match cond { - Condition::Cmp { op, left, right } => { - use mangle_ir::physical::CmpOp as MiCmpOp; - let my_op = match op { - MiCmpOp::Eq => CmpOp::Eq, - MiCmpOp::Neq => CmpOp::Neq, - MiCmpOp::Lt => CmpOp::Lt, - MiCmpOp::Le => CmpOp::Le, - MiCmpOp::Gt => CmpOp::Gt, - MiCmpOp::Ge => CmpOp::Ge, - }; - let left_slot = resolve_operand(left, ir, schema)?; - let right_slot = resolve_operand(right, ir, schema)?; - steps.push(Step::Cmp { - op: my_op, - left: left_slot, - right: right_slot, - }); - } - Condition::Negation { relation, args } => { - let rel_name = ir.resolve_name(*relation).to_string(); - let mut left_key_slots = Vec::new(); - let mut right_key_cols = Vec::new(); - let mut const_filters = Vec::new(); - - for (right_col, arg) in args.iter().enumerate() { - match arg { - Operand::Var(name_id) => { - let name = ir.resolve_name(*name_id); - if let Some(left_pos) = schema.iter().position(|v| v == name) { - // Shared variable: join on it. - left_key_slots.push(Slot::Col(left_pos)); - right_key_cols.push(right_col); - } - // Anonymous/wildcard var (not in schema) → no constraint. - } - Operand::Const(c) => { - const_filters.push((right_col, resolve_constant(c, ir))); - } - } - } - steps.push(Step::Antijoin { - rel: rel_name, - left_key_slots, - right_key_cols, - const_filters, - }); - } - Condition::Call { function, args } => { - let func_name = ir.resolve_name(*function).to_string(); - let arg_slots: Result> = args - .iter() - .map(|a| resolve_operand(a, ir, schema)) - .collect(); - steps.push(Step::CallFilter { - func: func_name, - args: arg_slots?, - }); - } - } + lower_cond(cond, false, ir, schema, steps)?; lower_inner(body, ir, schema, steps, head_rel) } diff --git a/src/dd/session.rs b/src/dd/session.rs index 356d73f..0e44645 100644 --- a/src/dd/session.rs +++ b/src/dd/session.rs @@ -80,14 +80,17 @@ pub enum Command { /// Snapshot a relation from the sink and send it back. /// Milestone A: the worker reads the mutex on behalf of the caller so /// the API is uniform; Milestone B will cursor the trace instead. + /// + /// Responds with `Err` while any rule has live evaluation errors. Query { rel: String, - resp: Sender>>, + resp: Sender>, String>>, }, /// Snapshot EVERY relation's current contents (EDB + IDB) and send them back. /// Used by the batch `DdBackend::evaluate` path. + /// Responds with `Err` while any rule has live evaluation errors. SnapshotAll { - resp: Sender>>>, + resp: Sender>>, String>>, }, /// Attempt to add new IDB rules by layering a fresh dataflow. /// @@ -150,6 +153,44 @@ fn drain_trace(trace: &mut RowTrace) -> Vec> { result } +/// Concatenate a dataflow's runtime error collections and arrange them. +/// +/// Always returns a trace (empty when no rule can error) so the caller can +/// treat every dataflow uniformly. +fn arrange_errors<'scope>( + error_colls: Vec>, + unit_coll: &VecCollection<'scope, u64, Row>, + probe: &timely::dataflow::ProbeHandle, +) -> RowTrace { + let errors = error_colls + .into_iter() + .reduce(|a, b| a.concat(b)) + .unwrap_or_else(|| unit_coll.clone().filter(|_| false)); + errors.distinct().probe_with(probe).arrange_by_self().trace +} + +/// `Err` with every distinct live evaluation error, if there are any. +/// +/// The interpreter fails the whole evaluation on the first such error; a +/// dataflow can't stop mid-flight, so instead reads are refused for as long +/// as any error row is live. +fn live_errors(error_traces: &mut [RowTrace]) -> std::result::Result<(), String> { + let mut msgs: Vec = error_traces + .iter_mut() + .flat_map(drain_trace) + .map(|row| match row.first() { + Some(Value::String(m)) => m.clone(), + other => format!("{other:?}"), + }) + .collect(); + if msgs.is_empty() { + return Ok(()); + } + msgs.sort(); + msgs.dedup(); + Err(format!("evaluation error: {}", msgs.join("\n"))) +} + // --------------------------------------------------------------------------- // DdSession // --------------------------------------------------------------------------- @@ -235,178 +276,193 @@ impl DdSession { let input_rels = Arc::clone(&input_rels); let edb_by_rel = Arc::clone(&edb_by_rel); - let (mut handles, mut traces, build_errors) = worker.dataflow::({ - let input_rels = Arc::clone(&input_rels); - let probe_ref = probe.clone(); - - move |scope| { - let mut handles: HashMap> = - HashMap::new(); - let mut rels: HashMap> = HashMap::new(); - // Any rule the DD builder can't translate is collected here and - // reported back to spawn() so it fails loudly rather than - // silently dropping the rule (which would yield wrong results). - let mut build_errors: Vec = Vec::new(); - - for rel in input_rels.iter() { - let (handle, coll) = scope.new_collection::(); - rels.insert(rel.clone(), coll); - handles.insert(rel.clone(), handle); - } + let (mut handles, mut traces, build_errors, error_trace) = worker + .dataflow::({ + let input_rels = Arc::clone(&input_rels); + let probe_ref = probe.clone(); + + move |scope| { + let mut handles: HashMap> = + HashMap::new(); + let mut rels: HashMap> = HashMap::new(); + // Any rule the DD builder can't translate is collected here and + // reported back to spawn() so it fails loudly rather than + // silently dropping the rule (which would yield wrong results). + let mut build_errors: Vec = Vec::new(); + // Runtime evaluation errors from every rule (see + // `build_rule`), arranged below into `error_trace`. + let mut error_colls: Vec> = Vec::new(); + + for rel in input_rels.iter() { + let (handle, coll) = scope.new_collection::(); + rels.insert(rel.clone(), coll); + handles.insert(rel.clone(), handle); + } - let unit_coll = scope.new_collection_from(vec![Row::empty()]).1; + let unit_coll = scope.new_collection_from(vec![Row::empty()]).1; - for stratum in strata_work.iter() { - if stratum.rules.is_empty() { - continue; - } + for stratum in strata_work.iter() { + if stratum.rules.is_empty() { + continue; + } - if stratum.is_recursive { - let head_preds: std::collections::HashSet = - stratum.rules.iter().map(|r| r.head_rel.clone()).collect(); - - let results: HashMap> = scope - .iterative::(|nested| { - let summary = Product::new(Default::default(), 1u32); - - let mut inner_rels: HashMap< - String, - VecCollection<'_, Product, Row>, - > = rels - .iter() - .map(|(k, v)| (k.clone(), v.clone().enter(nested))) - .collect(); - let inner_unit = unit_coll.clone().enter(nested); - - let mut vars: HashMap< - String, - VecVariable<'_, Product, Row, isize>, - > = HashMap::new(); - let mut var_colls: HashMap< - String, - VecCollection<'_, Product, Row>, - > = HashMap::new(); - for pred in &head_preds { - if let Some(seed) = inner_rels.remove(pred) { - let (var, coll) = VecVariable::new_from(seed, summary); - vars.insert(pred.clone(), var); - var_colls.insert(pred.clone(), coll); - } else { - let (var, coll) = VecVariable::new(nested, summary); - vars.insert(pred.clone(), var); - var_colls.insert(pred.clone(), coll); - } - } - for (pred, coll) in &var_colls { - inner_rels.insert(pred.clone(), coll.clone()); - } + if stratum.is_recursive { + let head_preds: std::collections::HashSet = + stratum.rules.iter().map(|r| r.head_rel.clone()).collect(); - let mut by_head: HashMap< - String, - Vec, Row>>, - > = HashMap::new(); - for rule in &stratum.rules { - match build_rule(rule, &inner_rels, &inner_unit) { - Ok(coll) => { - by_head - .entry(rule.head_rel.clone()) - .or_default() - .push(coll); - } - Err(e) => { - build_errors.push(format!( - "rule for `{}`: {e}", - rule.head_rel - )); + let results: HashMap> = + scope.iterative::(|nested| { + let summary = Product::new(Default::default(), 1u32); + + let mut inner_rels: HashMap< + String, + VecCollection<'_, Product, Row>, + > = rels + .iter() + .map(|(k, v)| (k.clone(), v.clone().enter(nested))) + .collect(); + let inner_unit = unit_coll.clone().enter(nested); + + let mut vars: HashMap< + String, + VecVariable<'_, Product, Row, isize>, + > = HashMap::new(); + let mut var_colls: HashMap< + String, + VecCollection<'_, Product, Row>, + > = HashMap::new(); + for pred in &head_preds { + if let Some(seed) = inner_rels.remove(pred) { + let (var, coll) = + VecVariable::new_from(seed, summary); + vars.insert(pred.clone(), var); + var_colls.insert(pred.clone(), coll); + } else { + let (var, coll) = VecVariable::new(nested, summary); + vars.insert(pred.clone(), var); + var_colls.insert(pred.clone(), coll); } } - } + for (pred, coll) in &var_colls { + inner_rels.insert(pred.clone(), coll.clone()); + } - let mut out: HashMap> = - HashMap::new(); - for (pred, var) in vars { - let curr = var_colls.remove(&pred).unwrap(); - let full = match by_head.remove(&pred) { - Some(colls) => { - let new_facts = colls - .into_iter() - .reduce(|a, b| a.concat(b)) - .unwrap(); - curr.concat(new_facts).distinct() + let mut by_head: HashMap< + String, + Vec, Row>>, + > = HashMap::new(); + let mut inner_errors = Vec::new(); + for rule in &stratum.rules { + match build_rule( + rule, + &inner_rels, + &inner_unit, + &mut inner_errors, + ) { + Ok(coll) => { + by_head + .entry(rule.head_rel.clone()) + .or_default() + .push(coll); + } + Err(e) => { + build_errors.push(format!( + "rule for `{}`: {e}", + rule.head_rel + )); + } } - None => curr.distinct(), - }; - var.set(full.clone()); - out.insert(pred, full.leave(scope)); - } - out - }); + } - for (pred, coll) in results { - match rels.entry(pred) { - std::collections::hash_map::Entry::Occupied(mut e) => { - let old = e.get().clone(); - *e.get_mut() = old.concat(coll).distinct(); - } - std::collections::hash_map::Entry::Vacant(e) => { - e.insert(coll); + let mut out: HashMap> = + HashMap::new(); + for (pred, var) in vars { + let curr = var_colls.remove(&pred).unwrap(); + let full = match by_head.remove(&pred) { + Some(colls) => { + let new_facts = colls + .into_iter() + .reduce(|a, b| a.concat(b)) + .unwrap(); + curr.concat(new_facts).distinct() + } + None => curr.distinct(), + }; + var.set(full.clone()); + out.insert(pred, full.leave(scope)); + } + error_colls.extend( + inner_errors.into_iter().map(|e| e.leave(scope)), + ); + out + }); + + for (pred, coll) in results { + match rels.entry(pred) { + std::collections::hash_map::Entry::Occupied(mut e) => { + let old = e.get().clone(); + *e.get_mut() = old.concat(coll).distinct(); + } + std::collections::hash_map::Entry::Vacant(e) => { + e.insert(coll); + } } } - } - } else { - let mut by_head: HashMap>> = - HashMap::new(); - - for rule in &stratum.rules { - match build_rule(rule, &rels, &unit_coll) { - Ok(coll) => { - by_head - .entry(rule.head_rel.clone()) - .or_default() - .push(coll); - } - Err(e) => { - build_errors - .push(format!("rule for `{}`: {e}", rule.head_rel)); + } else { + let mut by_head: HashMap>> = + HashMap::new(); + + for rule in &stratum.rules { + match build_rule(rule, &rels, &unit_coll, &mut error_colls) { + Ok(coll) => { + by_head + .entry(rule.head_rel.clone()) + .or_default() + .push(coll); + } + Err(e) => { + build_errors + .push(format!("rule for `{}`: {e}", rule.head_rel)); + } } } - } - for (head_rel, colls) in by_head { - let idb = colls - .into_iter() - .reduce(|a, b| a.concat(b)) - .unwrap() - .distinct(); - match rels.entry(head_rel) { - std::collections::hash_map::Entry::Occupied(mut e) => { - let old = e.get().clone(); - *e.get_mut() = old.concat(idb).distinct(); - } - std::collections::hash_map::Entry::Vacant(e) => { - e.insert(idb); + for (head_rel, colls) in by_head { + let idb = colls + .into_iter() + .reduce(|a, b| a.concat(b)) + .unwrap() + .distinct(); + match rels.entry(head_rel) { + std::collections::hash_map::Entry::Occupied(mut e) => { + let old = e.get().clone(); + *e.get_mut() = old.concat(idb).distinct(); + } + std::collections::hash_map::Entry::Vacant(e) => { + e.insert(idb); + } } } } } - } - // Arrange each output collection. - // - // `arrange_by_self` writes every update into an ordered, on-worker - // trace (TraceAgent). Queries cursor the trace at the current - // frontier; compaction keeps memory bounded after each commit. - // probe_with is called before arranging so the ProbeHandle - // registers this edge in the dataflow graph. - let mut traces = HashMap::new(); - for (rel_name, coll) in rels { - let arranged = coll.probe_with(&probe_ref).arrange_by_self(); - traces.insert(rel_name, arranged.trace); - } + // Arrange each output collection. + // + // `arrange_by_self` writes every update into an ordered, on-worker + // trace (TraceAgent). Queries cursor the trace at the current + // frontier; compaction keeps memory bounded after each commit. + // probe_with is called before arranging so the ProbeHandle + // registers this edge in the dataflow graph. + let mut traces = HashMap::new(); + for (rel_name, coll) in rels { + let arranged = coll.probe_with(&probe_ref).arrange_by_self(); + traces.insert(rel_name, arranged.trace); + } + let error_trace = arrange_errors(error_colls, &unit_coll, &probe_ref); - (handles, traces, build_errors) - } - }); + (handles, traces, build_errors, error_trace) + } + }); // If any rule failed to translate, abandon the (unrun) dataflow and // report the failure to spawn() instead of silently proceeding with @@ -444,6 +500,8 @@ impl DdSession { // diffs into the InputSession; Commit advances the epoch and steps // the worker until the probe clears, then acks. Shutdown exits. let mut epoch: u64 = 1; + // One error trace per dataflow (the initial one plus each layer). + let mut error_traces: Vec = vec![error_trace]; loop { match rx.recv() { @@ -470,22 +528,27 @@ impl DdSession { // and memory stays bounded (we never do time-travel queries). let frontier_elems = [epoch]; let frontier = AntichainRef::new(&frontier_elems); - for trace in traces.values_mut() { + for trace in traces.values_mut().chain(error_traces.iter_mut()) { trace.set_logical_compaction(frontier); trace.set_physical_compaction(frontier); } let _ = ack.send(()); } Command::Query { rel, resp } => { - let rows = traces.get_mut(&rel).map(drain_trace).unwrap_or_default(); - let _ = resp.send(rows); + let result = live_errors(&mut error_traces).map(|()| { + traces.get_mut(&rel).map(drain_trace).unwrap_or_default() + }); + let _ = resp.send(result); } Command::SnapshotAll { resp } => { - let mut out: HashMap>> = HashMap::new(); - for (rel, trace) in traces.iter_mut() { - out.insert(rel.clone(), drain_trace(trace)); - } - let _ = resp.send(out); + let result = live_errors(&mut error_traces).map(|()| { + let mut out: HashMap>> = HashMap::new(); + for (rel, trace) in traces.iter_mut() { + out.insert(rel.clone(), drain_trace(trace)); + } + out + }); + let _ = resp.send(result); } Command::AddRules { all_rule_sources, @@ -538,8 +601,8 @@ impl DdSession { traces.iter().map(|(k, v)| (k.clone(), v.clone())).collect(); let probe_ref = probe.clone(); - let (mut new_traces, layer_errors) = - worker.dataflow::(move |scope| { + let (mut new_traces, layer_errors, mut layer_error_trace) = worker + .dataflow::(move |scope| { // Import each existing trace as a VecCollection. let mut rels: HashMap> = imported_traces @@ -555,6 +618,7 @@ impl DdSession { let unit_coll = scope.new_collection_from(vec![Row::empty()]).1; let mut new_inner: HashMap = HashMap::new(); let mut layer_errors: Vec = Vec::new(); + let mut layer_error_colls = Vec::new(); for stratum in &to_layer { let mut by_head: HashMap< @@ -562,7 +626,12 @@ impl DdSession { Vec>, > = HashMap::new(); for rule in &stratum.rules { - match build_rule(rule, &rels, &unit_coll) { + match build_rule( + rule, + &rels, + &unit_coll, + &mut layer_error_colls, + ) { Ok(coll) => { by_head .entry(rule.head_rel.clone()) @@ -591,7 +660,9 @@ impl DdSession { } } - (new_inner, layer_errors) + let layer_error_trace = + arrange_errors(layer_error_colls, &unit_coll, &probe_ref); + (new_inner, layer_errors, layer_error_trace) }); // A rule failed to translate: reject the addition and keep @@ -608,13 +679,17 @@ impl DdSession { // Compact to the current frontier so history stays bounded. let frontier_elems = [epoch]; let frontier = AntichainRef::new(&frontier_elems); - for trace in new_traces.values_mut() { + for trace in new_traces + .values_mut() + .chain(std::iter::once(&mut layer_error_trace)) + { trace.set_logical_compaction(frontier); trace.set_physical_compaction(frontier); } // Merge new relation traces into the session. traces.extend(new_traces); + error_traces.push(layer_error_trace); let _ = ack.send(AddOutcome::Layered); } Command::Shutdown { ack } => { @@ -687,23 +762,34 @@ impl DdSession { /// /// Milestone A: reads the `Arc` sink via the worker (uniform API). /// Milestone B: will cursor a TraceAgent on the worker thread instead. - pub fn query(&self, rel: &str) -> Vec> { + /// + /// Fails while any rule has live evaluation errors (e.g. a `:string:*` + /// built-in applied to a non-string), mirroring the interpreter. + pub fn query(&self, rel: &str) -> Result>> { let (resp_tx, resp_rx) = bounded(1); let _ = self.tx.send(Command::Query { rel: rel.to_string(), resp: resp_tx, }); - resp_rx.recv().unwrap_or_default() + resp_rx + .recv() + .unwrap_or_else(|_| Ok(Vec::new())) + .map_err(anyhow::Error::msg) } /// Snapshot every relation (EDB + all derived IDB) at the current frontier. /// /// Used by the batch `DdBackend::evaluate` path: spawn a session, snapshot, /// then drop. Mirrors `query()` but returns all relations at once. - pub fn snapshot_all(&self) -> HashMap>> { + /// + /// Fails while any rule has live evaluation errors, like [`Self::query`]. + pub fn snapshot_all(&self) -> Result>>> { let (resp_tx, resp_rx) = bounded(1); let _ = self.tx.send(Command::SnapshotAll { resp: resp_tx }); - resp_rx.recv().unwrap_or_default() + resp_rx + .recv() + .unwrap_or_else(|_| Ok(HashMap::new())) + .map_err(anyhow::Error::msg) } /// Attempt to add a new IDB rule by layering a fresh dataflow on top of the @@ -844,7 +930,7 @@ mod tests { ]; let session = DdSession::spawn(&edb, &rules).expect("spawn"); - let mut results = session.query("path"); + let mut results = session.query("path").unwrap(); results.sort(); assert_eq!( results, @@ -865,14 +951,14 @@ mod tests { let session = DdSession::spawn(&edb, &rules).expect("spawn"); // Before insert: only a→b - let before = session.query("path"); + let before = session.query("path").unwrap(); assert_eq!(before.len(), 1); // Insert b→c, commit, then check session.insert("edge".to_string(), Row::from(&[v_str("b"), v_str("c")][..])); session.commit(); - let mut after = session.query("path"); + let mut after = session.query("path").unwrap(); after.sort(); assert_eq!( after, @@ -895,13 +981,13 @@ mod tests { let session = DdSession::spawn(&edb, &rules).expect("spawn"); // Both edges visible initially - assert_eq!(session.query("path").len(), 2); + assert_eq!(session.query("path").unwrap().len(), 2); // Retract b→c session.retract("edge".to_string(), Row::from(&[v_str("b"), v_str("c")][..])); session.commit(); - let after = session.query("path"); + let after = session.query("path").unwrap(); assert_eq!(after, vec![vec![v_str("a"), v_str("b")]]); } @@ -942,7 +1028,7 @@ mod tests { .add_idb(Some("view"), &edb, &all_rules) .expect("add_idb"); - let mut results = session.query("view"); + let mut results = session.query("view").unwrap(); results.sort(); // view should contain the source nodes of every path edge: a and b assert_eq!(results, vec![vec![v_str("a")], vec![v_str("b")]]); @@ -967,14 +1053,14 @@ mod tests { .expect("add_idb"); // Before insert: only a - let before = session.query("view"); + let before = session.query("view").unwrap(); assert_eq!(before, vec![vec![v_str("a")]]); // Insert b→c, commit — view should now also contain b. session.insert("edge".to_string(), Row::from(&[v_str("b"), v_str("c")][..])); session.commit(); - let mut after = session.query("view"); + let mut after = session.query("view").unwrap(); after.sort(); assert_eq!(after, vec![vec![v_str("a")], vec![v_str("b")]]); } @@ -1007,7 +1093,7 @@ mod tests { .add_idb(Some("view2"), &edb, &rules_v2) .expect("add_idb view2"); - let mut results = session.query("view2"); + let mut results = session.query("view2").unwrap(); results.sort(); assert_eq!(results, vec![vec![v_str("a")], vec![v_str("b")]]); } @@ -1033,7 +1119,7 @@ mod tests { .expect("add_idb fallback"); // After rebuild, path should include both original edges and the reflexive pairs. - let mut results = session.query("path"); + let mut results = session.query("path").unwrap(); results.sort(); assert!( results.contains(&vec![v_str("a"), v_str("b")]), @@ -1067,6 +1153,9 @@ mod tests { assert!(is_supported_call_filter(":string:ends_with")); assert!(is_supported_call_filter(":string:contains")); assert!(is_supported_call_filter(":match_prefix")); + // Check modes, reached only via negation. + assert!(is_supported_call_filter(":list:member")); + assert!(is_supported_call_filter(":match_field")); // A typo / not-yet-implemented builtin must be rejected. assert!(!is_supported_call_filter(":string:startswith")); assert!(!is_supported_call_filter(":string:matches")); @@ -1087,8 +1176,99 @@ mod tests { ]; let session = DdSession::spawn(&edb, &rules).expect("spawn with valid builtin"); - let mut results = session.query("hit"); + let mut results = session.query("hit").unwrap(); results.sort(); assert_eq!(results, vec![vec![v_str("alpha")]]); } + + /// A negated built-in stays correct as facts are inserted and retracted. + #[test] + fn session_negated_builtin_incremental() { + let edb: Vec<(String, Vec)> = vec![ + ("name".to_string(), vec![v_str("alpha")]), + ("name".to_string(), vec![v_str("beta")]), + ]; + let rules = vec![ + "Decl name(N).\nDecl miss(N).\nmiss(N) :- name(N), !:string:starts_with(N, \"al\")." + .to_string(), + ]; + + let session = DdSession::spawn(&edb, &rules).expect("spawn"); + assert_eq!(session.query("miss").unwrap(), vec![vec![v_str("beta")]]); + + session.insert("name".to_string(), Row::from(&[v_str("gamma")][..])); + session.insert("name".to_string(), Row::from(&[v_str("alps")][..])); + session.commit(); + let mut after_insert = session.query("miss").unwrap(); + after_insert.sort(); + assert_eq!( + after_insert, + vec![vec![v_str("beta")], vec![v_str("gamma")]] + ); + + session.retract("name".to_string(), Row::from(&[v_str("beta")][..])); + session.commit(); + assert_eq!(session.query("miss").unwrap(), vec![vec![v_str("gamma")]]); + } + + /// A built-in type error makes reads fail while the offending fact is + /// live, and retracting it makes reads succeed again. + #[test] + fn session_eval_error_tracks_offending_fact() { + let edb: Vec<(String, Vec)> = vec![ + ("name".to_string(), vec![v_str("alpha")]), + ("name".to_string(), vec![v_str("beta")]), + ]; + let rules = vec![ + "Decl name(N).\nDecl miss(N).\nmiss(N) :- name(N), !:string:starts_with(N, \"al\")." + .to_string(), + ]; + + let session = DdSession::spawn(&edb, &rules).expect("spawn"); + assert_eq!(session.query("miss").unwrap(), vec![vec![v_str("beta")]]); + + let bad = Row::from(&[Value::Number(5)][..]); + session.insert("name".to_string(), bad.clone()); + session.commit(); + let err = session + .query("miss") + .expect_err("type error should fail the read"); + assert!( + format!("{err:#}").contains(":string:starts_with: expected string arguments"), + "unexpected error: {err:#}" + ); + assert!(session.snapshot_all().is_err()); + + session.retract("name".to_string(), bad); + session.commit(); + assert_eq!(session.query("miss").unwrap(), vec![vec![v_str("beta")]]); + } + + /// A failing `let` function behaves like a built-in type error: reads + /// fail while the offending fact is live, and recover once it is retracted. + #[test] + fn session_let_error_tracks_offending_fact() { + let edb: Vec<(String, Vec)> = vec![("num".to_string(), vec![Value::Number(1)])]; + let rules = vec![ + "Decl num(X).\nDecl next(Y).\nnext(Y) :- num(X) |> let Y = fn:plus(X, 1).".to_string(), + ]; + + let session = DdSession::spawn(&edb, &rules).expect("spawn"); + assert_eq!(session.query("next").unwrap(), vec![vec![Value::Number(2)]]); + + let bad = Row::from(&[v_str("x")][..]); + session.insert("num".to_string(), bad.clone()); + session.commit(); + let err = session + .query("next") + .expect_err("fn:plus on a string should fail the read"); + assert!( + format!("{err:#}").contains("fn:plus: expected integer"), + "unexpected error: {err:#}" + ); + + session.retract("num".to_string(), bad); + session.commit(); + assert_eq!(session.query("next").unwrap(), vec![vec![Value::Number(2)]]); + } } diff --git a/src/engine.rs b/src/engine.rs index ce08ea4..16a9734 100644 --- a/src/engine.rs +++ b/src/engine.rs @@ -317,6 +317,9 @@ impl CompiledProgram { args, )?; } + Condition::Not(inner) => { + writeln!(w, "{}Not {:?}", prop_prefix, inner)?; + } } self.fprint_op(w, body, level + 1) @@ -472,7 +475,7 @@ impl Backend for DdBackend { // sees fully-derived state. Dropping the session shuts the worker down. let session = crate::dd::session::DdSession::spawn(edb, rule_sources).context("spawn dd session")?; - let facts = session.snapshot_all(); + let facts = session.snapshot_all()?; Ok(EvalStore { facts, provenance: vec![], @@ -758,11 +761,12 @@ impl Engine { /// Query a relation from the live DD session. /// - /// Returns an empty vec if no session is active or the relation is unknown. - pub fn query_live(&self, rel: &str) -> Vec> { + /// Returns an empty vec if no session is active or the relation is unknown, + /// and an error while any rule has live evaluation errors. + pub fn query_live(&self, rel: &str) -> Result>> { match &self.session { Some(s) => s.query(rel), - None => vec![], + None => Ok(vec![]), } } @@ -906,7 +910,7 @@ impl Engine { // rather than spawning (and settling) a second dataflow. if let Some(s) = &self.session { return Ok(EvalStore { - facts: s.snapshot_all(), + facts: s.snapshot_all()?, provenance: vec![], }); } @@ -1078,7 +1082,7 @@ mod tests { assert!(dd.has_session()); assert!(dd.add_fact("brand_new".into(), vec![Value::String("x".into())])); - let live = dd.query_live("brand_new"); + let live = dd.query_live("brand_new").unwrap(); assert_eq!(live, vec![vec![Value::String("x".into())]]); let mut interp = Engine::from_parts(vec![], rules, Box::new(InterpreterBackend)).unwrap(); @@ -1097,12 +1101,347 @@ mod tests { fn dd_declared_but_empty_relation_is_queryable() { let rules = vec!["Decl lonely(X).".to_string()]; let mut dd = Engine::from_parts(vec![], rules, Box::new(DdBackend)).unwrap(); - assert!(dd.query_live("lonely").is_empty()); + assert!(dd.query_live("lonely").unwrap().is_empty()); dd.add_rule("Decl seen(X).\nseen(X) :- lonely(X).".to_string()) .expect("rule over declared-but-empty relation"); - assert!(dd.query_live("seen").is_empty()); + assert!(dd.query_live("seen").unwrap().is_empty()); let store = dd.evaluate().unwrap(); assert!(store.scan("lonely").is_empty()); assert!(store.scan("seen").is_empty()); } + + // ----------------------------------------------------------------------- + // Negated built-in predicates (Condition::Not) + // + // Each case runs on both backends and is checked against a hand-written + // expectation, not just backend-vs-backend: before mangle-rs 0.9.1 both + // backends planned `!:builtin(..)` as a lookup of a never-populated + // relation, so they agreed with each other while both being wrong. + // ----------------------------------------------------------------------- + + use mangle_common::CompoundKind; + + fn num(n: i64) -> Value { + Value::Number(n) + } + fn st(s: &str) -> Value { + Value::String(s.to_string()) + } + fn nm(s: &str) -> Value { + Value::Name(s.to_string()) + } + fn list(vs: Vec) -> Value { + Value::Compound(CompoundKind::List, vs) + } + /// Struct layout is interleaved: [k1, v1, k2, v2, ...]. + fn strukt(fields: &[(&str, Value)]) -> Value { + let mut kvs = Vec::new(); + for (k, v) in fields { + kvs.push(nm(k)); + kvs.push(v.clone()); + } + Value::Compound(CompoundKind::Struct, kvs) + } + + fn facts(rel: &str, rows: Vec>) -> Vec<(String, Vec)> { + rows.into_iter().map(|r| (rel.to_string(), r)).collect() + } + + /// Evaluate `rules` on both backends and assert each yields exactly + /// `expected` for `rel`. + fn assert_both( + edb: &[(String, Vec)], + rules: &str, + rel: &str, + mut expected: Vec>, + ) { + expected.sort(); + let rules = vec![rules.to_string()]; + let backends: [(&str, &dyn Backend); 2] = + [("interpreter", &InterpreterBackend), ("dd", &DdBackend)]; + for (name, backend) in backends { + let store = backend + .evaluate(edb, &rules) + .unwrap_or_else(|e| panic!("{name} failed: {e:#}")); + let mut got = store.scan(rel).to_vec(); + got.sort(); + assert_eq!(got, expected, "{name} backend, relation {rel}"); + } + } + + #[test] + fn negated_cmp_builtins() { + let edb = facts( + "np_pair", + vec![ + vec![num(1), num(2)], + vec![num(4), num(4)], + vec![num(5), num(3)], + ], + ); + let lt = vec![vec![num(4), num(4)], vec![num(5), num(3)]]; + let le = vec![vec![num(5), num(3)]]; + let gt = vec![vec![num(1), num(2)], vec![num(4), num(4)]]; + let ge = vec![vec![num(1), num(2)]]; + for (op, expected) in [("lt", lt), ("le", le), ("gt", gt), ("ge", ge)] { + let rules = + format!("Decl np_pair(X, Y).\nnp_out(X, Y) :- np_pair(X, Y), !:{op}(X, Y)."); + assert_both(&edb, &rules, "np_out", expected); + } + } + + #[test] + fn negated_cmp_against_constant() { + let edb = facts("np_num", vec![vec![num(1)], vec![num(3)], vec![num(5)]]); + assert_both( + &edb, + "Decl np_num(X).\nnp_out(X) :- np_num(X), !:lt(X, 3).", + "np_out", + vec![vec![num(3)], vec![num(5)]], + ); + } + + #[test] + fn negated_time_and_duration_cmp() { + let edb = vec![ + ("np_t".to_string(), vec![Value::Time(10), Value::Time(20)]), + ("np_t".to_string(), vec![Value::Time(30), Value::Time(20)]), + ( + "np_d".to_string(), + vec![Value::Duration(5), Value::Duration(5)], + ), + ( + "np_d".to_string(), + vec![Value::Duration(9), Value::Duration(5)], + ), + ]; + assert_both( + &edb, + "Decl np_t(A, B).\nnp_tout(A) :- np_t(A, B), !:time:lt(A, B).", + "np_tout", + vec![vec![Value::Time(30)]], + ); + assert_both( + &edb, + "Decl np_d(A, B).\nnp_dout(A) :- np_d(A, B), !:duration:gt(A, B).", + "np_dout", + vec![vec![Value::Duration(5)]], + ); + } + + #[test] + fn negated_string_builtins() { + let edb = facts( + "np_s", + vec![vec![st("alpha")], vec![st("beta")], vec![st("gamma")]], + ); + let cases = [ + ( + r#"!:string:starts_with(S, "al")"#, + vec![st("beta"), st("gamma")], + ), + ( + r#"!:string:ends_with(S, "ta")"#, + vec![st("alpha"), st("gamma")], + ), + ( + r#"!:string:contains(S, "mm")"#, + vec![st("alpha"), st("beta")], + ), + ]; + for (premise, expected) in cases { + let rules = format!("Decl np_s(S).\nnp_out(S) :- np_s(S), {premise}."); + let expected = expected.into_iter().map(|v| vec![v]).collect(); + assert_both(&edb, &rules, "np_out", expected); + } + } + + /// `:match_prefix` requires the name to be strictly longer than the + /// prefix, so `/a` does not match itself and survives the negation. + #[test] + fn negated_match_prefix() { + let edb = facts( + "np_n", + vec![vec![nm("/a")], vec![nm("/a/b")], vec![nm("/c/d")]], + ); + assert_both( + &edb, + "Decl np_n(N).\nnp_out(N) :- np_n(N), !:match_prefix(N, /a).", + "np_out", + vec![vec![nm("/a")], vec![nm("/c/d")]], + ); + } + + /// Positive form of the same strictness rule. + #[test] + fn match_prefix_excludes_exact_match() { + let edb = facts( + "np_n", + vec![vec![nm("/a")], vec![nm("/a/b")], vec![nm("/c/d")]], + ); + assert_both( + &edb, + "Decl np_n(N).\nnp_out(N) :- np_n(N), :match_prefix(N, /a).", + "np_out", + vec![vec![nm("/a/b")]], + ); + } + + /// A non-list second argument makes the check false, so the negation + /// keeps the row. + #[test] + fn negated_list_member() { + let edb = [ + facts( + "np_holder", + vec![ + vec![list(vec![num(1), num(2)])], + vec![list(vec![])], + vec![st("not a list")], + ], + ), + facts("np_elem", vec![vec![num(1)], vec![num(3)]]), + ] + .concat(); + let rules = "Decl np_holder(L).\nDecl np_elem(E).\n\ + np_out(E, L) :- np_holder(L), np_elem(E), !:list:member(E, L)."; + assert_both( + &edb, + rules, + "np_out", + vec![ + vec![num(3), list(vec![num(1), num(2)])], + vec![num(1), list(vec![])], + vec![num(3), list(vec![])], + vec![num(1), st("not a list")], + vec![num(3), st("not a list")], + ], + ); + } + + /// Evaluate `rules` on both backends and assert each fails with an error + /// mentioning `needle`. + fn assert_both_err(edb: &[(String, Vec)], rules: &str, needle: &str) { + let rules = vec![rules.to_string()]; + let backends: [(&str, &dyn Backend); 2] = + [("interpreter", &InterpreterBackend), ("dd", &DdBackend)]; + for (name, backend) in backends { + match backend.evaluate(edb, &rules) { + Ok(_) => { + panic!("{name} backend succeeded, expected an error containing {needle:?}") + } + Err(e) => { + let msg = format!("{e:#}"); + assert!( + msg.contains(needle), + "{name} backend error {msg:?} lacks {needle:?}" + ); + } + } + } + } + + /// Following upstream mangle, a `:string:*` built-in on a non-string fails + /// the whole evaluation, in both the positive and negated forms. + #[test] + fn string_builtin_type_error_fails_both_backends() { + let edb = facts( + "np_s", + vec![vec![st("alpha")], vec![Value::Null], vec![num(5)]], + ); + for premise in [ + r#":string:starts_with(S, "al")"#, + r#"!:string:starts_with(S, "al")"#, + ] { + let rules = format!("Decl np_s(S).\nnp_out(S) :- np_s(S), {premise}."); + assert_both_err( + &edb, + &rules, + ":string:starts_with: expected string arguments", + ); + } + } + + /// The error also surfaces from inside a recursive stratum. + #[test] + fn builtin_type_error_in_recursive_rule_fails_both_backends() { + let edb = [ + facts("np_start", vec![vec![st("n1")]]), + facts( + "np_edge", + vec![vec![st("n1"), st("n2")], vec![st("n2"), num(3)]], + ), + ] + .concat(); + let rules = "Decl np_start(X).\nDecl np_edge(X, Y).\n\ + np_reach(X) :- np_start(X).\n\ + np_reach(Y) :- np_reach(X), np_edge(X, Y), :string:starts_with(Y, \"n\")."; + assert_both_err( + &edb, + rules, + ":string:starts_with: expected string arguments", + ); + } + + /// A failing function in a `let` fails the whole evaluation. + #[test] + fn let_function_error_fails_both_backends() { + let edb = facts("np_v", vec![vec![num(1)], vec![st("x")]]); + assert_both_err( + &edb, + "Decl np_v(X).\nnp_out(Y) :- np_v(X) |> let Y = fn:plus(X, 1).", + "fn:plus: expected integer", + ); + } + + /// The happy path still evaluates normally. + #[test] + fn let_function_ok() { + let edb = facts("np_v", vec![vec![num(1)], vec![num(41)]]); + assert_both( + &edb, + "Decl np_v(X).\nnp_out(Y) :- np_v(X) |> let Y = fn:plus(X, 1).", + "np_out", + vec![vec![num(2)], vec![num(42)]], + ); + } + + /// Same for `:match_prefix` on a non-name. + #[test] + fn match_prefix_type_error_fails_both_backends() { + let edb = facts("np_n", vec![vec![nm("/a/b")], vec![st("/a/c")]]); + assert_both_err( + &edb, + "Decl np_n(N).\nnp_out(N) :- np_n(N), !:match_prefix(N, /a).", + ":match_prefix: expected name arguments", + ); + } + + /// A missing field or a non-struct scrutinee makes the check false, so + /// the negation keeps the row. + #[test] + fn negated_match_field() { + let matching = strukt(&[("/kind", st("a"))]); + let other_value = strukt(&[("/kind", st("b"))]); + let missing_field = strukt(&[("/other", st("a"))]); + let edb = facts( + "np_obj", + vec![ + vec![matching], + vec![other_value.clone()], + vec![missing_field.clone()], + vec![st("not a struct")], + ], + ); + assert_both( + &edb, + "Decl np_obj(S).\nnp_out(S) :- np_obj(S), !:match_field(S, /kind, \"a\").", + "np_out", + vec![ + vec![other_value], + vec![missing_field], + vec![st("not a struct")], + ], + ); + } } diff --git a/src/repl.rs b/src/repl.rs index 965ac5f..acc1029 100644 --- a/src/repl.rs +++ b/src/repl.rs @@ -2256,7 +2256,13 @@ fn list_tuples( ) { let live: Vec>; let rows: &[Vec] = if engine.has_session() { - live = engine.query_live(rel); + live = match engine.query_live(rel) { + Ok(rows) => rows, + Err(e) => { + eprintln!("Error: {e:#}"); + return; + } + }; &live } else { store.scan(rel)