From a925fb0e02672a754a3e3cb3dbd15b010210bb55 Mon Sep 17 00:00:00 2001 From: Alexis Sellier Date: Thu, 15 Sep 2022 21:26:41 +0200 Subject: [PATCH] node: Create `Transport` layer Signed-off-by: Alexis Sellier --- node/src/client.rs | 18 +++---- node/src/lib.rs | 1 + node/src/protocol.rs | 87 +++++++++++++++------------------- node/src/test/peer.rs | 24 ++++++---- node/src/transport.rs | 106 ++++++++++++++++++++++++++++++++++++++++++ 5 files changed, 171 insertions(+), 65 deletions(-) create mode 100644 node/src/transport.rs diff --git a/node/src/client.rs b/node/src/client.rs index 60657f57..ae833774 100644 --- a/node/src/client.rs +++ b/node/src/client.rs @@ -9,6 +9,7 @@ use crate::collections::HashMap; use crate::crypto::Signer; use crate::protocol; use crate::storage::git::Storage; +use crate::transport::Transport; pub mod handle; @@ -86,16 +87,17 @@ impl Client { log::info!("Initializing client ({:?})..", network); + let protocol = protocol::Protocol::new( + config.protocol, + RefClock::from(time), + storage, + addresses, + signer, + rng, + ); self.reactor.run( &config.listen, - protocol::Protocol::new( - config.protocol, - RefClock::from(time), - storage, - addresses, - signer, - rng, - ), + Transport::new(protocol), self.events, self.commands, )?; diff --git a/node/src/lib.rs b/node/src/lib.rs index 810c5854..42241065 100644 --- a/node/src/lib.rs +++ b/node/src/lib.rs @@ -20,3 +20,4 @@ mod serde_ext; mod storage; #[cfg(test)] mod test; +mod transport; diff --git a/node/src/protocol.rs b/node/src/protocol.rs index d407906f..f42a5c8c 100644 --- a/node/src/protocol.rs +++ b/node/src/protocol.rs @@ -280,47 +280,7 @@ impl<'r, T: WriteStorage<'r>, S: address_book::Store, G: crypto::Signer> Protoco } } - //////////////////////////////////////////////////////////////////////////// - // Periodic tasks - //////////////////////////////////////////////////////////////////////////// - - /// Announce our inventory to all connected peers. - fn announce_inventory(&mut self) -> Result<(), storage::Error> { - let inv = Message::inventory(self.context.inventory_announcement()?, &self.context.signer); - - for addr in self.peers.negotiated().map(|(_, p)| p.addr) { - self.context.write(addr, inv.clone()); - } - Ok(()) - } - - fn prune_routing_entries(&mut self) { - // TODO - } - - fn maintain_connections(&mut self) { - // TODO: Connect to all potential seeds. - if self.peers.len() < TARGET_OUTBOUND_PEERS { - let delta = TARGET_OUTBOUND_PEERS - self.peers.len(); - - for _ in 0..delta { - // TODO: Connect to random peer. - } - } - } -} - -impl<'r, S, T, G> nakamoto::Protocol for Protocol -where - T: WriteStorage<'r> + 'static, - S: address_book::Store, - G: crypto::Signer, -{ - type Event = Event; - type Command = Command; - type DisconnectReason = DisconnectReason; - - fn initialize(&mut self, time: LocalTime) { + pub fn initialize(&mut self, time: LocalTime) { trace!("Init {}", time.as_secs()); self.start_time = time; @@ -332,13 +292,13 @@ where } } - fn tick(&mut self, now: nakamoto::LocalTime) { + pub fn tick(&mut self, now: nakamoto::LocalTime) { trace!("Tick +{}", now - self.start_time); self.context.clock.set(now); } - fn wake(&mut self) { + pub fn wake(&mut self) { let now = self.context.clock.local_time(); trace!("Wake +{}", now - self.start_time); @@ -373,7 +333,7 @@ where } } - fn command(&mut self, cmd: Self::Command) { + pub fn command(&mut self, cmd: Command) { debug!("Command {:?}", cmd); match cmd { @@ -468,7 +428,7 @@ where } } - fn attempted(&mut self, addr: &std::net::SocketAddr) { + pub fn attempted(&mut self, addr: &std::net::SocketAddr) { let ip = addr.ip(); let persistent = self.context.config.is_persistent(addr); let peer = self @@ -479,7 +439,7 @@ where peer.attempted(); } - fn connected( + pub fn connected( &mut self, addr: std::net::SocketAddr, _local_addr: &std::net::SocketAddr, @@ -512,10 +472,10 @@ where } } - fn disconnected( + pub fn disconnected( &mut self, addr: &std::net::SocketAddr, - reason: nakamoto::DisconnectReason, + reason: nakamoto::DisconnectReason, ) { let since = self.local_time(); let ip = addr.ip(); @@ -551,7 +511,7 @@ where } } - fn received_bytes(&mut self, addr: &std::net::SocketAddr, bytes: &[u8]) { + pub fn received_bytes(&mut self, addr: &std::net::SocketAddr, bytes: &[u8]) { let peer_ip = addr.ip(); let (peer, msgs) = if let Some(peer) = self.peers.get_mut(&peer_ip) { let decoder = peer.inbox(); @@ -602,6 +562,35 @@ where self.context.relay(msg, negotiated.clone()); } } + + //////////////////////////////////////////////////////////////////////////// + // Periodic tasks + //////////////////////////////////////////////////////////////////////////// + + /// Announce our inventory to all connected peers. + fn announce_inventory(&mut self) -> Result<(), storage::Error> { + let inv = Message::inventory(self.context.inventory_announcement()?, &self.context.signer); + + for addr in self.peers.negotiated().map(|(_, p)| p.addr) { + self.context.write(addr, inv.clone()); + } + Ok(()) + } + + fn prune_routing_entries(&mut self) { + // TODO + } + + fn maintain_connections(&mut self) { + // TODO: Connect to all potential seeds. + if self.peers.len() < TARGET_OUTBOUND_PEERS { + let delta = TARGET_OUTBOUND_PEERS - self.peers.len(); + + for _ in 0..delta { + // TODO: Connect to random peer. + } + } + } } impl Deref for Protocol { diff --git a/node/src/test/peer.rs b/node/src/test/peer.rs index 2191304c..180dcb58 100644 --- a/node/src/test/peer.rs +++ b/node/src/test/peer.rs @@ -15,15 +15,16 @@ use crate::protocol::message::*; use crate::protocol::*; use crate::storage::WriteStorage; use crate::test::crypto::MockSigner; +use crate::transport; use crate::*; -/// Protocol instantiation used for testing. -pub type Protocol = crate::protocol::Protocol, S, MockSigner>; +/// Transport instantiation used for testing. +pub type Transport = transport::Transport, S, MockSigner>; #[derive(Debug)] pub struct Peer { pub name: &'static str, - pub protocol: Protocol, + pub protocol: Transport, pub ip: net::IpAddr, pub rng: fastrand::Rng, pub local_time: LocalTime, @@ -32,7 +33,7 @@ pub struct Peer { initialized: bool, } -impl<'r, S> simulator::Peer> for Peer +impl<'r, S> simulator::Peer> for Peer where S: WriteStorage<'r> + 'static, { @@ -46,7 +47,7 @@ where } impl Deref for Peer { - type Target = Protocol; + type Target = Transport; fn deref(&self) -> &Self::Target { &self.protocol @@ -92,7 +93,14 @@ where let local_time = LocalTime::now(); let clock = RefClock::from(local_time); let signer = MockSigner::new(&mut rng); - let protocol = Protocol::new(config, clock, storage, addrs, signer, rng.clone()); + let protocol = Transport::new(Protocol::new( + config, + clock, + storage, + addrs, + signer, + rng.clone(), + )); let ip = ip.into(); let local_addr = net::SocketAddr::new(ip, rng.u16(..)); @@ -135,7 +143,7 @@ where } pub fn connect_from(&mut self, peer: &Self) { - let remote = simulator::Peer::>::addr(peer); + let remote = simulator::Peer::>::addr(peer); let local = net::SocketAddr::new(self.ip, self.rng.u16(..)); let git = format!("file:///{}.git", remote.ip()); let git = Url::from_bytes(git.as_bytes()).unwrap(); @@ -160,7 +168,7 @@ where } pub fn connect_to(&mut self, peer: &Self) { - let remote = simulator::Peer::>::addr(peer); + let remote = simulator::Peer::>::addr(peer); self.initialize(); self.protocol.attempted(&remote); diff --git a/node/src/transport.rs b/node/src/transport.rs new file mode 100644 index 00000000..f08c5b8b --- /dev/null +++ b/node/src/transport.rs @@ -0,0 +1,106 @@ +use std::net; +use std::ops::{Deref, DerefMut}; + +use nakamoto::LocalTime; +use nakamoto_net as nakamoto; +use nakamoto_net::{Io, Link}; + +use crate::address_book; +use crate::collections::HashMap; +use crate::crypto; +use crate::protocol::{Command, DisconnectReason, Event, Protocol}; +use crate::storage::WriteStorage; + +#[derive(Debug)] +struct Peer { + addr: net::SocketAddr, +} + +#[derive(Debug)] +pub struct Transport { + peers: HashMap, + protocol: Protocol, +} + +impl Transport { + pub fn new(protocol: Protocol) -> Self { + Self { + peers: HashMap::default(), + protocol, + } + } +} + +impl<'r, S, T, G> nakamoto::Protocol for Transport +where + T: WriteStorage<'r> + 'static, + S: address_book::Store, + G: crypto::Signer, +{ + type Event = Event; + type Command = Command; + type DisconnectReason = DisconnectReason; + + fn initialize(&mut self, time: LocalTime) { + self.protocol.initialize(time) + } + + fn tick(&mut self, now: nakamoto::LocalTime) { + self.protocol.tick(now) + } + + fn wake(&mut self) { + self.protocol.wake() + } + + fn command(&mut self, cmd: Self::Command) { + self.protocol.command(cmd) + } + + fn attempted(&mut self, addr: &std::net::SocketAddr) { + self.protocol.attempted(addr) + } + + fn connected( + &mut self, + addr: std::net::SocketAddr, + local_addr: &std::net::SocketAddr, + link: Link, + ) { + self.protocol.connected(addr, local_addr, link) + } + + fn disconnected( + &mut self, + addr: &std::net::SocketAddr, + reason: nakamoto::DisconnectReason, + ) { + self.protocol.disconnected(addr, reason) + } + + fn received_bytes(&mut self, addr: &std::net::SocketAddr, bytes: &[u8]) { + self.protocol.received_bytes(addr, bytes) + } +} + +impl Iterator for Transport { + type Item = Io; + + fn next(&mut self) -> Option { + self.protocol.next() + } +} + +impl Deref for Transport { + type Target = Protocol; + + fn deref(&self) -> &Self::Target { + &self.protocol + } +} + +impl DerefMut for Transport { + fn deref_mut(&mut self) -> &mut Self::Target { + &mut self.protocol + } +}