utility.rs 2.8 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889
  1. use std::collections::HashMap;
  2. use std::fs::OpenOptions;
  3. use std::io::prelude::*;
  4. use std::net::SocketAddr;
  5. use std::sync::atomic::AtomicU64;
  6. use std::sync::Arc;
  7. use std::time::{SystemTime, UNIX_EPOCH};
  8. use rand::seq::SliceRandom;
  9. use smol::{Executor, Task};
  10. //use crate::{net, serial, Channel, ClientProtocol, Result, SlabsManagerSafe};
  11. use crate::{net::messages as net, serial, Result};
  12. pub type ConnectionsMap = std::sync::Arc<
  13. async_std::sync::Mutex<HashMap<SocketAddr, async_channel::Sender<net::Message>>>,
  14. >;
  15. pub type AddrsStorage = std::sync::Arc<async_std::sync::Mutex<Vec<SocketAddr>>>;
  16. pub type Clock = std::sync::Arc<AtomicU64>;
  17. pub fn get_current_time() -> u64 {
  18. let start = SystemTime::now();
  19. let since_the_epoch = start
  20. .duration_since(UNIX_EPOCH)
  21. .expect("Incorrect system clock: time went backwards");
  22. let in_ms =
  23. since_the_epoch.as_secs() * 1000 + since_the_epoch.subsec_nanos() as u64 / 1_000_000;
  24. return in_ms;
  25. }
  26. pub fn save_to_addrs_store(stored_addrs: &Vec<SocketAddr>) -> Result<()> {
  27. let mut writer = OpenOptions::new()
  28. .write(true)
  29. .create(true)
  30. .open("addrs.dps")?;
  31. let buffer = serial::serialize(stored_addrs);
  32. writer.write_all(&buffer)?;
  33. Ok(())
  34. }
  35. pub fn load_stored_addrs() -> Result<Vec<SocketAddr>> {
  36. let mut reader = OpenOptions::new()
  37. .read(true)
  38. .write(true)
  39. .create(true)
  40. .open("addrs.dps")?;
  41. let mut buffer = Vec::new();
  42. reader.read_to_end(&mut buffer)?;
  43. if !buffer.is_empty() {
  44. let addrs: Vec<SocketAddr> = serial::deserialize(&buffer)?;
  45. Ok(addrs)
  46. } else {
  47. Ok(vec![])
  48. }
  49. }
  50. pub async fn start_connections_process(
  51. //slabman: SlabsManagerSafe,
  52. stored_addrs: Vec<SocketAddr>,
  53. connections: ConnectionsMap,
  54. _accept_addr: SocketAddr,
  55. _channel_secret: [u8; 32],
  56. executor: Arc<Executor<'_>>,
  57. ) -> Vec<Task<()>> {
  58. let mut tasks: Vec<Task<()>> = vec![];
  59. for _ in 0..10 {
  60. let connections_cloned = connections.clone();
  61. let stored_addrs_cloned = stored_addrs.clone();
  62. //let slabman_cloned = slabman.clone();
  63. //let channel_secret = channel_secret.clone();
  64. let task = executor.spawn(async move {
  65. loop {
  66. let addr = stored_addrs_cloned.choose(&mut rand::thread_rng()).unwrap();
  67. if !connections_cloned.lock().await.contains_key(addr) {
  68. /*let mut protocol =
  69. ClientProtocol::new(connections_cloned.clone(), slabman_cloned.clone());
  70. protocol
  71. .start(addr.clone(), accept_addr.clone(), &channel_secret)
  72. .await;
  73. */
  74. }
  75. }
  76. });
  77. tasks.push(task);
  78. }
  79. tasks
  80. }