main.rs 6.8 KB

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