diff --git a/interactive/src/server.rs b/interactive/src/server.rs index fe7c0ae30..62c242b53 100644 --- a/interactive/src/server.rs +++ b/interactive/src/server.rs @@ -994,7 +994,8 @@ impl Server { .probe_with(&mut probe); }); while probe.less_than(&epoch) { - worker.step(); + // Remote data or progress wakes idle workers; do not busy-poll. + worker.step_or_park(None); } worker.drop_dataflow(id); let mut rows: Vec<_> = acc @@ -1118,7 +1119,8 @@ impl Server { while self.programs.values().any(|p| p.probe.less_than(&epoch)) || self.bindings.iter().any(|b| b.probe.less_than(&epoch)) { - worker.step(); + // Remote data or progress wakes idle workers; do not busy-poll. + worker.step_or_park(None); } // Feedback: deliver each binding's buffered source changes into its diff --git a/interactive/tests/server_batch.rs b/interactive/tests/server_batch.rs index 588672ac1..fdcc36d29 100644 --- a/interactive/tests/server_batch.rs +++ b/interactive/tests/server_batch.rs @@ -42,3 +42,45 @@ fn feed_batch_validates_before_staging_one_epoch() { ); }); } + +/// Ticks and transient snapshots must wake for delayed remote rows and progress. +#[test] +fn delayed_peers_wake_ticks_and_snapshots() { + use interactive::server::RenderBackend; + use std::sync::{Arc, Barrier}; + use std::time::Duration; + + for backend in [RenderBackend::Vec, RenderBackend::Corgi] { + let barrier = Arc::new(Barrier::new(3)); + let guards = timely::execute(timely::Config::process(3), move |worker| { + let mut program = lower::lower_tree(parse::pipe::parse( + "let rows = input 0; export \"rows\" = rows;", + )); + program.optimize(); + let mut server = Server::with_backend(backend); + server.install(worker, "world", &program).unwrap(); + for epoch in 0..6 { + barrier.wait(); + if worker.index() == epoch % 3 { + std::thread::sleep(Duration::from_millis(20)); + let diff = if epoch < 3 { 1 } else { -1 }; + server.feed_batch("world", 0, vec![InputUpdate { + key: tup(&[(epoch % 3) as i64]), val: Value::unit(), diff, + }]).unwrap(); + } + server.tick(worker); + barrier.wait(); + if worker.index() == (epoch + 1) % 3 { + std::thread::sleep(Duration::from_millis(20)); + } + let rows = server.snapshot(worker, "rows").unwrap(); + let expected: Vec<_> = if worker.index() == 0 { + (0..3).filter(|&i| if epoch < 3 { i <= epoch } else { i > epoch - 3 }) + .map(|i| (tup(&[i as i64]), Value::unit(), 1)).collect() + } else { Vec::new() }; + assert_eq!(rows, expected); + } + }).unwrap(); + for result in guards.join() { result.unwrap(); } + } +}