Skip to content

Commit e55cb36

Browse files
committed
fix: CI
1 parent dc89235 commit e55cb36

7 files changed

Lines changed: 64 additions & 33 deletions

File tree

examples/memstore/src/web_server_api.rs

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,7 @@ async fn leader_id(data: web::Data<(HashStore, Raft)>) -> impl Responder {
4040
#[get("/leave")]
4141
async fn leave(data: web::Data<(HashStore, Raft)>) -> impl Responder {
4242
let raft = data.clone();
43-
raft.1.leave().await;
43+
raft.1.leave().await.unwrap();
4444
"OK".to_string()
4545
}
4646

@@ -107,21 +107,21 @@ async fn transfer_leader(
107107
) -> impl Responder {
108108
let raft = data.clone();
109109
let node_id: u64 = path.into_inner();
110-
raft.1.transfer_leader(node_id).await;
110+
raft.1.transfer_leader(node_id).await.unwrap();
111111
"OK".to_string()
112112
}
113113

114114
#[get("/campaign")]
115115
async fn campaign(data: web::Data<(HashStore, Raft)>) -> impl Responder {
116116
let raft = data.clone();
117-
raft.1.campaign().await;
117+
raft.1.campaign().await.unwrap();
118118
"OK".to_string()
119119
}
120120

121121
#[get("/demote/{term}/{leader_id}")]
122122
async fn demote(data: web::Data<(HashStore, Raft)>, path: web::Path<(u64, u64)>) -> impl Responder {
123123
let raft = data.clone();
124124
let (term, leader_id_) = path.into_inner();
125-
raft.1.demote(term, leader_id_).await;
125+
raft.1.demote(term, leader_id_).await.unwrap();
126126
"OK".to_string()
127127
}

harness/src/config.rs

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,8 @@
11
use raftify::{Config, RaftConfig};
22

3-
pub fn build_config() -> Config {
3+
pub fn build_config(node_id: u64) -> Config {
44
let raft_config = RaftConfig {
5+
id: node_id,
56
election_tick: 10,
67
heartbeat_tick: 3,
78
omit_heartbeat_log: true,

harness/src/raft.rs

Lines changed: 19 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -51,7 +51,7 @@ fn run_raft(
5151
should_be_leader: bool,
5252
) -> Result<JoinHandle<Result<()>>> {
5353
let peer = peers.get(node_id).unwrap();
54-
let mut cfg = build_config();
54+
let mut cfg = build_config(*node_id);
5555
cfg.initial_peers = if should_be_leader {
5656
None
5757
} else {
@@ -60,7 +60,7 @@ fn run_raft(
6060

6161
let store = HashStore::new();
6262
let logger = build_logger();
63-
let storage_pth = get_storage_path(cfg.log_dir.as_str(), 1);
63+
let storage_pth = get_storage_path(cfg.log_dir.as_str(), *node_id);
6464
ensure_directory_exist(storage_pth.as_str())?;
6565

6666
let storage = HeedStorage::create(
@@ -141,9 +141,9 @@ pub async fn spawn_extra_node(
141141
slog: build_logger(),
142142
});
143143

144-
let cfg = build_config();
144+
let cfg = build_config(node_id);
145145
let store = HashStore::new();
146-
let storage_pth = get_storage_path(cfg.log_dir.as_str(), 1);
146+
let storage_pth = get_storage_path(cfg.log_dir.as_str(), node_id);
147147
ensure_directory_exist(storage_pth.as_str())?;
148148

149149
let storage = HeedStorage::create(&storage_pth, &cfg, logger.clone())?;
@@ -172,11 +172,11 @@ pub async fn spawn_and_join_extra_node(
172172
.unwrap();
173173

174174
let node_id = join_ticket.reserved_id;
175-
let mut cfg = build_config();
175+
let mut cfg = build_config(node_id);
176176
cfg.initial_peers = Some(join_ticket.peers.clone().into());
177177
let store = HashStore::new();
178178

179-
let storage_pth = get_storage_path(cfg.log_dir.as_str(), 1);
179+
let storage_pth = get_storage_path(cfg.log_dir.as_str(), node_id);
180180
ensure_directory_exist(storage_pth.as_str())?;
181181

182182
let storage = HeedStorage::create(&storage_pth, &cfg, logger.clone())?;
@@ -189,8 +189,12 @@ pub async fn spawn_and_join_extra_node(
189189

190190
let raft_handle = tokio::spawn(raft.clone().run());
191191

192-
raft.add_peers(join_ticket.peers.clone()).await;
193-
raft.join_cluster(vec![join_ticket]).await;
192+
raft.add_peers(join_ticket.peers.clone())
193+
.await
194+
.expect("Failed to add peers");
195+
raft.join_cluster(vec![join_ticket])
196+
.await
197+
.expect("Failed to join cluster");
194198

195199
Ok(raft_handle)
196200
}
@@ -202,9 +206,14 @@ pub async fn join_nodes(rafts: Vec<&Raft>, raft_addrs: Vec<&str>, peer_addr: &st
202206
.await
203207
.unwrap();
204208

205-
raft.add_peers(join_ticket.peers.clone()).await;
209+
raft.add_peers(join_ticket.peers.clone())
210+
.await
211+
.expect("Failed to add peers");
206212
tickets.push(join_ticket);
207213
}
208214

209-
rafts[0].join_cluster(tickets).await;
215+
rafts[0]
216+
.join_cluster(tickets)
217+
.await
218+
.expect("Failed to join cluster");
210219
}

harness/src/utils.rs

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -175,6 +175,24 @@ pub fn kill_process_using_port(port: u16) {
175175
}
176176
}
177177

178+
pub fn cleanup_storage(log_dir: &str) {
179+
let storage_pth = Path::new(log_dir);
180+
181+
if fs::metadata(storage_pth).is_ok() {
182+
fs::remove_dir_all(storage_pth)
183+
.expect("Failed to remove storage directory");
184+
}
185+
186+
fs::create_dir_all(storage_pth)
187+
.expect("Failed to create storage directory");
188+
189+
for i in 1..=9 {
190+
let node_dir = storage_pth.join(format!("node-{}", i));
191+
fs::create_dir(&node_dir)
192+
.unwrap_or_else(|e| panic!("Failed to create directory {:?}: {}", node_dir, e));
193+
}
194+
}
195+
178196
pub fn kill_previous_raft_processes() {
179197
RAFT_PORTS.iter().for_each(|port| {
180198
kill_process_using_port(*port);

harness/tests/bootstrap.rs

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -4,11 +4,12 @@ use tokio::time::sleep;
44
use harness::{
55
constant::{ONE_NODE_EXAMPLE, RAFT_ADDRS, THREE_NODE_EXAMPLE},
66
raft::{build_raft_cluster, spawn_and_join_extra_node, wait_until_rafts_ready, Raft},
7-
utils::{kill_previous_raft_processes, load_peers, wait_for_until_cluster_size_increase},
7+
utils::{cleanup_storage, kill_previous_raft_processes, load_peers, wait_for_until_cluster_size_increase},
88
};
99

1010
#[tokio::test]
1111
pub async fn test_static_bootstrap() {
12+
cleanup_storage("./logs");
1213
kill_previous_raft_processes();
1314
let (tx_raft, rx_raft) = mpsc::channel::<(u64, Raft)>();
1415

@@ -23,14 +24,15 @@ pub async fn test_static_bootstrap() {
2324
wait_for_until_cluster_size_increase(raft_1.clone(), 3).await;
2425

2526
for (_, raft) in rafts.iter_mut() {
26-
raft.quit().await;
27+
raft.quit().await.expect("Failed to quit raft node");
2728
}
2829

2930
sleep(Duration::from_secs(1)).await;
3031
}
3132

3233
#[tokio::test]
3334
pub async fn test_dynamic_bootstrap() {
35+
cleanup_storage("./logs");
3436
kill_previous_raft_processes();
3537
let (tx_raft, rx_raft) = mpsc::channel::<(u64, Raft)>();
3638

@@ -64,7 +66,7 @@ pub async fn test_dynamic_bootstrap() {
6466
wait_for_until_cluster_size_increase(raft_1.clone(), 3).await;
6567

6668
for (_, raft) in rafts.iter_mut() {
67-
raft.quit().await;
69+
raft.quit().await.expect("Failed to quit raft node");
6870
}
6971
}
7072

harness/tests/data_replication.rs

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -6,11 +6,12 @@ use harness::{
66
constant::{RAFT_ADDRS, THREE_NODE_EXAMPLE},
77
raft::{build_raft_cluster, spawn_and_join_extra_node, wait_until_rafts_ready, Raft},
88
state_machine::LogEntry,
9-
utils::{kill_previous_raft_processes, load_peers, wait_for_until_cluster_size_increase},
9+
utils::{cleanup_storage, kill_previous_raft_processes, load_peers, wait_for_until_cluster_size_increase},
1010
};
1111

1212
#[tokio::test]
1313
pub async fn test_data_replication() {
14+
cleanup_storage("./logs");
1415
kill_previous_raft_processes();
1516

1617
let peers = load_peers(THREE_NODE_EXAMPLE).await.unwrap();
@@ -33,7 +34,7 @@ pub async fn test_data_replication() {
3334

3435
raft_1.propose(entry).await.unwrap();
3536

36-
sleep(Duration::from_secs(1)).await;
37+
sleep(Duration::from_secs(3)).await;
3738

3839
// Data should be replicated to all nodes.
3940
for (_, raft) in rafts.iter_mut() {
@@ -62,7 +63,7 @@ pub async fn test_data_replication() {
6263
let store = raft_4.state_machine().await.unwrap();
6364
let store_lk = store.0.read().unwrap();
6465

65-
// Data should be replicated to new joined node.
66+
// Data should be replicated to new member.
6667
assert_eq!(store_lk.get(&1).unwrap(), "test");
6768
std::mem::drop(store_lk);
6869

@@ -77,15 +78,14 @@ pub async fn test_data_replication() {
7778

7879
raft_1.propose(new_entry).await.unwrap();
7980

80-
// New entry data should be replicated to all nodes including new joined node.
81-
for (_, raft) in rafts.iter() {
82-
// stop
81+
// New entry data should be replicated to all nodes including new member.
82+
for (id, raft) in rafts.iter() {
8383
let store = raft.state_machine().await.unwrap();
8484
let store_lk = store.0.read().unwrap();
8585
assert_eq!(store_lk.get(&2).unwrap(), "test2");
8686
}
8787

8888
for (_, raft) in rafts.iter_mut() {
89-
raft.quit().await;
89+
raft.quit().await.expect("Failed to quit the raft node");
9090
}
9191
}

harness/tests/leader_election.rs

Lines changed: 9 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -5,13 +5,13 @@ use harness::{
55
constant::{FIVE_NODE_EXAMPLE, THREE_NODE_EXAMPLE},
66
raft::{build_raft_cluster, wait_until_rafts_ready, Raft},
77
utils::{
8-
kill_previous_raft_processes, load_peers, wait_for_until_cluster_size_decrease,
9-
wait_for_until_cluster_size_increase,
8+
cleanup_storage, kill_previous_raft_processes, load_peers, wait_for_until_cluster_size_decrease, wait_for_until_cluster_size_increase
109
},
1110
};
1211

1312
#[tokio::test]
1413
pub async fn test_leader_election_in_three_node_example() {
14+
cleanup_storage("./logs");
1515
kill_previous_raft_processes();
1616

1717
let (tx_raft, rx_raft) = mpsc::channel::<(u64, Raft)>();
@@ -28,7 +28,7 @@ pub async fn test_leader_election_in_three_node_example() {
2828

2929
sleep(Duration::from_secs(1)).await;
3030

31-
raft_1.leave().await;
31+
raft_1.leave().await.expect("Failed to leave");
3232

3333
sleep(Duration::from_secs(2)).await;
3434

@@ -52,15 +52,16 @@ pub async fn test_leader_election_in_three_node_example() {
5252
leader_id
5353
);
5454

55-
raft_2.quit().await;
55+
raft_2.quit().await.expect("Failed to quit");
5656
let raft_3 = rafts.get_mut(&3).unwrap();
57-
raft_3.quit().await;
57+
raft_3.quit().await.expect("Failed to quit");
5858
}
5959

6060
// TODO: Fix this test.
6161
#[tokio::test]
6262
#[ignore]
6363
pub async fn test_leader_election_in_five_node_example() {
64+
cleanup_storage("./logs");
6465
kill_previous_raft_processes();
6566

6667
let (tx_raft, rx_raft) = mpsc::channel::<(u64, Raft)>();
@@ -77,7 +78,7 @@ pub async fn test_leader_election_in_five_node_example() {
7778

7879
sleep(Duration::from_secs(1)).await;
7980

80-
raft_1.leave().await;
81+
raft_1.leave().await.expect("Failed to leave");
8182

8283
let raft_2 = rafts.get_mut(&2).unwrap();
8384

@@ -94,7 +95,7 @@ pub async fn test_leader_election_in_five_node_example() {
9495
);
9596

9697
let leader_raft = rafts.get_mut(&leader_id).unwrap();
97-
leader_raft.leave().await;
98+
leader_raft.leave().await.expect("Failed to leave");
9899

99100
let mut remaining_nodes = vec![2, 3, 4, 5];
100101
if let Some(pos) = remaining_nodes.iter().position(|&x| x == leader_id) {
@@ -115,6 +116,6 @@ pub async fn test_leader_election_in_five_node_example() {
115116

116117
for id in remaining_nodes {
117118
let raft = rafts.get_mut(&id).unwrap();
118-
raft.quit().await;
119+
raft.quit().await.expect("Failed to quit the raft node");
119120
}
120121
}

0 commit comments

Comments
 (0)