lib.rs 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326
  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::{HashMap, HashSet},
  20. sync::Arc,
  21. };
  22. use smol::lock::Mutex;
  23. use tracing::{debug, error, info};
  24. use darkfi::{
  25. net::settings::Settings,
  26. rpc::{
  27. jsonrpc::JsonSubscriber,
  28. server::{listen_and_serve, RequestHandler},
  29. settings::RpcSettings,
  30. },
  31. system::{ExecutorPtr, StoppableTask, StoppableTaskPtr},
  32. validator::{Validator, ValidatorConfig, ValidatorPtr},
  33. Error, Result,
  34. };
  35. use darkfi_sdk::crypto::keypair::Network;
  36. #[cfg(test)]
  37. mod tests;
  38. mod error;
  39. use error::{server_error, RpcError};
  40. /// JSON-RPC requests handler and methods
  41. mod rpc;
  42. use rpc::{management::ManagementRpcHandler, DefaultRpcHandler};
  43. /// Validator async tasks
  44. pub mod task;
  45. use task::{consensus::ConsensusInitTaskConfig, consensus_init_task, garbage_collect_task};
  46. /// P2P net protocols
  47. mod proto;
  48. use proto::{DarkfidP2pHandler, DarkfidP2pHandlerPtr};
  49. /// Miners registry
  50. mod registry;
  51. use registry::{DarkfiMinersRegistry, DarkfiMinersRegistryPtr};
  52. /// Atomic pointer to the DarkFi node
  53. pub type DarkfiNodePtr = Arc<DarkfiNode>;
  54. /// Structure representing a DarkFi node
  55. pub struct DarkfiNode {
  56. /// Validator(node) pointer
  57. validator: ValidatorPtr,
  58. /// P2P network protocols handler
  59. p2p_handler: DarkfidP2pHandlerPtr,
  60. /// Node miners registry pointer
  61. registry: DarkfiMinersRegistryPtr,
  62. /// A map of various subscribers exporting live info from the blockchain
  63. subscribers: HashMap<&'static str, JsonSubscriber>,
  64. /// Main JSON-RPC connection tracker
  65. rpc_connections: Mutex<HashSet<StoppableTaskPtr>>,
  66. /// Management JSON-RPC connection tracker
  67. management_rpc_connections: Mutex<HashSet<StoppableTaskPtr>>,
  68. }
  69. impl DarkfiNode {
  70. pub async fn new(
  71. validator: ValidatorPtr,
  72. p2p_handler: DarkfidP2pHandlerPtr,
  73. registry: DarkfiMinersRegistryPtr,
  74. subscribers: HashMap<&'static str, JsonSubscriber>,
  75. ) -> Result<DarkfiNodePtr> {
  76. Ok(Arc::new(Self {
  77. validator,
  78. p2p_handler,
  79. registry,
  80. subscribers,
  81. rpc_connections: Mutex::new(HashSet::new()),
  82. management_rpc_connections: Mutex::new(HashSet::new()),
  83. }))
  84. }
  85. }
  86. /// Atomic pointer to the DarkFi daemon
  87. pub type DarkfidPtr = Arc<Darkfid>;
  88. /// Structure representing a DarkFi daemon
  89. pub struct Darkfid {
  90. /// Darkfi node instance
  91. node: DarkfiNodePtr,
  92. /// `dnet` background task
  93. dnet_task: StoppableTaskPtr,
  94. /// Main JSON-RPC background task
  95. rpc_task: StoppableTaskPtr,
  96. /// Management JSON-RPC background task
  97. management_rpc_task: StoppableTaskPtr,
  98. /// Consensus protocol background task
  99. consensus_task: StoppableTaskPtr,
  100. /// Node garbage collection background task
  101. gc_task: StoppableTaskPtr,
  102. }
  103. impl Darkfid {
  104. /// Initialize a DarkFi daemon.
  105. ///
  106. /// Generates a new `DarkfiNode` for provided configuration,
  107. /// along with all the corresponding background tasks.
  108. pub async fn init(
  109. network: Network,
  110. sled_db: &sled_overlay::sled::Db,
  111. config: &ValidatorConfig,
  112. net_settings: &Settings,
  113. ex: &ExecutorPtr,
  114. ) -> Result<DarkfidPtr> {
  115. info!(target: "darkfid::Darkfid::init", "Initializing a Darkfi daemon...");
  116. // Initialize validator
  117. let validator = Validator::new(sled_db, config).await?;
  118. // Initialize P2P network
  119. let p2p_handler = DarkfidP2pHandler::init(net_settings, ex).await?;
  120. // Initialize the miners registry
  121. let registry = DarkfiMinersRegistry::init(network, &validator).await?;
  122. // Here we initialize various subscribers that can export live blockchain/consensus data.
  123. let mut subscribers = HashMap::new();
  124. subscribers.insert("blocks", JsonSubscriber::new("blockchain.subscribe_blocks"));
  125. subscribers.insert("txs", JsonSubscriber::new("blockchain.subscribe_txs"));
  126. subscribers.insert("proposals", JsonSubscriber::new("blockchain.subscribe_proposals"));
  127. subscribers.insert("dnet", JsonSubscriber::new("dnet.subscribe_events"));
  128. // Initialize node
  129. let node = DarkfiNode::new(validator, p2p_handler, registry, subscribers).await?;
  130. // Generate the background tasks
  131. let dnet_task = StoppableTask::new();
  132. let rpc_task = StoppableTask::new();
  133. let management_rpc_task = StoppableTask::new();
  134. let consensus_task = StoppableTask::new();
  135. let gc_task = StoppableTask::new();
  136. info!(target: "darkfid::Darkfid::init", "Darkfi daemon initialized successfully!");
  137. Ok(Arc::new(Self {
  138. node,
  139. dnet_task,
  140. rpc_task,
  141. management_rpc_task,
  142. consensus_task,
  143. gc_task,
  144. }))
  145. }
  146. /// Start the DarkFi daemon in the given executor, using the
  147. /// provided JSON-RPC settings and consensus initialization
  148. /// configuration.
  149. pub async fn start(
  150. &self,
  151. executor: &ExecutorPtr,
  152. rpc_settings: &RpcSettings,
  153. management_rpc_settings: &RpcSettings,
  154. stratum_rpc_settings: &Option<RpcSettings>,
  155. mm_rpc_settings: &Option<RpcSettings>,
  156. config: &ConsensusInitTaskConfig,
  157. ) -> Result<()> {
  158. info!(target: "darkfid::Darkfid::start", "Starting Darkfi daemon...");
  159. // Start the `dnet` task
  160. info!(target: "darkfid::Darkfid::start", "Starting dnet subs task");
  161. let dnet_sub_ = self.node.subscribers.get("dnet").unwrap().clone();
  162. let p2p_ = self.node.p2p_handler.p2p.clone();
  163. self.dnet_task.clone().start(
  164. async move {
  165. let dnet_sub = p2p_.dnet_subscribe().await;
  166. loop {
  167. let event = dnet_sub.receive().await;
  168. debug!(target: "darkfid::Darkfid::dnet_task", "Got dnet event: {event:?}");
  169. dnet_sub_.notify(vec![event.into()].into()).await;
  170. }
  171. },
  172. |res| async {
  173. match res {
  174. Ok(()) | Err(Error::DetachedTaskStopped) => { /* Do nothing */ }
  175. Err(e) => error!(target: "darkfid::Darkfid::start", "Failed starting dnet subs task: {e}"),
  176. }
  177. },
  178. Error::DetachedTaskStopped,
  179. executor.clone(),
  180. );
  181. // Start the main JSON-RPC task
  182. info!(target: "darkfid::Darkfid::start", "Starting main JSON-RPC server");
  183. let node_ = self.node.clone();
  184. self.rpc_task.clone().start(
  185. listen_and_serve::<DefaultRpcHandler>(rpc_settings.clone(), self.node.clone(), None, executor.clone()),
  186. |res| async move {
  187. match res {
  188. Ok(()) | Err(Error::RpcServerStopped) => <DarkfiNode as RequestHandler<DefaultRpcHandler>>::stop_connections(&node_).await,
  189. Err(e) => error!(target: "darkfid::Darkfid::start", "Failed starting main JSON-RPC server: {e}"),
  190. }
  191. },
  192. Error::RpcServerStopped,
  193. executor.clone(),
  194. );
  195. // Start the management JSON-RPC task
  196. info!(target: "darkfid::Darkfid::start", "Starting management JSON-RPC server");
  197. let node_ = self.node.clone();
  198. self.management_rpc_task.clone().start(
  199. listen_and_serve::<ManagementRpcHandler>(management_rpc_settings.clone(), self.node.clone(), None, executor.clone()),
  200. |res| async move {
  201. match res {
  202. Ok(()) | Err(Error::RpcServerStopped) => <DarkfiNode as RequestHandler<ManagementRpcHandler>>::stop_connections(&node_).await,
  203. Err(e) => error!(target: "darkfid::Darkfid::start", "Failed starting management JSON-RPC server: {e}"),
  204. }
  205. },
  206. Error::RpcServerStopped,
  207. executor.clone(),
  208. );
  209. // Start the miners registry
  210. info!(target: "darkfid::Darkfid::start", "Starting miners registry");
  211. self.node.registry.start(executor, &self.node, stratum_rpc_settings, mm_rpc_settings)?;
  212. // Start the P2P network
  213. info!(target: "darkfid::Darkfid::start", "Starting P2P network");
  214. self.node.p2p_handler.start(executor, &self.node).await?;
  215. // Generate the signal queue smol channel
  216. let (sender, receiver) = smol::channel::unbounded::<()>();
  217. // Start the consensus protocol
  218. info!(target: "darkfid::Darkfid::start", "Starting consensus protocol task");
  219. self.consensus_task.clone().start(
  220. consensus_init_task(
  221. self.node.clone(),
  222. config.clone(),
  223. sender,
  224. ),
  225. |res| async move {
  226. match res {
  227. Ok(()) | Err(Error::ConsensusTaskStopped) | Err(Error::MinerTaskStopped) => { /* Do nothing */ }
  228. Err(e) => error!(target: "darkfid::Darkfid::start", "Failed starting consensus initialization task: {e}"),
  229. }
  230. },
  231. Error::ConsensusTaskStopped,
  232. executor.clone(),
  233. );
  234. // Start the garbage collection task
  235. info!(target: "darkfid::Darkfid::start", "Starting garbage collection task");
  236. self.gc_task.clone().start(
  237. garbage_collect_task(receiver, self.node.clone()),
  238. |res| async {
  239. match res {
  240. Ok(()) | Err(Error::GarbageCollectionTaskStopped) => { /* Do nothing */ }
  241. Err(e) => {
  242. error!(target: "darkfid", "Failed starting garbage collection task: {e}")
  243. }
  244. }
  245. },
  246. Error::GarbageCollectionTaskStopped,
  247. executor.clone(),
  248. );
  249. info!(target: "darkfid::Darkfid::start", "Darkfi daemon started successfully!");
  250. Ok(())
  251. }
  252. /// Stop the DarkFi daemon.
  253. pub async fn stop(&self) -> Result<()> {
  254. info!(target: "darkfid::Darkfid::stop", "Terminating Darkfi daemon...");
  255. // Stop the `dnet` node
  256. info!(target: "darkfid::Darkfid::stop", "Stopping dnet subs task...");
  257. self.dnet_task.stop().await;
  258. // Stop the main JSON-RPC task
  259. info!(target: "darkfid::Darkfid::stop", "Stopping main JSON-RPC server...");
  260. self.rpc_task.stop().await;
  261. // Stop the management JSON-RPC task
  262. info!(target: "darkfid::Darkfid::stop", "Stopping management JSON-RPC server...");
  263. self.management_rpc_task.stop().await;
  264. // Stop the miners registry
  265. info!(target: "darkfid::Darkfid::stop", "Stopping miners registry...");
  266. self.node.registry.stop().await;
  267. // Stop the P2P network
  268. info!(target: "darkfid::Darkfid::stop", "Stopping P2P network protocols handler...");
  269. self.node.p2p_handler.stop().await;
  270. // Stop the garbage collection task
  271. info!(target: "darkfid::Darkfid::stop", "Stopping garbage collection task...");
  272. self.gc_task.stop().await;
  273. // Stop the consensus task
  274. info!(target: "darkfid::Darkfid::stop", "Stopping consensus task...");
  275. self.consensus_task.stop().await;
  276. // Flush sled database data
  277. info!(target: "darkfid::Darkfid::stop", "Flushing sled database...");
  278. let flushed_bytes =
  279. self.node.validator.read().await.blockchain.sled_db.flush_async().await?;
  280. info!(target: "darkfid::Darkfid::stop", "Flushed {flushed_bytes} bytes");
  281. info!(target: "darkfid::Darkfid::stop", "Darkfi daemon terminated successfully!");
  282. Ok(())
  283. }
  284. }