diff --git a/interactive/server/README.md b/interactive/server/README.md index aba99f27e..240e11531 100644 --- a/interactive/server/README.md +++ b/interactive/server/README.md @@ -12,14 +12,17 @@ defaults. Diagnostics are disabled by default so an idle server can park; `DDIR_DIAGNOSTICS=1` enables the diagnostics dataflow and its listener on `DDIR_DIAG_PORT` (default 51371). `DDIR_WORKERS` selects the number of worker threads (default 1), and `DDIR_BACKEND=vec|corgi` selects the renderer for installed -programs (default `vec`). +programs (default `vec`). With the Corgi backend, exports are kept as traces of +Corgi chunks rather than rows (less memory); `peek`, `tail`, `bind` and other +programs' imports still read them as rows. `DDIR_COLUMNAR_EXPORTS=0` keeps row traces. One backend is selected for the whole server. The current registry is a transitional row-speaking bridge: inputs and imports convert from `Value` rows to Corgi columns at a program boundary, and exports convert back before they become shareable traces. A Corgi program stays columnar between those boundaries, but a production Corgi server should replace the bridge with native -columnar inputs and traces. +columnar inputs and traces. Columnar exports are the first step on the export side: +the published trace is columnar, and rows are made only when it is read. Worker 0 admits one FIFO control stream and routes one ordered record to every worker. Small commands are replicated; framed input batches are partitioned by diff --git a/interactive/server/src/loop_.rs b/interactive/server/src/loop_.rs index 2d6af882c..0fb2ae8ec 100644 --- a/interactive/server/src/loop_.rs +++ b/interactive/server/src/loop_.rs @@ -102,6 +102,7 @@ pub fn run_worker( worker: &mut Worker, events: Option>, backend: RenderBackend, + columnar_exports: bool, ) { // Logging a park wakes the diagnostics dataflow, whose scheduling logs can // in turn wake it again. Keep idle servers genuinely idle unless an @@ -146,6 +147,7 @@ pub fn run_worker( let mut work_input = (worker.index() == 0).then_some(input); let mut server = Server::with_backend(backend); + server.set_columnar_exports(columnar_exports); let mut tails: HashMap = HashMap::new(); let mut responses: HashMap> = HashMap::new(); let mut next_token = 0u64; @@ -458,16 +460,14 @@ fn start_tail( return Err(format!("tail reqid {:?} is already active", reqid)); } let mut trace = server - .trace(name) + .published(name) .ok_or_else(|| format!("no trace {:?}", name))?; let dataflow_id = worker.next_dataflow_index(); let tag = reqid.to_string(); let mut probe = ProbeHandle::new(); let shutdown = worker.dataflow::(|scope| { - let (arranged, shutdown) = trace.import_core(scope.clone(), "TailImport"); - arranged - .as_collection(|key, val| (key.clone(), val.clone())) - .inner + let (rows, shutdown) = trace.import_rows(scope.clone(), "TailImport"); + rows .exchange(|_| 0u64) .inspect(move |((key, val), time, diff)| { if let Some(response) = response.as_ref() { diff --git a/interactive/server/src/main.rs b/interactive/server/src/main.rs index 7502d3a8d..3498df951 100644 --- a/interactive/server/src/main.rs +++ b/interactive/server/src/main.rs @@ -121,6 +121,9 @@ fn main() { .unwrap_or_else(|_| "vec".to_string()) .parse::() .unwrap_or_else(|error| panic!("DDIR_BACKEND: {error}")); + // Corgi-backend exports stay columnar (less memory; readers see rows); DDIR_COLUMNAR_EXPORTS=0 + // keeps the row traces instead. + let columnar_exports = std::env::var("DDIR_COLUMNAR_EXPORTS").map_or(true, |v| v != "0"); let (event_tx, event_rx) = channel::(); let (activation_tx, activation_rx) = sync_channel(1); @@ -144,7 +147,7 @@ fn main() { } else { None }; - control_loop::run_worker(worker, events, backend); + control_loop::run_worker(worker, events, backend, columnar_exports); // The live server intentionally keeps installed dataflows and // diagnostics around. Once every worker observes shutdown, remove diff --git a/interactive/src/backend/corgi.rs b/interactive/src/backend/corgi.rs index 7c09c4e3a..29a5fa1e6 100644 --- a/interactive/src/backend/corgi.rs +++ b/interactive/src/backend/corgi.rs @@ -343,12 +343,12 @@ pub fn render_tree<'s>( /// converts back (`FromCorgi`). Signature-compatible with /// [`vec::render_tree`](crate::backend::vec::render_tree) (hence the `vec::Col` alias), so a /// row-speaking driver switches backends by switching this one call. -pub fn render_tree_rows<'s>( +pub fn render_tree_corgi<'s>( s: &st::Scope, scope: Scope<'s, Time>, depth: usize, imports: Vec>, -) -> Vec> { +) -> Vec> { let corgi_imports: Vec> = crate::backend::vec::check_import_shapes(s, imports) .into_iter() .zip(&s.imports) @@ -380,6 +380,115 @@ pub fn render_tree_rows<'s>( }) .collect(); render_tree(s, scope, depth, corgi_imports) +} + +/// A columnar export: the program's corgi collection, left to the host time and arranged as +/// corgi chunks. Readers convert to rows only when they read ([`export_rows`]). +pub type ExportTrace = TraceAgent>>; + +/// Arrange an export (already at host time) as a columnar trace. +pub fn arrange_export<'s>(c: Collection<'s, u64, CorgiContainer>) -> Arranged<'s, ExportTrace> { + arrange_core::<_, CorgiContainer, _, differential_dataflow::trace::chunk::ChunkSpine>>( + c.inner, + CorgiPact, + "CorgiExport", + ChunkBatcher::, _>::new, + ) +} + +/// The batches an export's arrangement has produced on this worker, not yet taken. +pub type ExportTap = std::rc::Rc::Batch>>>; + +/// The updates in `batches`, as rows. +pub fn batch_rows(batches: Vec<::Batch>) -> Vec<((Row, Row), u64, Diff)> { + let mut out = Vec::new(); + for batch in batches { + for ch in batch.chunks.iter().filter(|c| c.len() > 0) { + let c = CorgiContainer { + keys: recover_key(ch.keys()), + vals: ch.vals().clone(), + times: ch.times().clone(), + diffs: ch.diffs().to_vec(), + }; + out.extend(c.into_updates()); + } + } + out +} + +/// Decode and consume at most `rows_per_batch` export rows at a time. +/// Row payloads (including nested lists) can still vary in size. +pub fn for_each_batch_rows( + batches: Vec<::Batch>, + rows_per_batch: usize, + mut consume: impl FnMut(Vec<((Row, Row), u64, Diff)>), +) { + assert!(rows_per_batch > 0, "export row batch size must be positive"); + let mut indices = Vec::new(); + for batch in batches { + for ch in batch.chunks.iter().filter(|c| c.len() > 0) { + let keys = recover_key(ch.keys()); + for start in (0..ch.len()).step_by(rows_per_batch) { + let end = start.saturating_add(rows_per_batch).min(ch.len()); + indices.clear(); + indices.extend(start..end); + let c = CorgiContainer { + keys: corgi::arrange::gather(&keys, &indices), + vals: corgi::arrange::gather(ch.vals(), &indices), + times: ch.times().gather(&indices), + diffs: ch.diffs()[start..end].to_vec(), + }; + consume(c.into_updates()); + } + } + } +} + +/// Record each batch the arrangement emits into `tap`, passing the stream through. +pub fn tap_export<'s>(a: &Arranged<'s, ExportTrace>, tap: ExportTap) -> timely::dataflow::Stream<'s, u64, Vec::Batch>>> { + a.stream.clone().unary(Pipeline, "ExportTap", move |_, _| { + move |input, output| { + input.for_each(|cap, data| { + tap.borrow_mut().extend(data.iter().filter_map(|b| b.inner.clone())); + output.session(&cap).give_container(data); + }); + } + }) +} + +/// The rows of an imported columnar export, as `((key, val), time, diff)` updates. +pub fn export_rows<'s>( + a: Arranged<'s, ExportTrace>, +) -> timely::dataflow::Stream<'s, u64, Vec<((Row, Row), u64, Diff)>> { + a.stream.unary(Pipeline, "ExportRows", |_, _| { + |input, output| { + input.for_each(|cap, data| { + let mut session = output.session(&cap); + for batch in data.iter() { + let Some(payload) = batch.inner.as_ref() else { continue }; + for ch in payload.chunks.iter().filter(|c| c.len() > 0) { + let c = CorgiContainer { + keys: recover_key(ch.keys()), + vals: ch.vals().clone(), + times: ch.times().clone(), + diffs: ch.diffs().to_vec(), + }; + session.give_container(&mut c.into_updates()); + } + } + }); + } + }) +} + +/// [`render_tree_corgi`] with each export converted back to rows (`FromCorgi`). +pub fn render_tree_rows<'s>( + s: &st::Scope, + scope: Scope<'s, Time>, + depth: usize, + imports: Vec>, +) -> Vec> { + render_tree_corgi(s, scope, depth, imports) .into_iter() .map(|c| { c.inner diff --git a/interactive/src/corgi/container.rs b/interactive/src/corgi/container.rs index 2ae94094b..af0a2c410 100644 --- a/interactive/src/corgi/container.rs +++ b/interactive/src/corgi/container.rs @@ -147,3 +147,16 @@ where CorgiContainer { keys, vals, times: self.times, diffs } } } + + +/// Leaving the program's iterative scope: each time keeps its outer coordinate (lane 0). +impl differential_dataflow::collection::containers::Leave + for CorgiContainer +{ + type OuterContainer = CorgiContainer; + fn leave(self) -> CorgiContainer { + let len = self.times.len(); + let lane = if self.times.width() == 0 { vec![0; len] } else { self.times.lane(0) }; + CorgiContainer { keys: self.keys, vals: self.vals, times: ColTimes::from_raw_lanes(vec![lane], len), diffs: self.diffs } + } +} diff --git a/interactive/src/server.rs b/interactive/src/server.rs index 62c242b53..499bdc742 100644 --- a/interactive/src/server.rs +++ b/interactive/src/server.rs @@ -106,6 +106,38 @@ impl std::str::FromStr for RenderBackend { /// arranged by key at the host time so any later install can `import` it. pub type ServerTrace = TraceAgent>; +/// A published export as its readers see it: a row trace, or with columnar exports +/// ([`Server::set_columnar_exports`]) a trace of corgi chunks. Either reads as rows. +#[derive(Clone)] +pub enum Published { + Rows(ServerTrace), + Columnar(crate::backend::corgi::ExportTrace), +} + +impl Published { + /// Import into `scope` as a stream of `((key, val), time, diff)` rows, with the + /// import's shutdown button. A columnar trace converts to rows as it is read. + pub fn import_rows<'s>( + &mut self, + scope: timely::dataflow::Scope<'s, OuterTime>, + name: &str, + ) -> ( + timely::dataflow::Stream<'s, OuterTime, Vec<((Value, Value), OuterTime, Diff)>>, + ShutdownButton>, + ) { + match self { + Published::Rows(trace) => { + let (arranged, shutdown) = trace.import_core(scope, name); + (arranged.as_collection(|k, v| (k.clone(), v.clone())).inner, shutdown) + } + Published::Columnar(trace) => { + let (arranged, shutdown) = trace.import_core(scope, name); + (crate::backend::corgi::export_rows(arranged), shutdown) + } + } + } +} + /// An input handle into an installed program's positional `input N`. type ServerInput = InputSession; @@ -353,6 +385,16 @@ struct Binding { pub struct Server { /// Published export name -> shareable trace. traces: HashMap, + /// With columnar exports: published export name -> trace of corgi chunks. Readers + /// (imports, binds, snapshots, subscriptions) see rows, converted as they read. + ctraces: HashMap, + /// With export taps: each columnar export's new batches on this worker, drained by + /// `take_changes` / `for_each_change_batch`: a standing change stream without a snapshot. + taps: HashMap, + /// `set_columnar_exports`. + columnar_exports: bool, + /// `set_export_taps`. + export_taps: bool, /// Installed program name -> its handles and lifecycle bookkeeping. programs: HashMap, /// Trace name -> number of installed programs importing it (the drop gate). @@ -376,6 +418,10 @@ impl Server { pub fn with_backend(backend: RenderBackend) -> Self { Server { traces: HashMap::new(), + ctraces: HashMap::new(), + taps: HashMap::new(), + columnar_exports: true, + export_taps: false, programs: HashMap::new(), importers: HashMap::new(), bindings: Vec::new(), @@ -389,11 +435,37 @@ impl Server { self.epoch } + /// Whether `name` (canonical) is published, as a row or a columnar trace. + fn is_published(&self, name: &str) -> bool { + self.traces.contains_key(name) || self.ctraces.contains_key(name) + } + + /// Keep corgi-backend exports columnar (default on): a program's exports leave its scope as + /// corgi containers and are arranged as corgi chunks, and readers convert to rows only when + /// they read (`snapshot`, imports, binds). Affects programs installed afterwards. + pub fn set_columnar_exports(&mut self, on: bool) { + self.columnar_exports = on; + } + + /// With columnar exports, also keep each export's new batches for [`Server::take_changes`] + /// and [`Server::for_each_change_batch`] (default off). Affects programs installed afterwards. + pub fn set_export_taps(&mut self, on: bool) { + self.export_taps = on; + } + /// Clone a trace reader for a transient peek or subscription dataflow. pub fn trace(&self, name: &str) -> Option { self.traces.get(&canonical_source_name(name)).cloned() } + /// A reader for published trace `name`, row or columnar, for a transient peek or + /// subscription dataflow. + pub fn published(&self, name: &str) -> Option { + let name = canonical_source_name(name); + self.traces.get(&name).cloned().map(Published::Rows) + .or_else(|| self.ctraces.get(&name).cloned().map(Published::Columnar)) + } + /// Return registry state without coupling a caller to stdout formatting. pub fn program_info(&self) -> Vec { let mut result: Vec<_> = self @@ -455,7 +527,7 @@ impl Server { for imp in &prog.root.imports { if let st::Source::Trace(t) = &imp.from { let key = canonical_source_name(t); - if !self.traces.contains_key(&key) { + if !self.is_published(&key) { if key != "clock" && Recipe::parse(&key).is_none() { return Err(format!( "program {:?} imports unknown trace {:?}; install its producer first", @@ -467,7 +539,7 @@ impl Server { } } for e in &prog.root.exports { - if self.traces.contains_key(&e.name) || generated.contains(&e.name) { + if self.is_published(&e.name) || generated.contains(&e.name) { return Err(format!("export name {:?} is already published; choose another name or drop its producer", e.name)); } } @@ -498,11 +570,15 @@ impl Server { let root = &prog.root; let traces = &mut self.traces; let backend = self.backend; + let columnar = backend == RenderBackend::Corgi && self.columnar_exports; + let tapped = columnar && self.export_taps; + let ctraces = &mut self.ctraces; // The id this dataflow will get; captured so `drop` can remove it. let dataflow_id = worker.next_dataflow_index(); - let (published, inputs): (Vec<(String, ServerTrace)>, Vec<(usize, ServerInput)>) = + #[allow(clippy::type_complexity)] + let (published, cpublished, inputs): (Vec<(String, ServerTrace)>, Vec<(String, crate::backend::corgi::ExportTrace, Option)>, Vec<(usize, ServerInput)>) = worker.dataflow::(|outer| { let mut inputs: Vec<(usize, ServerInput)> = Vec::new(); @@ -519,6 +595,11 @@ impl Server { st::Source::Trace(t) => { // The first binding point: resolve a named trace by importing it. let key = canonical_source_name(t); + // A columnar export is read as rows by its importers. + if let Some(ctrace) = ctraces.get_mut(&key) { + let rows = crate::backend::corgi::export_rows(ctrace.import(outer.clone())); + return differential_dataflow::AsCollection::as_collection(rows); + } let arranged = traces .get_mut(&key) .expect("validated above") @@ -531,21 +612,41 @@ impl Server { // Render the program body in its own iterative scope, then bring // every export back out to the host time. - let leaved: Vec> = outer + let (leaved, cleaved) = outer .iterative::, _, _>(|inner| { let entered: Vec<_> = outer_cols.iter().map(|c| c.clone().enter(inner)).collect(); + if columnar { + let exports = crate::backend::corgi::render_tree_corgi(root, inner.clone(), 0, entered); + return (Vec::new(), exports.into_iter().map(|c| c.leave(outer)).collect::>()); + } let exports = match backend { RenderBackend::Vec => render_tree(root, inner.clone(), 0, entered), RenderBackend::Corgi => { render_tree_rows(root, inner.clone(), 0, entered) } }; - exports + (exports .into_iter() .map(|c| c.leave(outer)) - .collect::>() + .collect::>(), Vec::new()) }); + let cpublished: Vec<_> = root + .exports + .iter() + .zip(cleaved) + .map(|(e, col)| { + use timely::dataflow::operators::Probe; + let arranged = crate::backend::corgi::arrange_export(col); + let tap: Option = tapped.then(Default::default); + let stream = match &tap { + Some(t) => crate::backend::corgi::tap_export(&arranged, t.clone()), + None => arranged.stream.clone(), + }; + stream.probe_with(&probe); + (e.name.clone(), arranged.trace, tap) + }) + .collect(); // The second binding point: probe and publish each export's trace. let published: Vec<(String, ServerTrace)> = root @@ -560,12 +661,16 @@ impl Server { }) .collect(); - (published, inputs) + (published, cpublished, inputs) }); for (export_name, trace) in published { self.traces.insert(export_name, trace); } + for (export_name, trace, tap) in cpublished { + if let Some(tap) = tap { self.taps.insert(export_name.clone(), tap); } + self.ctraces.insert(export_name, trace); + } for t in &import_names { *self.importers.entry(t.clone()).or_insert(0) += 1; } @@ -862,7 +967,7 @@ impl Server { ) -> Result<(), String> { let source = canonical_source_name(trace); let target = prog.to_string(); - if !self.traces.contains_key(&source) { + if !self.is_published(&source) { return Err(format!("no trace {:?}", source)); } let installed = self @@ -890,11 +995,11 @@ impl Server { let buffer_in = buffer.clone(); let mut probe = ProbeHandle::new(); let dataflow_id = worker.next_dataflow_index(); - let trace_handle = self.traces.get_mut(&source).expect("checked above"); + let mut trace = self.published(&source).expect("checked above"); let shutdown = worker.dataflow::(|scope| { - let (arranged, shutdown) = trace_handle.import_core(scope.clone(), "BindImport"); - arranged - .as_collection(|k, v| (k.clone(), v.clone())) + use timely::dataflow::operators::{Inspect, Probe}; + let (rows, shutdown) = trace.import_rows(scope.clone(), "BindImport"); + rows .inspect(move |((key, val), _time, diff)| { buffer_in .borrow_mut() @@ -970,18 +1075,15 @@ impl Server { let name = canonical_source_name(name); let epoch = self.epoch; - let mut trace = self - .trace(&name) - .ok_or_else(|| format!("no trace {:?}", name))?; + let mut trace = self.published(&name).ok_or_else(|| format!("no trace {:?}", name))?; let acc: Rc>> = Rc::new(RefCell::new(HashMap::new())); let acc_in = acc.clone(); let mut probe = ProbeHandle::new(); let id = worker.next_dataflow_index(); worker.dataflow::(|scope| { - trace - .import(scope.clone()) - .as_collection(|k, v| (k.clone(), v.clone())) - .inner + // The dataflow is dropped once drained, which releases the import. + let (rows, _shutdown) = trace.import_rows(scope.clone(), "SnapshotImport"); + rows .exchange(|_| 0u64) .inspect(move |((k, v), t, d)| { if *t < epoch { @@ -1008,6 +1110,32 @@ impl Server { Ok(rows) } + /// The changes to export `name` that this worker's share of its arrangement has produced + /// since the last call, as rows (with [`Server::set_export_taps`]). Each worker drains its own. + pub fn take_changes(&mut self, name: &str) -> Vec<((Value, Value), OuterTime, Diff)> { + match self.taps.get(name) { + Some(tap) => crate::backend::corgi::batch_rows(std::mem::take(&mut *tap.borrow_mut())), + None => Vec::new(), + } + } + + /// Drain this worker's tapped export changes in row batches no larger than + /// `rows_per_batch`. Callbacks run in export order; rows are consumed once. + /// Unlike `take_changes`, this does not materialize the entire export as rows. + /// The limit counts rows, not bytes (nested payloads may be large). + pub fn for_each_change_batch( + &mut self, + name: &str, + rows_per_batch: usize, + consume: impl FnMut(Vec<((Value, Value), OuterTime, Diff)>), + ) { + assert!(rows_per_batch > 0, "export row batch size must be positive"); + if let Some(tap) = self.taps.get(name) { + let batches = std::mem::take(&mut *tap.borrow_mut()); + crate::backend::corgi::for_each_batch_rows(batches, rows_per_batch, consume); + } + } + /// Drop installed program `name`, releasing its dataflow immediately. /// /// Refuses (changing nothing) if any trace the program publishes still has a @@ -1046,6 +1174,8 @@ impl Server { } for ex in &installed.exports { self.traces.remove(ex); + self.ctraces.remove(ex); + self.taps.remove(ex); } let id = installed.dataflow_id; // Drop the input handles first (closes the inputs while the operators @@ -1167,6 +1297,10 @@ impl Server { trace.set_logical_compaction(frontier.borrow()); trace.set_physical_compaction(frontier.borrow()); } + for trace in self.ctraces.values_mut() { + trace.set_logical_compaction(frontier.borrow()); + trace.set_physical_compaction(frontier.borrow()); + } } } diff --git a/interactive/tests/server_change_batches.rs b/interactive/tests/server_change_batches.rs new file mode 100644 index 000000000..1ee83ed75 --- /dev/null +++ b/interactive/tests/server_change_batches.rs @@ -0,0 +1,52 @@ +//! Bounded export consumption must preserve inserts/retractions and nested rows. +use interactive::ir::Value; +use interactive::server::{InputUpdate, RenderBackend, Server}; +use interactive::{lower, parse}; + +#[test] +fn bounded_exports_match_full_drain_and_are_consumed_once() { + timely::execute_directly(|worker| { + let mut program = lower::lower_tree(parse::pipe::parse(r#" + type Three = A u64 | B u64 | C u64; + let pairs = input 0 | key($0[0] ; $0[1]); + export "plain" = pairs; + export "lists" = pairs | collect; + export "tagged" = pairs | map($0 ; variant(Three, 2, $1[0])); + "#)); + program.optimize(); + let mut full = Server::with_backend(RenderBackend::Corgi); + let mut bounded = Server::with_backend(RenderBackend::Corgi); + for server in [&mut full, &mut bounded] { + server.set_columnar_exports(true); + server.set_export_taps(true); + } + full.install(worker, "full", &program).unwrap(); + bounded.install(worker, "bounded", &program).unwrap(); + for (epoch, limit) in [1, 7, 128, 4096].into_iter().enumerate() { + let diff = if epoch % 2 == 0 { 1 } else { -1 }; + let input: Vec<_> = (0..257).map(|i| InputUpdate { + key: Value::Tuple(vec![Value::Int(i % 11), Value::Int(i - 128)]), + val: Value::unit(), diff, + }).collect(); + full.feed_batch("full", 0, input.clone()).unwrap(); + bounded.feed_batch("bounded", 0, input).unwrap(); + full.tick(worker); + bounded.tick(worker); + for name in ["plain", "lists", "tagged"] { + let mut expected = full.take_changes(name); + let mut actual = Vec::new(); + bounded.for_each_change_batch(name, limit, |rows| { + assert!(!rows.is_empty() && rows.len() <= limit); + actual.extend(rows); + }); + assert!(!expected.is_empty(), "{name}, epoch {epoch}"); + expected.sort(); + actual.sort(); + assert_eq!(actual, expected, "{name}, epoch {epoch}, limit {limit}"); + bounded.for_each_change_batch(name, limit, |_| panic!("changes drained twice")); + assert!(bounded.take_changes(name).is_empty()); + } + } + bounded.for_each_change_batch("missing", 7, |_| panic!("missing export has changes")); + }); +} diff --git a/interactive/tests/server_columnar_exports.rs b/interactive/tests/server_columnar_exports.rs new file mode 100644 index 000000000..8d0f4d13f --- /dev/null +++ b/interactive/tests/server_columnar_exports.rs @@ -0,0 +1,72 @@ +//! Columnar exports (`Server::set_columnar_exports`) read the same as row exports: +//! through `snapshot`, through another program's `import`, and through `bind`. + +use interactive::ir::Value; +use interactive::server::{RenderBackend, Server}; +use interactive::{lower, parse}; + +fn tup(fields: &[i64]) -> Value { + Value::Tuple(fields.iter().map(|&n| Value::Int(n)).collect()) +} + +fn install(server: &mut Server, worker: &mut timely::worker::Worker, name: &str, src: &str) { + let mut program = lower::lower_tree(parse::pipe::parse(src)); + program.optimize(); + server.install(worker, name, &program).unwrap(); +} + +const PRODUCER: &str = r#" + type Three = A u64 | B u64 | C u64; + let pairs = input 0 | key($0[0] ; $0[1]); + export "plain" = pairs; + export "lists" = pairs | collect; + export "tagged" = pairs | map($0 ; variant(Three, 2, $1[0])); +"#; + +const CONSUMER: &str = r#" + let p = import "plain"; + export "shifted" = p | map($0 ; $1[0] + 1000); +"#; + +const COUNTER: &str = r#" + let seed = input 0; + let feedback = input 1; + let state = seed + feedback; + export "count" = state; + export "next" = (state | map($0[0] + 1 ;)) + (seed | negate); +"#; + +/// Everything the server can read back, after the same inputs and ticks. +fn run(worker: &mut timely::worker::Worker, columnar: bool) -> Vec> { + let mut server = Server::with_backend(RenderBackend::Corgi); + server.set_columnar_exports(columnar); + install(&mut server, worker, "producer", PRODUCER); + install(&mut server, worker, "consumer", CONSUMER); + install(&mut server, worker, "counter", COUNTER); + server.feed("counter", 0, tup(&[0]), Value::unit(), None, 1).unwrap(); + server.bind(worker, "next", "counter", 1).unwrap(); + let mut seen = Vec::new(); + for epoch in 0..4i64 { + for i in 0..40i64 { + let diff = if (i + epoch) % 3 == 0 { -1 } else { 1 }; + server.feed("producer", 0, tup(&[i % 7, i * epoch]), Value::unit(), None, diff).unwrap(); + } + server.tick(worker); + for name in ["plain", "lists", "tagged", "shifted", "count"] { + seen.push(server.snapshot(worker, name).unwrap()); + } + } + seen +} + +#[test] +fn columnar_exports_read_like_row_exports() { + timely::execute_directly(|worker| { + let rows = run(worker, false); + let columnar = run(worker, true); + assert!(rows.iter().any(|s| !s.is_empty())); + // The counter advanced through the binding in both. + assert_eq!(rows.last(), Some(&vec![(tup(&[3]), Value::unit(), 1)])); + assert_eq!(rows, columnar); + }); +}