main.rs 25 KB

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