use std::fs::File;
use std::io;
use std::path::Path;
use std::path::PathBuf;
use std::process::Child;
use std::process::Command;
use std::process::Stdio;
use std::sync::LazyLock;
use std::time::Duration;
use std::time::Instant;
use serde_json::Value;
use serde_json::json;
const WAIT: Duration = Duration::from_secs(30);
const POLL_INTERVAL: Duration = Duration::from_millis(100);
struct Node {
child: Child,
}
impl Drop for Node {
fn drop(&mut self) {
let _ = self.child.kill();
let _ = self.child.wait();
}
}
static KVSTORE_BIN: LazyLock<PathBuf> = LazyLock::new(|| {
let cargo = std::env::var_os("CARGO").expect("cargo sets CARGO for the tests it runs");
let out = Command::new(cargo)
.args(["build", "--package", "ezraft", "--example", "kvstore"])
.arg("--message-format=json-render-diagnostics")
.stderr(Stdio::inherit())
.output()
.expect("failed to run cargo");
assert!(out.status.success(), "building the kvstore example failed");
String::from_utf8(out.stdout)
.expect("cargo emits UTF-8")
.lines()
.filter_map(|line| serde_json::from_str::<Value>(line).ok())
.find(|msg| msg["reason"] == "compiler-artifact" && msg["target"]["name"] == "kvstore")
.and_then(|msg| msg["executable"].as_str().map(PathBuf::from))
.expect("cargo reported no executable for the kvstore example")
});
fn free_addr() -> String {
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
listener.local_addr().unwrap().to_string()
}
fn spawn_node(work_dir: &Path, addr: &str, seed: Option<&str>) -> io::Result<Node> {
let log = File::create(work_dir.join(format!("node-{}.log", addr.replace(':', "-"))))?;
let mut cmd = Command::new(&*KVSTORE_BIN);
cmd.current_dir(work_dir)
.args(["--addr", addr])
.stdout(log.try_clone()?)
.stderr(log)
.stdin(Stdio::null());
if let Some(seed) = seed {
cmd.args(["--seed", seed]);
}
Ok(Node { child: cmd.spawn()? })
}
async fn write(addr: &str, req: Value) -> io::Result<Value> {
let client = reqwest::Client::new();
let deadline = Instant::now() + WAIT;
loop {
let res = client.post(format!("http://{}/api/write", addr)).json(&req).send().await;
match res {
Ok(resp) if resp.status().is_success() => return resp.json().await.map_err(io::Error::other),
_ if Instant::now() < deadline => tokio::time::sleep(POLL_INTERVAL).await,
Ok(resp) => return Err(io::Error::other(format!("write to {} failed: {}", addr, resp.status()))),
Err(e) => return Err(io::Error::other(e)),
}
}
}
async fn wait_for_key(addr: &str, key: &str, expected: Option<&Value>) -> io::Result<()> {
let deadline = Instant::now() + WAIT;
let mut last = None;
loop {
if let Ok(resp) = reqwest::get(format!("http://{}/api/read?key={}", addr, key)).await {
let value = match resp.status() {
reqwest::StatusCode::NOT_FOUND => None,
status if status.is_success() => Some(resp.json().await.map_err(io::Error::other)?),
status => {
return Err(io::Error::other(format!(
"read {} from {} failed: {}",
key, addr, status
)));
}
};
if value.as_ref() == expected {
return Ok(());
}
last = Some(value);
}
if Instant::now() >= deadline {
return Err(io::Error::other(format!(
"{} never served {:?} for {}, last seen: {:?}",
addr, expected, key, last
)));
}
tokio::time::sleep(POLL_INTERVAL).await;
}
}
#[tokio::test(flavor = "multi_thread")]
async fn three_node_kvstore_cluster_writes_and_reads() -> io::Result<()> {
let work_dir = std::env::temp_dir().join(format!("ezraft-kvstore-{}", std::process::id()));
std::fs::create_dir_all(&work_dir)?;
let addr_a = free_addr();
let a = spawn_node(&work_dir, &addr_a, None)?;
assert_eq!(
json!({"value": null}),
write(&addr_a, json!({"Set": {"key": "k1", "value": "v1"}})).await?
);
let addr_b = free_addr();
let addr_c = free_addr();
let b = spawn_node(&work_dir, &addr_b, Some(&addr_a))?;
let c = spawn_node(&work_dir, &addr_c, Some(&addr_a))?;
wait_for_key(&addr_b, "k1", Some(&json!("v1"))).await?;
wait_for_key(&addr_c, "k1", Some(&json!("v1"))).await?;
assert_eq!(
json!({"value": null}),
write(&addr_b, json!({"Set": {"key": "k2", "value": "v2"}})).await?
);
assert_eq!(
json!({"value": "v1"}),
write(&addr_c, json!({"Delete": {"key": "k1"}})).await?
);
assert_eq!(
json!({"value": "v2"}),
write(&addr_a, json!({"Set": {"key": "k2", "value": "v2b"}})).await?
);
for addr in [&addr_a, &addr_b, &addr_c] {
wait_for_key(addr, "k2", Some(&json!("v2b"))).await?;
wait_for_key(addr, "k1", None).await?;
}
drop(a);
drop(b);
drop(c);
std::fs::remove_dir_all(&work_dir)?;
Ok(())
}