Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions interactive/server/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 <trace> <prog> <in#> unbind <trace> <prog> <in#>
Expand Down
143 changes: 139 additions & 4 deletions interactive/server/tests/multiworker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,7 @@ fn request_observing(
}
}

fn assert_backend(backend: &str) {
fn start_server(backend: &str, workers: usize) -> (ServerProcess, TcpStream, BufReader<TcpStream>) {
let listeners: Vec<_> = (0..3)
.map(|_| TcpListener::bind("127.0.0.1:0").unwrap())
.collect();
Expand All @@ -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.
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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);
}
Loading