Merge pull request #124 from Revertron/feature/encryption

Implemented P2P traffic encryption.
This commit is contained in:
Revertron
2021-05-30 00:55:20 +02:00
committed by GitHub
16 changed files with 2173 additions and 492 deletions
Generated
+1317
View File
File diff suppressed because it is too large Load Diff
+4 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "alfis" name = "alfis"
version = "0.5.9" version = "0.6.0"
authors = ["Revertron <alfis@revertron.com>"] authors = ["Revertron <alfis@revertron.com>"]
edition = "2018" edition = "2018"
build = "build.rs" build = "build.rs"
@@ -18,7 +18,9 @@ toml = "0.5.8"
digest = "0.9.0" digest = "0.9.0"
sha2 = "0.9.5" sha2 = "0.9.5"
ed25519-dalek = "1.0.1" ed25519-dalek = "1.0.1"
x25519-dalek = "1.1.1"
ecies-ed25519 = "0.5.1" ecies-ed25519 = "0.5.1"
chacha20poly1305 = "0.8.0"
signature = "1.3.0" signature = "1.3.0"
blakeout = "0.3.0" blakeout = "0.3.0"
num_cpus = "1.13.0" num_cpus = "1.13.0"
@@ -26,6 +28,7 @@ byteorder = "1.4.3"
serde = { version = "1.0.126", features = ["derive"] } serde = { version = "1.0.126", features = ["derive"] }
serde_json = "1.0.64" serde_json = "1.0.64"
bincode = "1.3.3" bincode = "1.3.3"
serde_cbor = "0.11.1"
base64 = "0.13.0" base64 = "0.13.0"
num-bigint = "0.4.0" num-bigint = "0.4.0"
num-traits = "0.2.14" num-traits = "0.2.14"
+10
View File
@@ -61,13 +61,23 @@ impl Block {
} }
} }
pub fn from_bytes(data: &[u8]) -> serde_cbor::Result<Self> {
serde_cbor::from_slice(data)
}
pub fn is_genesis(&self) -> bool { pub fn is_genesis(&self) -> bool {
self.index == 1 && self.index == 1 &&
matches!(Transaction::get_type(&self.transaction), TransactionType::Origin) && matches!(Transaction::get_type(&self.transaction), TransactionType::Origin) &&
self.prev_block_hash == Bytes::default() self.prev_block_hash == Bytes::default()
} }
/// Serializes block to CBOR for network
pub fn as_bytes(&self) -> Vec<u8> { pub fn as_bytes(&self) -> Vec<u8> {
serde_cbor::to_vec(&self).unwrap()
}
/// Serializes block to bincode format for hashing.
pub fn as_bytes_compact(&self) -> Vec<u8> {
bincode::serialize(&self).unwrap() bincode::serialize(&self).unwrap()
} }
+27 -1
View File
@@ -1010,8 +1010,10 @@ impl SignersCache {
pub mod tests { pub mod tests {
use log::LevelFilter; use log::LevelFilter;
use simplelog::{ColorChoice, ConfigBuilder, TerminalMode, TermLogger}; use simplelog::{ColorChoice, ConfigBuilder, TerminalMode, TermLogger};
#[allow(unused_imports)]
use log::{debug, error, info, trace, warn};
use crate::{Chain, Settings}; use crate::{Chain, Settings, Block};
fn init_logger() { fn init_logger() {
let config = ConfigBuilder::new() let config = ConfigBuilder::new()
@@ -1035,4 +1037,28 @@ pub mod tests {
chain.check_chain(u64::MAX); chain.check_chain(u64::MAX);
assert_eq!(chain.get_height(), 149); assert_eq!(chain.get_height(), 149);
} }
#[test]
pub fn check_serde() {
let settings = Settings::default();
let chain = Chain::new(&settings, "./tests/blockchain.db");
// Check the first block, its transaction doesn't have identity
let block = chain.get_block(1).unwrap();
let buf = serde_cbor::to_vec(&block).unwrap();
let block2: Block = serde_cbor::from_slice(&buf[..]).unwrap();
assert_eq!(block, block2);
// Check second block, it is common "full" block with domain
let block = chain.get_block(2).unwrap();
let buf = serde_cbor::to_vec(&block).unwrap();
let block2: Block = serde_cbor::from_slice(&buf[..]).unwrap();
assert_eq!(block, block2);
// Check block 36, it is an "empty" block, used to sign full blocks
let block = chain.get_block(36).unwrap();
let buf = serde_cbor::to_vec(&block).unwrap();
let block2: Block = serde_cbor::from_slice(&buf[..]).unwrap();
assert_eq!(block, block2);
}
} }
+2 -2
View File
@@ -9,7 +9,7 @@ pub fn check_block_hash(block: &Block) -> bool {
let mut copy: Block = block.clone(); let mut copy: Block = block.clone();
copy.hash = Bytes::default(); copy.hash = Bytes::default();
copy.signature = Bytes::default(); copy.signature = Bytes::default();
blakeout_data(&copy.as_bytes()) == block.hash blakeout_data(&copy.as_bytes_compact()) == block.hash
} }
/// Hashes data by given hasher /// Hashes data by given hasher
@@ -23,7 +23,7 @@ pub fn blakeout_data(data: &[u8]) -> Bytes {
pub fn check_block_signature(block: &Block) -> bool { pub fn check_block_signature(block: &Block) -> bool {
let mut copy = block.clone(); let mut copy = block.clone();
copy.signature = Bytes::default(); copy.signature = Bytes::default();
Keystore::check(&copy.as_bytes(), &copy.pub_key, &block.signature) Keystore::check(&copy.as_bytes_compact(), &copy.pub_key, &block.signature)
} }
/// Hashes some identity (domain in case of DNS). If you give it a public key, it will hash with it as well. /// Hashes some identity (domain in case of DNS). If you give it a public key, it will hash with it as well.
-5
View File
@@ -49,11 +49,6 @@ impl Transaction {
} }
} }
pub fn get_bytes(&self) -> Vec<u8> {
// Let it panic if something is not okay
serde_json::to_vec(&self).unwrap()
}
pub fn to_string(&self) -> String { pub fn to_string(&self) -> String {
// Let it panic if something is not okay // Let it panic if something is not okay
serde_json::to_string(&self).unwrap() serde_json::to_string(&self).unwrap()
-8
View File
@@ -1,7 +1,6 @@
use std::net::IpAddr; use std::net::IpAddr;
use std::num; use std::num;
use mio::Token;
use rand::Rng; use rand::Rng;
#[cfg(not(target_os = "macos"))] #[cfg(not(target_os = "macos"))]
use thread_priority::*; use thread_priority::*;
@@ -116,13 +115,6 @@ pub fn is_yggdrasil(addr: &IpAddr) -> bool {
false false
} }
/// Gets new token from old token, mutating the last
pub fn next(current: &mut Token) -> Token {
let next = current.0;
current.0 += 1;
Token(next)
}
/// Checks if this record has IP from Yggdrasil network /// Checks if this record has IP from Yggdrasil network
/// https://yggdrasil-network.github.io /// https://yggdrasil-network.github.io
pub fn is_yggdrasil_record(record: &DnsRecord) -> bool { pub fn is_yggdrasil_record(record: &DnsRecord) -> bool {
+66
View File
@@ -0,0 +1,66 @@
use chacha20poly1305::{ChaCha20Poly1305, Key, Nonce};
use chacha20poly1305::aead::{Aead, NewAead};
use std::fmt::{Debug, Formatter};
use std::fmt;
pub const ZERO_NONCE: [u8; 12] = [0u8; 12];
const FAILURE: &str = "encryption failure!";
/// A small wrap-up to use Chacha20 encryption for domain names.
#[derive(Clone)]
pub struct Chacha {
cipher: ChaCha20Poly1305,
nonce: [u8; 12]
}
impl Chacha {
pub fn new(key: &[u8], nonce: &[u8]) -> Self {
let key = Key::from_slice(key);
let cipher = ChaCha20Poly1305::new(key);
let mut buf = [0u8; 12];
buf.copy_from_slice(nonce);
Chacha { cipher, nonce: buf }
}
pub fn encrypt(&self, data: &[u8]) -> Vec<u8> {
let nonce = Nonce::from(self.nonce.clone());
self.cipher.encrypt(&nonce, data.as_ref()).expect(FAILURE)
}
pub fn decrypt(&self, data: &[u8]) -> Vec<u8> {
let nonce = Nonce::from(self.nonce.clone());
self.cipher.decrypt(&nonce, data.as_ref()).expect(FAILURE)
}
pub fn get_nonce(&self) -> &[u8; 12] {
&self.nonce
}
}
impl Debug for Chacha {
fn fmt(&self, fmt: &mut Formatter<'_>) -> fmt::Result {
fmt.write_str("ChaCha20Poly1305")
}
}
#[cfg(test)]
mod tests {
use crate::crypto::Chacha;
use crate::{to_hex};
#[test]
pub fn test_chacha() {
let buf = b"178135D209C697625E3EC71DA5C760382E54936F824EE5083908DA66B14ECE18";
let chacha1 = Chacha::new(b"178135D209C697625E3EC71DA5C76038", &buf[..12]);
let bytes1 = chacha1.encrypt(b"TEST");
println!("{}", to_hex(&bytes1));
let chacha2 = Chacha::new(b"178135D209C697625E3EC71DA5C76038", &buf[..12]);
let bytes2 = chacha2.decrypt(&bytes1);
assert_eq!(String::from_utf8(bytes2).unwrap(), "TEST");
let bytes2 = chacha2.encrypt(b"TEST");
assert_eq!(bytes1, bytes2);
}
}
+3
View File
@@ -1,3 +1,6 @@
mod crypto_box; mod crypto_box;
mod chacha;
pub use crypto_box::CryptoBox; pub use crypto_box::CryptoBox;
pub use chacha::Chacha;
pub use chacha::ZERO_NONCE;
+5 -1
View File
@@ -184,7 +184,11 @@ fn main() {
let miner: Arc<Mutex<Miner>> = Arc::new(Mutex::new(miner_obj)); let miner: Arc<Mutex<Miner>> = Arc::new(Mutex::new(miner_obj));
let mut network = Network::new(Arc::clone(&context)); let mut network = Network::new(Arc::clone(&context));
network.start().expect("Error starting network component"); thread::spawn(move || {
// Give UI some time to appear :)
thread::sleep(Duration::from_millis(1000));
network.start();
});
create_genesis_if_needed(&context, &miner); create_genesis_if_needed(&context, &miner);
if no_gui { if no_gui {
+2 -2
View File
@@ -280,7 +280,7 @@ impl Miner {
Some(mut block) => { Some(mut block) => {
let index = block.index; let index = block.index;
let mut context = context.lock().unwrap(); let mut context = context.lock().unwrap();
block.signature = Bytes::from_bytes(&job.keystore.sign(&block.as_bytes())); block.signature = Bytes::from_bytes(&job.keystore.sign(&block.as_bytes_compact()));
let mut success = false; let mut success = false;
if context.chain.check_new_block(&block) != BlockQuality::Good { if context.chain.check_new_block(&block) != BlockQuality::Good {
warn!("Error adding mined block!"); warn!("Error adding mined block!");
@@ -341,7 +341,7 @@ fn find_hash(context: Arc<Mutex<Context>>, mut block: Block, running: Arc<Atomic
block.nonce = nonce; block.nonce = nonce;
digest.reset(); digest.reset();
digest.update(&block.as_bytes()); digest.update(&block.as_bytes_compact());
let diff = hash_difficulty(digest.result()); let diff = hash_difficulty(digest.result());
if diff >= target_diff { if diff >= target_diff {
block.hash = Bytes::from_bytes(digest.result()); block.hash = Bytes::from_bytes(digest.result());
+10 -28
View File
@@ -7,8 +7,8 @@ use crate::Bytes;
#[derive(Debug, Serialize, Deserialize)] #[derive(Debug, Serialize, Deserialize)]
pub enum Message { pub enum Message {
Error, Error,
Hand { #[serde(default = "default_version")] app_version: String, origin: String, version: u32, public: bool, #[serde(default)] rand: String }, Hand { app_version: String, origin: String, version: u32, public: bool, rand_id: String, },
Shake { #[serde(default = "default_version")] app_version: String, origin: String, version: u32, ok: bool, height: u64 }, Shake { app_version: String, origin: String, version: u32, public: bool, rand_id: String, height: u64 },
Ping { height: u64, hash: Bytes }, Ping { height: u64, hash: Bytes },
Pong { height: u64, hash: Bytes }, Pong { height: u64, hash: Bytes },
Twin, Twin,
@@ -16,24 +16,23 @@ pub enum Message {
GetPeers, GetPeers,
Peers { peers: Vec<String> }, Peers { peers: Vec<String> },
GetBlock { index: u64 }, GetBlock { index: u64 },
Block { index: u64, block: String }, Block { index: u64, block: Vec<u8> },
} }
impl Message { impl Message {
pub fn from_bytes(bytes: Vec<u8>) -> Result<Self, ()> { pub fn from_bytes(bytes: Vec<u8>) -> Result<Self, ()> {
let text = String::from_utf8(bytes).unwrap_or(String::from("Error{}")); match serde_cbor::from_slice(bytes.as_slice()) {
match serde_json::from_str(&text) {
Ok(cmd) => Ok(cmd), Ok(cmd) => Ok(cmd),
Err(_) => Err(()) Err(_) => Err(())
} }
} }
pub fn hand(app_version: &str, origin: &str, version: u32, public: bool, rand: &str) -> Self { pub fn hand(app_version: &str, origin: &str, version: u32, public: bool, rand_id: &str) -> Self {
Message::Hand { app_version: app_version.to_owned(), origin: origin.to_owned(), version, public, rand: rand.to_owned() } Message::Hand { app_version: app_version.to_owned(), origin: origin.to_owned(), version, public, rand_id: rand_id.to_owned() }
} }
pub fn shake(app_version: &str, origin: &str, version: u32, ok: bool, height: u64) -> Self { pub fn shake(app_version: &str, origin: &str, version: u32, public: bool, rand_id: &str, height: u64) -> Self {
Message::Shake { app_version: app_version.to_owned(), origin: origin.to_owned(), version, ok, height } Message::Shake { app_version: app_version.to_owned(), origin: origin.to_owned(), version, public, rand_id: rand_id.to_owned(), height }
} }
pub fn ping(height: u64, hash: Bytes) -> Self { pub fn ping(height: u64, hash: Bytes) -> Self {
@@ -44,24 +43,7 @@ impl Message {
Message::Pong { height, hash } Message::Pong { height, hash }
} }
pub fn block(height: u64, str: String) -> Self { pub fn block(height: u64, block: Vec<u8>) -> Self {
Message::Block { index: height, block: str } Message::Block { index: height, block }
} }
} }
fn default_version() -> String {
String::from("0.0.0")
}
#[cfg(test)]
mod tests {
use crate::p2p::Message;
#[test]
pub fn test_hand() {
assert!(serde_json::from_str::<Message>("\"Error\"").is_ok());
assert!(serde_json::from_str::<Message>("{\"Hand\":{\"origin\":\"\",\"version\":1,\"public\":false,\"rand\":\"123\"}}").is_ok());
assert!(serde_json::from_str::<Message>("{\"Hand\":{\"origin\":\"\",\"version\":1,\"public\":false}}").is_ok());
}
}
+418 -176
View File
@@ -3,7 +3,7 @@ extern crate serde_json;
use std::{io, thread}; use std::{io, thread};
use std::cmp::max; use std::cmp::max;
use std::io::{Read, Write}; use std::io::{Read, Write, Error};
use std::net::{IpAddr, Shutdown, SocketAddr, SocketAddrV4}; use std::net::{IpAddr, Shutdown, SocketAddr, SocketAddrV4};
use std::sync::{Arc, Mutex}; use std::sync::{Arc, Mutex};
use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::atomic::{AtomicBool, Ordering};
@@ -15,25 +15,38 @@ use log::{debug, error, info, trace, warn};
use mio::{Events, Interest, Poll, Registry, Token}; use mio::{Events, Interest, Poll, Registry, Token};
use mio::event::Event; use mio::event::Event;
use mio::net::{TcpListener, TcpStream}; use mio::net::{TcpListener, TcpStream};
use rand::random; use rand::{random, RngCore, Rng};
use rand_old::prelude::thread_rng;
use x25519_dalek::{StaticSecret, PublicKey};
use crate::{Block, Context, p2p::Message, p2p::Peer, p2p::Peers, p2p::State}; use crate::{Block, Context, p2p::Message, p2p::Peer, p2p::Peers, p2p::State};
use crate::blockchain::types::BlockQuality; use crate::blockchain::types::BlockQuality;
use crate::commons::*; use crate::commons::*;
use crate::eventbus::{register, post}; use crate::eventbus::{register, post};
use crate::crypto::Chacha;
const SERVER: Token = Token(0); const SERVER: Token = Token(0);
pub struct Network { pub struct Network {
context: Arc<Mutex<Context>> context: Arc<Mutex<Context>>,
secret_key: StaticSecret,
public_key: PublicKey,
token: Token,
// States of peer connections, and some data to send when sockets become writable
peers: Peers,
} }
impl Network { impl Network {
pub fn new(context: Arc<Mutex<Context>>) -> Self { pub fn new(context: Arc<Mutex<Context>>) -> Self {
Network { context } // P2P encryption primitives
let mut thread_rng = thread_rng();
let secret_key = StaticSecret::new(&mut thread_rng);
let public_key = PublicKey::from(&secret_key);
let peers = Peers::new();
Network { context, secret_key, public_key, token: Token(1), peers }
} }
pub fn start(&mut self) -> Result<(), String> { pub fn start(&mut self) {
let (listen_addr, peers_addrs, yggdrasil_only) = { let (listen_addr, peers_addrs, yggdrasil_only) = {
let c = self.context.lock().unwrap(); let c = self.context.lock().unwrap();
(c.settings.net.listen.clone(), c.settings.net.peers.clone(), c.settings.net.yggdrasil_only) (c.settings.net.listen.clone(), c.settings.net.peers.clone(), c.settings.net.yggdrasil_only)
@@ -47,19 +60,12 @@ impl Network {
let mut server = TcpListener::bind(addr).expect("Can't bind to address"); let mut server = TcpListener::bind(addr).expect("Can't bind to address");
debug!("Started node listener on {}", server.local_addr().unwrap()); debug!("Started node listener on {}", server.local_addr().unwrap());
let mut events = Events::with_capacity(1024); let mut events = Events::with_capacity(64);
let mut poll = Poll::new().expect("Unable to create poll"); let mut poll = Poll::new().expect("Unable to create poll");
poll.registry().register(&mut server, SERVER, Interest::READABLE).expect("Error registering poll"); poll.registry().register(&mut server, SERVER, Interest::READABLE).expect("Error registering poll");
let context = Arc::clone(&self.context);
thread::spawn(move || {
// Give UI some time to appear :)
thread::sleep(Duration::from_millis(2000));
// Unique token for each incoming connection.
let mut unique_token = Token(SERVER.0 + 1);
// States of peer connections, and some data to send when sockets become writable
let mut peers = Peers::new();
// Starting peer connections to bootstrap nodes // Starting peer connections to bootstrap nodes
peers.connect_peers(&peers_addrs, &poll.registry(), &mut unique_token, yggdrasil_only); self.peers.connect_peers(&peers_addrs, &poll.registry(), &mut self.token, yggdrasil_only);
let mut ui_timer = Instant::now(); let mut ui_timer = Instant::now();
let mut log_timer = Instant::now(); let mut log_timer = Instant::now();
@@ -67,10 +73,10 @@ impl Network {
let mut connect_timer = Instant::now(); let mut connect_timer = Instant::now();
let mut last_events_time = Instant::now(); let mut last_events_time = Instant::now();
loop { loop {
if peers.get_peers_count() == 0 && bootstrap_timer.elapsed().as_secs() > 60 { if self.peers.get_peers_count() == 0 && bootstrap_timer.elapsed().as_secs() > 60 {
warn!("Restarting swarm connections..."); warn!("Restarting swarm connections...");
// Starting peer connections to bootstrap nodes // Starting peer connections to bootstrap nodes
peers.connect_peers(&peers_addrs, &poll.registry(), &mut unique_token, yggdrasil_only); self.peers.connect_peers(&peers_addrs, &poll.registry(), &mut self.token, yggdrasil_only);
bootstrap_timer = Instant::now(); bootstrap_timer = Instant::now();
last_events_time = Instant::now(); last_events_time = Instant::now();
} }
@@ -100,7 +106,7 @@ impl Network {
} }
} }
if peers.is_ignored(&address.ip()) { if self.peers.is_ignored(&address.ip()) {
debug!("Ignoring connection from banned {:?}", &address.ip()); debug!("Ignoring connection from banned {:?}", &address.ip());
continue; continue;
} }
@@ -113,26 +119,24 @@ impl Network {
} }
//debug!("Accepted connection from: {} to local IP: {}", address, local_ip); //debug!("Accepted connection from: {} to local IP: {}", address, local_ip);
let token = next(&mut unique_token); let token = self.next_token();
poll.registry().register(&mut stream, token, Interest::READABLE).expect("Error registering poll"); poll.registry().register(&mut stream, token, Interest::READABLE).expect("Error registering poll");
peers.add_peer(token, Peer::new(address, stream, State::Connected, true)); let peer = Peer::new(address, stream, State::Connected, true);
self.peers.add_peer(token, peer);
} }
Err(_) => {} Err(_) => {}
} }
match poll.registry().reregister(&mut server, SERVER, Interest::READABLE) { if let Err(e) = poll.registry().reregister(&mut server, SERVER, Interest::READABLE) {
Ok(_) => {}
Err(e) => {
panic!("Error reregistering server token!\n{}", e); panic!("Error reregistering server token!\n{}", e);
} }
} }
}
token => { token => {
if !handle_connection_event(Arc::clone(&context), &mut peers, &poll.registry(), &event) { if !self.handle_connection_event(&poll.registry(), &event) {
let _ = peers.close_peer(poll.registry(), &token); let _ = self.peers.close_peer(poll.registry(), &token);
let blocks = context.lock().unwrap().chain.get_height(); let blocks = self.context.lock().unwrap().chain.get_height();
let keys = context.lock().unwrap().chain.get_users_count(); let keys = self.context.lock().unwrap().chain.get_users_count();
let domains = context.lock().unwrap().chain.get_domains_count(); let domains = self.context.lock().unwrap().chain.get_domains_count();
post(crate::event::Event::NetworkStatus { blocks, domains, keys, nodes: peers.get_peers_active_count() }); post(crate::event::Event::NetworkStatus { blocks, domains, keys, nodes: self.peers.get_peers_active_count() });
} }
} }
} }
@@ -140,9 +144,9 @@ impl Network {
if !events.is_empty() { if !events.is_empty() {
last_events_time = Instant::now(); last_events_time = Instant::now();
} else if last_events_time.elapsed().as_secs() > MAX_IDLE_SECONDS { } else if last_events_time.elapsed().as_secs() > MAX_IDLE_SECONDS {
if peers.get_peers_count() > 0 { if self.peers.get_peers_count() > 0 {
warn!("Something is wrong with swarm connections, closing all."); warn!("Something is wrong with swarm connections, closing all.");
peers.close_all_peers(poll.registry()); self.peers.close_all_peers(poll.registry());
continue; continue;
} else { } else {
thread::sleep(POLL_TIMEOUT.unwrap()); thread::sleep(POLL_TIMEOUT.unwrap());
@@ -152,10 +156,10 @@ impl Network {
if ui_timer.elapsed().as_millis() > UI_REFRESH_DELAY_MS { if ui_timer.elapsed().as_millis() > UI_REFRESH_DELAY_MS {
// Send pings to idle peers // Send pings to idle peers
let (height, hash) = { let (height, hash) = {
let context = context.lock().unwrap(); let context = self.context.lock().unwrap();
let blocks = context.chain.get_height(); let blocks = context.chain.get_height();
let nodes = peers.get_peers_active_count(); let nodes = self.peers.get_peers_active_count();
let banned = peers.get_peers_banned_count(); let banned = self.peers.get_peers_banned_count();
let keys = context.chain.get_users_count(); let keys = context.chain.get_users_count();
let domains = context.chain.get_domains_count(); let domains = context.chain.get_domains_count();
@@ -170,12 +174,12 @@ impl Network {
log_timer = Instant::now(); log_timer = Instant::now();
} }
if nodes < MAX_NODES && connect_timer.elapsed().as_secs() >= 5 { if nodes < MAX_NODES && connect_timer.elapsed().as_secs() >= 5 {
peers.connect_new_peers(poll.registry(), &mut unique_token, yggdrasil_only); self.peers.connect_new_peers(poll.registry(), &mut self.token, yggdrasil_only);
connect_timer = Instant::now(); connect_timer = Instant::now();
} }
(blocks, context.chain.get_last_hash()) (blocks, context.chain.get_last_hash())
}; };
peers.update(poll.registry(), height, hash); self.peers.update(poll.registry(), height, hash);
ui_timer = Instant::now(); ui_timer = Instant::now();
} }
} }
@@ -184,26 +188,9 @@ impl Network {
} else { } else {
panic!("Network loop has broken prematurely!"); panic!("Network loop has broken prematurely!");
} }
});
Ok(())
} }
}
fn subscribe_to_bus(running: Arc<AtomicBool>) { fn handle_connection_event(&mut self, registry: &Registry, event: &Event) -> bool {
use crate::event::Event;
register(move |_uuid, e| {
match e {
Event::ActionQuit => {
running.store(false, Ordering::SeqCst);
return false;
}
_ => {}
}
true
});
}
fn handle_connection_event(context: Arc<Mutex<Context>>, peers: &mut Peers, registry: &Registry, event: &Event) -> bool {
if event.is_error() || (event.is_read_closed() && event.is_write_closed()) { if event.is_error() || (event.is_read_closed() && event.is_write_closed()) {
return false; return false;
} }
@@ -211,7 +198,7 @@ fn handle_connection_event(context: Arc<Mutex<Context>>, peers: &mut Peers, regi
if event.is_readable() { if event.is_readable() {
let data = { let data = {
let token = event.token(); let token = event.token();
match peers.get_mut_peer(&token) { match self.peers.get_mut_peer(&token) {
None => { None => {
error!("Error getting peer for connection {}", token.0); error!("Error getting peer for connection {}", token.0);
return false; return false;
@@ -221,20 +208,95 @@ fn handle_connection_event(context: Arc<Mutex<Context>>, peers: &mut Peers, regi
debug!("Node from {} disconnected", peer.get_addr().ip()); debug!("Node from {} disconnected", peer.get_addr().ip());
return false; return false;
} }
match peer.get_state().clone() {
State::Connected => {
let mut stream = peer.get_stream();
return match read_client_handshake(&mut stream) {
Ok(key) => {
let mut buf = [0u8; 32];
buf.copy_from_slice(key.as_slice());
let public_key: PublicKey = PublicKey::from(buf);
let shared = self.secret_key.diffie_hellman(&public_key);
let mut nonce = [0u8; 12];
let mut rng = rand::thread_rng();
rng.fill(&mut nonce);
let chacha = Chacha::new(shared.as_bytes(), &nonce);
registry.reregister(stream, event.token(), Interest::WRITABLE).unwrap();
std::mem::drop(stream);
peer.set_cipher(chacha);
peer.set_state(State::ServerHandshake);
info!("Client hello read successfully");
true
}
Err(e) => {
warn!("Error reading client handshake. {}", e);
false
}
}
}
State::ServerHandshake => {
let mut stream = peer.get_stream();
return match read_server_handshake(&mut stream) {
Ok(data) => {
if data.len() != 32 + 12 {
warn!("Server handshake of {} bytes instead of {}", data.len(), 32 + 12);
return false;
}
let mut buf = [0u8; 32];
buf.copy_from_slice(&data.as_slice()[0..32]);
let public_key: PublicKey = PublicKey::from(buf);
let mut nonce = [0u8; 12];
nonce.copy_from_slice(&data.as_slice()[32..]);
let shared = self.secret_key.diffie_hellman(&public_key);
let chacha = Chacha::new(shared.as_bytes(), &nonce);
registry.reregister(stream, event.token(), Interest::WRITABLE).unwrap();
std::mem::drop(stream);
peer.set_cipher(chacha);
peer.set_state(State::HandshakeFinished);
info!("Server hello read successfully");
true
}
Err(e) => {
warn!("Error reading server handshake. {}", e);
false
}
}
}
_ => {
let mut stream = peer.get_stream(); let mut stream = peer.get_stream();
read_message(&mut stream) read_message(&mut stream)
} }
} }
}
}
}; };
if data.is_ok() { if data.is_ok() {
let data = {
match self.peers.get_peer(&event.token()) {
Some(peer) => {
let data = data.unwrap(); let data = data.unwrap();
//info!("Decoding message {:?}", to_hex(data.as_slice()));
match decode_message(&data, peer.get_cipher()) {
Ok(data) => {
data
}
Err(_) => {
vec![]
}
}
}
None => {
vec![]
}
}
};
match Message::from_bytes(data) { match Message::from_bytes(data) {
Ok(message) => { Ok(message) => {
//let m = format!("{:?}", &message); let m = format!("{:?}", &message);
let new_state = handle_message(Arc::clone(&context), message, peers, &event.token()); let new_state = self.handle_message(message, &event.token());
let peer = peers.get_mut_peer(&event.token()).unwrap(); let peer = self.peers.get_mut_peer(&event.token()).unwrap();
//debug!("Got message from {}: {:?}", &peer.get_addr(), &m); debug!("Got message from {}: {:?}", &peer.get_addr(), &m);
let stream = peer.get_stream(); let stream = peer.get_stream();
match new_state { match new_state {
State::Message { data } => { State::Message { data } => {
@@ -243,31 +305,38 @@ fn handle_connection_event(context: Arc<Mutex<Context>>, peers: &mut Peers, regi
} }
State::Connecting => {} State::Connecting => {}
State::Connected => {} State::Connected => {}
State::ServerHandshake => {}
State::HandshakeFinished => {}
State::Idle { .. } => { State::Idle { .. } => {
peer.set_state(State::idle()); peer.set_state(State::idle());
} }
State::Error => {} State::Error => {}
State::Banned => { State::Banned => {
peers.ignore_peer(registry, &event.token()); self.peers.ignore_peer(registry, &event.token());
} }
State::Offline { .. } => { State::Offline { .. } => {
peer.set_state(State::offline()); peer.set_state(State::offline());
} }
State::Loop => { State::Loop => {
peer.set_state(State::Loop); peer.set_state(State::Loop);
peers.ignore_peer(registry, &event.token()); self.peers.ignore_peer(registry, &event.token());
} }
State::SendLoop => { State::SendLoop => {
registry.reregister(stream, event.token(), Interest::WRITABLE).unwrap(); registry.reregister(stream, event.token(), Interest::WRITABLE).unwrap();
peer.set_state(State::SendLoop); peer.set_state(State::SendLoop);
} }
State::Twin => { State::Twin => {
registry.reregister(stream, event.token(), Interest::WRITABLE).unwrap();
peer.set_state(State::Twin); peer.set_state(State::Twin);
// TODO set something in [Peers], maybe ignore this IP?
return false;
} }
} }
} }
Err(_) => { return false; } Err(_) => {
let peer = self.peers.get_peer(&event.token()).unwrap();
warn!("Error deserializing message from {}", &peer.get_addr());
return false;
}
} }
} else { } else {
return false; return false;
@@ -275,36 +344,50 @@ fn handle_connection_event(context: Arc<Mutex<Context>>, peers: &mut Peers, regi
} }
if event.is_writable() { if event.is_writable() {
//trace!("Socket {} is writable", event.token().0); let my_id = self.peers.get_my_id().to_owned();
let my_id = peers.get_my_id().to_owned(); match self.peers.get_mut_peer(&event.token()) {
match peers.get_mut_peer(&event.token()) {
None => {} None => {}
Some(peer) => { Some(peer) => {
match peer.get_state().clone() { match peer.get_state().clone() {
State::Connecting => { State::Connecting => {
if send_client_handshake(&mut peer.get_stream(), self.public_key.as_bytes()).is_err() {
return false;
}
peer.set_state(State::ServerHandshake);
}
State::ServerHandshake => {
if send_server_handshake(peer, self.public_key.as_bytes()).is_err() {
return false;
}
peer.set_state(State::HandshakeFinished);
info!("Server handshake sent");
}
State::HandshakeFinished => {
//debug!("Connected to peer {}, sending hello...", &peer.get_addr()); //debug!("Connected to peer {}, sending hello...", &peer.get_addr());
let data: String = { let data: Vec<u8> = {
let c = context.lock().unwrap(); let c = self.context.lock().unwrap();
let message = Message::hand(&c.app_version, &c.settings.origin, CHAIN_VERSION, c.settings.net.public, &my_id); let message = Message::hand(&c.app_version, &c.settings.origin, CHAIN_VERSION, c.settings.net.public, &my_id);
serde_json::to_string(&message).unwrap() info!("Sending: {:?}", &message);
encode_message(&message, peer.get_cipher()).unwrap()
}; };
send_message(peer.get_stream(), &data.into_bytes()).unwrap_or_else(|e| warn!("Error sending hello {}", e)); send_message(peer.get_stream(), &data).unwrap_or_else(|e| warn!("Error sending hello {}", e));
//debug!("Sent hello to {}", &peer.get_addr()); //debug!("Sent hello to {}", &peer.get_addr());
} }
State::Connected => {}
State::Message { data } => { State::Message { data } => {
//debug!("Sending data to {}: {}", &peer.get_addr(), &String::from_utf8(data.clone()).unwrap()); //debug!("Sending data to {}: {}", &peer.get_addr(), &String::from_utf8(data.clone()).unwrap());
let data = encode_bytes(&data, peer.get_cipher());
send_message(peer.get_stream(), &data).unwrap_or_else(|e| warn!("Error sending message {}", e)); send_message(peer.get_stream(), &data).unwrap_or_else(|e| warn!("Error sending message {}", e));
} }
State::Connected => {}
State::Idle { from } => { State::Idle { from } => {
debug!("Odd version of pings :)"); debug!("Odd version of pings :)");
if from.elapsed().as_secs() >= 30 { if from.elapsed().as_secs() >= 30 {
let data: String = { let data: Vec<u8> = {
let c = context.lock().unwrap(); let c = self.context.lock().unwrap();
let message = Message::ping(c.chain.get_height(), c.chain.get_last_hash()); let message = Message::ping(c.chain.get_height(), c.chain.get_last_hash());
serde_json::to_string(&message).unwrap() encode_message(&message, peer.get_cipher()).unwrap()
}; };
send_message(peer.get_stream(), &data.into_bytes()).unwrap_or_else(|e| warn!("Error sending ping {}", e)); send_message(peer.get_stream(), &data).unwrap_or_else(|e| warn!("Error sending ping {}", e));
} }
} }
State::Error => {} State::Error => {}
@@ -312,12 +395,12 @@ fn handle_connection_event(context: Arc<Mutex<Context>>, peers: &mut Peers, regi
State::Offline { .. } => {} State::Offline { .. } => {}
State::Loop => {} State::Loop => {}
State::SendLoop => { State::SendLoop => {
let data = serde_json::to_string(&Message::Loop).unwrap(); let data = encode_message(&Message::Loop, peer.get_cipher()).unwrap();
send_message(peer.get_stream(), &data.into_bytes()).unwrap_or_else(|e| warn!("Error sending loop {}", e)); send_message(peer.get_stream(), &data).unwrap_or_else(|e| warn!("Error sending loop {}", e));
} }
State::Twin => { State::Twin => {
let data = serde_json::to_string(&Message::Twin).unwrap(); let data = encode_message(&Message::Twin, peer.get_cipher()).unwrap();
send_message(peer.get_stream(), &data.into_bytes()).unwrap_or_else(|e| warn!("Error sending loop {}", e)); send_message(peer.get_stream(), &data).unwrap_or_else(|e| warn!("Error sending loop {}", e));
} }
} }
registry.reregister(peer.get_stream(), event.token(), Interest::READABLE).unwrap(); registry.reregister(peer.get_stream(), event.token(), Interest::READABLE).unwrap();
@@ -326,102 +409,50 @@ fn handle_connection_event(context: Arc<Mutex<Context>>, peers: &mut Peers, regi
} }
true true
}
fn read_message(stream: &mut TcpStream) -> Result<Vec<u8>, ()> {
let instant = Instant::now();
let data_size = match stream.read_u32::<BigEndian>() {
Ok(size) => { size as usize }
Err(e) => {
error!("Error reading from socket! {}", e);
0
}
};
//trace!("Payload size is {}", data_size);
if data_size > MAX_PACKET_SIZE || data_size == 0 {
return Err(());
} }
let mut buf = vec![0u8; data_size]; fn handle_message(&mut self, message: Message, token: &Token) -> State {
let mut bytes_read = 0; let (my_height, my_hash, my_origin, my_version, me_public) = {
loop { let context = self.context.lock().unwrap();
match stream.read(&mut buf[bytes_read..]) {
Ok(bytes) => {
bytes_read += bytes;
if bytes_read == data_size {
break;
}
}
// Would block "errors" are the OS's way of saying that the connection is not actually ready to perform this I/O operation.
Err(ref err) if would_block(err) => {
// We give every connection no more than 200ms to read a message
if instant.elapsed().as_millis() < MAX_READ_BLOCK_TIME {
// We need to sleep a bit, otherwise it can eat CPU
let delay = Duration::from_millis(2);
thread::sleep(delay);
continue;
} else {
break;
}
},
Err(ref err) if interrupted(err) => continue,
// Other errors we'll consider fatal.
Err(_) => {
debug!("Error reading message, only {}/{} bytes read", bytes_read, data_size);
return Err(())
},
}
}
if buf.len() == data_size {
Ok(buf)
} else {
Err(())
}
}
fn send_message(connection: &mut TcpStream, data: &Vec<u8>) -> io::Result<()> {
connection.write_u32::<BigEndian>(data.len() as u32)?;
connection.write_all(&data)?;
connection.flush()
}
fn handle_message(context: Arc<Mutex<Context>>, message: Message, peers: &mut Peers, token: &Token) -> State {
let (my_height, my_hash, my_origin, my_version) = {
let context = context.lock().unwrap();
// TODO cache it somewhere // TODO cache it somewhere
(context.chain.get_height(), context.chain.get_last_hash(), &context.settings.origin.clone(), CHAIN_VERSION) (context.chain.get_height(), context.chain.get_last_hash(), &context.settings.origin.clone(), CHAIN_VERSION, context.settings.net.public)
}; };
let my_id = self.peers.get_my_id().to_owned();
let answer = match message { let answer = match message {
Message::Hand { app_version, origin, version, public, rand} => { Message::Hand { app_version, origin, version, public, rand_id } => {
if peers.is_our_own_connect(&rand) { if self.peers.is_our_own_connect(&rand_id) {
warn!("Detected loop connect"); warn!("Detected loop connect");
State::SendLoop State::SendLoop
} else { } else {
if origin.eq(my_origin) && version == my_version { if origin.eq(my_origin) && version == my_version {
let peer = peers.get_mut_peer(token).unwrap(); let peer = self.peers.get_mut_peer(token).unwrap();
peer.set_public(public); peer.set_public(public);
peer.set_active(true); peer.set_active(true);
debug!("Incoming v{} on {}", &app_version, peer.get_addr().ip()); debug!("Incoming v{} on {}", &app_version, peer.get_addr().ip());
let app_version = context.lock().unwrap().app_version.clone(); let app_version = self.context.lock().unwrap().app_version.clone();
State::message(Message::shake(&app_version, &origin, version, true, my_height)) State::message(Message::shake(&app_version, &origin, version, me_public, &my_id, my_height))
} else { } else {
warn!("Handshake from unsupported chain or version"); warn!("Handshake from unsupported chain or version");
State::Banned State::Banned
} }
} }
} }
Message::Shake { app_version, origin, version, ok, height } => { Message::Shake { app_version, origin, version, public, rand_id, height } => {
if origin.ne(my_origin) || version != my_version { if origin.ne(my_origin) || version != my_version {
return State::Banned; return State::Banned;
} }
if ok { if self.peers.is_tween_connect(&rand_id) {
let nodes = peers.get_peers_active_count(); return State::Twin;
let peer = peers.get_mut_peer(token).unwrap(); }
let nodes = self.peers.get_peers_active_count();
let peer = self.peers.get_mut_peer(token).unwrap();
// TODO check rand_id whether we have this peers connection already
debug!("Outgoing v{} on {}", &app_version, peer.get_addr().ip()); debug!("Outgoing v{} on {}", &app_version, peer.get_addr().ip());
peer.set_height(height); peer.set_height(height);
peer.set_active(true); peer.set_active(true);
peer.set_public(public);
peer.reset_reconnects(); peer.reset_reconnects();
let mut context = context.lock().unwrap(); let mut context = self.context.lock().unwrap();
if peer.is_higher(my_height) { if peer.is_higher(my_height) {
context.chain.update_max_height(height); context.chain.update_max_height(height);
let event = crate::event::Event::Syncing { have: my_height, height: max(height, my_height) }; let event = crate::event::Event::Syncing { have: my_height, height: max(height, my_height) };
@@ -433,17 +464,14 @@ fn handle_message(context: Arc<Mutex<Context>>, message: Message, peers: &mut Pe
} else { } else {
State::idle() State::idle()
} }
} else {
State::Banned
}
} }
Message::Error => { State::Error } Message::Error => { State::Error }
Message::Ping { height, hash } => { Message::Ping { height, hash } => {
let peer = peers.get_mut_peer(token).unwrap(); let peer = self.peers.get_mut_peer(token).unwrap();
peer.set_height(height); peer.set_height(height);
peer.set_active(true); peer.set_active(true);
if peer.is_higher(my_height) { if peer.is_higher(my_height) {
let mut context = context.lock().unwrap(); let mut context = self.context.lock().unwrap();
context.chain.update_max_height(height); context.chain.update_max_height(height);
info!("Peer is higher, requesting block {} from {}", height, peer.get_addr().ip()); info!("Peer is higher, requesting block {} from {}", height, peer.get_addr().ip());
State::message(Message::GetBlock { index: my_height + 1 }) State::message(Message::GetBlock { index: my_height + 1 })
@@ -456,12 +484,12 @@ fn handle_message(context: Arc<Mutex<Context>>, message: Message, peers: &mut Pe
} }
} }
Message::Pong { height, hash } => { Message::Pong { height, hash } => {
let active_count = peers.get_peers_active_count(); let active_count = self.peers.get_peers_active_count();
let peer = peers.get_mut_peer(token).unwrap(); let peer = self.peers.get_mut_peer(token).unwrap();
peer.set_height(height); peer.set_height(height);
peer.set_active(true); peer.set_active(true);
if peer.is_higher(my_height) { if peer.is_higher(my_height) {
let mut context = context.lock().unwrap(); let mut context = self.context.lock().unwrap();
context.chain.update_max_height(height); context.chain.update_max_height(height);
info!("Peer is higher, requesting block {} from {}", height, peer.get_addr().ip()); info!("Peer is higher, requesting block {} from {}", height, peer.get_addr().ip());
State::message(Message::GetBlock { index: my_height + 1 }) State::message(Message::GetBlock { index: my_height + 1 })
@@ -480,52 +508,55 @@ fn handle_message(context: Arc<Mutex<Context>>, message: Message, peers: &mut Pe
} }
Message::GetPeers => { Message::GetPeers => {
let addr = { let addr = {
let peer = peers.get_mut_peer(token).unwrap(); let peer = self.peers.get_mut_peer(token).unwrap();
peer.set_active(true); peer.set_active(true);
peer.get_addr().clone() peer.get_addr().clone()
}; };
State::message(Message::Peers { peers: peers.get_peers_for_exchange(&addr) }) State::message(Message::Peers { peers: self.peers.get_peers_for_exchange(&addr) })
} }
Message::Peers { peers: new_peers } => { Message::Peers { peers: new_peers } => {
let peer = peers.get_mut_peer(token).unwrap(); let peer = self.peers.get_mut_peer(token).unwrap();
peer.set_active(true); peer.set_active(true);
peers.add_peers_from_exchange(new_peers); self.peers.add_peers_from_exchange(new_peers);
State::idle() State::idle()
} }
Message::GetBlock { index } => { Message::GetBlock { index } => {
let peer = peers.get_mut_peer(token).unwrap(); let peer = self.peers.get_mut_peer(token).unwrap();
peer.set_active(true); peer.set_active(true);
let context = context.lock().unwrap(); let context = self.context.lock().unwrap();
match context.chain.get_block(index) { match context.chain.get_block(index) {
Some(block) => State::message(Message::block(block.index, serde_json::to_string(&block).unwrap())), Some(block) => State::message(Message::block(block.index, block.as_bytes())),
None => State::Error None => State::Error
} }
} }
Message::Block { index, block } => { Message::Block { index, block } => {
let peer = peers.get_mut_peer(token).unwrap(); let peer = self.peers.get_mut_peer(token).unwrap();
peer.set_active(true); peer.set_active(true);
let block: Block = match serde_json::from_str(&block) { let block: Block = match Block::from_bytes(block.as_slice()) {
Ok(block) => block, Ok(block) => block,
Err(_) => return State::Banned Err(e) => {
warn!("Error deserializing block! {}", e);
return State::Banned
}
}; };
if index != block.index { if index != block.index {
return State::Banned; return State::Banned;
} }
info!("Received block {} with hash {:?}", block.index, &block.hash); info!("Received block {} with hash {:?}", block.index, &block.hash);
handle_block(context, peers, token, block) self.handle_block(token, block)
} }
Message::Twin => { State::Twin } Message::Twin => { State::Twin }
Message::Loop => { State::Loop } Message::Loop => { State::Loop }
}; };
answer answer
} }
fn handle_block(context: Arc<Mutex<Context>>, peers: &mut Peers, token: &Token, block: Block) -> State { fn handle_block(&mut self, token: &Token, block: Block) -> State {
let peers_count = peers.get_peers_active_count(); let peers_count = self.peers.get_peers_active_count();
let peer = peers.get_mut_peer(token).unwrap(); let peer = self.peers.get_mut_peer(token).unwrap();
peer.set_received_block(block.index); peer.set_received_block(block.index);
let mut context = context.lock().unwrap(); let mut context = self.context.lock().unwrap();
let max_height = context.chain.max_height(); let max_height = context.chain.max_height();
match context.chain.check_new_block(&block) { match context.chain.check_new_block(&block) {
BlockQuality::Good => { BlockQuality::Good => {
@@ -574,12 +605,223 @@ fn handle_block(context: Arc<Mutex<Context>>, peers: &mut Peers, token: &Token,
} else { } else {
debug!("Fork in not better than our block, dropping."); debug!("Fork in not better than our block, dropping.");
if let Some(block) = context.chain.get_block(block.index) { if let Some(block) = context.chain.get_block(block.index) {
return State::message(Message::block(block.index, serde_json::to_string(&block).unwrap())); return State::message(Message::block(block.index, block.as_bytes()));
} }
} }
} }
} }
State::idle() State::idle()
}
/// Gets new token from old token, mutating the last
pub fn next_token(&mut self) -> Token {
let current = self.token.0;
self.token.0 += 1;
Token(current)
}
}
fn subscribe_to_bus(running: Arc<AtomicBool>) {
use crate::event::Event;
register(move |_uuid, e| {
match e {
Event::ActionQuit => {
running.store(false, Ordering::SeqCst);
return false;
}
_ => {}
}
true
});
}
fn encode_bytes(data: &Vec<u8>, cipher: &Option<Chacha>) -> Vec<u8> {
match cipher {
None => { data.clone() }
Some(chacha) => {
chacha.encrypt(data.as_slice())
}
}
}
fn encode_message(message: &Message, cipher: &Option<Chacha>) -> Result<Vec<u8>, ()> {
match serde_cbor::to_vec(message) {
Ok(vec) => {
match cipher {
None => {
//info!("No cipher, not encoding message: {:?}", to_hex(&vec));
Ok(vec)
}
Some(chacha) => {
//info!("Encoding message: {:?}", to_hex(&vec));
Ok(chacha.encrypt(vec.as_slice()))
}
}
}
Err(e) => {
warn!("Could not encode message! {}", e);
Err(())
}
}
}
fn decode_message(data: &Vec<u8>, cipher: &Option<Chacha>) -> Result<Vec<u8>, Error> {
match cipher {
None => { Ok(data.clone()) }
Some(chacha) => {
Ok(chacha.decrypt(data.as_slice()))
}
}
}
fn read_message(stream: &mut TcpStream) -> Result<Vec<u8>, ()> {
let instant = Instant::now();
let data_size = match stream.read_u16::<BigEndian>() {
Ok(size) => { (size ^ 0xAAAA) as usize }
Err(e) => {
error!("Error reading from socket! {}", e);
0
}
};
trace!("Payload size is {}", data_size);
if data_size > MAX_PACKET_SIZE || data_size == 0 {
return Err(());
}
let mut buf = vec![0u8; data_size];
let mut bytes_read = 0;
loop {
match stream.read(&mut buf[bytes_read..]) {
Ok(bytes) => {
bytes_read += bytes;
if bytes_read == data_size {
break;
}
}
// Would block "errors" are the OS's way of saying that the connection is not actually ready to perform this I/O operation.
Err(ref err) if would_block(err) => {
// We give every connection no more than 200ms to read a message
if instant.elapsed().as_millis() < MAX_READ_BLOCK_TIME {
// We need to sleep a bit, otherwise it can eat CPU
let delay = Duration::from_millis(2);
thread::sleep(delay);
continue;
} else {
break;
}
},
Err(ref err) if interrupted(err) => continue,
// Other errors we'll consider fatal.
Err(_) => {
debug!("Error reading message, only {}/{} bytes read", bytes_read, data_size);
return Err(())
},
}
}
if buf.len() == data_size {
Ok(buf)
} else {
Err(())
}
}
/// Sends one byte [garbage_size], [random bytes], and [public_key]
fn send_client_handshake(stream: &mut TcpStream, public_key: &[u8]) -> io::Result<()> {
let mut rng = rand::thread_rng();
let packet_size: usize = rng.gen_range(64..255);
let mut buf = vec![0u8; packet_size];
rng.fill_bytes(&mut buf);
let garbage_size = packet_size - 33;
buf[0] = garbage_size as u8 ^ 0xA; // key length and 1 byte size
for i in 0..public_key.len() {
buf[i + garbage_size + 1] = public_key[i];
}
stream.write_all(buf.as_slice())?;
stream.flush()
}
fn read_client_handshake(stream: &mut TcpStream) -> Result<Vec<u8>, Error> {
// First, we read garbage size
let data_size = match stream.read_u8() {
Ok(size) => { (size ^ 0xA) as usize }
Err(e) => {
error!("Error reading from socket! {}", e);
return Err(e)
}
};
// Read the garbage
let mut buf = vec![0u8; data_size];
match stream.read_exact(&mut buf) {
Ok(_) => {}
Err(e) => { return Err(e); }
}
// Then we have public key for ECDH
let mut buf = vec![0u8; 32];
match stream.read_exact(&mut buf) {
Ok(_) => { Ok(buf) }
Err(e) => {
warn!("Error reading handshake!");
Err(e)
}
}
}
fn send_server_handshake(peer: &mut Peer, public_key: &[u8]) -> io::Result<()> {
let mut rng = rand::thread_rng();
let packet_size: usize = rng.gen_range(64..255);
let mut buf = vec![0u8; packet_size];
rng.fill_bytes(&mut buf);
let nonce = peer.get_nonce();
// We will write 1 byte size, garbage, public key, nonce
let garbage_size = packet_size - 1 - 32 - 12;
buf[0] = garbage_size as u8 ^ 0xA;
for i in 0..public_key.len() {
buf[i + garbage_size + 1] = public_key[i];
}
for i in 0..nonce.len() {
buf[i + garbage_size + 32 + 1] = nonce[i];
}
let stream = peer.get_stream();
stream.write_all(buf.as_slice())?;
stream.flush()
}
fn read_server_handshake(stream: &mut TcpStream) -> Result<Vec<u8>, Error> {
// First, we read garbage size
let data_size = match stream.read_u8() {
Ok(size) => { (size ^ 0xA) as usize }
Err(e) => {
error!("Error reading from socket! {}", e);
return Err(e)
}
};
// Read the garbage
let mut buf = vec![0u8; data_size];
match stream.read_exact(&mut buf) {
Ok(_) => {}
Err(e) => { return Err(e); }
}
// Then we have public key for ECDH, plus nonce 12 bytes
let mut buf = vec![0u8; 32 + 12];
match stream.read_exact(&mut buf) {
Ok(_) => { Ok(buf) }
Err(e) => {
warn!("Error reading handshake!");
Err(e)
}
}
}
fn send_message(connection: &mut TcpStream, data: &Vec<u8>) -> io::Result<()> {
let data_len = data.len() as u16;
//debug!("Sending {} bytes", data_len);
//debug!("Message: {:?}", to_hex(&data));
let mut buf: Vec<u8> = Vec::with_capacity(data.len() + 2);
buf.write_u16::<BigEndian>(data_len ^ 0xAAAA)?;
buf.write_all(&data)?;
connection.write_all(&buf)?;
connection.flush()
} }
fn would_block(err: &io::Error) -> bool { fn would_block(err: &io::Error) -> bool {
+18
View File
@@ -3,6 +3,7 @@ use std::collections::HashMap;
use mio::net::TcpStream; use mio::net::TcpStream;
use crate::p2p::State; use crate::p2p::State;
use crate::Block; use crate::Block;
use crate::crypto::Chacha;
#[derive(Debug)] #[derive(Debug)]
pub struct Peer { pub struct Peer {
@@ -16,6 +17,7 @@ pub struct Peer {
active: bool, active: bool,
reconnects: u32, reconnects: u32,
received_block: u64, received_block: u64,
cipher: Option<Chacha>,
fork: HashMap<u64, Block> fork: HashMap<u64, Block>
} }
@@ -32,10 +34,26 @@ impl Peer {
active: false, active: false,
reconnects: 0, reconnects: 0,
received_block: 0, received_block: 0,
cipher: None,
fork: HashMap::new() fork: HashMap::new()
} }
} }
pub fn set_cipher(&mut self, cipher: Chacha) {
self.cipher = Some(cipher);
}
pub fn get_cipher(&self) -> &Option<Chacha> {
&self.cipher
}
pub fn get_nonce(&self) -> &[u8; 12] {
match &self.cipher {
None => { &crate::crypto::ZERO_NONCE }
Some(chacha) => { chacha.get_nonce() }
}
}
pub fn get_addr(&self) -> SocketAddr { pub fn get_addr(&self) -> SocketAddr {
self.addr.clone() self.addr.clone()
} }
+22 -1
View File
@@ -12,7 +12,6 @@ use rand::seq::IteratorRandom;
use crate::{Bytes, commons}; use crate::{Bytes, commons};
use crate::commons::*; use crate::commons::*;
use crate::p2p::{Message, Peer, State}; use crate::p2p::{Message, Peer, State};
use crate::commons::next;
use std::io; use std::io;
const PING_PERIOD: u64 = 30; const PING_PERIOD: u64 = 30;
@@ -84,6 +83,12 @@ impl Peers {
State::Twin => { State::Twin => {
info!("Peer connection {} to {:?} is a twin", &token.0, &peer.get_addr()); info!("Peer connection {} to {:?} is a twin", &token.0, &peer.get_addr());
} }
State::ServerHandshake => {
info!("Peer connection {} from {:?} didn't shake hands", &token.0, &peer.get_addr());
}
State::HandshakeFinished => {
info!("Peer connection {} from {:?} shaked hands, but then failed", &token.0, &peer.get_addr());
}
} }
self.peers.remove(token); self.peers.remove(token);
@@ -193,6 +198,15 @@ impl Peers {
count count
} }
pub fn is_tween_connect(&self, id: &str) -> bool {
for (_, peer) in self.peers.iter() {
if peer.active() && peer.get_id() == id {
return true;
}
}
false
}
pub fn get_peers_banned_count(&self) -> usize { pub fn get_peers_banned_count(&self) -> usize {
self.ignored.len() self.ignored.len()
} }
@@ -401,6 +415,13 @@ impl Peers {
} }
} }
/// Gets new token from old token, mutating the last
pub fn next(current: &mut Token) -> Token {
let next = current.0;
current.0 += 1;
Token(next)
}
fn skip_private_addr(addr: &SocketAddr) -> bool { fn skip_private_addr(addr: &SocketAddr) -> bool {
if addr.ip().is_loopback() { if addr.ip().is_loopback() {
return true; return true;
+4 -2
View File
@@ -5,6 +5,8 @@ use crate::p2p::Message;
pub enum State { pub enum State {
Connecting, Connecting,
Connected, Connected,
ServerHandshake,
HandshakeFinished,
Idle { from: Instant }, Idle { from: Instant },
Message { data: Vec<u8> }, Message { data: Vec<u8> },
Error, Error,
@@ -25,8 +27,8 @@ impl State {
} }
pub fn message(message: Message) -> Self { pub fn message(message: Message) -> Self {
let response = serde_json::to_string(&message).unwrap(); let data = serde_cbor::to_vec(&message).unwrap();
State::Message {data: Vec::from(response.as_bytes()) } State::Message { data }
} }
pub fn is_idle(&self) -> bool { pub fn is_idle(&self) -> bool {