Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions changelog.d/ibd-track-retain.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
Changed

- **Satisfied-block pruning judges each in-flight hash once.** The reader
set is then replaced with that decision, so a confirm between two passes
cannot leave the initial-block-download reader and the assign loop
disagreeing about the next body.
16 changes: 16 additions & 0 deletions changelog.d/peer-lifecycle.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
Security

- Inbound eviction keeps a share of the longest-connected peers and
disconnects the newest peer in the largest netgroup. The netgroup is
fixed when the peer is accepted.
- A misbehavior disconnect refuses that address for one day, in memory
only. Rate-limit, oversize, and score-threshold exits record the same
refusal. A loopback peer is disconnected and is not recorded, so one
local failure does not block every other local connection. During
initial download, that death cools the dial even after a block body. A
netgroup that just lost an inbound slot waits ten minutes. The set
does not grow past its cap.
- During initial download, only a block this node requested moves the
stall clock or is queued. Other frames are rate-limited. Light
decodes do not wait on the reader.
- Findings write-ups: 060, 061, 062.
5 changes: 5 additions & 0 deletions changelog.d/sh-empty-ingest-seal.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
Fixed

- **`--sh-index` on an empty chain seals scripthash shards without reading
the empty ingest table.** That walk is 2^25 slots per shard and was
holding tip entry past the RPC cookie window.
95 changes: 91 additions & 4 deletions crates/rbitcoin-net/src/eviction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@ const PROTECT_NETGROUP: usize = 4;
const PROTECT_BLOCKS: usize = 4;
const PROTECT_TXS: usize = 4;
const PROTECT_MINPING: usize = 8;
/// Longest-connected inbound peers kept when slots are full.
const PROTECT_LONGEST: usize = 8;

/// Pick one inbound id to disconnect, or `None` if every candidate is protected.
pub fn select_inbound_eviction(mut cands: Vec<InboundEvictCandidate>) -> Option<u64> {
Expand Down Expand Up @@ -60,13 +62,37 @@ pub fn select_inbound_eviction(mut cands: Vec<InboundEvictCandidate>) -> Option<
return None;
}

// Prefer the longest-connected remaining peer (stable id tie-break).
// A share of the longest-connected peers stays. Evicting them is how a
// new inbound replaces the peers that have been useful the longest.
cands.sort_by(|a, b| {
a.connected_at
.cmp(&b.connected_at)
.then_with(|| a.id.cmp(&b.id))
});
Some(cands[0].id)
remove_first_k(&mut cands, PROTECT_LONGEST);
if cands.is_empty() {
return None;
}

// Newest peer in the largest netgroup. Group ids are the integers stored
// at accept; this compares those integers and does not read asmap.
let mut counts = std::collections::HashMap::<u64, usize>::new();
for c in &cands {
*counts.entry(c.netgroup).or_insert(0) += 1;
}
let (group, _) = counts
.into_iter()
.max_by(|a, b| a.1.cmp(&b.1).then_with(|| b.0.cmp(&a.0)))
.expect("at least one candidate");
cands
.iter()
.filter(|c| c.netgroup == group)
.max_by(|a, b| {
a.connected_at
.cmp(&b.connected_at)
.then_with(|| a.id.cmp(&b.id))
})
.map(|c| c.id)
}

fn ping_key(min_ping: Option<f64>) -> f64 {
Expand Down Expand Up @@ -134,6 +160,65 @@ mod tests {
}
}

#[test]
fn eviction_drops_the_newest_in_the_largest_netgroup() {
// Low ids are one-peer groups with a better ping, so the block, tx,
// ping, and netgroup protects consume them. The interesting peers are
// a size-3 group and one newer peer alone in another group.
let mut cands = Vec::new();
for i in 1..=40 {
cands.push(InboundEvictCandidate {
id: i,
connected_at: 1,
min_ping: Some(0.01),
last_block: 0,
last_tx: 0,
netgroup: 1_000 + i,
noban: false,
});
}
cands.push(InboundEvictCandidate {
id: 100,
connected_at: 10,
min_ping: Some(1.0),
last_block: 0,
last_tx: 0,
netgroup: 7,
noban: false,
});
cands.push(InboundEvictCandidate {
id: 101,
connected_at: 50,
min_ping: Some(1.0),
last_block: 0,
last_tx: 0,
netgroup: 7,
noban: false,
});
cands.push(InboundEvictCandidate {
id: 102,
connected_at: 200,
min_ping: Some(1.0),
last_block: 0,
last_tx: 0,
netgroup: 7,
noban: false,
});
cands.push(InboundEvictCandidate {
id: 103,
connected_at: 500,
min_ping: Some(1.0),
last_block: 0,
last_tx: 0,
netgroup: 9_000,
noban: false,
});
let victim = select_inbound_eviction(cands).expect("one inbound to evict");
assert_ne!(victim, 101, "the oldest peer in the largest group stays");
assert_ne!(victim, 103, "a newer peer in a smaller group stays");
assert_eq!(victim, 102, "evict the newest peer in the largest netgroup");
}

#[test]
fn eviction_protects_block_tx_ping_and_netgroup() {
// 4 block + 5 slow + 4 tx + 8 fast = 21; after protects, one slow remains.
Expand All @@ -150,8 +235,10 @@ mod tests {
for i in 13..21 {
cands.push(cand(i, 400 + i, Some(0.01), 0, 0));
}
let victim = select_inbound_eviction(cands).expect("one unprotected slow");
assert!((4..9).contains(&victim), "victim={victim}");
assert!(
select_inbound_eviction(cands).is_none(),
"block, tx, ping, netgroup, and longest-connected protects cover this set"
);
}

#[test]
Expand Down
12 changes: 8 additions & 4 deletions crates/rbitcoin-net/src/ibd/assign.rs
Original file line number Diff line number Diff line change
Expand Up @@ -161,7 +161,7 @@ pub(crate) fn clear_hash_inflight(
) {
inflight.remove(&hash);
for s in slots.iter_mut() {
s.in_flight.remove(&hash);
s.track_remove(&hash);
}
}

Expand All @@ -173,7 +173,7 @@ pub(crate) fn prune_satisfied_inflight(
) {
inflight.retain(|h, _| !hub.has_block(h));
for s in slots.iter_mut() {
s.in_flight.retain(|h| !hub.has_block(h));
s.track_retain(|h| !hub.has_block(h));
}
}

Expand Down Expand Up @@ -759,7 +759,7 @@ pub(crate) fn issue_batch(
}
let empty = st.slots[idx].in_flight.is_empty();
for &h in &batch {
st.slots[idx].in_flight.insert(h);
st.slots[idx].track_insert(h);
}
if empty {
st.slots[idx].rate.note_work_started(ibd_mono_ms());
Expand Down Expand Up @@ -1391,6 +1391,7 @@ pub(crate) fn cover_tip_holes(

#[cfg(test)]
pub(in crate::ibd) mod tests {
use super::super::peer_io::solicit_track;
use super::super::status::LoopStats;
use super::*;
use bitcoin::hashes::Hash;
Expand Down Expand Up @@ -1443,6 +1444,9 @@ pub(in crate::ibd) mod tests {
)),
cmd_tx,
in_flight: HashSet::new(),
requested: solicit_track().0,
solicited_bytes: solicit_track().1,
solicited_ms: solicit_track().2,
peer_height: 100,
connected_ms: 1,
first_data_ms: 0,
Expand Down Expand Up @@ -1517,7 +1521,7 @@ pub(in crate::ibd) mod tests {
let _ = getdata_asks(wire, BlockHash::all_zeros());
let mut slots = std::mem::take(&mut st.slots);
for s in &mut slots {
s.in_flight.clear();
s.track_clear();
s.rate = Default::default();
s.alive = alive.contains(&s.id);
}
Expand Down
52 changes: 47 additions & 5 deletions crates/rbitcoin-net/src/ibd/dial.rs
Original file line number Diff line number Diff line change
Expand Up @@ -462,7 +462,7 @@ pub(crate) fn release_peer_block_work(
let mut freed = Vec::new();
if let Some(s) = slots.iter_mut().find(|s| s.id == peer) {
s.alive = false;
for h in s.in_flight.drain() {
for h in s.track_drain() {
let empty = inflight
.get_mut(&h)
.map(|e| e.remove_peer(peer))
Expand Down Expand Up @@ -607,6 +607,17 @@ pub(crate) fn note_dead_without_block_bytes(
addr_cooldown.insert(addr, now + STALL_ADDR_COOLDOWN);
}

/// Misbehavior-threshold death. Cools the dial even after a block body was counted.
pub(crate) fn note_misbehavior_dead(
book: &mut AddrMan,
addr_cooldown: &mut HashMap<SocketAddr, Instant>,
addr: SocketAddr,
now: Instant,
) {
book.note_connect_failed(addr, false);
addr_cooldown.insert(addr, now + STALL_ADDR_COOLDOWN);
}

/// One stall rule: if a peer has outstanding block getdata and no **block**
/// progress for `stall`, disconnect it and free its work for reassignment.
///
Expand Down Expand Up @@ -650,16 +661,23 @@ pub(crate) fn disconnect_stalled_block_peers_at(
let mut freed = Vec::new();
let stall = stall.max(Duration::from_secs(30));
let stall_ms = stall.as_millis() as u64;
let stalled_peers: Vec<(usize, usize, SocketAddr)> = slots
let stalled_peers: Vec<(usize, usize, SocketAddr, u64)> = slots
.iter()
.filter(|s| s.alive && !s.in_flight.is_empty())
.filter(|s| s.rate.stalled(now_ms, stall_ms, true))
.map(|s| (s.id, s.in_flight.len(), s.addr))
.map(|s| {
(
s.id,
s.in_flight.len(),
s.addr,
s.bytes_rx_total.load(Ordering::Relaxed),
)
})
.collect();
for (id, n_work, addr) in stalled_peers {
for (id, n_work, addr, stream) in stalled_peers {
let cool = record_stall_kick(addr_cooldown, addr_strikes, addr, now);
warn!(
"ibd: peer[{id}] {addr} stalled (no block progress for {stall:?}, {n_work} in-flight) — disconnect + reassign (cooldown {cool:?})"
"ibd: peer[{id}] {addr} stalled (no block progress for {stall:?}, {n_work} in-flight, stream={stream}) — disconnect + reassign (cooldown {cool:?})"
);
if let Some(s) = slots.iter_mut().find(|s| s.id == id) {
let _ = s.cmd_tx.send(PeerCmd::Shutdown);
Expand Down Expand Up @@ -772,6 +790,7 @@ pub(crate) fn disconnect_relative_slow_block_peers_at(

#[cfg(test)]
mod tests {
use super::super::peer_io::solicit_track;
use super::*;
use bitcoin::hashes::Hash;
use bitcoin::BlockHash;
Expand Down Expand Up @@ -808,6 +827,9 @@ mod tests {
net: crate::NetAddr::from_socket(a),
cmd_tx,
in_flight: HashSet::new(),
requested: solicit_track().0,
solicited_bytes: solicit_track().1,
solicited_ms: solicit_track().2,
peer_height: 0,
connected_ms: 0,
first_data_ms: 0,
Expand Down Expand Up @@ -975,6 +997,26 @@ mod tests {
assert!(!cooldown.contains_key(&good));
}

#[test]
fn misbehavior_dead_cools_after_a_block_body() {
let mut book = AddrMan::new();
let lemon = addr(6);
book.note_connected(lemon);
let mut cooldown = HashMap::new();
let now = Instant::now();
note_misbehavior_dead(&mut book, &mut cooldown, lemon, now);
assert!(
book.flags(&lemon).failed_last_connect(),
"a misbehavior death is a failed connect"
);
assert!(cooldown.contains_key(&lemon));
let blocked = dial_blocked_addrs(&[], &cooldown, now);
assert!(
blocked.contains(&crate::NetAddr::Ip(lemon)),
"the dial stays cooled after a body was counted"
);
}

fn samp(id: usize, bps: u64, inflight: bool) -> RelativeSlowSample {
RelativeSlowSample {
peer_id: id,
Expand Down
6 changes: 6 additions & 0 deletions crates/rbitcoin-net/src/ibd/events/confirm_reject_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -451,6 +451,9 @@ fn apply_peer_event_body_and_control_surface() {
net: crate::NetAddr::from_socket(a),
cmd_tx,
in_flight: HashSet::new(),
requested: super::super::peer_io::solicit_track().0,
solicited_bytes: super::super::peer_io::solicit_track().1,
solicited_ms: super::super::peer_io::solicit_track().2,
peer_height: 10,
connected_ms: 1,
first_data_ms: 0,
Expand Down Expand Up @@ -745,6 +748,9 @@ fn apply_peer_event_block_framed_bq_horizon_and_headers_done() {
)),
cmd_tx,
in_flight: HashSet::new(),
requested: super::super::peer_io::solicit_track().0,
solicited_bytes: super::super::peer_io::solicit_track().1,
solicited_ms: super::super::peer_io::solicit_track().2,
peer_height: 5,
connected_ms: 1,
first_data_ms: 0,
Expand Down
Loading
Loading