main.rs 38 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2026 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::{BTreeMap, HashMap},
  20. env,
  21. ffi::CString,
  22. fs::{create_dir_all, remove_dir_all, File},
  23. io::{stdin, Write},
  24. str::FromStr,
  25. sync::{atomic::Ordering, 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 kvdb_overlay::{Batch, Database, Tree};
  37. use libc::mkfifo;
  38. use rand::rngs::OsRng;
  39. use smol::{fs, stream::StreamExt};
  40. use structopt_toml::StructOptToml;
  41. use tinyjson::JsonValue;
  42. use tracing::{debug, error, info, warn};
  43. use darkfi::{
  44. async_daemonize,
  45. event_graph::{
  46. proto::{EventPut, ProtocolEventGraph},
  47. Event, EventGraph, EventGraphConfig, 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. pasta_prelude::PrimeField,
  60. schnorr::{SchnorrPublic, SchnorrSecret, Signature},
  61. Keypair, PublicKey,
  62. };
  63. // =====================================================================
  64. // Taud consensus parameters.
  65. //
  66. // These define the EventGraph configuration that EVERY Taud node
  67. // in the network must agree on. Changing any of them is a hard fork.
  68. // They are passed verbatim to `EventGraph::new` at startup.
  69. // =====================================================================
  70. /// Epoch origin for DAG rotation (UTC midnight, 1 March 2025).
  71. /// Rotation boundaries are computed as offsets from this point.
  72. const TAUD_INITIAL_GENESIS: u64 = 1_740_787_200_000;
  73. /// DAG rotation period, in hours.
  74. const TAUD_HOURS_ROTATION: u64 = 0;
  75. /// Genesis payload. Two protocols MUST use distinct values; this
  76. /// also feeds into `RlnAppId::from_genesis` so RLN signals from one
  77. /// deployment never appear valid on another.
  78. const TAUD_GENESIS_CONTENTS: &[u8] = b"taud-v1";
  79. /// How many rotation periods to keep in the rolling DAG window.
  80. /// With `hours_rotation = 1` and `max_dags = 24`, this gives a
  81. /// 24-hour history window. Older events are evicted from kvdb.
  82. const TAUD_MAX_DAGS: usize = 1;
  83. /// Per-epoch limit printed by `--gen-rln-identity`.
  84. fn generated_rln_identity_user_msg_limit() -> u64 {
  85. darkfi::event_graph::rln::GENESIS_USER_MSG_LIMIT
  86. }
  87. mod jsonrpc;
  88. mod settings;
  89. use taud::{
  90. error::{TaudError, TaudResult},
  91. rln::{
  92. load_default_rln_identity, reserve_rln_message_id_in_store, RlnIdentity,
  93. RlnMessageReservation,
  94. },
  95. task_info::{TaskEvent, TaskInfo},
  96. util::pipe_write,
  97. };
  98. use crate::{
  99. jsonrpc::JsonRpcInterface,
  100. settings::{Args, CONFIG_FILE, CONFIG_FILE_CONTENTS},
  101. };
  102. struct Workspace {
  103. read_key: ChaChaBox,
  104. write_key: Option<darkfi_sdk::crypto::SecretKey>,
  105. write_pubkey: PublicKey,
  106. }
  107. impl Workspace {
  108. fn new() -> Self {
  109. let secret_key = SecretKey::generate(&mut OsRng);
  110. let keypair = Keypair::default();
  111. Self {
  112. read_key: ChaChaBox::new(&secret_key.public_key(), &secret_key),
  113. write_key: None,
  114. write_pubkey: keypair.public,
  115. }
  116. }
  117. }
  118. #[derive(Debug, Clone, SerialEncodable, SerialDecodable)]
  119. pub struct EncryptedTask {
  120. payload: String,
  121. }
  122. #[derive(SerialEncodable, SerialDecodable)]
  123. struct SignedTask {
  124. task: Vec<u8>,
  125. signature: Signature,
  126. }
  127. impl SignedTask {
  128. fn new(task: &TaskInfo, signature: Signature) -> Self {
  129. Self { task: serialize(task), signature }
  130. }
  131. }
  132. /// Sign then encrypt a task
  133. fn encrypt_sign_task(task: &TaskInfo, workspace: &Workspace) -> TaudResult<EncryptedTask> {
  134. debug!(target: "taud", "start encrypting task");
  135. if workspace.write_key.is_none() {
  136. error!(target: "taud", "You don't have write access")
  137. }
  138. let signature: Signature = workspace.write_key.as_ref().unwrap().sign(&serialize(task)[..]);
  139. let signed_task = SignedTask::new(task, signature);
  140. let nonce = ChaChaBox::generate_nonce(&mut OsRng);
  141. let payload = &serialize(&signed_task)[..];
  142. let mut payload = workspace.read_key.encrypt(&nonce, payload)?;
  143. let mut concat = vec![];
  144. concat.append(&mut nonce.as_slice().to_vec());
  145. concat.append(&mut payload);
  146. let payload = bs58::encode(concat.clone()).into_string();
  147. Ok(EncryptedTask { payload })
  148. }
  149. fn try_decrypt_task(
  150. encrypt_task: &EncryptedTask,
  151. chacha_box: &ChaChaBox,
  152. ) -> TaudResult<SignedTask> {
  153. debug!(target: "taud", "start decrypting task");
  154. let bytes = match bs58::decode(&encrypt_task.payload).into_vec() {
  155. Ok(v) => v,
  156. Err(_) => return Err(TaudError::DecryptionError("Error decoding payload".to_string())),
  157. };
  158. if bytes.len() < 25 {
  159. return Err(TaudError::DecryptionError("Invalid bytes length".to_string()))
  160. }
  161. // Try extracting the nonce
  162. let nonce = bytes[0..24].into();
  163. // Take the remaining ciphertext
  164. let message = &bytes[24..];
  165. // let nonce = encrypt_task.nonce.as_slice();
  166. let decrypted_task = chacha_box.decrypt(nonce, message)?;
  167. let signed_task = deserialize(&decrypted_task)?;
  168. Ok(signed_task)
  169. }
  170. fn parse_configured_workspaces(data: &toml::Value) -> Result<BTreeMap<String, Workspace>> {
  171. let mut ret = BTreeMap::new();
  172. let Some(table) = data.as_table() else { return Err(Error::ParseFailed("TOML not a map")) };
  173. let Some(workspace) = table.get("workspace") else { return Ok(ret) };
  174. let Some(workspace) = workspace.as_table() else {
  175. return Err(Error::ParseFailed("`workspace` not a map"))
  176. };
  177. for (name, items) in workspace {
  178. let mut ws = Workspace::new();
  179. if let Some(read_key) = items.get("read_key") {
  180. if let Some(read_key) = read_key.as_str() {
  181. let Ok(read_key_bytes) = bs58::decode(read_key).into_vec() else {
  182. return Err(Error::ParseFailed("Workspace secret not valid base58"))
  183. };
  184. if read_key_bytes.len() != 32 {
  185. return Err(Error::ParseFailed("Workspace read_key not 32 bytes long"))
  186. }
  187. let read_key_bytes: [u8; 32] = read_key_bytes.try_into().unwrap();
  188. let read_key = crypto_box::SecretKey::from(read_key_bytes);
  189. let public = read_key.public_key();
  190. ws.read_key = ChaChaBox::new(&public, &read_key);
  191. } else {
  192. return Err(Error::ParseFailed("Workspace read_key not a string"))
  193. }
  194. } else {
  195. return Err(Error::ParseFailed("Workspace read_key is not set"))
  196. }
  197. if let Some(write_pubkey) = items.get("write_public_key") {
  198. if let Some(write_pubkey) = write_pubkey.as_str() {
  199. if !write_pubkey.is_empty() {
  200. info!(target: "taud", "Found configured write_public_key for {name} workspace");
  201. let write_key = PublicKey::from_str(write_pubkey).unwrap();
  202. // let write_pubkey = write_pubkey.to_string();
  203. // let decoded_write_pubkey = bs58::decode(write_pubkey).into_vec().unwrap();
  204. ws.write_pubkey = write_key;
  205. }
  206. } else {
  207. return Err(Error::ParseFailed("Workspace write_public_key not a string"))
  208. }
  209. } else {
  210. return Err(Error::ParseFailed("Workspace write_public_key is not set"))
  211. }
  212. if let Some(write_key) = items.get("write_key") {
  213. if let Some(write_key) = write_key.as_str() {
  214. if !write_key.is_empty() {
  215. info!(target: "taud", "Found configured write_key for {name} workspace");
  216. let write_key = write_key.to_string();
  217. let write_key_bytes = bs58::decode(write_key).into_vec().unwrap();
  218. let secret = match darkfi_sdk::crypto::SecretKey::from_bytes(
  219. write_key_bytes.try_into().unwrap(),
  220. ) {
  221. Ok(key) => key,
  222. Err(e) => {
  223. error!(target: "taud", "Failed parsing write_key: {e}");
  224. return Err(Error::ParseFailed("Failed parsing write_key"))
  225. }
  226. };
  227. ws.write_key = Some(secret);
  228. }
  229. } else {
  230. return Err(Error::ParseFailed("Workspace write_key not a string"))
  231. }
  232. }
  233. if let Some(wrt_key) = ws.write_key.as_ref() {
  234. let pk = PublicKey::from_secret(*wrt_key);
  235. if pk != ws.write_pubkey {
  236. error!(target: "taud", "Wrong keypair for {name} workspace, the workspace is not added!");
  237. continue
  238. }
  239. }
  240. info!(target: "taud", "Configured NaCl box for workspace {name}");
  241. ret.insert(name.to_string(), ws);
  242. }
  243. Ok(ret)
  244. }
  245. async fn get_workspaces(settings: &Args) -> Result<BTreeMap<String, Workspace>> {
  246. let config_path = get_config_path(settings.config.clone(), CONFIG_FILE)?;
  247. let contents = fs::read_to_string(config_path).await?;
  248. let contents = match toml::from_str(&contents) {
  249. Ok(v) => v,
  250. Err(e) => {
  251. error!(target: "taud", "Failed parsing TOML config: {e}");
  252. return Err(Error::ParseFailed("Failed parsing TOML config"))
  253. }
  254. };
  255. let workspaces = parse_configured_workspaces(&contents)?;
  256. Ok(workspaces)
  257. }
  258. /// Atomically mark a message as seen.
  259. pub async fn mark_seen(
  260. kvdb: Database,
  261. seen: OnceLock<Tree>,
  262. event_id: &blake3::Hash,
  263. ) -> Result<()> {
  264. let tree = seen.get_or_init(|| kvdb.open_tree_default("tau_seen").unwrap());
  265. debug!(target: "taud", "Marking event {event_id} as seen");
  266. let mut batch = Batch::new();
  267. batch.insert(event_id.as_bytes(), &[]);
  268. Ok(kvdb.atomic_write(&[(tree, &batch)])?)
  269. }
  270. /// Check if a message was already marked seen.
  271. pub async fn is_seen(
  272. kvdb: Database,
  273. seen: OnceLock<Tree>,
  274. event_id: &blake3::Hash,
  275. ) -> Result<bool> {
  276. let tree = seen.get_or_init(|| kvdb.open_tree_default("tau_seen").unwrap());
  277. Ok(tree.contains_key(event_id.as_bytes())?)
  278. }
  279. #[allow(clippy::too_many_arguments)]
  280. async fn start_sync_loop(
  281. event_graph: EventGraphPtr,
  282. broadcast_rcv: smol::channel::Receiver<TaskInfo>,
  283. workspaces: Arc<BTreeMap<String, Workspace>>,
  284. kvdb: Database,
  285. settings: Args,
  286. p2p: P2pPtr,
  287. seen: OnceLock<Tree>,
  288. rln_identity: Arc<smol::lock::RwLock<Option<RlnIdentity>>>,
  289. ) -> TaudResult<()> {
  290. let incoming = event_graph.event_pub.clone().subscribe().await;
  291. loop {
  292. select! {
  293. // Process message from Tau client
  294. task_event = broadcast_rcv.recv().fuse() => {
  295. let tk = task_event.map_err(Error::from)?;
  296. if workspaces.contains_key(&tk.workspace) {
  297. let ws = workspaces.get(&tk.workspace).unwrap();
  298. let encrypted_task = encrypt_sign_task(&tk, ws)?;
  299. info!(target: "taud", "Send the task: ref: {}", tk.ref_id);
  300. // Build a DAG event and return it.
  301. let event = match Event::new(serialize_async(&encrypted_task).await, &event_graph).await {
  302. Ok(event) => event,
  303. Err(e) => {
  304. error!(target: "taud", "Failed creating new DAG event: {e}");
  305. continue
  306. }
  307. };
  308. let current_genesis = event_graph.current_genesis.read().await;
  309. let dag_name = current_genesis.header.timestamp.to_string();
  310. drop(current_genesis);
  311. // Build the RLN signal blob before touching the local DAG when RLN
  312. // is enabled. With RLN disabled, outbound events deliberately carry
  313. // no proof blob.
  314. let blob = if event_graph.rln_enabled() {
  315. let (rln_identity, mid) = {
  316. let mut active = rln_identity.write().await;
  317. match reserve_rln_message_id_in_store(
  318. &kvdb,
  319. &mut active,
  320. event.header.timestamp,
  321. )
  322. .await?
  323. {
  324. RlnMessageReservation::Reserved { identity, message_id } => {
  325. (identity, message_id)
  326. }
  327. RlnMessageReservation::MissingIdentity => {
  328. warn!(target: "taud", "No RLN identity registered; refusing to send. Run `tau rln register ...` to register.");
  329. continue
  330. }
  331. RlnMessageReservation::BudgetExhausted => {
  332. warn!(target: "taud", "RLN message budget exhausted for this epoch; dropping message to avoid slash");
  333. continue
  334. }
  335. }
  336. };
  337. match rln_identity.create_signal(&event, mid, &event_graph).await {
  338. Ok(blob) => serialize_async(&blob).await,
  339. Err(e) => {
  340. error!(target: "taud", "Failed creating RLN signal proof: {e}");
  341. continue
  342. }
  343. }
  344. } else {
  345. Vec::new()
  346. };
  347. if let Err(e) = event_graph.insert_signal_with_blob(&event, &blob, &dag_name).await {
  348. error!(target: "taud", "Failed inserting new event to DAG: {e}");
  349. } else {
  350. // Otherwise, broadcast it. Taud runs EventGraph with RLN disabled
  351. // by default, so the blob is empty unless RLN was enabled.
  352. if let Err(e) = p2p.broadcast(&EventPut(event, blob)).await {
  353. error!(target: "taud", "Event broadcast was not admitted: {e}");
  354. }
  355. }
  356. }
  357. }
  358. // Process message from the network. These should only be EncryptedTask.
  359. task_event = incoming.receive().fuse() => {
  360. let event_id = task_event.header.id();
  361. if is_seen(kvdb.clone(), seen.clone(), &event_id).await? {
  362. continue
  363. }
  364. mark_seen(kvdb.clone(), seen.clone(), &event_id).await?;
  365. // Try to deserialize the `Event`'s content into a `EncryptedTask`
  366. let enc_task: EncryptedTask = match deserialize_async_partial(task_event.content()).await {
  367. Ok((v, _)) => v,
  368. Err(e) => {
  369. error!(target: "taud", "[TAUD] Failed deserializing incoming EncryptedTask event: {e}");
  370. continue
  371. }
  372. };
  373. on_receive_task(&enc_task, &workspaces, &settings)
  374. .await?;
  375. }
  376. }
  377. }
  378. }
  379. /// Handle a received task, decrypt it, verify it, optionally write it
  380. /// to a named pipe and save it on disk.
  381. async fn on_receive_task(
  382. enc_task: &EncryptedTask,
  383. workspaces: &BTreeMap<String, Workspace>,
  384. settings: &Args,
  385. ) -> TaudResult<()> {
  386. for (ws_name, workspace) in workspaces.iter() {
  387. let signed_task = try_decrypt_task(enc_task, &workspace.read_key);
  388. if let Err(e) = signed_task {
  389. debug!(target: "taud", "Unable to decrypt the task: {e}");
  390. continue
  391. }
  392. if !workspace
  393. .write_pubkey
  394. .verify(&signed_task.as_ref().unwrap().task, &signed_task.as_ref().unwrap().signature)
  395. {
  396. error!(target: "taud", "Task is not verified: wrong write_public_key");
  397. error!(target: "taud", "Task is not saved");
  398. continue
  399. }
  400. let mut task: TaskInfo = deserialize(&signed_task.unwrap().task)?;
  401. info!(target: "taud", "Save the task: ref: {}", task.ref_id);
  402. task.workspace.clone_from(ws_name);
  403. let datastore_path = expand_path(&settings.datastore)?;
  404. // Push a notification to a fifo if set
  405. if settings.piped {
  406. // if we can't load the task then it's a new task.
  407. // otherwise it's a modification.
  408. match TaskInfo::load(&task.ref_id, &datastore_path) {
  409. Ok(loaded_task) => {
  410. let loaded_events = loaded_task.events;
  411. let mut events = task.events.clone();
  412. events.retain(|ev| !loaded_events.contains(ev));
  413. let file = settings.pipe_path.clone();
  414. let mut pipe_write = pipe_write(file)?;
  415. let mut task_clone = task.clone();
  416. task_clone.events = events;
  417. let json: JsonValue = (&task_clone).into();
  418. pipe_write.write_all(json.stringify().unwrap().as_bytes())?;
  419. }
  420. Err(_) => {
  421. let file = settings.pipe_path.clone();
  422. let mut pipe_write = pipe_write(file)?;
  423. let mut task_clone = task.clone();
  424. task_clone.events.push(TaskEvent::new(
  425. "add_task".to_string(),
  426. task_clone.owner.clone(),
  427. "".to_string(),
  428. ));
  429. let json: JsonValue = (&task_clone).into();
  430. pipe_write.write_all(json.stringify().unwrap().as_bytes())?;
  431. }
  432. }
  433. }
  434. task.save(&datastore_path)?;
  435. break
  436. }
  437. Ok(())
  438. }
  439. async_daemonize!(realmain);
  440. async fn realmain(settings: Args, executor: Arc<smol::Executor<'static>>) -> Result<()> {
  441. let datastore_path = expand_path(&settings.datastore)?;
  442. let nickname =
  443. if settings.nickname.is_some() { settings.nickname.clone() } else { env::var("USER").ok() };
  444. if settings.gen_rln_identity {
  445. let identity = RlnIdentity::new(&mut OsRng);
  446. let nullifier = bs58::encode(identity.nullifier.to_repr()).into_string();
  447. let trapdoor = bs58::encode(identity.trapdoor.to_repr()).into_string();
  448. // This value is part of the RLN commitment. It must match
  449. // the genesis budget used for pregenerated identities.
  450. let user_msg_limit = generated_rln_identity_user_msg_limit();
  451. println!("Generated a fresh RLN identity.\n");
  452. println!(
  453. "Current Taud registration accepts only identities whose commitments are in \
  454. the configured pregenerated set. Use this output for a genesis bundle or future \
  455. staked-registration testing; it will not register unless its commitment is \
  456. pregenerated.\n"
  457. );
  458. println!("Local account import command:\n");
  459. println!(" tau rln register <account_name> {nullifier} {trapdoor} {user_msg_limit}\n");
  460. println!(
  461. "Replace <account_name> with any local label you like (\"alice\", \"throwaway\", etc)."
  462. );
  463. println!(
  464. "Do not change user_msg_limit: it is part of the RLN commitment and must be \
  465. GENESIS_USER_MSG_LIMIT ({user_msg_limit}) for pregenerated genesis identities."
  466. );
  467. println!(
  468. "Keep the nullifier and trapdoor secret - they ARE the identity. \
  469. A `taud --gen-rln-identity` run is NOT idempotent; treat the \
  470. output like a freshly-minted password."
  471. );
  472. return Ok(())
  473. }
  474. if let Some(n_identities) = settings.gen_genesis_rln_identities {
  475. // We'll generate n_identities and hold them in a map
  476. // `k=commitment, v=(nullifier, trapdoor, used)`
  477. // We'll export the commitments to be used in the genesis event,
  478. // and the rest as a JSON file.
  479. let mut identities_map = HashMap::new();
  480. for _ in 0..n_identities {
  481. let identity = RlnIdentity::new(&mut OsRng);
  482. let commitment = identity.commitment();
  483. identities_map.insert(
  484. commitment.to_repr(),
  485. (identity.nullifier.to_repr(), identity.trapdoor.to_repr(), false),
  486. );
  487. }
  488. let mut commits = String::from(
  489. r#"
  490. use darkfi_sdk::{crypto::pasta_prelude::PrimeField, pasta::pallas};
  491. /// Return Taud's configured pregenerated RLN commitment set.
  492. pub fn pregenerated_identity_commitments() -> Vec<[u8; 32]> {
  493. TAUD_GENESIS_COMMITMENTS_REPR.to_vec()
  494. }
  495. /// Check whether an RLN commitment belongs to Taud's pregenerated set.
  496. pub fn is_pregenerated_commitment(commitment: &pallas::Base) -> bool {
  497. TAUD_GENESIS_COMMITMENTS_REPR.contains(&commitment.to_repr())
  498. }
  499. pub const TAUD_GENESIS_COMMITMENTS_REPR: &[[u8; 32]] = &[
  500. "#,
  501. );
  502. for commitment in identities_map.keys() {
  503. commits.push_str(&format!("{:?},\n", commitment));
  504. }
  505. commits.push_str("];\n");
  506. let mut file = File::create("genesis_commits.rs")?;
  507. file.write_all(commits.as_bytes())?;
  508. let mut file = File::create("taud_rln_commits.bin")?;
  509. let buf = serialize(&identities_map);
  510. file.write_all(&buf)?;
  511. return Ok(())
  512. }
  513. if settings.refresh {
  514. println!("Removing local data in: {datastore_path:?} (yes/no)? ");
  515. let mut confirm = String::new();
  516. stdin().read_line(&mut confirm).expect("Failed to read line");
  517. let confirm = confirm.to_lowercase();
  518. let confirm = confirm.trim();
  519. if confirm == "yes" || confirm == "y" {
  520. remove_dir_all(datastore_path).unwrap_or(());
  521. println!("Local data removed successfully.");
  522. } else {
  523. error!(target: "taud", "Unexpected Value: {confirm}");
  524. }
  525. return Ok(())
  526. }
  527. if nickname.is_none() {
  528. error!(target: "taud", "Provide a nickname in config file");
  529. return Ok(())
  530. }
  531. if settings.piped {
  532. let file = settings.pipe_path.clone();
  533. let path = CString::new(file).unwrap();
  534. unsafe { mkfifo(path.as_ptr(), 0o644) };
  535. }
  536. // mkdir datastore_path if not exists
  537. create_dir_all(datastore_path.clone())?;
  538. create_dir_all(datastore_path.join("month"))?;
  539. create_dir_all(datastore_path.join("task"))?;
  540. if settings.generate {
  541. println!("Generating a new workspace");
  542. loop {
  543. println!("Name for the new workspace: ");
  544. let mut workspace = String::new();
  545. stdin().read_line(&mut workspace).expect("Failed to read line");
  546. let workspace = workspace.to_lowercase();
  547. let workspace = workspace.trim();
  548. if workspace.is_empty() && workspace.len() < 3 {
  549. error!(target: "taud", "Wrong workspace try again");
  550. continue
  551. }
  552. // Encryption
  553. // Chachabox secret key (read_key) used for encrypting tasks.
  554. let secret_key = SecretKey::generate(&mut OsRng);
  555. let encoded = bs58::encode(secret_key.to_bytes());
  556. // Signature
  557. // Secret key (write_key) used for signing tasks.
  558. let keypair = Keypair::random(&mut OsRng);
  559. let sk = format!("{}", keypair.secret);
  560. // Public key (write_public_key) used for verifying tasks.
  561. let pk = format!("{}", keypair.public);
  562. println!("Please add the following to the config file:");
  563. println!("[workspace.\"{workspace}\"]");
  564. println!("read_key = \"{}\"", encoded.into_string());
  565. println!("write_key = \"{sk}\"");
  566. println!("write_public_key = \"{pk}\"");
  567. break
  568. }
  569. return Ok(())
  570. }
  571. let workspaces = Arc::new(get_workspaces(&settings).await?);
  572. let (workspace, _) = workspaces.first_key_value().unwrap();
  573. // let verified = Arc::new(Mutex::new(false));
  574. if workspaces.is_empty() {
  575. error!(target: "taud", "Please add at least one workspace to the config file.");
  576. println!("Run `$ taud --generate` to generate new workspace.");
  577. return Ok(())
  578. }
  579. info!(target: "taud", "Initializing taud node");
  580. let rln_enabled = settings.rln_enabled.unwrap_or(false);
  581. // Create datastore path if not there already.
  582. let datastore = expand_path(&settings.datastore)?;
  583. fs::create_dir_all(&datastore).await?;
  584. let zk_key_datastore = if rln_enabled {
  585. let zk_key_datastore = expand_path(&settings.zk_key_datastore)?;
  586. fs::create_dir_all(&zk_key_datastore).await?;
  587. Some(zk_key_datastore)
  588. } else {
  589. info!(target: "taud", "RLN disabled; skipping RLN key datastore setup");
  590. None
  591. };
  592. let replay_datastore = expand_path(&settings.replay_datastore)?;
  593. let replay_mode = settings.replay_mode;
  594. // let fast_mode = settings.fast_mode;
  595. info!(target: "taud", "Instantiating event DAG");
  596. let kvdb = Database::open_default(&datastore)?;
  597. let zk_key_db = if let Some(zk_key_datastore) = zk_key_datastore.as_ref() {
  598. info!(target: "taud", "Opening RLN key datastore");
  599. Some(match Database::open_default(zk_key_datastore) {
  600. Ok(v) => v,
  601. Err(e) => {
  602. error!(target: "taud", "Failed to open RLN key datastore `{zk_key_datastore:?}`: {e}");
  603. return Err(e.into());
  604. }
  605. })
  606. } else {
  607. None
  608. };
  609. let p2p_settings: darkfi::net::Settings =
  610. (env!("CARGO_PKG_NAME"), env!("CARGO_PKG_VERSION"), settings.net.clone()).try_into()?;
  611. let comms_timeout = p2p_settings.outbound_connect_timeout_max();
  612. let p2p = match P2p::new(p2p_settings, executor.clone()).await {
  613. Ok(p2p) => p2p,
  614. Err(e) => {
  615. error!("Unable to create P2P network: {e}");
  616. return Err(e);
  617. }
  618. };
  619. // Consensus config. Every node must use exactly these values.
  620. let eg_config = EventGraphConfig {
  621. initial_genesis: TAUD_INITIAL_GENESIS,
  622. hours_rotation: TAUD_HOURS_ROTATION,
  623. genesis_contents: TAUD_GENESIS_CONTENTS.to_vec(),
  624. rln_enabled,
  625. pregenerated_identity_commitments: if rln_enabled {
  626. taud::genesis_commits::pregenerated_identity_commitments()
  627. } else {
  628. Vec::new()
  629. },
  630. max_dags: Some(TAUD_MAX_DAGS),
  631. };
  632. let event_graph = match if let Some(zk_key_db) = zk_key_db.clone() {
  633. EventGraph::new_with_zk_key_db(
  634. p2p.clone(),
  635. kvdb.clone(),
  636. zk_key_db,
  637. replay_datastore.clone(),
  638. replay_mode,
  639. eg_config,
  640. executor.clone(),
  641. )
  642. .await
  643. } else {
  644. EventGraph::new(
  645. p2p.clone(),
  646. kvdb.clone(),
  647. replay_datastore.clone(),
  648. replay_mode,
  649. eg_config,
  650. executor.clone(),
  651. )
  652. .await
  653. } {
  654. Ok(v) => v,
  655. Err(e) => {
  656. error!("Event graph failed to start: {e}");
  657. return Err(e);
  658. }
  659. };
  660. // Set the active RLN account if any. When RLN is disabled, avoid
  661. // loading account state that cannot affect outbound messages.
  662. let rln_identity: Arc<smol::lock::RwLock<Option<RlnIdentity>>> =
  663. Arc::new(smol::lock::RwLock::new(if event_graph.rln_enabled() {
  664. let rln_identity = load_default_rln_identity(&kvdb).await?;
  665. if rln_identity.is_some() {
  666. info!(target: "taud", "Default RLN account set");
  667. }
  668. rln_identity
  669. } else {
  670. info!(target: "taud", "RLN disabled; skipping default RLN account load");
  671. None
  672. }));
  673. info!(target: "taud", "Registering EventGraph P2P protocol");
  674. let event_graph_ = Arc::clone(&event_graph);
  675. let registry = p2p.protocol_registry();
  676. registry
  677. .register(SESSION_DEFAULT, move |channel, _| {
  678. let event_graph_ = event_graph_.clone();
  679. async move { ProtocolEventGraph::init(event_graph_, channel).await.unwrap() }
  680. })
  681. .await;
  682. let (broadcast_snd, broadcast_rcv) = smol::channel::unbounded::<TaskInfo>();
  683. info!(target: "taud", "Starting P2P network");
  684. p2p.clone().start().await?;
  685. loop {
  686. if p2p.is_connected() {
  687. info!(target: "taud", "Got peer connection");
  688. // We'll attempt to sync for ever
  689. if !settings.skip_dag_sync {
  690. info!(target: "taud", "Syncing static DAG");
  691. match event_graph.static_sync().await {
  692. Ok(()) => info!(target: "taud", "Static synced successfully"),
  693. Err(e) => {
  694. error!(target: "taud", "Failed syncing static graph: {e}");
  695. sleep(comms_timeout).await;
  696. continue
  697. }
  698. }
  699. info!(target: "taud", "Syncing event DAG");
  700. match event_graph.sync_selected(1).await {
  701. Ok(()) => {
  702. info!(target: "taud", "Event DAG synced successfully!");
  703. break
  704. }
  705. Err(e) => {
  706. // TODO: Maybe at this point we should prune or something?
  707. // TODO: Or maybe just tell the user to delete the DAG from FS.
  708. error!(target: "taud", "Failed syncing DAG ({e}), retrying in {comms_timeout}s...");
  709. sleep(comms_timeout).await;
  710. }
  711. }
  712. } else {
  713. event_graph.synced.store(true, Ordering::Release);
  714. break
  715. }
  716. } else {
  717. info!(target: "taud", "Waiting for some P2P connections...");
  718. sleep(comms_timeout).await;
  719. }
  720. }
  721. let seen = OnceLock::new();
  722. seen.set(kvdb.open_tree_default("tau_seen").unwrap()).unwrap();
  723. ////////////////////
  724. // get history
  725. ////////////////////
  726. let dag_events = event_graph.order_events().await?;
  727. for event in dag_events.iter() {
  728. let event_id = event.header.id();
  729. // If it was seen, skip
  730. if is_seen(kvdb.clone(), seen.clone(), &event_id).await? {
  731. continue
  732. }
  733. mark_seen(kvdb.clone(), seen.clone(), &event_id).await?;
  734. // Try to deserialize it. (Here we skip errors)
  735. let Ok((enc_task, _)) = deserialize_async_partial(event.content()).await else { continue };
  736. // Potentially decrypt the privmsg
  737. on_receive_task(&enc_task, &workspaces, &settings).await.unwrap();
  738. }
  739. ////////////////////
  740. // Listner
  741. ////////////////////
  742. info!(target: "taud", "Starting sync loop task");
  743. let sync_loop_task = StoppableTask::new();
  744. sync_loop_task.clone().start(
  745. start_sync_loop(
  746. event_graph.clone(),
  747. broadcast_rcv,
  748. workspaces.clone(),
  749. kvdb.clone(),
  750. settings.clone(),
  751. p2p.clone(),
  752. seen.clone(),
  753. rln_identity.clone(),
  754. ),
  755. |res| async {
  756. match res {
  757. Ok(()) | Err(TaudError::Darkfi(Error::DetachedTaskStopped)) => { /* Do nothing */ }
  758. Err(e) => error!(target: "taud", "Failed stopping sync loop task: {e}"),
  759. }
  760. },
  761. TaudError::Darkfi(Error::DetachedTaskStopped),
  762. executor.clone(),
  763. );
  764. // ==============
  765. // p2p dnet setup
  766. // ==============
  767. info!(target: "taud", "Starting dnet subs task");
  768. let json_sub = JsonSubscriber::new("dnet.subscribe_events");
  769. let json_sub_ = json_sub.clone();
  770. let p2p_ = p2p.clone();
  771. let dnet_task = StoppableTask::new();
  772. dnet_task.clone().start(
  773. async move {
  774. let dnet_sub = p2p_.dnet_subscribe().await;
  775. loop {
  776. let event = dnet_sub.receive().await;
  777. debug!(target: "taud", "Got dnet event: {event:?}");
  778. json_sub_.notify(vec![event.into()].into()).await;
  779. }
  780. },
  781. |res| async {
  782. match res {
  783. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  784. Err(e) => {
  785. error!(target: "taud", "Failed stopping dnet subs task: {e}")
  786. }
  787. }
  788. },
  789. Error::DetachedTaskStopped,
  790. executor.clone(),
  791. );
  792. info!("Starting deg subs task");
  793. let deg_sub = JsonSubscriber::new("deg.subscribe_events");
  794. let deg_sub_ = deg_sub.clone();
  795. let event_graph_ = event_graph.clone();
  796. let deg_task = StoppableTask::new();
  797. deg_task.clone().start(
  798. async move {
  799. let deg_sub = event_graph_.deg_subscribe().await;
  800. loop {
  801. let event = deg_sub.receive().await;
  802. debug!(target: "taud", "Got deg event: {event:?}");
  803. deg_sub_.notify(vec![event.into()].into()).await;
  804. }
  805. },
  806. |res| async {
  807. match res {
  808. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  809. Err(e) => panic!("{e}"),
  810. }
  811. },
  812. Error::DetachedTaskStopped,
  813. executor.clone(),
  814. );
  815. //
  816. // RPC interface
  817. //
  818. let rpc_interface = Arc::new(JsonRpcInterface::new(
  819. datastore_path.clone(),
  820. broadcast_snd,
  821. nickname.unwrap(),
  822. workspace.to_string(),
  823. workspaces.clone(),
  824. p2p.clone(),
  825. event_graph.clone(),
  826. json_sub,
  827. deg_sub,
  828. kvdb.clone(),
  829. rln_identity.clone(),
  830. ));
  831. let rpc_task = StoppableTask::new();
  832. rpc_task.clone().start(
  833. listen_and_serve(settings.rpc.into(), rpc_interface.clone(), None, executor.clone()),
  834. |res| async move {
  835. match res {
  836. Ok(()) | Err(Error::RpcServerStopped) => rpc_interface.stop_connections().await,
  837. Err(e) => error!(target: "taud", "Failed stopping JSON-RPC server: {e}"),
  838. }
  839. },
  840. Error::RpcServerStopped,
  841. executor.clone(),
  842. );
  843. // Signal handling for graceful termination.
  844. let (signals_handler, signals_task) = SignalHandler::new(executor)?;
  845. signals_handler.wait_termination(signals_task).await?;
  846. info!(target: "taud", "Caught termination signal, cleaning up and exiting...");
  847. info!(target: "taud", "Stopping P2P network");
  848. p2p.stop().await;
  849. info!(target: "taud", "Stopping sync loop task...");
  850. sync_loop_task.stop().await;
  851. info!(target: "taud", "Stopping JSON-RPC server...");
  852. rpc_task.stop().await;
  853. dnet_task.stop().await;
  854. deg_task.stop().await;
  855. info!(target: "taud", "Flushing kvdb database...");
  856. kvdb.flush_default_mode_async().await?;
  857. if let Some(zk_key_db) = zk_key_db {
  858. info!(target: "taud", "Flushing RLN key kvdb database...");
  859. zk_key_db.flush_default_mode_async().await?;
  860. }
  861. info!(target: "taud", "Shut down successfully");
  862. Ok(())
  863. }
  864. #[cfg(test)]
  865. mod tests {
  866. use std::path::Path;
  867. use structopt::StructOpt;
  868. use darkfi::util::time::Timestamp;
  869. use taud::task_info::TaskInfo;
  870. use super::*;
  871. const TEST_DATA_PATH: &str = "/tmp/test_tau_ws_claim";
  872. /// Two workspaces that share the same read/write key, as is common in
  873. /// localnet testing where one workspace block is duplicated.
  874. fn shared_key_workspaces() -> BTreeMap<String, Workspace> {
  875. let read_key = SecretKey::generate(&mut OsRng);
  876. let chacha = ChaChaBox::new(&read_key.public_key(), &read_key);
  877. let write_key = darkfi_sdk::crypto::SecretKey::random(&mut OsRng);
  878. let write_pubkey = PublicKey::from_secret(write_key);
  879. let mut map = BTreeMap::new();
  880. map.insert(
  881. "darkfi-dev".to_string(),
  882. Workspace {
  883. read_key: ChaChaBox::new(&read_key.public_key(), &read_key),
  884. write_key: Some(write_key),
  885. write_pubkey,
  886. },
  887. );
  888. map.insert(
  889. "test".to_string(),
  890. Workspace { read_key: chacha, write_key: Some(write_key), write_pubkey },
  891. );
  892. map
  893. }
  894. #[test]
  895. fn shared_keys_task_is_claimed_by_first_workspace() -> TaudResult<()> {
  896. remove_dir_all(TEST_DATA_PATH).ok();
  897. create_dir_all(TEST_DATA_PATH).unwrap();
  898. create_dir_all(Path::new(TEST_DATA_PATH).join("task")).unwrap();
  899. create_dir_all(Path::new(TEST_DATA_PATH).join("month")).unwrap();
  900. let workspaces = shared_key_workspaces();
  901. let mut args = Args::from_iter_safe(vec!["taud".to_string()]).unwrap();
  902. args.datastore = TEST_DATA_PATH.to_string();
  903. let task = TaskInfo::new(
  904. "darkfi-dev".to_string(),
  905. "test_title",
  906. "test_desc",
  907. "NICK",
  908. None,
  909. None,
  910. Timestamp::current_time(),
  911. None,
  912. )?;
  913. let enc = encrypt_sign_task(&task, workspaces.get("darkfi-dev").unwrap())?;
  914. smol::block_on(async { on_receive_task(&enc, &workspaces, &args).await })?;
  915. // Even though both workspaces can decrypt the task, it must be
  916. // claimed exactly once and keep the originating workspace label,
  917. // not be overwritten by the later `test` workspace.
  918. let loaded = TaskInfo::load(&task.ref_id, Path::new(TEST_DATA_PATH))?;
  919. assert_eq!(loaded.workspace, "darkfi-dev");
  920. Ok(())
  921. }
  922. }