use std::collections::BTreeMap;
use std::collections::BTreeSet;
use std::io;
use std::io::Cursor;
use std::sync::Arc;
use std::sync::Mutex;
use std::time::Duration;
use async_trait::async_trait;
use ezraft::EzApp;
use ezraft::EzConfig;
use ezraft::EzEntry;
use ezraft::EzMeta;
use ezraft::EzRaft;
use ezraft::EzSnapshot;
use ezraft::EzSnapshotMeta;
use ezraft::EzStorage;
use ezraft::Loaded;
use ezraft::Persist;
use ezraft::admin::AdminClient;
use serde::Deserialize;
use serde::Serialize;
#[derive(Serialize, Deserialize, Debug, Clone, derive_more::Display)]
enum Request {
#[display("Set({key})")]
Set { key: String, value: String },
#[display("Get({key})")]
Get { key: String },
}
#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq)]
struct Response {
value: Option<String>,
}
fn set(key: &str, value: &str) -> Request {
Request::Set {
key: key.into(),
value: value.into(),
}
}
fn get(key: &str) -> Request {
Request::Get { key: key.into() }
}
#[derive(Default, Serialize, Deserialize)]
struct KvSm {
data: BTreeMap<String, String>,
}
#[async_trait]
impl EzApp for KvSm {
type Request = Request;
type Response = Response;
async fn apply(&mut self, req: Request) -> Response {
match req {
Request::Set { key, value } => {
self.data.insert(key, value);
Response { value: None }
}
Request::Get { key } => Response {
value: self.data.get(&key).cloned(),
},
}
}
type ReadRequest = String;
type ReadResponse = Option<String>;
fn read(&self, key: String) -> Option<String> {
self.data.get(&key).cloned()
}
}
#[derive(Default)]
struct Disk {
meta: EzMeta,
logs: BTreeMap<u64, Vec<u8>>,
snapshot: Option<(EzSnapshotMeta, Vec<u8>)>,
}
#[derive(Clone, Default)]
struct MemStorage {
disk: Arc<Mutex<Disk>>,
}
#[async_trait]
impl EzStorage<KvSm> for MemStorage {
async fn load(&mut self) -> io::Result<Loaded> {
let disk = self.disk.lock().unwrap();
let snapshot = disk.snapshot.as_ref().map(|(meta, data)| EzSnapshot {
meta: meta.clone(),
snapshot: Cursor::new(data.clone()),
});
Ok(Loaded {
meta: disk.meta.clone(),
snapshot,
})
}
async fn persist(&mut self, op: Persist<KvSm>) -> io::Result<()> {
let mut disk = self.disk.lock().unwrap();
match op {
Persist::Meta(meta) => disk.meta = meta,
Persist::LogEntry(entry) => {
disk.logs.insert(entry.log_id.1, serde_json::to_vec(&entry)?);
}
Persist::Snapshot(snapshot) => {
disk.snapshot = Some((snapshot.meta, snapshot.snapshot.into_inner()));
}
Persist::DeleteLogs { from, to } => disk.logs.retain(|&index, _| !(from..to).contains(&index)),
}
Ok(())
}
async fn read_logs(&mut self, start: u64, end: u64) -> io::Result<Vec<EzEntry<KvSm>>> {
let disk = self.disk.lock().unwrap();
(start..end)
.map(|index| {
let data =
disk.logs.get(&index).ok_or_else(|| io::Error::other(format!("missing log entry {}", index)))?;
Ok(serde_json::from_slice(data)?)
})
.collect()
}
}
fn free_addr() -> String {
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
listener.local_addr().unwrap().to_string()
}
fn config() -> EzConfig {
EzConfig {
heartbeat_interval: Duration::from_millis(100),
..EzConfig::default()
}
}
const WAIT: Option<Duration> = Some(Duration::from_secs(30));
fn expected_map(range: std::ops::Range<u32>) -> BTreeMap<String, String> {
range.map(|i| (format!("k{}", i), format!("v{}", i))).collect()
}
async fn http_write(addr: &str, req: Request) -> io::Result<Response> {
let resp = reqwest::Client::new()
.post(format!("http://{}/api/write", addr))
.json(&req)
.send()
.await
.map_err(io::Error::other)?;
let status = resp.status();
let text = resp.text().await.map_err(io::Error::other)?;
assert!(status.is_success(), "POST /api/write responded {}: {}", status, text);
Ok(serde_json::from_str(&text)?)
}
fn spawn_serve(node: &EzRaft<KvSm>) {
let node = node.clone();
tokio::spawn(async move { node.serve().await });
}
async fn founding_node() -> io::Result<(String, EzRaft<KvSm>)> {
let addr = free_addr();
let node = EzRaft::create(&addr, KvSm::default(), MemStorage::default(), config()).await?;
spawn_serve(&node);
node.inner()
.wait(WAIT)
.metrics(
|m| m.current_leader == Some(0) && m.committed_membership_config.log_id().is_some(),
"founding node leads, with its own membership committed",
)
.await
.map_err(io::Error::other)?;
Ok((addr, node))
}
async fn joined_voter(seed: &str) -> io::Result<(String, EzRaft<KvSm>)> {
let addr = free_addr();
let node = EzRaft::join(&addr, seed, KvSm::default(), MemStorage::default(), config()).await?;
spawn_serve(&node);
Ok((addr, node))
}
async fn joined_learner(seed: &str) -> io::Result<(String, EzRaft<KvSm>)> {
let addr = free_addr();
let node = EzRaft::join_as_learner(&addr, seed, KvSm::default(), MemStorage::default(), config()).await?;
spawn_serve(&node);
Ok((addr, node))
}
async fn wait_for_voters(node: &EzRaft<KvSm>, voters: BTreeSet<u64>, reason: &str) -> io::Result<()> {
node.inner()
.wait(WAIT)
.metrics(
|m| *m.committed_membership_config.membership().get_joint_config() == [voters.clone()],
reason,
)
.await
.map_err(io::Error::other)?;
Ok(())
}
async fn wait_for_applied(node: &EzRaft<KvSm>, leader: &EzRaft<KvSm>) -> io::Result<()> {
let target = leader.metrics().await.last_log_index.expect("the leader has written entries");
node.inner()
.wait(WAIT)
.metrics(
|m| m.last_applied.map(|log_id| log_id.index) >= Some(target),
"applied the leader's log",
)
.await
.map_err(io::Error::other)?;
Ok(())
}
async fn admin_answer(addr: &str, body: serde_json::Value) -> io::Result<Redirect> {
let resp = reqwest::Client::new()
.post(format!("http://{}/api/membership", addr))
.json(&body)
.send()
.await
.map_err(io::Error::other)?;
let status = resp.status();
let text = resp.text().await.map_err(io::Error::other)?;
assert!(
status.is_success(),
"POST /api/membership responded {}: {}",
status,
text
);
Ok(serde_json::from_str(&text)?)
}
type Redirect = Result<(), Option<String>>;
async fn admin_post(addr: &str, body: serde_json::Value) -> io::Result<()> {
admin_answer(addr, body)
.await?
.map_err(|leader| io::Error::other(format!("/api/membership redirected to {:?}", leader)))
}
#[tokio::test(flavor = "multi_thread")]
async fn join_promotes_to_voter_and_cluster_survives_leader_death() -> io::Result<()> {
let addr_a = free_addr();
let a = EzRaft::create(&addr_a, KvSm::default(), MemStorage::default(), config()).await?;
tokio::spawn({
let a = a.clone();
async move { a.serve().await }
});
a.inner()
.wait(WAIT)
.metrics(|m| m.current_leader == Some(0), "founding node leads")
.await
.map_err(io::Error::other)?;
let addr_b = free_addr();
let b = EzRaft::join(&addr_b, &addr_a, KvSm::default(), MemStorage::default(), config()).await?;
tokio::spawn({
let b = b.clone();
async move { b.serve().await }
});
let addr_c = free_addr();
let c = EzRaft::join(&addr_c, &addr_a, KvSm::default(), MemStorage::default(), config()).await?;
tokio::spawn({
let c = c.clone();
async move { c.serve().await }
});
let voters = BTreeSet::from([0, b.node_id(), c.node_id()]);
for node in [&a, &b, &c] {
wait_for_voters(node, voters.clone(), "every promoted node is a voter").await?;
}
assert_eq!(Response { value: None }, a.write(set("k1", "v1")).await?);
let value: Option<String> = reqwest::Client::new()
.post(format!("http://{}/api/read", addr_a))
.json("k1")
.send()
.await
.map_err(io::Error::other)?
.json()
.await
.map_err(io::Error::other)?;
assert_eq!(Some("v1".to_string()), value);
assert!(a.is_leader());
a.inner().shutdown().await.map_err(io::Error::other)?;
b.inner()
.wait(WAIT)
.metrics(
|m| matches!(m.current_leader, Some(id) if id != 0),
"a surviving node takes over",
)
.await
.map_err(io::Error::other)?;
assert_eq!(Response { value: None }, http_write(&addr_b, set("k2", "v2")).await?);
assert_eq!(
Response {
value: Some("v1".into())
},
http_write(&addr_b, get("k1")).await?
);
assert_eq!(
Response {
value: Some("v2".into())
},
http_write(&addr_b, get("k2")).await?
);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn learner_joins_and_stays_a_learner_until_promoted() -> io::Result<()> {
let addr_a = free_addr();
let a = EzRaft::create(&addr_a, KvSm::default(), MemStorage::default(), config()).await?;
tokio::spawn({
let a = a.clone();
async move { a.serve().await }
});
a.inner()
.wait(WAIT)
.metrics(|m| m.current_leader == Some(0), "founding node leads")
.await
.map_err(io::Error::other)?;
let addr_b = free_addr();
let b = EzRaft::join_as_learner(&addr_b, &addr_a, KvSm::default(), MemStorage::default(), config()).await?;
tokio::spawn({
let b = b.clone();
async move { b.serve().await }
});
let b_id = b.node_id();
for i in 1..3 {
a.write(set(&format!("k{}", i), &format!("v{}", i))).await?;
}
let target = a.metrics().await.last_log_index.expect("the leader has written entries");
b.inner()
.wait(WAIT)
.metrics(
|m| m.last_applied.map(|log_id| log_id.index) >= Some(target),
"learner applied the leader's log",
)
.await
.map_err(io::Error::other)?;
let metrics = a.metrics().await;
let membership = metrics.membership_config.membership();
assert_eq!(BTreeSet::from([0]), membership.voter_ids().collect::<BTreeSet<_>>());
assert_eq!(
BTreeSet::from([b_id]),
membership.learner_ids().collect::<BTreeSet<_>>()
);
assert_eq!(expected_map(1..3), b.read(|app| app.data.clone()).await?);
let unknown = a.promote(b_id + 1000).await.unwrap_err();
assert!(
unknown.to_string().contains(&format!("Learner {} not found", b_id + 1000)),
"unexpected error: {}",
unknown
);
a.promote(b_id).await?;
let voters = BTreeSet::from([0, b_id]);
for node in [&a, &b] {
wait_for_voters(node, voters.clone(), "the promoted learner is a voter everywhere").await?;
}
a.promote(b_id).await?;
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn snapshot_survives_restart() -> io::Result<()> {
let addr = free_addr();
let storage = MemStorage::default();
let a = EzRaft::create(&addr, KvSm::default(), storage.clone(), config()).await?;
a.inner()
.wait(WAIT)
.metrics(|m| m.current_leader == Some(0), "single node leads")
.await
.map_err(io::Error::other)?;
for i in 0..10 {
a.write(set(&format!("k{}", i), &format!("v{}", i))).await?;
}
a.inner().trigger().snapshot().await.map_err(io::Error::other)?;
a.inner()
.wait(WAIT)
.metrics(|m| m.snapshot.is_some(), "snapshot built")
.await
.map_err(io::Error::other)?;
{
let disk = storage.disk.lock().unwrap();
let (meta, data) = disk.snapshot.as_ref().expect("snapshot persisted to storage");
let snapshot_state: KvSm = serde_json::from_slice(data)?;
assert_eq!(expected_map(0..10), snapshot_state.data);
assert!(meta.last_log_id.is_some());
}
for i in 10..15 {
a.write(set(&format!("k{}", i), &format!("v{}", i))).await?;
}
a.inner().shutdown().await.map_err(io::Error::other)?;
drop(a);
let restarted = EzRaft::create(&addr, KvSm::default(), storage.clone(), config()).await?;
restarted
.inner()
.wait(WAIT)
.metrics(|m| m.current_leader == Some(0), "restarted node leads")
.await
.map_err(io::Error::other)?;
restarted
.inner()
.wait(WAIT)
.metrics(
|m| m.last_applied.map(|log_id| log_id.index) == m.last_log_index,
"log tail replayed",
)
.await
.map_err(io::Error::other)?;
restarted.linearizable().await?;
assert_eq!(expected_map(0..15), restarted.read(|app| app.data.clone()).await?);
assert_eq!(
Response {
value: Some("v14".into())
},
restarted.write(get("k14")).await?
);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn lagging_joiner_catches_up_from_snapshot() -> io::Result<()> {
let addr_a = free_addr();
let a = EzRaft::create(&addr_a, KvSm::default(), MemStorage::default(), config()).await?;
tokio::spawn({
let a = a.clone();
async move { a.serve().await }
});
a.inner()
.wait(WAIT)
.metrics(|m| m.current_leader == Some(0), "founding node leads")
.await
.map_err(io::Error::other)?;
for i in 0..10 {
a.write(set(&format!("k{}", i), &format!("v{}", i))).await?;
}
a.inner().trigger().snapshot().await.map_err(io::Error::other)?;
let snapshot_index = a
.inner()
.wait(WAIT)
.metrics(|m| m.snapshot.is_some(), "snapshot built")
.await
.map_err(io::Error::other)?
.snapshot
.unwrap()
.index;
a.inner().trigger().purge_log(snapshot_index).await.map_err(io::Error::other)?;
a.inner()
.wait(WAIT)
.metrics(
|m| m.purged.map(|log_id| log_id.index) == Some(snapshot_index),
"log purged up to the snapshot",
)
.await
.map_err(io::Error::other)?;
let addr_b = free_addr();
let b = EzRaft::join(&addr_b, &addr_a, KvSm::default(), MemStorage::default(), config()).await?;
tokio::spawn({
let b = b.clone();
async move { b.serve().await }
});
let voters = BTreeSet::from([0, b.node_id()]);
wait_for_voters(&b, voters, "snapshot-fed joiner promoted to voter").await?;
assert_eq!(expected_map(0..10), b.read(|app| app.data.clone()).await?);
assert_eq!(Response { value: None }, http_write(&addr_b, set("k10", "v10")).await?);
assert_eq!(
Response {
value: Some("v0".into())
},
http_write(&addr_b, get("k0")).await?
);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn automatic_snapshot_purges_old_logs() -> io::Result<()> {
let addr = free_addr();
let storage = MemStorage::default();
let config = EzConfig {
snapshot_interval: 5,
..config()
};
let a = EzRaft::create(&addr, KvSm::default(), storage.clone(), config).await?;
a.inner()
.wait(WAIT)
.metrics(|m| m.current_leader == Some(0), "single node leads")
.await
.map_err(io::Error::other)?;
for i in 0..12 {
a.write(set(&format!("k{}", i), &format!("v{}", i))).await?;
}
let metrics = a
.inner()
.wait(WAIT)
.metrics(
|m| m.snapshot.is_some() && m.purged.is_some(),
"snapshot built and log purged by the interval policy alone",
)
.await
.map_err(io::Error::other)?;
let disk = storage.disk.lock().unwrap();
let (meta, data) = disk.snapshot.as_ref().expect("snapshot persisted to storage");
assert!(meta.last_log_id.is_some());
let snapshot_state: KvSm = serde_json::from_slice(data)?;
assert!(!snapshot_state.data.is_empty());
let purged_index = metrics.purged.unwrap().index;
let min_kept = *disk.logs.keys().next().expect("entries above the purge point remain");
assert!(
min_kept > purged_index,
"min kept index {} must be above purged index {}",
min_kept,
purged_index
);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn demoted_voter_becomes_a_learner_and_keeps_replicating() -> io::Result<()> {
let (addr_a, a) = founding_node().await?;
let (_, b) = joined_voter(&addr_a).await?;
let (_, c) = joined_voter(&addr_a).await?;
let all = BTreeSet::from([0, b.node_id(), c.node_id()]);
wait_for_voters(&a, all, "every joined node is a voter").await?;
a.demote(c.node_id()).await?;
let voters = BTreeSet::from([0, b.node_id()]);
wait_for_voters(&a, voters.clone(), "the demoted node left the voter set").await?;
let metrics = a.metrics().await;
let membership = metrics.membership_config.membership();
assert_eq!(voters, membership.voter_ids().collect::<BTreeSet<_>>());
assert_eq!(
BTreeSet::from([c.node_id()]),
membership.learner_ids().collect::<BTreeSet<_>>()
);
assert_eq!(Response { value: None }, a.write(set("k1", "v1")).await?);
wait_for_applied(&c, &a).await?;
assert_eq!(expected_map(1..2), c.read(|app| app.data.clone()).await?);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn demote_refuses_to_empty_the_voter_set() -> io::Result<()> {
let (_, a) = founding_node().await?;
let err = a.demote(0).await.unwrap_err();
assert!(
err.to_string().contains("new membership cannot be empty"),
"unexpected error: {}",
err
);
let metrics = a.metrics().await;
assert_eq!(
BTreeSet::from([0]),
metrics.membership_config.membership().voter_ids().collect::<BTreeSet<_>>()
);
assert_eq!(Response { value: None }, a.write(set("k1", "v1")).await?);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn demoted_leader_keeps_leading_outside_the_quorum() -> io::Result<()> {
let (addr_a, a) = founding_node().await?;
let (addr_b, b) = joined_voter(&addr_a).await?;
let (_, c) = joined_voter(&addr_a).await?;
let all = BTreeSet::from([0, b.node_id(), c.node_id()]);
wait_for_voters(&a, all, "every joined node is a voter").await?;
assert!(a.is_leader());
a.demote(0).await?;
let voters = BTreeSet::from([b.node_id(), c.node_id()]);
wait_for_voters(&a, voters.clone(), "the demoted leader left the voter set").await?;
let metrics = a.metrics().await;
let membership = metrics.committed_membership_config.membership();
assert_eq!(voters, membership.voter_ids().collect::<BTreeSet<_>>());
assert_eq!(BTreeSet::from([0]), membership.learner_ids().collect::<BTreeSet<_>>());
assert!(a.is_leader());
assert_eq!(Response { value: None }, http_write(&addr_a, set("k1", "v1")).await?);
a.inner().shutdown().await.map_err(io::Error::other)?;
b.inner()
.wait(WAIT)
.metrics(
|m| matches!(m.current_leader, Some(id) if id != 0),
"a voter takes over from the demoted leader",
)
.await
.map_err(io::Error::other)?;
assert_eq!(Response { value: None }, http_write(&addr_b, set("k2", "v2")).await?);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn leadership_transfers_to_the_named_node() -> io::Result<()> {
let (addr_a, a) = founding_node().await?;
let (_, b) = joined_voter(&addr_a).await?;
let (addr_c, c) = joined_voter(&addr_a).await?;
let all = BTreeSet::from([0, b.node_id(), c.node_id()]);
wait_for_voters(&a, all, "every joined node is a voter").await?;
assert!(a.is_leader());
let target = c.node_id();
a.inner().trigger().transfer_leader(target).await.map_err(io::Error::other)?;
c.inner()
.wait(WAIT)
.metrics(|m| m.current_leader == Some(target), "the named node takes leadership")
.await
.map_err(io::Error::other)?;
assert_eq!(Response { value: None }, http_write(&addr_c, set("k1", "v1")).await?);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn removed_leader_hands_over_leadership() -> io::Result<()> {
let (addr_a, a) = founding_node().await?;
let (addr_b, b) = joined_voter(&addr_a).await?;
let (_, c) = joined_voter(&addr_a).await?;
let all = BTreeSet::from([0, b.node_id(), c.node_id()]);
wait_for_voters(&a, all, "every joined node is a voter").await?;
assert!(a.is_leader());
a.remove_node(0).await?;
b.inner()
.wait(WAIT)
.metrics(
|m| matches!(m.current_leader, Some(id) if id != 0),
"leadership moves off the removed node",
)
.await
.map_err(io::Error::other)?;
let voters = BTreeSet::from([b.node_id(), c.node_id()]);
wait_for_voters(&b, voters.clone(), "the removed leader left the cluster").await?;
let metrics = b.metrics().await;
let membership = metrics.committed_membership_config.membership();
assert_eq!(
voters.iter().copied().collect::<Vec<_>>(),
membership.nodes().map(|(id, _)| *id).collect::<Vec<_>>()
);
assert_eq!(Response { value: None }, http_write(&addr_b, set("k1", "v1")).await?);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn removed_nodes_leave_the_membership() -> io::Result<()> {
let (addr_a, a) = founding_node().await?;
let (_, b) = joined_voter(&addr_a).await?;
let (_, c) = joined_learner(&addr_a).await?;
wait_for_voters(&a, BTreeSet::from([0, b.node_id()]), "the joining voter is promoted").await?;
a.remove_node(c.node_id()).await?;
a.remove_node(b.node_id()).await?;
wait_for_voters(&a, BTreeSet::from([0]), "the removed voter left the voter set").await?;
let metrics = a.metrics().await;
let membership = metrics.membership_config.membership();
assert_eq!(vec![0], membership.nodes().map(|(id, _)| *id).collect::<Vec<_>>());
assert!(membership.learner_ids().next().is_none());
assert_eq!(Response { value: None }, a.write(set("k1", "v1")).await?);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn admin_api_changes_roles_and_removes_nodes() -> io::Result<()> {
let (addr_a, a) = founding_node().await?;
let (_, b) = joined_learner(&addr_a).await?;
let b_id = b.node_id();
admin_post(
&addr_a,
serde_json::json!({"op": "SetRole", "node_id": b_id, "role": "Voter"}),
)
.await?;
wait_for_voters(&a, BTreeSet::from([0, b_id]), "promoted over HTTP").await?;
admin_post(
&addr_a,
serde_json::json!({"op": "SetRole", "node_id": b_id, "role": "Learner"}),
)
.await?;
wait_for_voters(&a, BTreeSet::from([0]), "demoted over HTTP").await?;
let metrics = a.metrics().await;
assert_eq!(
BTreeSet::from([b_id]),
metrics.membership_config.membership().learner_ids().collect::<BTreeSet<_>>()
);
admin_post(&addr_a, serde_json::json!({"op": "Remove", "node_id": b_id})).await?;
a.inner()
.wait(WAIT)
.metrics(
|m| m.membership_config.membership().get_node(&b_id).is_none(),
"removed over HTTP",
)
.await
.map_err(io::Error::other)?;
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn admin_client_reads_the_metrics_of_the_node_it_asks() -> io::Result<()> {
let (addr_a, a) = founding_node().await?;
let (addr_b, b) = joined_voter(&addr_a).await?;
wait_for_voters(&a, BTreeSet::from([0, b.node_id()]), "the joining voter is promoted").await?;
assert!(!b.is_leader(), "b has to be a follower for this to be its own view");
wait_for_applied(&b, &a).await?;
assert_eq!(b.metrics().await, AdminClient::new(&addr_b).metrics::<KvSm>().await?);
let fresh = AdminClient::new(&addr_a).node_id().await?;
assert!(fresh > b.node_id(), "an id past every one handed out so far");
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn admin_api_redirects_a_follower_to_the_leader() -> io::Result<()> {
let (addr_a, a) = founding_node().await?;
let (addr_b, b) = joined_voter(&addr_a).await?;
let (_, c) = joined_learner(&addr_a).await?;
let c_id = c.node_id();
wait_for_voters(&a, BTreeSet::from([0, b.node_id()]), "the joining voter is promoted").await?;
assert!(!b.is_leader(), "b has to be a follower for this to be a redirect");
let promote = serde_json::json!({"op": "SetRole", "node_id": c_id, "role": "Voter"});
assert_eq!(Err(Some(addr_a.clone())), admin_answer(&addr_b, promote.clone()).await?);
assert!(
!a.metrics().await.membership_config.membership().voter_ids().any(|id| id == c_id),
"the redirected request must not have promoted anything"
);
admin_post(&addr_a, promote).await?;
wait_for_voters(
&a,
BTreeSet::from([0, b.node_id(), c_id]),
"the leader did what the follower would not",
)
.await?;
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn restart_after_a_snapshot_with_no_log_after_it() -> io::Result<()> {
let (addr_a, a) = founding_node().await?;
let addr_c = free_addr();
let disk_c = MemStorage::default();
let c = EzRaft::join_as_learner(&addr_c, &addr_a, KvSm::default(), disk_c.clone(), config()).await?;
for i in 0..10 {
a.write(set(&format!("k{}", i), &format!("v{}", i))).await?;
}
a.inner().trigger().snapshot().await.map_err(io::Error::other)?;
let snapshot_index = a
.inner()
.wait(WAIT)
.metrics(|m| m.snapshot.is_some(), "snapshot built")
.await
.map_err(io::Error::other)?
.snapshot
.unwrap()
.index;
a.inner().trigger().purge_log(snapshot_index).await.map_err(io::Error::other)?;
a.inner()
.wait(WAIT)
.metrics(
|m| m.purged.map(|log_id| log_id.index) == Some(snapshot_index),
"log purged up to the snapshot",
)
.await
.map_err(io::Error::other)?;
assert_eq!(
Some(snapshot_index),
a.metrics().await.last_log_index,
"the snapshot has to cover the whole log, or an append would follow it"
);
spawn_serve(&c);
wait_for_applied(&c, &a).await?;
assert_eq!(expected_map(0..10), c.read(|app| app.data.clone()).await?);
let meta = disk_c.disk.lock().unwrap().meta.clone();
assert!(
meta.last_purged <= meta.last_log_id,
"persisted an inverted log range: last_purged={:?} last_log_id={:?}",
meta.last_purged,
meta.last_log_id
);
c.inner().shutdown().await.map_err(io::Error::other)?;
drop(c);
EzRaft::join_as_learner(&free_addr(), &addr_a, KvSm::default(), disk_c, config()).await?;
Ok(())
}