main.rs 14 KB

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