utility.rs 5.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181
  1. use std::fs::OpenOptions;
  2. use std::io::prelude::*;
  3. use std::net::SocketAddr;
  4. use std::path::PathBuf;
  5. use std::sync::atomic::AtomicU64;
  6. use std::sync::Arc;
  7. use std::time::{SystemTime, UNIX_EPOCH};
  8. use async_executor::Executor;
  9. use futures::prelude::*;
  10. use log::*;
  11. use super::{aes::aes_encrypt, dbsql, net::messages, net::protocol_slab::ProtocolSlab,
  12. control_message::ControlMessage, control_message::ControlCommand, channel::Channel,
  13. slabs_manager::SlabsManagerSafe, control_message::MessagePayload};
  14. use crate::{ serial::{serialize, deserialize}, net::ChannelPtr, Result};
  15. pub type AddrsStorage = Arc<async_std::sync::Mutex<Vec<SocketAddr>>>;
  16. pub type Clock = 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 path = default_config_dir()?.join("addrs.add");
  28. let mut writer = OpenOptions::new().write(true).create(true).open(path)?;
  29. let buffer = serialize(stored_addrs);
  30. writer.write_all(&buffer)?;
  31. Ok(())
  32. }
  33. pub fn default_config_dir() -> Result<PathBuf> {
  34. let mut path = PathBuf::new();
  35. if let Some(home_dir) = dirs::home_dir() {
  36. path = home_dir;
  37. };
  38. let path = path.join(".darkpulse/");
  39. if !path.exists() {
  40. match std::fs::create_dir(&path) {
  41. Err(err) => {
  42. eprintln!("error: Creating config dir: {}", err);
  43. std::process::exit(-1);
  44. }
  45. Ok(()) => (),
  46. }
  47. }
  48. Ok(path)
  49. }
  50. pub fn load_stored_addrs() -> Result<Vec<SocketAddr>> {
  51. let path = default_config_dir()?.join("addrs.add");
  52. println!("{:?}", path);
  53. let mut reader = OpenOptions::new()
  54. .read(true)
  55. .write(true)
  56. .create(true)
  57. .open(path)?;
  58. let mut buffer = Vec::new();
  59. reader.read_to_end(&mut buffer)?;
  60. if !buffer.is_empty() {
  61. let addrs: Vec<SocketAddr> = deserialize(&buffer)?;
  62. Ok(addrs)
  63. } else {
  64. Ok(vec![])
  65. }
  66. }
  67. pub async fn pack_slab(
  68. channel_secret: &[u8; 32],
  69. username: String,
  70. message: String,
  71. control_command: ControlCommand,
  72. ) -> Result<messages::SlabMessage> {
  73. let nonce: [u8; 12] = rand::random();
  74. let timestamp = chrono::offset::Utc::now();
  75. let timestamp: i64 = timestamp.timestamp_millis() / 1000;
  76. let msg_payload = MessagePayload {
  77. nickname: username.clone(),
  78. text: message,
  79. timestamp,
  80. };
  81. let control_message = ControlMessage {
  82. control: control_command,
  83. payload: msg_payload,
  84. };
  85. let ser_message = serialize(&control_message);
  86. let ciphertext = aes_encrypt(channel_secret, &nonce, &ser_message[..])
  87. .expect("error during encrypting the message");
  88. let slab = messages::SlabMessage { nonce, ciphertext };
  89. Ok(slab)
  90. }
  91. pub fn setup_username(newname: Option<String>, db: &dbsql::Dbsql) -> Result<String> {
  92. let mut _username: String = String::new();
  93. match newname {
  94. Some(nm) => {
  95. _username = nm.clone();
  96. db.add_username(&nm).unwrap();
  97. }
  98. None => {
  99. _username = db.get_username()?;
  100. if _username.is_empty() {
  101. _username = String::from("username");
  102. }
  103. }
  104. }
  105. Ok(_username)
  106. }
  107. pub async fn read_line<R: AsyncBufRead + Unpin>(reader: &mut R) -> Result<String> {
  108. let mut buf = String::new();
  109. let _ = reader.read_line(&mut buf).await?;
  110. Ok(buf.trim().to_string())
  111. }
  112. pub fn choose_channel(db: &dbsql::Dbsql, channel_name: Option<String>) -> Result<Channel> {
  113. let channels = db.get_channels()?;
  114. let mut main_channel = Channel::gen_new(String::from("test_channel"));
  115. if channels.len() > 0 {
  116. match channel_name {
  117. Some(name) => {
  118. main_channel = channels
  119. .iter()
  120. .filter(|ch| ch.get_channel_name() == &name)
  121. .next()
  122. .expect(format!("there is no channel with the name {}: ", name).as_str())
  123. .clone();
  124. }
  125. None => {
  126. main_channel = channels.first().unwrap().clone();
  127. }
  128. }
  129. } else {
  130. error!("there are no channels available");
  131. db.add_channel(&main_channel)?;
  132. }
  133. Ok(main_channel)
  134. }
  135. pub async fn setup_network_channel(
  136. executor: Arc<Executor<'_>>,
  137. channel: ChannelPtr,
  138. slabman: SlabsManagerSafe,
  139. ) {
  140. let message_subsytem = channel.get_message_subsystem();
  141. message_subsytem
  142. .add_dispatch::<messages::SyncMessage>()
  143. .await;
  144. message_subsytem
  145. .add_dispatch::<messages::InvMessage>()
  146. .await;
  147. message_subsytem
  148. .add_dispatch::<messages::GetSlabsMessage>()
  149. .await;
  150. message_subsytem
  151. .add_dispatch::<messages::SlabMessage>()
  152. .await;
  153. let protocol_slab = ProtocolSlab::new(slabman, channel.clone()).await;
  154. protocol_slab.clone().start(executor.clone()).await;
  155. }