Skip to content

Commit 90dab5f

Browse files
authored
More reliable peers check (mimblewimble#3824)
* peer: unknown state for new peers, check peers state on every monitor (128 healthy non-connected + 128 defuncts + 128 unknown), mark peer as defunct when ping not passed, do not crash on toml parse with dns failure * p2p: cleanup before selection at monitor, add outbound to connected list only when there is not enough peers + disconnect extra peer immediately, reconnect to seeds at monitor to avoid stuck, update only defunct state to unknown when received existing peer address * p2p: reduced amount of total peers to check at monitor * p2p: do not check healthy and defunct peers more often than once per hour, store last connection attempt, do not ask for more peers when there is enough outbound * peer: update last_attempt when changing peer state to other than Banned * fix: log of peers amount to check
1 parent af0c1dc commit 90dab5f

6 files changed

Lines changed: 187 additions & 136 deletions

File tree

p2p/src/peers.rs

Lines changed: 26 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,7 @@ impl Peers {
6161
/// Adds the peer to our internal peer mapping. Note that the peer is still
6262
/// returned so the server can run it.
6363
pub fn add_connected(&self, peer: Arc<Peer>) -> Result<(), Error> {
64+
let enough_outbound = self.enough_outbound_peers();
6465
let peer_data: PeerData;
6566
{
6667
// Scope for peers vector lock - dont hold the peers lock while adding to lmdb
@@ -76,9 +77,12 @@ impl Peers {
7677
last_banned: 0,
7778
ban_reason: ReasonForBan::None,
7879
last_connected: Utc::now().timestamp(),
80+
last_attempt: Utc::now().timestamp(),
7981
};
80-
debug!("Adding newly connected peer {}.", peer_data.addr);
81-
peers.insert(peer_data.addr, peer);
82+
if !enough_outbound || !peer.info.is_outbound() {
83+
debug!("Adding newly connected peer {}.", peer_data.addr);
84+
peers.insert(peer_data.addr, peer);
85+
}
8286
}
8387
debug!("Saving newly connected peer {}.", peer_data.addr);
8488
if let Err(e) = self.save_peer(&peer_data) {
@@ -98,6 +102,7 @@ impl Peers {
98102
last_banned: Utc::now().timestamp(),
99103
ban_reason,
100104
last_connected: Utc::now().timestamp(),
105+
last_attempt: Utc::now().timestamp(),
101106
};
102107
debug!("Banning peer {}.", addr);
103108
self.save_peer(&peer_data)
@@ -142,6 +147,7 @@ impl Peers {
142147
}
143148
false
144149
}
150+
145151
/// Ban a peer, disconnecting it if we're currently connected
146152
pub fn ban_peer(&self, peer_addr: PeerAddr, ban_reason: ReasonForBan) -> Result<(), Error> {
147153
// Update the peer in peers db
@@ -261,6 +267,8 @@ impl Peers {
261267
break;
262268
}
263269
};
270+
// Mark peer as defunct after ping failure.
271+
let _ = self.update_state(p.info.addr, State::Defunct);
264272
p.stop();
265273
peers.remove(&p.info.addr);
266274
}
@@ -702,21 +710,25 @@ impl NetAdapter for Peers {
702710
trace!("Received {} peer addrs, saving.", peer_addrs.len());
703711
let mut to_save: Vec<PeerData> = Vec::new();
704712
for pa in peer_addrs {
705-
if let Ok(e) = self.exists_peer(pa) {
706-
if e {
713+
if let Ok(mut p) = self.get_peer(pa) {
714+
if p.flags != State::Defunct {
707715
continue;
708716
}
717+
p.flags = State::Unknown;
718+
to_save.push(p);
719+
} else {
720+
let peer = PeerData {
721+
addr: pa,
722+
capabilities: Capabilities::UNKNOWN,
723+
user_agent: "".to_string(),
724+
flags: State::Unknown,
725+
last_banned: 0,
726+
ban_reason: ReasonForBan::None,
727+
last_connected: 0,
728+
last_attempt: 0,
729+
};
730+
to_save.push(peer);
709731
}
710-
let peer = PeerData {
711-
addr: pa,
712-
capabilities: Capabilities::UNKNOWN,
713-
user_agent: "".to_string(),
714-
flags: State::Healthy,
715-
last_banned: 0,
716-
ban_reason: ReasonForBan::None,
717-
last_connected: Utc::now().timestamp(),
718-
};
719-
to_save.push(peer);
720732
}
721733
if let Err(e) = self.save_peers(to_save) {
722734
error!("Could not save received peer addresses: {:?}", e);

p2p/src/serv.rs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -187,6 +187,9 @@ impl Server {
187187
&self.handshake,
188188
self.peers.clone(),
189189
)?;
190+
if self.peers.enough_outbound_peers() {
191+
peer.stop();
192+
}
190193
let peer = Arc::new(peer);
191194
self.peers.add_connected(peer.clone())?;
192195
Ok(peer)

p2p/src/store.rs

Lines changed: 15 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -27,13 +27,14 @@ const STORE_SUBPATH: &str = "peers";
2727

2828
const PEER_PREFIX: u8 = b'P';
2929

30-
// Types of messages
30+
// Types of peers
3131
enum_from_primitive! {
3232
#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize)]
3333
pub enum State {
3434
Healthy = 0,
3535
Banned = 1,
3636
Defunct = 2,
37+
Unknown = 3,
3738
}
3839
}
3940

@@ -55,6 +56,8 @@ pub struct PeerData {
5556
pub ban_reason: ReasonForBan,
5657
/// Time when we last connected to this peer.
5758
pub last_connected: i64,
59+
/// Time when last connection attempt happened to this peer.
60+
pub last_attempt: i64,
5861
}
5962

6063
impl Writeable for PeerData {
@@ -67,7 +70,8 @@ impl Writeable for PeerData {
6770
[write_u8, self.flags as u8],
6871
[write_i64, self.last_banned],
6972
[write_i32, self.ban_reason as i32],
70-
[write_i64, self.last_connected]
73+
[write_i64, self.last_connected],
74+
[write_i64, self.last_attempt]
7175
);
7276
Ok(())
7377
}
@@ -81,12 +85,10 @@ impl Readable for PeerData {
8185
let (fl, lb, br) = ser_multiread!(reader, read_u8, read_i64, read_i32);
8286

8387
let lc = reader.read_i64();
84-
// this only works because each PeerData is read in its own vector and this
85-
// is the last data element
86-
let last_connected = match lc {
87-
Err(_) => Utc::now().timestamp(),
88-
Ok(lc) => lc,
89-
};
88+
let last_connected = lc.unwrap_or_else(|_| Utc::now().timestamp());
89+
90+
let la = reader.read_i64();
91+
let last_attempt = la.unwrap_or_else(|_| 0);
9092

9193
let user_agent = String::from_utf8(ua).map_err(|_| ser::Error::CorruptedData)?;
9294
let capabilities = Capabilities::from_bits_truncate(capab);
@@ -97,10 +99,11 @@ impl Readable for PeerData {
9799
addr,
98100
capabilities,
99101
user_agent,
100-
flags: flags,
102+
flags,
101103
last_banned: lb,
102104
ban_reason,
103105
last_connected,
106+
last_attempt,
104107
}),
105108
None => Err(ser::Error::CorruptedData),
106109
}
@@ -187,6 +190,7 @@ impl PeerStore {
187190

188191
/// Convenience method to load a peer data, update its status and save it
189192
/// back. If new state is Banned its last banned time will be updated too.
193+
/// If new state is Defunct last connection attempt will be updated too.
190194
pub fn update_state(&self, peer_addr: PeerAddr, new_state: State) -> Result<(), Error> {
191195
let batch = self.db.batch()?;
192196

@@ -197,6 +201,8 @@ impl PeerStore {
197201
peer.flags = new_state;
198202
if new_state == State::Banned {
199203
peer.last_banned = Utc::now().timestamp();
204+
} else {
205+
peer.last_attempt = Utc::now().timestamp();
200206
}
201207

202208
batch.put_ser(&peer_key(peer_addr)[..], &peer)?;

p2p/src/types.rs

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -183,10 +183,13 @@ impl<'de> Visitor<'de> for PeerAddrs {
183183
Ok(ip) => peers.push(PeerAddr(ip)),
184184
// If that fails it's probably a DNS record
185185
Err(_) => {
186-
let socket_addrs = entry.to_socket_addrs().map_err(|_| {
187-
serde::de::Error::custom(format!("Unable to resolve DNS: {}", entry))
188-
})?;
189-
peers.append(&mut socket_addrs.map(PeerAddr).collect());
186+
let socket_addrs: Result<std::vec::IntoIter<SocketAddr>, M::Error> =
187+
entry.to_socket_addrs().map_err(|_| {
188+
serde::de::Error::custom(format!("Unable to resolve DNS: {}", entry))
189+
});
190+
if let Ok(socket_addrs) = socket_addrs {
191+
peers.append(&mut socket_addrs.map(PeerAddr).collect());
192+
}
190193
}
191194
}
192195
}

0 commit comments

Comments
 (0)