main.rs 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2023 Dyne.org foundation
  4. *
  5. * This program is free software: you can redistribute it and/or modify
  6. * it under the terms of the GNU Affero General Public License as
  7. * published by the Free Software Foundation, either version 3 of the
  8. * License, or (at your option) any later version.
  9. *
  10. * This program is distributed in the hope that it will be useful,
  11. * but WITHOUT ANY WARRANTY; without even the implied warranty of
  12. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  13. * GNU Affero General Public License for more details.
  14. *
  15. * You should have received a copy of the GNU Affero General Public License
  16. * along with this program. If not, see <https://www.gnu.org/licenses/>.
  17. */
  18. use std::{
  19. collections::HashMap,
  20. env,
  21. ffi::CString,
  22. fs::{create_dir_all, remove_dir_all},
  23. io::{stdin, Write},
  24. path::Path,
  25. sync::{Arc, OnceLock},
  26. };
  27. use crypto_box::{
  28. aead::{Aead, AeadCore},
  29. ChaChaBox, SecretKey,
  30. };
  31. use darkfi_serial::{
  32. async_trait, deserialize, deserialize_async_partial, serialize, serialize_async,
  33. SerialDecodable, SerialEncodable,
  34. };
  35. use futures::{select, FutureExt};
  36. use libc::mkfifo;
  37. use log::{debug, error, info};
  38. use rand::rngs::OsRng;
  39. use smol::{fs, lock::RwLock, stream::StreamExt};
  40. use structopt_toml::StructOptToml;
  41. use tinyjson::JsonValue;
  42. use darkfi::{
  43. async_daemonize,
  44. event_graph::{
  45. proto::{EventPut, ProtocolEventGraph},
  46. Event, EventGraph, EventGraphPtr, NULL_ID,
  47. },
  48. net::{P2p, P2pPtr, SESSION_ALL},
  49. rpc::{
  50. jsonrpc::JsonSubscriber,
  51. server::{listen_and_serve, RequestHandler},
  52. },
  53. system::{sleep, StoppableTask},
  54. util::path::expand_path,
  55. Error, Result,
  56. };
  57. mod jsonrpc;
  58. mod settings;
  59. use taud::{
  60. error::{TaudError, TaudResult},
  61. task_info::{TaskEvent, TaskInfo},
  62. util::pipe_write,
  63. };
  64. use crate::{
  65. jsonrpc::JsonRpcInterface,
  66. settings::{Args, CONFIG_FILE, CONFIG_FILE_CONTENTS},
  67. };
  68. fn get_workspaces(settings: &Args) -> Result<HashMap<String, ChaChaBox>> {
  69. let mut workspaces = HashMap::new();
  70. for workspace in settings.workspaces.iter() {
  71. let workspace: Vec<&str> = workspace.split(':').collect();
  72. let (workspace, secret) = (workspace[0], workspace[1]);
  73. let bytes: [u8; 32] = bs58::decode(secret)
  74. .into_vec()?
  75. .try_into()
  76. .map_err(|_| Error::ParseFailed("Parse secret key failed"))?;
  77. let secret = crypto_box::SecretKey::from(bytes);
  78. let public = secret.public_key();
  79. let chacha_box = crypto_box::ChaChaBox::new(&public, &secret);
  80. workspaces.insert(workspace.to_string(), chacha_box);
  81. }
  82. Ok(workspaces)
  83. }
  84. #[derive(Debug, Clone, SerialEncodable, SerialDecodable)]
  85. pub struct EncryptedTask {
  86. payload: String,
  87. }
  88. fn encrypt_task(
  89. task: &TaskInfo,
  90. chacha_box: &ChaChaBox,
  91. rng: &mut OsRng,
  92. ) -> TaudResult<EncryptedTask> {
  93. debug!("start encrypting task");
  94. let nonce = ChaChaBox::generate_nonce(rng);
  95. let payload = &serialize(task)[..];
  96. let mut payload = chacha_box.encrypt(&nonce, payload)?;
  97. let mut concat = vec![];
  98. concat.append(&mut nonce.as_slice().to_vec());
  99. concat.append(&mut payload);
  100. let payload = bs58::encode(concat.clone()).into_string();
  101. Ok(EncryptedTask { payload })
  102. }
  103. fn try_decrypt_task(encrypt_task: &EncryptedTask, chacha_box: &ChaChaBox) -> TaudResult<TaskInfo> {
  104. debug!("start decrypting task");
  105. let bytes = match bs58::decode(&encrypt_task.payload).into_vec() {
  106. Ok(v) => v,
  107. Err(_) => return Err(TaudError::DecryptionError("Error decoding payload".to_string())),
  108. };
  109. if bytes.len() < 25 {
  110. return Err(TaudError::DecryptionError("Invalid bytes length".to_string()))
  111. }
  112. // Try extracting the nonce
  113. let nonce = match bytes[0..24].try_into() {
  114. Ok(v) => v,
  115. Err(_) => return Err(TaudError::DecryptionError("Invalid nonce".to_string())),
  116. };
  117. // Take the remaining ciphertext
  118. let message = &bytes[24..];
  119. // let nonce = encrypt_task.nonce.as_slice();
  120. let decrypted_task = chacha_box.decrypt(nonce, message)?;
  121. let task = deserialize(&decrypted_task)?;
  122. Ok(task)
  123. }
  124. #[allow(clippy::too_many_arguments)]
  125. async fn start_sync_loop(
  126. event_graph: EventGraphPtr,
  127. broadcast_rcv: smol::channel::Receiver<TaskInfo>,
  128. workspaces: Arc<HashMap<String, ChaChaBox>>,
  129. datastore_path: std::path::PathBuf,
  130. piped: bool,
  131. p2p: P2pPtr,
  132. last_sent: RwLock<blake3::Hash>,
  133. seen: OnceLock<sled::Tree>,
  134. ) -> TaudResult<()> {
  135. let incoming = event_graph.event_sub.clone().subscribe().await;
  136. let seen_events = seen.get().unwrap();
  137. loop {
  138. select! {
  139. task_event = broadcast_rcv.recv().fuse() => {
  140. let tk = task_event.map_err(Error::from)?;
  141. if workspaces.contains_key(&tk.workspace) {
  142. let chacha_box = workspaces.get(&tk.workspace).unwrap();
  143. let encrypted_task = encrypt_task(&tk, chacha_box, &mut OsRng)?;
  144. info!(target: "tau", "Send the task: ref: {}", tk.ref_id);
  145. // Build a DAG event and return it.
  146. let event = Event::new(
  147. serialize_async(&encrypted_task).await,
  148. event_graph.clone(),
  149. )
  150. .await;
  151. // Update the last sent event.
  152. // let event_id = event.id();
  153. // *last_sent.write().await = event_id;
  154. // If it fails for some reason, for now, we just note it
  155. // and pass.
  156. if let Err(e) = event_graph.dag_insert(event.clone()).await {
  157. error!("[IRC CLIENT] Failed inserting new event to DAG: {}", e);
  158. } else {
  159. // We sent this, so it should be considered seen.
  160. // TODO: should we save task on send or on receive?
  161. // on receive better because it's garanteed your event is out there
  162. // debug!("Marking event {} as seen", event_id);
  163. // seen.get().unwrap().insert(event_id.as_bytes(), &[]).unwrap();
  164. // Otherwise, broadcast it
  165. p2p.broadcast(&EventPut(event)).await;
  166. }
  167. }
  168. }
  169. task_event = incoming.receive().fuse() => {
  170. let event_id = task_event.id();
  171. if *last_sent.read().await == event_id {
  172. continue
  173. }
  174. if seen_events.contains_key(event_id.as_bytes()).unwrap() {
  175. continue
  176. }
  177. // Try to deserialize the `Event`'s content into a `Privmsg`
  178. let enc_task: EncryptedTask = match deserialize_async_partial(task_event.content()).await {
  179. Ok((v, _)) => v,
  180. Err(e) => {
  181. error!("[TAUD] Failed deserializing incoming EncryptedTask event: {}", e);
  182. continue
  183. }
  184. };
  185. on_receive_task(&enc_task, &datastore_path, &workspaces, piped)
  186. .await?;
  187. }
  188. }
  189. }
  190. }
  191. async fn on_receive_task(
  192. task: &EncryptedTask,
  193. datastore_path: &Path,
  194. workspaces: &HashMap<String, ChaChaBox>,
  195. piped: bool,
  196. ) -> TaudResult<()> {
  197. for (workspace, chacha_box) in workspaces.iter() {
  198. let task = try_decrypt_task(task, chacha_box);
  199. if let Err(e) = task {
  200. debug!("unable to decrypt the task: {}", e);
  201. continue
  202. }
  203. let mut task = task.unwrap();
  204. info!(target: "tau", "Save the task: ref: {}", task.ref_id);
  205. task.workspace = workspace.clone();
  206. if piped {
  207. // if we can't load the task then it's a new task.
  208. // otherwise it's a modification.
  209. match TaskInfo::load(&task.ref_id, datastore_path) {
  210. Ok(loaded_task) => {
  211. let loaded_events = loaded_task.events;
  212. let mut events = task.events.clone();
  213. events.retain(|ev| !loaded_events.contains(ev));
  214. let file = "/tmp/tau_pipe";
  215. let mut pipe_write = pipe_write(file)?;
  216. let mut task_clone = task.clone();
  217. task_clone.events = events;
  218. let json: JsonValue = (&task_clone).into();
  219. pipe_write.write_all(json.stringify().unwrap().as_bytes())?;
  220. }
  221. Err(_) => {
  222. let file = "/tmp/tau_pipe";
  223. let mut pipe_write = pipe_write(file)?;
  224. let mut task_clone = task.clone();
  225. task_clone.events.push(TaskEvent::new(
  226. "add_task".to_string(),
  227. task_clone.owner.clone(),
  228. "".to_string(),
  229. ));
  230. let json: JsonValue = (&task_clone).into();
  231. pipe_write.write_all(json.stringify().unwrap().as_bytes())?;
  232. }
  233. }
  234. }
  235. task.save(datastore_path)?;
  236. }
  237. Ok(())
  238. }
  239. async_daemonize!(realmain);
  240. async fn realmain(settings: Args, executor: Arc<smol::Executor<'static>>) -> Result<()> {
  241. let datastore_path = expand_path(&settings.datastore)?;
  242. let nickname =
  243. if settings.nickname.is_some() { settings.nickname.clone() } else { env::var("USER").ok() };
  244. if settings.refresh {
  245. println!("Removing local data in: {:?} (yes/no)? ", datastore_path);
  246. let mut confirm = String::new();
  247. stdin().read_line(&mut confirm).expect("Failed to read line");
  248. let confirm = confirm.to_lowercase();
  249. let confirm = confirm.trim();
  250. if confirm == "yes" || confirm == "y" {
  251. remove_dir_all(datastore_path).unwrap_or(());
  252. println!("Local data removed successfully.");
  253. } else {
  254. error!("Unexpected Value: {}", confirm);
  255. }
  256. return Ok(())
  257. }
  258. if nickname.is_none() {
  259. error!("Provide a nickname in config file");
  260. return Ok(())
  261. }
  262. if settings.piped {
  263. let file = "/tmp/tau_pipe";
  264. let path = CString::new(file).unwrap();
  265. unsafe { mkfifo(path.as_ptr(), 0o644) };
  266. }
  267. // mkdir datastore_path if not exists
  268. create_dir_all(datastore_path.clone())?;
  269. create_dir_all(datastore_path.join("month"))?;
  270. create_dir_all(datastore_path.join("task"))?;
  271. if settings.generate {
  272. println!("Generating a new workspace");
  273. loop {
  274. println!("Name for the new workspace: ");
  275. let mut workspace = String::new();
  276. stdin().read_line(&mut workspace).expect("Failed to read line");
  277. let workspace = workspace.to_lowercase();
  278. let workspace = workspace.trim();
  279. if workspace.is_empty() && workspace.len() < 3 {
  280. error!("Wrong workspace try again");
  281. continue
  282. }
  283. let secret_key = SecretKey::generate(&mut OsRng);
  284. let encoded = bs58::encode(secret_key.to_bytes());
  285. println!("workspace: {}:{}", workspace, encoded.into_string());
  286. println!("Please add it to the config file.");
  287. break
  288. }
  289. return Ok(())
  290. }
  291. let workspaces = Arc::new(get_workspaces(&settings)?);
  292. if workspaces.is_empty() {
  293. error!("Please add at least one workspace to the config file.");
  294. println!("Run `$ taud --generate` to generate new workspace.");
  295. return Ok(())
  296. }
  297. info!("Initializing taud node");
  298. // Create datastore path if not there already.
  299. let datastore = expand_path(&settings.datastore)?;
  300. fs::create_dir_all(&datastore).await?;
  301. info!("Instantiating event DAG");
  302. let sled_db = sled::open(datastore)?;
  303. let p2p = P2p::new(settings.net.into(), executor.clone()).await;
  304. let event_graph =
  305. EventGraph::new(p2p.clone(), sled_db.clone(), "taud_dag", 0, executor.clone()).await?;
  306. info!("Registering EventGraph P2P protocol");
  307. let event_graph_ = Arc::clone(&event_graph);
  308. let registry = p2p.protocol_registry();
  309. registry
  310. .register(SESSION_ALL, move |channel, _| {
  311. let event_graph_ = event_graph_.clone();
  312. async move { ProtocolEventGraph::init(event_graph_, channel).await.unwrap() }
  313. })
  314. .await;
  315. let (broadcast_snd, broadcast_rcv) = smol::channel::unbounded::<TaskInfo>();
  316. info!(target: "taud", "Starting P2P network");
  317. p2p.clone().start().await?;
  318. info!(target: "taud", "Waiting for some P2P connections...");
  319. sleep(5).await;
  320. // We'll attempt to sync 5 times
  321. if !settings.skip_dag_sync {
  322. for i in 1..=6 {
  323. info!("Syncing event DAG (attempt #{})", i);
  324. match event_graph.dag_sync().await {
  325. Ok(()) => break,
  326. Err(e) => {
  327. if i == 6 {
  328. error!("Failed syncing DAG. Exiting.");
  329. p2p.stop().await;
  330. return Err(Error::DagSyncFailed)
  331. } else {
  332. // TODO: Maybe at this point we should prune or something?
  333. // TODO: Or maybe just tell the user to delete the DAG from FS.
  334. error!("Failed syncing DAG ({}), retrying in 10s...", e);
  335. sleep(10).await;
  336. }
  337. }
  338. }
  339. }
  340. }
  341. ////////////////////
  342. // Listner
  343. ////////////////////
  344. info!(target: "taud", "Starting sync loop task");
  345. let last_sent = RwLock::new(NULL_ID);
  346. let seen = OnceLock::new();
  347. seen.set(sled_db.open_tree("tau_db").unwrap()).unwrap();
  348. ////////////////////
  349. // get history
  350. ////////////////////
  351. let dag_events = event_graph.order_events().await;
  352. let seen_events = seen.get().unwrap();
  353. for event_id in dag_events.iter() {
  354. // If it was seen, skip
  355. if seen_events.contains_key(event_id.as_bytes()).unwrap() {
  356. continue
  357. }
  358. // Get the event from the DAG
  359. let event = event_graph.dag_get(event_id).await.unwrap().unwrap();
  360. // Try to deserialize it. (Here we skip errors)
  361. let Ok((enc_task, _)) = deserialize_async_partial(event.content()).await else { continue };
  362. // Potentially decrypt the privmsg
  363. on_receive_task(&enc_task, &datastore_path, &workspaces, false).await.unwrap();
  364. debug!("Marking event {} as seen", event_id);
  365. seen_events.insert(event_id.as_bytes(), &[]).unwrap();
  366. }
  367. let sync_loop_task = StoppableTask::new();
  368. sync_loop_task.clone().start(
  369. start_sync_loop(
  370. event_graph.clone(),
  371. broadcast_rcv,
  372. workspaces.clone(),
  373. datastore_path.clone(),
  374. settings.piped,
  375. p2p.clone(),
  376. last_sent,
  377. seen.clone(),
  378. ),
  379. |res| async {
  380. match res {
  381. Ok(()) | Err(TaudError::Darkfi(Error::DetachedTaskStopped)) => { /* Do nothing */ }
  382. Err(e) => error!(target: "taud", "Failed starting sync loop task: {}", e),
  383. }
  384. },
  385. TaudError::Darkfi(Error::DetachedTaskStopped),
  386. executor.clone(),
  387. );
  388. // ==============
  389. // p2p dnet setup
  390. // ==============
  391. info!(target: "taud", "Starting dnet subs task");
  392. let json_sub = JsonSubscriber::new("dnet.subscribe_events");
  393. let json_sub_ = json_sub.clone();
  394. let p2p_ = p2p.clone();
  395. let dnet_task = StoppableTask::new();
  396. dnet_task.clone().start(
  397. async move {
  398. let dnet_sub = p2p_.dnet_subscribe().await;
  399. loop {
  400. let event = dnet_sub.receive().await;
  401. debug!("Got dnet event: {:?}", event);
  402. json_sub_.notify(vec![event.into()]).await;
  403. }
  404. },
  405. |res| async {
  406. match res {
  407. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  408. Err(e) => {
  409. error!(target: "taud", "Failed starting dnet subs task: {}", e)
  410. }
  411. }
  412. },
  413. Error::DetachedTaskStopped,
  414. executor.clone(),
  415. );
  416. //
  417. // RPC interface
  418. //
  419. let rpc_interface = Arc::new(JsonRpcInterface::new(
  420. datastore_path.clone(),
  421. broadcast_snd,
  422. nickname.unwrap(),
  423. workspaces.clone(),
  424. p2p.clone(),
  425. json_sub,
  426. ));
  427. let rpc_task = StoppableTask::new();
  428. rpc_task.clone().start(
  429. listen_and_serve(settings.rpc_listen, rpc_interface.clone(), None, executor.clone()),
  430. |res| async move {
  431. match res {
  432. Ok(()) | Err(Error::RpcServerStopped) => rpc_interface.stop_connections().await,
  433. Err(e) => error!(target: "taud", "Failed starting JSON-RPC server: {}", e),
  434. }
  435. },
  436. Error::RpcServerStopped,
  437. executor.clone(),
  438. );
  439. // Signal handling for graceful termination.
  440. let (signals_handler, signals_task) = SignalHandler::new(executor)?;
  441. signals_handler.wait_termination(signals_task).await?;
  442. info!("Caught termination signal, cleaning up and exiting...");
  443. info!(target: "taud", "Stopping JSON-RPC server...");
  444. rpc_task.stop().await;
  445. info!(target: "taud", "Stopping sync loop task...");
  446. sync_loop_task.stop().await;
  447. p2p.stop().await;
  448. Ok(())
  449. }