From 50f158985c6854f44a89cbc349a43d1c6658d994 Mon Sep 17 00:00:00 2001 From: Frank McSherry Date: Thu, 24 Sep 2026 16:28:47 -0400 Subject: [PATCH 1/6] DDIR: floating-point math (fsqrt, fexp, fln, fpow, fmin/fmax, IEEE compares, fint) Adds the F64 operations a procedural-generation layer needs, spelled like the existing `fadd`/`fneg`: `f` + the Rust `f64` method name. - F64 -> F64: fabs, fsqrt, fexp, fln, ffloor, fceil, fround, fsin, fcos, ftan (`UnOp::F64Fn`), each the Rust method of the same name. - F64 -> Int: fint(x), Rust's `x as i64` (truncate toward zero, NaN is 0, saturating). - Binary: fpow (powf), fpowi (powi, Int exponent), fmin/fmax (skip a NaN operand, else total order, so signed zeros are deterministic), and the IEEE comparisons feq fne flt fle fgt fge returning Int 0/1. The generic `== < ...` were already right for negative F64 values: the payload is an order-preserving encoding, so they are `f64::total_cmp`. They differ from IEEE only on -0.0 vs 0.0 and NaN, which is what the f-comparisons are for. Corgi backend: fabs, fmin/fmax and the IEEE comparisons have columnar kernels built from existing corgi ops. The rest have no corgi kernel, so a term using one compiles to `Kernel::Rows`: the columns are untranscoded, `ir::eval` runs a row at a time, and the result is transcoded back. The output shape is still corgi's typing, of the term with each row-only op replaced by a same-shaped columnar stand-in. This applies to map, filter, flatmap, enter_at, and join projections. tests/programs/f64_math.ddp exercises every op, including negatives, both zeros, infinities and NaN, in a map, filter, flatmap and join, and is in the corgi-vs-vec gate. The tour gains a small float example. Co-Authored-By: Claude Opus 5.5 (1M context) --- interactive/examples/programs/tour.ddp | 6 + interactive/src/backend/corgi.rs | 24 +- interactive/src/corgi/join.rs | 2 +- interactive/src/corgi/logic.rs | 330 ++++++++++++++++++++---- interactive/src/ir.rs | 67 +++++ interactive/src/parse/mod.rs | 18 +- interactive/src/parse/pipe.rs | 12 + interactive/tests/corgi_backend.rs | 7 + interactive/tests/programs/f64_math.ddp | 47 ++++ 9 files changed, 450 insertions(+), 63 deletions(-) create mode 100644 interactive/tests/programs/f64_math.ddp diff --git a/interactive/examples/programs/tour.ddp b/interactive/examples/programs/tour.ddp index 3d2a3bedf..6ef8bb744 100644 --- a/interactive/examples/programs/tour.ddp +++ b/interactive/examples/programs/tour.ddp @@ -13,6 +13,11 @@ let scored = edges | map($0 ; $1[0] * 2, -$1[0], if($1[0] > 2, 1, 0), hash(97, $0[0])) | filter(not($1[1] == 0 - 99)); +-- Floats: explicit `float`, IEEE arithmetic and math, IEEE compares, `fint` (truncating) back. +let floated = edges + | map($0 ; fmul(fdiv(float(2786), float(10)), fpow(float($1[0]), fdiv(float(1), float(4))))) + | map($0 ; $1[0], fint($1[0]), fgt(fsqrt($1[0]), float(18))); + -- Sums: tagged intro, `case` with binder + captured var + default, `istag`. let tagged = edges | map($0 ; if($0[0] < $1[0], Fwd($1[0]), Bwd($1[0]))) @@ -36,6 +41,7 @@ reach: { } export "scored" = scored | arrange | inspect(total); +export "floats" = floated | arrange | inspect(total); export "sums" = tagged | arrange | inspect(total); export "folded" = folded | arrange | inspect(total); export "nested" = nested | arrange | inspect(nest); diff --git a/interactive/src/backend/corgi.rs b/interactive/src/backend/corgi.rs index 48b065e7a..8eadefe32 100644 --- a/interactive/src/backend/corgi.rs +++ b/interactive/src/backend/corgi.rs @@ -1,6 +1,7 @@ //! The corgi rendering substrate: corgi columns are the native representation on dataflow edges, //! arrangements are chains of sorted columnar chunks (`ChunkSpine`, cursor-less), and -//! scalar logic runs columnar via `eval_graph`. The row-wise `backend::vec` remains useful for +//! scalar logic runs columnar via `eval_graph` (a term using an F64 function corgi has no kernel +//! for runs a row at a time instead; see `logic::Kernel`). The row-wise `backend::vec` remains useful for //! comparison, but its representation choices do not define corgi's physical semantics. //! //! All `Backend` methods are corgi-native: `linear` folds a `LinearOp` chain over each container @@ -30,8 +31,8 @@ use crate::corgi::exchange::CorgiPact; use crate::corgi::join::CorgiJoinBackend; use crate::corgi::reduce::CorgiReduceBackend; use differential_dataflow::operators::int_proxy::{ProxyJoinTactic, ProxyReduceTactic}; -use crate::corgi::logic::{compile_flatmap, compile_predicate, compile_projection, compile_scalar, shape_of_row}; -use corgi::{Graph, NumOp, Shape}; +use crate::corgi::logic::{compile_flatmap, compile_predicate, compile_projection, compile_scalar, shape_of_row, Kernel}; +use corgi::Shape; use crate::ir::{Diff, LinearOp, Projection, Reducer, Time, Value as DValue}; use crate::scope_ir as st; @@ -47,13 +48,13 @@ type CTrace = differential_dataflow::trace::chunk::ChunkSpine)>, + compiled: Option<(Shape, Shape, Kernel)>, } impl Plan { /// The graph for a container of these shapes, compiling on first use. A type error is a /// panic with corgi's message: a program that typechecks never reaches it. - fn graph(&mut self, what: &str, kshape: Shape, vshape: Shape, compile: impl FnOnce(&Shape, &Shape) -> Result, String>) -> &Graph { + fn graph(&mut self, what: &str, kshape: Shape, vshape: Shape, compile: impl FnOnce(&Shape, &Shape) -> Result) -> &Kernel { if self.compiled.is_none() { let g = compile(&kshape, &vshape).unwrap_or_else(|e| panic!("{what}: type error at shapes ({kshape}, {vshape}): {e}")); self.compiled = Some((kshape.clone(), vshape.clone(), g)); @@ -66,8 +67,9 @@ impl Plan { /// Apply a `LinearOp` chain to one corgi container (the corgi-native compute per batch). /// Project = corgi `eval_graph`; Filter = corgi mask + `gather`; FlatMap = `eval_graph` to a list -/// column + a structural explode; Negate = Rust. Every term is columnar — there is no row-wise -/// path inside the dataflow — and each op's graph is compiled once (`plans`). An empty batch +/// column + a structural explode; Negate = Rust. Every term is columnar except one that uses an +/// F64 function corgi has no kernel for (`logic::row_only`), which `ir::eval` runs a row at a +/// time; either way each op's kernel is compiled once (`plans`). An empty batch /// passes through untouched: it carries no shape to compile against and no rows to compute. /// The two data<->time ops are columnar and total: EnterAt reads its delay field as a column and /// joins it into `times` in place; LiftIter reads the iteration coordinate out of `times` and @@ -87,14 +89,14 @@ fn apply_ops(mut c: CC, ops: &[LinearOp], level: usize, plans: &mut [Plan]) -> C c = match op { LinearOp::Project(p) => { let g = plan.graph("map", kshape, vshape, |k, v| compile_projection(&p.key, &p.val, k, v)); - let mut cols = corgi::eval_graph(g, CValue::Prod(vec![c.keys, c.vals])).into_prod("linear project").unwrap(); + let mut cols = g.eval(CValue::Prod(vec![c.keys, c.vals])).into_prod("linear project").unwrap(); let vals = cols.pop().unwrap(); let keys = cols.pop().unwrap(); CorgiContainer { keys, vals, times: c.times, diffs: c.diffs } } LinearOp::Filter(cond) => { let g = plan.graph("filter", kshape, vshape, |k, v| compile_predicate(cond, k, v)); - let mask = corgi::eval_graph(g, CValue::Prod(vec![c.keys.clone(), c.vals.clone()])).into_u64("filter mask").unwrap(); + let mask = g.eval(CValue::Prod(vec![c.keys.clone(), c.vals.clone()])).into_u64("filter mask").unwrap(); let keep: Vec = (0..mask.len()).filter(|&i| mask[i] != 0).collect(); let keys = gather(&c.keys, &keep); let vals = gather(&c.vals, &keep); @@ -116,7 +118,7 @@ fn apply_ops(mut c: CC, ops: &[LinearOp], level: usize, plans: &mut [Plan]) -> C // `level-1` and identity everywhere else (u64's minimum is 0), so the delta // never has to be built. The epoch is lane 0, so PointStamp index `level-1` is // lane `level`, and the join is one lane-wise max. - let raw = corgi::eval_graph(g, CValue::Prod(vec![c.keys.clone(), c.vals.clone()])) + let raw = g.eval(CValue::Prod(vec![c.keys.clone(), c.vals.clone()])) .into_u64("enter_at delay") .unwrap(); let delays: Vec = raw.iter().map(|r| 256 * (64 - r.leading_zeros() as u64)).collect(); @@ -150,7 +152,7 @@ fn apply_ops(mut c: CC, ops: &[LinearOp], level: usize, plans: &mut [Plan]) -> C // bounds gives both the within-row position (DDIR's `$1[0]`) and a repeat map // carrying key/time/diff across. No per-row eval, no transcode. let (bounds, elems) = - corgi::eval_graph(g, CValue::Prod(vec![c.keys.clone(), c.vals])).into_list("flatmap list").unwrap(); + g.eval(CValue::Prod(vec![c.keys.clone(), c.vals])).into_list("flatmap list").unwrap(); let ends: Vec = bounds.to_vec(); let total = ends.last().copied().unwrap_or(0); let (mut reps, mut pos) = (Vec::with_capacity(total), Vec::with_capacity(total)); diff --git a/interactive/src/corgi/join.rs b/interactive/src/corgi/join.rs index fdaa222a7..1fb24f2c0 100644 --- a/interactive/src/corgi/join.rs +++ b/interactive/src/corgi/join.rs @@ -201,7 +201,7 @@ impl ProxyJoinBackend, CBatch> for CorgiJoinBackend< } else { gather_lanes(&vals1, &tag1, &off1) }; let proj = compile_join_projection(&self.key, &self.val, &shape_of_value(&kc), &shape_of_value(&v0), &shape_of_value(&v1)) .unwrap_or_else(|e| panic!("join projection: type error: {e}")); - let projected = corgi::eval_graph(&proj, CValue::Prod(vec![kc, v0, v1])); + let projected = proj.eval(CValue::Prod(vec![kc, v0, v1])); let mut cols = projected.into_prod("corgi join projection").unwrap(); let nv = cols.pop().unwrap(); let nk = cols.pop().unwrap(); diff --git a/interactive/src/corgi/logic.rs b/interactive/src/corgi/logic.rs index 6600ce0d6..f34fc2c9e 100644 --- a/interactive/src/corgi/logic.rs +++ b/interactive/src/corgi/logic.rs @@ -14,7 +14,7 @@ //! reported with corgi's message. Ordered compares are signed-correct (`ToSigned`); `hash` is corgi's structural //! `Op::Hash`, the same function `ir::eval` folds row-wise. -use crate::ir::{BinOp, SumTy, Term, UnOp, Value as DValue}; +use crate::ir::{BinOp, F64Fn, SumTy, Term, UnOp, Value as DValue}; use corgi::{ArithOp, BinOp as CBinOp, Builder, CmpOp, Graph, Kind, NumOp, Op, Pred, Shape, Value as CValue}; @@ -326,6 +326,58 @@ pub fn compile( let payload = b.add(ArithOp::ToSigned, vec![f]); b.add(Op::Inject(0, vec![Shape::Prim(64)]), vec![payload]) } + BinOp::F64Min | BinOp::F64Max | BinOp::F64Eq | BinOp::F64Ne | BinOp::F64Lt | BinOp::F64Le | BinOp::F64Gt | BinOp::F64Ge => { + let expected = Shape::Sum(vec![Shape::Prim(64)]); + if shape_of_term(l, env_shapes, None)? != expected || shape_of_term(r, env_shapes, None)? != expected { + return Err(format!("{op:?} expects two F64 newtypes; use float(int)")); + } + let (x, y) = (float_leaf(b, lid), float_leaf(b, rid)); + match op { + BinOp::F64Min | BinOp::F64Max => { + // The total-order pick, then a NaN operand yields the other operand + // (and two NaNs the second), as `ir::eval` does. + let pick = if matches!(op, BinOp::F64Min) { CmpOp::Min } else { CmpOp::Max }; + let p = pair(b, x, y); + let total = b.add(pick, vec![p]); + let y_nan = is_nan(b, y); + let choices = b.tuple(vec![y_nan, x, total]); + let unless_y = b.add(Op::Select, vec![choices]); + let x_nan = is_nan(b, x); + let choices = b.tuple(vec![x_nan, y, unless_y]); + let f = b.add(Op::Select, vec![choices]); + let payload = b.add(ArithOp::ToSigned, vec![f]); + b.add(Op::Inject(0, vec![Shape::Prim(64)]), vec![payload]) + } + _ => { + // IEEE: the total order once `-0.0` is folded onto `0.0`, and false + // whenever an operand is NaN; `fne` is the negation of `feq`. + let (cx, cy) = (fold_negative_zero(b, x, anchor), fold_negative_zero(b, y, anchor)); + let (p, pred) = match op { + BinOp::F64Eq | BinOp::F64Ne => (pair(b, cx, cy), Pred::Eq), + BinOp::F64Lt => (pair(b, cx, cy), Pred::Lt), + BinOp::F64Le => (pair(b, cx, cy), Pred::Le), + BinOp::F64Gt => (pair(b, cy, cx), Pred::Lt), + _ => (pair(b, cy, cx), Pred::Le), + }; + let rel = b.add(CmpOp::Rel(pred), vec![p]); + let (x_nan, y_nan) = (is_nan(b, x), is_nan(b, y)); + let p = pair(b, x_nan, y_nan); + let any_nan = b.add(CmpOp::Max, vec![p]); + let zero = b.add(Op::Lit(CValue::u64(vec![0])), vec![anchor]); + let p = pair(b, any_nan, zero); + let ordered = b.add(CmpOp::Rel(Pred::Eq), vec![p]); + let p = pair(b, rel, ordered); + let holds = b.add(CmpOp::Min, vec![p]); + if matches!(op, BinOp::F64Ne) { + let p = pair(b, holds, zero); + b.add(CmpOp::Rel(Pred::Eq), vec![p]) + } else { + holds + } + } + } + } + BinOp::F64Pow | BinOp::F64PowI => return Err(format!("{op:?} has no columnar kernel; it runs a row at a time (see `row_only`)")), BinOp::Eq | BinOp::Ne => { // Cross-shape structural compare folds to a constant (Eq→0, Ne→1) over `anchor`; // same-shape emits a real corgi `Rel`. @@ -513,6 +565,16 @@ pub fn compile( let payload = b.add(ArithOp::ToSigned, vec![negative]); b.add(Op::Inject(0, vec![Shape::Prim(64)]), vec![payload]) } + // |x| is the larger of x and -x in the total order: that clears the sign bit, + // NaN included, exactly as `f64::abs`. + UnOp::F64Fn(F64Fn::Abs) => { + if shape != Shape::Sum(vec![Shape::Prim(64)]) { return Err("fabs expects an F64 newtype".into()); } + let f = float_leaf(b, id); + let abs = float_abs(b, f); + let payload = b.add(ArithOp::ToSigned, vec![abs]); + b.add(Op::Inject(0, vec![Shape::Prim(64)]), vec![payload]) + } + UnOp::F64Fn(_) | UnOp::F64ToInt => return Err(format!("{op:?} has no columnar kernel; it runs a row at a time (see `row_only`)")), // `truthy` is "nonzero Int": scalars compare against zero; non-`Int` values // are never truthy, so their `not` folds to the constant 1 (the cross-shape // `Eq` fold's precedent). @@ -622,6 +684,41 @@ pub fn compile( } } +/// Corgi's float encoding of an F64 constant: the total-order key, as a `U64` leaf holds it. +fn float_key(f: f64) -> u64 { + let bits = f.to_bits(); + if bits >> 63 == 1 { !bits } else { bits ^ (1 << 63) } +} + +/// An F64 newtype column -> its corgi float leaf (total-order key). The DDIR payload is the +/// signed form of that key, and `ToSigned` is the involution between the two. +fn float_leaf(b: &mut Builder, newtype: usize) -> usize { + let payload = b.add(Op::Unwrap, vec![newtype]); + b.add(ArithOp::ToSigned, vec![payload]) +} + +/// `|x|` on a float leaf: the larger of `x` and `-x` in the total order. +fn float_abs(b: &mut Builder, x: usize) -> usize { + let negated = b.add(ArithOp::Neg(Kind::F, 64), vec![x]); + let p = b.tuple(vec![x, negated]); + b.add(CmpOp::Max, vec![p]) +} + +/// A 0/1 mask: is the float leaf NaN? The NaNs are exactly the keys whose magnitude exceeds +inf. +fn is_nan(b: &mut Builder, x: usize) -> usize { + let abs = float_abs(b, x); + b.add(CmpOp::Gt(float_key(f64::INFINITY)), vec![abs]) +} + +/// Replace `-0.0` by `0.0` in a float leaf. Their keys are adjacent, so this adds the mask. +fn fold_negative_zero(b: &mut Builder, x: usize, anchor: usize) -> usize { + let negative_zero = b.add(Op::Lit(CValue::u64(vec![float_key(-0.0)])), vec![anchor]); + let p = b.tuple(vec![x, negative_zero]); + let is_negative_zero = b.add(CmpOp::Rel(Pred::Eq), vec![p]); + let p = b.tuple(vec![x, is_negative_zero]); + b.add(ArithOp::Bin(CBinOp::Add, Kind::U, 64), vec![p]) +} + /// Compile a `Fold` step into a closed corgi sub-graph. Without capture the body's input is /// `Prod([acc, elem])`; with it, `Prod([acc, (ctx, elem)])` where `ctx` is the captured /// environment (its fields come first, so `Var(i)` and outer `Bound`s resolve as they do in @@ -651,81 +748,151 @@ fn compile_fold_body(step: &Term, ctx: Option<&[Shape]>, init_shape: &Shape, ele Ok(bb.finish(out)) } -/// Compile a term in the row environment `Var(0)=key` (shape `kshape`), `Var(1)=val` (`vshape`) — -/// the environment every `LinearOp` reads. The graph's input is `Prod([key, val])`, and it is -/// typechecked once here, so an `Ok` graph runs on every batch of these shapes. -fn compile_over_kv(term: &Term, kshape: &Shape, vshape: &Shape) -> Res> { +/// A compiled term, ready to run on a batch of columns. +pub enum Kernel { + /// Every op in the terms has a corgi kernel: one columnar graph. + Columnar(Graph), + /// Some op has none (see [`row_only`]). The terms run through `ir::eval` a row at a time: + /// the input columns are untranscoded from `inputs`, and the results transcoded to `output`. + /// The output shape is still corgi's typing, of the terms with each row-only op replaced by a + /// columnar op of the same shape ([`columnar_stand_in`]), so both paths type a term alike. + Rows { terms: Vec, inputs: Vec, output: Shape }, +} + +impl Kernel { + /// Run on `Prod(inputs)`. One term yields its column; several yield `Prod` of their columns. + pub fn eval(&self, input: CValue) -> CValue { + match self { + Kernel::Columnar(g) => corgi::eval_graph(g, input), + Kernel::Rows { terms, inputs, output } => { + let rows = untranscode(input, &Shape::Prod(inputs.clone())); + let results: Vec = rows + .into_iter() + .map(|row| { + let DValue::Tuple(mut env) = row else { unreachable!("untranscode of a Prod is a Tuple") }; + match &terms[..] { + [term] => crate::ir::eval(term, &mut env), + _ => DValue::Tuple(terms.iter().map(|t| crate::ir::eval(t, &mut env)).collect()), + } + }) + .collect(); + transcode(&results, output) + } + } + } +} + +/// Does `t` use an op with no corgi kernel? These are the F64 functions other than `fabs`, `fint`, +/// `fpow` and `fpowi`: libm calls and rounding, which corgi's arithmetic does not offer. +pub fn row_only(t: &Term) -> bool { + match t { + Term::Var(_) | Term::Bound(_) | Term::Int(_) => false, + Term::Unary(UnOp::F64Fn(f), _) if *f != F64Fn::Abs => true, + Term::Unary(UnOp::F64ToInt, _) | Term::Binary(BinOp::F64Pow | BinOp::F64PowI, _, _) => true, + Term::Tuple(fs) | Term::List(fs) | Term::Hash(fs) => fs.iter().any(row_only), + Term::Spread(inner) | Term::Proj(inner, _) | Term::Unary(_, inner) => row_only(inner), + Term::Inject { tag, payload, .. } => row_only(tag) || row_only(payload), + Term::Case { scrutinee, arms, default } => { + row_only(scrutinee) || arms.iter().any(row_only) || default.as_deref().is_some_and(row_only) + } + Term::Fold { list, init, step } => row_only(list) || row_only(init) || row_only(step), + Term::If { cond, then, els } => row_only(cond) || row_only(then) || row_only(els), + Term::Binary(_, l, r) => row_only(l) || row_only(r), + } +} + +/// `t` with each row-only op replaced by a columnar op that demands the same operand shapes and +/// yields the same result shape: F64 -> F64 by `fneg`, `fint` by a compare of two `fneg`s, +/// `fpow` by `fadd`, and `fpowi(x, n)` by `fadd(x, float(n))`. Only its typing is used. +fn columnar_stand_in(t: &Term) -> Term { + let bx = |t: &Term| Box::new(columnar_stand_in(t)); + match t { + Term::Var(_) | Term::Bound(_) | Term::Int(_) => t.clone(), + Term::Unary(UnOp::F64Fn(f), x) if *f != F64Fn::Abs => Term::Unary(UnOp::F64Neg, bx(x)), + Term::Unary(UnOp::F64ToInt, x) => { + let x = Term::Unary(UnOp::F64Neg, bx(x)); + Term::Binary(BinOp::Lt, Box::new(x.clone()), Box::new(x)) + } + Term::Binary(BinOp::F64Pow, l, r) => Term::Binary(BinOp::F64Add, bx(l), bx(r)), + Term::Binary(BinOp::F64PowI, l, r) => Term::Binary(BinOp::F64Add, bx(l), Box::new(Term::Unary(UnOp::ToF64, bx(r)))), + Term::Tuple(fs) => Term::Tuple(fs.iter().map(columnar_stand_in).collect()), + Term::List(fs) => Term::List(fs.iter().map(columnar_stand_in).collect()), + Term::Hash(fs) => Term::Hash(fs.iter().map(columnar_stand_in).collect()), + Term::Spread(inner) => Term::Spread(bx(inner)), + Term::Proj(inner, i) => Term::Proj(bx(inner), *i), + Term::Unary(op, inner) => Term::Unary(*op, bx(inner)), + Term::Inject { tag, payload, sum } => Term::Inject { tag: bx(tag), payload: bx(payload), sum: sum.clone() }, + Term::Case { scrutinee, arms, default } => Term::Case { + scrutinee: bx(scrutinee), + arms: arms.iter().map(columnar_stand_in).collect(), + default: default.as_deref().map(bx), + }, + Term::Fold { list, init, step } => Term::Fold { list: bx(list), init: bx(init), step: bx(step) }, + Term::If { cond, then, els } => Term::If { cond: bx(cond), then: bx(then), els: bx(els) }, + Term::Binary(op, l, r) => Term::Binary(*op, bx(l), bx(r)), + } +} + +/// Compile `terms` over the environment `Var(i)` = field `i` of the input, whose shapes are +/// `shapes`. The graph's input is `Prod(shapes)`; its output is the one term's column, or `Prod` +/// of the terms' columns. Typechecked once here, so an `Ok` kernel runs on every batch of these +/// shapes. Returns the kernel and its output shape. +fn lower(terms: &[&Term], shapes: &[Shape]) -> Res<(Kernel, Shape)> { + let rows = terms.iter().any(|t| row_only(t)); + let stand_ins: Vec = if rows { terms.iter().map(|t| columnar_stand_in(t)).collect() } else { Vec::new() }; + let typed: Vec<&Term> = if rows { stand_ins.iter().collect() } else { terms.to_vec() }; let mut b = Builder::::default(); let input = b.input(); - let var_k = b.add(Op::Field(0), vec![input]); - let var_v = b.add(Op::Field(1), vec![input]); - let out = compile(term, &mut b, &[var_k, var_v], &[kshape.clone(), vshape.clone()], input, None)?; + let env: Vec = (0..shapes.len()).map(|i| b.add(Op::Field(i), vec![input])).collect(); + let outs = typed.iter().map(|t| compile(t, &mut b, &env, shapes, input, None)).collect::>>()?; + let out = if let [one] = outs[..] { one } else { b.tuple(outs) }; let g = b.finish(out); - corgi::shape_of(&g, &Shape::Prod(vec![kshape.clone(), vshape.clone()]))?; - Ok(g) + let output = corgi::shape_of(&g, &Shape::Prod(shapes.to_vec()))?; + let kernel = if rows { + Kernel::Rows { terms: terms.iter().map(|t| (*t).clone()).collect(), inputs: shapes.to_vec(), output: output.clone() } + } else { + Kernel::Columnar(g) + }; + Ok((kernel, output)) } /// Compile a `FlatMap`'s list term → a corgi `List` column, one list per input row. A term that is /// not list-shaped is the type error: the backend explodes the column structurally. -pub fn compile_flatmap(list_term: &Term, kshape: &Shape, vshape: &Shape) -> Res> { - match shape_of_term(list_term, &[kshape.clone(), vshape.clone()], None)? { - Shape::List(_) => compile_over_kv(list_term, kshape, vshape), - other => Err(format!("flatmap over a non-list: {other}")), +pub fn compile_flatmap(list_term: &Term, kshape: &Shape, vshape: &Shape) -> Res { + match lower(&[list_term], &[kshape.clone(), vshape.clone()])? { + (k, Shape::List(_)) => Ok(k), + (_, other) => Err(format!("flatmap over a non-list: {other}")), } } /// Compile a scalar term (`EnterAt`'s delay field) → a `U64` column; a non-integer term is the /// type error (the delay is read as one integer per row). -pub fn compile_scalar(term: &Term, kshape: &Shape, vshape: &Shape) -> Res> { - match shape_of_term(term, &[kshape.clone(), vshape.clone()], None)? { - Shape::Prim(_) => compile_over_kv(term, kshape, vshape), - other => Err(format!("enter_at delay is not an integer: {other}")), +pub fn compile_scalar(term: &Term, kshape: &Shape, vshape: &Shape) -> Res { + match lower(&[term], &[kshape.clone(), vshape.clone()])? { + (k, Shape::Prim(_)) => Ok(k), + (_, other) => Err(format!("enter_at delay is not an integer: {other}")), } } /// Compile a `Filter` predicate → a mask column (nonzero keeps the row). A predicate must be an /// `Int`; any other shape is a type error, as it is in the row backend. -pub fn compile_predicate(cond: &Term, kshape: &Shape, vshape: &Shape) -> Res> { - let g = compile_over_kv(cond, kshape, vshape)?; - match corgi::shape_of(&g, &Shape::Prod(vec![kshape.clone(), vshape.clone()]))? { - Shape::Prim(_) => Ok(g), - other => Err(format!("a filter predicate must be an Int, got {other}")), +pub fn compile_predicate(cond: &Term, kshape: &Shape, vshape: &Shape) -> Res { + match lower(&[cond], &[kshape.clone(), vshape.clone()])? { + (k, Shape::Prim(_)) => Ok(k), + (_, other) => Err(format!("a filter predicate must be an Int, got {other}")), } } /// Compile a join projection: key/val Terms over `Var(0)=key`, `Var(1)=val0`, `Var(2)=val1` (with /// their shapes). Input `Prod([key, val0, val1])`; output `Prod([newkey, newval])`. -pub fn compile_join_projection(key: &Term, val: &Term, kshape: &Shape, v0shape: &Shape, v1shape: &Shape) -> Res> { - let mut b = Builder::::default(); - let input = b.input(); - let var_k = b.add(Op::Field(0), vec![input]); - let var_0 = b.add(Op::Field(1), vec![input]); - let var_1 = b.add(Op::Field(2), vec![input]); - let env = [var_k, var_0, var_1]; - let shapes = [kshape.clone(), v0shape.clone(), v1shape.clone()]; - let nk = compile(key, &mut b, &env, &shapes, input, None)?; - let nv = compile(val, &mut b, &env, &shapes, input, None)?; - let out = b.tuple(vec![nk, nv]); - let g = b.finish(out); - corgi::shape_of(&g, &Shape::Prod(shapes.to_vec()))?; - Ok(g) +pub fn compile_join_projection(key: &Term, val: &Term, kshape: &Shape, v0shape: &Shape, v1shape: &Shape) -> Res { + Ok(lower(&[key, val], &[kshape.clone(), v0shape.clone(), v1shape.clone()])?.0) } /// Compile a DDIR `Projection` over `Var(0)=key` (`kshape`), `Var(1)=val` (`vshape`). /// Input `Prod([key, val])`; output `Prod([newkey, newval])`. -pub fn compile_projection(key: &Term, val: &Term, kshape: &Shape, vshape: &Shape) -> Res> { - let mut b = Builder::::default(); - let input = b.input(); - let var_k = b.add(Op::Field(0), vec![input]); - let var_v = b.add(Op::Field(1), vec![input]); - let env = [var_k, var_v]; - let shapes = [kshape.clone(), vshape.clone()]; - let nk = compile(key, &mut b, &env, &shapes, input, None)?; - let nv = compile(val, &mut b, &env, &shapes, input, None)?; - let out = b.tuple(vec![nk, nv]); - let g = b.finish(out); - corgi::shape_of(&g, &Shape::Prod(shapes.to_vec()))?; - Ok(g) +pub fn compile_projection(key: &Term, val: &Term, kshape: &Shape, vshape: &Shape) -> Res { + Ok(lower(&[key, val], &[kshape.clone(), vshape.clone()])?.0) } #[cfg(test)] @@ -766,6 +933,69 @@ mod tests { } } + /// Every pair of these F64 values: signed zeros, infinities, and both signs of NaN. + fn float_pairs() -> Vec> { + let specials = [f64::NEG_INFINITY, -2.5, -1.0, -0.0, 0.0, 0.5, 1.0, 2.0, 3.7, f64::INFINITY, f64::NAN, -f64::NAN]; + specials.iter().flat_map(|&a| specials.iter().map(move |&b| vec![V::f64_value(a), V::f64_value(b)])).collect() + } + + /// The F64 ops with a columnar kernel agree with `ir::eval` on every special pair. + #[test] + fn columnar_float_math_agrees() { + let f = sum(vec![u64s()]); + for source in ["fabs($0)", "fmin($0, $1)", "fmax($0, $1)", "feq($0, $1)", "fne($0, $1)", + "flt($0, $1)", "fle($0, $1)", "fgt($0, $1)", "fge($0, $1)"] { + let term = crate::parse::pipe::parse_term(source); + assert!(!row_only(&term), "{source} should have a columnar kernel"); + agrees_with_rows(&term, &[f.clone(), f.clone()], &float_pairs()); + } + } + + /// The F64 ops without one run a row at a time, typed as their columnar stand-ins; the + /// kernel's output round-trips through the transcode layer at that shape. + #[test] + fn row_wise_float_math_agrees() { + let f = sum(vec![u64s()]); + for source in ["fsqrt($0)", "fexp($0)", "fln($0)", "ffloor($0)", "fceil($0)", "fround($0)", + "fsin($0)", "fcos($0)", "ftan($0)", "fint($0)", "fpow($0, $1)", "fpowi($0, fint($1))", + "tuple(fadd(fexp($0), $1), flt(fln($0), $1))"] { + let term = crate::parse::pipe::parse_term(source); + assert!(row_only(&term), "{source} should run row-wise"); + let shapes = [f.clone(), f.clone()]; + let (kernel, output) = lower(&[&term], &shapes).unwrap(); + assert!(matches!(kernel, Kernel::Rows { .. })); + let rows = float_pairs(); + let cols = (0..2).map(|i| transcode(&rows.iter().map(|r| r[i].clone()).collect::>(), &shapes[i])).collect(); + let actual = untranscode(kernel.eval(CValue::Prod(cols)), &output); + let expected: Vec<_> = rows.iter().map(|r| crate::ir::eval(&term, &mut r.clone())).collect(); + assert_eq!(actual, expected, "{source}"); + } + // A row-only op is still typed: its operands must be F64 newtypes. + assert!(lower(&[&crate::parse::pipe::parse_term("fsqrt($0)")], &[u64s()]).is_err()); + assert!(lower(&[&crate::parse::pipe::parse_term("fpow($0, $0)")], &[u64s()]).is_err()); + } + + /// Spot checks of the chosen semantics, independent of either backend. + #[test] + fn float_math_semantics() { + let f = |x: f64| V::f64_value(x); + let ev = |src: &str, a: f64, b: f64| crate::ir::eval(&crate::parse::pipe::parse_term(src), &mut vec![f(a), f(b)]); + assert_eq!(ev("fint($0)", -2.7, 0.0), V::Int(-2)); + assert_eq!(ev("fint($0)", f64::NAN, 0.0), V::Int(0)); + assert_eq!(ev("fint($0)", 1e300, 0.0), V::Int(i64::MAX)); + assert_eq!(ev("flt($0, $1)", -2.0, -1.0), V::Int(1)); + assert_eq!(ev("$0 < $1", -2.0, -1.0), V::Int(1)); + assert_eq!(ev("feq($0, $1)", -0.0, 0.0), V::Int(1)); + assert_eq!(ev("$0 == $1", -0.0, 0.0), V::Int(0)); + assert_eq!(ev("fne($0, $1)", f64::NAN, f64::NAN), V::Int(1)); + assert_eq!(ev("fmax($0, $1)", f64::NAN, 1.0), f(1.0)); + assert_eq!(ev("fmin($0, $1)", 0.0, -0.0).as_f64().to_bits(), (-0.0f64).to_bits()); + assert!(ev("fln($0)", -1.0, 0.0).as_f64().is_nan()); + assert_eq!(ev("fln($0)", 0.0, 0.0), f(f64::NEG_INFINITY)); + assert_eq!(ev("fpow($0, $1)", 4.0, 0.5), f(2.0)); + assert!(ev("fpow($0, $1)", -8.0, 0.5).as_f64().is_nan()); + } + #[test] fn declared_constructor_types_an_empty_list() { let term = Term::Inject { diff --git a/interactive/src/ir.rs b/interactive/src/ir.rs index 46dda7249..96156e589 100644 --- a/interactive/src/ir.rs +++ b/interactive/src/ir.rs @@ -145,6 +145,43 @@ pub enum UnOp { ToF64, /// Floating-point negation; does not reinterpret integer arithmetic. F64Neg, + /// A one-argument F64 -> F64 function (`fsqrt`, `fexp`, ...), with Rust's `f64` semantics. + F64Fn(F64Fn), + /// F64 -> Int, truncating toward zero and saturating, exactly Rust's `x as i64`: + /// NaN is 0, and values beyond the `i64` range clamp to `i64::MIN`/`i64::MAX`. + F64ToInt, +} + +/// The one-argument F64 -> F64 functions. Each is the Rust `f64` method of the same name, so a +/// program ported from Rust gets bit-identical results on the same platform. `Abs` is exact +/// everywhere; the transcendental ones (`Exp`, `Ln`, `Sin`, `Cos`, `Tan`) call the platform libm +/// and are only as reproducible across platforms as it is. +#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +pub enum F64Fn { + Abs, Sqrt, Exp, Ln, Floor, Ceil, Round, Sin, Cos, Tan, +} + +impl F64Fn { + /// The surface names, `f` + the Rust method name (`ln`, not `log`). + pub const ALL: [(&'static str, F64Fn); 10] = [ + ("fabs", F64Fn::Abs), ("fsqrt", F64Fn::Sqrt), ("fexp", F64Fn::Exp), ("fln", F64Fn::Ln), + ("ffloor", F64Fn::Floor), ("fceil", F64Fn::Ceil), ("fround", F64Fn::Round), + ("fsin", F64Fn::Sin), ("fcos", F64Fn::Cos), ("ftan", F64Fn::Tan), + ]; + pub fn apply(self, x: f64) -> f64 { + match self { + F64Fn::Abs => x.abs(), + F64Fn::Sqrt => x.sqrt(), + F64Fn::Exp => x.exp(), + F64Fn::Ln => x.ln(), + F64Fn::Floor => x.floor(), + F64Fn::Ceil => x.ceil(), + F64Fn::Round => x.round(), + F64Fn::Sin => x.sin(), + F64Fn::Cos => x.cos(), + F64Fn::Tan => x.tan(), + } + } } #[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)] @@ -155,6 +192,20 @@ pub enum BinOp { /// Concatenation of two lists with the same element type. Append, F64Add, F64Sub, F64Mul, F64Div, + /// `x.powf(y)`. + F64Pow, + /// `x.powi(n)` with an Int exponent (saturated to `i32`). Rust's `powi` is repeated + /// multiplication, which need not round like `powf`; this is here so ported code can match. + F64PowI, + /// Minimum and maximum that skip a NaN operand (as Rust's `f64::min`/`max` do), and otherwise + /// follow the total order, so `fmin(-0.0, 0.0)` is `-0.0` and `fmax` of them is `0.0` — a + /// deterministic choice where Rust leaves signed zeros unspecified. Two NaNs give the second. + F64Min, F64Max, + /// IEEE comparisons, returning `Int` 0/1 like the generic ones: every comparison with a NaN is + /// false except `F64Ne`, which is true, and `-0.0 == 0.0`. The generic `== < ...` instead + /// order F64 values by the total order (`f64::total_cmp`), which is right for negative numbers + /// but distinguishes the two zeros and places NaNs at the ends. + F64Eq, F64Ne, F64Lt, F64Le, F64Gt, F64Ge, Eq, Ne, Lt, Le, Gt, Ge, And, Or, } @@ -364,6 +415,8 @@ fn eval_unary(op: UnOp, v: Value) -> Value { UnOp::Neg => Value::Int(-v.as_int()), UnOp::ToF64 => Value::f64_value(v.as_int() as f64), UnOp::F64Neg => Value::f64_value(-v.as_f64()), + UnOp::F64Fn(f) => Value::f64_value(f.apply(v.as_f64())), + UnOp::F64ToInt => Value::Int(v.as_f64() as i64), UnOp::Not => Value::Int((!v.truthy()) as i64), UnOp::IsTag(t) => Value::Int(matches!(&v, Value::Variant(tag, _) if *tag == t) as i64), UnOp::Len => match v { @@ -388,6 +441,20 @@ fn eval_binary(op: BinOp, l: Value, r: Value) -> Value { BinOp::F64Sub => Value::f64_value(l.as_f64() - r.as_f64()), BinOp::F64Mul => Value::f64_value(l.as_f64() * r.as_f64()), BinOp::F64Div => Value::f64_value(l.as_f64() / r.as_f64()), + BinOp::F64Pow => Value::f64_value(l.as_f64().powf(r.as_f64())), + BinOp::F64PowI => Value::f64_value(l.as_f64().powi(r.as_int().clamp(i32::MIN as i64, i32::MAX as i64) as i32)), + BinOp::F64Min | BinOp::F64Max => { + let (x, y) = (l.as_f64(), r.as_f64()); + let pick_x = if x.is_nan() { false } else if y.is_nan() { true } + else if matches!(op, BinOp::F64Min) { x.total_cmp(&y).is_le() } else { x.total_cmp(&y).is_ge() }; + if pick_x { l } else { r } + } + BinOp::F64Eq => b(l.as_f64() == r.as_f64()), + BinOp::F64Ne => b(l.as_f64() != r.as_f64()), + BinOp::F64Lt => b(l.as_f64() < r.as_f64()), + BinOp::F64Le => b(l.as_f64() <= r.as_f64()), + BinOp::F64Gt => b(l.as_f64() > r.as_f64()), + BinOp::F64Ge => b(l.as_f64() >= r.as_f64()), // Comparisons are structural, using the derived `Ord`/`Eq` on `Value`. BinOp::Eq => b(l == r), BinOp::Ne => b(l != r), diff --git a/interactive/src/parse/mod.rs b/interactive/src/parse/mod.rs index 5f4f106f4..b5dc99655 100644 --- a/interactive/src/parse/mod.rs +++ b/interactive/src/parse/mod.rs @@ -10,7 +10,7 @@ pub mod pipe; -use crate::ir::{BinOp, Projection, Reducer, SumTy, Term, UnOp}; +use crate::ir::{BinOp, F64Fn, Projection, Reducer, SumTy, Term, UnOp}; #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub enum Expr { @@ -78,6 +78,22 @@ pub(crate) fn build_builtin(name: &str, args: &mut Vec) -> Term { assert_eq!(args.len(), 1, "{name}(value)"); Term::Unary(if name == "float" { UnOp::ToF64 } else { UnOp::F64Neg }, Box::new(args.remove(0))) } + "fint" => { assert_eq!(args.len(), 1, "fint(value)"); Term::Unary(UnOp::F64ToInt, Box::new(args.remove(0))) } + name if F64Fn::ALL.iter().any(|(n, _)| *n == name) => { + assert_eq!(args.len(), 1, "{name}(value)"); + let f = F64Fn::ALL.iter().find(|(n, _)| *n == name).unwrap().1; + Term::Unary(UnOp::F64Fn(f), Box::new(args.remove(0))) + } + "fpow" | "fpowi" | "fmin" | "fmax" | "feq" | "fne" | "flt" | "fle" | "fgt" | "fge" => { + assert_eq!(args.len(), 2, "{name}(a, b)"); + let b = Box::new(args.remove(1)); let a = Box::new(args.remove(0)); + let op = match name { + "fpow" => BinOp::F64Pow, "fpowi" => BinOp::F64PowI, "fmin" => BinOp::F64Min, "fmax" => BinOp::F64Max, + "feq" => BinOp::F64Eq, "fne" => BinOp::F64Ne, "flt" => BinOp::F64Lt, "fle" => BinOp::F64Le, + "fgt" => BinOp::F64Gt, _ => BinOp::F64Ge, + }; + Term::Binary(op, a, b) + } "fadd" | "fsub" | "fmul" | "fdiv" => { assert_eq!(args.len(), 2, "{name}(a, b)"); let b = Box::new(args.remove(1)); let a = Box::new(args.remove(0)); diff --git a/interactive/src/parse/pipe.rs b/interactive/src/parse/pipe.rs index 2396b6eae..58dbed238 100644 --- a/interactive/src/parse/pipe.rs +++ b/interactive/src/parse/pipe.rs @@ -55,6 +55,18 @@ //! Nominal type names are erased: `fneg` and binary floating operators also //! accept a user's single-variant integer newtype, treating its payload as //! encoded f64 bits. +//! - Floating-point math, each the Rust `f64` method of the same name: +//! `fabs`, `fsqrt`, `fexp`, `fln`, `ffloor`, `fceil`, `fround`, `fsin`, +//! `fcos`, `ftan` (F64 -> F64); `fpow(x, y)` (`powf`), `fpowi(x, n)` (`powi`, +//! Int exponent). `fint(x)` is F64 -> Int as Rust's `x as i64`: truncate +//! toward zero, NaN is 0, out of range saturates. `fmin`/`fmax` skip a NaN +//! operand and otherwise follow the total order (`fmin(-0.0, 0.0)` is +//! `-0.0`). `feq fne flt fle fgt fge` are IEEE comparisons returning Int 0/1: +//! NaN compares false (`fne` true) and `-0.0` equals `0.0`, where the generic +//! `== < …` use the total order. There are no float literals: write +//! `fdiv(float(2786), float(10))` for 278.6. The Corgi backend computes +//! `fabs`, `fmin`/`fmax` and the comparisons in columns; a term using any +//! other of these is evaluated a row at a time. //! - Products: `tuple(a, …)`; index with `v[i]` or `proj(v, i)`; `len(v)`. //! - Lists: `list(a, …)`, `append(a, b)` (concatenation); eliminated by //! `flatmap` / `collect` / `fold`. A declared constructor supplies the element diff --git a/interactive/tests/corgi_backend.rs b/interactive/tests/corgi_backend.rs index 62fca1f3a..61cb8bc10 100644 --- a/interactive/tests/corgi_backend.rs +++ b/interactive/tests/corgi_backend.rs @@ -63,6 +63,12 @@ fn inputs_for(prog: &str) -> Vec> { &[2, 5], &[2, -2], ])], + // f64_math: (key, a, b) with x = a / 4, y = b / 2: negatives, zero, and a key with + // two rows; then (key, n) integer exponents for the join, one key unmatched. + "f64_math" => vec![ + rows(&[&[1, 10, 1], &[2, -6, 3], &[3, 0, -1], &[4, 9, 0], &[5, -1, -4], &[5, 7, 5]]), + rows(&[&[1, 2], &[2, 3], &[3, -1], &[5, 0]]), + ], // tour: edges (with a cycle and a chord) + roots. "tour" => vec![ rows(&[&[1, 2], &[2, 3], &[3, 1], &[3, 4], &[5, 2]]), @@ -135,6 +141,7 @@ fn serializing(n: usize) -> timely::Config { #[test] fn pair_keys() { assert_backends_agree("pair_keys"); } #[test] fn signed_min() { assert_backends_agree("signed_min"); } #[test] fn spread_values() { assert_backends_agree("spread_values"); } +#[test] fn f64_math() { assert_backends_agree("f64_math"); } /// A filter predicate must be an `Int`: both backends reject a tuple rather than one of them /// keeping nothing. diff --git a/interactive/tests/programs/f64_math.ddp b/interactive/tests/programs/f64_math.ddp new file mode 100644 index 000000000..ae8973285 --- /dev/null +++ b/interactive/tests/programs/f64_math.ddp @@ -0,0 +1,47 @@ +-- F64 math, checked against the vec backend. `fabs`, `fmin`/`fmax` and the IEEE +-- comparisons have columnar kernels; the libm functions, rounding, `fint`, `fpow` +-- and `fpowi` run a row at a time, in a map, a filter, a flatmap and a join. +-- Values include negatives, zeros of both signs, infinities and NaN. + +-- (key ; x, y) with x = a / 4 and y = b / 2. +let xs = input 0 | key($0[0] ; fdiv(float($0[1]), float(4)), fdiv(float($0[2]), float(2))); + +let unary = xs + | map($0 ; fabs($1[0]), fsqrt($1[0]), fexp($1[1]), fln($1[0]), ffloor($1[0]), fceil($1[0]), + fround($1[0]), fsin($1[0]), fcos($1[0]), ftan($1[0]), fint($1[0]), fint(fmul($1[0], float(1000)))); + +-- Negative bases with fractional exponents are NaN; `fpowi` is repeated multiplication. +let binary = xs + | map($0 ; fpow($1[0], $1[1]), fpow(fabs($1[0]), fdiv(float(1), float(4))), fpow($1[0], fdiv(float(1), float(3))), + fpowi($1[0], 3), fpowi($1[0], 0 - 2), fmin($1[0], $1[1]), fmax($1[0], $1[1])); + +-- IEEE compares beside the generic (total-order) ones: they agree on these finite values, +-- negatives included. +let compares = xs + | map($0 ; feq($1[0], $1[1]), fne($1[0], $1[1]), flt($1[0], $1[1]), fle($1[0], $1[1]), + fgt($1[0], $1[1]), fge($1[0], $1[1]), $1[0] < $1[1], $1[0] >= $1[1]); + +-- The special values: -0.0, +inf, NaN (0/0), and a NaN from fln/fsqrt of a negative. +let specials = xs + | map($0 ; fneg(float(0)), fdiv(float(1), float(0)), fdiv(float(0), float(0)), $1[0]) + | map($0 ; feq($1[0], float(0)), $1[0] == float(0), fle(float(0), $1[0]), float(0) <= $1[0], + flt($1[2], $1[3]), fgt($1[2], $1[3]), fne($1[2], $1[2]), feq($1[2], $1[2]), + fmin($1[2], $1[3]), fmax($1[3], $1[2]), fmin($1[0], float(0)), fmax(float(0), $1[0]), + fln(float(0)), fln(fneg(float(1))), fsqrt(fneg(float(1))), fexp($1[1]), fexp(fneg($1[1])), + fint($1[1]), fint(fneg($1[1])), fint($1[2]), fabs(fneg($1[2])), fabs($1[0]), flt($1[3], $1[1])); + +-- A row-at-a-time filter and flatmap. +let big = xs | filter(fgt(fexp($1[0]), float(2))); +let spread = xs | flatmap(list(fsqrt(fabs($1[0])), fexp($1[1]))); + +-- A join whose projection runs a row at a time. +let scale = input 1 | key($0[0] ; $0[1]); +let scaled = xs | join(scale, ($0 ; fpowi($1[0], $2[0]), fint(fpow(fabs($1[1]), float($2[0]))))); + +export "unary" = unary | arrange | inspect(total); +export "binary" = binary | arrange | inspect(total); +export "compares" = compares | arrange | inspect(total); +export "specials" = specials | arrange | inspect(total); +export "big" = big | arrange | inspect(total); +export "spread" = spread | arrange | inspect(total); +export "scaled" = scaled | arrange | inspect(total); From c9a5cb8ddc3f6ecec69886b6ea1475c88141d6e3 Mon Sep 17 00:00:00 2001 From: Frank McSherry Date: Thu, 24 Sep 2026 17:54:17 -0400 Subject: [PATCH 2/6] DDIR: registered functions, called by name from terms An embedding program registers a pure Rust function with `ir::register` (name, argument shapes, result shape, body), and any program parsed afterwards calls it like a builtin: `grow($0[0])`. The parser resolves a name that is not a builtin against the registry (a registered name may not shadow a builtin), producing `Term::Call`. - Row backend: `ir::eval` evaluates the arguments and calls the body. - Corgi backend: a call has no columnar kernel, so a term containing one runs through the row-at-a-time fallback (`Kernel::Rows`, from the float work). It is typed as a tuple of its arguments' stand-ins (so they still typecheck) projected to a literal of the declared result shape. The contract is purity: the same arguments must give the same value, since the dataflow re-evaluates terms to retract what they produced. tests/programs/registered.ddp calls three registered functions (a list-valued `grow`, a float-valued `blend`, a nested-shape `describe`) in a map, a filter, a flatmap, a join projection and a refinement fixpoint, and is in the corgi-vs-vec gate; the test also checks the fixpoint's contents. Co-Authored-By: Claude Opus 5.5 (1M context) --- interactive/src/corgi/logic.rs | 30 +++++++++++- interactive/src/ir.rs | 45 +++++++++++++++++ interactive/src/parse/mod.rs | 14 +++++- interactive/tests/corgi_backend.rs | 59 +++++++++++++++++++++++ interactive/tests/programs/registered.ddp | 31 ++++++++++++ 5 files changed, 176 insertions(+), 3 deletions(-) create mode 100644 interactive/tests/programs/registered.ddp diff --git a/interactive/src/corgi/logic.rs b/interactive/src/corgi/logic.rs index f34fc2c9e..89317dec5 100644 --- a/interactive/src/corgi/logic.rs +++ b/interactive/src/corgi/logic.rs @@ -171,7 +171,7 @@ fn mentions_env(t: &Term, depth: usize) -> bool { Term::Var(_) => true, Term::Bound(k) => *k >= depth, Term::Int(_) => false, - Term::Tuple(fs) | Term::List(fs) | Term::Hash(fs) => fs.iter().any(|f| mentions_env(f, depth)), + Term::Tuple(fs) | Term::List(fs) | Term::Hash(fs) | Term::Call(_, fs) => fs.iter().any(|f| mentions_env(f, depth)), Term::Spread(inner) | Term::Proj(inner, _) | Term::Unary(_, inner) => mentions_env(inner, depth), Term::Inject { tag, payload, .. } => mentions_env(tag, depth) || mentions_env(payload, depth), Term::Case { scrutinee, arms, default } => { @@ -249,6 +249,7 @@ pub fn compile( expected: Option<&Shape>, ) -> Res { match term { + Term::Call(name, _) => Err(format!("`{name}` is a registered function; it runs a row at a time (see `row_only`)")), Term::Var(i) => env.get(*i).copied().ok_or_else(|| format!("`${i}` is not in scope here")), Term::Bound(k) => { env.len().checked_sub(1 + *k).map(|i| env[i]).ok_or_else(|| format!("binder `^{k}` is not in scope here")) @@ -788,7 +789,7 @@ pub fn row_only(t: &Term) -> bool { match t { Term::Var(_) | Term::Bound(_) | Term::Int(_) => false, Term::Unary(UnOp::F64Fn(f), _) if *f != F64Fn::Abs => true, - Term::Unary(UnOp::F64ToInt, _) | Term::Binary(BinOp::F64Pow | BinOp::F64PowI, _, _) => true, + Term::Unary(UnOp::F64ToInt, _) | Term::Binary(BinOp::F64Pow | BinOp::F64PowI, _, _) | Term::Call(..) => true, Term::Tuple(fs) | Term::List(fs) | Term::Hash(fs) => fs.iter().any(row_only), Term::Spread(inner) | Term::Proj(inner, _) | Term::Unary(_, inner) => row_only(inner), Term::Inject { tag, payload, .. } => row_only(tag) || row_only(payload), @@ -814,6 +815,16 @@ fn columnar_stand_in(t: &Term) -> Term { Term::Binary(BinOp::Lt, Box::new(x.clone()), Box::new(x)) } Term::Binary(BinOp::F64Pow, l, r) => Term::Binary(BinOp::F64Add, bx(l), bx(r)), + // A registered function: its arguments still compile (so they typecheck), and the + // result is a literal of the declared result shape. + Term::Call(name, args) => match crate::ir::lookup(name) { + Some(f) => { + let mut fields: Vec = args.iter().map(columnar_stand_in).collect(); + fields.push(literal_of_shape(&f.result)); + Term::Proj(Box::new(Term::Tuple(fields)), args.len()) + } + None => t.clone(), + }, Term::Binary(BinOp::F64PowI, l, r) => Term::Binary(BinOp::F64Add, bx(l), Box::new(Term::Unary(UnOp::ToF64, bx(r)))), Term::Tuple(fs) => Term::Tuple(fs.iter().map(columnar_stand_in).collect()), Term::List(fs) => Term::List(fs.iter().map(columnar_stand_in).collect()), @@ -833,6 +844,21 @@ fn columnar_stand_in(t: &Term) -> Term { } } +/// A term whose corgi shape is `shape`: zeros, first lanes, one-element lists. +fn literal_of_shape(shape: &Shape) -> Term { + match shape { + Shape::Prim(_) => Term::Int(0), + Shape::Unit => Term::Tuple(Vec::new()), + Shape::Prod(fs) => Term::Tuple(fs.iter().map(literal_of_shape).collect()), + Shape::List(e) => Term::List(vec![literal_of_shape(e)]), + Shape::Sum(lanes) => Term::Inject { + tag: Box::new(Term::Int(0)), + payload: Box::new(literal_of_shape(&lanes[0])), + sum: SumTy::Declared(lanes.clone()), + }, + } +} + /// Compile `terms` over the environment `Var(i)` = field `i` of the input, whose shapes are /// `shapes`. The graph's input is `Prod(shapes)`; its output is the one term's column, or `Prod` /// of the terms' columns. Typechecked once here, so an `Ok` kernel runs on every batch of these diff --git a/interactive/src/ir.rs b/interactive/src/ir.rs index 96156e589..4b2b23ab3 100644 --- a/interactive/src/ir.rs +++ b/interactive/src/ir.rs @@ -112,6 +112,11 @@ pub enum Term { /// (the raw non-negative hash if `bound <= 0`), mixed from the keys. /// The building block for generators derived from `iota`/`clock`. Hash(Vec), + /// `name(args…)` for a function the embedding program registered (see + /// [`register`]): a pure Rust function from argument values to a value. + /// DDIR only moves its arguments and result; what it computes is the + /// embedder's. The columnar backend runs it a row at a time. + Call(String, Vec), } /// The sum type an `Inject` builds into. `Declared` carries the full lane shapes of a `type` @@ -297,6 +302,11 @@ pub fn structural_hash(v: &Value) -> u64 { /// around sub-evaluation, so `env` is restored on return. pub fn eval(term: &Term, env: &mut Vec) -> Value { match term { + Term::Call(name, args) => { + let f = lookup(name).unwrap_or_else(|| panic!("call to unregistered function `{name}`")); + let vals: Vec = args.iter().map(|a| eval(a, env)).collect(); + (f.body)(&vals) + } Term::Var(i) => env[*i].clone(), Term::Bound(k) => env[env.len() - 1 - *k].clone(), Term::Int(n) => Value::Int(*n), @@ -465,3 +475,38 @@ fn eval_binary(op: BinOp, l: Value, r: Value) -> Value { BinOp::And | BinOp::Or => unreachable!("logical ops short-circuit in eval"), } } + +// --------------------------------------------------------------------------- +// Registered functions: how an embedding program extends the scalar language. + +/// A function an embedding program supplies to DDIR programs, called by name. +/// +/// The contract is purity: the same arguments must always give the same +/// result, because differential dataflow re-evaluates terms when it retracts +/// what they produced, and a retraction must cancel exactly. The declared +/// shapes are what the columnar backend types the call as; values that do not +/// match them are an error of the embedder's. +pub struct Function { + pub name: String, + pub args: Vec, + pub result: corgi::Shape, + pub body: Box Value + Send + Sync>, +} + +fn functions() -> &'static std::sync::RwLock>> { + static REGISTRY: std::sync::OnceLock>>> = std::sync::OnceLock::new(); + REGISTRY.get_or_init(Default::default) +} + +/// Register `f` under its name, for programs parsed afterwards (the parser +/// resolves calls against the registry). Registering a name again replaces it; +/// a name may not shadow a builtin. +pub fn register(f: Function) { + assert!(!crate::parse::is_builtin(&f.name), "`{}` is a builtin and cannot be registered", f.name); + functions().write().unwrap().insert(f.name.clone(), std::sync::Arc::new(f)); +} + +/// The registered function called `name`, if any. +pub fn lookup(name: &str) -> Option> { + functions().read().unwrap().get(name).cloned() +} diff --git a/interactive/src/parse/mod.rs b/interactive/src/parse/mod.rs index b5dc99655..0ddf9313f 100644 --- a/interactive/src/parse/mod.rs +++ b/interactive/src/parse/mod.rs @@ -107,6 +107,18 @@ pub(crate) fn build_builtin(name: &str, args: &mut Vec) -> Term { } "if" => { assert_eq!(args.len(), 3, "if(cond, then, els)"); let els = Box::new(args.remove(2)); let then = Box::new(args.remove(1)); let cond = Box::new(args.remove(0)); Term::If { cond, then, els } } "hash" => { assert!(args.len() >= 2, "hash(bound, key, ...)"); Term::Hash(std::mem::take(args)) } - other => panic!("Unknown scalar builtin: {}", other), + other if crate::ir::lookup(other).is_some() => { + let f = crate::ir::lookup(other).unwrap(); + assert_eq!(args.len(), f.args.len(), "{other} takes {} arguments", f.args.len()); + Term::Call(other.to_string(), std::mem::take(args)) + } + other => panic!("Unknown scalar builtin (and no function registered by that name): {}", other), } } + +/// Whether `name` is a builtin of the scalar language (registered functions may not shadow one). +pub fn is_builtin(name: &str) -> bool { + const NAMES: &[&str] = &["tuple", "list", "inject", "variant", "case", "fold", "proj", "len", "istag", "not", "float", "fneg", "fint", + "fpow", "fpowi", "fmin", "fmax", "feq", "fne", "flt", "fle", "fgt", "fge", "fadd", "fsub", "fmul", "fdiv", "or", "idiv", "append", "if", "hash"]; + NAMES.contains(&name) || F64Fn::ALL.iter().any(|(n, _)| *n == name) +} diff --git a/interactive/tests/corgi_backend.rs b/interactive/tests/corgi_backend.rs index 61cb8bc10..aacd9477e 100644 --- a/interactive/tests/corgi_backend.rs +++ b/interactive/tests/corgi_backend.rs @@ -65,6 +65,7 @@ fn inputs_for(prog: &str) -> Vec> { ])], // f64_math: (key, a, b) with x = a / 4, y = b / 2: negatives, zero, and a key with // two rows; then (key, n) integer exponents for the join, one key unmatched. + "registered" => vec![rows(&[&[1], &[3], &[-4]])], "f64_math" => vec![ rows(&[&[1, 10, 1], &[2, -6, 3], &[3, 0, -1], &[4, 9, 0], &[5, -1, -4], &[5, 7, 5]]), rows(&[&[1, 2], &[2, 3], &[3, -1], &[5, 0]]), @@ -88,6 +89,7 @@ fn assert_backends_agree(prog: &str) { } else { format!("{}/examples/programs/{prog}.ddp", env!("CARGO_MANIFEST_DIR")) }; + register_test_functions(); let src = interactive::load_program(&path); let mut tree = lower::lower_tree(parse::pipe::parse(&src)); tree.optimize(); @@ -142,6 +144,63 @@ fn serializing(n: usize) -> timely::Config { #[test] fn signed_min() { assert_backends_agree("signed_min"); } #[test] fn spread_values() { assert_backends_agree("spread_values"); } #[test] fn f64_math() { assert_backends_agree("f64_math"); } +#[test] fn registered() { + assert_backends_agree("registered"); + // And the program does what it says: the refinement from 1 and 3 reaches every + // number from 1 to 39 (grow stops at 20, so its last children are 38 and 39), and + // -4 has no children. + register_test_functions(); + let src = interactive::load_program(&format!("{}/tests/programs/registered.ddp", env!("CARGO_MANIFEST_DIR"))); + let tree = lower::lower_tree(parse::pipe::parse(&src)); + let out = evaluate(RenderBackend::Corgi, timely::Config::process(2), &tree, &inputs_for("registered")); + let mut cells: Vec = out["cells"].iter().map(|((k, _), _)| match k { Value::Tuple(f) => f[0].as_int(), _ => panic!() }).collect(); + cells.sort(); + let mut want: Vec = (1..40).collect(); + want.insert(0, -4); + assert_eq!(cells, want); + assert_eq!(out["joined"].len(), 40); +} + +/// The functions `registered.ddp` calls. Registering again is idempotent, so every +/// test may do it. +fn register_test_functions() { + use corgi::Shape; + use interactive::ir::{register, Function}; + let int = || Shape::Prim(64); + let float = || Shape::Sum(vec![Shape::Prim(64)]); + // A cell's children: two, while the cell is small and positive; none after. + register(Function { + name: "grow".into(), + args: vec![int()], + result: Shape::List(Box::new(Shape::Prod(vec![int(), int()]))), + body: Box::new(|a| { + let n = a[0].as_int(); + let kids = if (1..20).contains(&n) { vec![2 * n, 2 * n + 1] } else { vec![] }; + Value::List(kids.into_iter().map(|k| Value::Tuple(vec![Value::Int(k), Value::Int(n)])).collect()) + }), + }); + // A float from a pair of ints. + register(Function { + name: "blend".into(), + args: vec![Shape::Prod(vec![int(), int()])], + result: float(), + body: Box::new(|a| { + let Value::Tuple(p) = &a[0] else { panic!("blend expects a pair") }; + Value::f64_value((p[0].as_int() as f64).sqrt() - p[1].as_int() as f64 / 3.0) + }), + }); + // A nested shape: (n * 3, [divisors of n under 5]). + register(Function { + name: "describe".into(), + args: vec![int()], + result: Shape::Prod(vec![Shape::Prod(vec![int()]), Shape::List(Box::new(int()))]), + body: Box::new(|a| { + let n = a[0].as_int(); + let divisors = (1..5).filter(|d| n % d == 0).map(Value::Int).collect(); + Value::Tuple(vec![Value::Tuple(vec![Value::Int(3 * n)]), Value::List(divisors)]) + }), + }); +} /// A filter predicate must be an `Int`: both backends reject a tuple rather than one of them /// keeping nothing. diff --git a/interactive/tests/programs/registered.ddp b/interactive/tests/programs/registered.ddp new file mode 100644 index 000000000..f32c8de57 --- /dev/null +++ b/interactive/tests/programs/registered.ddp @@ -0,0 +1,31 @@ +-- Registered functions: an embedding program supplies Rust functions by name +-- (`ir::register`), and a program calls them like builtins. The test harness +-- registers `grow`, `blend` and `describe` (see `corgi_backend.rs`). Every call +-- runs a row at a time in the corgi backend; the results must match the vec +-- backend's in a map, a filter, a flatmap, a join, and a refinement fixpoint +-- (the shape of a procedural-generation layer: a record's children come from a +-- rule, until no record has any). + +-- Seeds: (n ;). +let seeds = input 0 | key($0[0] ;); + +-- Refine: each cell's children are `grow(cell)`, until `grow` returns none. +refine: { + let kids = cells | flatmap(grow($0[0])) | map($1[1][0] ;); + var cells = seeds + kids | distinct; +} +let cells = refine::cells; + +-- A float from a pair, a nested shape, and a filter on a registered result. +let blended = cells | map($0 ; blend(tuple($0[0], $0[0] + 1))); +let described = cells | map($0 ; describe($0[0])); +let big = described | filter($1[0][0][0] > 40); + +-- A join whose projection calls a registered function on both sides. +let joined = cells | join(described, ($0 ; blend(tuple($0[0], $2[0][0][0])))); + +export "cells" = cells | arrange | inspect(total); +export "blended" = blended | arrange | inspect(total); +export "described" = described | arrange | inspect(total); +export "big" = big | arrange | inspect(total); +export "joined" = joined | arrange | inspect(total); From 9a7d3bc7617046d9c094dc5979bd54623dfe1537 Mon Sep 17 00:00:00 2001 From: Frank McSherry Date: Sat, 26 Sep 2026 09:19:26 -0400 Subject: [PATCH 3/6] Corgi row fallback: move values instead of cloning them transcode/untranscode (the conversions around Kernel::Rows) cloned every value they read. They now take their input by value and move it: transcode_owned builds columns from owned rows, and untranscode moves out of the columns. Worldgen DDIR tour, one worker: 16.7 -> 12.6 s. Co-Authored-By: Claude Opus 5.5 (1M context) --- interactive/src/corgi/container.rs | 4 +- interactive/src/corgi/logic.rs | 111 +++++++++++++++-------------- 2 files changed, 59 insertions(+), 56 deletions(-) diff --git a/interactive/src/corgi/container.rs b/interactive/src/corgi/container.rs index 9a5c830c9..2ae94094b 100644 --- a/interactive/src/corgi/container.rs +++ b/interactive/src/corgi/container.rs @@ -16,7 +16,7 @@ use differential_dataflow::collection::containers::{Enter, Leave, Negate, Result use differential_dataflow::difference::Abelian; use crate::corgi::col_times::{ColTimes, LaneSummary, Lanes}; -use crate::corgi::logic::{transcode, untranscode}; +use crate::corgi::logic::{transcode_owned, untranscode}; use crate::ir::Value as DValue; type Row = DValue; @@ -72,7 +72,7 @@ impl CorgiContainer { times.push(&time); diffs.push(diff); } - CorgiContainer { keys: transcode(&keys_rows, kshape), vals: transcode(&vals_rows, vshape), times, diffs } + CorgiContainer { keys: transcode_owned(keys_rows, kshape), vals: transcode_owned(vals_rows, vshape), times, diffs } } /// Test convenience: build a container from row updates, pinning the shapes from the first diff --git a/interactive/src/corgi/logic.rs b/interactive/src/corgi/logic.rs index 89317dec5..a01d0fd1f 100644 --- a/interactive/src/corgi/logic.rs +++ b/interactive/src/corgi/logic.rs @@ -39,75 +39,69 @@ pub fn shape_of_row(row: &DValue) -> Res { /// AoS rows -> SoA corgi columns, directed by `shape`. A row that does not fit the shape is a /// panic: the shape was pinned from a row of this collection, so a misfit is an ingest error. +/// Borrowing form: one clone of the rows, then [`transcode_owned`]. pub fn transcode(rows: &[DValue], shape: &Shape) -> CValue { + transcode_owned(rows.to_vec(), shape) +} + +/// As [`transcode`], consuming the rows: each field moves into its column, nothing is cloned. +pub fn transcode_owned(rows: Vec, shape: &Shape) -> CValue { match shape { Shape::Prim(_) => CValue::u64(rows.iter().map(|r| r.as_int() as u64).collect()), Shape::Unit => CValue::Unit(rows.len()), - Shape::Prod(fs) => CValue::Prod( - fs.iter() - .enumerate() - .map(|(i, fsi)| { - let sub: Vec = rows - .iter() - .map(|r| match r { - DValue::Tuple(xs) => xs[i].clone(), - other => panic!("transcode: expected Tuple, got {other:?}"), - }) - .collect(); - transcode(&sub, fsi) - }) - .collect(), - ), + Shape::Prod(fs) => { + // Transpose by moving: one Vec per field, filled by draining each tuple. + let mut cols: Vec> = fs.iter().map(|_| Vec::with_capacity(rows.len())).collect(); + for r in rows { + match r { + DValue::Tuple(xs) => { + assert!(xs.len() >= fs.len(), "transcode: a {}-tuple for a {}-field shape", xs.len(), fs.len()); + for (col, x) in cols.iter_mut().zip(xs) { col.push(x) } + } + other => panic!("transcode: expected Tuple, got {other:?}"), + } + } + CValue::Prod(cols.into_iter().zip(fs).map(|(c, fsi)| transcode_owned(c, fsi)).collect()) + } Shape::List(elem) => { // List column = per-row END offsets + a flattened element column. let mut ends = Vec::with_capacity(rows.len()); let mut flat: Vec = Vec::new(); - let mut acc = 0usize; for r in rows { match r { DValue::List(xs) => { - acc += xs.len(); - ends.push(acc); - flat.extend(xs.iter().cloned()); + flat.extend(xs); + ends.push(flat.len()); } other => panic!("transcode: expected List, got {other:?}"), } } - CValue::List(ends.into(), Box::new(transcode(&flat, elem))) + CValue::List(ends.into(), Box::new(transcode_owned(flat, elem))) } Shape::Sum(lanes) => { // Per-row tag, plus one packed lane per variant (its arm's rows in row order; a // variant no row uses is an empty column of its declared shape). - let tags: Vec = rows - .iter() - .map(|r| match r { - DValue::Variant(t, _) => *t as usize, + let mut tags: Vec = Vec::with_capacity(rows.len()); + let mut payloads: Vec> = lanes.iter().map(|_| Vec::new()).collect(); + for r in rows { + match r { + DValue::Variant(t, p) => { + let t = t as usize; + if t >= lanes.len() { panic!("transcode: tag {t} is outside the declared {}-variant sum", lanes.len()) } + tags.push(t); + payloads[t].push(*p); + } other => panic!("transcode: expected Variant, got {other:?}"), - }) - .collect(); - if let Some(t) = tags.iter().find(|&&t| t >= lanes.len()) { - panic!("transcode: tag {t} is outside the declared {}-variant sum", lanes.len()); + } } - let lane_vals: Vec = lanes - .iter() - .enumerate() - .map(|(tag, lshape)| { - let payloads: Vec = rows - .iter() - .filter_map(|r| match r { - DValue::Variant(t, p) if *t as usize == tag => Some((**p).clone()), - _ => None, - }) - .collect(); - transcode(&payloads, lshape) - }) - .collect(); + let lane_vals = payloads.into_iter().zip(lanes).map(|(p, lshape)| transcode_owned(p, lshape)).collect(); CValue::sum(tags, lane_vals) } } } -/// SoA corgi columns -> AoS rows, directed by `shape`. Inverse of [`transcode`]. +/// SoA corgi columns -> AoS rows, directed by `shape`. Inverse of [`transcode`]. Each value +/// moves into its row; nothing is cloned. pub fn untranscode(col: CValue, shape: &Shape) -> Vec { match shape { Shape::Prim(_) => col.into_u64("untranscode").unwrap().into_iter().map(|x| DValue::Int(x as i64)).collect(), @@ -115,38 +109,47 @@ pub fn untranscode(col: CValue, shape: &Shape) -> Vec { Shape::Prod(fs) => { let cols = col.into_prod("untranscode").unwrap(); let n = if cols.is_empty() { 0 } else { cols[0].len() }; - let per_field: Vec> = - cols.into_iter().zip(fs.iter()).map(|(c, fsi)| untranscode(c, fsi)).collect(); - (0..n).map(|i| DValue::Tuple(per_field.iter().map(|f| f[i].clone()).collect())).collect() + let mut per_field: Vec> = + cols.into_iter().zip(fs.iter()).map(|(c, fsi)| untranscode(c, fsi).into_iter()).collect(); + (0..n).map(|_| DValue::Tuple(per_field.iter_mut().map(|f| f.next().expect("field columns agree in length")).collect())).collect() } Shape::List(elem) => { // Inverse of transcode's List: per-row END offsets + a flattened element column → one - // `List` per row, slicing the untranscoded flat column by each row's span. + // `List` per row, taking each row's span off the front of the untranscoded column. let (bounds, vals) = match col { CValue::List(b, vals) => (b, *vals), other => panic!("untranscode: expected List, got {other:?}"), }; - let flat = untranscode(vals, elem); + let mut flat = untranscode(vals, elem).into_iter(); let ends: Vec = bounds.to_vec(); let mut out = Vec::with_capacity(ends.len()); let mut start = 0usize; for end in ends { - out.push(DValue::List(flat[start..end].to_vec())); + out.push(DValue::List(flat.by_ref().take(end - start).collect())); start = end; } out } Shape::Sum(lanes) => { - // Inverse of transcode's Sum: untranscode each lane, then for each row pull its payload + // Inverse of transcode's Sum: untranscode each lane, then for each row take its payload // from its lane at the recorded within-lane OFFSET (robust to row reordering from a - // prior gather/merge — not a sequential cursor). + // prior gather/merge — not a sequential cursor). Rows may share a payload, so each is + // cloned except at its last use, where it moves. let (tags, variant_vals) = col.into_sum("untranscode").unwrap(); - let lane_rows: Vec> = + let mut lane_rows: Vec> = variant_vals.into_iter().zip(lanes.iter()).map(|(v, ls)| untranscode(v, ls)).collect(); + let mut uses: Vec> = lane_rows.iter().map(|l| vec![0; l.len()]).collect(); + for r in 0..tags.len() { uses[tags.tag_at(r)][tags.offset_at(r)] += 1 } (0..tags.len()) .map(|r| { let (tag, off) = (tags.tag_at(r), tags.offset_at(r)); - DValue::Variant(tag as u32, Box::new(lane_rows[tag][off].clone())) + uses[tag][off] -= 1; + let payload = if uses[tag][off] == 0 { + std::mem::replace(&mut lane_rows[tag][off], DValue::Int(0)) + } else { + lane_rows[tag][off].clone() + }; + DValue::Variant(tag as u32, Box::new(payload)) }) .collect() } @@ -777,7 +780,7 @@ impl Kernel { } }) .collect(); - transcode(&results, output) + transcode_owned(results, output) } } } From 4acb9d371583fbbf7d563c00ce3b9a7df77c925e Mon Sep 17 00:00:00 2001 From: Frank McSherry Date: Sat, 26 Sep 2026 09:44:53 -0400 Subject: [PATCH 4/6] Registered functions lower to corgi host kernels Term::Call compiles to NumOp::Host (corgi #38) over the tuple of its arguments, so the term around a call stays columnar instead of falling back to Kernel::Rows. A function's columnar body comes from ir::register_kernel; otherwise its row body runs behind RowKernel, which converts only the call's arguments and result. kernel_of keeps one Arc per function, so every call site shares one kernel and CSE can merge equal calls. Argument shapes are checked when the program is typed at install, by the host op. Pins corgi at the merge of #38. Worldgen DDIR tour at 4 workers: 8.6 -> 5.0 s with its six hottest rules as columnar kernels; a 1 km move 263 -> 60 ms. Co-Authored-By: Claude Opus 5.5 (1M context) --- interactive/Cargo.toml | 2 +- interactive/src/corgi/logic.rs | 12 +++++++--- interactive/src/ir.rs | 40 ++++++++++++++++++++++++++++++++++ 3 files changed, 50 insertions(+), 4 deletions(-) diff --git a/interactive/Cargo.toml b/interactive/Cargo.toml index fa2d72a96..cbac3f3c7 100644 --- a/interactive/Cargo.toml +++ b/interactive/Cargo.toml @@ -14,7 +14,7 @@ workspace = true [dependencies] columnar = { workspace = true } # The columnar kernels for the interpreted backend, pinned by git rev. -corgi = { git = "https://github.com/frankmcsherry/WIP", rev = "be003988f17e3998a8abba3e168a2520a46d10e4", features = ["serde"] } +corgi = { git = "https://github.com/frankmcsherry/WIP", rev = "657fc963cccde58ff4b8fcad9fac1331bb2c5b1c", features = ["serde"] } differential-dataflow = { workspace = true } serde = { version = "1.0", features = ["derive"] } smallvec = "1.15.1" diff --git a/interactive/src/corgi/logic.rs b/interactive/src/corgi/logic.rs index a01d0fd1f..5c9129aa6 100644 --- a/interactive/src/corgi/logic.rs +++ b/interactive/src/corgi/logic.rs @@ -252,7 +252,13 @@ pub fn compile( expected: Option<&Shape>, ) -> Res { match term { - Term::Call(name, _) => Err(format!("`{name}` is a registered function; it runs a row at a time (see `row_only`)")), + Term::Call(name, args) => { + // One host-kernel node over the tuple of its arguments; the kernel checks their shapes. + let kernel = crate::ir::kernel_of(name).ok_or_else(|| format!("`{name}` is not a registered function"))?; + let fields = args.iter().map(|a| compile(a, b, env, env_shapes, anchor, None)).collect::>>()?; + let tuple = b.tuple(fields); + Ok(b.add(NumOp::Host(corgi::HostOp(kernel)), vec![tuple])) + } Term::Var(i) => env.get(*i).copied().ok_or_else(|| format!("`${i}` is not in scope here")), Term::Bound(k) => { env.len().checked_sub(1 + *k).map(|i| env[i]).ok_or_else(|| format!("binder `^{k}` is not in scope here")) @@ -792,8 +798,8 @@ pub fn row_only(t: &Term) -> bool { match t { Term::Var(_) | Term::Bound(_) | Term::Int(_) => false, Term::Unary(UnOp::F64Fn(f), _) if *f != F64Fn::Abs => true, - Term::Unary(UnOp::F64ToInt, _) | Term::Binary(BinOp::F64Pow | BinOp::F64PowI, _, _) | Term::Call(..) => true, - Term::Tuple(fs) | Term::List(fs) | Term::Hash(fs) => fs.iter().any(row_only), + Term::Unary(UnOp::F64ToInt, _) | Term::Binary(BinOp::F64Pow | BinOp::F64PowI, _, _) => true, + Term::Tuple(fs) | Term::List(fs) | Term::Hash(fs) | Term::Call(_, fs) => fs.iter().any(row_only), Term::Spread(inner) | Term::Proj(inner, _) | Term::Unary(_, inner) => row_only(inner), Term::Inject { tag, payload, .. } => row_only(tag) || row_only(payload), Term::Case { scrutinee, arms, default } => { diff --git a/interactive/src/ir.rs b/interactive/src/ir.rs index 4b2b23ab3..2ce1ed517 100644 --- a/interactive/src/ir.rs +++ b/interactive/src/ir.rs @@ -506,6 +506,46 @@ pub fn register(f: Function) { functions().write().unwrap().insert(f.name.clone(), std::sync::Arc::new(f)); } +fn kernels() -> &'static std::sync::RwLock>> { + static KERNELS: std::sync::OnceLock>>> = std::sync::OnceLock::new(); + KERNELS.get_or_init(Default::default) +} + +/// Give the registered function `name` a columnar body: the corgi backend calls `kernel` on whole +/// columns instead of `body` a row at a time. Its declared shapes must be the function's. +pub fn register_kernel(name: &str, kernel: std::sync::Arc) { + let f = lookup(name).unwrap_or_else(|| panic!("register_kernel: `{name}` is not a registered function")); + assert_eq!(kernel.input(), &corgi::Shape::Prod(f.args.clone()), "register_kernel: `{name}` input shape"); + assert_eq!(kernel.output(), &f.result, "register_kernel: `{name}` output shape"); + kernels().write().unwrap().insert(name.to_string(), kernel); +} + +/// The columnar body of `name`: its registered kernel, or its row body behind an adapter that +/// converts only the call's arguments and result. +pub fn kernel_of(name: &str) -> Option> { + if let Some(k) = kernels().read().unwrap().get(name) { return Some(k.clone()) } + let f = lookup(name)?; + let k: std::sync::Arc = std::sync::Arc::new(RowKernel { input: corgi::Shape::Prod(f.args.clone()), f }); + // Keep it, so every call site shares one kernel (and CSE can merge equal calls). + Some(kernels().write().unwrap().entry(name.to_string()).or_insert(k).clone()) +} + +/// A row-at-a-time function as a host kernel. +struct RowKernel { f: std::sync::Arc, input: corgi::Shape } +impl corgi::HostKernel for RowKernel { + fn name(&self) -> &str { &self.f.name } + fn input(&self) -> &corgi::Shape { &self.input } + fn output(&self) -> &corgi::Shape { &self.f.result } + fn eval(&self, input: corgi::Value) -> Result { + let rows = crate::corgi::logic::untranscode(input, &self.input); + let out: Vec = rows.into_iter().map(|r| match r { + Value::Tuple(args) => (self.f.body)(&args), + other => (self.f.body)(std::slice::from_ref(&other)), + }).collect(); + Ok(crate::corgi::logic::transcode_owned(out, &self.f.result)) + } +} + /// The registered function called `name`, if any. pub fn lookup(name: &str) -> Option> { functions().read().unwrap().get(name).cloned() From 0e9b95246f5d091173dcc3a1a74a8f637a4664b7 Mon Sep 17 00:00:00 2001 From: Frank McSherry Date: Sun, 27 Sep 2026 09:16:09 -0400 Subject: [PATCH 5/6] Registered calls: no-argument calls pass Unit; test a columnar kernel; review tidy - A call with no arguments passes its kernel a Unit over the anchor (ir::call_input): an empty product carries no row count, and corgi (#38) rejects it as a kernel input. register_kernel and RowKernel use the same input shape. - registered.ddp adds a zero-argument call (seven) and a function with a columnar kernel (double, registered with register_kernel); the gate checks both against the vec backend at 1-4 workers and over serializing channels. - Arc::clone for kernel handles; float-lowering locals renamed so no binding shadows an unrelated one. Clippy warnings for interactive: 165 -> 164. Co-Authored-By: Claude Opus 5.5 (1M context) --- interactive/src/corgi/logic.rs | 49 +++++++++++++---------- interactive/src/ir.rs | 14 +++++-- interactive/tests/corgi_backend.rs | 27 +++++++++++++ interactive/tests/programs/registered.ddp | 16 +++++--- 4 files changed, 75 insertions(+), 31 deletions(-) diff --git a/interactive/src/corgi/logic.rs b/interactive/src/corgi/logic.rs index 5c9129aa6..5206f823b 100644 --- a/interactive/src/corgi/logic.rs +++ b/interactive/src/corgi/logic.rs @@ -255,9 +255,14 @@ pub fn compile( Term::Call(name, args) => { // One host-kernel node over the tuple of its arguments; the kernel checks their shapes. let kernel = crate::ir::kernel_of(name).ok_or_else(|| format!("`{name}` is not a registered function"))?; - let fields = args.iter().map(|a| compile(a, b, env, env_shapes, anchor, None)).collect::>>()?; - let tuple = b.tuple(fields); - Ok(b.add(NumOp::Host(corgi::HostOp(kernel)), vec![tuple])) + // A call with no arguments passes a `Unit` over the anchor (`ir::call_input`). + let input = if args.is_empty() { + b.add(Op::Unit, vec![anchor]) + } else { + let fields = args.iter().map(|a| compile(a, b, env, env_shapes, anchor, None)).collect::>>()?; + b.tuple(fields) + }; + Ok(b.add(NumOp::Host(corgi::HostOp(kernel)), vec![input])) } Term::Var(i) => env.get(*i).copied().ok_or_else(|| format!("`${i}` is not in scope here")), Term::Bound(k) => { @@ -322,8 +327,8 @@ pub fn compile( } BinOp::Append => { let p = pair(b, lid, rid); b.add(Op::Append, vec![p]) } BinOp::F64Add | BinOp::F64Sub | BinOp::F64Mul | BinOp::F64Div => { - let expected = Shape::Sum(vec![Shape::Prim(64)]); - if shape_of_term(l, env_shapes, None)? != expected || shape_of_term(r, env_shapes, None)? != expected { + let f64_shape = Shape::Sum(vec![Shape::Prim(64)]); + if shape_of_term(l, env_shapes, None)? != f64_shape || shape_of_term(r, env_shapes, None)? != f64_shape { return Err("floating arithmetic expects two F64 newtypes; use float(int)".into()); } let l = b.add(Op::Unwrap, vec![lid]); @@ -337,8 +342,8 @@ pub fn compile( b.add(Op::Inject(0, vec![Shape::Prim(64)]), vec![payload]) } BinOp::F64Min | BinOp::F64Max | BinOp::F64Eq | BinOp::F64Ne | BinOp::F64Lt | BinOp::F64Le | BinOp::F64Gt | BinOp::F64Ge => { - let expected = Shape::Sum(vec![Shape::Prim(64)]); - if shape_of_term(l, env_shapes, None)? != expected || shape_of_term(r, env_shapes, None)? != expected { + let f64_shape = Shape::Sum(vec![Shape::Prim(64)]); + if shape_of_term(l, env_shapes, None)? != f64_shape || shape_of_term(r, env_shapes, None)? != f64_shape { return Err(format!("{op:?} expects two F64 newtypes; use float(int)")); } let (x, y) = (float_leaf(b, lid), float_leaf(b, rid)); @@ -353,8 +358,8 @@ pub fn compile( let choices = b.tuple(vec![y_nan, x, total]); let unless_y = b.add(Op::Select, vec![choices]); let x_nan = is_nan(b, x); - let choices = b.tuple(vec![x_nan, y, unless_y]); - let f = b.add(Op::Select, vec![choices]); + let choices_x = b.tuple(vec![x_nan, y, unless_y]); + let f = b.add(Op::Select, vec![choices_x]); let payload = b.add(ArithOp::ToSigned, vec![f]); b.add(Op::Inject(0, vec![Shape::Prim(64)]), vec![payload]) } @@ -371,16 +376,16 @@ pub fn compile( }; let rel = b.add(CmpOp::Rel(pred), vec![p]); let (x_nan, y_nan) = (is_nan(b, x), is_nan(b, y)); - let p = pair(b, x_nan, y_nan); - let any_nan = b.add(CmpOp::Max, vec![p]); + let nans = pair(b, x_nan, y_nan); + let any_nan = b.add(CmpOp::Max, vec![nans]); let zero = b.add(Op::Lit(CValue::u64(vec![0])), vec![anchor]); - let p = pair(b, any_nan, zero); - let ordered = b.add(CmpOp::Rel(Pred::Eq), vec![p]); - let p = pair(b, rel, ordered); - let holds = b.add(CmpOp::Min, vec![p]); + let no_nan = pair(b, any_nan, zero); + let ordered = b.add(CmpOp::Rel(Pred::Eq), vec![no_nan]); + let both = pair(b, rel, ordered); + let holds = b.add(CmpOp::Min, vec![both]); if matches!(op, BinOp::F64Ne) { - let p = pair(b, holds, zero); - b.add(CmpOp::Rel(Pred::Eq), vec![p]) + let negate = pair(b, holds, zero); + b.add(CmpOp::Rel(Pred::Eq), vec![negate]) } else { holds } @@ -723,10 +728,10 @@ fn is_nan(b: &mut Builder, x: usize) -> usize { /// Replace `-0.0` by `0.0` in a float leaf. Their keys are adjacent, so this adds the mask. fn fold_negative_zero(b: &mut Builder, x: usize, anchor: usize) -> usize { let negative_zero = b.add(Op::Lit(CValue::u64(vec![float_key(-0.0)])), vec![anchor]); - let p = b.tuple(vec![x, negative_zero]); - let is_negative_zero = b.add(CmpOp::Rel(Pred::Eq), vec![p]); - let p = b.tuple(vec![x, is_negative_zero]); - b.add(ArithOp::Bin(CBinOp::Add, Kind::U, 64), vec![p]) + let probe = b.tuple(vec![x, negative_zero]); + let is_negative_zero = b.add(CmpOp::Rel(Pred::Eq), vec![probe]); + let sum = b.tuple(vec![x, is_negative_zero]); + b.add(ArithOp::Bin(CBinOp::Add, Kind::U, 64), vec![sum]) } /// Compile a `Fold` step into a closed corgi sub-graph. Without capture the body's input is @@ -815,7 +820,7 @@ pub fn row_only(t: &Term) -> bool { /// yields the same result shape: F64 -> F64 by `fneg`, `fint` by a compare of two `fneg`s, /// `fpow` by `fadd`, and `fpowi(x, n)` by `fadd(x, float(n))`. Only its typing is used. fn columnar_stand_in(t: &Term) -> Term { - let bx = |t: &Term| Box::new(columnar_stand_in(t)); + let bx = |term: &Term| Box::new(columnar_stand_in(term)); match t { Term::Var(_) | Term::Bound(_) | Term::Int(_) => t.clone(), Term::Unary(UnOp::F64Fn(f), x) if *f != F64Fn::Abs => Term::Unary(UnOp::F64Neg, bx(x)), diff --git a/interactive/src/ir.rs b/interactive/src/ir.rs index 2ce1ed517..9d3fe49fc 100644 --- a/interactive/src/ir.rs +++ b/interactive/src/ir.rs @@ -515,7 +515,7 @@ fn kernels() -> &'static std::sync::RwLock) { let f = lookup(name).unwrap_or_else(|| panic!("register_kernel: `{name}` is not a registered function")); - assert_eq!(kernel.input(), &corgi::Shape::Prod(f.args.clone()), "register_kernel: `{name}` input shape"); + assert_eq!(kernel.input(), &call_input(&f.args), "register_kernel: `{name}` input shape"); assert_eq!(kernel.output(), &f.result, "register_kernel: `{name}` output shape"); kernels().write().unwrap().insert(name.to_string(), kernel); } @@ -523,11 +523,17 @@ pub fn register_kernel(name: &str, kernel: std::sync::Arc /// The columnar body of `name`: its registered kernel, or its row body behind an adapter that /// converts only the call's arguments and result. pub fn kernel_of(name: &str) -> Option> { - if let Some(k) = kernels().read().unwrap().get(name) { return Some(k.clone()) } + if let Some(k) = kernels().read().unwrap().get(name) { return Some(std::sync::Arc::clone(k)) } let f = lookup(name)?; - let k: std::sync::Arc = std::sync::Arc::new(RowKernel { input: corgi::Shape::Prod(f.args.clone()), f }); + let k: std::sync::Arc = std::sync::Arc::new(RowKernel { input: call_input(&f.args), f }); // Keep it, so every call site shares one kernel (and CSE can merge equal calls). - Some(kernels().write().unwrap().entry(name.to_string()).or_insert(k).clone()) + Some(std::sync::Arc::clone(kernels().write().unwrap().entry(name.to_string()).or_insert(k))) +} + +/// The column a call passes its kernel: the tuple of its arguments, or `Unit` for a call with none +/// (an empty product carries no row count, and corgi rejects it as a kernel input). +pub fn call_input(args: &[corgi::Shape]) -> corgi::Shape { + if args.is_empty() { corgi::Shape::Unit } else { corgi::Shape::Prod(args.to_vec()) } } /// A row-at-a-time function as a host kernel. diff --git a/interactive/tests/corgi_backend.rs b/interactive/tests/corgi_backend.rs index aacd9477e..3bafd8ddb 100644 --- a/interactive/tests/corgi_backend.rs +++ b/interactive/tests/corgi_backend.rs @@ -200,6 +200,33 @@ fn register_test_functions() { Value::Tuple(vec![Value::Tuple(vec![Value::Int(3 * n)]), Value::List(divisors)]) }), }); + // No arguments: the call passes a `Unit` column, so it still runs once per row. + register(Function { + name: "seven".into(), + args: vec![], + result: int(), + body: Box::new(|_| Value::Int(7)), + }); + // A columnar body: the corgi backend runs `Double` on whole columns, the vec backend the row + // body; the gate checks they agree. + register(Function { + name: "double".into(), + args: vec![int()], + result: int(), + body: Box::new(|a| Value::Int(2 * a[0].as_int())), + }); + struct Double(Shape, Shape); + impl corgi::HostKernel for Double { + fn name(&self) -> &str { "double" } + fn input(&self) -> &Shape { &self.0 } + fn output(&self) -> &Shape { &self.1 } + fn eval(&self, input: corgi::Value) -> Result { + let args = input.into_prod("double")?; + let xs = args[0].as_u64("double")?; + Ok(corgi::Value::u64(xs.iter().map(|&x| (2 * x as i64) as u64).collect())) + } + } + interactive::ir::register_kernel("double", std::sync::Arc::new(Double(Shape::Prod(vec![int()]), int()))); } /// A filter predicate must be an `Int`: both backends reject a tuple rather than one of them diff --git a/interactive/tests/programs/registered.ddp b/interactive/tests/programs/registered.ddp index f32c8de57..b5e90600b 100644 --- a/interactive/tests/programs/registered.ddp +++ b/interactive/tests/programs/registered.ddp @@ -1,10 +1,10 @@ -- Registered functions: an embedding program supplies Rust functions by name -- (`ir::register`), and a program calls them like builtins. The test harness --- registers `grow`, `blend` and `describe` (see `corgi_backend.rs`). Every call --- runs a row at a time in the corgi backend; the results must match the vec --- backend's in a map, a filter, a flatmap, a join, and a refinement fixpoint --- (the shape of a procedural-generation layer: a record's children come from a --- rule, until no record has any). +-- registers `grow`, `blend`, `describe`, `seven` (no arguments) and `double` +-- (with a columnar kernel; see `corgi_backend.rs`). The corgi backend runs calls +-- as host kernels; the results must match the vec backend's in a map, a filter, +-- a flatmap, a join, and a refinement fixpoint (the shape of a procedural- +-- generation layer: a record's children come from a rule, until none has any). -- Seeds: (n ;). let seeds = input 0 | key($0[0] ;); @@ -21,6 +21,10 @@ let blended = cells | map($0 ; blend(tuple($0[0], $0[0] + 1))); let described = cells | map($0 ; describe($0[0])); let big = described | filter($1[0][0][0] > 40); +-- A call with no arguments, and one with a columnar kernel. +let sevens = cells | map($0 ; seven() + $0[0]); +let doubled = cells | map($0 ; double($0[0])); + -- A join whose projection calls a registered function on both sides. let joined = cells | join(described, ($0 ; blend(tuple($0[0], $2[0][0][0])))); @@ -29,3 +33,5 @@ export "blended" = blended | arrange | inspect(total); export "described" = described | arrange | inspect(total); export "big" = big | arrange | inspect(total); export "joined" = joined | arrange | inspect(total); +export "sevens" = sevens | arrange | inspect(total); +export "doubled" = doubled | arrange | inspect(total); From b1ec0e0eb867a3d8452f0952ed42936ecbbbbcf2 Mon Sep 17 00:00:00 2001 From: Frank McSherry Date: Sun, 27 Sep 2026 10:20:29 -0400 Subject: [PATCH 6/6] Float functions as host kernels; drop the row fallback; review fixes The F64 functions that are not compositions of corgi ops (fsqrt, fexp, fln, ffloor, fceil, fround, fsin, fcos, ftan, fint, fpow, fpowi) are now built-in host kernels over the float column: decode each total-order key, apply the function as ir::eval does, encode. One kernel per function (CSE merges equal calls). With no row-only op left, Kernel, row_only, columnar_stand_in, literal_of_shape and Kernel::Rows go; lowering returns a Graph again, and the backend and join call eval_graph as on master-next. Review fixes: - register drops any kernel of the name being replaced (the row adapter held the old body; a register_kernel kernel could have the old shapes). - The row backend checks each call's arguments against the declared shapes (Value::has_shape), so a mis-shaped call fails on both backends. - Keywords (min, map, key, count, ...) cannot be registered: the lexer's table is one function, keyword(), and is_builtin consults it. - transcode asserts a tuple has exactly its shape's fields (was >=). - RowKernel's unreachable arm is gone; Term::Call and the float docs updated. Tests: every F64 op agrees with ir::eval on all special-value pairs (NaNs of both signs); a composite term over non-NaN pairs (corgi's float Add returns the positive NaN where ir::eval keeps an operand's sign, as fadd alone did before); re-registration; argument shapes on both backends. Co-Authored-By: Claude Opus 5.5 (1M context) --- interactive/src/backend/corgi.rs | 24 ++- interactive/src/corgi/join.rs | 2 +- interactive/src/corgi/logic.rs | 270 ++++++++++++----------------- interactive/src/ir.rs | 23 ++- interactive/src/parse/mod.rs | 2 +- interactive/src/parse/pipe.rs | 40 +++-- interactive/tests/corgi_backend.rs | 38 +++- 7 files changed, 206 insertions(+), 193 deletions(-) diff --git a/interactive/src/backend/corgi.rs b/interactive/src/backend/corgi.rs index 8eadefe32..48b065e7a 100644 --- a/interactive/src/backend/corgi.rs +++ b/interactive/src/backend/corgi.rs @@ -1,7 +1,6 @@ //! The corgi rendering substrate: corgi columns are the native representation on dataflow edges, //! arrangements are chains of sorted columnar chunks (`ChunkSpine`, cursor-less), and -//! scalar logic runs columnar via `eval_graph` (a term using an F64 function corgi has no kernel -//! for runs a row at a time instead; see `logic::Kernel`). The row-wise `backend::vec` remains useful for +//! scalar logic runs columnar via `eval_graph`. The row-wise `backend::vec` remains useful for //! comparison, but its representation choices do not define corgi's physical semantics. //! //! All `Backend` methods are corgi-native: `linear` folds a `LinearOp` chain over each container @@ -31,8 +30,8 @@ use crate::corgi::exchange::CorgiPact; use crate::corgi::join::CorgiJoinBackend; use crate::corgi::reduce::CorgiReduceBackend; use differential_dataflow::operators::int_proxy::{ProxyJoinTactic, ProxyReduceTactic}; -use crate::corgi::logic::{compile_flatmap, compile_predicate, compile_projection, compile_scalar, shape_of_row, Kernel}; -use corgi::Shape; +use crate::corgi::logic::{compile_flatmap, compile_predicate, compile_projection, compile_scalar, shape_of_row}; +use corgi::{Graph, NumOp, Shape}; use crate::ir::{Diff, LinearOp, Projection, Reducer, Time, Value as DValue}; use crate::scope_ir as st; @@ -48,13 +47,13 @@ type CTrace = differential_dataflow::trace::chunk::ChunkSpine, + compiled: Option<(Shape, Shape, Graph)>, } impl Plan { /// The graph for a container of these shapes, compiling on first use. A type error is a /// panic with corgi's message: a program that typechecks never reaches it. - fn graph(&mut self, what: &str, kshape: Shape, vshape: Shape, compile: impl FnOnce(&Shape, &Shape) -> Result) -> &Kernel { + fn graph(&mut self, what: &str, kshape: Shape, vshape: Shape, compile: impl FnOnce(&Shape, &Shape) -> Result, String>) -> &Graph { if self.compiled.is_none() { let g = compile(&kshape, &vshape).unwrap_or_else(|e| panic!("{what}: type error at shapes ({kshape}, {vshape}): {e}")); self.compiled = Some((kshape.clone(), vshape.clone(), g)); @@ -67,9 +66,8 @@ impl Plan { /// Apply a `LinearOp` chain to one corgi container (the corgi-native compute per batch). /// Project = corgi `eval_graph`; Filter = corgi mask + `gather`; FlatMap = `eval_graph` to a list -/// column + a structural explode; Negate = Rust. Every term is columnar except one that uses an -/// F64 function corgi has no kernel for (`logic::row_only`), which `ir::eval` runs a row at a -/// time; either way each op's kernel is compiled once (`plans`). An empty batch +/// column + a structural explode; Negate = Rust. Every term is columnar — there is no row-wise +/// path inside the dataflow — and each op's graph is compiled once (`plans`). An empty batch /// passes through untouched: it carries no shape to compile against and no rows to compute. /// The two data<->time ops are columnar and total: EnterAt reads its delay field as a column and /// joins it into `times` in place; LiftIter reads the iteration coordinate out of `times` and @@ -89,14 +87,14 @@ fn apply_ops(mut c: CC, ops: &[LinearOp], level: usize, plans: &mut [Plan]) -> C c = match op { LinearOp::Project(p) => { let g = plan.graph("map", kshape, vshape, |k, v| compile_projection(&p.key, &p.val, k, v)); - let mut cols = g.eval(CValue::Prod(vec![c.keys, c.vals])).into_prod("linear project").unwrap(); + let mut cols = corgi::eval_graph(g, CValue::Prod(vec![c.keys, c.vals])).into_prod("linear project").unwrap(); let vals = cols.pop().unwrap(); let keys = cols.pop().unwrap(); CorgiContainer { keys, vals, times: c.times, diffs: c.diffs } } LinearOp::Filter(cond) => { let g = plan.graph("filter", kshape, vshape, |k, v| compile_predicate(cond, k, v)); - let mask = g.eval(CValue::Prod(vec![c.keys.clone(), c.vals.clone()])).into_u64("filter mask").unwrap(); + let mask = corgi::eval_graph(g, CValue::Prod(vec![c.keys.clone(), c.vals.clone()])).into_u64("filter mask").unwrap(); let keep: Vec = (0..mask.len()).filter(|&i| mask[i] != 0).collect(); let keys = gather(&c.keys, &keep); let vals = gather(&c.vals, &keep); @@ -118,7 +116,7 @@ fn apply_ops(mut c: CC, ops: &[LinearOp], level: usize, plans: &mut [Plan]) -> C // `level-1` and identity everywhere else (u64's minimum is 0), so the delta // never has to be built. The epoch is lane 0, so PointStamp index `level-1` is // lane `level`, and the join is one lane-wise max. - let raw = g.eval(CValue::Prod(vec![c.keys.clone(), c.vals.clone()])) + let raw = corgi::eval_graph(g, CValue::Prod(vec![c.keys.clone(), c.vals.clone()])) .into_u64("enter_at delay") .unwrap(); let delays: Vec = raw.iter().map(|r| 256 * (64 - r.leading_zeros() as u64)).collect(); @@ -152,7 +150,7 @@ fn apply_ops(mut c: CC, ops: &[LinearOp], level: usize, plans: &mut [Plan]) -> C // bounds gives both the within-row position (DDIR's `$1[0]`) and a repeat map // carrying key/time/diff across. No per-row eval, no transcode. let (bounds, elems) = - g.eval(CValue::Prod(vec![c.keys.clone(), c.vals])).into_list("flatmap list").unwrap(); + corgi::eval_graph(g, CValue::Prod(vec![c.keys.clone(), c.vals])).into_list("flatmap list").unwrap(); let ends: Vec = bounds.to_vec(); let total = ends.last().copied().unwrap_or(0); let (mut reps, mut pos) = (Vec::with_capacity(total), Vec::with_capacity(total)); diff --git a/interactive/src/corgi/join.rs b/interactive/src/corgi/join.rs index 1fb24f2c0..fdaa222a7 100644 --- a/interactive/src/corgi/join.rs +++ b/interactive/src/corgi/join.rs @@ -201,7 +201,7 @@ impl ProxyJoinBackend, CBatch> for CorgiJoinBackend< } else { gather_lanes(&vals1, &tag1, &off1) }; let proj = compile_join_projection(&self.key, &self.val, &shape_of_value(&kc), &shape_of_value(&v0), &shape_of_value(&v1)) .unwrap_or_else(|e| panic!("join projection: type error: {e}")); - let projected = proj.eval(CValue::Prod(vec![kc, v0, v1])); + let projected = corgi::eval_graph(&proj, CValue::Prod(vec![kc, v0, v1])); let mut cols = projected.into_prod("corgi join projection").unwrap(); let nv = cols.pop().unwrap(); let nk = cols.pop().unwrap(); diff --git a/interactive/src/corgi/logic.rs b/interactive/src/corgi/logic.rs index 5206f823b..affb513f9 100644 --- a/interactive/src/corgi/logic.rs +++ b/interactive/src/corgi/logic.rs @@ -55,7 +55,7 @@ pub fn transcode_owned(rows: Vec, shape: &Shape) -> CValue { for r in rows { match r { DValue::Tuple(xs) => { - assert!(xs.len() >= fs.len(), "transcode: a {}-tuple for a {}-field shape", xs.len(), fs.len()); + assert!(xs.len() == fs.len(), "transcode: a {}-tuple for a {}-field shape", xs.len(), fs.len()); for (col, x) in cols.iter_mut().zip(xs) { col.push(x) } } other => panic!("transcode: expected Tuple, got {other:?}"), @@ -392,7 +392,23 @@ pub fn compile( } } } - BinOp::F64Pow | BinOp::F64PowI => return Err(format!("{op:?} has no columnar kernel; it runs a row at a time (see `row_only`)")), + BinOp::F64Pow | BinOp::F64PowI => { + let f64_shape = Shape::Sum(vec![Shape::Prim(64)]); + let exponent = if matches!(op, BinOp::F64Pow) { f64_shape.clone() } else { Shape::Prim(64) }; + if shape_of_term(l, env_shapes, None)? != f64_shape || shape_of_term(r, env_shapes, None)? != exponent { + return Err(format!("{op:?} expects an F64 newtype and an {}", if matches!(op, BinOp::F64Pow) { "F64" } else { "Int" })); + } + let x = float_leaf(b, lid); + let (kind, y) = if matches!(op, BinOp::F64Pow) { + (float_kernels::FloatOp::Pow, float_leaf(b, rid)) + } else { + (float_kernels::FloatOp::PowI, rid) + }; + let args = b.tuple(vec![x, y]); + let z = b.add(NumOp::Host(float_kernels::op(kind)), vec![args]); + let payload = b.add(ArithOp::ToSigned, vec![z]); + b.add(Op::Inject(0, vec![Shape::Prim(64)]), vec![payload]) + } BinOp::Eq | BinOp::Ne => { // Cross-shape structural compare folds to a constant (Eq→0, Ne→1) over `anchor`; // same-shape emits a real corgi `Rel`. @@ -589,7 +605,20 @@ pub fn compile( let payload = b.add(ArithOp::ToSigned, vec![abs]); b.add(Op::Inject(0, vec![Shape::Prim(64)]), vec![payload]) } - UnOp::F64Fn(_) | UnOp::F64ToInt => return Err(format!("{op:?} has no columnar kernel; it runs a row at a time (see `row_only`)")), + // The rest (sqrt, exp, ln, rounding, trig) and `fint` are host kernels over the + // float leaf: decode each key, apply the function, encode the result. + UnOp::F64Fn(f) => { + if shape != Shape::Sum(vec![Shape::Prim(64)]) { return Err(format!("{op:?} expects an F64 newtype")); } + let x = float_leaf(b, id); + let y = b.add(NumOp::Host(float_kernels::op(float_kernels::FloatOp::Fn(*f))), vec![x]); + let payload = b.add(ArithOp::ToSigned, vec![y]); + b.add(Op::Inject(0, vec![Shape::Prim(64)]), vec![payload]) + } + UnOp::F64ToInt => { + if shape != Shape::Sum(vec![Shape::Prim(64)]) { return Err("fint expects an F64 newtype".into()); } + let x = float_leaf(b, id); + b.add(NumOp::Host(float_kernels::op(float_kernels::FloatOp::ToInt)), vec![x]) + } // `truthy` is "nonzero Int": scalars compare against zero; non-`Int` values // are never truthy, so their `not` folds to the constant 1 (the cross-shape // `Eq` fold's precedent). @@ -763,142 +792,24 @@ fn compile_fold_body(step: &Term, ctx: Option<&[Shape]>, init_shape: &Shape, ele Ok(bb.finish(out)) } -/// A compiled term, ready to run on a batch of columns. -pub enum Kernel { - /// Every op in the terms has a corgi kernel: one columnar graph. - Columnar(Graph), - /// Some op has none (see [`row_only`]). The terms run through `ir::eval` a row at a time: - /// the input columns are untranscoded from `inputs`, and the results transcoded to `output`. - /// The output shape is still corgi's typing, of the terms with each row-only op replaced by a - /// columnar op of the same shape ([`columnar_stand_in`]), so both paths type a term alike. - Rows { terms: Vec, inputs: Vec, output: Shape }, -} - -impl Kernel { - /// Run on `Prod(inputs)`. One term yields its column; several yield `Prod` of their columns. - pub fn eval(&self, input: CValue) -> CValue { - match self { - Kernel::Columnar(g) => corgi::eval_graph(g, input), - Kernel::Rows { terms, inputs, output } => { - let rows = untranscode(input, &Shape::Prod(inputs.clone())); - let results: Vec = rows - .into_iter() - .map(|row| { - let DValue::Tuple(mut env) = row else { unreachable!("untranscode of a Prod is a Tuple") }; - match &terms[..] { - [term] => crate::ir::eval(term, &mut env), - _ => DValue::Tuple(terms.iter().map(|t| crate::ir::eval(t, &mut env)).collect()), - } - }) - .collect(); - transcode_owned(results, output) - } - } - } -} - -/// Does `t` use an op with no corgi kernel? These are the F64 functions other than `fabs`, `fint`, -/// `fpow` and `fpowi`: libm calls and rounding, which corgi's arithmetic does not offer. -pub fn row_only(t: &Term) -> bool { - match t { - Term::Var(_) | Term::Bound(_) | Term::Int(_) => false, - Term::Unary(UnOp::F64Fn(f), _) if *f != F64Fn::Abs => true, - Term::Unary(UnOp::F64ToInt, _) | Term::Binary(BinOp::F64Pow | BinOp::F64PowI, _, _) => true, - Term::Tuple(fs) | Term::List(fs) | Term::Hash(fs) | Term::Call(_, fs) => fs.iter().any(row_only), - Term::Spread(inner) | Term::Proj(inner, _) | Term::Unary(_, inner) => row_only(inner), - Term::Inject { tag, payload, .. } => row_only(tag) || row_only(payload), - Term::Case { scrutinee, arms, default } => { - row_only(scrutinee) || arms.iter().any(row_only) || default.as_deref().is_some_and(row_only) - } - Term::Fold { list, init, step } => row_only(list) || row_only(init) || row_only(step), - Term::If { cond, then, els } => row_only(cond) || row_only(then) || row_only(els), - Term::Binary(_, l, r) => row_only(l) || row_only(r), - } -} - -/// `t` with each row-only op replaced by a columnar op that demands the same operand shapes and -/// yields the same result shape: F64 -> F64 by `fneg`, `fint` by a compare of two `fneg`s, -/// `fpow` by `fadd`, and `fpowi(x, n)` by `fadd(x, float(n))`. Only its typing is used. -fn columnar_stand_in(t: &Term) -> Term { - let bx = |term: &Term| Box::new(columnar_stand_in(term)); - match t { - Term::Var(_) | Term::Bound(_) | Term::Int(_) => t.clone(), - Term::Unary(UnOp::F64Fn(f), x) if *f != F64Fn::Abs => Term::Unary(UnOp::F64Neg, bx(x)), - Term::Unary(UnOp::F64ToInt, x) => { - let x = Term::Unary(UnOp::F64Neg, bx(x)); - Term::Binary(BinOp::Lt, Box::new(x.clone()), Box::new(x)) - } - Term::Binary(BinOp::F64Pow, l, r) => Term::Binary(BinOp::F64Add, bx(l), bx(r)), - // A registered function: its arguments still compile (so they typecheck), and the - // result is a literal of the declared result shape. - Term::Call(name, args) => match crate::ir::lookup(name) { - Some(f) => { - let mut fields: Vec = args.iter().map(columnar_stand_in).collect(); - fields.push(literal_of_shape(&f.result)); - Term::Proj(Box::new(Term::Tuple(fields)), args.len()) - } - None => t.clone(), - }, - Term::Binary(BinOp::F64PowI, l, r) => Term::Binary(BinOp::F64Add, bx(l), Box::new(Term::Unary(UnOp::ToF64, bx(r)))), - Term::Tuple(fs) => Term::Tuple(fs.iter().map(columnar_stand_in).collect()), - Term::List(fs) => Term::List(fs.iter().map(columnar_stand_in).collect()), - Term::Hash(fs) => Term::Hash(fs.iter().map(columnar_stand_in).collect()), - Term::Spread(inner) => Term::Spread(bx(inner)), - Term::Proj(inner, i) => Term::Proj(bx(inner), *i), - Term::Unary(op, inner) => Term::Unary(*op, bx(inner)), - Term::Inject { tag, payload, sum } => Term::Inject { tag: bx(tag), payload: bx(payload), sum: sum.clone() }, - Term::Case { scrutinee, arms, default } => Term::Case { - scrutinee: bx(scrutinee), - arms: arms.iter().map(columnar_stand_in).collect(), - default: default.as_deref().map(bx), - }, - Term::Fold { list, init, step } => Term::Fold { list: bx(list), init: bx(init), step: bx(step) }, - Term::If { cond, then, els } => Term::If { cond: bx(cond), then: bx(then), els: bx(els) }, - Term::Binary(op, l, r) => Term::Binary(*op, bx(l), bx(r)), - } -} - -/// A term whose corgi shape is `shape`: zeros, first lanes, one-element lists. -fn literal_of_shape(shape: &Shape) -> Term { - match shape { - Shape::Prim(_) => Term::Int(0), - Shape::Unit => Term::Tuple(Vec::new()), - Shape::Prod(fs) => Term::Tuple(fs.iter().map(literal_of_shape).collect()), - Shape::List(e) => Term::List(vec![literal_of_shape(e)]), - Shape::Sum(lanes) => Term::Inject { - tag: Box::new(Term::Int(0)), - payload: Box::new(literal_of_shape(&lanes[0])), - sum: SumTy::Declared(lanes.clone()), - }, - } -} - /// Compile `terms` over the environment `Var(i)` = field `i` of the input, whose shapes are /// `shapes`. The graph's input is `Prod(shapes)`; its output is the one term's column, or `Prod` /// of the terms' columns. Typechecked once here, so an `Ok` kernel runs on every batch of these /// shapes. Returns the kernel and its output shape. -fn lower(terms: &[&Term], shapes: &[Shape]) -> Res<(Kernel, Shape)> { - let rows = terms.iter().any(|t| row_only(t)); - let stand_ins: Vec = if rows { terms.iter().map(|t| columnar_stand_in(t)).collect() } else { Vec::new() }; - let typed: Vec<&Term> = if rows { stand_ins.iter().collect() } else { terms.to_vec() }; +fn lower(terms: &[&Term], shapes: &[Shape]) -> Res<(Graph, Shape)> { let mut b = Builder::::default(); let input = b.input(); let env: Vec = (0..shapes.len()).map(|i| b.add(Op::Field(i), vec![input])).collect(); - let outs = typed.iter().map(|t| compile(t, &mut b, &env, shapes, input, None)).collect::>>()?; + let outs = terms.iter().map(|t| compile(t, &mut b, &env, shapes, input, None)).collect::>>()?; let out = if let [one] = outs[..] { one } else { b.tuple(outs) }; let g = b.finish(out); let output = corgi::shape_of(&g, &Shape::Prod(shapes.to_vec()))?; - let kernel = if rows { - Kernel::Rows { terms: terms.iter().map(|t| (*t).clone()).collect(), inputs: shapes.to_vec(), output: output.clone() } - } else { - Kernel::Columnar(g) - }; - Ok((kernel, output)) + Ok((g, output)) } /// Compile a `FlatMap`'s list term → a corgi `List` column, one list per input row. A term that is /// not list-shaped is the type error: the backend explodes the column structurally. -pub fn compile_flatmap(list_term: &Term, kshape: &Shape, vshape: &Shape) -> Res { +pub fn compile_flatmap(list_term: &Term, kshape: &Shape, vshape: &Shape) -> Res> { match lower(&[list_term], &[kshape.clone(), vshape.clone()])? { (k, Shape::List(_)) => Ok(k), (_, other) => Err(format!("flatmap over a non-list: {other}")), @@ -907,7 +818,7 @@ pub fn compile_flatmap(list_term: &Term, kshape: &Shape, vshape: &Shape) -> Res< /// Compile a scalar term (`EnterAt`'s delay field) → a `U64` column; a non-integer term is the /// type error (the delay is read as one integer per row). -pub fn compile_scalar(term: &Term, kshape: &Shape, vshape: &Shape) -> Res { +pub fn compile_scalar(term: &Term, kshape: &Shape, vshape: &Shape) -> Res> { match lower(&[term], &[kshape.clone(), vshape.clone()])? { (k, Shape::Prim(_)) => Ok(k), (_, other) => Err(format!("enter_at delay is not an integer: {other}")), @@ -916,7 +827,7 @@ pub fn compile_scalar(term: &Term, kshape: &Shape, vshape: &Shape) -> Res Res { +pub fn compile_predicate(cond: &Term, kshape: &Shape, vshape: &Shape) -> Res> { match lower(&[cond], &[kshape.clone(), vshape.clone()])? { (k, Shape::Prim(_)) => Ok(k), (_, other) => Err(format!("a filter predicate must be an Int, got {other}")), @@ -925,16 +836,76 @@ pub fn compile_predicate(cond: &Term, kshape: &Shape, vshape: &Shape) -> Res Res { +pub fn compile_join_projection(key: &Term, val: &Term, kshape: &Shape, v0shape: &Shape, v1shape: &Shape) -> Res> { Ok(lower(&[key, val], &[kshape.clone(), v0shape.clone(), v1shape.clone()])?.0) } /// Compile a DDIR `Projection` over `Var(0)=key` (`kshape`), `Var(1)=val` (`vshape`). /// Input `Prod([key, val])`; output `Prod([newkey, newval])`. -pub fn compile_projection(key: &Term, val: &Term, kshape: &Shape, vshape: &Shape) -> Res { +pub fn compile_projection(key: &Term, val: &Term, kshape: &Shape, vshape: &Shape) -> Res> { Ok(lower(&[key, val], &[kshape.clone(), vshape.clone()])?.0) } +/// The F64 functions without a composition of corgi ops, as host kernels over the float leaf +/// (corgi's total-order key): decode each key, apply the function as `ir::eval` does, encode. +/// One kernel per function, shared by every call site, so CSE merges equal calls. +mod float_kernels { + use std::sync::{Arc, OnceLock}; + use corgi::{HostKernel, HostOp, Shape, Value}; + use crate::ir::F64Fn; + + #[derive(Clone, Copy, PartialEq, Eq, Hash, Debug)] + pub enum FloatOp { Fn(F64Fn), ToInt, Pow, PowI } + + fn decode(key: u64) -> f64 { + f64::from_bits(if key >> 63 == 1 { key ^ (1 << 63) } else { !key }) + } + fn encode(x: f64) -> u64 { super::float_key(x) } + + struct Kernel { name: String, op: FloatOp, input: Shape, output: Shape } + impl HostKernel for Kernel { + fn name(&self) -> &str { &self.name } + fn input(&self) -> &Shape { &self.input } + fn output(&self) -> &Shape { &self.output } + fn eval(&self, input: Value) -> Result { + let out: Vec = match self.op { + FloatOp::Fn(f) => input.as_u64(&self.name)?.iter().map(|&k| encode(f.apply(decode(k)))).collect(), + FloatOp::ToInt => input.as_u64(&self.name)?.iter().map(|&k| decode(k) as i64 as u64).collect(), + FloatOp::Pow | FloatOp::PowI => { + let args = input.into_prod(&self.name)?; + let (x, y) = (args[0].as_u64(&self.name)?, args[1].as_u64(&self.name)?); + x.iter().zip(y).map(|(&a, &b)| { + let a = decode(a); + encode(if self.op == FloatOp::Pow { + a.powf(decode(b)) + } else { + a.powi((b as i64).clamp(i32::MIN as i64, i32::MAX as i64) as i32) + }) + }).collect() + } + }; + Ok(Value::u64(out)) + } + } + + /// The host op for `op`: one `Arc` per function for the life of the process. + pub fn op(op: FloatOp) -> HostOp { + static KERNELS: OnceLock>>> = OnceLock::new(); + let mut map = KERNELS.get_or_init(Default::default).lock().unwrap(); + let k = map.entry(op).or_insert_with(|| { + let (p, pair) = (Shape::Prim(64), Shape::Prod(vec![Shape::Prim(64), Shape::Prim(64)])); + let (name, input) = match op { + FloatOp::Fn(f) => (format!("{f:?}").to_lowercase(), p.clone()), + FloatOp::ToInt => ("fint".into(), p.clone()), + FloatOp::Pow => ("fpow".into(), pair.clone()), + FloatOp::PowI => ("fpowi".into(), pair), + }; + Arc::new(Kernel { name: format!("f64:{name}"), op, input, output: p }) + }); + HostOp(Arc::clone(k)) + } +} + #[cfg(test)] mod tests { use super::*; @@ -979,40 +950,29 @@ mod tests { specials.iter().flat_map(|&a| specials.iter().map(move |&b| vec![V::f64_value(a), V::f64_value(b)])).collect() } - /// The F64 ops with a columnar kernel agree with `ir::eval` on every special pair. + /// Every F64 op agrees with `ir::eval` on every special pair: the ones composed of corgi ops + /// and the ones that are host kernels (`float_kernels`), alone and inside larger terms. #[test] fn columnar_float_math_agrees() { let f = sum(vec![u64s()]); for source in ["fabs($0)", "fmin($0, $1)", "fmax($0, $1)", "feq($0, $1)", "fne($0, $1)", - "flt($0, $1)", "fle($0, $1)", "fgt($0, $1)", "fge($0, $1)"] { + "flt($0, $1)", "fle($0, $1)", "fgt($0, $1)", "fge($0, $1)", + "fsqrt($0)", "fexp($0)", "fln($0)", "ffloor($0)", "fceil($0)", "fround($0)", + "fsin($0)", "fcos($0)", "ftan($0)", "fint($0)", "fpow($0, $1)", "fpowi($0, fint($1))"] { let term = crate::parse::pipe::parse_term(source); - assert!(!row_only(&term), "{source} should have a columnar kernel"); agrees_with_rows(&term, &[f.clone(), f.clone()], &float_pairs()); } - } - - /// The F64 ops without one run a row at a time, typed as their columnar stand-ins; the - /// kernel's output round-trips through the transcode layer at that shape. - #[test] - fn row_wise_float_math_agrees() { - let f = sum(vec![u64s()]); - for source in ["fsqrt($0)", "fexp($0)", "fln($0)", "ffloor($0)", "fceil($0)", "fround($0)", - "fsin($0)", "fcos($0)", "ftan($0)", "fint($0)", "fpow($0, $1)", "fpowi($0, fint($1))", - "tuple(fadd(fexp($0), $1), flt(fln($0), $1))"] { - let term = crate::parse::pipe::parse_term(source); - assert!(row_only(&term), "{source} should run row-wise"); - let shapes = [f.clone(), f.clone()]; - let (kernel, output) = lower(&[&term], &shapes).unwrap(); - assert!(matches!(kernel, Kernel::Rows { .. })); - let rows = float_pairs(); - let cols = (0..2).map(|i| transcode(&rows.iter().map(|r| r[i].clone()).collect::>(), &shapes[i])).collect(); - let actual = untranscode(kernel.eval(CValue::Prod(cols)), &output); - let expected: Vec<_> = rows.iter().map(|r| crate::ir::eval(&term, &mut r.clone())).collect(); - assert_eq!(actual, expected, "{source}"); + // Inside a larger term. Over non-NaN pairs only: corgi's float `Add` returns the positive + // quiet NaN where `ir::eval` keeps a NaN operand's sign (as `fadd` alone does, before + // this change); the kernels themselves agree on NaNs above. + let finite: Vec<_> = float_pairs().into_iter().filter(|r| r.iter().all(|v| !v.as_f64().is_nan())).collect(); + let term = crate::parse::pipe::parse_term("tuple(fadd(fexp($0), $1), flt(fln($0), $1), fpowi(fsqrt($0), fint($1)))"); + agrees_with_rows(&term, &[f.clone(), f.clone()], &finite); + // The kernels are typed: their operands must be F64 newtypes (and `fpowi`'s exponent an Int). + for bad in ["fsqrt($0)", "fint($0)", "fpow($0, $0)"] { + assert!(lower(&[&crate::parse::pipe::parse_term(bad)], &[u64s()]).is_err(), "{bad}"); } - // A row-only op is still typed: its operands must be F64 newtypes. - assert!(lower(&[&crate::parse::pipe::parse_term("fsqrt($0)")], &[u64s()]).is_err()); - assert!(lower(&[&crate::parse::pipe::parse_term("fpow($0, $0)")], &[u64s()]).is_err()); + assert!(lower(&[&crate::parse::pipe::parse_term("fpowi($0, $0)")], &[f.clone()]).is_err()); } /// Spot checks of the chosen semantics, independent of either backend. diff --git a/interactive/src/ir.rs b/interactive/src/ir.rs index 9d3fe49fc..c8dfa210a 100644 --- a/interactive/src/ir.rs +++ b/interactive/src/ir.rs @@ -115,7 +115,9 @@ pub enum Term { /// `name(args…)` for a function the embedding program registered (see /// [`register`]): a pure Rust function from argument values to a value. /// DDIR only moves its arguments and result; what it computes is the - /// embedder's. The columnar backend runs it a row at a time. + /// embedder's. The corgi backend runs it as a host kernel: the function's + /// columnar kernel if one is registered (`register_kernel`), otherwise its + /// row body over just the call's arguments. Call(String, Vec), } @@ -161,7 +163,7 @@ pub enum UnOp { /// program ported from Rust gets bit-identical results on the same platform. `Abs` is exact /// everywhere; the transcendental ones (`Exp`, `Ln`, `Sin`, `Cos`, `Tan`) call the platform libm /// and are only as reproducible across platforms as it is. -#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)] pub enum F64Fn { Abs, Sqrt, Exp, Ln, Floor, Ceil, Round, Sin, Cos, Tan, } @@ -305,6 +307,11 @@ pub fn eval(term: &Term, env: &mut Vec) -> Value { Term::Call(name, args) => { let f = lookup(name).unwrap_or_else(|| panic!("call to unregistered function `{name}`")); let vals: Vec = args.iter().map(|a| eval(a, env)).collect(); + // The corgi backend rejects mis-shaped arguments when it types the program; the row + // backend has no typing pass, so it checks each call rather than run the body on them. + for (i, (v, s)) in vals.iter().zip(&f.args).enumerate() { + assert!(v.has_shape(s), "`{name}`: argument {i} is {v:?}, not of the declared shape {s}"); + } (f.body)(&vals) } Term::Var(i) => env[*i].clone(), @@ -502,7 +509,10 @@ fn functions() -> &'static std::sync::RwLock &corgi::Shape { &self.f.result } fn eval(&self, input: corgi::Value) -> Result { let rows = crate::corgi::logic::untranscode(input, &self.input); - let out: Vec = rows.into_iter().map(|r| match r { - Value::Tuple(args) => (self.f.body)(&args), - other => (self.f.body)(std::slice::from_ref(&other)), + // The input is `Prod(args)` or `Unit` (`call_input`); both untranscode to tuples. + let out: Vec = rows.into_iter().map(|r| { + let Value::Tuple(args) = r else { unreachable!("a call's arguments untranscode to a tuple") }; + (self.f.body)(&args) }).collect(); Ok(crate::corgi::logic::transcode_owned(out, &self.f.result)) } diff --git a/interactive/src/parse/mod.rs b/interactive/src/parse/mod.rs index 0ddf9313f..cd0694eef 100644 --- a/interactive/src/parse/mod.rs +++ b/interactive/src/parse/mod.rs @@ -120,5 +120,5 @@ pub(crate) fn build_builtin(name: &str, args: &mut Vec) -> Term { pub fn is_builtin(name: &str) -> bool { const NAMES: &[&str] = &["tuple", "list", "inject", "variant", "case", "fold", "proj", "len", "istag", "not", "float", "fneg", "fint", "fpow", "fpowi", "fmin", "fmax", "feq", "fne", "flt", "fle", "fgt", "fge", "fadd", "fsub", "fmul", "fdiv", "or", "idiv", "append", "if", "hash"]; - NAMES.contains(&name) || F64Fn::ALL.iter().any(|(n, _)| *n == name) + NAMES.contains(&name) || F64Fn::ALL.iter().any(|(n, _)| *n == name) || pipe::is_keyword(name) } diff --git a/interactive/src/parse/pipe.rs b/interactive/src/parse/pipe.rs index 58dbed238..fa6c3823b 100644 --- a/interactive/src/parse/pipe.rs +++ b/interactive/src/parse/pipe.rs @@ -64,9 +64,9 @@ //! `-0.0`). `feq fne flt fle fgt fge` are IEEE comparisons returning Int 0/1: //! NaN compares false (`fne` true) and `-0.0` equals `0.0`, where the generic //! `== < …` use the total order. There are no float literals: write -//! `fdiv(float(2786), float(10))` for 278.6. The Corgi backend computes -//! `fabs`, `fmin`/`fmax` and the comparisons in columns; a term using any -//! other of these is evaluated a row at a time. +//! `fdiv(float(2786), float(10))` for 278.6. The Corgi backend computes all +//! of these in columns: `fabs`, `fmin`/`fmax` and the comparisons from corgi +//! ops, the rest as host kernels over the float column. //! - Products: `tuple(a, …)`; index with `v[i]` or `proj(v, i)`; `len(v)`. //! - Lists: `list(a, …)`, `append(a, b)` (concatenation); eliminated by //! `flatmap` / `collect` / `fold`. A declared constructor supplies the element @@ -156,19 +156,7 @@ fn tokenize(input: &str) -> Vec { c if c.is_ascii_alphabetic() || c == '_' => { let mut ident = String::new(); while let Some(&c) = chars.peek() { if c.is_ascii_alphanumeric() || c == '_' { ident.push(c); chars.next(); } else { break; } } - tokens.push(match ident.as_str() { - "let" => Token::Let, "var" => Token::Var, "export" => Token::Export, - "type" => Token::Type, "case" => Token::Case, "fold" => Token::Fold, - "input" => Token::Input, "import" => Token::Import, - "key" => Token::Key, "map" => Token::Map, - "join" => Token::Join, "min" => Token::Min, "distinct" => Token::Distinct, - "count" => Token::Count, "collect" => Token::Collect, - "flatmap" => Token::FlatMap, - "arrange" => Token::Arrange, "negate" => Token::Negate, - "filter" => Token::Filter, "enter_at" => Token::EnterAt, "inspect" => Token::Inspect, - "lift_iter" => Token::LiftIter, - _ => Token::Ident(ident), - }); + tokens.push(keyword(&ident).unwrap_or(Token::Ident(ident))); }, other => panic!("Unexpected character: {:?}", other), } @@ -177,6 +165,26 @@ fn tokenize(input: &str) -> Vec { tokens } +/// The lexer's keywords: an identifier spelled like one lexes as the keyword, never as a name. +fn keyword(ident: &str) -> Option { + Some(match ident { + "let" => Token::Let, "var" => Token::Var, "export" => Token::Export, + "type" => Token::Type, "case" => Token::Case, "fold" => Token::Fold, + "input" => Token::Input, "import" => Token::Import, + "key" => Token::Key, "map" => Token::Map, + "join" => Token::Join, "min" => Token::Min, "distinct" => Token::Distinct, + "count" => Token::Count, "collect" => Token::Collect, + "flatmap" => Token::FlatMap, + "arrange" => Token::Arrange, "negate" => Token::Negate, + "filter" => Token::Filter, "enter_at" => Token::EnterAt, "inspect" => Token::Inspect, + "lift_iter" => Token::LiftIter, + _ => return None, + }) +} + +/// Whether `ident` is a keyword of the lexer (so it can never be used as a name). +pub(crate) fn is_keyword(ident: &str) -> bool { keyword(ident).is_some() } + struct Parser { tokens: Vec, pos: usize, diff --git a/interactive/tests/corgi_backend.rs b/interactive/tests/corgi_backend.rs index 3bafd8ddb..7d3cd04ce 100644 --- a/interactive/tests/corgi_backend.rs +++ b/interactive/tests/corgi_backend.rs @@ -161,7 +161,7 @@ fn serializing(n: usize) -> timely::Config { assert_eq!(out["joined"].len(), 40); } -/// The functions `registered.ddp` calls. Registering again is idempotent, so every +/// The functions `registered.ddp` calls. Registering again replaces (here with the same bodies), so every /// test may do it. fn register_test_functions() { use corgi::Shape; @@ -242,3 +242,39 @@ fn filter_requires_an_int_predicate() { assert!(result.is_err(), "{backend:?} accepted a tuple filter predicate"); } } + +/// A call whose argument does not have the declared shape is rejected by both backends: corgi +/// when it types the program, the row backend when it makes the call. +#[test] +fn call_argument_shapes_are_checked_by_both_backends() { + register_test_functions(); + // `blend` takes a pair; pass it an Int. + let mut tree = lower::lower_tree(parse::pipe::parse(r#"export "result" = input 0 | map($0 ; blend($0[0]));"#)); + tree.optimize(); + let inputs = vec![rows(&[&[1], &[2]])]; + for backend in [RenderBackend::Vec, RenderBackend::Corgi] { + let (tree, inputs) = (tree.clone(), inputs.clone()); + let result = std::panic::catch_unwind(move || evaluate(backend, timely::Config::process(1), &tree, &inputs)); + assert!(result.is_err(), "{backend:?} ran a call with a mis-shaped argument"); + } +} + +/// Registering a name again replaces its kernel too: the row adapter for the old body, or a +/// columnar kernel registered for it, no longer runs. +#[test] +fn reregistering_a_function_drops_its_old_kernel() { + use corgi::Shape; + use interactive::ir::{kernel_of, register, register_kernel, Function}; + let one = |k: i64| Function { name: "reregistered".into(), args: vec![Shape::Prim(64)], result: Shape::Prim(64), body: Box::new(move |_| Value::Int(k)) }; + register(one(1)); + let first = kernel_of("reregistered").unwrap(); + assert!(std::sync::Arc::ptr_eq(&first, &kernel_of("reregistered").unwrap()), "one kernel per registration"); + register(one(2)); + let second = kernel_of("reregistered").unwrap(); + assert!(!std::sync::Arc::ptr_eq(&first, &second), "re-registration kept the old kernel"); + register_kernel("reregistered", first); + register(one(3)); + assert!(!std::sync::Arc::ptr_eq(&second, &kernel_of("reregistered").unwrap())); + // A keyword can never be called, so it cannot be registered. + assert!(std::panic::catch_unwind(|| register(Function { name: "min".into(), args: vec![], result: Shape::Prim(64), body: Box::new(|_| Value::Int(0)) })).is_err()); +}