From e9ae5897f77097ba6214865084a9beca8b92de67 Mon Sep 17 00:00:00 2001 From: "Dr. Maxim Orlovsky" Date: Mon, 14 Nov 2022 10:41:44 +0100 Subject: [PATCH] Complete transcoder * Encode/encrypt data sent to the remote peer * Remove unused Transport * Rename Decoder -> Deserializer --- radicle-node/src/client.rs | 7 +- .../src/{decoder.rs => deserializer.rs} | 20 +-- radicle-node/src/lib.rs | 5 +- radicle-node/src/service.rs | 1 + radicle-node/src/transport.rs | 111 -------------- radicle-node/src/wire.rs | 138 +++++++++++++----- radicle-node/src/wire/message.rs | 4 +- radicle-node/src/wire/transcoder.rs | 60 +++++++- 8 files changed, 175 insertions(+), 171 deletions(-) rename radicle-node/src/{decoder.rs => deserializer.rs} (81%) delete mode 100644 radicle-node/src/transport.rs diff --git a/radicle-node/src/client.rs b/radicle-node/src/client.rs index 8747d4c2..938a5ca6 100644 --- a/radicle-node/src/client.rs +++ b/radicle-node/src/client.rs @@ -9,8 +9,7 @@ use radicle::crypto::Signer; use crate::clock::RefClock; use crate::profile::Profile; use crate::service::routing; -use crate::transport::Transport; -use crate::wire::transcoder::PlainTranscoder; +use crate::wire::transcoder::NoHandshake; use crate::wire::Wire; use crate::{address, service}; @@ -131,11 +130,9 @@ impl Client { rng, ); - let transcode = PlainTranscoder::default(); - self.reactor.run( &config.listen, - Transport::new(Wire::new(service, transcode)), + Wire::<_, _, _, _, NoHandshake>::new(service), self.events, self.commands, )?; diff --git a/radicle-node/src/decoder.rs b/radicle-node/src/deserializer.rs similarity index 81% rename from radicle-node/src/decoder.rs rename to radicle-node/src/deserializer.rs index 3984e589..7507b03e 100644 --- a/radicle-node/src/decoder.rs +++ b/radicle-node/src/deserializer.rs @@ -4,16 +4,16 @@ use std::marker::PhantomData; use crate::service::message::Message; use crate::wire; -/// Message stream decoder. +/// Message stream deserializer. /// /// Used to for example turn a byte stream into network messages. #[derive(Debug)] -pub struct Decoder { +pub struct Deserializer { unparsed: Vec, item: PhantomData, } -impl From> for Decoder { +impl From> for Deserializer { fn from(unparsed: Vec) -> Self { Self { unparsed, @@ -22,7 +22,7 @@ impl From> for Decoder { } } -impl Decoder { +impl Deserializer { /// Create a new stream decoder. pub fn new(capacity: usize) -> Self { Self { @@ -37,7 +37,7 @@ impl Decoder { } /// Decode and return the next message. Returns [`None`] if nothing was decoded. - pub fn decode_next(&mut self) -> Result, wire::Error> { + pub fn deserialize_next(&mut self) -> Result, wire::Error> { let mut reader = io::Cursor::new(self.unparsed.as_mut_slice()); match D::decode(&mut reader) { @@ -53,7 +53,7 @@ impl Decoder { } } -impl io::Write for Decoder { +impl io::Write for Deserializer { fn write(&mut self, buf: &[u8]) -> io::Result { self.input(buf); @@ -65,11 +65,11 @@ impl io::Write for Decoder { } } -impl Iterator for Decoder { +impl Iterator for Deserializer { type Item = Result; fn next(&mut self) -> Option { - self.decode_next().transpose() + self.deserialize_next().transpose() } } @@ -85,7 +85,7 @@ mod test { fn prop_decode_next(chunk_size: usize) { let mut bytes = vec![]; let mut msgs = vec![]; - let mut decoder = Decoder::::new(8); + let mut decoder = Deserializer::::new(8); let chunk_size = 1 + chunk_size % MSG_HELLO.len() + MSG_BYE.len(); @@ -95,7 +95,7 @@ mod test { for chunk in bytes.as_slice().chunks(chunk_size) { decoder.input(chunk); - while let Some(msg) = decoder.decode_next().unwrap() { + while let Some(msg) = decoder.deserialize_next().unwrap() { msgs.push(msg); } } diff --git a/radicle-node/src/lib.rs b/radicle-node/src/lib.rs index 47f95c3b..1119eca9 100644 --- a/radicle-node/src/lib.rs +++ b/radicle-node/src/lib.rs @@ -2,7 +2,7 @@ pub mod address; pub mod client; pub mod clock; pub mod control; -pub mod decoder; +pub mod deserializer; pub mod logger; pub mod service; pub mod sql; @@ -10,7 +10,6 @@ pub mod sql; pub mod test; #[cfg(test)] pub mod tests; -pub mod transport; pub mod wire; pub use nakamoto_net::{Io, Link, LocalDuration, LocalTime}; @@ -19,7 +18,7 @@ pub use radicle::{collections, crypto, git, hash, identity, node, profile, rad, pub mod prelude { pub use crate::clock::Timestamp; pub use crate::crypto::{PublicKey, Signature, Signer}; - pub use crate::decoder::Decoder; + pub use crate::deserializer::Deserializer; pub use crate::hash::Digest; pub use crate::identity::{Did, Id}; pub use crate::service::filter::Filter; diff --git a/radicle-node/src/service.rs b/radicle-node/src/service.rs index a1e0be1a..e311e37e 100644 --- a/radicle-node/src/service.rs +++ b/radicle-node/src/service.rs @@ -501,6 +501,7 @@ where peer.attempted(); } + // TODO: Split into two functions: `connected` and `negotiated` pub fn connected( &mut self, addr: std::net::SocketAddr, diff --git a/radicle-node/src/transport.rs b/radicle-node/src/transport.rs deleted file mode 100644 index ba081285..00000000 --- a/radicle-node/src/transport.rs +++ /dev/null @@ -1,111 +0,0 @@ -use std::net; -use std::ops::{Deref, DerefMut}; - -use nakamoto::LocalTime; -use nakamoto_net as nakamoto; -use nakamoto_net::{Io, Link}; - -use crate::address; -use crate::collections::HashMap; -use crate::crypto; -use crate::service::routing; -use crate::service::{Command, DisconnectReason, Event, Service}; -use crate::storage::WriteStorage; -use crate::wire::transcoder::Transcode; -use crate::wire::Wire; - -#[derive(Debug)] -struct Peer { - _addr: net::SocketAddr, -} - -#[derive(Debug)] -pub struct Transport { - _peers: HashMap, - inner: Wire, -} - -impl Transport { - pub fn new(inner: Wire) -> Self { - Self { - _peers: HashMap::default(), - inner, - } - } -} - -impl nakamoto::Protocol for Transport -where - R: routing::Store, - W: WriteStorage + 'static, - S: address::Store, - G: crypto::Signer, - T: Transcode, -{ - type Event = Event; - type Command = Command; - type DisconnectReason = DisconnectReason; - - fn initialize(&mut self, time: LocalTime) { - self.inner.initialize(time) - } - - fn tick(&mut self, now: nakamoto::LocalTime) { - self.inner.tick(now) - } - - fn wake(&mut self) { - self.inner.wake() - } - - fn command(&mut self, cmd: Self::Command) { - self.inner.command(cmd) - } - - fn attempted(&mut self, addr: &std::net::SocketAddr) { - self.inner.attempted(addr) - } - - fn connected( - &mut self, - addr: std::net::SocketAddr, - local_addr: &std::net::SocketAddr, - link: Link, - ) { - self.inner.connected(addr, local_addr, link) - } - - fn disconnected( - &mut self, - addr: &std::net::SocketAddr, - reason: nakamoto::DisconnectReason, - ) { - self.inner.disconnected(addr, reason) - } - - fn received_bytes(&mut self, addr: &std::net::SocketAddr, bytes: &[u8]) { - self.inner.received_bytes(addr, bytes) - } -} - -impl Iterator for Transport { - type Item = Io; - - fn next(&mut self) -> Option { - self.inner.next() - } -} - -impl Deref for Transport { - type Target = Service; - - fn deref(&self) -> &Self::Target { - &self.inner - } -} - -impl DerefMut for Transport { - fn deref_mut(&mut self) -> &mut Self::Target { - &mut self.inner - } -} diff --git a/radicle-node/src/wire.rs b/radicle-node/src/wire.rs index e7ae91cc..1320fc28 100644 --- a/radicle-node/src/wire.rs +++ b/radicle-node/src/wire.rs @@ -1,20 +1,20 @@ pub mod message; pub mod transcoder; -use std::collections::{BTreeMap, HashMap}; +use std::collections::{BTreeMap, HashMap, VecDeque}; use std::convert::TryFrom; use std::net; -use std::ops::{Deref, DerefMut}; +use std::ops::Deref; use std::string::FromUtf8Error; use std::{io, mem}; use byteorder::{NetworkEndian, ReadBytesExt, WriteBytesExt}; use nakamoto_net as nakamoto; -use nakamoto_net::Link; +use nakamoto_net::{Link, LocalTime}; use crate::address; use crate::crypto::{PublicKey, Signature, Signer, Unverified}; -use crate::decoder::Decoder; +use crate::deserializer::Deserializer; use crate::git; use crate::git::fmt; use crate::hash::Digest; @@ -26,7 +26,7 @@ use crate::service::{filter, routing}; use crate::storage::refs::Refs; use crate::storage::refs::SignedRefs; use crate::storage::WriteStorage; -use crate::wire::transcoder::Transcode; +use crate::wire::transcoder::{Handshake, HandshakeResult, Transcode}; /// The default type we use to represent sizes on the wire. /// @@ -425,51 +425,119 @@ impl Decode for node::Features { } #[derive(Debug)] -pub struct Wire { - inboxes: HashMap, - inner: service::Service, - #[allow(dead_code)] - transcoder: T, +pub struct Inbox { + pub transcoder: T, + pub deserializer: Deserializer, } -impl Wire { - pub fn new(inner: service::Service, transcoder: T) -> Self { +#[derive(Debug)] +pub struct Wire { + handshakes: HashMap, + handshake_queue: VecDeque<(net::SocketAddr, Vec)>, + inboxes: HashMap>, + inner: service::Service, +} + +impl Wire { + pub fn new(inner: service::Service) -> Self { Self { + handshakes: HashMap::new(), + handshake_queue: Default::default(), inboxes: HashMap::new(), inner, - transcoder, } } } -impl Wire +impl nakamoto::Protocol for Wire where R: routing::Store, S: address::Store, W: WriteStorage + 'static, G: Signer, - T: Transcode, + H: Handshake, { - pub fn connected(&mut self, addr: net::SocketAddr, local_addr: &net::SocketAddr, link: Link) { - self.inboxes.insert(addr, Decoder::new(256)); + type Event = service::Event; + type Command = service::Command; + type DisconnectReason = service::DisconnectReason; + + fn initialize(&mut self, time: LocalTime) { + self.inner.initialize(time) + } + + fn tick(&mut self, now: nakamoto::LocalTime) { + self.inner.tick(now) + } + + fn wake(&mut self) { + self.inner.wake() + } + + fn command(&mut self, cmd: Self::Command) { + self.inner.command(cmd) + } + + fn attempted(&mut self, addr: &std::net::SocketAddr) { + self.inner.attempted(addr) + } + + fn connected(&mut self, addr: net::SocketAddr, local_addr: &net::SocketAddr, link: Link) { + self.handshakes.insert(addr, H::new()); self.inner.connected(addr, local_addr, link) } - pub fn disconnected( + fn disconnected( &mut self, addr: &net::SocketAddr, reason: nakamoto::DisconnectReason, ) { + self.handshakes.remove(addr); self.inboxes.remove(addr); self.inner.disconnected(addr, &reason) } - pub fn received_bytes(&mut self, addr: &net::SocketAddr, bytes: &[u8]) { - if let Some(inbox) = self.inboxes.get_mut(addr) { - inbox.input(bytes); + fn received_bytes(&mut self, addr: &net::SocketAddr, raw_bytes: &[u8]) { + if let Some(handshake) = self.handshakes.remove(addr) { + debug_assert!(!self.inboxes.contains_key(addr)); + + match handshake.step(raw_bytes) { + HandshakeResult::Next(handshake, reply) => { + self.handshakes.insert(*addr, handshake); + if !reply.is_empty() { + self.handshake_queue.push_back((*addr, reply)); + } + return; + } + HandshakeResult::Complete(transcoder, reply) => { + log::debug!("handshake with peer {} is complete", addr); + if !reply.is_empty() { + self.handshake_queue.push_back((*addr, reply)); + } + self.inboxes.insert( + *addr, + Inbox { + transcoder, + deserializer: Deserializer::new(256), + }, + ); + } + HandshakeResult::Error(err) => { + log::error!("invalid handshake input. Details: {}", err); + return; + } + } + } + + if let Some(Inbox { + transcoder, + deserializer, + }) = self.inboxes.get_mut(addr) + { + let bytes = transcoder.decrypt(raw_bytes); + deserializer.input(&bytes); loop { - match inbox.decode_next() { + match deserializer.deserialize_next() { Ok(Some(msg)) => self.inner.received_message(addr, msg), Ok(None) => break, @@ -487,10 +555,14 @@ where } } -impl Iterator for Wire { +impl Iterator for Wire { type Item = nakamoto::Io; fn next(&mut self) -> Option { + if let Some((addr, handshake_data)) = self.handshake_queue.pop_front() { + return Some(nakamoto::Io::Write(addr, handshake_data)); + } + match self.inner.next() { Some(Io::Write(addr, msgs)) => { let mut buf = Vec::new(); @@ -500,7 +572,11 @@ impl Iterator for Wire { msg.encode(&mut buf) .expect("writing to an in-memory buffer doesn't fail"); } - Some(nakamoto::Io::Write(addr, buf)) + let Inbox { transcoder, .. } = self.inboxes.get_mut(&addr).expect( + "broken handshake implementation: data sent before handshake was complete", + ); + let data = transcoder.encrypt(buf); + Some(nakamoto::Io::Write(addr, data)) } Some(Io::Event(e)) => Some(nakamoto::Io::Event(e)), Some(Io::Connect(a)) => Some(nakamoto::Io::Connect(a)), @@ -512,20 +588,6 @@ impl Iterator for Wire { } } -impl Deref for Wire { - type Target = service::Service; - - fn deref(&self) -> &Self::Target { - &self.inner - } -} - -impl DerefMut for Wire { - fn deref_mut(&mut self) -> &mut Self::Target { - &mut self.inner - } -} - #[cfg(test)] mod tests { use super::*; diff --git a/radicle-node/src/wire/message.rs b/radicle-node/src/wire/message.rs index b9ff2ecc..1824cdf9 100644 --- a/radicle-node/src/wire/message.rs +++ b/radicle-node/src/wire/message.rs @@ -360,7 +360,7 @@ mod tests { use super::*; use quickcheck_macros::quickcheck; - use crate::decoder::Decoder; + use crate::deserializer::Deserializer; use crate::wire::{self, Encode}; #[test] @@ -412,7 +412,7 @@ mod tests { #[test] fn prop_message_decoder() { fn property(items: Vec) { - let mut decoder = Decoder::::new(8); + let mut decoder = Deserializer::::new(8); for item in &items { item.encode(&mut decoder).unwrap(); diff --git a/radicle-node/src/wire/transcoder.rs b/radicle-node/src/wire/transcoder.rs index 6ca59517..94e34197 100644 --- a/radicle-node/src/wire/transcoder.rs +++ b/radicle-node/src/wire/transcoder.rs @@ -1,6 +1,62 @@ -pub trait Transcode {} +use std::convert::Infallible; +// TODO: Implement Try trait once stabilized +/// Result of a state-machine transition. +pub enum HandshakeResult { + Next(H, Vec), + Complete(T, Vec), + Error(H::Error), +} + +pub trait Handshake: Sized { + /// Errors which may happen during the handshake. + type Error: std::error::Error; + /// Underlying transcoder. + type Transcoder: Transcode; + + /// Create a new handshake state-machine. + fn new() -> Self; + /// Advance the state-machine to the next state. + fn step(self, input: &[u8]) -> HandshakeResult; +} + +#[derive(Debug, Default)] +pub struct NoHandshake; + +impl Handshake for NoHandshake { + type Error = Infallible; + type Transcoder = PlainTranscoder; + + fn new() -> Self { + NoHandshake + } + + fn step(self, _input: &[u8]) -> HandshakeResult { + HandshakeResult::Complete(PlainTranscoder, vec![]) + } +} + +/// Trait allowing transcoding a stream using some form of stream encryption +/// and/or encoding. +pub trait Transcode { + /// Decodes data received from the remote peer and updates the internal state + /// of the transcoder. + fn decrypt(&mut self, data: &[u8]) -> Vec; + + /// Encodes data before sending it to the remote peer. + fn encrypt(&mut self, data: Vec) -> Vec; +} + +/// Transcoder which does nothing. #[derive(Debug, Default)] pub struct PlainTranscoder; -impl Transcode for PlainTranscoder {} +impl Transcode for PlainTranscoder { + fn decrypt(&mut self, data: &[u8]) -> Vec { + data.to_vec() + } + + fn encrypt(&mut self, data: Vec) -> Vec { + data + } +}