use std::io::{BufRead, Write};
use std::process::{Child, Stdio};
use super::*;
const RESPONSE_PREFIX: &str = "UQA_SEQUENCE_POSITION_RESPONSE ";
struct Peer {
child: Child,
responses: std::sync::mpsc::Receiver<String>,
reader: Option<std::thread::JoinHandle<()>>,
}
impl Peer {
fn start(path: &std::path::Path) -> Self {
let mut child = std::process::Command::new(std::env::current_exe().unwrap())
.args([
"--ignored",
"--exact",
"row_locks::cross_process::file::sequence_positions::tests::peer::sequence_position_peer",
"--nocapture",
"--test-threads=1",
])
.env("UQA_SEQUENCE_POSITION_TEST_PATH", path)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::inherit())
.spawn()
.unwrap();
let stdout = child.stdout.take().unwrap();
let (send, responses) = std::sync::mpsc::channel();
let reader = std::thread::spawn(move || {
for line in std::io::BufReader::new(stdout).lines() {
let Ok(line) = line else { break };
if let Some((_, response)) = line.split_once(RESPONSE_PREFIX) {
if send.send(response.to_owned()).is_err() {
break;
}
}
}
});
let mut peer = Self {
child,
responses,
reader: Some(reader),
};
assert_eq!(peer.response(), "ready");
peer
}
fn request(&mut self, request: &str) -> String {
let input = self.child.stdin.as_mut().unwrap();
writeln!(input, "{request}").unwrap();
input.flush().unwrap();
self.response()
}
fn response(&mut self) -> String {
self.responses
.recv_timeout(std::time::Duration::from_secs(60))
.unwrap_or_else(|error| {
panic!(
"sequence position peer response failed: {error}; status {:?}",
self.child.try_wait()
)
})
}
fn close(mut self) {
assert_eq!(self.request("close"), "closed");
self.child.wait().unwrap();
}
fn terminate(&mut self) {
self.child.kill().unwrap();
self.child.wait().unwrap();
}
}
impl Drop for Peer {
fn drop(&mut self) {
let _ = self.child.kill();
let _ = self.child.wait();
if let Some(reader) = self.reader.take() {
let _ = reader.join();
}
}
}
fn respond(value: impl std::fmt::Display) {
println!("{RESPONSE_PREFIX}{value}");
std::io::stdout().flush().unwrap();
}
#[test]
#[ignore = "subprocess entry point for the sequence position tests"]
fn sequence_position_peer() {
let path = std::env::var_os("UQA_SEQUENCE_POSITION_TEST_PATH").unwrap();
let coordinator = FileLockCoordinator::open(std::path::Path::new(&path)).unwrap();
respond("ready");
for line in std::io::stdin().lock().lines() {
let line = line.unwrap();
let mut command = line.split_whitespace();
let operation = command.next().unwrap();
let mut number = || command.next().map(|word| word.parse::<i64>().unwrap());
match operation {
"record" => {
let sequence = number().unwrap() as u32;
let recorded = Locked::new(&coordinator).record(sequence, number().unwrap());
respond(if recorded { "recorded" } else { "full" });
}
"read" => match Locked::new(&coordinator).recorded(number().unwrap() as u32) {
Some((current, fresh)) => respond(format!("{current} {fresh}")),
None => respond("none"),
},
"advance" => {
let sequence = number().unwrap() as u32;
for _ in 0..number().unwrap() {
let positions = Locked::new(&coordinator);
let current = positions
.recorded(sequence)
.map_or(0, |(current, _)| current);
assert!(positions.record(sequence, current + 1));
}
respond("advanced");
}
"close" => break,
_ => panic!("unexpected sequence position peer command {line}"),
}
}
drop(coordinator);
respond("closed");
}
#[test]
fn processes_read_and_advance_the_same_positions() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("shared.db");
let coordinator = FileLockCoordinator::open(&path).unwrap();
assert!(Locked::new(&coordinator).record(1, 10));
let mut peer = Peer::start(&path);
assert_eq!(peer.request("read 1"), "10 true");
assert_eq!(peer.request("record 1 11"), "recorded");
assert_eq!(peer.request("record 2 70"), "recorded");
let positions = Locked::new(&coordinator);
assert_eq!(positions.recorded(1), Some((11, true)));
assert_eq!(positions.recorded(2), Some((70, true)));
positions.remove(2);
drop(positions);
assert_eq!(peer.request("read 2"), "none");
}
#[test]
fn concurrent_processes_never_lose_an_advance() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("advance.db");
let coordinator = FileLockCoordinator::open(&path).unwrap();
let mut peer = Peer::start(&path);
let input = peer.child.stdin.as_mut().unwrap();
writeln!(input, "advance 5 2000").unwrap();
input.flush().unwrap();
for _ in 0..2000 {
let positions = Locked::new(&coordinator);
let current = positions.recorded(5).map_or(0, |(current, _)| current);
assert!(positions.record(5, current + 1));
}
assert_eq!(peer.response(), "advanced");
assert_eq!(Locked::new(&coordinator).recorded(5), Some((4000, true)));
}
#[test]
fn the_positions_of_a_killed_process_stay_while_another_process_is_attached() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("survivor.db");
let coordinator = FileLockCoordinator::open(&path).unwrap();
assert!(Locked::new(&coordinator).record(9, 1));
let mut peer = Peer::start(&path);
assert_eq!(peer.request("record 1 10"), "recorded");
peer.terminate();
assert_eq!(Locked::new(&coordinator).recorded(1), Some((10, true)));
}
#[test]
fn the_positions_of_a_killed_process_are_discarded_when_no_process_remains() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("killed.db");
let mut peer = Peer::start(&path);
assert_eq!(peer.request("record 1 10"), "recorded");
peer.terminate();
let coordinator = FileLockCoordinator::open(&path).unwrap();
assert_eq!(Locked::new(&coordinator).recorded(1), None);
}
#[test]
fn the_positions_of_a_closed_process_are_kept_for_the_next_run() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("closed.db");
let mut peer = Peer::start(&path);
assert_eq!(peer.request("record 1 10"), "recorded");
peer.close();
let coordinator = FileLockCoordinator::open(&path).unwrap();
assert_eq!(Locked::new(&coordinator).recorded(1), Some((10, false)));
}
#[test]
fn the_last_process_to_close_marks_the_positions_of_every_process_clean() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("last.db");
let coordinator = FileLockCoordinator::open(&path).unwrap();
assert!(Locked::new(&coordinator).record(1, 10));
let mut peer = Peer::start(&path);
assert_eq!(peer.request("record 2 20"), "recorded");
drop(coordinator);
assert_eq!(peer.request("read 1"), "10 true");
peer.close();
let coordinator = FileLockCoordinator::open(&path).unwrap();
assert!(coordinator.stored_sequence_header().unwrap().unwrap().clean);
let positions = Locked::new(&coordinator);
assert_eq!(positions.recorded(1), Some((10, false)));
assert_eq!(positions.recorded(2), Some((20, false)));
}