mod.rs 7.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273
  1. use std::future::Future;
  2. use std::sync::atomic::{AtomicUsize, Ordering};
  3. use std::sync::Arc;
  4. use std::time::{Duration, Instant};
  5. use async_trait::async_trait;
  6. use crossbeam_channel::{Receiver, Sender};
  7. use once_cell::sync::OnceCell;
  8. use kvraft::{
  9. GetArgs, KVServer, PutAppendArgs, PutAppendEnum, UniqueId, UniqueKVOp,
  10. };
  11. use ruaft::{
  12. AppendEntriesArgs, AppendEntriesReply, InstallSnapshotArgs,
  13. InstallSnapshotReply, Raft, RemoteRaft, ReplicableCommand, RequestVoteArgs,
  14. RequestVoteReply,
  15. };
  16. use crate::Persister;
  17. type RaftId = usize;
  18. pub struct EventHandle {
  19. pub from: RaftId,
  20. pub to: RaftId,
  21. sender: futures_channel::oneshot::Sender<std::io::Result<()>>,
  22. }
  23. struct EventStub {
  24. receiver: futures_channel::oneshot::Receiver<std::io::Result<()>>,
  25. }
  26. fn create_event_pair(from: RaftId, to: RaftId) -> (EventHandle, EventStub) {
  27. let (sender, receiver) = futures_channel::oneshot::channel();
  28. (EventHandle { from, to, sender }, EventStub { receiver })
  29. }
  30. impl EventHandle {
  31. pub fn unblock(self) {
  32. self.sender.send(Ok(())).unwrap();
  33. }
  34. pub fn reply_error(self, e: std::io::Error) {
  35. self.sender.send(Err(e)).unwrap();
  36. }
  37. pub fn reply_interrupted_error(self) {
  38. self.reply_error(std::io::Error::from(std::io::ErrorKind::Interrupted))
  39. }
  40. }
  41. impl EventStub {
  42. pub async fn wait(self) -> std::io::Result<()> {
  43. self.receiver.await.unwrap_or(Ok(()))
  44. }
  45. }
  46. pub enum RaftRpcEvent<T> {
  47. RequestVoteRequest(RequestVoteArgs),
  48. RequestVoteResponse(RequestVoteArgs, RequestVoteReply),
  49. AppendEntriesRequest(AppendEntriesArgs<T>),
  50. AppendEntriesResponse(AppendEntriesArgs<T>, AppendEntriesReply),
  51. InstallSnapshotRequest(InstallSnapshotArgs),
  52. InstallSnapshotResponse(InstallSnapshotArgs, InstallSnapshotReply),
  53. }
  54. struct InterceptingRpcClient<T> {
  55. from: RaftId,
  56. to: RaftId,
  57. target: OnceCell<Raft<T>>,
  58. event_queue: Sender<(RaftRpcEvent<T>, EventHandle)>,
  59. }
  60. impl<T> InterceptingRpcClient<T> {
  61. async fn intercept(&self, event: RaftRpcEvent<T>) -> std::io::Result<()> {
  62. let (handle, stub) = create_event_pair(self.from, self.to);
  63. let _ = self.event_queue.send((event, handle));
  64. stub.wait().await
  65. }
  66. pub fn set_raft(&self, target: Raft<T>) {
  67. self.target
  68. .set(target)
  69. .map_err(|_| ())
  70. .expect("Raft should only be set once");
  71. }
  72. }
  73. #[async_trait]
  74. impl<T: ReplicableCommand> RemoteRaft<T> for &InterceptingRpcClient<T> {
  75. async fn request_vote(
  76. &self,
  77. args: RequestVoteArgs,
  78. ) -> std::io::Result<RequestVoteReply> {
  79. let event_result = self
  80. .intercept(RaftRpcEvent::RequestVoteRequest(args.clone()))
  81. .await;
  82. if let Err(e) = event_result {
  83. return Err(e);
  84. };
  85. let reply = self.target.wait().process_request_vote(args.clone());
  86. self.intercept(RaftRpcEvent::RequestVoteResponse(args, reply.clone()))
  87. .await
  88. .map(|_| reply)
  89. }
  90. async fn append_entries(
  91. &self,
  92. args: AppendEntriesArgs<T>,
  93. ) -> std::io::Result<AppendEntriesReply> {
  94. let args_clone = args.clone();
  95. let event_result = self
  96. .intercept(RaftRpcEvent::AppendEntriesRequest(args_clone))
  97. .await;
  98. if let Err(e) = event_result {
  99. return Err(e);
  100. };
  101. let reply = self.target.wait().process_append_entries(args.clone());
  102. self.intercept(RaftRpcEvent::AppendEntriesResponse(args, reply.clone()))
  103. .await
  104. .map(|_| reply)
  105. }
  106. async fn install_snapshot(
  107. &self,
  108. args: InstallSnapshotArgs,
  109. ) -> std::io::Result<InstallSnapshotReply> {
  110. let event_result = self
  111. .intercept(RaftRpcEvent::InstallSnapshotRequest(args.clone()))
  112. .await;
  113. if let Err(e) = event_result {
  114. return Err(e);
  115. };
  116. let reply = self.target.wait().process_install_snapshot(args.clone());
  117. self.intercept(RaftRpcEvent::InstallSnapshotResponse(
  118. args,
  119. reply.clone(),
  120. ))
  121. .await
  122. .map(|_| reply)
  123. }
  124. }
  125. pub struct EventQueue<T> {
  126. pub receiver: Receiver<(RaftRpcEvent<T>, EventHandle)>,
  127. }
  128. fn make_grid_clients<T>(
  129. server_count: usize,
  130. ) -> (EventQueue<T>, Vec<Vec<InterceptingRpcClient<T>>>) {
  131. let (sender, receiver) = crossbeam_channel::unbounded();
  132. let mut all_clients = vec![];
  133. for from in 0..server_count {
  134. let mut clients = vec![];
  135. for to in 0..server_count {
  136. let interceptor = InterceptingRpcClient {
  137. from,
  138. to,
  139. target: Default::default(),
  140. event_queue: sender.clone(),
  141. };
  142. clients.push(interceptor);
  143. }
  144. all_clients.push(clients);
  145. }
  146. (EventQueue { receiver }, all_clients)
  147. }
  148. pub struct Config {
  149. pub event_queue: EventQueue<UniqueKVOp>,
  150. pub kv_servers: Vec<Arc<KVServer>>,
  151. seq: AtomicUsize,
  152. }
  153. impl Config {
  154. pub fn find_leader(&self) -> Option<&KVServer> {
  155. let start = Instant::now();
  156. while start.elapsed() < Duration::from_secs(1) {
  157. if let Some(kv_server) = self
  158. .kv_servers
  159. .iter()
  160. .find(|kv_server| kv_server.raft().get_state().1)
  161. {
  162. return Some(kv_server.as_ref());
  163. }
  164. }
  165. None
  166. }
  167. pub async fn put(&self, key: String, value: String) -> Result<(), ()> {
  168. let kv_server = self.find_leader().unwrap();
  169. let result = kv_server
  170. .put_append(PutAppendArgs {
  171. key,
  172. value,
  173. op: PutAppendEnum::Put,
  174. unique_id: UniqueId {
  175. clerk_id: 1,
  176. sequence_id: self.seq.fetch_add(1, Ordering::Relaxed)
  177. as u64,
  178. },
  179. })
  180. .await;
  181. result.result.map_err(|_| ())
  182. }
  183. pub fn spawn_put(
  184. self: &Arc<Self>,
  185. key: String,
  186. value: String,
  187. ) -> impl Future<Output = Result<(), ()>> {
  188. let this = self.clone();
  189. async move { this.put(key, value).await }
  190. }
  191. pub async fn get(&self, key: String) -> Result<String, ()> {
  192. let kv_server = self.find_leader().unwrap();
  193. let result = kv_server.get(GetArgs { key }).await;
  194. result.result.map(|v| v.unwrap_or_default()).map_err(|_| ())
  195. }
  196. pub fn spawn_get(
  197. self: &Arc<Self>,
  198. key: String,
  199. ) -> impl Future<Output = Result<String, ()>> {
  200. let this = self.clone();
  201. async move { this.get(key).await }
  202. }
  203. }
  204. pub fn make_config(server_count: usize, max_state: Option<usize>) -> Config {
  205. let (event_queue, clients) = make_grid_clients(server_count);
  206. let persister = Arc::new(Persister::new());
  207. let mut kv_servers = vec![];
  208. let clients: Vec<Vec<&'static InterceptingRpcClient<UniqueKVOp>>> = clients
  209. .into_iter()
  210. .map(|v| {
  211. v.into_iter()
  212. .map(|c| {
  213. let c = Box::leak(Box::new(c));
  214. &*c
  215. })
  216. .collect()
  217. })
  218. .collect();
  219. for (index, client_vec) in clients.iter().enumerate() {
  220. let kv_server = KVServer::new(
  221. client_vec.to_vec(),
  222. index,
  223. persister.clone(),
  224. max_state,
  225. );
  226. kv_servers.push(kv_server);
  227. }
  228. for clients in clients.iter() {
  229. for j in 0..server_count {
  230. clients[j].set_raft(kv_servers[j].raft().clone());
  231. }
  232. }
  233. Config {
  234. event_queue,
  235. kv_servers,
  236. seq: AtomicUsize::new(0),
  237. }
  238. }