node: Upgrade `netservices` and `io-reactor`
There was an issue with the old version due to the use of `RawFd` as peer IDs, since they are not unique for the lifetime of a process. With this change, we use the new `ResourceId` type, though these are not immediately available, as they are generated when the resource is registered.
This commit is contained in:
parent
2e781b1efd
commit
25ca4c8b92
|
|
@ -80,9 +80,9 @@ checksum = "0942ffc6dcaadf03badf6e6a2d0228460359d5e34b57ccdc720b7382dfbd5ec5"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "amplify"
|
name = "amplify"
|
||||||
version = "4.0.0"
|
version = "4.5.0"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "f26966af46e0d200e8bf2b7f16230997c1c3f2d141bc27ccc091c012ed527b58"
|
checksum = "8629db306c0bbeb0a402e2918bdcf0026b5ddb24c46460f3bf5410b350d98710"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"amplify_derive",
|
"amplify_derive",
|
||||||
"amplify_num",
|
"amplify_num",
|
||||||
|
|
@ -92,9 +92,9 @@ dependencies = [
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "amplify_derive"
|
name = "amplify_derive"
|
||||||
version = "3.0.1"
|
version = "4.0.0"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "c87df0f28e6eb1f2d355f29ba6793fa9ca643967528609608d5cbd70bd68f9d1"
|
checksum = "759dcbfaf94d838367a86d493ec34ccc8aa6fe365cb7880d6bf89006de24d9c1"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"amplify_syn",
|
"amplify_syn",
|
||||||
"proc-macro2",
|
"proc-macro2",
|
||||||
|
|
@ -1677,9 +1677,9 @@ dependencies = [
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "io-reactor"
|
name = "io-reactor"
|
||||||
version = "0.2.1"
|
version = "0.3.0"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "850495a01d6f0b6d29adf5849a75b5a82449d9181bdd642cd681ce3d6f14b58c"
|
checksum = "2457e8fb1b1f298809fcd93cd15d485f52ef565f7ad47970583c7a80ae6c7e78"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"amplify",
|
"amplify",
|
||||||
"crossbeam-channel",
|
"crossbeam-channel",
|
||||||
|
|
@ -1748,9 +1748,9 @@ checksum = "478ee9e62aaeaf5b140bd4138753d1f109765488581444218d3ddda43234f3e8"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "libc"
|
name = "libc"
|
||||||
version = "0.2.150"
|
version = "0.2.152"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "89d92a4743f9a61002fae18374ed11e7973f530cb3a3255fb354818118b2203c"
|
checksum = "13e3bf6590cbc649f4d1a3eefc9d5d6eb746f5200ffb04e5e142700b8faa56e7"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "libgit2-sys"
|
name = "libgit2-sys"
|
||||||
|
|
@ -1908,9 +1908,9 @@ dependencies = [
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "netservices"
|
name = "netservices"
|
||||||
version = "0.4.0"
|
version = "0.5.0"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "ee13a6ce51c79cf719cea8a3f0d3584098e529c726634bc1197cbbb6967b7de2"
|
checksum = "e227ff39744a414d32b9b473266c1739a950b89cd03694a58b624f4bde926c10"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"amplify",
|
"amplify",
|
||||||
"cyphernet",
|
"cyphernet",
|
||||||
|
|
|
||||||
|
|
@ -19,12 +19,12 @@ colored = { version = "1.9.0" }
|
||||||
crossbeam-channel = { version = "0.5.6" }
|
crossbeam-channel = { version = "0.5.6" }
|
||||||
cyphernet = { version = "0.4.1", features = ["tor", "dns", "ed25519", "p2p-ed25519"] }
|
cyphernet = { version = "0.4.1", features = ["tor", "dns", "ed25519", "p2p-ed25519"] }
|
||||||
fastrand = { version = "2.0.0" }
|
fastrand = { version = "2.0.0" }
|
||||||
io-reactor = { version = "0.2.1", features = ["popol"] }
|
io-reactor = { version = "0.3.0", features = ["popol"] }
|
||||||
lexopt = { version = "0.2.1" }
|
lexopt = { version = "0.2.1" }
|
||||||
libc = { version = "0.2.137" }
|
libc = { version = "0.2.137" }
|
||||||
log = { version = "0.4.17", features = ["std"] }
|
log = { version = "0.4.17", features = ["std"] }
|
||||||
localtime = { version = "1.2.0" }
|
localtime = { version = "1.2.0" }
|
||||||
netservices = { version = "0.4.0", features = ["io-reactor", "socket2"] }
|
netservices = { version = "0.5.0", features = ["io-reactor", "socket2"] }
|
||||||
nonempty = { version = "0.8.1", features = ["serialize"] }
|
nonempty = { version = "0.8.1", features = ["serialize"] }
|
||||||
once_cell = { version = "1.13" }
|
once_cell = { version = "1.13" }
|
||||||
qcheck = { version = "1", default-features = false, optional = true }
|
qcheck = { version = "1", default-features = false, optional = true }
|
||||||
|
|
|
||||||
|
|
@ -4,8 +4,7 @@
|
||||||
//! The handshake itself is implemented in the external [`cyphernet`] and [`netservices`] crates.
|
//! The handshake itself is implemented in the external [`cyphernet`] and [`netservices`] crates.
|
||||||
use std::collections::hash_map::Entry;
|
use std::collections::hash_map::Entry;
|
||||||
use std::collections::VecDeque;
|
use std::collections::VecDeque;
|
||||||
use std::os::unix::io::AsRawFd;
|
use std::os::unix::io::{AsRawFd, RawFd};
|
||||||
use std::os::unix::prelude::RawFd;
|
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::{io, net, time};
|
use std::{io, net, time};
|
||||||
|
|
||||||
|
|
@ -19,7 +18,7 @@ use localtime::LocalTime;
|
||||||
use netservices::resource::{ListenerEvent, NetAccept, NetTransport, SessionEvent};
|
use netservices::resource::{ListenerEvent, NetAccept, NetTransport, SessionEvent};
|
||||||
use netservices::session::{ProtocolArtifact, Socks5Session};
|
use netservices::session::{ProtocolArtifact, Socks5Session};
|
||||||
use netservices::{NetConnection, NetProtocol, NetReader, NetWriter};
|
use netservices::{NetConnection, NetProtocol, NetReader, NetWriter};
|
||||||
use reactor::Timestamp;
|
use reactor::{ResourceId, Timestamp};
|
||||||
|
|
||||||
use radicle::collections::RandomMap;
|
use radicle::collections::RandomMap;
|
||||||
use radicle::node::NodeId;
|
use radicle::node::NodeId;
|
||||||
|
|
@ -141,15 +140,28 @@ impl Streams {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// The initial state of an outbound peer before handshake is completed.
|
||||||
|
#[derive(Debug)]
|
||||||
|
struct Outbound {
|
||||||
|
/// Resource ID, if registered.
|
||||||
|
id: Option<ResourceId>,
|
||||||
|
/// Remote address.
|
||||||
|
addr: NetAddr<HostName>,
|
||||||
|
/// Remote Node ID.
|
||||||
|
nid: NodeId,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The initial state of an inbound peer before handshake is completed.
|
||||||
|
#[derive(Debug)]
|
||||||
|
struct Inbound {
|
||||||
|
/// Resource ID, if registered.
|
||||||
|
id: Option<ResourceId>,
|
||||||
|
/// Remote address.
|
||||||
|
addr: NetAddr<HostName>,
|
||||||
|
}
|
||||||
|
|
||||||
/// Peer connection state machine.
|
/// Peer connection state machine.
|
||||||
enum Peer {
|
enum Peer {
|
||||||
/// The initial state of an inbound peer before handshake is completed.
|
|
||||||
Inbound { addr: NetAddr<HostName> },
|
|
||||||
/// The initial state of an outbound peer before handshake is completed.
|
|
||||||
Outbound {
|
|
||||||
addr: NetAddr<HostName>,
|
|
||||||
nid: NodeId,
|
|
||||||
},
|
|
||||||
/// The state after handshake is completed.
|
/// The state after handshake is completed.
|
||||||
/// Peers in this state are handled by the underlying service.
|
/// Peers in this state are handled by the underlying service.
|
||||||
Connected {
|
Connected {
|
||||||
|
|
@ -171,8 +183,6 @@ enum Peer {
|
||||||
impl std::fmt::Debug for Peer {
|
impl std::fmt::Debug for Peer {
|
||||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||||
match self {
|
match self {
|
||||||
Self::Inbound { addr } => write!(f, "Inbound({addr})"),
|
|
||||||
Self::Outbound { nid, .. } => write!(f, "Outbound({nid})"),
|
|
||||||
Self::Connected { link, nid, .. } => write!(f, "Connected({link:?}, {nid})"),
|
Self::Connected { link, nid, .. } => write!(f, "Connected({link:?}, {nid})"),
|
||||||
Self::Disconnecting { .. } => write!(f, "Disconnecting"),
|
Self::Disconnecting { .. } => write!(f, "Disconnecting"),
|
||||||
}
|
}
|
||||||
|
|
@ -183,58 +193,19 @@ impl Peer {
|
||||||
/// Return the peer's id, if any.
|
/// Return the peer's id, if any.
|
||||||
fn id(&self) -> Option<&NodeId> {
|
fn id(&self) -> Option<&NodeId> {
|
||||||
match self {
|
match self {
|
||||||
Peer::Outbound { nid, .. }
|
Peer::Connected { nid, .. } | Peer::Disconnecting { nid: Some(nid), .. } => Some(nid),
|
||||||
| Peer::Connected { nid, .. }
|
|
||||||
| Peer::Disconnecting { nid: Some(nid), .. } => Some(nid),
|
|
||||||
Peer::Inbound { .. } => None,
|
|
||||||
Peer::Disconnecting { nid: None, .. } => None,
|
Peer::Disconnecting { nid: None, .. } => None,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Return a new inbound connecting peer.
|
/// Connected peer.
|
||||||
fn inbound(addr: NetAddr<HostName>) -> Self {
|
fn connected(nid: NodeId, addr: NetAddr<HostName>, link: Link) -> Self {
|
||||||
Self::Inbound { addr }
|
Self::Connected {
|
||||||
}
|
link,
|
||||||
|
|
||||||
/// Return a new outbound connecting peer.
|
|
||||||
fn outbound(addr: NetAddr<HostName>, nid: NodeId) -> Self {
|
|
||||||
Self::Outbound { addr, nid }
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Switch to connected state.
|
|
||||||
fn connected(&mut self, nid: NodeId) -> (NetAddr<HostName>, Link) {
|
|
||||||
if let Self::Inbound { addr } = self {
|
|
||||||
let link = Link::Inbound;
|
|
||||||
let addr = addr.clone();
|
|
||||||
|
|
||||||
*self = Self::Connected {
|
|
||||||
link,
|
|
||||||
addr: addr.clone(),
|
|
||||||
nid,
|
|
||||||
inbox: Deserializer::default(),
|
|
||||||
streams: Streams::new(link),
|
|
||||||
};
|
|
||||||
(addr, link)
|
|
||||||
} else if let Self::Outbound {
|
|
||||||
addr,
|
addr,
|
||||||
nid: expected,
|
nid,
|
||||||
} = self
|
inbox: Deserializer::default(),
|
||||||
{
|
streams: Streams::new(link),
|
||||||
assert_eq!(nid, *expected);
|
|
||||||
|
|
||||||
let link = Link::Outbound;
|
|
||||||
let addr = addr.clone();
|
|
||||||
|
|
||||||
*self = Self::Connected {
|
|
||||||
link,
|
|
||||||
addr: addr.clone(),
|
|
||||||
nid,
|
|
||||||
inbox: Deserializer::default(),
|
|
||||||
streams: Streams::new(link),
|
|
||||||
};
|
|
||||||
(addr, link)
|
|
||||||
} else {
|
|
||||||
panic!("Peer::connected: session for {nid} is already established");
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -243,13 +214,6 @@ impl Peer {
|
||||||
if let Self::Connected { nid, streams, .. } = self {
|
if let Self::Connected { nid, streams, .. } = self {
|
||||||
streams.shutdown();
|
streams.shutdown();
|
||||||
|
|
||||||
*self = Self::Disconnecting {
|
|
||||||
nid: Some(*nid),
|
|
||||||
reason,
|
|
||||||
};
|
|
||||||
} else if let Self::Inbound { .. } = self {
|
|
||||||
*self = Self::Disconnecting { nid: None, reason };
|
|
||||||
} else if let Self::Outbound { nid, .. } = self {
|
|
||||||
*self = Self::Disconnecting {
|
*self = Self::Disconnecting {
|
||||||
nid: Some(*nid),
|
nid: Some(*nid),
|
||||||
reason,
|
reason,
|
||||||
|
|
@ -260,54 +224,53 @@ impl Peer {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
struct Peers(RandomMap<RawFd, Peer>);
|
/// Holds connected peers.
|
||||||
|
struct Peers(RandomMap<ResourceId, Peer>);
|
||||||
|
|
||||||
impl Peers {
|
impl Peers {
|
||||||
fn get_mut(&mut self, fd: &RawFd) -> Option<&mut Peer> {
|
fn get_mut(&mut self, id: &ResourceId) -> Option<&mut Peer> {
|
||||||
self.0.get_mut(fd)
|
self.0.get_mut(id)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn entry(&mut self, fd: RawFd) -> Entry<RawFd, Peer> {
|
fn entry(&mut self, id: ResourceId) -> Entry<ResourceId, Peer> {
|
||||||
self.0.entry(fd)
|
self.0.entry(id)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn insert(&mut self, fd: RawFd, peer: Peer) {
|
fn insert(&mut self, id: ResourceId, peer: Peer) {
|
||||||
if self.0.insert(fd, peer).is_some() {
|
if self.0.insert(id, peer).is_some() {
|
||||||
log::warn!(target: "wire", "Replacing existing peer fd={fd}");
|
log::warn!(target: "wire", "Replacing existing peer id={id}");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn remove(&mut self, fd: &RawFd) -> Option<Peer> {
|
fn remove(&mut self, id: &ResourceId) -> Option<Peer> {
|
||||||
self.0.remove(fd)
|
self.0.remove(id)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn lookup(&self, node_id: &NodeId) -> Option<(RawFd, &Peer)> {
|
fn lookup(&self, node_id: &NodeId) -> Option<(ResourceId, &Peer)> {
|
||||||
self.0
|
self.0
|
||||||
.iter()
|
.iter()
|
||||||
.find(|(_, peer)| peer.id() == Some(node_id))
|
.find(|(_, peer)| peer.id() == Some(node_id))
|
||||||
.map(|(fd, peer)| (*fd, peer))
|
.map(|(fd, peer)| (*fd, peer))
|
||||||
}
|
}
|
||||||
|
|
||||||
fn lookup_mut(&mut self, node_id: &NodeId) -> Option<(RawFd, &mut Peer)> {
|
fn lookup_mut(&mut self, node_id: &NodeId) -> Option<(ResourceId, &mut Peer)> {
|
||||||
self.0
|
self.0
|
||||||
.iter_mut()
|
.iter_mut()
|
||||||
.find(|(_, peer)| peer.id() == Some(node_id))
|
.find(|(_, peer)| peer.id() == Some(node_id))
|
||||||
.map(|(fd, peer)| (*fd, peer))
|
.map(|(fd, peer)| (*fd, peer))
|
||||||
}
|
}
|
||||||
|
|
||||||
fn active(&self) -> impl Iterator<Item = (RawFd, &NodeId)> {
|
fn active(&self) -> impl Iterator<Item = (ResourceId, &NodeId)> {
|
||||||
self.0.iter().filter_map(|(fd, peer)| match peer {
|
self.0.iter().filter_map(|(id, peer)| match peer {
|
||||||
Peer::Inbound { .. } => None,
|
Peer::Connected { nid, .. } => Some((*id, nid)),
|
||||||
Peer::Outbound { nid, .. } => Some((*fd, nid)),
|
|
||||||
Peer::Connected { nid, .. } => Some((*fd, nid)),
|
|
||||||
Peer::Disconnecting { .. } => None,
|
Peer::Disconnecting { .. } => None,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
fn connected(&self) -> impl Iterator<Item = (RawFd, &NodeId)> {
|
fn connected(&self) -> impl Iterator<Item = (ResourceId, &NodeId)> {
|
||||||
self.0.iter().filter_map(|(fd, peer)| {
|
self.0.iter().filter_map(|(id, peer)| {
|
||||||
if let Peer::Connected { nid: id, .. } = peer {
|
if let Peer::Connected { nid, .. } = peer {
|
||||||
Some((*fd, id))
|
Some((*id, nid))
|
||||||
} else {
|
} else {
|
||||||
None
|
None
|
||||||
}
|
}
|
||||||
|
|
@ -325,7 +288,11 @@ pub struct Wire<D, S, G: Signer + Ecdh> {
|
||||||
signer: G,
|
signer: G,
|
||||||
/// Internal queue of actions to send to the reactor.
|
/// Internal queue of actions to send to the reactor.
|
||||||
actions: VecDeque<Action<G>>,
|
actions: VecDeque<Action<G>>,
|
||||||
/// Peer sessions.
|
/// Outbound attempted peers without a session.
|
||||||
|
outbound: RandomMap<RawFd, Outbound>,
|
||||||
|
/// Inbound peers without a session.
|
||||||
|
inbound: RandomMap<RawFd, Inbound>,
|
||||||
|
/// Peer (established) sessions.
|
||||||
peers: Peers,
|
peers: Peers,
|
||||||
/// SOCKS5 proxy address.
|
/// SOCKS5 proxy address.
|
||||||
proxy: net::SocketAddr,
|
proxy: net::SocketAddr,
|
||||||
|
|
@ -351,6 +318,8 @@ where
|
||||||
signer,
|
signer,
|
||||||
proxy,
|
proxy,
|
||||||
actions: VecDeque::new(),
|
actions: VecDeque::new(),
|
||||||
|
inbound: RandomMap::default(),
|
||||||
|
outbound: RandomMap::default(),
|
||||||
peers: Peers(RandomMap::default()),
|
peers: Peers(RandomMap::default()),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -359,19 +328,19 @@ where
|
||||||
self.actions.push_back(Action::RegisterListener(socket));
|
self.actions.push_back(Action::RegisterListener(socket));
|
||||||
}
|
}
|
||||||
|
|
||||||
fn disconnect(&mut self, fd: RawFd, reason: DisconnectReason) {
|
fn disconnect(&mut self, id: ResourceId, reason: DisconnectReason) {
|
||||||
match self.peers.get_mut(&fd) {
|
match self.peers.get_mut(&id) {
|
||||||
Some(Peer::Disconnecting { .. }) => {
|
Some(Peer::Disconnecting { .. }) => {
|
||||||
log::error!(target: "wire", "Peer (fd={fd}) is already disconnecting");
|
log::error!(target: "wire", "Peer (id={id}) is already disconnecting");
|
||||||
}
|
}
|
||||||
Some(peer) => {
|
Some(peer) => {
|
||||||
log::debug!(target: "wire", "Disconnecting peer (fd={fd}): {reason}");
|
log::debug!(target: "wire", "Disconnecting peer (id={id}): {reason}");
|
||||||
|
|
||||||
peer.disconnecting(reason);
|
peer.disconnecting(reason);
|
||||||
self.actions.push_back(Action::UnregisterTransport(fd));
|
self.actions.push_back(Action::UnregisterTransport(id));
|
||||||
}
|
}
|
||||||
None => {
|
None => {
|
||||||
log::error!(target: "wire", "Unknown peer (fd={fd}) cannot be disconnected");
|
log::error!(target: "wire", "Unknown peer (id={id}) cannot be disconnected");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -467,22 +436,20 @@ where
|
||||||
|
|
||||||
fn handle_listener_event(
|
fn handle_listener_event(
|
||||||
&mut self,
|
&mut self,
|
||||||
_sock: net::SocketAddr,
|
_: ResourceId, // Nb. This is the ID of the listener socket.
|
||||||
event: ListenerEvent<WireSession<G>>,
|
event: ListenerEvent<WireSession<G>>,
|
||||||
_: Timestamp,
|
_: Timestamp,
|
||||||
) {
|
) {
|
||||||
match event {
|
match event {
|
||||||
ListenerEvent::Accepted(connection) => {
|
ListenerEvent::Accepted(connection) => {
|
||||||
let addr = connection.remote_addr();
|
let addr = connection.remote_addr();
|
||||||
log::debug!(target: "wire", "Accepting inbound peer connection from {addr}..");
|
let fd = connection.as_raw_fd();
|
||||||
|
log::debug!(target: "wire", "Accepting inbound connection from {addr} (fd={fd})..");
|
||||||
self.peers
|
|
||||||
.insert(connection.as_raw_fd(), Peer::inbound(addr.clone().into()));
|
|
||||||
|
|
||||||
// If the service doesn't want to accept this connection,
|
// If the service doesn't want to accept this connection,
|
||||||
// we drop the connection here, which disconnects the socket.
|
// we drop the connection here, which disconnects the socket.
|
||||||
if !self.service.accepted(NetAddr::from(addr.clone()).into()) {
|
if !self.service.accepted(NetAddr::from(addr.clone()).into()) {
|
||||||
log::debug!(target: "wire", "Dropping inbound connection from {addr}..");
|
log::debug!(target: "wire", "Rejecting inbound connection from {addr} (fd={fd})..");
|
||||||
drop(connection);
|
drop(connection);
|
||||||
|
|
||||||
return;
|
return;
|
||||||
|
|
@ -496,6 +463,14 @@ where
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
|
self.inbound.insert(
|
||||||
|
fd,
|
||||||
|
Inbound {
|
||||||
|
id: None,
|
||||||
|
addr: addr.into(),
|
||||||
|
},
|
||||||
|
);
|
||||||
self.actions
|
self.actions
|
||||||
.push_back(reactor::Action::RegisterTransport(transport))
|
.push_back(reactor::Action::RegisterTransport(transport))
|
||||||
}
|
}
|
||||||
|
|
@ -505,45 +480,62 @@ where
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn handle_registered(&mut self, fd: RawFd, id: ResourceId) {
|
||||||
|
if let Some(outbound) = self.outbound.get_mut(&fd) {
|
||||||
|
log::debug!(target: "wire", "Outbound peer resource registered with id={id} (fd={fd})");
|
||||||
|
outbound.id = Some(id);
|
||||||
|
} else if let Some(inbound) = self.inbound.get_mut(&fd) {
|
||||||
|
log::debug!(target: "wire", "Inbound peer resource registered with id={id} (fd={fd})");
|
||||||
|
inbound.id = Some(id);
|
||||||
|
} else {
|
||||||
|
log::error!(target: "wire", "Unknown peer registered with fd={fd} and id={id}");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fn handle_transport_event(
|
fn handle_transport_event(
|
||||||
&mut self,
|
&mut self,
|
||||||
fd: RawFd,
|
id: ResourceId,
|
||||||
event: SessionEvent<WireSession<G>>,
|
event: SessionEvent<WireSession<G>>,
|
||||||
_: Timestamp,
|
_: Timestamp,
|
||||||
) {
|
) {
|
||||||
match event {
|
match event {
|
||||||
SessionEvent::Established(ProtocolArtifact { state, .. }) => {
|
SessionEvent::Established(fd, ProtocolArtifact { state, .. }) => {
|
||||||
// SAFETY: With the NoiseXK protocol, there is always a remote static key.
|
// SAFETY: With the NoiseXK protocol, there is always a remote static key.
|
||||||
let id: NodeId = state.remote_static_key.unwrap();
|
let nid: NodeId = state.remote_static_key.unwrap();
|
||||||
|
|
||||||
log::debug!(target: "wire", "Session established with {id} (fd={fd})");
|
log::debug!(target: "wire", "Session established with {nid} (id={id}) (fd={fd})");
|
||||||
|
|
||||||
let conflicting = self
|
let conflicting = self
|
||||||
.peers
|
.peers
|
||||||
.active()
|
.active()
|
||||||
.filter(|(other, d)| **d == id && *other != fd)
|
.filter(|(other, d)| **d == nid && *other != id)
|
||||||
.map(|(fd, _)| fd)
|
.map(|(id, _)| id)
|
||||||
.collect::<Vec<_>>();
|
.collect::<Vec<_>>();
|
||||||
|
|
||||||
for fd in conflicting {
|
for id in conflicting {
|
||||||
log::warn!(
|
log::warn!(
|
||||||
target: "wire", "Closing conflicting session with {id} (fd={fd})"
|
target: "wire", "Closing conflicting session with {nid} (id={id})"
|
||||||
);
|
);
|
||||||
self.disconnect(
|
self.disconnect(
|
||||||
fd,
|
id,
|
||||||
DisconnectReason::Dial(Arc::new(io::Error::from(
|
DisconnectReason::Dial(Arc::new(io::Error::from(
|
||||||
io::ErrorKind::AlreadyExists,
|
io::ErrorKind::AlreadyExists,
|
||||||
))),
|
))),
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
let Some(peer) = self.peers.get_mut(&fd) else {
|
let (addr, link) = if let Some(peer) = self.inbound.remove(&fd) {
|
||||||
log::error!(target: "wire", "Session not found for fd {fd}");
|
(peer.addr, Link::Inbound)
|
||||||
|
} else if let Some(peer) = self.outbound.remove(&fd) {
|
||||||
|
assert_eq!(nid, peer.nid);
|
||||||
|
(peer.addr, Link::Outbound)
|
||||||
|
} else {
|
||||||
|
log::error!(target: "wire", "Session for {nid} (id={id}) not found");
|
||||||
return;
|
return;
|
||||||
};
|
};
|
||||||
let (addr, link) = peer.connected(id);
|
self.peers
|
||||||
|
.insert(id, Peer::connected(nid, addr.clone(), link));
|
||||||
self.service.connected(id, addr.into(), link);
|
self.service.connected(nid, addr.into(), link);
|
||||||
}
|
}
|
||||||
SessionEvent::Data(data) => {
|
SessionEvent::Data(data) => {
|
||||||
if let Some(Peer::Connected {
|
if let Some(Peer::Connected {
|
||||||
|
|
@ -551,7 +543,7 @@ where
|
||||||
inbox,
|
inbox,
|
||||||
streams,
|
streams,
|
||||||
..
|
..
|
||||||
}) = self.peers.get_mut(&fd)
|
}) = self.peers.get_mut(&id)
|
||||||
{
|
{
|
||||||
inbox.input(&data);
|
inbox.input(&data);
|
||||||
|
|
||||||
|
|
@ -631,7 +623,7 @@ where
|
||||||
log::debug!(target: "wire", "Dropping read buffer for {nid} with {} bytes", inbox.unparsed().count());
|
log::debug!(target: "wire", "Dropping read buffer for {nid} with {} bytes", inbox.unparsed().count());
|
||||||
}
|
}
|
||||||
self.disconnect(
|
self.disconnect(
|
||||||
fd,
|
id,
|
||||||
DisconnectReason::Session(session::Error::Misbehavior),
|
DisconnectReason::Session(session::Error::Misbehavior),
|
||||||
);
|
);
|
||||||
break;
|
break;
|
||||||
|
|
@ -639,11 +631,11 @@ where
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
log::warn!(target: "wire", "Dropping message from unconnected peer (fd={fd})");
|
log::warn!(target: "wire", "Dropping message from unconnected peer (id={id})");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
SessionEvent::Terminated(err) => {
|
SessionEvent::Terminated(err) => {
|
||||||
self.disconnect(fd, DisconnectReason::Connection(Arc::new(err)));
|
self.disconnect(id, DisconnectReason::Connection(Arc::new(err)));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -702,36 +694,35 @@ where
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn handover_listener(&mut self, _listener: Self::Listener) {
|
fn handover_listener(&mut self, _id: ResourceId, _listener: Self::Listener) {
|
||||||
panic!("Wire::handover_listener: listener handover is not supported");
|
panic!("Wire::handover_listener: listener handover is not supported");
|
||||||
}
|
}
|
||||||
|
|
||||||
fn handover_transport(&mut self, transport: Self::Transport) {
|
fn handover_transport(&mut self, id: ResourceId, transport: Self::Transport) {
|
||||||
let fd = transport.as_raw_fd();
|
let fd = transport.as_raw_fd();
|
||||||
log::debug!(target: "wire", "Received transport handover (fd={fd})");
|
|
||||||
|
|
||||||
match self.peers.entry(fd) {
|
match self.peers.entry(id) {
|
||||||
Entry::Occupied(e) => {
|
Entry::Occupied(e) => {
|
||||||
match e.get() {
|
match e.get() {
|
||||||
Peer::Disconnecting {
|
Peer::Disconnecting { nid, reason, .. } => {
|
||||||
nid: id, reason, ..
|
log::debug!(target: "wire", "Received transport handover for disconnecting peer with id={id} (fd={fd})");
|
||||||
} => {
|
|
||||||
// Disconnect TCP stream.
|
// Disconnect TCP stream.
|
||||||
drop(transport);
|
drop(transport);
|
||||||
|
|
||||||
// If there is no ID, the service is not aware of the peer.
|
// If there is no NID, the service is not aware of the peer.
|
||||||
if let Some(id) = id {
|
if let Some(nid) = nid {
|
||||||
self.service.disconnected(*id, reason);
|
self.service.disconnected(*nid, reason);
|
||||||
}
|
}
|
||||||
e.remove();
|
e.remove();
|
||||||
}
|
}
|
||||||
_ => {
|
Peer::Connected { nid, .. } => {
|
||||||
panic!("Wire::handover_transport: Unexpected peer with fd {fd} handed over from the reactor");
|
panic!("Wire::handover_transport: Unexpected handover of connected peer {} with id={id} (fd={fd})", nid);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Entry::Vacant(_) => {
|
Entry::Vacant(_) => {
|
||||||
panic!("Wire::handover_transport: Unknown peer with fd {fd} handed over");
|
panic!("Wire::handover_transport: Unknown peer with id={id} (fd={fd}) handed over");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -795,11 +786,13 @@ where
|
||||||
}) {
|
}) {
|
||||||
Ok(transport) => {
|
Ok(transport) => {
|
||||||
self.service.attempted(node_id, addr.clone());
|
self.service.attempted(node_id, addr.clone());
|
||||||
// TODO: Keep track of peer address for when peer disconnects before
|
self.outbound.insert(
|
||||||
// handshake is complete.
|
|
||||||
self.peers.insert(
|
|
||||||
transport.as_raw_fd(),
|
transport.as_raw_fd(),
|
||||||
Peer::outbound(addr.to_inner(), node_id),
|
Outbound {
|
||||||
|
id: None,
|
||||||
|
nid: node_id,
|
||||||
|
addr: addr.to_inner(),
|
||||||
|
},
|
||||||
);
|
);
|
||||||
self.actions
|
self.actions
|
||||||
.push_back(reactor::Action::RegisterTransport(transport));
|
.push_back(reactor::Action::RegisterTransport(transport));
|
||||||
|
|
@ -813,8 +806,8 @@ where
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Io::Disconnect(nid, reason) => {
|
Io::Disconnect(nid, reason) => {
|
||||||
if let Some((fd, Peer::Connected { .. })) = self.peers.lookup(&nid) {
|
if let Some((id, Peer::Connected { .. })) = self.peers.lookup(&nid) {
|
||||||
self.disconnect(fd, reason);
|
self.disconnect(id, reason);
|
||||||
} else {
|
} else {
|
||||||
log::warn!(target: "wire", "Peer {nid} is not connected: ignoring disconnect");
|
log::warn!(target: "wire", "Peer {nid} is not connected: ignoring disconnect");
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue