main.rs 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344
  1. use async_executor::Executor;
  2. use async_std::sync::{Arc, Mutex};
  3. use easy_parallel::Parallel;
  4. use std::{
  5. fs::File,
  6. io::{stdin, stdout, Read, Write},
  7. };
  8. use log::debug;
  9. use simplelog::WriteLogger;
  10. use url::Url;
  11. use termion::{event::Key, input::TermRead, raw::IntoRawMode};
  12. use darkfi::{
  13. net,
  14. net::Settings,
  15. util::cli::{get_log_config, get_log_level},
  16. };
  17. use crate::{
  18. dchatmsg::{Dchatmsg, DchatmsgsBuffer},
  19. error::{Error, MissingSpecifier, Result},
  20. protocol_dchat::ProtocolDchat,
  21. };
  22. pub mod dchatmsg;
  23. pub mod error;
  24. pub mod protocol_dchat;
  25. struct Dchat {
  26. p2p: net::P2pPtr,
  27. recv_msgs: DchatmsgsBuffer,
  28. input: String,
  29. display: DisplayMode,
  30. }
  31. enum DisplayMode {
  32. Normal,
  33. Editing,
  34. Inbox,
  35. MessageSent,
  36. SendFailed(Error),
  37. }
  38. impl Dchat {
  39. fn new(
  40. p2p: net::P2pPtr,
  41. recv_msgs: DchatmsgsBuffer,
  42. input: String,
  43. display: DisplayMode,
  44. ) -> Self {
  45. Self { p2p, recv_msgs, input, display }
  46. }
  47. async fn menu(&mut self) -> Result<()> {
  48. debug!(target: "dchat", "Dchat::menu() [START]");
  49. let stdout = stdout();
  50. let mut stdout = stdout.lock().into_raw_mode().unwrap();
  51. let mut stdin = stdin();
  52. loop {
  53. self.render().await?;
  54. for k in stdin.by_ref().keys() {
  55. match &self.display {
  56. DisplayMode::Normal => match k.unwrap() {
  57. Key::Char('q') => return Ok(()),
  58. Key::Char('i') => {
  59. self.display = DisplayMode::Inbox;
  60. break
  61. }
  62. Key::Char('s') => {
  63. self.display = DisplayMode::Editing;
  64. break
  65. }
  66. _ => {}
  67. },
  68. DisplayMode::Editing => match k.unwrap() {
  69. Key::Char('q') => return Ok(()),
  70. Key::Char('\n') => {
  71. match self.send().await {
  72. Ok(_) => {
  73. self.display = DisplayMode::MessageSent;
  74. }
  75. Err(e) => {
  76. self.display = DisplayMode::SendFailed(e);
  77. }
  78. }
  79. break
  80. }
  81. Key::Char(c) => {
  82. self.input.push(c);
  83. }
  84. Key::Esc => {
  85. self.display = DisplayMode::Normal;
  86. break
  87. }
  88. _ => {}
  89. },
  90. DisplayMode::MessageSent => match k.unwrap() {
  91. Key::Char('q') => return Ok(()),
  92. Key::Esc => {
  93. self.display = DisplayMode::Normal;
  94. break
  95. }
  96. _ => {}
  97. },
  98. DisplayMode::Inbox => match k.unwrap() {
  99. Key::Char('q') => return Ok(()),
  100. _ => {}
  101. },
  102. DisplayMode::SendFailed(_) => match k.unwrap() {
  103. Key::Char('q') => return Ok(()),
  104. Key::Esc => {
  105. self.display = DisplayMode::Normal;
  106. break
  107. }
  108. _ => {}
  109. },
  110. }
  111. }
  112. stdout.flush()?;
  113. }
  114. }
  115. async fn render(&mut self) -> Result<()> {
  116. debug!(target: "dchat", "Dchat::render() [START]");
  117. let stdout = stdout();
  118. let mut stdout = stdout.lock().into_raw_mode().unwrap();
  119. match &self.display {
  120. DisplayMode::Normal => {
  121. write!(
  122. stdout,
  123. "{}{}{}Welcome to dchat. {} s: send message {} i: inbox {} q: quit {}",
  124. termion::clear::All,
  125. termion::style::Bold,
  126. termion::cursor::Goto(1, 2),
  127. termion::cursor::Goto(1, 3),
  128. termion::cursor::Goto(1, 4),
  129. termion::cursor::Goto(1, 5),
  130. termion::cursor::Goto(1, 6)
  131. )?;
  132. stdout.flush()?;
  133. }
  134. DisplayMode::Editing => {
  135. write!(
  136. stdout,
  137. "{}{}{}enter your msg.{} esc: stop editing {} enter: send {}",
  138. termion::clear::All,
  139. termion::style::Bold,
  140. termion::cursor::Goto(1, 2),
  141. termion::cursor::Goto(1, 3),
  142. termion::cursor::Goto(1, 4),
  143. termion::cursor::Goto(1, 5)
  144. )?;
  145. stdout.flush()?;
  146. }
  147. DisplayMode::Inbox => {
  148. let msgs = self.recv_msgs.lock().await;
  149. for i in msgs.iter() {
  150. if !i.message.is_empty() {
  151. write!(
  152. stdout,
  153. "{}{}{}received msg: {}",
  154. termion::clear::All,
  155. termion::style::Bold,
  156. termion::cursor::Goto(1, 2),
  157. i.message
  158. )?;
  159. } else {
  160. write!(
  161. stdout,
  162. "{}{}{}inbox is empty",
  163. termion::clear::All,
  164. termion::style::Bold,
  165. termion::cursor::Goto(1, 2),
  166. )?;
  167. }
  168. }
  169. stdout.flush()?;
  170. }
  171. DisplayMode::MessageSent => {
  172. write!(
  173. stdout,
  174. "{}{}{}message sent! {} esc: return to main menu {}",
  175. termion::clear::All,
  176. termion::style::Bold,
  177. termion::cursor::Goto(1, 2),
  178. termion::cursor::Goto(1, 3),
  179. termion::cursor::Goto(1, 4),
  180. )?;
  181. stdout.flush()?;
  182. }
  183. DisplayMode::SendFailed(e) => {
  184. write!(
  185. stdout,
  186. "{}{}{}send message failed! reason: {} {} esc: return to main menu {}",
  187. termion::clear::All,
  188. termion::style::Bold,
  189. termion::cursor::Goto(1, 2),
  190. e,
  191. termion::cursor::Goto(1, 3),
  192. termion::cursor::Goto(1, 4),
  193. )?;
  194. stdout.flush()?;
  195. }
  196. }
  197. Ok(())
  198. }
  199. async fn register_protocol(&self, msgs: DchatmsgsBuffer) -> Result<()> {
  200. debug!(target: "dchat", "Dchat::register_protocol() [START]");
  201. let registry = self.p2p.protocol_registry();
  202. registry
  203. .register(net::SESSION_ALL, move |channel, _p2p| {
  204. let msgs2 = msgs.clone();
  205. async move { ProtocolDchat::init(channel, msgs2).await }
  206. })
  207. .await;
  208. debug!(target: "dchat", "Dchat::register_protocol() [STOP]");
  209. Ok(())
  210. }
  211. async fn start(&self, ex: Arc<Executor<'_>>) -> Result<()> {
  212. debug!(target: "dchat", "Dchat::start() [START]");
  213. let ex2 = ex.clone();
  214. self.register_protocol(self.recv_msgs.clone()).await?;
  215. self.p2p.clone().start(ex.clone()).await?;
  216. ex2.spawn(self.p2p.clone().run(ex.clone())).detach();
  217. debug!(target: "dchat", "Dchat::start() [STOP]");
  218. Ok(())
  219. }
  220. async fn send(&self) -> Result<()> {
  221. let message = self.input.clone();
  222. let dchatmsg = Dchatmsg { message };
  223. self.p2p.broadcast(dchatmsg).await?;
  224. Ok(())
  225. }
  226. }
  227. // inbound
  228. fn alice() -> Result<Settings> {
  229. let seed = Url::parse("tcp://127.0.0.1:55555").unwrap();
  230. let inbound = Url::parse("tcp://127.0.0.1:55554").unwrap();
  231. let ext_addr = Url::parse("tcp://127.0.0.1:55554").unwrap();
  232. let settings = Settings {
  233. inbound: Some(inbound),
  234. outbound_connections: 0,
  235. manual_attempt_limit: 0,
  236. seed_query_timeout_seconds: 8,
  237. connect_timeout_seconds: 10,
  238. channel_handshake_seconds: 4,
  239. channel_heartbeat_seconds: 10,
  240. outbound_retry_seconds: 1200,
  241. external_addr: Some(ext_addr),
  242. peers: Vec::new(),
  243. seeds: vec![seed],
  244. node_id: String::new(),
  245. };
  246. Ok(settings)
  247. }
  248. // outbound
  249. fn bob() -> Result<Settings> {
  250. let seed = Url::parse("tcp://127.0.0.1:55555").unwrap();
  251. let oc = 5;
  252. let settings = Settings {
  253. inbound: None,
  254. outbound_connections: oc,
  255. manual_attempt_limit: 0,
  256. seed_query_timeout_seconds: 8,
  257. connect_timeout_seconds: 10,
  258. channel_handshake_seconds: 4,
  259. channel_heartbeat_seconds: 10,
  260. outbound_retry_seconds: 1200,
  261. external_addr: None,
  262. peers: Vec::new(),
  263. seeds: vec![seed],
  264. node_id: String::new(),
  265. };
  266. Ok(settings)
  267. }
  268. #[async_std::main]
  269. async fn main() -> Result<()> {
  270. let log_level = get_log_level(1);
  271. let log_config = get_log_config();
  272. let log_path = "/tmp/dchat.log";
  273. let file = File::create(log_path).unwrap();
  274. WriteLogger::init(log_level, log_config, file)?;
  275. // TODO:: proper error handling
  276. let settings: Result<Settings> = match std::env::args().nth(1) {
  277. Some(id) => match id.as_str() {
  278. "a" => alice(),
  279. "b" => bob(),
  280. _ => {
  281. println!("you must specify either a or b");
  282. Err(MissingSpecifier.into())
  283. }
  284. },
  285. None => {
  286. println!("you must specify either a or b");
  287. Err(MissingSpecifier.into())
  288. }
  289. };
  290. let p2p = net::P2p::new(settings?.into()).await;
  291. let nthreads = num_cpus::get();
  292. let (signal, shutdown) = async_channel::unbounded::<()>();
  293. let ex = Arc::new(Executor::new());
  294. let ex2 = ex.clone();
  295. let msgs: DchatmsgsBuffer = Arc::new(Mutex::new(vec![Dchatmsg { message: String::new() }]));
  296. let mut dchat = Dchat::new(p2p, msgs, String::new(), DisplayMode::Normal);
  297. let (_, result) = Parallel::new()
  298. .each(0..nthreads, |_| smol::future::block_on(ex.run(shutdown.recv())))
  299. .finish(|| {
  300. smol::future::block_on(async move {
  301. dchat.start(ex2).await?;
  302. dchat.menu().await?;
  303. drop(signal);
  304. Ok(())
  305. })
  306. });
  307. result
  308. }