darkfid.rs 4.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141
  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 subscriber = GatewayClient::start_subscriber(sub_addr).await?;
  23. let slabstore = client.get_slabstore();
  24. let _ = executor.spawn(GatewayClient::subscribe(subscriber, slabstore));
  25. // TEST
  26. let _slab = Slab::new("testcoin".to_string(), vec![0, 0, 0, 0]);
  27. client.put_slab(_slab).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. }
  64. // $ cargo test --bin darkfid
  65. // run 10 clients simultaneously
  66. #[cfg(test)]
  67. mod test {
  68. #[test]
  69. fn test_darkfid_client() {
  70. use std::path::Path;
  71. use drk::service::GatewayClient;
  72. use drk::slab::Slab;
  73. use log::*;
  74. use rand::Rng;
  75. use simplelog::*;
  76. let logger_config = ConfigBuilder::new().set_time_format_str("%T%.6f").build();
  77. CombinedLogger::init(vec![
  78. TermLogger::new(LevelFilter::Debug, logger_config, TerminalMode::Mixed).unwrap(),
  79. WriteLogger::new(
  80. LevelFilter::Debug,
  81. Config::default(),
  82. std::fs::File::create(Path::new("/tmp/dar.log")).unwrap(),
  83. ),
  84. ])
  85. .unwrap();
  86. let mut thread_pools: Vec<std::thread::JoinHandle<()>> = vec![];
  87. for _ in 1..11 {
  88. let thread = std::thread::spawn(|| {
  89. smol::future::block_on(async move {
  90. let mut rng = rand::thread_rng();
  91. let rnd: u32 = rng.gen();
  92. let mut client = GatewayClient::new(
  93. "127.0.0.1:3333".parse().unwrap(),
  94. Path::new(&format!("slabstore_{}.db", rnd)),
  95. )
  96. .unwrap();
  97. client.start().await.unwrap();
  98. let _slab = Slab::new("testcoin".to_string(), rnd.to_le_bytes().to_vec());
  99. client.put_slab(_slab).await.unwrap();
  100. std::thread::sleep(std::time::Duration::from_secs(3));
  101. let last_index = client.slabstore.get_last_index().unwrap();
  102. info!("last index: {}", last_index);
  103. })
  104. });
  105. thread_pools.push(thread);
  106. }
  107. for t in thread_pools {
  108. t.join().unwrap();
  109. }
  110. }
  111. }