darkd.rs 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385
  1. #[macro_use]
  2. extern crate clap;
  3. use async_executor::Executor;
  4. use async_native_tls::TlsAcceptor;
  5. use async_std::sync::Mutex;
  6. use easy_parallel::Parallel;
  7. use ff::Field;
  8. use http_types::{Request, Response, StatusCode};
  9. use log::*;
  10. use rand::rngs::OsRng;
  11. use sapvi::serial;
  12. use serde_json::json;
  13. use smol::Async;
  14. use std::net::SocketAddr;
  15. use std::net::TcpListener;
  16. use std::sync::Arc;
  17. use sapvi::{net, Result};
  18. /// Listens for incoming connections and serves them.
  19. async fn listen(
  20. executor: Arc<Executor<'_>>,
  21. rpc: Arc<RpcInterface>,
  22. listener: Async<TcpListener>,
  23. tls: Option<TlsAcceptor>,
  24. ) -> Result<()> {
  25. // Format the full host address.
  26. let host = match &tls {
  27. None => format!("http://{}", listener.get_ref().local_addr()?),
  28. Some(_) => format!("https://{}", listener.get_ref().local_addr()?),
  29. };
  30. println!("Listening on {}", host);
  31. loop {
  32. // Accept the next connection.
  33. let (stream, _) = listener.accept().await?;
  34. // Spawn a background task serving this connection.
  35. let task = match &tls {
  36. None => {
  37. let stream = async_dup::Arc::new(stream);
  38. let rpc = rpc.clone();
  39. executor.spawn(async move {
  40. if let Err(err) = async_h1::accept(stream, move |req| {
  41. let rpc = rpc.clone();
  42. rpc.serve(req)
  43. })
  44. .await
  45. {
  46. println!("Connection error: {:#?}", err);
  47. }
  48. })
  49. }
  50. Some(tls) => {
  51. // In case of HTTPS, establish a secure TLS connection first.
  52. match tls.accept(stream).await {
  53. Ok(stream) => {
  54. let _stream = async_dup::Arc::new(async_dup::Mutex::new(stream));
  55. executor.spawn(async move {
  56. /*if let Err(err) = async_h1::accept(stream, serve).await {
  57. println!("Connection error: {:#?}", err);
  58. }*/
  59. unimplemented!();
  60. })
  61. }
  62. Err(err) => {
  63. println!("Failed to establish secure TLS connection: {:#?}", err);
  64. continue;
  65. }
  66. }
  67. }
  68. };
  69. // Detach the task to let it run in the background.
  70. task.detach();
  71. }
  72. }
  73. struct RpcInterface {
  74. p2p: Arc<net::P2p>,
  75. started: Mutex<bool>,
  76. stop_send: async_channel::Sender<()>,
  77. stop_recv: async_channel::Receiver<()>,
  78. }
  79. impl RpcInterface {
  80. fn new(p2p: Arc<net::P2p>) -> Arc<Self> {
  81. let (stop_send, stop_recv) = async_channel::unbounded::<()>();
  82. Arc::new(Self {
  83. p2p,
  84. started: Mutex::new(false),
  85. stop_send,
  86. stop_recv,
  87. })
  88. }
  89. // add new methods to handle wallet commands
  90. async fn serve(self: Arc<Self>, mut req: Request) -> http_types::Result<Response> {
  91. info!("RPC serving {}", req.url());
  92. let request = req.body_string().await?;
  93. let mut io = jsonrpc_core::IoHandler::new();
  94. io.add_sync_method("say_hello", |_| {
  95. Ok(jsonrpc_core::Value::String("Hello World!".into()))
  96. });
  97. let self2 = self.clone();
  98. io.add_method("get_info", move |_| {
  99. let self2 = self2.clone();
  100. async move {
  101. Ok(json!({
  102. "started": *self2.started.lock().await,
  103. "connections": self2.p2p.connections_count().await
  104. }))
  105. }
  106. });
  107. let stop_send = self.stop_send.clone();
  108. io.add_method("stop", move |_| {
  109. let stop_send = stop_send.clone();
  110. async move {
  111. let _ = stop_send.send(()).await;
  112. Ok(jsonrpc_core::Value::Null)
  113. }
  114. });
  115. io.add_method("key_gen", move |_| {
  116. //let stop_send = stop_send.clone();
  117. async move {
  118. //let _ = stop_send.send(()).await;
  119. let secret: jubjub::Fr = jubjub::Fr::random(&mut OsRng);
  120. let public = zcash_primitives::constants::SPENDING_KEY_GENERATOR * secret;
  121. let pubkey = serial::serialize(&public);
  122. let privkey = serial::serialize(&secret);
  123. //println!("{:?}", pubkey);
  124. //println!("{:?}", privkey);
  125. Ok(jsonrpc_core::Value::Null)
  126. }
  127. });
  128. let response = io
  129. .handle_request_sync(&request)
  130. .ok_or(sapvi::Error::BadOperationType)?;
  131. let mut res = Response::new(StatusCode::Ok);
  132. res.insert_header("Content-Type", "text/plain");
  133. res.set_body(response);
  134. Ok(res)
  135. }
  136. async fn wait_for_quit(self: Arc<Self>) -> Result<()> {
  137. Ok(self.stop_recv.recv().await?)
  138. }
  139. }
  140. async fn start(executor: Arc<Executor<'_>>, options: ProgramOptions) -> Result<()> {
  141. let p2p = net::P2p::new(options.network_settings);
  142. let rpc = RpcInterface::new(p2p.clone());
  143. let http = listen(
  144. executor.clone(),
  145. rpc.clone(),
  146. Async::<TcpListener>::bind(([127, 0, 0, 1], options.rpc_port))?,
  147. None,
  148. );
  149. let http_task = executor.spawn(http);
  150. *rpc.started.lock().await = true;
  151. p2p.clone().start(executor.clone()).await?;
  152. p2p.run(executor).await?;
  153. rpc.wait_for_quit().await?;
  154. http_task.cancel().await;
  155. Ok(())
  156. }
  157. /*
  158. async fn start2(executor: Arc<Executor<'_>>, options: ProgramOptions) -> Result<()> {
  159. let connections = Arc::new(Mutex::new(HashMap::new()));
  160. let stored_addrs = Arc::new(Mutex::new(Vec::new()));
  161. let executor2 = executor.clone();
  162. let stored_addrs2 = stored_addrs.clone();
  163. let mut server_task = None;
  164. if let Some(accept_addr) = options.accept_addr {
  165. let accept_addr = accept_addr.clone();
  166. let protocol = ServerProtocol::new(connections.clone(), accept_addr, stored_addrs2);
  167. server_task = Some(executor.spawn(async move {
  168. protocol.start(executor2).await?;
  169. Ok::<(), sapvi::Error>(())
  170. }));
  171. }
  172. let mut seed_protocols = Vec::with_capacity(options.seed_addrs.len());
  173. // Normally we query this from a server
  174. let accept_addr = options.accept_addr.clone();
  175. for seed_addr in options.seed_addrs.iter() {
  176. let protocol = SeedProtocol::new(seed_addr.clone(), accept_addr, stored_addrs.clone());
  177. protocol.clone().start(executor.clone()).await;
  178. seed_protocols.push(protocol);
  179. }
  180. debug!("Waiting for seed node queries to finish...");
  181. for seed_protocol in seed_protocols {
  182. seed_protocol.await_finish().await;
  183. }
  184. debug!("Seed nodes queried.");
  185. let mut client_slots = vec![];
  186. for i in 0..options.connection_slots {
  187. debug!("Starting connection slot {}", i);
  188. let client = Channel::new(
  189. connections.clone(),
  190. accept_addr.clone(),
  191. stored_addrs.clone(),
  192. );
  193. client.clone().start(executor.clone()).await;
  194. client_slots.push(client);
  195. }
  196. for remote_addr in options.manual_connects {
  197. debug!("Starting connection (manual) to {}", remote_addr);
  198. let client = Channel::new(
  199. connections.clone(),
  200. accept_addr.clone(),
  201. stored_addrs.clone(),
  202. );
  203. client
  204. .clone()
  205. .start_manual(remote_addr, executor.clone())
  206. .await;
  207. client_slots.push(client);
  208. }
  209. let rpc = RpcInterface::new();
  210. let http = listen(
  211. executor.clone(),
  212. rpc.clone(),
  213. Async::<TcpListener>::bind(([127, 0, 0, 1], 8000))?,
  214. None,
  215. );
  216. let http_task = executor.spawn(http);
  217. rpc.stop_recv.recv().await?;
  218. http_task.cancel().await;
  219. match server_task {
  220. None => {}
  221. Some(server_task) => {
  222. server_task.cancel().await;
  223. }
  224. }
  225. Ok(())
  226. }
  227. */
  228. struct ProgramOptions {
  229. network_settings: net::Settings,
  230. log_path: Box<std::path::PathBuf>,
  231. rpc_port: u16,
  232. }
  233. impl ProgramOptions {
  234. fn load() -> Result<ProgramOptions> {
  235. let app = clap_app!(dfi =>
  236. (version: "0.1.0")
  237. (author: "Amir Taaki <amir@dyne.org>")
  238. (about: "Dark node")
  239. (@arg ACCEPT: -a --accept +takes_value "Accept address")
  240. (@arg SEED_NODES: -s --seeds ... "Seed nodes")
  241. (@arg CONNECTS: -c --connect ... "Manual connections")
  242. (@arg CONNECT_SLOTS: --slots +takes_value "Connection slots")
  243. (@arg LOG_PATH: --log +takes_value "Logfile path")
  244. (@arg RPC_PORT: -r --rpc +takes_value "RPC port")
  245. )
  246. .get_matches();
  247. let accept_addr = if let Some(accept_addr) = app.value_of("ACCEPT") {
  248. Some(accept_addr.parse()?)
  249. } else {
  250. None
  251. };
  252. let mut seed_addrs: Vec<SocketAddr> = vec![];
  253. if let Some(seeds) = app.values_of("SEED_NODES") {
  254. for seed in seeds {
  255. seed_addrs.push(seed.parse()?);
  256. }
  257. }
  258. let mut manual_connects: Vec<SocketAddr> = vec![];
  259. if let Some(connections) = app.values_of("CONNECTS") {
  260. for connect in connections {
  261. manual_connects.push(connect.parse()?);
  262. }
  263. }
  264. let connection_slots = if let Some(connection_slots) = app.value_of("CONNECT_SLOTS") {
  265. connection_slots.parse()?
  266. } else {
  267. 0
  268. };
  269. let log_path = Box::new(
  270. if let Some(log_path) = app.value_of("LOG_PATH") {
  271. std::path::Path::new(log_path)
  272. } else {
  273. std::path::Path::new("/tmp/darkfid.log")
  274. }
  275. .to_path_buf(),
  276. );
  277. let rpc_port = if let Some(rpc_port) = app.value_of("RPC_PORT") {
  278. rpc_port.parse()?
  279. } else {
  280. 8000
  281. };
  282. Ok(ProgramOptions {
  283. network_settings: net::Settings {
  284. inbound: accept_addr,
  285. outbound_connections: connection_slots,
  286. external_addr: accept_addr,
  287. peers: manual_connects,
  288. seeds: seed_addrs,
  289. ..Default::default()
  290. },
  291. log_path,
  292. rpc_port,
  293. })
  294. }
  295. }
  296. fn main() -> Result<()> {
  297. use simplelog::*;
  298. let options = ProgramOptions::load()?;
  299. let logger_config = ConfigBuilder::new().set_time_format_str("%T%.6f").build();
  300. CombinedLogger::init(vec![
  301. TermLogger::new(LevelFilter::Debug, logger_config, TerminalMode::Mixed).unwrap(),
  302. WriteLogger::new(
  303. LevelFilter::Debug,
  304. Config::default(),
  305. std::fs::File::create(options.log_path.as_path()).unwrap(),
  306. ),
  307. ])
  308. .unwrap();
  309. let ex = Arc::new(Executor::new());
  310. let (signal, shutdown) = async_channel::unbounded::<()>();
  311. let ex2 = ex.clone();
  312. let (_, result) = Parallel::new()
  313. // Run four executor threads.
  314. .each(0..3, |_| smol::future::block_on(ex.run(shutdown.recv())))
  315. // Run the main future on the current thread.
  316. .finish(|| {
  317. smol::future::block_on(async move {
  318. start(ex2, options).await?;
  319. drop(signal);
  320. Ok::<(), sapvi::Error>(())
  321. })
  322. });
  323. result
  324. }