main.rs 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2023 Dyne.org foundation
  4. *
  5. * This program is free software: you can redistribute it and/or modify
  6. * it under the terms of the GNU Affero General Public License as
  7. * published by the Free Software Foundation, either version 3 of the
  8. * License, or (at your option) any later version.
  9. *
  10. * This program is distributed in the hope that it will be useful,
  11. * but WITHOUT ANY WARRANTY; without even the implied warranty of
  12. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  13. * GNU Affero General Public License for more details.
  14. *
  15. * You should have received a copy of the GNU Affero General Public License
  16. * along with this program. If not, see <https://www.gnu.org/licenses/>.
  17. */
  18. use std::{collections::HashMap, fs::create_dir_all};
  19. use async_std::{
  20. stream::StreamExt,
  21. sync::{Arc, Mutex},
  22. task,
  23. };
  24. use chrono::{Duration, Utc};
  25. use irc::ClientSubMsg;
  26. use log::{debug, error, info};
  27. use rand::rngs::OsRng;
  28. use structopt_toml::StructOptToml;
  29. use tinyjson::JsonValue;
  30. use darkfi::{
  31. async_daemonize,
  32. event_graph::{
  33. events_queue::EventsQueue,
  34. model::{Model, ModelPtr},
  35. protocol_event::{ProtocolEvent, Seen},
  36. view::View,
  37. },
  38. net,
  39. rpc::{jsonrpc::JsonSubscriber, server::listen_and_serve},
  40. system::{StoppableTask, Subscriber, SubscriberPtr},
  41. util::{async_util::sleep, file::save_json_file, path::expand_path, time::Timestamp},
  42. Error, Result,
  43. };
  44. pub mod crypto;
  45. pub mod irc;
  46. pub mod privmsg;
  47. pub mod rpc;
  48. pub mod settings;
  49. use crate::{
  50. crypto::KeyPair,
  51. irc::{IrcConfig, IrcServer},
  52. privmsg::PrivMsgEvent,
  53. rpc::JsonRpcInterface,
  54. settings::{Args, ChannelInfo, CONFIG_FILE, CONFIG_FILE_CONTENTS},
  55. };
  56. async fn parse_signals(
  57. sighup_sub: SubscriberPtr<Args>,
  58. client_sub: SubscriberPtr<ClientSubMsg>,
  59. ) -> Result<()> {
  60. debug!("Started signal parsing handler");
  61. let subscription = sighup_sub.subscribe().await;
  62. loop {
  63. let args = subscription.receive().await;
  64. let new_config = IrcConfig::new(&args)?;
  65. client_sub.notify(ClientSubMsg::Config(new_config)).await;
  66. }
  67. }
  68. // Removes events older than one week ,then sleeps untill next midnight
  69. async fn remove_old_events(model: ModelPtr<PrivMsgEvent>) -> Result<()> {
  70. loop {
  71. let now = Utc::now();
  72. // clocks are valid, safe to unwrap
  73. let next_midnight = (now + Duration::days(1)).date_naive().and_hms_opt(0, 0, 0).unwrap();
  74. let duration = next_midnight.signed_duration_since(now.naive_utc()).to_std().unwrap();
  75. let week_old_datetime =
  76. (now - Duration::weeks(1)).date_naive().and_hms_opt(0, 0, 0).unwrap();
  77. let timestamp = week_old_datetime.timestamp() as u64;
  78. model.lock().await.remove_old_events(Timestamp(timestamp))?;
  79. info!("Removing old events");
  80. sleep(duration.as_secs() + 1).await;
  81. }
  82. }
  83. async_daemonize!(realmain);
  84. async fn realmain(settings: Args, executor: Arc<smol::Executor<'_>>) -> Result<()> {
  85. let datastore_path = expand_path(&settings.datastore)?;
  86. // mkdir datastore_path if not exists
  87. create_dir_all(datastore_path.clone())?;
  88. // Signal handling for config reload and graceful termination.
  89. let (signals_handler, signals_task) = SignalHandler::new()?;
  90. let client_sub = Subscriber::new();
  91. task::spawn(parse_signals(signals_handler.sighup_sub.clone(), client_sub.clone()));
  92. ////////////////////
  93. // Generate new keypair and exit
  94. ////////////////////
  95. if settings.gen_keypair {
  96. let secret_key = crypto_box::SecretKey::generate(&mut OsRng);
  97. let public_key = secret_key.public_key();
  98. let secret = bs58::encode(secret_key.to_bytes()).into_string();
  99. let public = bs58::encode(public_key.as_bytes()).into_string();
  100. let kp = KeyPair { secret, public };
  101. if settings.output.is_some() {
  102. let datastore = expand_path(&settings.output.unwrap())?;
  103. let kp_enc = JsonValue::Object(HashMap::from([
  104. ("public".to_string(), JsonValue::String(kp.public)),
  105. ("secret".to_string(), JsonValue::String(kp.secret)),
  106. ]));
  107. save_json_file(&datastore, &kp_enc, false)?;
  108. } else {
  109. println!("Generated keypair:\n{}", kp);
  110. }
  111. return Ok(())
  112. }
  113. if settings.secret.is_some() {
  114. let secret = settings.secret.clone().unwrap();
  115. let bytes: [u8; 32] = bs58::decode(secret).into_vec()?.try_into().unwrap();
  116. let secret = crypto_box::SecretKey::from(bytes);
  117. let pubkey = secret.public_key();
  118. let pub_encoded = bs58::encode(pubkey.as_bytes()).into_string();
  119. if settings.output.is_some() {
  120. let datastore = expand_path(&settings.output.unwrap())?;
  121. save_json_file(&datastore, &JsonValue::String(pub_encoded), false)?;
  122. } else {
  123. println!("Public key recoverd: {}", pub_encoded);
  124. }
  125. return Ok(())
  126. }
  127. if settings.gen_secret {
  128. let secret_key = crypto_box::SecretKey::generate(&mut OsRng);
  129. let encoded = bs58::encode(secret_key.to_bytes());
  130. println!("{}", encoded.into_string());
  131. return Ok(())
  132. }
  133. ////////////////////
  134. // Initialize the base structures
  135. ////////////////////
  136. let events_queue = EventsQueue::<PrivMsgEvent>::new();
  137. let model = Arc::new(Mutex::new(Model::new(events_queue.clone())));
  138. let view = Arc::new(Mutex::new(View::new(events_queue.clone())));
  139. let model_clone = model.clone();
  140. let model_clone2 = model.clone();
  141. {
  142. // Temporarly load model and check if the loaded head is not
  143. // older than one week (already removed from other node's tree)
  144. let now = Utc::now();
  145. let now_datetime = (now - Duration::weeks(1)).date_naive().and_hms_opt(0, 0, 0).unwrap();
  146. let timestamp = Timestamp(now_datetime.timestamp() as u64);
  147. let mut loaded_model = Model::new(events_queue.clone());
  148. loaded_model.load_tree(&datastore_path)?;
  149. if loaded_model
  150. .get_event(&loaded_model.get_head_hash())
  151. .is_some_and(|event| event.timestamp >= timestamp)
  152. {
  153. model.lock().await.load_tree(&datastore_path)?;
  154. }
  155. }
  156. ////////////////////
  157. // P2p setup
  158. ////////////////////
  159. // Buffers
  160. let seen_event = Seen::new();
  161. let seen_inv = Seen::new();
  162. // Check the version
  163. let net_settings = settings.net.clone();
  164. // New p2p
  165. let p2p = net::P2p::new(net_settings.into()).await;
  166. // Register the protocol_event
  167. let registry = p2p.protocol_registry();
  168. registry
  169. .register(net::SESSION_ALL, move |channel, p2p| {
  170. let seen_event = seen_event.clone();
  171. let seen_inv = seen_inv.clone();
  172. let model = model.clone();
  173. async move { ProtocolEvent::init(channel, p2p, model, seen_event, seen_inv).await }
  174. })
  175. .await;
  176. // ==============
  177. // p2p dnet setup
  178. // ==============
  179. info!(target: "darkirc", "Starting dnet subs task");
  180. let json_sub = JsonSubscriber::new("dnet.subscribe_events");
  181. let json_sub_ = json_sub.clone();
  182. let p2p_ = p2p.clone();
  183. let dnet_task = StoppableTask::new();
  184. dnet_task.clone().start(
  185. async move {
  186. let dnet_sub = p2p_.dnet_subscribe().await;
  187. loop {
  188. let event = dnet_sub.receive().await;
  189. debug!("Got dnet event: {:?}", event);
  190. json_sub_.notify(vec![event.into()]).await;
  191. }
  192. },
  193. |res| async {
  194. match res {
  195. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  196. Err(e) => {
  197. error!(target: "darkirc", "Failed starting remove old events task: {}", e)
  198. }
  199. }
  200. },
  201. Error::DetachedTaskStopped,
  202. executor.clone(),
  203. );
  204. ////////////////////
  205. // RPC interface setup
  206. ////////////////////
  207. let rpc_listen_addr = settings.rpc_listen.clone();
  208. info!(target: "darkirc", "Starting JSON-RPC server on {}", rpc_listen_addr);
  209. let rpc_interface = Arc::new(JsonRpcInterface {
  210. addr: rpc_listen_addr.clone(),
  211. p2p: p2p.clone(),
  212. dnet_sub: json_sub,
  213. });
  214. let rpc_task = StoppableTask::new();
  215. rpc_task.clone().start(
  216. listen_and_serve(rpc_listen_addr, rpc_interface, executor.clone()),
  217. |res| async {
  218. match res {
  219. Ok(()) | Err(Error::RPCServerStopped) => { /* Do nothing */ }
  220. Err(e) => error!(target: "darkirc", "Failed starting JSON-RPC server: {}", e),
  221. }
  222. },
  223. Error::RPCServerStopped,
  224. executor.clone(),
  225. );
  226. ////////////////////
  227. // Start P2P network
  228. ////////////////////
  229. info!(target: "darkirc", "Starting P2P network");
  230. p2p.clone().start(executor.clone()).await?;
  231. StoppableTask::new().start(
  232. p2p.clone().run(executor.clone()),
  233. |res| async {
  234. match res {
  235. Ok(()) | Err(Error::P2PNetworkStopped) => { /* Do nothing */ }
  236. Err(e) => error!(target: "darkirc", "Failed starting P2P network: {}", e),
  237. }
  238. },
  239. Error::P2PNetworkStopped,
  240. executor.clone(),
  241. );
  242. ////////////////////
  243. // IRC server
  244. ////////////////////
  245. info!(target: "darkirc", "Starting IRC server");
  246. let irc_server = IrcServer::new(
  247. settings.clone(),
  248. p2p.clone(),
  249. model_clone.clone(),
  250. view.clone(),
  251. client_sub,
  252. )
  253. .await?;
  254. let irc_server_task = StoppableTask::new();
  255. let executor_ = executor.clone();
  256. irc_server_task.clone().start(
  257. // Weird hack to prevent lifetimes hell
  258. async move { irc_server.start(executor_).await },
  259. |res| async {
  260. match res {
  261. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  262. Err(e) => error!(target: "darkirc", "Failed starting IRC server: {}", e),
  263. }
  264. },
  265. Error::DetachedTaskStopped,
  266. executor.clone(),
  267. );
  268. // Reset root task
  269. info!(target: "darkirc", "Starting remove old events task");
  270. let remove_old_events_task = StoppableTask::new();
  271. remove_old_events_task.clone().start(
  272. remove_old_events(model_clone2),
  273. |res| async {
  274. match res {
  275. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  276. Err(e) => {
  277. error!(target: "darkirc", "Failed starting remove old events task: {}", e)
  278. }
  279. }
  280. },
  281. Error::DetachedTaskStopped,
  282. executor,
  283. );
  284. // Wait for termination signal
  285. signals_handler.wait_termination(signals_task).await?;
  286. info!("Caught termination signal, cleaning up and exiting...");
  287. model_clone.lock().await.save_tree(&datastore_path)?;
  288. info!(target: "darkirc", "Stopping dnet subs task...");
  289. dnet_task.stop().await;
  290. info!(target: "darkirc", "Stopping JSON-RPC server...");
  291. rpc_task.stop().await;
  292. info!(target: "darkirc", "Stopping P2P network");
  293. p2p.stop().await;
  294. info!(target: "darkirc", "Stopping IRC server...");
  295. irc_server_task.stop().await;
  296. info!(target: "darkirc", "Stopping remove old events task...");
  297. remove_old_events_task.stop().await;
  298. Ok(())
  299. }