Skip to content

Commit 74d9370

Browse files
committed
fixup! Make the DDIR server parallel, timer-free, and Corgi-capable (#853)
1 parent 88f4224 commit 74d9370

4 files changed

Lines changed: 39 additions & 8 deletions

File tree

‎interactive/src/corgi/bytes.rs‎

Lines changed: 24 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -177,11 +177,32 @@ mod test {
177177
]
178178
}
179179

180+
/// The declared shapes of each family — what a program's schema would pin (a variant carries
181+
/// only its tag, so a family with sums cannot be pinned from a row).
182+
fn shapes_of(name: &str) -> (corgi::Shape, corgi::Shape) {
183+
use corgi::Shape::{List, Prim, Prod, Sum};
184+
let u = || Prim(64);
185+
let pair = || Prod(vec![u(), u()]);
186+
match name {
187+
"scalars" => (u(), u()),
188+
"tuples" => (pair(), Prod(vec![u()])),
189+
"lists" => (u(), List(Box::new(u()))),
190+
"variants" => (Sum(vec![u(), pair()]), u()),
191+
"nested" => (Prod(vec![u(), Sum(vec![u(), pair()])]), List(Box::new(u()))),
192+
other => panic!("no shapes for family {other}"),
193+
}
194+
}
195+
196+
fn container_of(name: &str, updates: Vec<((DValue, DValue), Time, Diff)>) -> CorgiContainer<Time, Diff> {
197+
let (k, v) = shapes_of(name);
198+
CorgiContainer::from_updates(updates, &k, &v)
199+
}
200+
180201
/// Every update survives the round trip, with its time and diff, for every shape family.
181202
#[test]
182203
fn round_trip_preserves_updates() {
183204
for (name, updates) in shape_families() {
184-
let c = CorgiContainer::<Time, Diff>::from_updates(updates.clone());
205+
let c = container_of(name, updates.clone());
185206
let back = round_trip(&c);
186207
assert_eq!(back.into_updates(), updates, "{name} did not survive the round trip");
187208
}
@@ -202,7 +223,7 @@ mod test {
202223
#[test]
203224
fn round_trip_preserves_column_hashes() {
204225
for (name, updates) in shape_families() {
205-
let c = CorgiContainer::<Time, Diff>::from_updates(updates);
226+
let c = container_of(name, updates);
206227
let (keys, vals) = (c.keys.clone(), c.vals.clone());
207228
let back = round_trip(&c);
208229
assert_eq!(corgi::arrange::hash_rows(&back.keys), corgi::arrange::hash_rows(&keys), "{name} keys");
@@ -218,7 +239,7 @@ mod test {
218239
let updates: Vec<_> = (0..1000i64)
219240
.map(|i| ((DValue::Int(i), DValue::Int(i * 2)), time(0, &[]), 1))
220241
.collect();
221-
let c = CorgiContainer::<Time, Diff>::from_updates(updates);
242+
let c = CorgiContainer::<Time, Diff>::from_updates_pinned(updates);
222243
assert_eq!(corgi::bytes::length_in_bytes(&c.keys), 24 + 8 * 1000);
223244
assert_eq!(corgi::bytes::length_in_bytes(&c.vals), 24 + 8 * 1000);
224245
}

‎interactive/src/corgi/container.rs‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,16 @@ impl<T: Clone + 'static, R: Clone + 'static> CorgiContainer<T, R> {
6969
CorgiContainer { keys: transcode(&keys_rows, kshape), vals: transcode(&vals_rows, vshape), times, diffs }
7070
}
7171

72+
/// Test convenience: build a container from row updates, pinning the shapes from the first
73+
/// row (what the ingest operator does with the first batch it sees).
74+
#[cfg(test)]
75+
pub(crate) fn from_updates_pinned(updates: Vec<((Row, Row), T, R)>) -> Self {
76+
use crate::corgi::logic::shape_of_row;
77+
let Some(((k, v), _, _)) = updates.first() else { return Self::default() };
78+
let (ks, vs) = (shape_of_row(k).unwrap(), shape_of_row(v).unwrap());
79+
Self::from_updates(updates, &ks, &vs)
80+
}
81+
7282
/// Read the container back to DDIR row updates — the **egress boundary** transcode (once).
7383
/// corgi `Value` is self-describing, so shapes come from `shape_of_value`.
7484
pub fn into_updates(self) -> Vec<((Row, Row), T, R)> {

‎interactive/src/corgi/exchange.rs‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -111,7 +111,7 @@ impl<T: Clone + 'static, R: Clone + 'static> Distributor<CorgiContainer<T, R>> f
111111
}
112112
let peers = pushers.len();
113113

114-
let ids = corgi::hash(&container.keys).into_u64("corgi exchange: key hash");
114+
let ids = corgi::hash(&container.keys).into_u64("corgi exchange: key hash").unwrap();
115115
self.counting_sort(&ids, peers);
116116

117117
// Whole-container fast path. When every row shares a destination — a batch narrower than
@@ -213,7 +213,7 @@ mod test {
213213

214214
/// Partition `updates` across `peers` destinations and read each destination back as rows.
215215
fn partition(updates: Vec<((DValue, DValue), Time, Diff)>, peers: usize) -> Vec<Vec<((DValue, DValue), Time, Diff)>> {
216-
let mut container = CorgiContainer::<Time, Diff>::from_updates(updates);
216+
let mut container = CorgiContainer::<Time, Diff>::from_updates_pinned(updates);
217217
let mut pushers: Vec<Collect<Time, Diff>> = (0..peers).map(|_| Collect::default()).collect();
218218
let mut distributor = CorgiDistributor::<Time, Diff>::default();
219219
distributor.partition(&mut container, &timely::progress::Stamp::from_elem(0u64), &mut pushers);
@@ -299,7 +299,7 @@ mod test {
299299
/// either way so no row is sent twice.
300300
#[test]
301301
fn the_input_is_consumed() {
302-
let mut container = CorgiContainer::<Time, Diff>::from_updates(scalar_updates(50));
302+
let mut container = CorgiContainer::<Time, Diff>::from_updates_pinned(scalar_updates(50));
303303
let mut pushers: Vec<Collect<Time, Diff>> = (0..4).map(|_| Collect::default()).collect();
304304
let mut distributor = CorgiDistributor::<Time, Diff>::default();
305305
distributor.partition(&mut container, &timely::progress::Stamp::from_elem(0u64), &mut pushers);

‎interactive/src/corgi/reduce.rs‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,7 @@ type CBatch<T> = Rc<ChunkBatch<CorgiChunk<T, Diff>>>;
5353
/// columns, allowing Corgi to swizzle their buffers in place when unshared.
5454
fn signed_order_view(value: CValue) -> CValue {
5555
match value {
56-
value @ CValue::Prim(_) => NumOp::from(ArithOp::ToSigned).eval(value),
56+
value @ CValue::Prim(_) => NumOp::from(ArithOp::ToSigned).eval(value).expect("ToSigned on a leaf"),
5757
CValue::Prod(fields) => {
5858
CValue::Prod(fields.into_iter().map(signed_order_view).collect())
5959
}
@@ -62,7 +62,7 @@ fn signed_order_view(value: CValue) -> CValue {
6262
within,
6363
variants
6464
.into_iter()
65-
.map(|variant| variant.map(signed_order_view))
65+
.map(signed_order_view)
6666
.collect(),
6767
),
6868
CValue::List(bounds, values) => {

0 commit comments

Comments
 (0)