diff --git a/interactive/server/README.md b/interactive/server/README.md index e3dcc9881..9e5bc7209 100644 --- a/interactive/server/README.md +++ b/interactive/server/README.md @@ -118,6 +118,16 @@ every replay. Identity is convention, not enforcement: we are not defending against adversarial clients yet, and server-side attribution is deliberately deferred until a deployment needs it. +## Requests as input rows + +Request parameters can arrive as ordinary input rows. For example, a row +`(request_id, node)` can join a named graph export by `node`, retaining +`request_id` in the output. +Insert the request, `tick`, `peek` its answer, then retract the request and +`tick` again. An empty result is complete when the tick completes; it does not +need an output row as a completion marker. Keeping the binding instead gives +a standing query whose answer follows later graph changes. + ## Feedback: `bind` bind unbind diff --git a/interactive/server/tests/multiworker.rs b/interactive/server/tests/multiworker.rs index 3ff80f4eb..0c23dfefa 100644 --- a/interactive/server/tests/multiworker.rs +++ b/interactive/server/tests/multiworker.rs @@ -77,7 +77,7 @@ fn request_observing( } } -fn assert_backend(backend: &str) { +fn start_server(backend: &str, workers: usize) -> (ServerProcess, TcpStream, BufReader) { let listeners: Vec<_> = (0..3) .map(|_| TcpListener::bind("127.0.0.1:0").unwrap()) .collect(); @@ -88,7 +88,7 @@ fn assert_backend(backend: &str) { drop(listeners); let child = Command::new(env!("CARGO_BIN_EXE_ddir_server")) - .env("DDIR_WORKERS", "4") + .env("DDIR_WORKERS", workers.to_string()) .env("DDIR_BACKEND", backend) // The retired polling/tick knob must not reintroduce wall-clock // progress if it remains in an old deployment environment. @@ -116,8 +116,13 @@ fn assert_backend(backend: &str) { stream .set_read_timeout(Some(Duration::from_secs(10))) .unwrap(); - let mut writer = stream.try_clone().unwrap(); - let mut reader = BufReader::new(stream); + let writer = stream.try_clone().unwrap(); + let reader = BufReader::new(stream); + (server, writer, reader) +} + +fn assert_backend(backend: &str) { + let (server, mut writer, mut reader) = start_server(backend, 4); request( &mut writer, @@ -209,3 +214,133 @@ fn vec_commands_replay_on_four_workers_without_duplicating_input() { fn corgi_commands_replay_on_four_workers_without_duplicating_input() { assert_backend("corgi"); } + +/// Two consumers join request rows against the same named graph through TCP, +/// while both the requests and graph change. +fn assert_shared_import_requests(backend: &str, workers: usize) { + let (server, mut writer, mut reader) = start_server(backend, workers); + request( + &mut writer, + &mut reader, + "g", + "g load graph begin\nlet e = input 0;\nexport \"edges\" = e | arrange;\ng end-load\n", + ); + for name in ["a", "b"] { + request( + &mut writer, + &mut reader, + name, + &format!( + "{name} load {name} begin\nlet q = input 0;\nlet e = import \"edges\";\nexport \"{name}.answer\" = (q | key($0[1] ; $0[0])) | join(e, ($1[0] ; $2[0]));\n{name} end-load\n" + ), + ); + } + request( + &mut writer, + &mut reader, + "f", + "f feed graph 0 begin\n1 val=2\n1 val=3\n2 val=4\nf end-feed\n", + ); + request(&mut writer, &mut reader, "t", "t tick\n"); + for name in ["a", "b"] { + assert!(request( + &mut writer, + &mut reader, + "p", + &format!("p peek {name}.answer\n") + ) + .is_empty()); + } + // Distinct request identities with equal bindings must not collapse. A + // missing key completes normally with no output rows for its request id. + request( + &mut writer, + &mut reader, + "f", + "f feed a 0 begin\n10,1\n11,1\n12,99\nf end-feed\n", + ); + request(&mut writer, &mut reader, "f", "f feed b 0 20,2\n"); + request(&mut writer, &mut reader, "t", "t tick\n"); + let answer = |rid, node| format!("diff=1 key=Tuple([Int({rid})]) val=Tuple([Int({node})])"); + assert_eq!( + request(&mut writer, &mut reader, "p", "p peek a.answer\n"), + vec![answer(10, 2), answer(10, 3), answer(11, 2), answer(11, 3)] + ); + assert_eq!( + request(&mut writer, &mut reader, "p", "p peek b.answer\n"), + vec![answer(20, 4)] + ); + + // Retract one request without removing the other equal binding. Graph + // changes maintain active requests in both independently installed plans. + request(&mut writer, &mut reader, "f", "f feed a 0 10,1 diff=-1\n"); + request( + &mut writer, + &mut reader, + "f", + "f feed graph 0 begin\n1 val=3 diff=-1\n1 val=5\n2 val=4 diff=-1\n2 val=6\nf end-feed\n", + ); + request(&mut writer, &mut reader, "t", "t tick\n"); + assert_eq!( + request(&mut writer, &mut reader, "p", "p peek a.answer\n"), + vec![answer(11, 2), answer(11, 5)] + ); + assert_eq!( + request(&mut writer, &mut reader, "p", "p peek b.answer\n"), + vec![answer(20, 6)] + ); + request( + &mut writer, + &mut reader, + "f", + "f feed a 0 begin\n11,1 diff=-1\n12,99 diff=-1\nf end-feed\n", + ); + request(&mut writer, &mut reader, "f", "f feed b 0 20,2 diff=-1\n"); + request(&mut writer, &mut reader, "t", "t tick\n"); + for name in ["a", "b"] { + assert!(request( + &mut writer, + &mut reader, + "p", + &format!("p peek {name}.answer\n") + ) + .is_empty()); + } + // Reuse an id with a different parameter against the changed graph. No + // program reinstall and no retained answer from its previous binding. + request(&mut writer, &mut reader, "f", "f feed a 0 10,2\n"); + request(&mut writer, &mut reader, "t", "t tick\n"); + assert_eq!( + request(&mut writer, &mut reader, "p", "p peek a.answer\n"), + vec![answer(10, 6)] + ); + request(&mut writer, &mut reader, "f", "f feed a 0 10,2 diff=-1\n"); + request(&mut writer, &mut reader, "t", "t tick\n"); + assert!(request(&mut writer, &mut reader, "p", "p peek a.answer\n").is_empty()); + request(&mut writer, &mut reader, "d", "d drop a\n"); + request(&mut writer, &mut reader, "d", "d drop b\n"); + request(&mut writer, &mut reader, "d", "d drop graph\n"); + drop(reader); + drop(writer); + server.stop(); +} + +#[test] +fn vec_requests_follow_shared_graph_changes_on_one_worker() { + assert_shared_import_requests("vec", 1); +} + +#[test] +fn vec_requests_follow_shared_graph_changes_on_four_workers() { + assert_shared_import_requests("vec", 4); +} + +#[test] +fn corgi_requests_follow_shared_graph_changes_on_one_worker() { + assert_shared_import_requests("corgi", 1); +} + +#[test] +fn corgi_requests_follow_shared_graph_changes_on_four_workers() { + assert_shared_import_requests("corgi", 4); +}