Skip to content

Commit f4a8439

Browse files
committed
fix: CI
1 parent dc89235 commit f4a8439

9 files changed

Lines changed: 70 additions & 38 deletions

File tree

binding/python/examples/main.py

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -37,8 +37,9 @@ def load_peers() -> Peers:
3737
)
3838

3939

40-
def build_config(initial_peers: Peers) -> Config:
40+
def build_config(node_id: int, initial_peers: Peers) -> Config:
4141
raft_cfg = RaftConfig(
42+
id=node_id,
4243
election_tick=10,
4344
heartbeat_tick=3,
4445
)
@@ -106,12 +107,12 @@ async def main():
106107

107108
initial_peers = load_peers()
108109

109-
cfg = build_config(initial_peers)
110+
node_id = initial_peers.get_node_id_by_addr(raft_addr)
111+
112+
cfg = build_config(node_id, initial_peers)
110113
logger = Logger(setup_logger())
111114
store = HashStore()
112115

113-
node_id = initial_peers.get_node_id_by_addr(raft_addr)
114-
115116
tasks = []
116117
raft = Raft.bootstrap(node_id, raft_addr, store, cfg, logger)
117118
tasks.append(raft.run())

binding/python/tests/harness/raft_server.py

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,8 +6,9 @@
66
RAFTS: dict[int, Raft] = {}
77

88

9-
def build_config(initial_peers: Peers) -> Config:
9+
def build_config(node_id: int, initial_peers: Peers) -> Config:
1010
raft_cfg = RaftConfig(
11+
id=node_id,
1112
election_tick=10,
1213
heartbeat_tick=3,
1314
)
@@ -23,7 +24,7 @@ def build_config(initial_peers: Peers) -> Config:
2324

2425
async def run_raft(node_id: int, initial_peers: Peers):
2526
peer = initial_peers.get(node_id)
26-
cfg = build_config(initial_peers)
27+
cfg = build_config(node_id, initial_peers)
2728

2829
store = HashStore()
2930
logger = Slogger.default()

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: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -175,6 +175,16 @@ 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).expect("Failed to remove storage directory");
183+
}
184+
185+
fs::create_dir_all(storage_pth).expect("Failed to create storage directory");
186+
}
187+
178188
pub fn kill_previous_raft_processes() {
179189
RAFT_PORTS.iter().for_each(|port| {
180190
kill_process_using_port(*port);

harness/tests/bootstrap.rs

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -4,11 +4,15 @@ 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::{
8+
cleanup_storage, kill_previous_raft_processes, load_peers,
9+
wait_for_until_cluster_size_increase,
10+
},
811
};
912

1013
#[tokio::test]
1114
pub async fn test_static_bootstrap() {
15+
cleanup_storage("./logs");
1216
kill_previous_raft_processes();
1317
let (tx_raft, rx_raft) = mpsc::channel::<(u64, Raft)>();
1418

@@ -23,14 +27,15 @@ pub async fn test_static_bootstrap() {
2327
wait_for_until_cluster_size_increase(raft_1.clone(), 3).await;
2428

2529
for (_, raft) in rafts.iter_mut() {
26-
raft.quit().await;
30+
raft.quit().await.expect("Failed to quit raft node");
2731
}
2832

2933
sleep(Duration::from_secs(1)).await;
3034
}
3135

3236
#[tokio::test]
3337
pub async fn test_dynamic_bootstrap() {
38+
cleanup_storage("./logs");
3439
kill_previous_raft_processes();
3540
let (tx_raft, rx_raft) = mpsc::channel::<(u64, Raft)>();
3641

@@ -64,7 +69,7 @@ pub async fn test_dynamic_bootstrap() {
6469
wait_for_until_cluster_size_increase(raft_1.clone(), 3).await;
6570

6671
for (_, raft) in rafts.iter_mut() {
67-
raft.quit().await;
72+
raft.quit().await.expect("Failed to quit raft node");
6873
}
6974
}
7075

harness/tests/data_replication.rs

Lines changed: 9 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -6,11 +6,15 @@ 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::{
10+
cleanup_storage, kill_previous_raft_processes, load_peers,
11+
wait_for_until_cluster_size_increase,
12+
},
1013
};
1114

1215
#[tokio::test]
1316
pub async fn test_data_replication() {
17+
cleanup_storage("./logs");
1418
kill_previous_raft_processes();
1519

1620
let peers = load_peers(THREE_NODE_EXAMPLE).await.unwrap();
@@ -33,7 +37,7 @@ pub async fn test_data_replication() {
3337

3438
raft_1.propose(entry).await.unwrap();
3539

36-
sleep(Duration::from_secs(1)).await;
40+
sleep(Duration::from_secs(3)).await;
3741

3842
// Data should be replicated to all nodes.
3943
for (_, raft) in rafts.iter_mut() {
@@ -62,7 +66,7 @@ pub async fn test_data_replication() {
6266
let store = raft_4.state_machine().await.unwrap();
6367
let store_lk = store.0.read().unwrap();
6468

65-
// Data should be replicated to new joined node.
69+
// Data should be replicated to new member.
6670
assert_eq!(store_lk.get(&1).unwrap(), "test");
6771
std::mem::drop(store_lk);
6872

@@ -77,15 +81,14 @@ pub async fn test_data_replication() {
7781

7882
raft_1.propose(new_entry).await.unwrap();
7983

80-
// New entry data should be replicated to all nodes including new joined node.
84+
// New entry data should be replicated to all nodes including new member.
8185
for (_, raft) in rafts.iter() {
82-
// stop
8386
let store = raft.state_machine().await.unwrap();
8487
let store_lk = store.0.read().unwrap();
8588
assert_eq!(store_lk.get(&2).unwrap(), "test2");
8689
}
8790

8891
for (_, raft) in rafts.iter_mut() {
89-
raft.quit().await;
92+
raft.quit().await.expect("Failed to quit the raft node");
9093
}
9194
}

harness/tests/leader_election.rs

Lines changed: 10 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -5,13 +5,14 @@ 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,
9+
wait_for_until_cluster_size_decrease, wait_for_until_cluster_size_increase,
1010
},
1111
};
1212

1313
#[tokio::test]
1414
pub async fn test_leader_election_in_three_node_example() {
15+
cleanup_storage("./logs");
1516
kill_previous_raft_processes();
1617

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

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

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

3334
sleep(Duration::from_secs(2)).await;
3435

@@ -52,15 +53,16 @@ pub async fn test_leader_election_in_three_node_example() {
5253
leader_id
5354
);
5455

55-
raft_2.quit().await;
56+
raft_2.quit().await.expect("Failed to quit");
5657
let raft_3 = rafts.get_mut(&3).unwrap();
57-
raft_3.quit().await;
58+
raft_3.quit().await.expect("Failed to quit");
5859
}
5960

6061
// TODO: Fix this test.
6162
#[tokio::test]
6263
#[ignore]
6364
pub async fn test_leader_election_in_five_node_example() {
65+
cleanup_storage("./logs");
6466
kill_previous_raft_processes();
6567

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

7880
sleep(Duration::from_secs(1)).await;
7981

80-
raft_1.leave().await;
82+
raft_1.leave().await.expect("Failed to leave");
8183

8284
let raft_2 = rafts.get_mut(&2).unwrap();
8385

@@ -94,7 +96,7 @@ pub async fn test_leader_election_in_five_node_example() {
9496
);
9597

9698
let leader_raft = rafts.get_mut(&leader_id).unwrap();
97-
leader_raft.leave().await;
99+
leader_raft.leave().await.expect("Failed to leave");
98100

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

116118
for id in remaining_nodes {
117119
let raft = rafts.get_mut(&id).unwrap();
118-
raft.quit().await;
120+
raft.quit().await.expect("Failed to quit the raft node");
119121
}
120122
}

0 commit comments

Comments
 (0)