darkfid.rs 2.3 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980
  1. use async_executor::Executor;
  2. use async_std::sync::Arc;
  3. use easy_parallel::Parallel;
  4. use std::net::SocketAddr;
  5. use drk::service::{ClientProgramOptions, GatewayClient};
  6. use drk::{slab::Slab, Result};
  7. fn setup_addr(address: Option<SocketAddr>, default: SocketAddr) -> SocketAddr {
  8. match address {
  9. Some(addr) => addr,
  10. None => default,
  11. }
  12. }
  13. async fn start(executor: Arc<Executor<'_>>, options: ClientProgramOptions) -> Result<()> {
  14. let connect_addr: SocketAddr = setup_addr(options.connect_addr, "127.0.0.1:3333".parse()?);
  15. let sub_addr: SocketAddr = setup_addr(options.sub_addr, "127.0.0.1:4444".parse()?);
  16. let slabstore_path = options.slabstore_path.as_path();
  17. // create gateway client
  18. let mut client = GatewayClient::new(connect_addr, slabstore_path)?;
  19. // start gateway client
  20. client.start().await?;
  21. // start subscribe to gateway publisher
  22. let slabstore = client.get_slabstore();
  23. let subscriber_task = executor.spawn(GatewayClient::subscribe(slabstore, sub_addr));
  24. // TEST
  25. let _slab = Slab::new("testcoin".to_string(), vec![0, 0, 0, 0]);
  26. // client.put_slab(_slab).await?;
  27. subscriber_task.cancel().await;
  28. Ok(())
  29. }
  30. fn main() -> Result<()> {
  31. use simplelog::*;
  32. let ex = Arc::new(Executor::new());
  33. let (signal, shutdown) = async_channel::unbounded::<()>();
  34. let options = ClientProgramOptions::load()?;
  35. let logger_config = ConfigBuilder::new().set_time_format_str("%T%.6f").build();
  36. let debug_level = if options.verbose {
  37. LevelFilter::Debug
  38. } else {
  39. LevelFilter::Off
  40. };
  41. CombinedLogger::init(vec![
  42. TermLogger::new(debug_level, logger_config, TerminalMode::Mixed).unwrap(),
  43. WriteLogger::new(
  44. LevelFilter::Debug,
  45. Config::default(),
  46. std::fs::File::create(options.log_path.as_path()).unwrap(),
  47. ),
  48. ])
  49. .unwrap();
  50. let ex2 = ex.clone();
  51. let (_, result) = Parallel::new()
  52. // Run four executor threads.
  53. .each(0..3, |_| smol::future::block_on(ex.run(shutdown.recv())))
  54. // Run the main future on the current thread.
  55. .finish(|| {
  56. smol::future::block_on(async move {
  57. start(ex2, options).await?;
  58. drop(signal);
  59. Ok::<(), drk::Error>(())
  60. })
  61. });
  62. result
  63. }