node: Handle connection state discrepancies

There are cases where the service state doesn't match the state
of the underlying wire protocol. We remedy this by transitioning peers
to a "connected" state when a message is received, if they are in a
"connecting" state.

The long term solution to this will likely be to merge the service and
protocol layer so that there are no inconsistencies.

The other case happens when a persistent peer is in "disconnected" state
and attempts an inbound connection to us.
This commit is contained in:
cloudhead 2024-06-18 14:41:56 +02:00
parent a1c9a927c2
commit 47c9f792e6
No known key found for this signature in database
2 changed files with 58 additions and 38 deletions

View File

@ -1314,11 +1314,17 @@ where
}
} else {
match self.sessions.entry(remote) {
Entry::Occupied(e) => {
warn!(
Entry::Occupied(mut e) => {
// In this scenario, it's possible that our peer is persistent, and
// disconnected. We get an inbound connection before we attempt a re-connection,
// and therefore we treat it as a regular inbound connection.
let peer = e.get_mut();
debug!(
target: "service",
"Connecting peer {remote} already has a session open ({})", e.get()
"Connecting peer {remote} already has a session open ({peer})"
);
peer.to_connected(self.clock);
self.outbox.write_all(peer, msgs);
}
Entry::Vacant(e) => {
if let HostName::Ip(ip) = addr.host {
@ -1772,15 +1778,39 @@ where
}
message.log(log::Level::Debug, remote, Link::Inbound);
trace!(target: "service", "Received message {:?} from {}", &message, peer.id);
let connected = match &mut peer.state {
session::State::Disconnected { .. } => {
debug!(target: "service", "Ignoring message from disconnected peer {}", peer.id);
return Ok(());
}
// In case of a discrepancy between the service state and the state of the underlying
// wire protocol, we may receive a message from a peer that we consider not fully connected
// at the service level. To remedy this, we simply transition the peer to a connected state.
//
// This is not ideal, but until the wire protocol and service are unified, it's the simplest
// solution to converge towards the same state.
session::State::Attempted { .. } | session::State::Initial => {
debug!(target: "service", "Received unexpected message from connecting peer {}", peer.id);
debug!(target: "service", "Transitioning peer {} to 'connected' state", peer.id);
match (&mut peer.state, message) {
peer.to_connected(self.clock);
None
}
session::State::Connected {
ping, latencies, ..
} => Some((ping, latencies)),
};
trace!(target: "service", "Received message {message:?} from {remote}");
match message {
// Process a peer announcement.
(session::State::Connected { .. }, Message::Announcement(ann)) => {
let relayer = peer.id;
Message::Announcement(ann) => {
let relayer = remote;
let relayer_addr = peer.addr.clone();
if let Some(id) = self.handle_announcement(&relayer, &relayer_addr, &ann)? {
if let Some(id) = self.handle_announcement(relayer, &relayer_addr, &ann)? {
if self.config.is_relay() {
if let AnnouncementMessage::Inventory(_) = ann.message {
if let Err(e) = self
@ -1797,7 +1827,7 @@ where
}
}
}
(session::State::Connected { .. }, Message::Subscribe(subscribe)) => {
Message::Subscribe(subscribe) => {
// Filter announcements by interest.
match self
.db
@ -1829,11 +1859,10 @@ where
}
peer.subscribe = Some(subscribe);
}
(session::State::Connected { .. }, Message::Info(info)) => {
let remote = peer.id;
self.handle_info(remote, &info)?;
Message::Info(info) => {
self.handle_info(*remote, &info)?;
}
(session::State::Connected { .. }, Message::Ping(Ping { ponglen, .. })) => {
Message::Ping(Ping { ponglen, .. }) => {
// Ignore pings which ask for too much data.
if ponglen > Ping::MAX_PONG_ZEROES {
return Ok(());
@ -1845,33 +1874,24 @@ where
},
);
}
(
session::State::Connected {
ping, latencies, ..
},
Message::Pong { zeroes },
) => {
if let session::PingState::AwaitingResponse {
len: ponglen,
since,
} = *ping
{
if (ponglen as usize) == zeroes.len() {
*ping = session::PingState::Ok;
// Keep track of peer latency.
latencies.push_back(self.clock - since);
if latencies.len() > MAX_LATENCIES {
latencies.pop_front();
Message::Pong { zeroes } => {
if let Some((ping, latencies)) = connected {
if let session::PingState::AwaitingResponse {
len: ponglen,
since,
} = *ping
{
if (ponglen as usize) == zeroes.len() {
*ping = session::PingState::Ok;
// Keep track of peer latency.
latencies.push_back(self.clock - since);
if latencies.len() > MAX_LATENCIES {
latencies.pop_front();
}
}
}
}
}
(session::State::Attempted { .. } | session::State::Initial, msg) => {
debug!(target: "service", "Ignoring unexpected message {:?} from connecting peer {}", msg, peer.id);
}
(session::State::Disconnected { .. }, msg) => {
debug!(target: "service", "Ignoring {:?} from disconnected peer {}", msg, peer.id);
}
}
Ok(())
}

View File

@ -230,8 +230,8 @@ impl Session {
pub fn to_connected(&mut self, since: LocalTime) {
self.last_active = since;
let State::Attempted = &self.state else {
panic!("Session::to_connected: can only transition to 'connected' state from 'attempted' state");
if let State::Connected { .. } = &self.state {
panic!("Session::to_connected: session is already in 'connected' state");
};
self.state = State::Connected {
since,