Pārlūkot izejas kodu

darkpulse: cargo fmt.

parazyd 4 gadi atpakaļ
vecāks
revīzija
b94a55ef4e

+ 14 - 22
src/bin/darkpulse.rs

@@ -1,14 +1,16 @@
-use async_std::sync::{Arc, Mutex};
 use async_executor::Executor;
+use async_std::sync::{Arc, Mutex};
 use easy_parallel::Parallel;
 use log::*;
 use smol::Unblock;
 
-use drk::darkpulse::{
-    dbsql, messages, utility, CiphertextHash, CliOption, ControlCommand, MemPool,
-    SlabsManager,
+use drk::{
+    darkpulse::{
+        dbsql, messages, utility, CiphertextHash, CliOption, ControlCommand, MemPool, SlabsManager,
+    },
+    net::P2p,
+    Result,
 };
-use drk::{Result, net::P2p};
 
 async fn on_receive_slab(
     p2p: Arc<P2p>,
@@ -16,11 +18,8 @@ async fn on_receive_slab(
 ) -> Result<()> {
     loop {
         let slab = slab_rx.recv().await?;
-        p2p.broadcast(messages::InvMessage {
-            slabs_hash: vec![slab],
-        })
-        .await?;
-        }
+        p2p.broadcast(messages::InvMessage { slabs_hash: vec![slab] }).await?;
+    }
 }
 
 async fn start(executor: Arc<Executor<'_>>, options: CliOption, db: dbsql::Dbsql) -> Result<()> {
@@ -34,10 +33,7 @@ async fn start(executor: Arc<Executor<'_>>, options: CliOption, db: dbsql::Dbsql
 
     // choose a channel
     if let Some(new_channel) = options.new_channel {
-        info!(
-            "channel added with the name {}",
-            new_channel.get_channel_name()
-        );
+        info!("channel added with the name {}", new_channel.get_channel_name());
         db.add_channel(&new_channel).unwrap();
     }
     let main_channel = utility::choose_channel(&db, options.channel_name)?;
@@ -56,7 +52,7 @@ async fn start(executor: Arc<Executor<'_>>, options: CliOption, db: dbsql::Dbsql
             let network_channel = subscribtion.receive().await.unwrap();
             utility::setup_network_channel(executor2.clone(), network_channel, slabman.clone())
                 .await;
-            }
+        }
     });
 
     let receive_slab = executor.spawn(on_receive_slab(p2p.clone(), slab_rx.clone()));
@@ -83,7 +79,7 @@ async fn start(executor: Arc<Executor<'_>>, options: CliOption, db: dbsql::Dbsql
                     String::from("Hello"),
                     ControlCommand::Message,
                 )
-                    .await?;
+                .await?;
 
                 p2p.broadcast(slab).await?;
             }
@@ -120,11 +116,7 @@ pub fn main() -> Result<()> {
 
     let cli_option = CliOption::get()?;
 
-    let debug_level = if cli_option.verbose {
-        LevelFilter::Debug
-    } else {
-        LevelFilter::Off
-    };
+    let debug_level = if cli_option.verbose { LevelFilter::Debug } else { LevelFilter::Off };
 
     CombinedLogger::init(vec![
         TermLogger::new(debug_level, Config::default(), TerminalMode::Mixed, ColorChoice::Always),
@@ -134,7 +126,7 @@ pub fn main() -> Result<()> {
             std::fs::File::create("/tmp/darkpulsenode.log").unwrap(),
         ),
     ])
-        .unwrap();
+    .unwrap();
 
     let mut db = dbsql::Dbsql::new()?;
     db.start()?;

+ 4 - 3
src/darkpulse/aes.rs

@@ -1,6 +1,7 @@
-
-use aes_gcm::aead::{generic_array::GenericArray, Aead, NewAead};
-use aes_gcm::Aes256Gcm;
+use aes_gcm::{
+    aead::{generic_array::GenericArray, Aead, NewAead},
+    Aes256Gcm,
+};
 
 pub type AesKey = [u8; 32];
 pub type Plaintext = Vec<u8>;

+ 4 - 23
src/darkpulse/channel.rs

@@ -1,4 +1,3 @@
-
 use bs58;
 use rand::Rng;
 use sha2::{Digest, Sha256};
@@ -20,35 +19,20 @@ impl Channel {
         address: String,
         id: u32,
     ) -> Channel {
-        Channel {
-            channel_secret,
-            channel_name,
-            address,
-            id: Some(id),
-        }
+        Channel { channel_secret, channel_name, address, id: Some(id) }
     }
 
     pub fn gen_new(channel_name: String) -> Channel {
         let channel_secret = rand::thread_rng().gen::<[u8; 32]>();
         let address = Self::gen_address(channel_secret);
-        Channel {
-            channel_secret,
-            channel_name,
-            address,
-            id: None,
-        }
+        Channel { channel_secret, channel_name, address, id: None }
     }
 
     pub fn gen_new_with_addr(channel_name: String, channel_address: String) -> Result<Channel> {
         let mut decoded = bs58::decode(channel_address.clone()).into_vec()?;
         let mut channel_secret: [u8; 32] = [0; 32];
         channel_secret.copy_from_slice(&mut decoded[4..36]);
-        Ok(Channel {
-            channel_secret,
-            channel_name,
-            address: channel_address,
-            id: None,
-        })
+        Ok(Channel { channel_secret, channel_name, address: channel_address, id: None })
     }
 
     pub fn gen_address(channel_secret: [u8; 32]) -> String {
@@ -99,10 +83,7 @@ mod tests {
         let channel_address = channel.get_channel_address();
         let channel2 = Channel::gen_new_with_addr(String::from("test"), channel_address.clone())?;
         assert_eq!(channel.get_channel_secret(), channel2.get_channel_secret());
-        assert_eq!(
-            channel.get_channel_address(),
-            channel2.get_channel_address()
-        );
+        assert_eq!(channel.get_channel_address(), channel2.get_channel_address());
         Ok(())
     }
 }

+ 3 - 18
src/darkpulse/cli_option.rs

@@ -28,12 +28,7 @@ impl CliOption {
                     .long("accept")
                     .required(true),
             )
-            .arg(
-                Arg::with_name("slots")
-                    .value_name("SLOTS")
-                    .long("slots")
-                    .help("Connection slots"),
-            )
+            .arg(Arg::with_name("slots").value_name("SLOTS").long("slots").help("Connection slots"))
             .arg(
                 Arg::with_name("verbose")
                     .takes_value(false)
@@ -142,10 +137,7 @@ impl CliOption {
             if let Some(channeladdress) = newch.value_of("address") {
                 new_channel_address = String::from(channeladdress);
             }
-            new_channel = Some(Channel::gen_new_with_addr(
-                new_channel_name,
-                new_channel_address,
-            )?);
+            new_channel = Some(Channel::gen_new_with_addr(new_channel_name, new_channel_address)?);
         }
 
         let verbose = mat.is_present("verbose");
@@ -171,14 +163,7 @@ impl CliOption {
             seeds: seed_addresses,
         };
 
-        Ok(CliOption {
-            network_settings,
-            username,
-            channel_name,
-            new_channel,
-            verbose,
-            log_path,
-        })
+        Ok(CliOption { network_settings, username, channel_name, new_channel, verbose, log_path })
     }
 
     fn collect_addrs(addrs: Vec<&str>) -> Vec<SocketAddr> {

+ 5 - 5
src/darkpulse/control_message.rs

@@ -1,6 +1,9 @@
 use std::io;
 
-use crate::{serial::{Decodable, Encodable}, Result};
+use crate::{
+    serial::{Decodable, Encodable},
+    Result,
+};
 
 #[derive(Copy, Clone)]
 pub enum ControlCommand {
@@ -56,9 +59,6 @@ impl Decodable for ControlMessage {
             1 => ControlCommand::Leave,
             _ => ControlCommand::Message,
         };
-        Ok(Self {
-            control,
-            payload: Decodable::decode(&mut d)?,
-        })
+        Ok(Self { control, payload: Decodable::decode(&mut d)? })
     }
 }

+ 8 - 24
src/darkpulse/dbsql.rs

@@ -1,13 +1,9 @@
-use std::collections::HashMap;
-use std::convert::TryInto;
-use std::fs::File;
-use std::io::prelude::*;
+use std::{collections::HashMap, convert::TryInto, fs::File, io::prelude::*};
 
 use rusqlite::{params, Connection};
 
+use super::{utility::default_config_dir, Channel, CiphertextHash, SlabMessage};
 use crate::Result;
-use super::{SlabMessage, CiphertextHash,utility::default_config_dir,
-Channel};
 
 #[derive(Debug)]
 pub struct Dbsql {
@@ -20,10 +16,7 @@ impl Dbsql {
         let path = default_config_dir()?.join("data.db");
         let connection = Connection::open(path)?;
         let username = String::new();
-        Ok(Dbsql {
-            connection,
-            username,
-        })
+        Ok(Dbsql { connection, username })
     }
 
     pub fn start(&mut self) -> Result<()> {
@@ -42,10 +35,8 @@ impl Dbsql {
     }
 
     pub fn add_username(&self, username: &String) -> Result<()> {
-        self.connection.execute(
-            "INSERT OR IGNORE INTO node (username) VALUES (?1)",
-            params![username],
-        )?;
+        self.connection
+            .execute("INSERT OR IGNORE INTO node (username) VALUES (?1)", params![username])?;
         Ok(())
     }
 
@@ -86,10 +77,8 @@ impl Dbsql {
     }
 
     pub fn delete_channel(&self, channel_name: &String) -> Result<()> {
-        self.connection.execute(
-            "DELETE FROM channel WHERE channel_name = (?1)",
-            params![channel_name,],
-        )?;
+        self.connection
+            .execute("DELETE FROM channel WHERE channel_name = (?1)", params![channel_name,])?;
         Ok(())
     }
 
@@ -138,12 +127,7 @@ impl Dbsql {
                 .expect("error when converting vector to slice with size [u8; 32]");
 
             let address = row.get(3)?;
-            Ok(Channel::new(
-                    channel_name,
-                    channel_secret,
-                    address,
-                    channel_id,
-            ))
+            Ok(Channel::new(channel_name, channel_secret, address, channel_id))
         })?;
 
         for channel in channel_iter {

+ 3 - 5
src/darkpulse/mod.rs

@@ -12,12 +12,10 @@ use async_std::sync::{Arc, Mutex};
 pub type CiphertextHash = [u8; 32];
 pub type MemPool = Arc<Mutex<Vec<(CiphertextHash, net::messages::SlabMessage)>>>;
 
-
-pub use aes::{aes_decrypt, Ciphertext, Plaintext, aes_encrypt};
+pub use aes::{aes_decrypt, aes_encrypt, Ciphertext, Plaintext};
 pub use channel::Channel;
 pub use cli_option::CliOption;
-pub use control_message::{ControlMessage, ControlCommand, MessagePayload};
+pub use control_message::{ControlCommand, ControlMessage, MessagePayload};
+pub use dbsql::Dbsql;
 pub use net::{messages, messages::SlabMessage, protocol_slab::ProtocolSlab};
-pub use dbsql::Dbsql; 
 pub use slabs_manager::{SlabsManager, SlabsManagerSafe};
-

+ 5 - 12
src/darkpulse/net/messages.rs

@@ -1,11 +1,11 @@
 use std::io;
 
 use crate::{
+    darkpulse::Ciphertext,
     net::messages::Message,
     serial::{Decodable, Encodable},
-    Result
+    Result,
 };
-use crate::darkpulse::Ciphertext;
 
 #[derive(Clone)]
 pub struct GetSlabsMessage {
@@ -58,9 +58,7 @@ impl Encodable for GetSlabsMessage {
 
 impl Decodable for GetSlabsMessage {
     fn decode<D: io::Read>(mut d: D) -> Result<Self> {
-        Ok(Self {
-            slabs_hash: Decodable::decode(&mut d)?,
-        })
+        Ok(Self { slabs_hash: Decodable::decode(&mut d)? })
     }
 }
 
@@ -75,10 +73,7 @@ impl Encodable for SlabMessage {
 
 impl Decodable for SlabMessage {
     fn decode<D: io::Read>(mut d: D) -> Result<Self> {
-        Ok(Self {
-            nonce: Decodable::decode(&mut d)?,
-            ciphertext: Decodable::decode(&mut d)?,
-        })
+        Ok(Self { nonce: Decodable::decode(&mut d)?, ciphertext: Decodable::decode(&mut d)? })
     }
 }
 
@@ -92,9 +87,7 @@ impl Encodable for InvMessage {
 
 impl Decodable for InvMessage {
     fn decode<D: io::Read>(mut d: D) -> Result<Self> {
-        Ok(Self {
-            slabs_hash: Decodable::decode(&mut d)?,
-        })
+        Ok(Self { slabs_hash: Decodable::decode(&mut d)? })
     }
 }
 

+ 0 - 1
src/darkpulse/net/mod.rs

@@ -1,3 +1,2 @@
 pub mod messages;
 pub mod protocol_slab;
-

+ 18 - 27
src/darkpulse/net/protocol_slab.rs

@@ -3,11 +3,18 @@ use std::sync::Arc;
 use log::*;
 use smol::Executor;
 
-use crate::error::Result as NetResult;
-use crate::serial::deserialize;
-use crate::net::{message_subscriber::MessageSubscription, protocols::ProtocolJobsManager,
-protocols::ProtocolJobsManagerPtr ,ChannelPtr};
-use crate::darkpulse::{aes_decrypt, messages, ControlCommand, ControlMessage, SlabsManagerSafe, CiphertextHash};
+use crate::{
+    darkpulse::{
+        aes_decrypt, messages, CiphertextHash, ControlCommand, ControlMessage, SlabsManagerSafe,
+    },
+    error::Result as NetResult,
+    net::{
+        message_subscriber::MessageSubscription,
+        protocols::{ProtocolJobsManager, ProtocolJobsManagerPtr},
+        ChannelPtr,
+    },
+    serial::deserialize,
+};
 
 pub struct ProtocolSlab {
     channel: ChannelPtr,
@@ -62,23 +69,11 @@ impl ProtocolSlab {
         debug!(target: "net", "ProtocolSlab::start() [START]");
         self.jobsman.clone().start(executor.clone());
 
-        self.jobsman
-            .clone()
-            .spawn(self.clone().handle_receive_sync(), executor.clone())
-            .await;
-        self.jobsman
-            .clone()
-            .spawn(self.clone().handle_receive_inv(), executor.clone())
-            .await;
+        self.jobsman.clone().spawn(self.clone().handle_receive_sync(), executor.clone()).await;
+        self.jobsman.clone().spawn(self.clone().handle_receive_inv(), executor.clone()).await;
 
-        self.jobsman
-            .clone()
-            .spawn(self.clone().handle_receive_get_slabs(), executor.clone())
-            .await;
-        self.jobsman
-            .clone()
-            .spawn(self.clone().handle_receive_slab(), executor)
-            .await;
+        self.jobsman.clone().spawn(self.clone().handle_receive_get_slabs(), executor.clone()).await;
+        self.jobsman.clone().spawn(self.clone().handle_receive_slab(), executor).await;
 
         let _ = self.channel.send(messages::SyncMessage {}).await;
 
@@ -90,9 +85,7 @@ impl ProtocolSlab {
         loop {
             let _sync_msg = self.sync_sub.receive().await?;
             let slab_hashs = self.slabman.lock().await.get_slabs_hash();
-            let inv_msg = messages::InvMessage {
-                slabs_hash: slab_hashs.clone(),
-            };
+            let inv_msg = messages::InvMessage { slabs_hash: slab_hashs.clone() };
             self.channel.send(inv_msg).await?;
             info!("receive sync message!");
         }
@@ -109,9 +102,7 @@ impl ProtocolSlab {
                     list_of_hash.push(slab.clone());
                 }
             }
-            let getslabs_msg = messages::GetSlabsMessage {
-                slabs_hash: list_of_hash,
-            };
+            let getslabs_msg = messages::GetSlabsMessage { slabs_hash: list_of_hash };
             self.channel.send(getslabs_msg).await?;
             info!("receive inv message!");
         }

+ 4 - 8
src/darkpulse/slabs_manager.rs

@@ -1,11 +1,10 @@
-use std::collections::HashMap;
-use std::sync::Arc;
+use std::{collections::HashMap, sync::Arc};
 
 use log::*;
 use sha2::{Digest, Sha256};
 
-use crate:: Result;
-use super::{dbsql, net::messages::SlabMessage, channel::Channel, CiphertextHash, aes::Ciphertext};
+use super::{aes::Ciphertext, channel::Channel, dbsql, net::messages::SlabMessage, CiphertextHash};
+use crate::Result;
 
 pub fn cipher_hash(ciphertext: &Ciphertext) -> CiphertextHash {
     let mut cipher_hash = [0u8; 32];
@@ -83,10 +82,7 @@ impl SlabsManager {
         let default_slabs: HashMap<CiphertextHash, SlabMessage> = HashMap::new();
 
         if let Some(channel_id) = self.main_channel.get_channel_id() {
-            self.slabs = self
-                .db
-                .get_channel_slabs(channel_id.clone())
-                .unwrap_or(default_slabs);
+            self.slabs = self.db.get_channel_slabs(channel_id.clone()).unwrap_or(default_slabs);
         }
     }
 

+ 28 - 42
src/darkpulse/utility.rs

@@ -1,20 +1,26 @@
-use std::fs::OpenOptions;
-use std::io::prelude::*;
-use std::net::SocketAddr;
-use std::path::PathBuf;
-use std::sync::atomic::AtomicU64;
-use std::sync::Arc;
-use std::time::{SystemTime, UNIX_EPOCH};
+use std::{
+    fs::OpenOptions,
+    io::prelude::*,
+    net::SocketAddr,
+    path::PathBuf,
+    sync::{atomic::AtomicU64, Arc},
+    time::{SystemTime, UNIX_EPOCH},
+};
 
 use async_executor::Executor;
 use futures::prelude::*;
 use log::*;
 
-use super::{aes_encrypt, Dbsql, messages, ProtocolSlab, ControlMessage, ControlCommand, Channel,
-SlabsManagerSafe, MessagePayload};
+use super::{
+    aes_encrypt, messages, Channel, ControlCommand, ControlMessage, Dbsql, MessagePayload,
+    ProtocolSlab, SlabsManagerSafe,
+};
 
-
-use crate::{ serial::{serialize, deserialize}, net::ChannelPtr, Result};
+use crate::{
+    net::ChannelPtr,
+    serial::{deserialize, serialize},
+    Result,
+};
 
 pub type AddrsStorage = Arc<async_std::sync::Mutex<Vec<SocketAddr>>>;
 
@@ -22,12 +28,11 @@ pub type Clock = Arc<AtomicU64>;
 
 pub fn get_current_time() -> u64 {
     let start = SystemTime::now();
-    let since_the_epoch = start
-        .duration_since(UNIX_EPOCH)
-        .expect("Incorrect system clock: time went backwards");
+    let since_the_epoch =
+        start.duration_since(UNIX_EPOCH).expect("Incorrect system clock: time went backwards");
     let in_ms =
         since_the_epoch.as_secs() * 1000 + since_the_epoch.subsec_nanos() as u64 / 1_000_000;
-    return in_ms;
+    return in_ms
 }
 
 pub fn save_to_addrs_store(stored_addrs: &Vec<SocketAddr>) -> Result<()> {
@@ -62,11 +67,7 @@ pub fn default_config_dir() -> Result<PathBuf> {
 pub fn load_stored_addrs() -> Result<Vec<SocketAddr>> {
     let path = default_config_dir()?.join("addrs.add");
     println!("{:?}", path);
-    let mut reader = OpenOptions::new()
-        .read(true)
-        .write(true)
-        .create(true)
-        .open(path)?;
+    let mut reader = OpenOptions::new().read(true).write(true).create(true).open(path)?;
     let mut buffer = Vec::new();
     reader.read_to_end(&mut buffer)?;
     if !buffer.is_empty() {
@@ -87,16 +88,9 @@ pub async fn pack_slab(
     let timestamp = chrono::offset::Utc::now();
     let timestamp: i64 = timestamp.timestamp_millis() / 1000;
 
-    let msg_payload = MessagePayload {
-        nickname: username.clone(),
-        text: message,
-        timestamp,
-    };
+    let msg_payload = MessagePayload { nickname: username.clone(), text: message, timestamp };
 
-    let control_message = ControlMessage {
-        control: control_command,
-        payload: msg_payload,
-    };
+    let control_message = ControlMessage { control: control_command, payload: msg_payload };
 
     let ser_message = serialize(&control_message);
 
@@ -143,7 +137,7 @@ pub fn choose_channel(db: &Dbsql, channel_name: Option<String>) -> Result<Channe
                     .next()
                     .expect(format!("there is no channel with the name {}: ", name).as_str())
                     .clone();
-                }
+            }
             None => {
                 main_channel = channels.first().unwrap().clone();
             }
@@ -162,18 +156,10 @@ pub async fn setup_network_channel(
 ) {
     let message_subsytem = channel.get_message_subsystem();
 
-    message_subsytem
-        .add_dispatch::<messages::SyncMessage>()
-        .await;
-    message_subsytem
-        .add_dispatch::<messages::InvMessage>()
-        .await;
-    message_subsytem
-        .add_dispatch::<messages::GetSlabsMessage>()
-        .await;
-    message_subsytem
-        .add_dispatch::<messages::SlabMessage>()
-        .await;
+    message_subsytem.add_dispatch::<messages::SyncMessage>().await;
+    message_subsytem.add_dispatch::<messages::InvMessage>().await;
+    message_subsytem.add_dispatch::<messages::GetSlabsMessage>().await;
+    message_subsytem.add_dispatch::<messages::SlabMessage>().await;
 
     let protocol_slab = ProtocolSlab::new(slabman, channel.clone()).await;
     protocol_slab.clone().start(executor.clone()).await;