main.rs 7.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253
  1. use darkfi::{
  2. error::{Error, Result},
  3. rpc::{jsonrpc, jsonrpc::JsonResult},
  4. util::{
  5. async_util,
  6. cli::{log_config, spawn_config, Config},
  7. join_config_path,
  8. },
  9. };
  10. use async_std::sync::Arc;
  11. use easy_parallel::Parallel;
  12. use log::{debug, info};
  13. use serde::{Deserialize, Serialize};
  14. use serde_json::{json, Value};
  15. use simplelog::*;
  16. use smol::Executor;
  17. use std::{
  18. collections::{HashMap, HashSet},
  19. fs::File,
  20. io,
  21. io::Read,
  22. path::PathBuf,
  23. };
  24. use termion::{async_stdin, event::Key, input::TermRead, raw::IntoRawMode};
  25. use tui::{
  26. backend::{Backend, TermionBackend},
  27. Terminal,
  28. };
  29. use url::Url;
  30. use dnetview::{
  31. model::{Connection, IdList, InfoList, NodeInfo},
  32. options::ProgramOptions,
  33. ui,
  34. view::{IdListView, InfoListView},
  35. Model, View,
  36. };
  37. const CONFIG_FILE_CONTENTS: &[u8] = include_bytes!("../dnetview_config.toml");
  38. #[derive(Clone, Serialize, Deserialize, Debug)]
  39. pub struct DnvConfig {
  40. pub nodes: Vec<IrcNode>,
  41. }
  42. #[derive(Clone, Serialize, Deserialize, Debug)]
  43. pub struct IrcNode {
  44. pub node_id: String,
  45. //pub rpc_url: String,
  46. }
  47. struct Map {
  48. url: Url,
  49. }
  50. impl Map {
  51. pub fn new(url: Url) -> Self {
  52. Self { url }
  53. }
  54. async fn request(&self, r: jsonrpc::JsonRequest) -> Result<Value> {
  55. let reply: JsonResult = match jsonrpc::send_request(&self.url, json!(r), None).await {
  56. Ok(v) => v,
  57. Err(e) => return Err(e),
  58. };
  59. match reply {
  60. JsonResult::Resp(r) => {
  61. debug!(target: "RPC", "<-- {}", serde_json::to_string(&r)?);
  62. Ok(r.result)
  63. }
  64. JsonResult::Err(e) => {
  65. debug!(target: "RPC", "<-- {}", serde_json::to_string(&e)?);
  66. Err(Error::JsonRpcError(e.error.message.to_string()))
  67. }
  68. JsonResult::Notif(n) => {
  69. debug!(target: "RPC", "<-- {}", serde_json::to_string(&n)?);
  70. Err(Error::JsonRpcError("Unexpected reply".to_string()))
  71. }
  72. }
  73. }
  74. // --> {"jsonrpc": "2.0", "method": "ping", "params": [], "id": 42}
  75. // <-- {"jsonrpc": "2.0", "result": "pong", "id": 42}
  76. async fn _ping(&self) -> Result<Value> {
  77. let req = jsonrpc::request(json!("ping"), json!([]));
  78. Ok(self.request(req).await?)
  79. }
  80. //--> {"jsonrpc": "2.0", "method": "poll", "params": [], "id": 42}
  81. // <-- {"jsonrpc": "2.0", "result": {"nodeID": [], "nodeinfo" [], "id": 42}
  82. async fn get_info(&self) -> Result<Value> {
  83. let req = jsonrpc::request(json!("get_info"), json!([]));
  84. Ok(self.request(req).await?)
  85. }
  86. }
  87. #[async_std::main]
  88. async fn main() -> Result<()> {
  89. let options = ProgramOptions::load()?;
  90. let verbosity_level = options.app.occurrences_of("verbose");
  91. let (lvl, cfg) = log_config(verbosity_level)?;
  92. let file = File::create(&*options.log_path).unwrap();
  93. WriteLogger::init(lvl, cfg, file)?;
  94. info!("Log level: {}", lvl);
  95. let config_path = join_config_path(&PathBuf::from("dnetview_config.toml"))?;
  96. spawn_config(&config_path, CONFIG_FILE_CONTENTS)?;
  97. let config = Config::<DnvConfig>::load(config_path)?;
  98. let stdout = io::stdout().into_raw_mode()?;
  99. let backend = TermionBackend::new(stdout);
  100. let mut terminal = Terminal::new(backend)?;
  101. terminal.clear()?;
  102. let info_list = InfoList::new();
  103. let ids = HashSet::new();
  104. let id_list = IdList::new(ids);
  105. let model = Arc::new(Model::new(id_list, info_list));
  106. let nthreads = num_cpus::get();
  107. let (signal, shutdown) = async_channel::unbounded::<()>();
  108. let ex = Arc::new(Executor::new());
  109. let ex2 = ex.clone();
  110. let (_, result) = Parallel::new()
  111. .each(0..nthreads, |_| smol::future::block_on(ex.run(shutdown.recv())))
  112. .finish(|| {
  113. smol::future::block_on(async move {
  114. run_rpc(&config, ex2.clone(), model.clone()).await?;
  115. render(&mut terminal, model.clone(), config).await?;
  116. drop(signal);
  117. Ok::<(), darkfi::Error>(())
  118. })
  119. });
  120. result
  121. }
  122. async fn run_rpc(config: &DnvConfig, ex: Arc<Executor<'_>>, model: Arc<Model>) -> Result<()> {
  123. for node in config.nodes.clone() {
  124. let client = Map::new(Url::parse(&node.node_id)?);
  125. ex.spawn(poll(client, model.clone())).detach();
  126. }
  127. Ok(())
  128. }
  129. async fn poll(client: Map, model: Arc<Model>) -> Result<()> {
  130. debug!("Attemping to poll: {}", client.url);
  131. loop {
  132. let reply = client.get_info().await?;
  133. debug!("{:?}", reply);
  134. if reply.as_object().is_some() && !reply.as_object().unwrap().is_empty() {
  135. let external_addr = reply.as_object().unwrap().get("external_addr");
  136. let session_inbound = reply.as_object().unwrap().get("session_inbound");
  137. let si_key =
  138. session_inbound.unwrap().as_object().unwrap().get("key").unwrap().as_u64().unwrap();
  139. let session_manual = reply.as_object().unwrap().get("session_manual");
  140. let sm_key =
  141. session_manual.unwrap().as_object().unwrap().get("key").unwrap().as_u64().unwrap();
  142. let session_outbound = reply.as_object().unwrap().get("session_outbound");
  143. let so_key =
  144. session_manual.unwrap().as_object().unwrap().get("key").unwrap().as_u64().unwrap();
  145. let channel_state = reply.as_object().unwrap().get("state").unwrap().as_str().unwrap();
  146. let slots = reply.as_object().unwrap().get("slots");
  147. let session_in = Connection::new(si_key.to_string(), channel_state.to_string());
  148. let session_man = Connection::new(sm_key.to_string(), channel_state.to_string());
  149. let session_out = Connection::new(so_key.to_string(), channel_state.to_string());
  150. let mut outconnects = Vec::new();
  151. let mut inconnects = Vec::new();
  152. let mut manconnects = Vec::new();
  153. outconnects.push(session_out);
  154. inconnects.push(session_in);
  155. manconnects.push(session_man);
  156. let infos =
  157. NodeInfo { outbound: outconnects, manual: manconnects, inbound: inconnects };
  158. let mut node_info = HashMap::new();
  159. // TODO: fix this. this key should be global identifier for each connection
  160. // right now we are using si_key, which is the inbound session key.
  161. node_info.insert(si_key, infos);
  162. for (si_key, value) in node_info.clone() {
  163. model.id_list.node_id.lock().await.insert(si_key.to_string().clone());
  164. model.info_list.infos.lock().await.insert(si_key.to_string(), value);
  165. }
  166. } else {
  167. // TODO: error handling
  168. debug!("Reply is empty");
  169. }
  170. async_util::sleep(2).await;
  171. }
  172. }
  173. async fn render<B: Backend>(
  174. terminal: &mut Terminal<B>,
  175. model: Arc<Model>,
  176. config: DnvConfig,
  177. ) -> io::Result<()> {
  178. let mut asi = async_stdin();
  179. terminal.clear()?;
  180. let id_list = IdListView::new(HashSet::new());
  181. let info_list = InfoListView::new(HashMap::new());
  182. let mut view = View::new(id_list.clone(), info_list.clone());
  183. view.id_list.state.select(Some(0));
  184. view.info_list.index = 0;
  185. loop {
  186. view.update(model.info_list.infos.lock().await.clone());
  187. terminal.draw(|f| {
  188. ui::ui(f, view.clone());
  189. })?;
  190. for k in asi.by_ref().keys() {
  191. match k.unwrap() {
  192. Key::Char('q') => {
  193. terminal.clear()?;
  194. return Ok(());
  195. }
  196. Key::Char('j') => {
  197. view.id_list.next();
  198. view.info_list.next().await;
  199. }
  200. Key::Char('k') => {
  201. view.id_list.previous();
  202. view.info_list.previous().await;
  203. }
  204. _ => (),
  205. }
  206. }
  207. }
  208. }