main.rs 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393
  1. use std::net::SocketAddr;
  2. use async_executor::Executor;
  3. use async_std::sync::{Arc, Mutex};
  4. use async_trait::async_trait;
  5. use easy_parallel::Parallel;
  6. use futures_lite::future;
  7. use lazy_init::Lazy;
  8. use log::{error, info};
  9. use rand::Rng;
  10. use serde_derive::Deserialize;
  11. use simplelog::{ColorChoice, TermLogger, TerminalMode};
  12. use structopt::StructOpt;
  13. use structopt_toml::StructOptToml;
  14. use url::Url;
  15. use darkfi::{
  16. async_daemonize,
  17. blockchain::{NullifierStore, RootStore},
  18. cli_desc,
  19. consensus2::{
  20. proto::{
  21. ProtocolParticipant, ProtocolProposal, ProtocolSync, ProtocolSyncConsensus, ProtocolTx,
  22. ProtocolVote,
  23. },
  24. state::ValidatorStatePtr,
  25. task::{block_sync_task, proposal_task},
  26. util::Timestamp,
  27. ValidatorState, MAINNET_GENESIS_HASH_BYTES, TESTNET_GENESIS_HASH_BYTES,
  28. },
  29. net,
  30. net::P2pPtr,
  31. node::{Client, State},
  32. rpc::{
  33. jsonrpc,
  34. jsonrpc::{
  35. ErrorCode::{InvalidParams, MethodNotFound},
  36. JsonRequest, JsonResult,
  37. },
  38. rpcserver2::{listen_and_serve, RequestHandler},
  39. },
  40. util::{
  41. cli::{log_config, spawn_config},
  42. expand_path,
  43. path::get_config_path,
  44. },
  45. wallet::walletdb::{init_wallet, WalletPtr},
  46. Error, Result,
  47. };
  48. mod error;
  49. use error::{server_error, RpcError};
  50. const CONFIG_FILE: &str = "darkfid_config.toml";
  51. const CONFIG_FILE_CONTENTS: &str = include_str!("../darkfid_config.toml");
  52. #[derive(Clone, Debug, Deserialize, StructOpt, StructOptToml)]
  53. #[serde(default)]
  54. #[structopt(name = "darkfid", about = cli_desc!())]
  55. struct Args {
  56. #[structopt(short, long)]
  57. /// Configuration file to use
  58. config: Option<String>,
  59. #[structopt(long, default_value = "testnet")]
  60. /// Chain to use (testnet, mainnet)
  61. chain: String,
  62. #[structopt(long)]
  63. /// Participate in consensus
  64. consensus: bool,
  65. #[structopt(long, default_value = "~/.config/darkfi/darkfid_wallet.db")]
  66. /// Path to wallet database
  67. wallet_path: String,
  68. #[structopt(long, default_value = "changeme")]
  69. /// Password for the wallet database
  70. wallet_pass: String,
  71. #[structopt(long, default_value = "~/.config/darkfi/darkfid_blockchain")]
  72. /// Path to blockchain database
  73. database: String,
  74. #[structopt(long, default_value = "tcp://127.0.0.1:5397")]
  75. /// JSON-RPC listen URL
  76. rpc_listen: Url,
  77. #[structopt(long)]
  78. /// P2P accept address for the consensus protocol
  79. consensus_p2p_accept: Option<SocketAddr>,
  80. #[structopt(long)]
  81. /// P2P external address for the consensus protocol
  82. consensus_p2p_external: Option<SocketAddr>,
  83. #[structopt(long, default_value = "8")]
  84. /// Connection slots for the consensus protocol
  85. consensus_slots: u32,
  86. #[structopt(long)]
  87. /// Connect to peer for the consensus protocol (repeatable flag)
  88. consensus_peer: Vec<SocketAddr>,
  89. #[structopt(long)]
  90. /// Connect to seed for the consensus protocol (repeatable flag)
  91. consensus_seed: Vec<SocketAddr>,
  92. #[structopt(long)]
  93. /// P2P accept address for the syncing protocol
  94. sync_p2p_accept: Option<SocketAddr>,
  95. #[structopt(long)]
  96. /// P2P external address for the syncing protocol
  97. sync_p2p_external: Option<SocketAddr>,
  98. #[structopt(long, default_value = "8")]
  99. /// Connection slots for the syncing protocol
  100. sync_slots: u32,
  101. #[structopt(long)]
  102. /// Connect to peer for the syncing protocol (repeatable flag)
  103. sync_peer: Vec<SocketAddr>,
  104. #[structopt(long)]
  105. /// Connect to seed for the syncing protocol (repeatable flag)
  106. sync_seed: Vec<SocketAddr>,
  107. #[structopt(short, parse(from_occurrences))]
  108. /// Increase verbosity (-vvv supported)
  109. verbose: u8,
  110. #[structopt(short)]
  111. /// Genesis time
  112. genesis_time: i64,
  113. }
  114. pub struct Darkfid {
  115. synced: Mutex<bool>, // AtomicBool is weird in Arc
  116. client: Client,
  117. consensus_p2p: Option<P2pPtr>,
  118. sync_p2p: Option<P2pPtr>,
  119. validator_state: ValidatorStatePtr,
  120. state: Arc<Mutex<State>>,
  121. }
  122. // JSON-RPC methods
  123. mod rpc_blockchain;
  124. mod rpc_misc;
  125. mod rpc_wallet;
  126. #[async_trait]
  127. impl RequestHandler for Darkfid {
  128. async fn handle_request(&self, req: JsonRequest) -> JsonResult {
  129. if !req.params.is_array() {
  130. return jsonrpc::error(InvalidParams, None, req.id).into()
  131. }
  132. let params = req.params.as_array().unwrap();
  133. match req.method.as_str() {
  134. Some("ping") => return self.pong(req.id, params).await,
  135. Some("blockchain.get_slot") => return self.get_slot(req.id, params).await,
  136. Some("wallet.keygen") => return self.keygen(req.id, params).await,
  137. Some("wallet.get_key") => return self.get_key(req.id, params).await,
  138. Some("wallet.export_keypair") => return self.export_keypair(req.id, params).await,
  139. Some("wallet.import_keypair") => return self.import_keypair(req.id, params).await,
  140. Some("wallet.set_default_address") => {
  141. return self.set_default_address(req.id, params).await
  142. }
  143. Some(_) | None => return jsonrpc::error(MethodNotFound, None, req.id).into(),
  144. }
  145. }
  146. }
  147. impl Darkfid {
  148. pub async fn new(
  149. db: &sled::Db,
  150. wallet: WalletPtr,
  151. validator_state: ValidatorStatePtr,
  152. consensus_p2p: Option<P2pPtr>,
  153. sync_p2p: Option<P2pPtr>,
  154. ) -> Result<Self> {
  155. // Initialize Client
  156. let client = Client::new(wallet).await?;
  157. let tree = client.get_tree().await?;
  158. let merkle_roots = RootStore::new(db)?;
  159. let nullifiers = NullifierStore::new(db)?;
  160. // Initialize State
  161. let state = Arc::new(Mutex::new(State {
  162. tree,
  163. merkle_roots,
  164. nullifiers,
  165. cashier_pubkeys: vec![],
  166. faucet_pubkeys: vec![],
  167. mint_vk: Lazy::new(),
  168. burn_vk: Lazy::new(),
  169. }));
  170. Ok(Self {
  171. synced: Mutex::new(false),
  172. client,
  173. consensus_p2p,
  174. sync_p2p,
  175. validator_state,
  176. state,
  177. })
  178. }
  179. }
  180. async_daemonize!(realmain);
  181. async fn realmain(args: Args, ex: Arc<Executor<'_>>) -> Result<()> {
  182. // We use this handler to block this function after detaching all
  183. // tasks, and to catch a shutdown signal, where we can clean up and
  184. // exit gracefully.
  185. let (signal, shutdown) = async_channel::bounded::<()>(1);
  186. ctrlc_async::set_async_handler(async move {
  187. signal.send(()).await.unwrap();
  188. })
  189. .unwrap();
  190. // Initialize or load wallet
  191. let wallet = init_wallet(&args.wallet_path, &args.wallet_pass).await?;
  192. // Initialize or open sled database
  193. let db_path = format!("{}/{}", expand_path(&args.database)?.to_str().unwrap(), args.chain);
  194. let sled_db = sled::open(&db_path)?;
  195. // Initialize validator state
  196. // TODO: genesis_ts should be some hardcoded constant
  197. let genesis_ts = Timestamp(args.genesis_time);
  198. let genesis_data = match args.chain.as_str() {
  199. "mainnet" => *MAINNET_GENESIS_HASH_BYTES,
  200. "testnet" => *TESTNET_GENESIS_HASH_BYTES,
  201. x => {
  202. error!("Unsupported chain `{}`", x);
  203. return Err(Error::UnsupportedChain)
  204. }
  205. };
  206. // TODO: Is this ok?
  207. let mut rng = rand::thread_rng();
  208. let id: u64 = rng.gen();
  209. // Initialize validator state
  210. let state = ValidatorState::new(&sled_db, id, genesis_ts, genesis_data)?;
  211. let sync_p2p = {
  212. info!("Registering block sync P2P protocols...");
  213. let sync_network_settings = net::Settings {
  214. inbound: args.sync_p2p_accept,
  215. outbound_connections: args.sync_slots,
  216. external_addr: args.sync_p2p_external,
  217. peers: args.sync_peer.clone(),
  218. seeds: args.sync_seed.clone(),
  219. ..Default::default()
  220. };
  221. let p2p = net::P2p::new(sync_network_settings).await;
  222. let registry = p2p.protocol_registry();
  223. let _state = state.clone();
  224. registry
  225. .register(net::SESSION_ALL, move |channel, p2p| {
  226. let state = _state.clone();
  227. async move { ProtocolSync::init(channel, state, p2p, args.consensus).await.unwrap() }
  228. })
  229. .await;
  230. let _state = state.clone();
  231. registry
  232. .register(net::SESSION_ALL, move |channel, p2p| {
  233. let state = _state.clone();
  234. async move { ProtocolTx::init(channel, state, p2p).await.unwrap() }
  235. })
  236. .await;
  237. Some(p2p)
  238. };
  239. // P2P network settings for the consensus protocol
  240. let consensus_p2p = {
  241. if !args.consensus {
  242. None
  243. } else {
  244. info!("Registering consensus P2P protocols...");
  245. let consensus_network_settings = net::Settings {
  246. inbound: args.consensus_p2p_accept,
  247. outbound_connections: args.consensus_slots,
  248. external_addr: args.consensus_p2p_external,
  249. peers: args.consensus_peer.clone(),
  250. seeds: args.consensus_seed.clone(),
  251. ..Default::default()
  252. };
  253. let p2p = net::P2p::new(consensus_network_settings).await;
  254. let registry = p2p.protocol_registry();
  255. let _state = state.clone();
  256. registry
  257. .register(net::SESSION_ALL, move |channel, p2p| {
  258. let state = _state.clone();
  259. async move { ProtocolParticipant::init(channel, state, p2p).await.unwrap() }
  260. })
  261. .await;
  262. let _state = state.clone();
  263. registry
  264. .register(net::SESSION_ALL, move |channel, p2p| {
  265. let state = _state.clone();
  266. async move { ProtocolProposal::init(channel, state, p2p).await.unwrap() }
  267. })
  268. .await;
  269. let _state = state.clone();
  270. let _sync_p2p = sync_p2p.clone().unwrap();
  271. registry
  272. .register(net::SESSION_ALL, move |channel, p2p| {
  273. let state = _state.clone();
  274. let __sync_p2p = _sync_p2p.clone();
  275. async move {
  276. ProtocolVote::init(channel, state, __sync_p2p, p2p)
  277. .await
  278. .unwrap()
  279. }
  280. })
  281. .await;
  282. let _state = state.clone();
  283. registry
  284. .register(net::SESSION_ALL, move |channel, p2p| {
  285. let state = _state.clone();
  286. async move { ProtocolSyncConsensus::init(channel, state, p2p).await.unwrap() }
  287. })
  288. .await;
  289. Some(p2p)
  290. }
  291. };
  292. // Initialize program state
  293. let darkfid =
  294. Darkfid::new(&sled_db, wallet, state.clone(), consensus_p2p.clone(), sync_p2p.clone())
  295. .await?;
  296. let darkfid = Arc::new(darkfid);
  297. // JSON-RPC server
  298. info!("Starting JSON-RPC server");
  299. ex.spawn(listen_and_serve(args.rpc_listen, darkfid.clone())).detach();
  300. info!("Starting sync P2P network");
  301. sync_p2p.clone().unwrap().start(ex.clone()).await?;
  302. let _ex = ex.clone();
  303. let _sync_p2p = sync_p2p.clone();
  304. ex.spawn(async move {
  305. if let Err(e) = _sync_p2p.unwrap().run(_ex).await {
  306. error!("Failed starting sync P2P network: {}", e);
  307. }
  308. })
  309. .detach();
  310. match block_sync_task(sync_p2p.clone().unwrap(), state.clone()).await {
  311. Ok(()) => *darkfid.synced.lock().await = true,
  312. Err(e) => error!("Failed syncing blockchain: {}", e),
  313. }
  314. // Consensus protocol
  315. if args.consensus {
  316. info!("Starting consensus P2P network");
  317. consensus_p2p.clone().unwrap().start(ex.clone()).await?;
  318. let _ex = ex.clone();
  319. let _consensus_p2p = consensus_p2p.clone();
  320. ex.spawn(async move {
  321. if let Err(e) = _consensus_p2p.unwrap().run(_ex).await {
  322. error!("Failed starting consensus P2P network: {}", e);
  323. }
  324. })
  325. .detach();
  326. info!("Starting consensus protocol task");
  327. ex.spawn(proposal_task(consensus_p2p.unwrap(), state)).detach();
  328. }
  329. // Wait for SIGINT
  330. shutdown.recv().await?;
  331. print!("\r");
  332. info!("Caught termination signal, cleaning up and exiting...");
  333. info!("Flushing database...");
  334. let flushed_bytes = sled_db.flush_async().await?;
  335. info!("Flushed {} bytes", flushed_bytes);
  336. Ok(())
  337. }