lib.rs 12 KB

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