main.rs 8.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296
  1. use async_std::sync::{Arc, Mutex};
  2. use std::{
  3. env,
  4. fs::{create_dir_all, remove_dir_all},
  5. io::stdin,
  6. path::Path,
  7. };
  8. use async_executor::Executor;
  9. use crypto_box::{
  10. aead::{Aead, AeadCore},
  11. SalsaBox, SecretKey,
  12. };
  13. use futures::{select, FutureExt};
  14. use fxhash::FxHashMap;
  15. use log::{debug, error, info, warn};
  16. use smol::future;
  17. use structopt_toml::StructOptToml;
  18. use darkfi::{
  19. async_daemonize, net,
  20. raft::{NetMsg, ProtocolRaft, Raft, RaftSettings},
  21. rpc::server::listen_and_serve,
  22. util::{
  23. cli::{get_log_config, get_log_level, spawn_config},
  24. expand_path,
  25. path::get_config_path,
  26. serial::{deserialize, serialize, SerialDecodable, SerialEncodable},
  27. },
  28. Error, Result,
  29. };
  30. mod error;
  31. mod jsonrpc;
  32. mod month_tasks;
  33. mod settings;
  34. mod task_info;
  35. mod util;
  36. use crate::{
  37. error::TaudResult,
  38. jsonrpc::JsonRpcInterface,
  39. settings::{Args, CONFIG_FILE, CONFIG_FILE_CONTENTS},
  40. task_info::TaskInfo,
  41. };
  42. fn get_workspaces(settings: &Args) -> Result<FxHashMap<String, SalsaBox>> {
  43. let mut workspaces = FxHashMap::default();
  44. for workspace in settings.workspaces.iter() {
  45. let workspace: Vec<&str> = workspace.split(':').collect();
  46. let (workspace, secret) = (workspace[0], workspace[1]);
  47. let bytes: [u8; 32] = bs58::decode(secret)
  48. .into_vec()?
  49. .try_into()
  50. .map_err(|_| Error::ParseFailed("Parse secret key failed"))?;
  51. let secret = crypto_box::SecretKey::from(bytes);
  52. let public = secret.public_key();
  53. let salsa_box = crypto_box::SalsaBox::new(&public, &secret);
  54. workspaces.insert(workspace.to_string(), salsa_box);
  55. }
  56. Ok(workspaces)
  57. }
  58. #[derive(Debug, Clone, SerialEncodable, SerialDecodable)]
  59. pub struct EncryptedTask {
  60. nonce: Vec<u8>,
  61. payload: Vec<u8>,
  62. }
  63. fn encrypt_task(
  64. task: &TaskInfo,
  65. salsa_box: &SalsaBox,
  66. rng: &mut crypto_box::rand_core::OsRng,
  67. ) -> TaudResult<EncryptedTask> {
  68. debug!("start encrypting task");
  69. let nonce = SalsaBox::generate_nonce(rng);
  70. let payload = &serialize(task)[..];
  71. let payload = salsa_box.encrypt(&nonce, payload)?;
  72. let nonce = nonce.to_vec();
  73. Ok(EncryptedTask { nonce, payload })
  74. }
  75. fn decrypt_task(encrypt_task: &EncryptedTask, salsa_box: &SalsaBox) -> TaudResult<TaskInfo> {
  76. debug!("start decrypting task");
  77. let nonce = encrypt_task.nonce.as_slice();
  78. let decrypted_task = salsa_box.decrypt(nonce.into(), &encrypt_task.payload[..])?;
  79. let task = deserialize(&decrypted_task)?;
  80. Ok(task)
  81. }
  82. async fn start_sync_loop(
  83. broadcast_rcv: async_channel::Receiver<TaskInfo>,
  84. raft_msgs_sender: async_channel::Sender<EncryptedTask>,
  85. commits_recv: async_channel::Receiver<EncryptedTask>,
  86. datastore_path: std::path::PathBuf,
  87. workspaces: FxHashMap<String, SalsaBox>,
  88. mut rng: crypto_box::rand_core::OsRng,
  89. ) -> TaudResult<()> {
  90. loop {
  91. select! {
  92. task = broadcast_rcv.recv().fuse() => {
  93. let tk = task.map_err(Error::from)?;
  94. if workspaces.contains_key(&tk.workspace) {
  95. let salsa_box = workspaces.get(&tk.workspace).unwrap();
  96. let encrypted_task = encrypt_task(&tk, salsa_box, &mut rng)?;
  97. info!(target: "tau", "Send the task: ref: {}", tk.ref_id);
  98. raft_msgs_sender.send(encrypted_task).await.map_err(Error::from)?;
  99. }
  100. }
  101. task = commits_recv.recv().fuse() => {
  102. let task = task.map_err(Error::from)?;
  103. on_receive_task(&task,&datastore_path, &workspaces)
  104. .await?;
  105. }
  106. }
  107. }
  108. }
  109. async fn on_receive_task(
  110. task: &EncryptedTask,
  111. datastore_path: &Path,
  112. workspaces: &FxHashMap<String, SalsaBox>,
  113. ) -> TaudResult<()> {
  114. for (workspace, salsa_box) in workspaces.iter() {
  115. let task = decrypt_task(&task, &salsa_box);
  116. if let Err(e) = task {
  117. info!("unable to decrypt the task: {}", e);
  118. continue
  119. }
  120. let mut task = task.unwrap();
  121. info!(target: "tau", "Save the task: ref: {}", task.ref_id);
  122. task.workspace = workspace.clone();
  123. task.save(&datastore_path)?;
  124. }
  125. Ok(())
  126. }
  127. async_daemonize!(realmain);
  128. async fn realmain(settings: Args, executor: Arc<Executor<'_>>) -> Result<()> {
  129. let datastore_path = expand_path(&settings.datastore)?;
  130. let nickname =
  131. if settings.nickname.is_some() { settings.nickname.clone() } else { env::var("USER").ok() };
  132. if settings.refresh {
  133. println!("Removing local data in: {:?} (yes/no)? ", datastore_path);
  134. let mut confirm = String::new();
  135. stdin().read_line(&mut confirm).expect("Failed to read line");
  136. let confirm = confirm.to_lowercase();
  137. let confirm = confirm.trim();
  138. if confirm == "yes" || confirm == "y" {
  139. remove_dir_all(datastore_path).unwrap_or(());
  140. println!("Local data removed successfully.");
  141. } else {
  142. error!("Unexpected Value: {}", confirm);
  143. }
  144. return Ok(())
  145. }
  146. if nickname.is_none() {
  147. error!("Provide a nickname in config file");
  148. return Ok(())
  149. }
  150. // mkdir datastore_path if not exists
  151. create_dir_all(datastore_path.clone())?;
  152. create_dir_all(datastore_path.join("month"))?;
  153. create_dir_all(datastore_path.join("task"))?;
  154. let rng = crypto_box::rand_core::OsRng;
  155. if settings.generate {
  156. println!("Generating a new workspace");
  157. loop {
  158. println!("Name for the new workspace: ");
  159. let mut workspace = String::new();
  160. stdin().read_line(&mut workspace).ok().expect("Failed to read line");
  161. let workspace = workspace.to_lowercase();
  162. let workspace = workspace.trim();
  163. if workspace.is_empty() && workspace.len() < 3 {
  164. error!("Wrong workspace try again");
  165. continue
  166. }
  167. let mut rng = crypto_box::rand_core::OsRng;
  168. let secret_key = SecretKey::generate(&mut rng);
  169. let encoded = bs58::encode(secret_key.as_bytes());
  170. println!("workspace: {}:{}", workspace, encoded.into_string());
  171. println!("Please add it to the config file.");
  172. break
  173. }
  174. return Ok(())
  175. }
  176. let workspaces = get_workspaces(&settings)?;
  177. if workspaces.is_empty() {
  178. error!("Please add at least one workspace to the config file.");
  179. println!("Run `$ taud --generate` to generate new workspace.");
  180. return Ok(())
  181. }
  182. //
  183. // Raft
  184. //
  185. let seen_net_msgs = Arc::new(Mutex::new(FxHashMap::default()));
  186. let datastore_raft = datastore_path.join("tau.db");
  187. let raft_settings = RaftSettings { datastore_path: datastore_raft, ..RaftSettings::default() };
  188. let mut raft = Raft::<EncryptedTask>::new(raft_settings, seen_net_msgs.clone())?;
  189. let raft_id = raft.id();
  190. let (broadcast_snd, broadcast_rcv) = async_channel::unbounded::<TaskInfo>();
  191. //
  192. // P2p setup
  193. //
  194. let mut net_settings = settings.net.clone();
  195. net_settings.app_version = Some(option_env!("CARGO_PKG_VERSION").unwrap_or("").to_string());
  196. let (p2p_send_channel, p2p_recv_channel) = async_channel::unbounded::<NetMsg>();
  197. let p2p = net::P2p::new(net_settings.into()).await;
  198. let p2p = p2p.clone();
  199. let registry = p2p.protocol_registry();
  200. registry
  201. .register(net::SESSION_ALL, move |channel, p2p| {
  202. let raft_id = raft_id.clone();
  203. let sender = p2p_send_channel.clone();
  204. let seen_net_msgs_cloned = seen_net_msgs.clone();
  205. async move {
  206. ProtocolRaft::init(raft_id, channel, sender, p2p, seen_net_msgs_cloned).await
  207. }
  208. })
  209. .await;
  210. p2p.clone().start(executor.clone()).await?;
  211. executor.spawn(p2p.clone().run(executor.clone())).detach();
  212. //
  213. // RPC interface
  214. //
  215. let rpc_interface = Arc::new(JsonRpcInterface::new(
  216. datastore_path.clone(),
  217. broadcast_snd,
  218. nickname.unwrap(),
  219. workspaces.clone(),
  220. p2p.clone(),
  221. ));
  222. executor.spawn(listen_and_serve(settings.rpc_listen.clone(), rpc_interface)).detach();
  223. //
  224. // Waiting Exit signal
  225. //
  226. let (signal, shutdown) = async_channel::bounded::<()>(1);
  227. ctrlc::set_handler(move || {
  228. warn!(target: "tau", "Catch exit signal");
  229. // cleaning up tasks running in the background
  230. if let Err(e) = async_std::task::block_on(signal.send(())) {
  231. error!("Error on sending exit signal: {}", e);
  232. }
  233. })
  234. .unwrap();
  235. executor
  236. .spawn(start_sync_loop(
  237. broadcast_rcv,
  238. raft.sender(),
  239. raft.receiver(),
  240. datastore_path,
  241. workspaces,
  242. rng,
  243. ))
  244. .detach();
  245. raft.run(p2p.clone(), p2p_recv_channel.clone(), executor.clone(), shutdown.clone()).await?;
  246. Ok(())
  247. }