gatewayd.rs 3.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113
  1. use std::net::SocketAddr;
  2. use std::str;
  3. use std::sync::Arc;
  4. use std::fs::OpenOptions;
  5. use std::io::Read;
  6. use std::{fs, path::PathBuf};
  7. use toml;
  8. use drk::blockchain::{rocks::columns, Rocks, RocksColumn};
  9. use drk::cli::{ServiceCli, GatewaydConfig};
  10. use drk::service::GatewayService;
  11. use drk::util::join_config_path;
  12. use drk::Result;
  13. extern crate clap;
  14. use async_executor::Executor;
  15. use easy_parallel::Parallel;
  16. async fn start(executor: Arc<Executor<'_>>, config: Arc<&GatewaydConfig>) -> Result<()> {
  17. let accept_addr: SocketAddr = config.accept_url.parse()?;
  18. let pub_addr: SocketAddr = config.publisher_url.parse()?;
  19. let database_path = config.database_path.clone();
  20. let database_path = join_config_path(&PathBuf::from(database_path))?;
  21. let rocks = Rocks::new(&database_path)?;
  22. let rocks_slabstore_column = RocksColumn::<columns::Slabs>::new(rocks);
  23. let gateway = GatewayService::new(accept_addr, pub_addr, rocks_slabstore_column)?;
  24. gateway.start(executor.clone()).await?;
  25. Ok(())
  26. }
  27. fn set_default() -> Result<GatewaydConfig> {
  28. let config_file = GatewaydConfig {
  29. accept_url: String::from("127.0.0.1:3333"),
  30. publisher_url: String::from("127.0.0.1:4444"),
  31. database_path: String::from("gatewayd.db"),
  32. log_path: String::from("/tmp/gatewayd.log"),
  33. };
  34. Ok(config_file)
  35. }
  36. fn main() -> Result<()> {
  37. use simplelog::*;
  38. let ex = Arc::new(Executor::new());
  39. let (signal, shutdown) = async_channel::unbounded::<()>();
  40. let config_path = PathBuf::from("gatewayd.toml");
  41. let path = join_config_path(&config_path).unwrap();
  42. let mut file = OpenOptions::new()
  43. .read(true)
  44. .write(true)
  45. .create(true)
  46. .open(&path)?;
  47. let mut buffer: Vec<u8> = vec![];
  48. file.read_to_end(&mut buffer)?;
  49. if buffer.is_empty() {
  50. // set the default setting
  51. let config_file = set_default()?;
  52. let config_file = toml::to_string(&config_file)?;
  53. fs::write(&path, &config_file)?;
  54. }
  55. // reload the config
  56. let toml = fs::read(&path)?;
  57. let str_buff = str::from_utf8(&toml)?;
  58. // read from config file
  59. let config: GatewaydConfig = toml::from_str(str_buff)?;
  60. let config_pointer = Arc::new(&config);
  61. let options = ServiceCli::load()?;
  62. let logger_config = ConfigBuilder::new().set_time_format_str("%T%.6f").build();
  63. let debug_level = if options.verbose {
  64. LevelFilter::Debug
  65. } else {
  66. LevelFilter::Off
  67. };
  68. let log_path = config.log_path.clone();
  69. CombinedLogger::init(vec![
  70. TermLogger::new(debug_level, logger_config, TerminalMode::Mixed).unwrap(),
  71. WriteLogger::new(
  72. LevelFilter::Debug,
  73. Config::default(),
  74. std::fs::File::create(log_path).unwrap(),
  75. ),
  76. ])
  77. .unwrap();
  78. let ex2 = ex.clone();
  79. let (_, result) = Parallel::new()
  80. // Run four executor threads.
  81. .each(0..3, |_| smol::future::block_on(ex.run(shutdown.recv())))
  82. // Run the main future on the current thread.
  83. .finish(|| {
  84. smol::future::block_on(async move {
  85. start(ex2, config_pointer).await?;
  86. drop(signal);
  87. Ok::<(), drk::Error>(())
  88. })
  89. });
  90. result
  91. }