consensus.rs 22 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714
  1. use async_std::{
  2. sync::{Arc, Mutex},
  3. task,
  4. };
  5. use std::{cmp::min, collections::HashMap, path::PathBuf, time::Duration};
  6. use async_executor::Executor;
  7. use futures::{select, FutureExt};
  8. use log::{debug, error, info, warn};
  9. use rand::{rngs::OsRng, Rng, RngCore};
  10. use url::Url;
  11. use crate::{
  12. net,
  13. util::serial::{deserialize, serialize, Decodable, Encodable},
  14. Error, Result,
  15. };
  16. use super::{
  17. primitives::{
  18. BroadcastMsgRequest, Channel, Log, LogRequest, LogResponse, Logs, MapLength, NetMsg,
  19. NetMsgMethod, NodeId, Role, Sender, SyncRequest, SyncResponse, VoteRequest, VoteResponse,
  20. },
  21. DataStore,
  22. };
  23. // In milliseconds
  24. const HEARTBEATTIMEOUT: u64 = 500;
  25. const TIMEOUT: u64 = 6000;
  26. const TIMEOUT_NODES: u64 = 1000;
  27. const SYNC_TIMEOUT_FOR_EACH_ATTEMPT: u64 = 1000;
  28. const SYNC_ATTEMPTS: u64 = 8;
  29. async fn load_node_ids_loop(
  30. nodes: Arc<Mutex<HashMap<NodeId, Url>>>,
  31. p2p: net::P2pPtr,
  32. role: Role,
  33. ) -> Result<()> {
  34. if role == Role::Listener {
  35. return Ok(())
  36. }
  37. loop {
  38. debug!(target: "raft", "Loading node ids from p2p hosts",);
  39. task::sleep(Duration::from_millis(TIMEOUT_NODES)).await;
  40. let hosts = p2p.hosts().clone();
  41. let nodes_ip = hosts.load_all().await.clone();
  42. for ip in nodes_ip.iter() {
  43. (*nodes.lock().await).insert(NodeId::from(ip.clone()), ip.clone());
  44. }
  45. }
  46. }
  47. async fn p2p_send_loop(receiver: async_channel::Receiver<NetMsg>, p2p: net::P2pPtr) -> Result<()> {
  48. loop {
  49. let msg: NetMsg = receiver.recv().await?;
  50. if let Err(e) = p2p.broadcast(msg).await {
  51. error!(target: "raft", "error occurred during broadcasting a msg: {}", e);
  52. continue
  53. }
  54. }
  55. }
  56. pub struct Raft<T> {
  57. // this will be derived from the ip
  58. pub id: Option<NodeId>,
  59. // these four vars should be on local storage
  60. current_term: u64,
  61. voted_for: Option<NodeId>,
  62. logs: Logs,
  63. commit_length: u64,
  64. role: Role,
  65. current_leader: Option<NodeId>,
  66. votes_received: Vec<NodeId>,
  67. sent_length: MapLength,
  68. acked_length: MapLength,
  69. nodes: Arc<Mutex<HashMap<NodeId, Url>>>,
  70. last_term: u64,
  71. sender: Sender,
  72. msgs_channel: Channel<T>,
  73. commits_channel: Channel<T>,
  74. datastore: DataStore<T>,
  75. seen_msgs: Arc<Mutex<Vec<u64>>>,
  76. }
  77. impl<T: Decodable + Encodable + Clone> Raft<T> {
  78. pub fn new(
  79. addr: Option<Url>,
  80. db_path: PathBuf,
  81. seen_msgs: Arc<Mutex<Vec<u64>>>,
  82. ) -> Result<Self> {
  83. if db_path.to_str().is_none() {
  84. error!(target: "raft", "datastore path is incorrect");
  85. return Err(Error::ParseFailed("unable to parse pathbuf to str"))
  86. };
  87. let datastore = DataStore::new(db_path.to_str().unwrap())?;
  88. // load from sled datastore
  89. let current_term = datastore.current_term.get_last()?.unwrap_or(0);
  90. let voted_for = datastore.voted_for.get_last()?.flatten();
  91. let logs = Logs(datastore.logs.get_all()?);
  92. let commit_length = datastore.commits.get_all()?.len() as u64;
  93. // broadcasting channels
  94. let msgs_channel = async_channel::unbounded::<T>();
  95. let commits_channel = async_channel::unbounded::<T>();
  96. let sender = async_channel::unbounded::<NetMsg>();
  97. let id = addr.map(NodeId::from);
  98. let role = if id.is_some() { Role::Follower } else { Role::Listener };
  99. Ok(Self {
  100. id,
  101. current_term,
  102. voted_for,
  103. logs,
  104. commit_length,
  105. role,
  106. current_leader: None,
  107. votes_received: vec![],
  108. sent_length: MapLength(HashMap::new()),
  109. acked_length: MapLength(HashMap::new()),
  110. nodes: Arc::new(Mutex::new(HashMap::new())),
  111. last_term: 0,
  112. sender,
  113. msgs_channel,
  114. commits_channel,
  115. datastore,
  116. seen_msgs,
  117. })
  118. }
  119. pub async fn start(
  120. &mut self,
  121. p2p: net::P2pPtr,
  122. p2p_recv_channel: async_channel::Receiver<NetMsg>,
  123. executor: Arc<Executor<'_>>,
  124. stop_signal: async_channel::Receiver<()>,
  125. ) -> Result<()> {
  126. let p2p_send_task = executor.spawn(p2p_send_loop(self.sender.1.clone(), p2p.clone()));
  127. let load_ips_task =
  128. executor.spawn(load_node_ids_loop(self.nodes.clone(), p2p.clone(), self.role.clone()));
  129. let mut synced = false;
  130. // Sync listener node
  131. if self.role == Role::Listener {
  132. let last_term =
  133. if !self.logs.0.is_empty() { self.logs.0.last().unwrap().term } else { 0 };
  134. let sync_request = SyncRequest { logs_len: self.logs.len(), last_term };
  135. info!("Start Syncing...");
  136. for _ in 0..SYNC_ATTEMPTS {
  137. if synced {
  138. break
  139. }
  140. self.send(None, &serialize(&sync_request), NetMsgMethod::SyncRequest, None).await?;
  141. synced = self
  142. .waiting_for_sync(
  143. executor.clone(),
  144. p2p_recv_channel.clone(),
  145. stop_signal.clone(),
  146. )
  147. .await?;
  148. }
  149. if synced {
  150. info!("SYNCED SUCCESSFULLY!!");
  151. }
  152. }
  153. let mut rng = rand::thread_rng();
  154. let broadcast_msg_rv = self.msgs_channel.1.clone();
  155. if !synced && self.role == Role::Listener {
  156. error!("SYNCING FAILED!!");
  157. } else {
  158. loop {
  159. let timeout: Duration = if self.role == Role::Leader {
  160. Duration::from_millis(HEARTBEATTIMEOUT)
  161. } else {
  162. Duration::from_millis(rng.gen_range(0..HEARTBEATTIMEOUT) + TIMEOUT)
  163. };
  164. let result: Result<()>;
  165. select! {
  166. m = p2p_recv_channel.recv().fuse() => result = self.handle_method(m?).await,
  167. m = broadcast_msg_rv.recv().fuse() => result = self.broadcast_msg(&m?,None).await,
  168. _ = task::sleep(timeout).fuse() => {
  169. result = if self.role == Role::Leader {
  170. self.send_heartbeat().await
  171. }else {
  172. self.send_vote_request().await
  173. };
  174. },
  175. _ = stop_signal.recv().fuse() => break,
  176. }
  177. match result {
  178. Ok(_) => {}
  179. Err(e) => warn!(target: "raft", "warn: {}", e),
  180. }
  181. }
  182. }
  183. warn!(target: "raft", "Raft Terminating...");
  184. load_ips_task.cancel().await;
  185. p2p_send_task.cancel().await;
  186. self.datastore.flush().await?;
  187. Ok(())
  188. }
  189. pub fn get_commits_channel(&self) -> async_channel::Receiver<T> {
  190. self.commits_channel.1.clone()
  191. }
  192. pub fn get_msgs_channel(&self) -> async_channel::Sender<T> {
  193. self.msgs_channel.0.clone()
  194. }
  195. async fn broadcast_msg(&mut self, msg: &T, msg_id: Option<u64>) -> Result<()> {
  196. if self.role == Role::Leader {
  197. let msg = serialize(msg);
  198. let log = Log { msg, term: self.current_term };
  199. self.push_log(&log)?;
  200. self.acked_length.insert(&self.id.clone().unwrap(), self.logs.len());
  201. } else {
  202. let b_msg = BroadcastMsgRequest(serialize(msg));
  203. self.send(
  204. self.current_leader.clone(),
  205. &serialize(&b_msg),
  206. NetMsgMethod::BroadcastRequest,
  207. msg_id,
  208. )
  209. .await?;
  210. }
  211. debug!(target: "raft", "Role: {:?}, broadcast a msg id: {:?} ", self.role, msg_id);
  212. Ok(())
  213. }
  214. async fn handle_method(&mut self, msg: NetMsg) -> Result<()> {
  215. match msg.method {
  216. NetMsgMethod::LogResponse => {
  217. let lr: LogResponse = deserialize(&msg.payload)?;
  218. self.receive_log_response(lr).await?;
  219. }
  220. NetMsgMethod::LogRequest => {
  221. let lr: LogRequest = deserialize(&msg.payload)?;
  222. self.receive_log_request(lr).await?;
  223. }
  224. NetMsgMethod::VoteResponse => {
  225. let vr: VoteResponse = deserialize(&msg.payload)?;
  226. self.receive_vote_response(vr).await?;
  227. }
  228. NetMsgMethod::VoteRequest => {
  229. let vr: VoteRequest = deserialize(&msg.payload)?;
  230. self.receive_vote_request(vr).await?;
  231. }
  232. NetMsgMethod::BroadcastRequest => {
  233. let vr: BroadcastMsgRequest = deserialize(&msg.payload)?;
  234. let d: T = deserialize(&vr.0)?;
  235. self.broadcast_msg(&d, Some(msg.id)).await?;
  236. }
  237. NetMsgMethod::SyncRequest => {
  238. debug!(target: "raft", "Receive sync request");
  239. let sr: SyncRequest = deserialize(&msg.payload)?;
  240. self.receive_sync_request(&sr, msg.id).await?;
  241. }
  242. NetMsgMethod::SyncResponse => {}
  243. }
  244. debug!(target: "raft", "Role: {:?} receive msg id: {} recipient_id: {:?} method: {:?} ",
  245. self.role, msg.id, &msg.recipient_id.is_some(), &msg.method);
  246. Ok(())
  247. }
  248. async fn receive_sync_request(&self, sr: &SyncRequest, msg_id: u64) -> Result<()> {
  249. if self.role == Role::Leader {
  250. let mut wipe = false;
  251. let logs = if sr.logs_len == 0 {
  252. self.logs.clone()
  253. } else if self.logs.len() >= sr.logs_len &&
  254. self.logs.get(sr.logs_len - 1)?.term == sr.last_term
  255. {
  256. self.logs.slice_from(sr.logs_len).unwrap()
  257. } else {
  258. wipe = true;
  259. self.logs.clone()
  260. };
  261. let sync_response = SyncResponse {
  262. logs,
  263. commit_length: self.commit_length,
  264. leader_id: self.id.clone().unwrap(),
  265. wipe,
  266. };
  267. debug!(target: "raft", "Send sync response");
  268. self.send(None, &serialize(&sync_response), NetMsgMethod::SyncResponse, None).await?;
  269. } else {
  270. self.send(
  271. self.current_leader.clone(),
  272. &serialize(sr),
  273. NetMsgMethod::SyncRequest,
  274. Some(msg_id),
  275. )
  276. .await?;
  277. }
  278. Ok(())
  279. }
  280. async fn receive_sync_response(&mut self, sr: &SyncResponse) -> Result<()> {
  281. debug!(target: "raft", "Receive sync response");
  282. if sr.wipe {
  283. self.set_commit_length(&0)?;
  284. self.push_logs(&sr.logs)?;
  285. } else {
  286. for log in sr.logs.0.iter() {
  287. self.push_log(log)?;
  288. }
  289. }
  290. if !self.logs.is_empty() {
  291. self.set_current_term(&self.logs.0.last().unwrap().term.clone())?;
  292. }
  293. for i in self.commit_length..sr.commit_length {
  294. self.push_commit(&self.logs.get(i)?.msg).await?;
  295. }
  296. self.set_commit_length(&sr.commit_length)?;
  297. self.current_leader = Some(sr.leader_id.clone());
  298. Ok(())
  299. }
  300. async fn send(
  301. &self,
  302. recipient_id: Option<NodeId>,
  303. payload: &[u8],
  304. method: NetMsgMethod,
  305. msg_id: Option<u64>,
  306. ) -> Result<()> {
  307. let random_id = if msg_id.is_some() { msg_id.unwrap() } else { OsRng.next_u64() };
  308. debug!(target: "raft","Role: {:?} send a msg id: {} recipient_id: {:?} method: {:?} ",
  309. self.role, random_id, &recipient_id.is_some(), &method);
  310. let net_msg = NetMsg { id: random_id, recipient_id, payload: payload.to_vec(), method };
  311. self.seen_msgs.lock().await.push(random_id);
  312. self.sender.0.send(net_msg).await?;
  313. Ok(())
  314. }
  315. async fn waiting_for_sync(
  316. &mut self,
  317. executor: Arc<Executor<'_>>,
  318. p2p_recv_channel: async_channel::Receiver<NetMsg>,
  319. stop_signal: async_channel::Receiver<()>,
  320. ) -> Result<bool> {
  321. let (timeout_s, timeout_r) = async_channel::unbounded::<()>();
  322. executor
  323. .spawn(async move {
  324. task::sleep(Duration::from_millis(SYNC_TIMEOUT_FOR_EACH_ATTEMPT)).await;
  325. timeout_s.send(()).await.unwrap_or(());
  326. })
  327. .detach();
  328. loop {
  329. select! {
  330. msg = p2p_recv_channel.recv().fuse() => {
  331. let msg = msg?;
  332. if msg.method == NetMsgMethod::SyncResponse {
  333. let sr: SyncResponse = deserialize(&msg.payload)?;
  334. self.receive_sync_response(&sr).await?;
  335. return Ok(true)
  336. }},
  337. _ = stop_signal.recv().fuse() => break,
  338. _ = timeout_r.recv().fuse() => break,
  339. }
  340. }
  341. Ok(false)
  342. }
  343. async fn send_heartbeat(&self) -> Result<()> {
  344. if self.role == Role::Leader {
  345. let nodes = self.nodes.lock().await;
  346. let nodes_cloned = nodes.clone();
  347. drop(nodes);
  348. for node in nodes_cloned.iter() {
  349. self.update_logs(node.0).await?;
  350. }
  351. }
  352. Ok(())
  353. }
  354. async fn send_vote_request(&mut self) -> Result<()> {
  355. if self.role == Role::Listener {
  356. return Ok(())
  357. }
  358. let self_id = self.id.clone().unwrap();
  359. self.set_current_term(&(self.current_term + 1))?;
  360. self.role = Role::Candidate;
  361. self.set_voted_for(&Some(self_id.clone()))?;
  362. self.votes_received = vec![];
  363. self.votes_received.push(self_id.clone());
  364. self.reset_last_term();
  365. let request = VoteRequest {
  366. node_id: self_id,
  367. current_term: self.current_term,
  368. log_length: self.logs.len(),
  369. last_term: self.last_term,
  370. };
  371. let payload = serialize(&request);
  372. self.send(None, &payload, NetMsgMethod::VoteRequest, None).await
  373. }
  374. async fn receive_vote_request(&mut self, vr: VoteRequest) -> Result<()> {
  375. if self.role == Role::Listener {
  376. return Ok(())
  377. }
  378. if vr.current_term > self.current_term {
  379. self.set_current_term(&vr.current_term)?;
  380. self.set_voted_for(&None)?;
  381. self.role = Role::Follower;
  382. }
  383. self.reset_last_term();
  384. // check the logs of the candidate
  385. let vote_ok = (vr.last_term > self.last_term) ||
  386. (vr.last_term == self.last_term && vr.log_length >= self.logs.len());
  387. // slef.voted_for equal to vr.node_id or is None or voted to someone else
  388. let vote = if let Some(voted_for) = self.voted_for.as_ref() {
  389. *voted_for == vr.node_id
  390. } else {
  391. true
  392. };
  393. let mut response = VoteResponse {
  394. node_id: self.id.clone().unwrap(),
  395. current_term: self.current_term,
  396. ok: false,
  397. };
  398. if vr.current_term == self.current_term && vote_ok && vote {
  399. self.set_voted_for(&Some(vr.node_id.clone()))?;
  400. response.set_ok(true);
  401. }
  402. let payload = serialize(&response);
  403. self.send(Some(vr.node_id), &payload, NetMsgMethod::VoteResponse, None).await
  404. }
  405. async fn receive_vote_response(&mut self, vr: VoteResponse) -> Result<()> {
  406. if self.role == Role::Listener {
  407. return Ok(())
  408. }
  409. if self.role == Role::Candidate && vr.current_term == self.current_term && vr.ok {
  410. self.votes_received.push(vr.node_id);
  411. let nodes = self.nodes.lock().await;
  412. let nodes_cloned = nodes.clone();
  413. drop(nodes);
  414. if self.votes_received.len() >= (nodes_cloned.len() / 2) {
  415. self.role = Role::Leader;
  416. self.current_leader = Some(self.id.clone().unwrap());
  417. for node in nodes_cloned.iter() {
  418. self.sent_length.insert(node.0, self.logs.len());
  419. self.acked_length.insert(node.0, 0);
  420. }
  421. }
  422. } else if vr.current_term > self.current_term {
  423. self.set_current_term(&vr.current_term)?;
  424. self.role = Role::Follower;
  425. self.set_voted_for(&None)?;
  426. }
  427. Ok(())
  428. }
  429. async fn update_logs(&self, node_id: &NodeId) -> Result<()> {
  430. let prefix_len = match self.sent_length.get(node_id) {
  431. Ok(len) => len,
  432. Err(_) => {
  433. // return if failed to index
  434. return Ok(())
  435. }
  436. };
  437. let suffix: Logs = if self.logs.slice_from(prefix_len).is_some() {
  438. self.logs.slice_from(prefix_len).unwrap()
  439. } else {
  440. return Ok(())
  441. };
  442. let mut prefix_term = 0;
  443. if prefix_len > 0 {
  444. prefix_term = self.logs.get(prefix_len - 1)?.term;
  445. }
  446. let request = LogRequest {
  447. leader_id: self.id.clone().unwrap(),
  448. current_term: self.current_term,
  449. prefix_len,
  450. prefix_term,
  451. commit_length: self.commit_length,
  452. suffix,
  453. };
  454. let payload = serialize(&request);
  455. self.send(Some(node_id.clone()), &payload, NetMsgMethod::LogRequest, None).await
  456. }
  457. async fn receive_log_request(&mut self, lr: LogRequest) -> Result<()> {
  458. if lr.current_term > self.current_term {
  459. self.set_current_term(&lr.current_term)?;
  460. self.set_voted_for(&None)?;
  461. }
  462. if lr.current_term == self.current_term {
  463. if self.role != Role::Listener {
  464. self.role = Role::Follower;
  465. }
  466. self.current_leader = Some(lr.leader_id.clone());
  467. }
  468. let mut ok = (self.logs.len() >= lr.prefix_len) &&
  469. (lr.prefix_len == 0 || self.logs.get(lr.prefix_len - 1)?.term == lr.prefix_term);
  470. let mut ack = 0;
  471. if lr.current_term == self.current_term && ok {
  472. self.append_log(lr.prefix_len, lr.commit_length, &lr.suffix).await?;
  473. ack = lr.prefix_len + lr.suffix.len();
  474. } else {
  475. ok = false;
  476. }
  477. if self.role == Role::Listener {
  478. return Ok(())
  479. }
  480. let response = LogResponse {
  481. node_id: self.id.clone().unwrap(),
  482. current_term: self.current_term,
  483. ack,
  484. ok,
  485. };
  486. let payload = serialize(&response);
  487. self.send(Some(lr.leader_id.clone()), &payload, NetMsgMethod::LogResponse, None).await
  488. }
  489. async fn receive_log_response(&mut self, lr: LogResponse) -> Result<()> {
  490. if lr.current_term == self.current_term && self.role == Role::Leader {
  491. if lr.ok && lr.ack >= self.acked_length.get(&lr.node_id)? {
  492. self.sent_length.insert(&lr.node_id, lr.ack);
  493. self.acked_length.insert(&lr.node_id, lr.ack);
  494. self.commit_log().await?;
  495. } else if self.sent_length.get(&lr.node_id)? > 0 {
  496. self.sent_length.insert(&lr.node_id, self.sent_length.get(&lr.node_id)? - 1);
  497. }
  498. } else if lr.current_term > self.current_term {
  499. self.set_current_term(&lr.current_term)?;
  500. if self.role != Role::Listener {
  501. self.role = Role::Follower;
  502. }
  503. self.set_voted_for(&None)?;
  504. }
  505. Ok(())
  506. }
  507. fn reset_last_term(&mut self) {
  508. self.last_term = 0;
  509. if let Some(log) = self.logs.0.last() {
  510. self.last_term = log.term;
  511. }
  512. }
  513. fn acks(&self, nodes: HashMap<NodeId, Url>, length: u64) -> HashMap<NodeId, Url> {
  514. nodes
  515. .into_iter()
  516. .filter(|n| {
  517. let len = self.acked_length.get(&n.0);
  518. len.is_ok() && len.unwrap() >= length
  519. })
  520. .collect()
  521. }
  522. async fn commit_log(&mut self) -> Result<()> {
  523. let nodes_ptr = self.nodes.lock().await;
  524. let min_acks = ((nodes_ptr.len() + 1) / 2) as usize;
  525. let nodes = nodes_ptr.clone();
  526. drop(nodes_ptr);
  527. let mut ready: Vec<u64> = vec![];
  528. for len in 1..(self.logs.len() + 1) {
  529. if self.acks(nodes.clone(), len).len() >= min_acks {
  530. ready.push(len);
  531. }
  532. }
  533. if ready.is_empty() {
  534. return Ok(())
  535. }
  536. let max_ready = *ready.iter().max().unwrap();
  537. if max_ready > self.commit_length && self.logs.get(max_ready - 1)?.term == self.current_term
  538. {
  539. for i in self.commit_length..max_ready {
  540. self.push_commit(&self.logs.get(i)?.msg).await?;
  541. }
  542. self.set_commit_length(&max_ready)?;
  543. }
  544. Ok(())
  545. }
  546. async fn append_log(
  547. &mut self,
  548. prefix_len: u64,
  549. leader_commit: u64,
  550. suffix: &Logs,
  551. ) -> Result<()> {
  552. if !suffix.is_empty() && self.logs.len() > prefix_len {
  553. let index = min(self.logs.len(), prefix_len + suffix.len()) - 1;
  554. if self.logs.get(index)?.term != suffix.get(index - prefix_len)?.term {
  555. self.push_logs(&self.logs.slice_to(prefix_len))?;
  556. }
  557. }
  558. if prefix_len + suffix.len() > self.logs.len() {
  559. for i in (self.logs.len() - prefix_len)..suffix.len() {
  560. self.push_log(&suffix.get(i)?)?;
  561. }
  562. }
  563. if leader_commit > self.commit_length {
  564. for i in self.commit_length..leader_commit {
  565. self.push_commit(&self.logs.get(i)?.msg).await?;
  566. }
  567. self.set_commit_length(&leader_commit)?;
  568. }
  569. Ok(())
  570. }
  571. fn set_commit_length(&mut self, i: &u64) -> Result<()> {
  572. self.commit_length = *i;
  573. Ok(())
  574. }
  575. fn set_current_term(&mut self, i: &u64) -> Result<()> {
  576. self.current_term = *i;
  577. self.datastore.current_term.insert(i)
  578. }
  579. fn set_voted_for(&mut self, i: &Option<NodeId>) -> Result<()> {
  580. self.voted_for = i.clone();
  581. self.datastore.voted_for.insert(i)
  582. }
  583. async fn push_commit(&mut self, commit: &[u8]) -> Result<()> {
  584. let commit: T = deserialize(commit)?;
  585. self.commits_channel.0.send(commit.clone()).await?;
  586. self.datastore.commits.insert(&commit)
  587. }
  588. fn push_log(&mut self, log: &Log) -> Result<()> {
  589. self.logs.push(log);
  590. self.datastore.logs.insert(log)
  591. }
  592. fn push_logs(&mut self, logs: &Logs) -> Result<()> {
  593. self.logs = logs.clone();
  594. self.datastore.logs.wipe_insert_all(&logs.to_vec())
  595. }
  596. }