From cc204609db73957691c43aeadb8551981d36ed01 Mon Sep 17 00:00:00 2001 From: Frank McSherry Date: Wed, 2 Sep 2026 17:11:40 -0400 Subject: [PATCH] DDIR: adopt corgi's copy-audit API (frankmcsherry/WIP#18) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Bumps the corgi pin from #16 to #18 (`de0f2ac`), which also picks up #17. Four API changes, all of which delete code at the call site: * `corgi::hash` returns the ids (`Vec`) rather than a `Value` that every caller here unwrapped in the same expression. Four sites lose `.into_u64(..).unwrap()`. * `arrange::hash_rows` is gone — it was a second, disagreeing implementation of the same fold (it mixed leaf width into its seed, so it called a `u8 5` and a `u64 5` different values, which the canonical hash does not). The two round-trip test calls become `corgi::hash`. * `Bounds::Offsets` holds an `Arc`, so a `Value` clone no longer copies a list's partition. Construction goes through `Bounds::offsets`, and the three places that hand-rolled a `Vec` of end offsets by matching both variants now call `Bounds::to_vec`, which is public and does it. * `Value::Sum` carries one `Tags` — the lane assignment — instead of a tag column and an offset column, with a two-word `Const` form for the one-tag case. `signed_order_view` moves the assignment through untouched; `untranscode` reads `tag_at`/`offset_at` per row. And one adoption beyond compiling: `sort_consolidate` was building the two `i`/`i+1` index columns to ask for an adjacent comparison. `compare_adjacent` names the pattern instead, so the index columns are not built and corgi reads both sides densely. That is on the chunk-merge path. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01PPtSwqXTuH4wefbQd2F8pA --- interactive/Cargo.toml | 2 +- interactive/src/backend/corgi.rs | 5 +---- interactive/src/corgi/bytes.rs | 4 ++-- interactive/src/corgi/chunk.rs | 16 ++++++---------- interactive/src/corgi/exchange.rs | 2 +- interactive/src/corgi/logic.rs | 17 ++++++++--------- interactive/src/corgi/reduce.rs | 16 ++++++---------- 7 files changed, 25 insertions(+), 37 deletions(-) diff --git a/interactive/Cargo.toml b/interactive/Cargo.toml index 1673d1440..adb8ea20a 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 = "578ef3fe79fdaa813ba0c721c08f60323f22d420", features = ["serde"] } +corgi = { git = "https://github.com/frankmcsherry/WIP", rev = "de0f2ac91d31ac63641035e5a5bb1b5640fe9424", features = ["serde"] } differential-dataflow = { workspace = true } mimalloc = "0.1.48" serde = { version = "1.0", features = ["derive"] } diff --git a/interactive/src/backend/corgi.rs b/interactive/src/backend/corgi.rs index 4422b235c..0b48a19ef 100644 --- a/interactive/src/backend/corgi.rs +++ b/interactive/src/backend/corgi.rs @@ -203,10 +203,7 @@ fn apply_ops(mut c: CC, ops: &[LinearOp], level: usize, plans: &mut [Plan]) -> C // 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(); - let ends: Vec = match &bounds { - corgi::Bounds::Offsets(v) => v.clone(), - corgi::Bounds::Stride(k, rows) => (1..=*rows).map(|i| i * k).collect(), - }; + 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)); let mut start = 0usize; diff --git a/interactive/src/corgi/bytes.rs b/interactive/src/corgi/bytes.rs index 3275a9dfe..c7f848c80 100644 --- a/interactive/src/corgi/bytes.rs +++ b/interactive/src/corgi/bytes.rs @@ -226,8 +226,8 @@ mod test { let c = container_of(name, updates); let (keys, vals) = (c.keys.clone(), c.vals.clone()); let back = round_trip(&c); - assert_eq!(corgi::arrange::hash_rows(&back.keys), corgi::arrange::hash_rows(&keys), "{name} keys"); - assert_eq!(corgi::arrange::hash_rows(&back.vals), corgi::arrange::hash_rows(&vals), "{name} vals"); + assert_eq!(corgi::hash(&back.keys), corgi::hash(&keys), "{name} keys"); + assert_eq!(corgi::hash(&back.vals), corgi::hash(&vals), "{name} vals"); } } diff --git a/interactive/src/corgi/chunk.rs b/interactive/src/corgi/chunk.rs index a7aea6b13..79eb4562d 100644 --- a/interactive/src/corgi/chunk.rs +++ b/interactive/src/corgi/chunk.rs @@ -28,7 +28,7 @@ use timely::progress::frontier::AntichainRef; use differential_dataflow::difference::Semigroup; use differential_dataflow::trace::chunk::{pack, Chunk, ChunkBatch}; -use corgi::arrange::{compare_at, compare_idx, gather, gather_lanes, group_bounds, sort_perm}; +use corgi::arrange::{compare_adjacent, compare_at, gather, gather_lanes, group_bounds, sort_perm}; use corgi::Value as CValue; use columnar::Columnar; @@ -348,7 +348,7 @@ where /// (summing diffs, dropping zeros). Returns a sorted+consolidated `(keys, vals, times, diffs)`. /// /// Multi-record: one columnar `sort_perm` (discrimination sort) orders by `(key, val)`, one batched -/// `compare_idx` flags adjacent-equal runs; only the small per-run *time* tiebreak is a Rust sort +/// `compare_adjacent` flags adjacent-equal runs; only the small per-run *time* tiebreak is a Rust sort /// (time is not a corgi type). No per-pair `compare_at`. fn sort_consolidate(keys: CValue, vals: CValue, times: Vec, diffs: Vec) -> (CValue, CValue, Vec, Vec) where @@ -366,13 +366,9 @@ where let times_s: Vec = perm.iter().map(|&i| times[i].clone()).collect(); let diffs_s: Vec = perm.iter().map(|&i| diffs[i].clone()).collect(); // Batched adjacent-equality over the kv-sorted column: `adj[m] == 0` iff `kv_s[m] == kv_s[m+1]`. - let adj: Vec = if n > 1 { - let left: Vec = (0..n - 1).collect(); - let right: Vec = (1..n).collect(); - compare_idx(&kv_s, &kv_s, &left, &right) - } else { - Vec::new() - }; + // Naming the pattern rather than writing out the two index columns: corgi reads both sides + // densely, and the `i`/`i+1` index vectors this used to build are not built at all. + let adj: Vec = compare_adjacent(&kv_s); // Walk maximal equal-`(key,val)` runs; within each, order by time and consolidate equal times. let (mut keep, mut ot, mut od) = (Vec::new(), Vec::new(), Vec::new()); @@ -497,7 +493,7 @@ pub fn present_key(keys: CValue) -> CValue { if corgi::arrange::leaf_slice(&keys).is_some() { return keys; } - let hashes = corgi::hash(&keys).into_u64("present_key").unwrap(); + let hashes = corgi::hash(&keys); CValue::Prod(vec![CValue::u64(hashes), keys]) } diff --git a/interactive/src/corgi/exchange.rs b/interactive/src/corgi/exchange.rs index aaa875ff9..110ea0fef 100644 --- a/interactive/src/corgi/exchange.rs +++ b/interactive/src/corgi/exchange.rs @@ -111,7 +111,7 @@ impl Distributor> f } let peers = pushers.len(); - let ids = corgi::hash(&container.keys).into_u64("corgi exchange: key hash").unwrap(); + let ids = corgi::hash(&container.keys); self.counting_sort(&ids, peers); // Whole-container fast path. When every row shares a destination — a batch narrower than diff --git a/interactive/src/corgi/logic.rs b/interactive/src/corgi/logic.rs index 5e20aafb2..434539b71 100644 --- a/interactive/src/corgi/logic.rs +++ b/interactive/src/corgi/logic.rs @@ -130,10 +130,7 @@ pub fn untranscode(col: CValue, shape: &Shape) -> Vec { other => panic!("untranscode: expected List, got {other:?}"), }; let flat = untranscode(vals, elem); - let ends: Vec = match &bounds { - corgi::Bounds::Offsets(v) => v.clone(), - corgi::Bounds::Stride(k, rows) => (1..=*rows).map(|i| i * k).collect(), - }; + let ends: Vec = bounds.to_vec(); let mut out = Vec::with_capacity(ends.len()); let mut start = 0usize; for end in ends { @@ -146,12 +143,14 @@ pub fn untranscode(col: CValue, shape: &Shape) -> Vec { // Inverse of transcode's Sum: untranscode each lane, then for each row pull its payload // from its lane at the recorded within-lane OFFSET (robust to row reordering from a // prior gather/merge — not a sequential cursor). - let (tags, offsets, variant_vals) = col.into_sum("untranscode").unwrap(); + let (tags, variant_vals) = col.into_sum("untranscode").unwrap(); let lane_rows: Vec> = variant_vals.into_iter().zip(lanes.iter()).map(|(v, ls)| untranscode(v, ls)).collect(); - tags.iter() - .zip(&offsets) - .map(|(&tag, &off)| DValue::Variant(tag as u32, Box::new(lane_rows[tag][off].clone()))) + (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())) + }) .collect() } } @@ -702,7 +701,7 @@ mod tests { /// this is what fails. fn hash_agrees(rows: Vec, shape: Shape) { let col = transcode(&rows, &shape); - let columnar = corgi::hash(&col).into_u64("hash").unwrap(); + let columnar = corgi::hash(&col); let row_wise: Vec = rows.iter().map(crate::ir::structural_hash).collect(); assert_eq!(columnar, row_wise, "hash disagrees (shape {shape:?})"); } diff --git a/interactive/src/corgi/reduce.rs b/interactive/src/corgi/reduce.rs index 85e247bd5..63e5281b2 100644 --- a/interactive/src/corgi/reduce.rs +++ b/interactive/src/corgi/reduce.rs @@ -57,14 +57,10 @@ fn signed_order_view(value: CValue) -> CValue { CValue::Prod(fields) => { CValue::Prod(fields.into_iter().map(signed_order_view).collect()) } - CValue::Sum(tags, within, variants) => CValue::Sum( - tags, - within, - variants - .into_iter() - .map(signed_order_view) - .collect(), - ), + CValue::Sum(tags, variants) => { + // the lane assignment is untouched — only the payload lanes are swizzled. + CValue::Sum(tags, variants.into_iter().map(signed_order_view).collect()) + } CValue::List(bounds, values) => { CValue::List(bounds, Box::new(signed_order_view(*values))) } @@ -200,7 +196,7 @@ fn ids(col: &CValue) -> Vec { if let Some(sl) = corgi::arrange::leaf_slice(col) { return sl.to_vec(); } - corgi::hash(col).into_u64("ids").unwrap() + corgi::hash(col) } /// Concatenate the records of the `changed` keys across a run of chunks into parallel @@ -542,7 +538,7 @@ where // next batch's `List` where the two are concatenated. `gather` at no indices // is the empty column of that shape. let elems = gather(&self.in_vals, &elem_reps); - let col = CValue::List(Bounds::Offsets(bracket_ends), Box::new(elems)); + let col = CValue::List(Bounds::offsets(bracket_ends), Box::new(elems)); out_ids = ids(&col); self.register_vals(col, &out_ids); }