main.rs 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293
  1. use async_std::sync::Arc;
  2. use std::{fs::File, io, io::Read, path::PathBuf};
  3. use easy_parallel::Parallel;
  4. use fxhash::{FxHashMap, FxHashSet};
  5. use log::{debug, info};
  6. use serde_json::{json, Value};
  7. use simplelog::*;
  8. use smol::Executor;
  9. use termion::{async_stdin, event::Key, input::TermRead, raw::IntoRawMode};
  10. use tui::{
  11. backend::{Backend, TermionBackend},
  12. Terminal,
  13. };
  14. use url::Url;
  15. use darkfi::{
  16. error::{Error, Result},
  17. rpc::{jsonrpc, jsonrpc::JsonResult},
  18. util::{
  19. async_util,
  20. cli::{log_config, spawn_config, Config},
  21. join_config_path,
  22. },
  23. };
  24. use dnetview::{
  25. config::{DnvConfig, CONFIG_FILE_CONTENTS},
  26. model::{Channel, IdList, InboundInfo, InfoList, ManualInfo, NodeInfo, OutboundInfo, Slot},
  27. options::ProgramOptions,
  28. ui,
  29. view::{IdListView, InfoListView},
  30. Model, View,
  31. };
  32. struct DNetView {
  33. url: Url,
  34. name: String,
  35. }
  36. impl DNetView {
  37. pub fn new(url: Url, name: String) -> Self {
  38. Self { url, name }
  39. }
  40. async fn request(&self, r: jsonrpc::JsonRequest) -> Result<Value> {
  41. let reply: JsonResult = match jsonrpc::send_request(&self.url, json!(r), None).await {
  42. Ok(v) => v,
  43. Err(e) => return Err(e),
  44. };
  45. match reply {
  46. JsonResult::Resp(r) => {
  47. debug!(target: "RPC", "<-- {}", serde_json::to_string(&r)?);
  48. Ok(r.result)
  49. }
  50. JsonResult::Err(e) => {
  51. debug!(target: "RPC", "<-- {}", serde_json::to_string(&e)?);
  52. Err(Error::JsonRpcError(e.error.message.to_string()))
  53. }
  54. JsonResult::Notif(n) => {
  55. debug!(target: "RPC", "<-- {}", serde_json::to_string(&n)?);
  56. Err(Error::JsonRpcError("Unexpected reply".to_string()))
  57. }
  58. }
  59. }
  60. // --> {"jsonrpc": "2.0", "method": "ping", "params": [], "id": 42}
  61. // <-- {"jsonrpc": "2.0", "result": "pong", "id": 42}
  62. async fn _ping(&self) -> Result<Value> {
  63. let req = jsonrpc::request(json!("ping"), json!([]));
  64. Ok(self.request(req).await?)
  65. }
  66. //--> {"jsonrpc": "2.0", "method": "poll", "params": [], "id": 42}
  67. // <-- {"jsonrpc": "2.0", "result": {"nodeID": [], "nodeinfo" [], "id": 42}
  68. async fn get_info(&self) -> Result<Value> {
  69. let req = jsonrpc::request(json!("get_info"), json!([]));
  70. Ok(self.request(req).await?)
  71. }
  72. }
  73. #[async_std::main]
  74. async fn main() -> Result<()> {
  75. let options = ProgramOptions::load()?;
  76. let verbosity_level = options.app.occurrences_of("verbose");
  77. let (lvl, cfg) = log_config(verbosity_level)?;
  78. let file = File::create(&*options.log_path).unwrap();
  79. WriteLogger::init(lvl, cfg, file)?;
  80. info!("Log level: {}", lvl);
  81. let config_path = join_config_path(&PathBuf::from("dnetview_config.toml"))?;
  82. spawn_config(&config_path, CONFIG_FILE_CONTENTS)?;
  83. let config = Config::<DnvConfig>::load(config_path)?;
  84. let stdout = io::stdout().into_raw_mode()?;
  85. let backend = TermionBackend::new(stdout);
  86. let mut terminal = Terminal::new(backend)?;
  87. terminal.clear()?;
  88. let info_list = InfoList::new();
  89. let ids = FxHashSet::default();
  90. let id_list = IdList::new(ids);
  91. let model = Arc::new(Model::new(id_list, info_list));
  92. let nthreads = num_cpus::get();
  93. let (signal, shutdown) = async_channel::unbounded::<()>();
  94. let ex = Arc::new(Executor::new());
  95. let ex2 = ex.clone();
  96. let (_, result) = Parallel::new()
  97. .each(0..nthreads, |_| smol::future::block_on(ex.run(shutdown.recv())))
  98. .finish(|| {
  99. smol::future::block_on(async move {
  100. run_rpc(&config, ex2.clone(), model.clone()).await?;
  101. render(&mut terminal, model.clone()).await?;
  102. drop(signal);
  103. Ok::<(), darkfi::Error>(())
  104. })
  105. });
  106. result
  107. }
  108. async fn run_rpc(config: &DnvConfig, ex: Arc<Executor<'_>>, model: Arc<Model>) -> Result<()> {
  109. for node in config.nodes.clone() {
  110. let client = DNetView::new(Url::parse(&node.rpc_url)?, node.name);
  111. ex.spawn(poll(client, model.clone())).detach();
  112. }
  113. Ok(())
  114. }
  115. // TODO: clean up into seperate functions.
  116. // TODO: replace if/else with match where possible
  117. // TODO: test unwraps will never ever crash
  118. async fn poll(client: DNetView, model: Arc<Model>) -> Result<()> {
  119. loop {
  120. let reply = client.get_info().await?;
  121. if reply.as_object().is_some() && !reply.as_object().unwrap().is_empty() {
  122. // TODO: we are ignoring this value for now
  123. let _ext_addr = reply.as_object().unwrap().get("external_addr");
  124. let inbound_obj = &reply.as_object().unwrap()["session_inbound"];
  125. let manual_obj = &reply.as_object().unwrap()["session_manual"];
  126. let outbound_obj = &reply.as_object().unwrap()["session_outbound"];
  127. let mut iconnects = Vec::new();
  128. let mut mconnects = Vec::new();
  129. let mut oconnects = Vec::new();
  130. let mut slots = Vec::new();
  131. let mut addrs = Vec::new();
  132. let mut msgs = Vec::new();
  133. // parse inbound connection data
  134. let i_connected = &inbound_obj["connected"];
  135. if i_connected.as_object().unwrap().is_empty() {
  136. // channel is empty. initialize with empty values
  137. let connected = "Empty".to_string();
  138. let msg = "Null".to_string();
  139. let status = "Null".to_string();
  140. let channel = Channel::new(msg, status);
  141. let is_empty = true;
  142. let iinfo = InboundInfo::new(is_empty, connected, channel);
  143. iconnects.push(iinfo);
  144. } else {
  145. // channel is not empty. initialize with whole values
  146. let ic = i_connected.as_object().unwrap();
  147. for k in ic.keys() {
  148. let node = ic.get(k);
  149. let addr = k.to_string();
  150. let msg = node.unwrap().get("last_msg").unwrap().as_str().unwrap().to_string();
  151. let status =
  152. node.unwrap().get("last_status").unwrap().as_str().unwrap().to_string();
  153. let channel = Channel::new(msg.clone(), status);
  154. let is_empty = false;
  155. let iinfo = InboundInfo::new(is_empty, addr.clone(), channel);
  156. iconnects.push(iinfo);
  157. addrs.push(addr);
  158. msgs.push(msg.clone());
  159. }
  160. }
  161. // parse manual connection data
  162. let minfo: ManualInfo = serde_json::from_value(manual_obj.clone())?;
  163. mconnects.push(minfo);
  164. // parse outbound connection data
  165. let outbound_slots = &outbound_obj["slots"];
  166. for slot in outbound_slots.as_array().unwrap() {
  167. if slot["channel"].is_null() {
  168. // channel is empty. initialize with empty values
  169. let is_empty = true;
  170. let state = &slot["state"];
  171. let msg = "Null".to_string();
  172. let status = "Null".to_string();
  173. let channel = Channel::new(msg, status);
  174. let new_slot = Slot::new(
  175. is_empty,
  176. String::new(),
  177. channel,
  178. state.as_str().unwrap().to_string(),
  179. );
  180. slots.push(new_slot.clone())
  181. } else {
  182. // channel is not empty. initialize with whole values
  183. let is_empty = false;
  184. let addr = &slot["addr"];
  185. let state = &slot["state"];
  186. let channel: Channel = serde_json::from_value(slot["channel"].clone())?;
  187. let new_slot = Slot::new(
  188. is_empty,
  189. addr.as_str().unwrap().to_string(),
  190. channel.clone(),
  191. state.as_str().unwrap().to_string(),
  192. );
  193. slots.push(new_slot);
  194. addrs.push(addr.as_str().unwrap().to_string());
  195. msgs.push(channel.last_msg.clone());
  196. }
  197. }
  198. // create node_info
  199. let is_empty = is_empty_outbound(slots.clone());
  200. let oinfo = OutboundInfo::new(is_empty, slots.clone());
  201. oconnects.push(oinfo);
  202. let infos = NodeInfo { outbound: oconnects, manual: mconnects, inbound: iconnects };
  203. let mut node_info = FxHashMap::default();
  204. let node_name = &client.name.as_str();
  205. node_info.insert(&node_name, infos.clone());
  206. // insert into model
  207. for (key, value) in node_info.clone() {
  208. model.id_list.node_id.lock().await.insert(key.to_string().clone());
  209. model.info_list.infos.lock().await.insert(key.to_string(), value);
  210. }
  211. } else {
  212. // TODO: error handling
  213. //debug!("Reply is empty");
  214. }
  215. async_util::sleep(2).await;
  216. }
  217. }
  218. fn is_empty_outbound(slots: Vec<Slot>) -> bool {
  219. return slots.iter().all(|slot| slot.is_empty)
  220. }
  221. async fn render<B: Backend>(terminal: &mut Terminal<B>, model: Arc<Model>) -> io::Result<()> {
  222. let mut asi = async_stdin();
  223. terminal.clear()?;
  224. let id_list = IdListView::new(FxHashSet::default());
  225. let info_list = InfoListView::new(FxHashMap::default());
  226. let mut view = View::new(id_list.clone(), info_list.clone());
  227. view.id_list.state.select(Some(0));
  228. view.info_list.index = 0;
  229. loop {
  230. view.update(model.info_list.infos.lock().await.clone());
  231. terminal.draw(|f| {
  232. ui::ui(f, view.clone());
  233. })?;
  234. for k in asi.by_ref().keys() {
  235. match k.unwrap() {
  236. Key::Char('q') => {
  237. terminal.clear()?;
  238. return Ok(())
  239. }
  240. Key::Char('j') => {
  241. view.id_list.next();
  242. }
  243. Key::Char('k') => {
  244. view.id_list.previous();
  245. }
  246. _ => (),
  247. }
  248. }
  249. }
  250. }