peer.rs 1.3 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455
  1. use std::net::SocketAddr;
  2. use async_executor::Executor;
  3. use async_std::sync::Arc;
  4. use clap::Parser;
  5. use easy_parallel::Parallel;
  6. use simplelog::{ColorChoice, Config, LevelFilter, TermLogger, TerminalMode};
  7. use sraft::{Raft, RaftRpc};
  8. #[derive(Parser)]
  9. struct Args {
  10. #[clap(long, short)]
  11. peer: Vec<SocketAddr>,
  12. #[clap(long, short)]
  13. id: u64,
  14. #[clap(long, short)]
  15. listen: SocketAddr,
  16. }
  17. #[async_std::main]
  18. async fn main() {
  19. let args = Args::parse();
  20. TermLogger::init(LevelFilter::Debug, Config::default(), TerminalMode::Mixed, ColorChoice::Auto)
  21. .unwrap();
  22. let mut raft = Raft::new(args.id);
  23. for (k, v) in args.peer.iter().enumerate() {
  24. raft.peers.insert(k as u64, *v);
  25. }
  26. let raft_rpc = RaftRpc(args.listen);
  27. let ex = Arc::new(Executor::new());
  28. let (_signal, shutdown) = async_channel::unbounded::<()>();
  29. Parallel::new()
  30. .each(0..4, |_| smol::future::block_on(ex.run(shutdown.recv())))
  31. //
  32. .add(|| {
  33. smol::future::block_on(async move {
  34. raft_rpc.start().await;
  35. });
  36. Ok(())
  37. })
  38. //
  39. .finish(|| {
  40. smol::future::block_on(async move {
  41. raft.start().await;
  42. })
  43. });
  44. }