consensus.rs 21 KB

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