main.rs 6.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224
  1. use async_executor::Executor;
  2. use async_std::sync::{Arc, Mutex};
  3. use easy_parallel::Parallel;
  4. use std::{error, fs::File, io::stdin};
  5. use log::debug;
  6. use simplelog::WriteLogger;
  7. use url::Url;
  8. use darkfi::{net, net::Settings, rpc::server::listen_and_serve};
  9. use crate::{
  10. dchat_error::ErrorMissingSpecifier,
  11. dchatmsg::{DchatMsg, DchatMsgsBuffer},
  12. protocol_dchat::ProtocolDchat,
  13. rpc::JsonRpcInterface,
  14. };
  15. pub mod dchat_error;
  16. pub mod dchatmsg;
  17. pub mod protocol_dchat;
  18. pub mod rpc;
  19. pub type Error = Box<dyn error::Error>;
  20. pub type Result<T> = std::result::Result<T, Error>;
  21. struct Dchat {
  22. p2p: net::P2pPtr,
  23. recv_msgs: DchatMsgsBuffer,
  24. }
  25. impl Dchat {
  26. fn new(p2p: net::P2pPtr, recv_msgs: DchatMsgsBuffer) -> Self {
  27. Self { p2p, recv_msgs }
  28. }
  29. async fn menu(&self) -> Result<()> {
  30. let mut buffer = String::new();
  31. let stdin = stdin();
  32. loop {
  33. println!(
  34. "Welcome to dchat.
  35. s: send message
  36. i: inbox
  37. q: quit "
  38. );
  39. stdin.read_line(&mut buffer)?;
  40. // Remove trailing \n
  41. buffer.pop();
  42. match buffer.as_str() {
  43. "q" => return Ok(()),
  44. "s" => {
  45. // Remove trailing s
  46. buffer.pop();
  47. stdin.read_line(&mut buffer)?;
  48. match self.send(buffer.clone()).await {
  49. Ok(_) => {
  50. println!("you sent: {}", buffer);
  51. }
  52. Err(e) => {
  53. println!("send failed for reason: {}", e);
  54. }
  55. }
  56. buffer.clear();
  57. }
  58. "i" => {
  59. let msgs = self.recv_msgs.lock().await;
  60. if msgs.is_empty() {
  61. println!("inbox is empty")
  62. } else {
  63. println!("received:");
  64. for i in msgs.iter() {
  65. if !i.msg.is_empty() {
  66. println!("{}", i.msg);
  67. }
  68. }
  69. }
  70. buffer.clear();
  71. }
  72. _ => {}
  73. }
  74. }
  75. }
  76. async fn register_protocol(&self, msgs: DchatMsgsBuffer) -> Result<()> {
  77. debug!(target: "dchat", "Dchat::register_protocol() [START]");
  78. let registry = self.p2p.protocol_registry();
  79. registry
  80. .register(!net::SESSION_SEED, move |channel, _p2p| {
  81. let msgs2 = msgs.clone();
  82. async move { ProtocolDchat::init(channel, msgs2).await }
  83. })
  84. .await;
  85. debug!(target: "dchat", "Dchat::register_protocol() [STOP]");
  86. Ok(())
  87. }
  88. async fn start(&mut self, ex: Arc<Executor<'_>>) -> Result<()> {
  89. debug!(target: "dchat", "Dchat::start() [START]");
  90. let ex2 = ex.clone();
  91. self.register_protocol(self.recv_msgs.clone()).await?;
  92. self.p2p.clone().start(ex.clone()).await?;
  93. ex2.spawn(self.p2p.clone().run(ex.clone())).detach();
  94. self.menu().await?;
  95. self.p2p.stop().await;
  96. debug!(target: "dchat", "Dchat::start() [STOP]");
  97. Ok(())
  98. }
  99. async fn send(&self, msg: String) -> Result<()> {
  100. let dchatmsg = DchatMsg { msg };
  101. self.p2p.broadcast(dchatmsg).await?;
  102. Ok(())
  103. }
  104. }
  105. #[derive(Clone, Debug)]
  106. struct AppSettings {
  107. accept_addr: Url,
  108. net: Settings,
  109. }
  110. impl AppSettings {
  111. pub fn new(accept_addr: Url, net: Settings) -> Self {
  112. Self { accept_addr, net }
  113. }
  114. }
  115. fn alice() -> Result<AppSettings> {
  116. let log_level = simplelog::LevelFilter::Debug;
  117. let log_config = simplelog::Config::default();
  118. let log_path = "/tmp/alice.log";
  119. let file = File::create(log_path).unwrap();
  120. WriteLogger::init(log_level, log_config, file)?;
  121. let seed = Url::parse("tcp://127.0.0.1:50515").unwrap();
  122. let inbound = Url::parse("tcp://127.0.0.1:51554").unwrap();
  123. let ext_addr = Url::parse("tcp://127.0.0.1:51554").unwrap();
  124. let net = Settings {
  125. inbound: vec![inbound],
  126. external_addr: vec![ext_addr],
  127. seeds: vec![seed],
  128. ..Default::default()
  129. };
  130. let accept_addr = Url::parse("tcp://127.0.0.1:55054").unwrap();
  131. let settings = AppSettings::new(accept_addr, net);
  132. Ok(settings)
  133. }
  134. fn bob() -> Result<AppSettings> {
  135. let log_level = simplelog::LevelFilter::Debug;
  136. let log_config = simplelog::Config::default();
  137. let log_path = "/tmp/bob.log";
  138. let file = File::create(log_path).unwrap();
  139. WriteLogger::init(log_level, log_config, file)?;
  140. let seed = Url::parse("tcp://127.0.0.1:50515").unwrap();
  141. let net = Settings {
  142. inbound: vec![],
  143. outbound_connections: 5,
  144. seeds: vec![seed],
  145. ..Default::default()
  146. };
  147. let accept_addr = Url::parse("tcp://127.0.0.1:51054").unwrap();
  148. let settings = AppSettings::new(accept_addr, net);
  149. Ok(settings)
  150. }
  151. #[async_std::main]
  152. async fn main() -> Result<()> {
  153. let settings: Result<AppSettings> = match std::env::args().nth(1) {
  154. Some(id) => match id.as_str() {
  155. "a" => alice(),
  156. "b" => bob(),
  157. _ => Err(ErrorMissingSpecifier.into()),
  158. },
  159. None => Err(ErrorMissingSpecifier.into()),
  160. };
  161. let settings = settings?.clone();
  162. let p2p = net::P2p::new(settings.net).await;
  163. let nthreads = num_cpus::get();
  164. let (signal, shutdown) = async_channel::unbounded::<()>();
  165. let ex = Arc::new(Executor::new());
  166. let ex2 = ex.clone();
  167. let ex3 = ex2.clone();
  168. let msgs: DchatMsgsBuffer = Arc::new(Mutex::new(vec![DchatMsg { msg: String::new() }]));
  169. let mut dchat = Dchat::new(p2p.clone(), msgs);
  170. let accept_addr = settings.accept_addr.clone();
  171. let rpc = Arc::new(JsonRpcInterface { addr: accept_addr.clone(), p2p });
  172. ex.spawn(async move { listen_and_serve(accept_addr.clone(), rpc).await }).detach();
  173. let (_, result) = Parallel::new()
  174. .each(0..nthreads, |_| smol::future::block_on(ex2.run(shutdown.recv())))
  175. .finish(|| {
  176. smol::future::block_on(async move {
  177. dchat.start(ex3).await?;
  178. drop(signal);
  179. Ok(())
  180. })
  181. });
  182. result
  183. }