main.rs 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2024 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. sync::{Arc, OnceLock},
  25. };
  26. use crypto_box::{
  27. aead::{Aead, AeadCore},
  28. ChaChaBox, SecretKey,
  29. };
  30. use darkfi_serial::{
  31. async_trait, deserialize, deserialize_async_partial, serialize, serialize_async,
  32. SerialDecodable, SerialEncodable,
  33. };
  34. use futures::{select, FutureExt};
  35. use libc::mkfifo;
  36. use log::{debug, error, info};
  37. use rand::rngs::OsRng;
  38. use ring::{
  39. rand::SystemRandom,
  40. signature::{Ed25519KeyPair, KeyPair, Signature, UnparsedPublicKey, ED25519},
  41. };
  42. use sled_overlay::sled;
  43. use smol::{fs, stream::StreamExt};
  44. use structopt_toml::StructOptToml;
  45. use tinyjson::JsonValue;
  46. use darkfi::{
  47. async_daemonize,
  48. event_graph::{
  49. proto::{EventPut, ProtocolEventGraph},
  50. Event, EventGraph, EventGraphPtr,
  51. },
  52. net::{session::SESSION_DEFAULT, P2p, P2pPtr},
  53. rpc::{
  54. jsonrpc::JsonSubscriber,
  55. server::{listen_and_serve, RequestHandler},
  56. },
  57. system::{sleep, StoppableTask},
  58. util::path::{expand_path, get_config_path},
  59. Error, Result,
  60. };
  61. mod jsonrpc;
  62. mod settings;
  63. use taud::{
  64. error::{TaudError, TaudResult},
  65. task_info::{TaskEvent, TaskInfo},
  66. util::pipe_write,
  67. };
  68. use crate::{
  69. jsonrpc::JsonRpcInterface,
  70. settings::{Args, CONFIG_FILE, CONFIG_FILE_CONTENTS},
  71. };
  72. struct Workspace {
  73. read_key: ChaChaBox,
  74. write_key: Option<Ed25519KeyPair>,
  75. write_pubkey: UnparsedPublicKey<Vec<u8>>,
  76. }
  77. impl Workspace {
  78. fn new() -> Self {
  79. let secret_key = SecretKey::generate(&mut OsRng);
  80. Self {
  81. read_key: ChaChaBox::new(&secret_key.public_key(), &secret_key),
  82. write_key: None,
  83. write_pubkey: UnparsedPublicKey::new(&ED25519, vec![0]),
  84. }
  85. }
  86. }
  87. #[derive(Debug, Clone, SerialEncodable, SerialDecodable)]
  88. pub struct EncryptedTask {
  89. payload: String,
  90. }
  91. #[derive(SerialEncodable, SerialDecodable)]
  92. struct SignedTask {
  93. task: Vec<u8>,
  94. signature: Vec<u8>,
  95. }
  96. impl SignedTask {
  97. fn new(task: &TaskInfo, signature: Signature) -> Self {
  98. Self { task: serialize(task), signature: signature.as_ref().to_vec() }
  99. }
  100. }
  101. /// Sign then encrypt a task
  102. fn encrypt_sign_task(task: &TaskInfo, workspace: &Workspace) -> TaudResult<EncryptedTask> {
  103. debug!(target: "taud", "start encrypting task");
  104. if workspace.write_key.is_none() {
  105. error!(target: "taud", "You don't have write access")
  106. }
  107. let signature: Signature = workspace.write_key.as_ref().unwrap().sign(&serialize(task)[..]);
  108. let signed_task = SignedTask::new(task, signature);
  109. let nonce = ChaChaBox::generate_nonce(&mut OsRng);
  110. let payload = &serialize(&signed_task)[..];
  111. let mut payload = workspace.read_key.encrypt(&nonce, payload)?;
  112. let mut concat = vec![];
  113. concat.append(&mut nonce.as_slice().to_vec());
  114. concat.append(&mut payload);
  115. let payload = bs58::encode(concat.clone()).into_string();
  116. Ok(EncryptedTask { payload })
  117. }
  118. fn try_decrypt_task(
  119. encrypt_task: &EncryptedTask,
  120. chacha_box: &ChaChaBox,
  121. ) -> TaudResult<SignedTask> {
  122. debug!(target: "taud", "start decrypting task");
  123. let bytes = match bs58::decode(&encrypt_task.payload).into_vec() {
  124. Ok(v) => v,
  125. Err(_) => return Err(TaudError::DecryptionError("Error decoding payload".to_string())),
  126. };
  127. if bytes.len() < 25 {
  128. return Err(TaudError::DecryptionError("Invalid bytes length".to_string()))
  129. }
  130. // Try extracting the nonce
  131. let nonce = bytes[0..24].into();
  132. // Take the remaining ciphertext
  133. let message = &bytes[24..];
  134. // let nonce = encrypt_task.nonce.as_slice();
  135. let decrypted_task = chacha_box.decrypt(nonce, message)?;
  136. let signed_task = deserialize(&decrypted_task)?;
  137. Ok(signed_task)
  138. }
  139. fn parse_configured_workspaces(data: &toml::Value) -> Result<HashMap<String, Workspace>> {
  140. let mut ret = HashMap::new();
  141. let Some(table) = data.as_table() else { return Err(Error::ParseFailed("TOML not a map")) };
  142. let Some(workspace) = table.get("workspace") else { return Ok(ret) };
  143. let Some(workspace) = workspace.as_table() else {
  144. return Err(Error::ParseFailed("`workspace` not a map"))
  145. };
  146. for (name, items) in workspace {
  147. let mut ws = Workspace::new();
  148. if let Some(read_key) = items.get("read_key") {
  149. if let Some(read_key) = read_key.as_str() {
  150. let Ok(read_key_bytes) = bs58::decode(read_key).into_vec() else {
  151. return Err(Error::ParseFailed("Workspace secret not valid base58"))
  152. };
  153. if read_key_bytes.len() != 32 {
  154. return Err(Error::ParseFailed("Workspace read_key not 32 bytes long"))
  155. }
  156. let read_key_bytes: [u8; 32] = read_key_bytes.try_into().unwrap();
  157. let read_key = crypto_box::SecretKey::from(read_key_bytes);
  158. let public = read_key.public_key();
  159. ws.read_key = ChaChaBox::new(&public, &read_key);
  160. } else {
  161. return Err(Error::ParseFailed("Workspace read_key not a string"))
  162. }
  163. } else {
  164. return Err(Error::ParseFailed("Workspace read_key is not set"))
  165. }
  166. if let Some(write_pubkey) = items.get("write_public_key") {
  167. if let Some(write_pubkey) = write_pubkey.as_str() {
  168. if !write_pubkey.is_empty() {
  169. info!(target: "taud", "Found configured write_public_key for {} workspace", name);
  170. let write_pubkey = write_pubkey.to_string();
  171. let decoded_write_pubkey = bs58::decode(write_pubkey).into_vec().unwrap();
  172. ws.write_pubkey = UnparsedPublicKey::new(&ED25519, decoded_write_pubkey);
  173. }
  174. } else {
  175. return Err(Error::ParseFailed("Workspace write_public_key not a string"))
  176. }
  177. } else {
  178. return Err(Error::ParseFailed("Workspace write_public_key is not set"))
  179. }
  180. if let Some(write_key) = items.get("write_key") {
  181. if let Some(write_key) = write_key.as_str() {
  182. if !write_key.is_empty() {
  183. info!(target: "taud", "Found configured write_key for {} workspace", name);
  184. let write_key = write_key.to_string();
  185. let pkcs8_bytes = bs58::decode(write_key).into_vec().unwrap();
  186. let ed25519 = match Ed25519KeyPair::from_pkcs8(pkcs8_bytes.as_ref()) {
  187. Ok(key) => key,
  188. Err(e) => {
  189. error!(target: "taud", "Failed parsing write_key: {}", e);
  190. return Err(Error::ParseFailed("Failed parsing write_key"))
  191. }
  192. };
  193. ws.write_key = Some(ed25519);
  194. }
  195. } else {
  196. return Err(Error::ParseFailed("Workspace write_key not a string"))
  197. }
  198. }
  199. if let Some(wrt_key) = ws.write_key.as_ref() {
  200. if wrt_key.public_key().as_ref() != ws.write_pubkey.as_ref() {
  201. error!(target: "taud", "Wrong keypair for {} workspace, the workspace is not added!", name);
  202. continue
  203. }
  204. }
  205. info!(target: "taud", "Configured NaCl box for workspace {}", name);
  206. ret.insert(name.to_string(), ws);
  207. }
  208. Ok(ret)
  209. }
  210. async fn get_workspaces(settings: &Args) -> Result<HashMap<String, Workspace>> {
  211. let config_path = get_config_path(settings.config.clone(), CONFIG_FILE)?;
  212. let contents = fs::read_to_string(config_path).await?;
  213. let contents = match toml::from_str(&contents) {
  214. Ok(v) => v,
  215. Err(e) => {
  216. error!(target: "taud", "Failed parsing TOML config: {}", e);
  217. return Err(Error::ParseFailed("Failed parsing TOML config"))
  218. }
  219. };
  220. let workspaces = parse_configured_workspaces(&contents)?;
  221. Ok(workspaces)
  222. }
  223. /// Atomically mark a message as seen.
  224. pub async fn mark_seen(
  225. sled_db: sled::Db,
  226. seen: OnceLock<sled::Tree>,
  227. event_id: &blake3::Hash,
  228. ) -> Result<()> {
  229. let db = seen.get_or_init(|| sled_db.open_tree("tau_seen").unwrap());
  230. debug!(target: "taud", "Marking event {} as seen", event_id);
  231. let mut batch = sled::Batch::default();
  232. batch.insert(event_id.as_bytes(), &[]);
  233. Ok(db.apply_batch(batch)?)
  234. }
  235. /// Check if a message was already marked seen.
  236. pub async fn is_seen(
  237. sled_db: sled::Db,
  238. seen: OnceLock<sled::Tree>,
  239. event_id: &blake3::Hash,
  240. ) -> Result<bool> {
  241. let db = seen.get_or_init(|| sled_db.open_tree("tau_seen").unwrap());
  242. Ok(db.contains_key(event_id.as_bytes())?)
  243. }
  244. #[allow(clippy::too_many_arguments)]
  245. async fn start_sync_loop(
  246. event_graph: EventGraphPtr,
  247. broadcast_rcv: smol::channel::Receiver<TaskInfo>,
  248. workspaces: Arc<HashMap<String, Workspace>>,
  249. sled_db: sled::Db,
  250. settings: Args,
  251. p2p: P2pPtr,
  252. seen: OnceLock<sled::Tree>,
  253. ) -> TaudResult<()> {
  254. let incoming = event_graph.event_pub.clone().subscribe().await;
  255. loop {
  256. select! {
  257. // Process message from Tau client
  258. task_event = broadcast_rcv.recv().fuse() => {
  259. let tk = task_event.map_err(Error::from)?;
  260. if workspaces.contains_key(&tk.workspace) {
  261. let ws = workspaces.get(&tk.workspace).unwrap();
  262. let encrypted_task = encrypt_sign_task(&tk, ws)?;
  263. info!(target: "taud", "Send the task: ref: {}", tk.ref_id);
  264. // Build a DAG event and return it.
  265. let event = Event::new(
  266. serialize_async(&encrypted_task).await,
  267. &event_graph,
  268. )
  269. .await;
  270. // If it fails for some reason, for now, we just note it
  271. // and pass.
  272. if let Err(e) = event_graph.dag_insert(&[event.clone()]).await {
  273. error!(target: "taud", "Failed inserting new event to DAG: {}", e);
  274. } else {
  275. // Otherwise, broadcast it
  276. let self_version = p2p.settings().read().await.app_version.clone();
  277. let connected_peers = p2p.hosts().peers();
  278. let mut peers_with_matched_version = vec![];
  279. let mut peers_with_different_version = vec![];
  280. for peer in connected_peers {
  281. let peer_version = peer.version.lock().await.clone();
  282. if let Some(ref peer_version) = peer_version {
  283. if self_version == peer_version.version {
  284. peers_with_matched_version.push(peer)
  285. } else {
  286. peers_with_different_version.push(peer)
  287. }
  288. }
  289. }
  290. if !peers_with_matched_version.is_empty() {
  291. p2p.broadcast_to(&EventPut(event.clone()), &peers_with_matched_version).await;
  292. }
  293. if !peers_with_different_version.is_empty() {
  294. let mut event = event;
  295. event.timestamp /= 1000;
  296. p2p.broadcast_to(&EventPut(event), &peers_with_different_version).await;
  297. }
  298. // p2p.broadcast(&EventPut(event)).await;
  299. }
  300. }
  301. }
  302. // Process message from the network. These should only be EncryptedTask.
  303. task_event = incoming.receive().fuse() => {
  304. let event_id = task_event.id();
  305. if is_seen(sled_db.clone(), seen.clone(), &event_id).await? {
  306. continue
  307. }
  308. mark_seen(sled_db.clone(), seen.clone(), &event_id).await?;
  309. // Try to deserialize the `Event`'s content into a `EncryptedTask`
  310. let enc_task: EncryptedTask = match deserialize_async_partial(task_event.content()).await {
  311. Ok((v, _)) => v,
  312. Err(e) => {
  313. error!(target: "taud", "[TAUD] Failed deserializing incoming EncryptedTask event: {}", e);
  314. continue
  315. }
  316. };
  317. on_receive_task(&enc_task, &workspaces, &settings)
  318. .await?;
  319. }
  320. }
  321. }
  322. }
  323. /// Handle a received task, decrypt it, verify it, optionally write it
  324. /// to a named pipe and save it on disk.
  325. async fn on_receive_task(
  326. enc_task: &EncryptedTask,
  327. workspaces: &HashMap<String, Workspace>,
  328. settings: &Args,
  329. ) -> TaudResult<()> {
  330. for (ws_name, workspace) in workspaces.iter() {
  331. let signed_task = try_decrypt_task(enc_task, &workspace.read_key);
  332. if let Err(e) = signed_task {
  333. debug!(target: "taud", "Unable to decrypt the task: {}", e);
  334. continue
  335. }
  336. if workspace
  337. .write_pubkey
  338. .verify(&signed_task.as_ref().unwrap().task, &signed_task.as_ref().unwrap().signature)
  339. .is_err()
  340. {
  341. // *verified.lock().await = false;
  342. error!(target: "taud", "Task is not verified: wrong write_public_key");
  343. error!(target: "taud", "Task is not saved");
  344. continue
  345. }
  346. // else {
  347. // *verified.lock().await = true;
  348. // }
  349. let mut task: TaskInfo = deserialize(&signed_task.unwrap().task)?;
  350. info!(target: "taud", "Save the task: ref: {}", task.ref_id);
  351. task.workspace.clone_from(ws_name);
  352. let datastore_path = expand_path(&settings.datastore)?;
  353. if settings.piped {
  354. // if we can't load the task then it's a new task.
  355. // otherwise it's a modification.
  356. match TaskInfo::load(&task.ref_id, &datastore_path) {
  357. Ok(loaded_task) => {
  358. let loaded_events = loaded_task.events;
  359. let mut events = task.events.clone();
  360. events.retain(|ev| !loaded_events.contains(ev));
  361. let file = settings.pipe_path.clone();
  362. let mut pipe_write = pipe_write(file)?;
  363. let mut task_clone = task.clone();
  364. task_clone.events = events;
  365. let json: JsonValue = (&task_clone).into();
  366. pipe_write.write_all(json.stringify().unwrap().as_bytes())?;
  367. }
  368. Err(_) => {
  369. let file = settings.pipe_path.clone();
  370. let mut pipe_write = pipe_write(file)?;
  371. let mut task_clone = task.clone();
  372. task_clone.events.push(TaskEvent::new(
  373. "add_task".to_string(),
  374. task_clone.owner.clone(),
  375. "".to_string(),
  376. ));
  377. let json: JsonValue = (&task_clone).into();
  378. pipe_write.write_all(json.stringify().unwrap().as_bytes())?;
  379. }
  380. }
  381. }
  382. task.save(&datastore_path)?;
  383. }
  384. Ok(())
  385. }
  386. async_daemonize!(realmain);
  387. async fn realmain(settings: Args, executor: Arc<smol::Executor<'static>>) -> Result<()> {
  388. let datastore_path = expand_path(&settings.datastore)?;
  389. let nickname =
  390. if settings.nickname.is_some() { settings.nickname.clone() } else { env::var("USER").ok() };
  391. if settings.refresh {
  392. println!("Removing local data in: {:?} (yes/no)? ", datastore_path);
  393. let mut confirm = String::new();
  394. stdin().read_line(&mut confirm).expect("Failed to read line");
  395. let confirm = confirm.to_lowercase();
  396. let confirm = confirm.trim();
  397. if confirm == "yes" || confirm == "y" {
  398. remove_dir_all(datastore_path).unwrap_or(());
  399. println!("Local data removed successfully.");
  400. } else {
  401. error!(target: "taud", "Unexpected Value: {}", confirm);
  402. }
  403. return Ok(())
  404. }
  405. if nickname.is_none() {
  406. error!(target: "taud", "Provide a nickname in config file");
  407. return Ok(())
  408. }
  409. if settings.piped {
  410. let file = settings.pipe_path.clone();
  411. let path = CString::new(file).unwrap();
  412. unsafe { mkfifo(path.as_ptr(), 0o644) };
  413. }
  414. // mkdir datastore_path if not exists
  415. create_dir_all(datastore_path.clone())?;
  416. create_dir_all(datastore_path.join("month"))?;
  417. create_dir_all(datastore_path.join("task"))?;
  418. if settings.generate {
  419. println!("Generating a new workspace");
  420. loop {
  421. println!("Name for the new workspace: ");
  422. let mut workspace = String::new();
  423. stdin().read_line(&mut workspace).expect("Failed to read line");
  424. let workspace = workspace.to_lowercase();
  425. let workspace = workspace.trim();
  426. if workspace.is_empty() && workspace.len() < 3 {
  427. error!(target: "taud", "Wrong workspace try again");
  428. continue
  429. }
  430. // Chachabox secret key (read_key) used for encrypting tasks.
  431. let secret_key = SecretKey::generate(&mut OsRng);
  432. let encoded = bs58::encode(secret_key.to_bytes());
  433. // Ed25519 secret key (write_key) used for signing tasks.
  434. let rng = SystemRandom::new();
  435. let pkcs8_bytes = Ed25519KeyPair::generate_pkcs8(&rng).unwrap();
  436. // openssl genpkey -algorithm ED25519
  437. let kp = Ed25519KeyPair::from_pkcs8(pkcs8_bytes.as_ref()).unwrap();
  438. let sk = bs58::encode(pkcs8_bytes).into_string();
  439. // Ed25519 public key (write_public_key) used for verifying tasks.
  440. let peer_public_key_bytes = kp.public_key().as_ref();
  441. let pk = bs58::encode(peer_public_key_bytes).into_string();
  442. println!("Please add the following to the config file:");
  443. println!("[workspace.\"{}\"]", workspace);
  444. println!("read_key = \"{}\"", encoded.into_string());
  445. println!("write_key = \"{}\"", sk);
  446. println!("write_public_key = \"{}\"", pk);
  447. break
  448. }
  449. return Ok(())
  450. }
  451. let workspaces = Arc::new(get_workspaces(&settings).await?);
  452. // let verified = Arc::new(Mutex::new(false));
  453. if workspaces.is_empty() {
  454. error!(target: "taud", "Please add at least one workspace to the config file.");
  455. println!("Run `$ taud --generate` to generate new workspace.");
  456. return Ok(())
  457. }
  458. info!(target: "taud", "Initializing taud node");
  459. // Create datastore path if not there already.
  460. let datastore = expand_path(&settings.datastore)?;
  461. fs::create_dir_all(&datastore).await?;
  462. let replay_datastore = expand_path(&settings.replay_datastore)?;
  463. let replay_mode = settings.replay_mode;
  464. info!(target: "taud", "Instantiating event DAG");
  465. let sled_db = sled::open(datastore)?;
  466. let p2p = P2p::new(settings.net.clone().into(), executor.clone()).await?;
  467. let event_graph = EventGraph::new(
  468. p2p.clone(),
  469. sled_db.clone(),
  470. replay_datastore,
  471. replay_mode,
  472. "taud_dag",
  473. 0,
  474. executor.clone(),
  475. )
  476. .await?;
  477. info!(target: "taud", "Registering EventGraph P2P protocol");
  478. let event_graph_ = Arc::clone(&event_graph);
  479. let registry = p2p.protocol_registry();
  480. registry
  481. .register(SESSION_DEFAULT, move |channel, _| {
  482. let event_graph_ = event_graph_.clone();
  483. async move { ProtocolEventGraph::init(event_graph_, channel).await.unwrap() }
  484. })
  485. .await;
  486. let (broadcast_snd, broadcast_rcv) = smol::channel::unbounded::<TaskInfo>();
  487. info!(target: "taud", "Starting P2P network");
  488. p2p.clone().start().await?;
  489. info!(target: "taud", "Waiting for some P2P connections...");
  490. sleep(5).await;
  491. // We'll attempt to sync {sync_attempts} times
  492. if !settings.skip_dag_sync {
  493. for i in 1..=settings.sync_attempts {
  494. info!(target: "taud", "Syncing event DAG (attempt #{})", i);
  495. match event_graph.dag_sync().await {
  496. Ok(()) => break,
  497. Err(e) => {
  498. if i == settings.sync_attempts {
  499. error!(target: "taud", "Failed syncing DAG. Exiting.");
  500. p2p.stop().await;
  501. return Err(Error::DagSyncFailed)
  502. } else {
  503. // TODO: Maybe at this point we should prune or something?
  504. // TODO: Or maybe just tell the user to delete the DAG from FS.
  505. error!(
  506. "Failed syncing DAG ({}), retrying in {}s...",
  507. e, settings.sync_timeout
  508. );
  509. sleep(settings.sync_timeout.into()).await;
  510. }
  511. }
  512. }
  513. }
  514. } else {
  515. *event_graph.synced.write().await = true;
  516. }
  517. let seen = OnceLock::new();
  518. seen.set(sled_db.open_tree("tau_seen").unwrap()).unwrap();
  519. ////////////////////
  520. // get history
  521. ////////////////////
  522. let dag_events = event_graph.order_events().await;
  523. for event in dag_events.iter() {
  524. let event_id = event.id();
  525. // If it was seen, skip
  526. if is_seen(sled_db.clone(), seen.clone(), &event_id).await? {
  527. continue
  528. }
  529. mark_seen(sled_db.clone(), seen.clone(), &event_id).await?;
  530. // Try to deserialize it. (Here we skip errors)
  531. let Ok((enc_task, _)) = deserialize_async_partial(event.content()).await else { continue };
  532. // Potentially decrypt the privmsg
  533. on_receive_task(&enc_task, &workspaces, &settings).await.unwrap();
  534. }
  535. ////////////////////
  536. // Listner
  537. ////////////////////
  538. info!(target: "taud", "Starting sync loop task");
  539. let sync_loop_task = StoppableTask::new();
  540. sync_loop_task.clone().start(
  541. start_sync_loop(
  542. event_graph.clone(),
  543. broadcast_rcv,
  544. workspaces.clone(),
  545. sled_db.clone(),
  546. settings.clone(),
  547. p2p.clone(),
  548. seen.clone(),
  549. ),
  550. |res| async {
  551. match res {
  552. Ok(()) | Err(TaudError::Darkfi(Error::DetachedTaskStopped)) => { /* Do nothing */ }
  553. Err(e) => error!(target: "taud", "Failed stopping sync loop task: {}", e),
  554. }
  555. },
  556. TaudError::Darkfi(Error::DetachedTaskStopped),
  557. executor.clone(),
  558. );
  559. // ==============
  560. // p2p dnet setup
  561. // ==============
  562. info!(target: "taud", "Starting dnet subs task");
  563. let json_sub = JsonSubscriber::new("dnet.subscribe_events");
  564. let json_sub_ = json_sub.clone();
  565. let p2p_ = p2p.clone();
  566. let dnet_task = StoppableTask::new();
  567. dnet_task.clone().start(
  568. async move {
  569. let dnet_sub = p2p_.dnet_subscribe().await;
  570. loop {
  571. let event = dnet_sub.receive().await;
  572. debug!(target: "taud", "Got dnet event: {:?}", event);
  573. json_sub_.notify(vec![event.into()].into()).await;
  574. }
  575. },
  576. |res| async {
  577. match res {
  578. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  579. Err(e) => {
  580. error!(target: "taud", "Failed stopping dnet subs task: {}", e)
  581. }
  582. }
  583. },
  584. Error::DetachedTaskStopped,
  585. executor.clone(),
  586. );
  587. info!("Starting deg subs task");
  588. let deg_sub = JsonSubscriber::new("deg.subscribe_events");
  589. let deg_sub_ = deg_sub.clone();
  590. let event_graph_ = event_graph.clone();
  591. let deg_task = StoppableTask::new();
  592. deg_task.clone().start(
  593. async move {
  594. let deg_sub = event_graph_.deg_subscribe().await;
  595. loop {
  596. let event = deg_sub.receive().await;
  597. debug!(target: "taud", "Got deg event: {:?}", event);
  598. deg_sub_.notify(vec![event.into()].into()).await;
  599. }
  600. },
  601. |res| async {
  602. match res {
  603. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  604. Err(e) => panic!("{}", e),
  605. }
  606. },
  607. Error::DetachedTaskStopped,
  608. executor.clone(),
  609. );
  610. //
  611. // RPC interface
  612. //
  613. let rpc_interface = Arc::new(JsonRpcInterface::new(
  614. datastore_path.clone(),
  615. broadcast_snd,
  616. nickname.unwrap(),
  617. workspaces.clone(),
  618. p2p.clone(),
  619. event_graph.clone(),
  620. json_sub,
  621. deg_sub,
  622. ));
  623. let rpc_task = StoppableTask::new();
  624. rpc_task.clone().start(
  625. listen_and_serve(settings.rpc_listen, rpc_interface.clone(), None, executor.clone()),
  626. |res| async move {
  627. match res {
  628. Ok(()) | Err(Error::RpcServerStopped) => rpc_interface.stop_connections().await,
  629. Err(e) => error!(target: "taud", "Failed stopping JSON-RPC server: {}", e),
  630. }
  631. },
  632. Error::RpcServerStopped,
  633. executor.clone(),
  634. );
  635. // Signal handling for graceful termination.
  636. let (signals_handler, signals_task) = SignalHandler::new(executor)?;
  637. signals_handler.wait_termination(signals_task).await?;
  638. info!(target: "taud", "Caught termination signal, cleaning up and exiting...");
  639. info!(target: "taud", "Stopping P2P network");
  640. p2p.stop().await;
  641. info!(target: "taud", "Stopping sync loop task...");
  642. sync_loop_task.stop().await;
  643. info!(target: "taud", "Stopping JSON-RPC server...");
  644. rpc_task.stop().await;
  645. dnet_task.stop().await;
  646. deg_task.stop().await;
  647. info!(target: "taud", "Flushing sled database...");
  648. let flushed_bytes = sled_db.flush_async().await?;
  649. info!(target: "taud", "Flushed {} bytes", flushed_bytes);
  650. info!(target: "taud", "Shut down successfully");
  651. Ok(())
  652. }