diff --git a/changelog.d/ibd-track-retain.md b/changelog.d/ibd-track-retain.md new file mode 100644 index 000000000..1b04a58ff --- /dev/null +++ b/changelog.d/ibd-track-retain.md @@ -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. diff --git a/changelog.d/peer-lifecycle.md b/changelog.d/peer-lifecycle.md new file mode 100644 index 000000000..53b95049b --- /dev/null +++ b/changelog.d/peer-lifecycle.md @@ -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. diff --git a/changelog.d/sh-empty-ingest-seal.md b/changelog.d/sh-empty-ingest-seal.md new file mode 100644 index 000000000..62a682939 --- /dev/null +++ b/changelog.d/sh-empty-ingest-seal.md @@ -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. diff --git a/crates/rbitcoin-net/src/eviction.rs b/crates/rbitcoin-net/src/eviction.rs index 412889357..1eba468fc 100644 --- a/crates/rbitcoin-net/src/eviction.rs +++ b/crates/rbitcoin-net/src/eviction.rs @@ -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) -> Option { @@ -60,13 +62,37 @@ pub fn select_inbound_eviction(mut cands: Vec) -> 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::::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 { @@ -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. @@ -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] diff --git a/crates/rbitcoin-net/src/ibd/assign.rs b/crates/rbitcoin-net/src/ibd/assign.rs index 2943b6209..73e20336d 100644 --- a/crates/rbitcoin-net/src/ibd/assign.rs +++ b/crates/rbitcoin-net/src/ibd/assign.rs @@ -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); } } @@ -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)); } } @@ -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()); @@ -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; @@ -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, @@ -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); } diff --git a/crates/rbitcoin-net/src/ibd/dial.rs b/crates/rbitcoin-net/src/ibd/dial.rs index d568013cb..dc7e1ffa9 100644 --- a/crates/rbitcoin-net/src/ibd/dial.rs +++ b/crates/rbitcoin-net/src/ibd/dial.rs @@ -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)) @@ -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, + 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. /// @@ -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); @@ -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; @@ -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, @@ -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, diff --git a/crates/rbitcoin-net/src/ibd/events/confirm_reject_tests.rs b/crates/rbitcoin-net/src/ibd/events/confirm_reject_tests.rs index 005ce9ad6..1db64b17e 100644 --- a/crates/rbitcoin-net/src/ibd/events/confirm_reject_tests.rs +++ b/crates/rbitcoin-net/src/ibd/events/confirm_reject_tests.rs @@ -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, @@ -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, diff --git a/crates/rbitcoin-net/src/ibd/events/mod.rs b/crates/rbitcoin-net/src/ibd/events/mod.rs index 37c736ea2..21011e6ed 100644 --- a/crates/rbitcoin-net/src/ibd/events/mod.rs +++ b/crates/rbitcoin-net/src/ibd/events/mod.rs @@ -4,7 +4,9 @@ use super::assign::clear_hash_inflight; use super::assign_plan::{ remove_from_ordered, should_enqueue_header, want_headers_beyond_soft_cap, }; -use super::dial::{disconnect_peer, note_dead_without_block_bytes, release_peer_block_work}; +use super::dial::{ + disconnect_peer, note_dead_without_block_bytes, note_misbehavior_dead, release_peer_block_work, +}; use super::exit::{ header_lag_behind_peers, should_advance_locator_after_known_batch, should_log_empty_headers_lag, should_rerequest_headers_on_empty_lag, @@ -39,7 +41,7 @@ pub(crate) fn disconnect_all_peers(st: &mut IbdWorkState) { } st.inflight.clear(); for s in &mut st.slots { - s.in_flight.clear(); + s.track_clear(); s.alive = false; } st.slots.clear(); @@ -553,15 +555,14 @@ fn apply_block_framed( payload: Vec, ) { let wire_bytes = payload.len(); - note_block_rx(&mut st.slots, peer, wire_bytes); - // Unsolicited wire is not a body we asked for. Drop it before any copy. + // Unsolicited wire is not a body we asked for. Drop it before any copy + // and before it can move this peer's progress clock. let requested = st.inflight.contains_key(&hash); - if requested { - super::assign::note_block_len(st, wire_bytes); - } if !requested { return; } + note_block_rx(&mut st.slots, peer, wire_bytes); + super::assign::note_block_len(st, wire_bytes); clear_hash_inflight(&mut st.slots, &mut st.inflight, hash); if st.body.is_rejected(&hash) || hub.has_block(&hash) { return; @@ -649,7 +650,7 @@ fn apply_notfound(st: &mut IbdWorkState, peer: usize, hashes: Vec) { let mut freed = Vec::new(); if let Some(s) = st.slots.iter_mut().find(|s| s.id == peer) { for h in &hashes { - s.in_flight.remove(h); + s.track_remove(h); let empty = st .inflight .get_mut(h) @@ -671,13 +672,17 @@ fn apply_peer_dead(st: &mut IbdWorkState, peer_book: &mut AddrMan, peer: usize, warn!("ibd: peer[{peer}] dead: {reason}"); super::header_walk::forget_walk_peer(st, peer); if let Some(s) = st.slots.iter().find(|s| s.id == peer) { - note_dead_without_block_bytes( - peer_book, - &mut st.addr_cooldown, - s.addr, - s.first_data_ms, - Instant::now(), - ); + if reason == "peer misbehavior threshold" { + note_misbehavior_dead(peer_book, &mut st.addr_cooldown, s.addr, Instant::now()); + } else { + note_dead_without_block_bytes( + peer_book, + &mut st.addr_cooldown, + s.addr, + s.first_data_ms, + Instant::now(), + ); + } let lat = s.first_data_ms.saturating_sub(s.connected_ms); peer_book.apply_ibd_dead_speed( s.addr, @@ -1070,6 +1075,48 @@ pub(crate) fn parent_height( None } +#[cfg(test)] +mod misbehavior_dead_tests { + use super::super::assign::tests::dummy_slot; + use super::super::state::IbdWorkState; + use super::apply_peer_dead; + use crate::seeds::AddrMan; + + #[test] + fn threshold_dead_cools_after_a_body_and_other_deaths_do_not() { + let mut slot = dummy_slot(1); + slot.first_data_ms = 42; + let addr = slot.addr; + let mut st = IbdWorkState::new(vec![slot], None, Some(0)); + let mut book = AddrMan::new(); + book.note_connected(addr); + apply_peer_dead( + &mut st, + &mut book, + 1, + "peer misbehavior threshold".to_string(), + ); + assert!( + st.addr_cooldown.contains_key(&addr), + "misbehavior cools the dial after a block body" + ); + assert!(book.flags(&addr).failed_last_connect()); + + let mut other = dummy_slot(2); + other.first_data_ms = 42; + let other_addr = other.addr; + let mut st = IbdWorkState::new(vec![other], None, Some(0)); + let mut book = AddrMan::new(); + book.note_connected(other_addr); + apply_peer_dead(&mut st, &mut book, 2, "bye".to_string()); + assert!( + !st.addr_cooldown.contains_key(&other_addr), + "a body already counted keeps an ordinary death off the cooldown" + ); + assert!(!book.flags(&other_addr).failed_last_connect()); + } +} + #[cfg(test)] mod confirm_reject_tests; #[cfg(test)] diff --git a/crates/rbitcoin-net/src/ibd/header_walk.rs b/crates/rbitcoin-net/src/ibd/header_walk.rs index d709c94dd..384a829ae 100644 --- a/crates/rbitcoin-net/src/ibd/header_walk.rs +++ b/crates/rbitcoin-net/src/ibd/header_walk.rs @@ -2859,7 +2859,7 @@ fn renote_stored_path(st: &IbdWorkState, hub: &ChainHub) { mod tests { use super::*; use crate::ibd::events::apply_peer_event; - use crate::ibd::peer_io::{PeerCmd, PeerEvent, PeerSlot}; + use crate::ibd::peer_io::{solicit_track, PeerCmd, PeerEvent, PeerSlot}; use crate::ibd::state::IbdWorkState; use crate::seeds::AddrMan; use bitcoin::block::Header; @@ -2898,6 +2898,9 @@ mod tests { net: crate::NetAddr::from_socket(addr), cmd_tx, in_flight: Default::default(), + requested: solicit_track().0, + solicited_bytes: solicit_track().1, + solicited_ms: solicit_track().2, peer_height: 50_000, connected_ms: 1, first_data_ms: 0, diff --git a/crates/rbitcoin-net/src/ibd/peer_io.rs b/crates/rbitcoin-net/src/ibd/peer_io.rs index c2fd8dd34..623466b58 100644 --- a/crates/rbitcoin-net/src/ibd/peer_io.rs +++ b/crates/rbitcoin-net/src/ibd/peer_io.rs @@ -22,9 +22,9 @@ use bitcoin::BlockHash; use std::collections::HashSet; use std::net::SocketAddr; use std::sync::atomic::{AtomicU64, Ordering}; -use std::sync::Arc; +use std::sync::{Arc, Mutex}; use std::time::Instant; -use tokio::sync::mpsc; +use tokio::sync::{mpsc, OwnedSemaphorePermit, Semaphore}; use tokio::task::JoinHandle; /// Block and witness-block entries from an `inv` or `notfound`. @@ -111,77 +111,110 @@ pub(crate) struct PeerSlot { pub cmd_tx: mpsc::UnboundedSender, /// Hashes currently requested from this peer. pub in_flight: HashSet, + /// Same hashes, shared with the reader. One extra set per live peer so a + /// block check does not take the IBD work state. + pub requested: Arc>>, + /// Bytes of blocks this peer was asked for. Stall-clock input. + pub solicited_bytes: Arc, + /// Mono ms of the last solicited block the reader accepted. + pub solicited_ms: Arc, /// Peer's `version.start_height` (best-effort network tip signal). pub peer_height: u32, /// Mono ms when the slot became live (post-handshake). pub connected_ms: u64, /// First block-payload mono ms (0 = none yet). IBD main thread only. pub first_data_ms: u64, - /// All streamed wire bytes (EWMA input). Reader-only `fetch_add`. + /// All streamed wire bytes. Reader-only `fetch_add`. Not the stall clock. pub bytes_rx_total: Arc, pub rate: PeerRate, pub alive: bool, pub task: JoinHandle<()>, } -impl Drop for PeerSlot { - fn drop(&mut self) { - let _ = self.cmd_tx.send(PeerCmd::Shutdown); - self.task.abort(); - } +/// Empty side-set and solicited counters for a new [`PeerSlot`]. +pub(crate) fn solicit_track() -> ( + Arc>>, + Arc, + Arc, +) { + ( + Arc::new(Mutex::new(HashSet::new())), + Arc::new(AtomicU64::new(0)), + Arc::new(AtomicU64::new(0)), + ) } -/// Monotonic milliseconds for IBD stall clocks (process-relative). -pub(crate) fn ibd_mono_ms() -> u64 { - static T0: std::sync::OnceLock = std::sync::OnceLock::new(); - T0.get_or_init(Instant::now).elapsed().as_millis() as u64 -} +/// In-flight light decodes (inv, addr, tx, cmpct) per IBD peer. Not a knob. +const LIGHT_DECODE_PERMITS: usize = 32; -pub(crate) fn note_stream_bytes(counter: &AtomicU64, n: u64) { - if n == 0 { - return; +impl PeerSlot { + fn requested_set(&self) -> std::sync::MutexGuard<'_, HashSet> { + self.requested.lock().unwrap_or_else(|e| e.into_inner()) } - counter.fetch_add(n, Ordering::Relaxed); -} -pub(crate) fn sample_peer_rates(slots: &mut [PeerSlot], now_ms: u64) { - for s in slots { - if !s.alive { - continue; - } - let bytes = s.bytes_rx_total.load(Ordering::Relaxed); - s.rate.sample(now_ms, bytes, !s.in_flight.is_empty()); + pub(crate) fn track_insert(&mut self, hash: BlockHash) { + self.in_flight.insert(hash); + self.requested_set().insert(hash); } -} -pub(crate) fn note_block_progress(slots: &mut [PeerSlot], peer: usize) { - if let Some(s) = slots.iter_mut().find(|s| s.id == peer) { - s.rate.note_rx(ibd_mono_ms()); + pub(crate) fn track_remove(&mut self, hash: &BlockHash) -> bool { + let removed = self.in_flight.remove(hash); + self.requested_set().remove(hash); + removed + } + + pub(crate) fn track_clear(&mut self) { + self.in_flight.clear(); + self.requested_set().clear(); + } + + pub(crate) fn track_retain(&mut self, mut keep: impl FnMut(&BlockHash) -> bool) { + // Judge each hash once, before taking the reader mutex. The shared + // set then matches that decision. `has_block` does not run while the + // IBD reader is blocked in `note_solicited_block`. + self.in_flight.retain(|h| keep(h)); + let mut set = self.requested_set(); + set.clear(); + set.extend(self.in_flight.iter().copied()); + } + + pub(crate) fn track_drain(&mut self) -> Vec { + let drained: Vec = self.in_flight.drain().collect(); + self.requested_set().clear(); + drained } } -pub(crate) fn note_block_rx(slots: &mut [PeerSlot], peer: usize, wire_bytes: usize) { - if let Some(s) = slots.iter_mut().find(|s| s.id == peer) { - let now = ibd_mono_ms(); - s.rate.note_rx(now); - if wire_bytes > 0 && s.first_data_ms == 0 { - s.first_data_ms = now; - } +/// `true` when `hash` was requested: record `len` and the read time. +/// An absent hash leaves both counters unchanged. +pub(crate) fn note_solicited_block( + requested: &Mutex>, + bytes: &AtomicU64, + ms: &AtomicU64, + hash: &BlockHash, + len: usize, +) -> bool { + let mut set = requested.lock().unwrap_or_else(|e| e.into_inner()); + // One request, one unmetered body. A resend is not new progress. + if !set.remove(hash) { + return false; } + bytes.fetch_add(len as u64, Ordering::Relaxed); + ms.store(ibd_mono_ms(), Ordering::Relaxed); + true } -fn note_read_progress(bytes_io: &AtomicU64, prog_mark: &mut usize, buffered: usize) { - let delta = buffered.saturating_sub(*prog_mark); - note_stream_bytes(bytes_io, delta as u64); - *prog_mark = buffered; +/// One light-decode slot, or `None` without waiting. +pub(crate) fn try_light_decode_permit(sem: &Arc) -> Option { + Arc::clone(sem).try_acquire_owned().ok() } -fn decoy_hook(rate: &mut PeerRateLimiter, ban_score: &mut u32, n: usize) -> Result<(), NetError> { - if crate::peer_dos::decoy_stays(rate, ban_score, n, BAN_SCORE_THRESHOLD) { - Ok(()) - } else { - Err(NetError::Protocol("peer misbehavior threshold")) +fn misbehavior_disconnect(rate: &mut PeerRateLimiter, ban_score: &mut u32, n: usize) -> bool { + if rate.note(n) { + return false; } + *ban_score = ban_score.saturating_add(RATE_LIMIT_BAN_SCORE); + *ban_score >= BAN_SCORE_THRESHOLD } fn relay_headers(id: usize, sinks: &PeerEventSinks, headers: Vec
) { @@ -227,23 +260,95 @@ fn apply_decoded_message( NetworkMessage::Addr(list) => relay_addr(id, sinks, &list), NetworkMessage::AddrV2(list) => relay_addrv2(id, sinks, &list), NetworkMessage::SendAddrV2 => {} - // Blocks must not reach decode (handled before the off-thread decode). NetworkMessage::Block(_) => {} NetworkMessage::Inv(inv) => relay_blocks_inv(id, sinks, &inv), _other => {} } } -fn on_heavy_err(id: usize, sinks: &PeerEventSinks, err: NetError) { - sinks.send_body(PeerEvent::Dead { - peer: id, - reason: err.to_string(), - }); +impl Drop for PeerSlot { + fn drop(&mut self) { + let _ = self.cmd_tx.send(PeerCmd::Shutdown); + self.task.abort(); + } } -fn is_eof_or_reset(err: &std::io::Error) -> bool { - err.kind() == std::io::ErrorKind::UnexpectedEof - || err.kind() == std::io::ErrorKind::ConnectionReset +/// Monotonic milliseconds for IBD stall clocks (process-relative). +pub(crate) fn ibd_mono_ms() -> u64 { + static T0: std::sync::OnceLock = std::sync::OnceLock::new(); + T0.get_or_init(Instant::now).elapsed().as_millis() as u64 +} + +pub(crate) fn note_stream_bytes(counter: &AtomicU64, n: u64) { + if n == 0 { + return; + } + counter.fetch_add(n, Ordering::Relaxed); +} + +pub(crate) fn sample_peer_rates(slots: &mut [PeerSlot], now_ms: u64) { + for s in slots { + if !s.alive { + continue; + } + // Read-time credit. Decoys stay in `bytes_rx_total` and do not move this. + let mark = s.solicited_ms.load(Ordering::Relaxed); + if mark > s.rate.progress_ms { + s.rate.progress_ms = mark; + } + let bytes = s.solicited_bytes.load(Ordering::Relaxed); + s.rate.sample(now_ms, bytes, !s.in_flight.is_empty()); + } +} + +pub(crate) fn note_block_progress(slots: &mut [PeerSlot], peer: usize) { + if let Some(s) = slots.iter_mut().find(|s| s.id == peer) { + s.rate.note_rx(ibd_mono_ms()); + } +} + +pub(crate) fn note_block_rx(slots: &mut [PeerSlot], peer: usize, wire_bytes: usize) { + if let Some(s) = slots.iter_mut().find(|s| s.id == peer) { + let now = ibd_mono_ms(); + s.rate.note_rx(now); + if wire_bytes > 0 && s.first_data_ms == 0 { + s.first_data_ms = now; + } + } +} + +/// Light-decode tasks die with the reader. A dropped `JoinHandle` would detach. +struct LightDecodes { + tasks: Vec>, +} + +impl Drop for LightDecodes { + fn drop(&mut self) { + for task in &self.tasks { + task.abort(); + } + } +} + +impl LightDecodes { + fn push(&mut self, task: JoinHandle<()>) { + self.tasks.retain(|t| !t.is_finished()); + self.tasks.push(task); + } +} + +fn note_read_progress(bytes_io: &AtomicU64, prog_mark: &mut usize, buffered: usize) { + let delta = buffered.saturating_sub(*prog_mark); + note_stream_bytes(bytes_io, delta as u64); + *prog_mark = buffered; +} + +fn decoy_hook(rate: &mut PeerRateLimiter, ban_score: &mut u32, n: usize) -> Result<(), NetError> { + if crate::peer_dos::decoy_stays(rate, ban_score, n, BAN_SCORE_THRESHOLD) { + Ok(()) + } else { + Err(NetError::Protocol("peer misbehavior threshold")) + } } /// Reader state for one IBD peer. The socket loop only pulls frames. @@ -251,6 +356,11 @@ struct IbdReadCtx { id: usize, out_tx: mpsc::UnboundedSender, sinks: PeerEventSinks, + requested: Arc>>, + solicited_bytes: Arc, + solicited_ms: Arc, + light_decode: Arc, + light: LightDecodes, /// Decoys and unknown types only. Requested block bodies stay off /// this window so a fast peer is not clipped at the tip-follow cap. rate: PeerRateLimiter, @@ -258,6 +368,51 @@ struct IbdReadCtx { ban_score: u32, } +fn send_threshold_dead(sinks: &PeerEventSinks, id: usize) { + sinks.send_body(PeerEvent::Dead { + peer: id, + reason: "peer misbehavior threshold".to_string(), + }); +} + +fn is_eof_or_reset(err: &std::io::Error) -> bool { + err.kind() == std::io::ErrorKind::UnexpectedEof + || err.kind() == std::io::ErrorKind::ConnectionReset +} + +async fn decode_light_frame( + id: usize, + sinks: PeerEventSinks, + frame: FramedMessage, + permit: OwnedSemaphorePermit, +) { + let _permit = permit; + match frame.try_decode() { + Ok(msg) => apply_decoded_message(id, &sinks, msg), + Err(e) => { + sinks.send_body(PeerEvent::Dead { + peer: id, + reason: e.to_string(), + }); + } + } +} + +fn on_heavy_decoded( + id: usize, + sinks: &PeerEventSinks, + msg: bitcoin::p2p::message::RawNetworkMessage, +) { + apply_decoded_message(id, sinks, msg); +} + +fn on_heavy_err(id: usize, sinks: &PeerEventSinks, err: NetError) { + sinks.send_body(PeerEvent::Dead { + peer: id, + reason: err.to_string(), + }); +} + impl IbdReadCtx { /// `true` stops the reader. fn handle(&mut self, frame: Result) -> bool { @@ -273,8 +428,12 @@ impl IbdReadCtx { return false; } if frame.is_block() { - self.on_block(frame); - return false; + return self.on_block(frame); + } + // Headers and notfound stay on the heavy decode pool. + // Light frames take a permit without waiting. + if !frame.decode_is_cpu_heavy() { + return self.on_light(frame); } self.on_heavy(frame); false @@ -286,14 +445,28 @@ impl IbdReadCtx { } } - fn on_block(&self, frame: FramedMessage) { + fn on_block(&mut self, frame: FramedMessage) -> bool { match frame.block_hash_from_header() { Some(hash) if frame.payload.len() >= 80 => { - self.sinks.send_body(PeerEvent::BlockFramed { - peer: self.id, - hash, - payload: frame.payload, - }); + // Unsolicited bodies stay off the channel. + // Requested bodies are not rate-capped. + let n = frame.payload.len(); + if note_solicited_block( + &self.requested, + &self.solicited_bytes, + &self.solicited_ms, + &hash, + n, + ) { + self.sinks.send_body(PeerEvent::BlockFramed { + peer: self.id, + hash, + payload: frame.payload, + }); + false + } else { + self.score_overflow(n) + } } // Re-request when a hash is known but the header bytes are short. // `block_hash_from_header` rejects that today; the arm stays so the @@ -303,13 +476,29 @@ impl IbdReadCtx { peer: self.id, hash, }); + false } None => { rbitcoin_log::debug!( "ibd: peer[{}] block frame without usable header hash", self.id ); + false + } + } + } + + fn on_light(&mut self, frame: FramedMessage) -> bool { + let n = frame.payload.len(); + match try_light_decode_permit(&self.light_decode) { + Some(permit) => { + let sinks = self.sinks.clone(); + let id = self.id; + let task = tokio::spawn(decode_light_frame(id, sinks, frame, permit)); + self.light.push(task); + false } + None => self.score_overflow(n), } } @@ -320,7 +509,7 @@ impl IbdReadCtx { let sinks_err = self.sinks.clone(); spawn_decode_then_with_err( frame, - move |msg| apply_decoded_message(id, &sinks_ok, msg), + move |msg| on_heavy_decoded(id, &sinks_ok, msg), move |err| on_heavy_err(id, &sinks_err, err), ); } @@ -350,15 +539,12 @@ impl IbdReadCtx { self.logged_invalid_v2 = true; rbitcoin_log::debug!("{}", crate::v2::v2_invalid_message_type_log()); } - if self.rate.note(contents_len) { - return false; - } - self.ban_score = self.ban_score.saturating_add(RATE_LIMIT_BAN_SCORE); - if self.ban_score >= BAN_SCORE_THRESHOLD { - self.sinks.send_body(PeerEvent::Dead { - peer: self.id, - reason: "peer misbehavior threshold".to_string(), - }); + self.score_overflow(contents_len) + } + + fn score_overflow(&mut self, n: usize) -> bool { + if misbehavior_disconnect(&mut self.rate, &mut self.ban_score, n) { + send_threshold_dead(&self.sinks, self.id); true } else { false @@ -366,6 +552,7 @@ impl IbdReadCtx { } } +#[allow(clippy::too_many_arguments)] // call-site args stay unbundled async fn read_ibd_peer( id: usize, magic: Magic, @@ -373,12 +560,21 @@ async fn read_ibd_peer( out_tx: mpsc::UnboundedSender, sinks_r: PeerEventSinks, bytes_io: Arc, + requested: Arc>>, + solicited_bytes: Arc, + solicited_ms: Arc, + light_decode: Arc, ) { let mut prog_mark = 0usize; let mut ctx = IbdReadCtx { id, out_tx, sinks: sinks_r, + requested, + solicited_bytes, + solicited_ms, + light_decode, + light: LightDecodes { tasks: Vec::new() }, rate: PeerRateLimiter::default_limits(), logged_invalid_v2: false, ban_score: 0, @@ -515,6 +711,11 @@ pub(crate) async fn spawn_peer( let (out_tx, out_rx) = mpsc::unbounded_channel::(); let bytes_rx_total = Arc::new(AtomicU64::new(0)); let bytes_io = Arc::clone(&bytes_rx_total); + let (requested, solicited_bytes, solicited_ms) = solicit_track(); + let req_r = Arc::clone(&requested); + let sol_bytes_r = Arc::clone(&solicited_bytes); + let sol_ms_r = Arc::clone(&solicited_ms); + let light_decode = Arc::new(Semaphore::new(LIGHT_DECODE_PERMITS)); // Parent owns concurrent read + write tasks. Aborting the parent (PeerSlot // Drop / stall disconnect) must abort both children — plain JoinHandle drop @@ -537,7 +738,18 @@ pub(crate) async fn spawn_peer( // before getheaders; those writes raced Core's pipeline and peers closed // (ordered=0 / inflight=0 / never archive). let sinks_r = sinks.clone(); - let reader_task = tokio::spawn(read_ibd_peer(id, magic, reader, out_tx, sinks_r, bytes_io)); + let reader_task = tokio::spawn(read_ibd_peer( + id, + magic, + reader, + out_tx, + sinks_r, + bytes_io, + req_r, + sol_bytes_r, + sol_ms_r, + light_decode, + )); // Let the reader poll once before we accept write work (getheaders). tokio::task::yield_now().await; @@ -561,6 +773,9 @@ pub(crate) async fn spawn_peer( net: addr, cmd_tx, in_flight: HashSet::new(), + requested, + solicited_bytes, + solicited_ms, peer_height, connected_ms: ibd_mono_ms(), first_data_ms: 0, @@ -656,6 +871,9 @@ mod tests { )), cmd_tx, in_flight: HashSet::new(), + requested: Arc::new(Mutex::new(HashSet::new())), + solicited_bytes: Arc::new(AtomicU64::new(0)), + solicited_ms: Arc::new(AtomicU64::new(0)), peer_height: 100, connected_ms: 1, first_data_ms: 0, @@ -666,6 +884,166 @@ mod tests { } } + #[test] + fn usable_dial_and_services_filters() { + assert!(!services_useful_for_ibd(ServiceFlags::NETWORK)); + assert!(!services_useful_for_ibd(ServiceFlags::NETWORK_LIMITED)); + assert!(!services_useful_for_ibd(ServiceFlags::P2P_V2)); + assert!(!services_useful_for_ibd(ServiceFlags::NONE)); + assert!(services_useful_for_ibd( + ServiceFlags::NETWORK | ServiceFlags::P2P_V2 + )); + assert!(services_useful_for_ibd( + ServiceFlags::NETWORK_LIMITED | ServiceFlags::P2P_V2 + )); + + let good = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(1, 2, 3, 4)), 8333); + assert!(usable_dial_addr(&good)); + let zero_port = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(1, 2, 3, 4)), 0); + assert!(!usable_dial_addr(&zero_port)); + let unspec = SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 8333); + assert!(!usable_dial_addr(&unspec)); + let bcast = SocketAddr::new(IpAddr::V4(Ipv4Addr::BROADCAST), 8333); + assert!(!usable_dial_addr(&bcast)); + let multi = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(224, 0, 0, 1)), 8333); + assert!(!usable_dial_addr(&multi)); + let v6_good = SocketAddr::new(IpAddr::V6(Ipv6Addr::LOCALHOST), 8333); + assert!(usable_dial_addr(&v6_good)); + let v6_unspec = SocketAddr::new(IpAddr::V6(Ipv6Addr::UNSPECIFIED), 8333); + assert!(!usable_dial_addr(&v6_unspec)); + } + + #[test] + fn stream_bytes_sample_and_first_data() { + let mut s = dummy_slot(7); + note_stream_bytes(&s.bytes_rx_total, 0); + assert_eq!(s.bytes_rx_total.load(Ordering::Relaxed), 0); + note_stream_bytes(&s.bytes_rx_total, 100); + note_stream_bytes(&s.bytes_rx_total, 50); + assert_eq!(s.bytes_rx_total.load(Ordering::Relaxed), 150); + + s.in_flight.insert(BlockHash::from_byte_array([1u8; 32])); + sample_peer_rates(std::slice::from_mut(&mut s), 0); + sample_peer_rates(std::slice::from_mut(&mut s), 5_000); + assert!(s.rate.bps().is_some()); + + while ibd_mono_ms() == 0 { + std::thread::sleep(std::time::Duration::from_millis(1)); + } + note_block_rx(std::slice::from_mut(&mut s), 7, 0); + assert_eq!(s.first_data_ms, 0); + note_block_progress(std::slice::from_mut(&mut s), 7); + note_block_rx(std::slice::from_mut(&mut s), 7, 1000); + assert!(s.first_data_ms > 0); + note_block_progress(std::slice::from_mut(&mut s), 99); + note_block_rx(std::slice::from_mut(&mut s), 99, 1); + assert!(ibd_mono_ms() > 0); + } + + #[test] + fn track_retain_judges_each_hash_once() { + let mut slot = dummy_slot(3); + let keep_hash = BlockHash::from_byte_array([1u8; 32]); + let drop_hash = BlockHash::from_byte_array([2u8; 32]); + slot.track_insert(keep_hash); + slot.track_insert(drop_hash); + let mut calls = 0usize; + slot.track_retain(|hash| { + calls += 1; + *hash != drop_hash + }); + assert_eq!(calls, 2, "one judgment per in-flight hash"); + assert!(slot.in_flight.contains(&keep_hash)); + assert!(!slot.in_flight.contains(&drop_hash)); + let requested = slot.requested.lock().unwrap().clone(); + assert_eq!(requested, slot.in_flight); + } + + #[test] + fn unsolicited_block_does_not_refresh_progress() { + let hash = BlockHash::from_byte_array([7u8; 32]); + let requested = Mutex::new(HashSet::new()); + let bytes = AtomicU64::new(0); + let ms = AtomicU64::new(0); + assert!( + !note_solicited_block(&requested, &bytes, &ms, &hash, 80), + "an unsolicited block does not count" + ); + assert_eq!(bytes.load(Ordering::Relaxed), 0); + assert_eq!(ms.load(Ordering::Relaxed), 0); + + while ibd_mono_ms() == 0 { + std::thread::sleep(std::time::Duration::from_millis(1)); + } + requested.lock().unwrap().insert(hash); + let len = crate::ibd::rate::PROGRESS_STEP as usize; + assert!( + note_solicited_block(&requested, &bytes, &ms, &hash, len), + "a requested block records its bytes" + ); + assert_eq!(bytes.load(Ordering::Relaxed), len as u64); + let marked = ms.load(Ordering::Relaxed); + assert!(marked > 0, "a requested block stamps the read time"); + + let mut quiet = PeerRate::default(); + quiet.sample(0, 0, true); + quiet.sample(10_000, 0, true); + assert_eq!( + quiet.progress_ms, 0, + "zero solicited bytes over 10s leave the stall clock" + ); + let mut moved = PeerRate::default(); + moved.sample(0, 0, true); + moved.sample(1_000, bytes.load(Ordering::Relaxed), true); + assert_eq!(moved.progress_ms, 1_000); + + let mut slot = dummy_slot(4); + slot.in_flight.insert(hash); + assert_eq!(slot.solicited_bytes.load(Ordering::Relaxed), 0); + assert_eq!(slot.solicited_ms.load(Ordering::Relaxed), 0); + sample_peer_rates(std::slice::from_mut(&mut slot), 0); + note_stream_bytes(&slot.bytes_rx_total, 8_000_000); + sample_peer_rates(std::slice::from_mut(&mut slot), 10_000); + assert_eq!( + slot.rate.progress_ms, 0, + "unsolicited stream bytes do not refresh progress" + ); + + slot.solicited_ms.store(marked, Ordering::Relaxed); + sample_peer_rates(std::slice::from_mut(&mut slot), 0); + assert_eq!(slot.rate.progress_ms, marked); + } + + #[test] + fn light_decode_permit_does_not_wait() { + let sem = Arc::new(Semaphore::new(1)); + let first = try_light_decode_permit(&sem); + assert!(first.is_some(), "one permit is available immediately"); + assert!( + try_light_decode_permit(&sem).is_none(), + "a second permit does not wait" + ); + drop(first); + assert!(try_light_decode_permit(&sem).is_some()); + } + + #[test] + fn solicited_replay_does_not_refresh_progress() { + let (req, bytes, ms) = solicit_track(); + let hash = BlockHash::from_byte_array([9u8; 32]); + req.lock().unwrap().insert(hash); + assert!(note_solicited_block(&req, &bytes, &ms, &hash, 100)); + let stamped = ms.load(Ordering::Relaxed); + assert_eq!(bytes.load(Ordering::Relaxed), 100); + assert!(!req.lock().unwrap().contains(&hash)); + assert!( + !note_solicited_block(&req, &bytes, &ms, &hash, 50), + "a resend of a solicited block is not another request" + ); + assert_eq!(bytes.load(Ordering::Relaxed), 100); + assert_eq!(ms.load(Ordering::Relaxed), stamped); + } + fn framed(command: [u8; 12], payload: Vec) -> FramedMessage { FramedMessage { magic: Magic::from(bitcoin::Network::Regtest), @@ -674,7 +1052,9 @@ mod tests { } } - fn test_ctx() -> ( + fn test_ctx( + permits: usize, + ) -> ( IbdReadCtx, mpsc::UnboundedReceiver, mpsc::UnboundedReceiver, @@ -690,6 +1070,11 @@ mod tests { body: body_tx, ctrl: ctrl_tx, }, + requested: Arc::new(Mutex::new(HashSet::new())), + solicited_bytes: Arc::new(AtomicU64::new(0)), + solicited_ms: Arc::new(AtomicU64::new(0)), + light_decode: Arc::new(Semaphore::new(permits)), + light: LightDecodes { tasks: Vec::new() }, rate: PeerRateLimiter::default_limits(), logged_invalid_v2: false, ban_score: 0, @@ -697,6 +1082,11 @@ mod tests { (ctx, out_rx, body_rx, ctrl_rx) } + fn oversized_block() -> FramedMessage { + let n = (crate::peer_dos::DEFAULT_MAX_BYTES_PER_SEC as usize) + 1; + framed(*b"block\0\0\0\0\0\0\0", vec![0u8; n]) + } + fn one_header() -> Header { use bitcoin::block::Version; use bitcoin::CompactTarget; @@ -710,6 +1100,10 @@ mod tests { } } + fn recv_body(rx: &mut mpsc::UnboundedReceiver) -> PeerEvent { + rx.try_recv().expect("body event") + } + #[test] fn read_progress_counts_the_delta_from_the_mark() { let bytes = AtomicU64::new(0); @@ -718,6 +1112,7 @@ mod tests { note_read_progress(&bytes, &mut mark, 25); assert_eq!(bytes.load(Ordering::Relaxed), 25); assert_eq!(mark, 25); + // The reader stores the latest buffered size, even when it shrinks. note_read_progress(&bytes, &mut mark, 4); assert_eq!(bytes.load(Ordering::Relaxed), 25); assert_eq!(mark, 4); @@ -738,12 +1133,12 @@ mod tests { } #[test] - fn ping_and_block_frames_stay_on_their_channels() { - let (mut ctx, mut out_rx, mut body_rx, mut ctrl_rx) = test_ctx(); + fn ping_and_solicited_block_stay_on_their_channels() { + let (mut ctx, mut out_rx, mut body_rx, mut ctrl_rx) = test_ctx(1); let nonce = 0x0102_0304_0506_0708u64; assert!(!ctx.handle(Ok(framed( *b"ping\0\0\0\0\0\0\0\0", - nonce.to_le_bytes().to_vec(), + nonce.to_le_bytes().to_vec() )))); assert!(matches!(out_rx.try_recv(), Ok(NetworkMessage::Pong(n)) if n == nonce)); assert!(!ctx.handle(Ok(framed(*b"ping\0\0\0\0\0\0\0\0", vec![1, 2, 3])))); @@ -751,74 +1146,112 @@ mod tests { let block = framed(*b"block\0\0\0\0\0\0\0", vec![7u8; 80]); let hash = block.block_hash_from_header().expect("80-byte header"); + ctx.requested.lock().unwrap().insert(hash); assert!(!ctx.handle(Ok(block))); - match body_rx.try_recv() { - Ok(PeerEvent::BlockFramed { + match recv_body(&mut body_rx) { + PeerEvent::BlockFramed { peer, hash: got, payload, - }) => { + } => { assert_eq!(peer, 3); assert_eq!(got, hash); assert_eq!(payload.len(), 80); } - _ => panic!("a block frame must be forwarded"), + _ => panic!("solicited block must be framed"), } + assert_eq!(ctx.solicited_bytes.load(Ordering::Relaxed), 80); assert!(ctrl_rx.try_recv().is_err()); + + let again = framed(*b"block\0\0\0\0\0\0\0", vec![7u8; 80]); + assert!(!ctx.handle(Ok(again)), "one resend still fits the window"); + assert!(body_rx.try_recv().is_err(), "a resend is not a second body"); + assert_eq!(ctx.solicited_bytes.load(Ordering::Relaxed), 80); + assert!(!ctx.handle(Ok(framed(*b"block\0\0\0\0\0\0\0", vec![1, 2, 3])))); assert!(body_rx.try_recv().is_err(), "a short block has no hash"); } + #[test] + fn second_unsolicited_block_past_the_window_stops_the_reader() { + let (mut ctx, _out, mut body_rx, _ctrl) = test_ctx(1); + assert!(!ctx.handle(Ok(oversized_block()))); + assert!(body_rx.try_recv().is_err(), "the first overflow stays"); + assert!(ctx.handle(Ok(oversized_block()))); + match recv_body(&mut body_rx) { + PeerEvent::Dead { peer, reason } => { + assert_eq!(peer, 3); + assert_eq!(reason, "peer misbehavior threshold"); + } + _ => panic!("second overflow must stop"), + } + } + + #[test] + fn light_permit_exhaustion_scores_like_other_frames() { + let (mut ctx, _out, mut body_rx, _ctrl) = test_ctx(0); + let n = (crate::peer_dos::DEFAULT_MAX_BYTES_PER_SEC as usize) + 1; + let frame = framed(*b"inv\0\0\0\0\0\0\0\0\0", vec![0u8; n]); + assert!(!ctx.handle(Ok(frame))); + assert!(body_rx.try_recv().is_err()); + let frame = framed(*b"inv\0\0\0\0\0\0\0\0\0", vec![0u8; n]); + assert!(ctx.handle(Ok(frame))); + match recv_body(&mut body_rx) { + PeerEvent::Dead { reason, .. } => assert_eq!(reason, "peer misbehavior threshold"), + _ => panic!("a second permit miss past the window must stop"), + } + } + #[test] fn invalid_v2_and_io_errors_stop_on_the_shipped_reasons() { - let (mut ctx, _out, mut body_rx, _ctrl) = test_ctx(); + let (mut ctx, _out, mut body_rx, _ctrl) = test_ctx(1); let n = (crate::peer_dos::DEFAULT_MAX_BYTES_PER_SEC as usize) + 1; assert!(!ctx.handle(Err(NetError::InvalidV2Type { contents_len: n }))); assert!(ctx.handle(Err(NetError::InvalidV2Type { contents_len: n }))); - match body_rx.try_recv() { - Ok(PeerEvent::Dead { reason, .. }) => assert_eq!(reason, "peer misbehavior threshold"), + match recv_body(&mut body_rx) { + PeerEvent::Dead { reason, .. } => assert_eq!(reason, "peer misbehavior threshold"), _ => panic!("second invalid type must stop"), } - let (mut ctx, _out, mut body_rx, _ctrl) = test_ctx(); + let (mut ctx, _out, mut body_rx, _ctrl) = test_ctx(1); let eof = std::io::Error::new(std::io::ErrorKind::UnexpectedEof, "closed"); assert!(ctx.handle(Err(NetError::Io(eof)))); - match body_rx.try_recv() { - Ok(PeerEvent::Dead { reason, .. }) => assert!(reason.starts_with("eof:"), "{reason}"), + match recv_body(&mut body_rx) { + PeerEvent::Dead { reason, .. } => assert!(reason.starts_with("eof:"), "{reason}"), _ => panic!("eof must be a dead peer"), } - let (mut ctx, _out, mut body_rx, _ctrl) = test_ctx(); + let (mut ctx, _out, mut body_rx, _ctrl) = test_ctx(1); let reset = std::io::Error::new(std::io::ErrorKind::ConnectionReset, "reset"); assert!(ctx.handle(Err(NetError::Io(reset)))); - match body_rx.try_recv() { - Ok(PeerEvent::Dead { reason, .. }) => assert!(reason.starts_with("eof:"), "{reason}"), + match recv_body(&mut body_rx) { + PeerEvent::Dead { reason, .. } => assert!(reason.starts_with("eof:"), "{reason}"), _ => panic!("reset must be a dead peer"), } - let (mut ctx, _out, mut body_rx, _ctrl) = test_ctx(); + let (mut ctx, _out, mut body_rx, _ctrl) = test_ctx(1); assert!(ctx.handle(Err(NetError::Protocol("bye")))); - match body_rx.try_recv() { - Ok(PeerEvent::Dead { reason, .. }) => assert!(reason.contains("bye"), "{reason}"), + match recv_body(&mut body_rx) { + PeerEvent::Dead { reason, .. } => assert!(reason.contains("bye"), "{reason}"), _ => panic!("protocol error must be a dead peer"), } } #[test] - fn heavy_frames_deliver_or_die_without_stopping_the_read() { - use bitcoin::consensus::encode::Encodable; + fn light_and_heavy_frames_deliver_or_die_without_stopping_the_read() { let rt = tokio::runtime::Builder::new_current_thread() .enable_all() .build() .unwrap(); rt.block_on(async { - let (mut ctx, _out, mut body_rx, mut ctrl_rx) = test_ctx(); + use bitcoin::consensus::encode::Encodable; + let (mut ctx, _out, mut body_rx, mut ctrl_rx) = test_ctx(2); let hash = BlockHash::from_byte_array([4u8; 32]); let inv = bitcoin::consensus::encode::serialize(&vec![Inventory::Block(hash)]); assert!(!ctx.handle(Ok(framed(*b"inv\0\0\0\0\0\0\0\0\0", inv)))); let ev = tokio::time::timeout(std::time::Duration::from_secs(2), ctrl_rx.recv()) .await - .expect("inv") + .expect("light inv") .expect("ctrl open"); match ev { PeerEvent::BlocksInv { peer, hashes } => { @@ -828,13 +1261,37 @@ mod tests { _ => panic!("block inv must reach ctrl"), } + let mut over = Vec::new(); + bitcoin::consensus::encode::VarInt((crate::codec::MAX_INV_SIZE as u64) + 1) + .consensus_encode(&mut over) + .unwrap(); + assert!(!ctx.handle(Ok(framed(*b"inv\0\0\0\0\0\0\0\0\0", over)))); + let ev = tokio::time::timeout(std::time::Duration::from_secs(2), body_rx.recv()) + .await + .expect("oversize inv") + .expect("body open"); + match ev { + PeerEvent::Dead { reason, .. } => { + assert!(reason.contains("message too large"), "{reason}") + } + _ => panic!("oversize inv must die in the decoder"), + } + + // Core's headers payload is count, then each 80-byte header and a 0 tx count. let mut headers = Vec::new(); bitcoin::consensus::encode::VarInt(1) .consensus_encode(&mut headers) .unwrap(); one_header().consensus_encode(&mut headers).unwrap(); headers.push(0); - assert!(!ctx.handle(Ok(framed(*b"headers\0\0\0\0\0", headers)))); + let headers_frame = framed(*b"headers\0\0\0\0\0", headers); + match headers_frame.clone().try_decode() { + Ok(msg) => { + assert!(matches!(msg.payload(), NetworkMessage::Headers(h) if h.len() == 1)) + } + Err(e) => panic!("headers fixture must decode, got {e}"), + } + assert!(!ctx.handle(Ok(headers_frame))); let ev = tokio::time::timeout(std::time::Duration::from_secs(2), ctrl_rx.recv()) .await .expect("headers") @@ -862,7 +1319,7 @@ mod tests { }); } - fn event_sinks() -> ( + fn sinks() -> ( PeerEventSinks, mpsc::UnboundedReceiver, mpsc::UnboundedReceiver, @@ -879,115 +1336,92 @@ mod tests { ) } + fn apply(msg: NetworkMessage, sinks: &PeerEventSinks) { + apply_decoded_message( + 3, + sinks, + bitcoin::p2p::message::RawNetworkMessage::new( + Magic::from(bitcoin::Network::Regtest), + msg, + ), + ); + } + #[test] fn decoded_messages_fan_out_blocks_and_addrs_only() { use bitcoin::p2p::address::{AddrV2, AddrV2Message, Address}; - let (sinks, mut body_rx, mut ctrl_rx) = event_sinks(); - let magic = Magic::from(bitcoin::Network::Regtest); - let apply = |msg| { - apply_decoded_message( - 3, - &sinks, - bitcoin::p2p::message::RawNetworkMessage::new(magic, msg), - ); - }; - apply(NetworkMessage::Headers(vec![one_header()])); + + let (sinks, mut body_rx, mut ctrl_rx) = sinks(); + apply(NetworkMessage::Headers(vec![one_header()]), &sinks); assert!(matches!( ctrl_rx.try_recv(), Ok(PeerEvent::Headers { headers, .. }) if headers.len() == 1 )); - apply(NetworkMessage::NotFound(vec![])); + apply(NetworkMessage::Headers(vec![]), &sinks); + assert!(matches!(ctrl_rx.try_recv(), Ok(PeerEvent::Headers { .. }))); + + apply(NetworkMessage::NotFound(vec![]), &sinks); + apply( + NetworkMessage::NotFound(vec![Inventory::Transaction( + bitcoin::Txid::from_byte_array([3u8; 32]), + )]), + &sinks, + ); assert!(body_rx.try_recv().is_err()); let block = BlockHash::from_byte_array([1u8; 32]); - apply(NetworkMessage::NotFound(vec![Inventory::Block(block)])); + apply( + NetworkMessage::NotFound(vec![Inventory::Block(block)]), + &sinks, + ); assert!(matches!( body_rx.try_recv(), Ok(PeerEvent::NotFound { hashes, .. }) if hashes == vec![block] )); - apply(NetworkMessage::Inv(vec![Inventory::Block(block)])); + + apply(NetworkMessage::Inv(vec![]), &sinks); + assert!(ctrl_rx.try_recv().is_err()); + apply(NetworkMessage::Inv(vec![Inventory::Block(block)]), &sinks); assert!(matches!( ctrl_rx.try_recv(), Ok(PeerEvent::BlocksInv { hashes, .. }) if hashes == vec![block] )); + let good = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(1, 2, 3, 4)), 8333); let flags = ServiceFlags::NETWORK | ServiceFlags::P2P_V2; - apply(NetworkMessage::Addr(vec![(1, Address::new(&good, flags))])); + apply(NetworkMessage::Addr(vec![]), &sinks); + assert!(ctrl_rx.try_recv().is_err()); + apply( + NetworkMessage::Addr(vec![(1, Address::new(&good, flags))]), + &sinks, + ); assert!(matches!(ctrl_rx.try_recv(), Ok(PeerEvent::Addrs { .. }))); - apply(NetworkMessage::AddrV2(vec![AddrV2Message { - time: 1, - services: flags, - addr: AddrV2::Ipv4(Ipv4Addr::new(9, 9, 9, 9)), - port: 8333, - }])); + + apply(NetworkMessage::AddrV2(vec![]), &sinks); + assert!(ctrl_rx.try_recv().is_err()); + apply( + NetworkMessage::AddrV2(vec![AddrV2Message { + time: 1, + services: flags, + addr: AddrV2::Ipv4(Ipv4Addr::new(9, 9, 9, 9)), + port: 8333, + }]), + &sinks, + ); assert!(matches!(ctrl_rx.try_recv(), Ok(PeerEvent::Addrs { .. }))); - apply(NetworkMessage::SendAddrV2); - apply(NetworkMessage::Block(bitcoin::Block { - header: one_header(), - txdata: vec![], - })); - apply(NetworkMessage::Verack); - apply(NetworkMessage::Addr(vec![])); - apply(NetworkMessage::Inv(vec![])); + + apply(NetworkMessage::SendAddrV2, &sinks); + apply( + NetworkMessage::Block(bitcoin::Block { + header: one_header(), + txdata: vec![], + }), + &sinks, + ); + apply(NetworkMessage::Verack, &sinks); assert!(body_rx.try_recv().is_err()); assert!(ctrl_rx.try_recv().is_err()); } - #[test] - fn usable_dial_and_services_filters() { - assert!(!services_useful_for_ibd(ServiceFlags::NETWORK)); - assert!(!services_useful_for_ibd(ServiceFlags::NETWORK_LIMITED)); - assert!(!services_useful_for_ibd(ServiceFlags::P2P_V2)); - assert!(!services_useful_for_ibd(ServiceFlags::NONE)); - assert!(services_useful_for_ibd( - ServiceFlags::NETWORK | ServiceFlags::P2P_V2 - )); - assert!(services_useful_for_ibd( - ServiceFlags::NETWORK_LIMITED | ServiceFlags::P2P_V2 - )); - - let good = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(1, 2, 3, 4)), 8333); - assert!(usable_dial_addr(&good)); - let zero_port = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(1, 2, 3, 4)), 0); - assert!(!usable_dial_addr(&zero_port)); - let unspec = SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 8333); - assert!(!usable_dial_addr(&unspec)); - let bcast = SocketAddr::new(IpAddr::V4(Ipv4Addr::BROADCAST), 8333); - assert!(!usable_dial_addr(&bcast)); - let multi = SocketAddr::new(IpAddr::V4(Ipv4Addr::new(224, 0, 0, 1)), 8333); - assert!(!usable_dial_addr(&multi)); - let v6_good = SocketAddr::new(IpAddr::V6(Ipv6Addr::LOCALHOST), 8333); - assert!(usable_dial_addr(&v6_good)); - let v6_unspec = SocketAddr::new(IpAddr::V6(Ipv6Addr::UNSPECIFIED), 8333); - assert!(!usable_dial_addr(&v6_unspec)); - } - - #[test] - fn stream_bytes_sample_and_first_data() { - let mut s = dummy_slot(7); - note_stream_bytes(&s.bytes_rx_total, 0); - assert_eq!(s.bytes_rx_total.load(Ordering::Relaxed), 0); - note_stream_bytes(&s.bytes_rx_total, 100); - note_stream_bytes(&s.bytes_rx_total, 50); - assert_eq!(s.bytes_rx_total.load(Ordering::Relaxed), 150); - - s.in_flight.insert(BlockHash::from_byte_array([1u8; 32])); - sample_peer_rates(std::slice::from_mut(&mut s), 0); - sample_peer_rates(std::slice::from_mut(&mut s), 5_000); - assert!(s.rate.bps().is_some()); - - while ibd_mono_ms() == 0 { - std::thread::sleep(std::time::Duration::from_millis(1)); - } - note_block_rx(std::slice::from_mut(&mut s), 7, 0); - assert_eq!(s.first_data_ms, 0); - note_block_progress(std::slice::from_mut(&mut s), 7); - note_block_rx(std::slice::from_mut(&mut s), 7, 1000); - assert!(s.first_data_ms > 0); - note_block_progress(std::slice::from_mut(&mut s), 99); - note_block_rx(std::slice::from_mut(&mut s), 99, 1); - assert!(ibd_mono_ms() > 0); - } - #[test] fn block_inventory_hashes_keeps_blocks_only() { let block = BlockHash::from_byte_array([1u8; 32]); diff --git a/crates/rbitcoin-net/src/peer.rs b/crates/rbitcoin-net/src/peer.rs index 0858a1a5b..f250aa15a 100644 --- a/crates/rbitcoin-net/src/peer.rs +++ b/crates/rbitcoin-net/src/peer.rs @@ -140,6 +140,35 @@ async fn accept_received_from_peer( }) } +/// One-day in-memory refusal. Noban peers are not recorded. +/// A loopback address is every local connection, so Core disconnects that +/// peer and does not discourage the address. +fn note_threshold_refusal(session: Option<&crate::peers::LivePeer>) { + let Some(s) = session.filter(|s| !s.session_noban()) else { + return; + }; + if local_addr_skips_discourage(s.addr.ip()) { + return; + } + if let Some(hub) = s.peer_hub() { + hub.note_misbehavior_addr(s.addr.ip()); + } +} + +fn local_addr_skips_discourage(ip: std::net::IpAddr) -> bool { + match ip { + std::net::IpAddr::V4(v) => v.is_loopback() || v.octets()[0] == 0, + std::net::IpAddr::V6(v) => v.is_loopback(), + } +} + +/// Disconnect error for a score that is already at the misbehavior line. +/// Records the address, including when the caller is not [`punish_disconnect`]. +fn threshold_disconnect(session: Option<&crate::peers::LivePeer>) -> NetError { + note_threshold_refusal(session); + NetError::Protocol("peer misbehavior threshold") +} + fn punish_disconnect(ban_score: &mut u32, session: Option<&crate::peers::LivePeer>) { if let Some(s) = session.filter(|s| s.session_noban()) { rbitcoin_log::info!("Warning: not punishing noban peer {}!", s.id); @@ -148,6 +177,7 @@ fn punish_disconnect(ban_score: &mut u32, session: Option<&crate::peers::LivePee *ban_score = ban_score.saturating_add(BAN_SCORE_THRESHOLD); if let Some(s) = session { s.request_disconnect(); + note_threshold_refusal(Some(s)); } } @@ -1623,6 +1653,9 @@ pub async fn peer_session_with( // Any socket Io means the peer is gone — exit cleanly so // unregister runs inside the Core disconnect_nodes 5s wait. Err(NetError::Io(_)) => return Ok(()), + Err(NetError::Protocol("peer misbehavior threshold")) => { + return Err(threshold_disconnect(session.as_deref())); + } Err(NetError::MessageTooLarge(n)) => { follow.ban_score = follow.ban_score.saturating_add(OVERSIZE_BAN_SCORE); rbitcoin_log::warn!( @@ -1630,7 +1663,7 @@ pub async fn peer_session_with( follow.ban_score ); if follow.ban_score >= BAN_SCORE_THRESHOLD { - return Err(NetError::Protocol("peer misbehavior threshold")); + return Err(threshold_disconnect(session.as_deref())); } return Err(NetError::MessageTooLarge(n)); } @@ -1654,7 +1687,7 @@ pub async fn peer_session_with( follow.ban_score ); if follow.ban_score >= BAN_SCORE_THRESHOLD { - return Err(NetError::Protocol("peer misbehavior threshold")); + return Err(threshold_disconnect(session.as_deref())); } } continue; @@ -1672,7 +1705,7 @@ pub async fn peer_session_with( follow.ban_score ); if follow.ban_score >= BAN_SCORE_THRESHOLD { - return Err(NetError::Protocol("peer misbehavior threshold")); + return Err(threshold_disconnect(session.as_deref())); } continue; } @@ -1711,7 +1744,7 @@ pub async fn peer_session_with( "{}", misbehavior_disconnect_log(&peer_s, follow.ban_score) ); - return Err(NetError::Protocol("peer misbehavior threshold")); + return Err(threshold_disconnect(session.as_deref())); } } tip = tip_rx.recv() => { diff --git a/crates/rbitcoin-net/src/peer_tests.rs b/crates/rbitcoin-net/src/peer_tests.rs index 711d66b22..b2e5f9eb7 100644 --- a/crates/rbitcoin-net/src/peer_tests.rs +++ b/crates/rbitcoin-net/src/peer_tests.rs @@ -2976,6 +2976,177 @@ fn inbound_peer( peers.register(addr, addr, &ver, true, crate::peers::PeerConnType::Inbound) } +#[test] +fn inbound_netgroup_is_fixed_at_accept() { + let peers = crate::peers::PeerHub::new(); + let mk = |ip: [u8; 4]| { + let addr = std::net::SocketAddr::from((ip, 1)); + let ver = bitcoin::p2p::message_network::VersionMessage { + version: 70016, + services: bitcoin::p2p::ServiceFlags::NETWORK, + timestamp: 0, + receiver: bitcoin::p2p::address::Address::new(&addr, bitcoin::p2p::ServiceFlags::NONE), + sender: bitcoin::p2p::address::Address::new(&addr, bitcoin::p2p::ServiceFlags::NONE), + nonce: u64::from(ip[3]), + user_agent: "/rbitcoin:test/".into(), + start_height: 0, + relay: true, + }; + peers.register(addr, addr, &ver, true, crate::peers::PeerConnType::Inbound) + }; + let a = mk([1, 2, 3, 4]); + let b = mk([1, 2, 9, 9]); + let c = mk([1, 3, 0, 1]); + assert_eq!(a.netgroup(), b.netgroup(), "same /16 is one group"); + assert_ne!( + a.netgroup(), + c.netgroup(), + "a different /16 is another group" + ); + assert_eq!(a.netgroup(), crate::eviction::eviction_netgroup(a.addr)); +} + +#[test] +fn misbehavior_disconnect_refuses_the_same_address() { + let peers = crate::peers::PeerHub::new(); + let now = 1_700_000_000u64; + peers.set_mock_now(now); + let addr = std::net::SocketAddr::from(([9, 9, 9, 9], 8333)); + let ver = bitcoin::p2p::message_network::VersionMessage { + version: 70016, + services: bitcoin::p2p::ServiceFlags::NETWORK, + timestamp: 0, + receiver: bitcoin::p2p::address::Address::new(&addr, bitcoin::p2p::ServiceFlags::NONE), + sender: bitcoin::p2p::address::Address::new(&addr, bitcoin::p2p::ServiceFlags::NONE), + nonce: 9, + user_agent: "/rbitcoin:test/".into(), + start_height: 0, + relay: true, + }; + let peer = peers.register(addr, addr, &ver, true, crate::peers::PeerConnType::Inbound); + let mut score = 0u32; + punish_disconnect(&mut score, Some(peer.as_ref())); + let other_port = std::net::SocketAddr::from(([9, 9, 9, 9], 9999)); + assert!( + peers.inbound_discouraged(other_port), + "a misbehavior disconnect refuses that address" + ); + peers.set_mock_now(now + crate::peers::PeerHub::DISCOURAGE_TTL_SECS); + assert!( + !peers.inbound_discouraged(addr), + "the refusal ends after a day" + ); +} + +#[test] +fn misbehavior_disconnect_does_not_refuse_loopback() { + let peers = crate::peers::PeerHub::new(); + peers.set_mock_now(1_700_000_000); + let addr = std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, 8333)); + let ver = bitcoin::p2p::message_network::VersionMessage { + version: 70016, + services: bitcoin::p2p::ServiceFlags::NETWORK, + timestamp: 0, + receiver: bitcoin::p2p::address::Address::new(&addr, bitcoin::p2p::ServiceFlags::NONE), + sender: bitcoin::p2p::address::Address::new(&addr, bitcoin::p2p::ServiceFlags::NONE), + nonce: 1, + user_agent: "/rbitcoin:test/".into(), + start_height: 0, + relay: true, + }; + let peer = peers.register(addr, addr, &ver, true, crate::peers::PeerConnType::Inbound); + let mut score = 0u32; + punish_disconnect(&mut score, Some(peer.as_ref())); + assert!( + score >= BAN_SCORE_THRESHOLD, + "a loopback peer is still disconnected" + ); + let other_port = std::net::SocketAddr::from((std::net::Ipv4Addr::LOCALHOST, 18444)); + assert!( + !peers.inbound_discouraged(other_port), + "one loopback disconnect must not refuse every local connection" + ); +} + +#[test] +fn threshold_exit_refuses_the_address_without_punish_disconnect() { + let peers = crate::peers::PeerHub::new(); + let now = 1_700_000_100u64; + peers.set_mock_now(now); + let addr = std::net::SocketAddr::from(([8, 8, 4, 4], 8333)); + let ver = bitcoin::p2p::message_network::VersionMessage { + version: 70016, + services: bitcoin::p2p::ServiceFlags::NETWORK, + timestamp: 0, + receiver: bitcoin::p2p::address::Address::new(&addr, bitcoin::p2p::ServiceFlags::NONE), + sender: bitcoin::p2p::address::Address::new(&addr, bitcoin::p2p::ServiceFlags::NONE), + nonce: 4, + user_agent: "/rbitcoin:test/".into(), + start_height: 0, + relay: true, + }; + let peer = peers.register(addr, addr, &ver, true, crate::peers::PeerConnType::Inbound); + let err = threshold_disconnect(Some(peer.as_ref())); + assert!(matches!( + err, + crate::error::NetError::Protocol("peer misbehavior threshold") + )); + let other_port = std::net::SocketAddr::from(([8, 8, 4, 4], 9999)); + assert!( + peers.inbound_discouraged(other_port), + "a rate-limit or score-threshold exit refuses that address" + ); + + peers.set_noban(true); + let noban = peers.register( + std::net::SocketAddr::from(([1, 2, 3, 4], 8333)), + std::net::SocketAddr::from(([1, 2, 3, 4], 8333)), + &ver, + true, + crate::peers::PeerConnType::Inbound, + ); + let _ = threshold_disconnect(Some(noban.as_ref())); + assert!( + !peers.inbound_discouraged(std::net::SocketAddr::from(([1, 2, 3, 4], 1))), + "noban is not recorded" + ); +} + +#[test] +fn evicted_netgroup_waits_less_than_a_day_and_the_set_is_capped() { + let peers = crate::peers::PeerHub::new(); + let now = 1_800_000_000u64; + peers.set_mock_now(now); + let group = crate::eviction::eviction_netgroup("8.8.1.1:1".parse().unwrap()); + peers.note_slot_evict(group); + let same = "8.8.9.9:8333".parse().unwrap(); + assert!( + peers.inbound_discouraged(same), + "a netgroup that just lost a slot is refused" + ); + peers.set_mock_now(now + crate::peers::PeerHub::NETGROUP_SLOT_WAIT_SECS); + assert!( + !peers.inbound_discouraged(same), + "the netgroup wait is shorter than a day" + ); + peers.set_mock_now(now); + for i in 0..crate::peers::PeerHub::DISCOURAGE_CAP { + let ip = std::net::Ipv4Addr::from(i as u32); + peers.note_misbehavior_addr(std::net::IpAddr::V4(ip)); + } + let extra = std::net::SocketAddr::from(([255, 255, 255, 254], 1)); + peers.note_misbehavior_addr(extra.ip()); + assert!( + !peers.inbound_discouraged(extra), + "past the cap the set does not grow" + ); + let first = std::net::SocketAddr::from((std::net::Ipv4Addr::from(0u32), 1)); + assert!( + peers.inbound_discouraged(first), + "rows already stored stay until they expire" + ); +} + #[tokio::test] async fn inv_getdata_charges_send_budget() { use bitcoin::hashes::Hash; diff --git a/crates/rbitcoin-net/src/peers.rs b/crates/rbitcoin-net/src/peers.rs index 592636c75..b991984e6 100644 --- a/crates/rbitcoin-net/src/peers.rs +++ b/crates/rbitcoin-net/src/peers.rs @@ -319,6 +319,9 @@ pub struct LivePeer { addr_token_ms: AtomicU64, /// Unix seconds when this session was registered. connected_at: AtomicU64, + /// Netgroup fixed at accept. Inbound eviction compares this integer + /// and does not read asmap. + netgroup: u64, /// Skip INV for mempool txs with `accept_gen < floor` (post-verack privacy). inv_gen_floor: AtomicU64, /// Age-INV due-log cursor (`due_secs`, `accept_gen`). @@ -928,6 +931,10 @@ impl LivePeer { self.connected_at.load(Ordering::Relaxed) } + pub(crate) fn netgroup(&self) -> u64 { + self.netgroup + } + pub fn peer_hub(&self) -> Option> { self.owner.upgrade() } @@ -1309,6 +1316,10 @@ pub struct PeerHub { mempool: Mutex>>, /// Core `-whitelist` / `-whitebind` grants (`getpeerinfo.permissions`). net_perms: Mutex, + /// Misbehavior addresses → unix second the refusal ends. In memory only. + discouraged_addrs: Mutex>, + /// Netgroup → unix second a slot-loss refusal ends. + discouraged_groups: Mutex>, } fn canonical_bind(addr: SocketAddr) -> SocketAddr { @@ -1400,6 +1411,8 @@ impl PeerHub { asmap: Mutex::new(None), mempool: Mutex::new(None), net_perms: Mutex::new(crate::net_permissions::NetPermTable::default()), + discouraged_addrs: Mutex::new(HashMap::new()), + discouraged_groups: Mutex::new(HashMap::new()), }) } @@ -1863,6 +1876,78 @@ impl PeerHub { } } + /// One day. Not a configuration flag. + pub(crate) const DISCOURAGE_TTL_SECS: u64 = 24 * 60 * 60; + /// Shorter than [`Self::DISCOURAGE_TTL_SECS`]. A netgroup that just lost + /// an inbound slot cannot refill it immediately. + pub(crate) const NETGROUP_SLOT_WAIT_SECS: u64 = 10 * 60; + /// The in-memory set does not grow past this many addresses. + pub(crate) const DISCOURAGE_CAP: usize = 10_000; + + fn sweep_discouraged(&self, now: u64) { + self.discouraged_addrs + .lock() + .unwrap_or_else(|e| e.into_inner()) + .retain(|_, until| *until > now); + self.discouraged_groups + .lock() + .unwrap_or_else(|e| e.into_inner()) + .retain(|_, until| *until > now); + } + + /// Record `ip` until [`Self::DISCOURAGE_TTL_SECS`] from now. At the cap a + /// new address is ignored; an address already stored is refreshed. + pub(crate) fn note_misbehavior_addr(&self, ip: IpAddr) { + let now = self.now_secs(); + let until = now.saturating_add(Self::DISCOURAGE_TTL_SECS); + let mut g = self + .discouraged_addrs + .lock() + .unwrap_or_else(|e| e.into_inner()); + if !g.contains_key(&ip) && g.len() >= Self::DISCOURAGE_CAP { + return; + } + g.insert(ip, until); + } + + /// The netgroup that just lost an inbound slot waits + /// [`Self::NETGROUP_SLOT_WAIT_SECS`]. + pub(crate) fn note_slot_evict(&self, group: u64) { + let until = self + .now_secs() + .saturating_add(Self::NETGROUP_SLOT_WAIT_SECS); + let mut g = self + .discouraged_groups + .lock() + .unwrap_or_else(|e| e.into_inner()); + if !g.contains_key(&group) && g.len() >= Self::DISCOURAGE_CAP { + return; + } + g.insert(group, until); + } + + /// True when this inbound address or its accept-time netgroup is still + /// refused. Expired rows are dropped here, which is the only sweep. + pub(crate) fn inbound_discouraged(&self, addr: SocketAddr) -> bool { + let now = self.now_secs(); + self.sweep_discouraged(now); + if self + .discouraged_addrs + .lock() + .unwrap_or_else(|e| e.into_inner()) + .get(&addr.ip()) + .is_some_and(|until| *until > now) + { + return true; + } + let group = crate::eviction::eviction_netgroup(addr); + self.discouraged_groups + .lock() + .unwrap_or_else(|e| e.into_inner()) + .get(&group) + .is_some_and(|until| *until > now) + } + pub fn now_secs(&self) -> u64 { let mock = self.mock_now.load(Ordering::Acquire); if mock != 0 { @@ -2121,6 +2206,7 @@ impl PeerHub { addr_tokens: Mutex::new(crate::peer::ADDR_RELAY_BURST), addr_token_ms: AtomicU64::new(0), connected_at: AtomicU64::new(connected_at), + netgroup: crate::eviction::eviction_netgroup(endpoint.addr), inv_gen_floor: AtomicU64::new(0), age_inv_seen_due: AtomicU64::new(0), age_inv_seen_gen: AtomicU64::new(0), @@ -2606,7 +2692,7 @@ impl PeerHub { min_ping: minping, last_block: p.last_block.load(Ordering::Relaxed), last_tx: p.last_transaction.load(Ordering::Relaxed), - netgroup: crate::eviction::eviction_netgroup(p.addr), + netgroup: p.netgroup, noban: p.session_noban(), } }) @@ -2614,6 +2700,9 @@ impl PeerHub { let Some(id) = crate::eviction::select_inbound_eviction(cands) else { return false; }; + if let Some(p) = self.get(id) { + self.note_slot_evict(p.netgroup()); + } rbitcoin_log::info!("p2p: evict inbound peer={id} (inbound full)"); self.disconnect_id(id) } diff --git a/crates/rbitcoin-net/src/service.rs b/crates/rbitcoin-net/src/service.rs index 50d4aacbf..c41269536 100644 --- a/crates/rbitcoin-net/src/service.rs +++ b/crates/rbitcoin-net/src/service.rs @@ -484,6 +484,11 @@ fn spawn_inbound_accept( let accept = tokio::time::timeout(Duration::from_millis(200), listener.accept()).await; match accept { Ok(Ok((stream, peer_addr))) => { + if peers.inbound_discouraged(peer_addr) { + rbitcoin_log::info!("p2p: reject discouraged inbound {peer_addr}"); + drop(stream); + continue; + } let permit = match inbound_sem.clone().try_acquire_owned() { Ok(p) => p, Err(_) => { diff --git a/crates/rbitcoin-store/src/scripthash.rs b/crates/rbitcoin-store/src/scripthash.rs index eefa3c629..67d055e4f 100644 --- a/crates/rbitcoin-store/src/scripthash.rs +++ b/crates/rbitcoin-store/src/scripthash.rs @@ -2790,6 +2790,12 @@ impl ScriptHashTable { /// packed chain. Clear ingest keys this shard's main now owns. Ingest body /// bytes stay; only the head slot is soft-cleared. fn drop_ingest_covered_by_packed_main(&self, shard: usize) -> Result<(), StoreError> { + // Mainnet ingest is 2^25 slots. Known-empty has no row that can hide + // the packed chain. Walking it once per shard holds tip entry past + // the RPC cookie window. + if self.ingest.lock().unwrap().is_known_empty() { + return Ok(()); + } let Some(slot) = self.sorted_main.get(shard) else { return Ok(()); }; diff --git a/crates/rbitcoin-store/src/scripthash_tests.rs b/crates/rbitcoin-store/src/scripthash_tests.rs index d34e4064a..60826397d 100644 --- a/crates/rbitcoin-store/src/scripthash_tests.rs +++ b/crates/rbitcoin-store/src/scripthash_tests.rs @@ -1395,6 +1395,27 @@ fn pack_one_shard() { } } +/// Mainnet ingest is 2^25 slots. Sealing a shard must not walk that table when +/// occupancy is already known to be zero: one pass per shard holds tip entry +/// past the RPC cookie window. A truncated ingest file makes a walk fail. +#[test] +fn empty_ingest_shard_seal_does_not_walk_slots() { + let dir = tmp(); + let t = ScriptHashTable::create(dir.path()).unwrap(); + assert!(t.head_is_empty()); + let ingest = ingest_path(dir.path()); + std::fs::OpenOptions::new() + .write(true) + .open(&ingest) + .unwrap() + .set_len(0) + .unwrap(); + let session = t.pack_shard_session(0).unwrap(); + let pack = session.finish_pack().unwrap(); + t.publish_packed_shard(0, pack) + .expect("known-empty ingest is not scanned"); +} + fn shard0_key(i: u8) -> [u8; 32] { let mut k = [0u8; 32]; k[0] = i & 0x3f; diff --git a/crates/rbitcoin-test/tests/integration_multinode.rs b/crates/rbitcoin-test/tests/integration_multinode.rs index 5c7ef55a3..228d790e7 100644 --- a/crates/rbitcoin-test/tests/integration_multinode.rs +++ b/crates/rbitcoin-test/tests/integration_multinode.rs @@ -1454,21 +1454,41 @@ fn pin_select_node_to_evict_ranking() { noban: false, } } + let mut covered = Vec::new(); + for i in 0..4 { + covered.push(cand(i, 100 + i, Some(0.05), 1000 + i, 0)); + } + for i in 4..9 { + covered.push(cand(i, 200 + i, Some(0.5), 0, 0)); + } + for i in 9..13 { + covered.push(cand(i, 300 + i, Some(0.05), 0, 1000 + i)); + } + for i in 13..21 { + covered.push(cand(i, 400 + i, Some(0.01), 0, 0)); + } + assert!( + select_inbound_eviction(covered).is_none(), + "longest-connected protection covers the remaining slow peers" + ); + + // Ten slow peers: netgroup and the fourth block slot take one, the eight + // longest-connected stay, and the newest slow peer is the victim. let mut cands = Vec::new(); for i in 0..4 { cands.push(cand(i, 100 + i, Some(0.05), 1000 + i, 0)); } - for i in 4..9 { + for i in 4..14 { cands.push(cand(i, 200 + i, Some(0.5), 0, 0)); } - for i in 9..13 { + for i in 14..18 { cands.push(cand(i, 300 + i, Some(0.05), 0, 1000 + i)); } - for i in 13..21 { + for i in 18..26 { 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_eq!(victim, 13, "evict the newest slow peer, victim={victim}"); } #[tokio::test] diff --git a/docs/core-functional.md b/docs/core-functional.md index 323797f88..6094d324c 100644 --- a/docs/core-functional.md +++ b/docs/core-functional.md @@ -180,6 +180,7 @@ Shim-only scripts that used to be `run` stay skip: `feature_help.py` still exists so `rpc_getblockstats.py`'s rename-file needle can run; the rest of that script is `submitblock` + archive `getblockstats`. `feature_port.py` is `run`: `-bind`/`-port` become `--listen` on those sockets. +`p2p_eviction.py` is skip (`core-net-policy`): the 21-peer set in that script is fully protected, including a share of the longest-connected peers, so the extra inbound is rejected. Core accepts it and evicts one slow peer. `feature_filelock.py` is skip (`harness`): the node exclusive-locks `{datadir}` (unit + shim InitError tests). Core `-wallet` SQLite flock and `-blocksdir` `blocks/.lock` are not product. `rpc_orphans.py` is `run`: hidden diff --git a/docs/external_findings/052-livera-review-index.md b/docs/external_findings/052-livera-review-index.md index bde016c61..ad6ed73e2 100644 --- a/docs/external_findings/052-livera-review-index.md +++ b/docs/external_findings/052-livera-review-index.md @@ -9,11 +9,11 @@ steps. | C1 | critical | Unbounded parent-request tracker | fixed | `parent_req_stops_at_per_peer_cap`, `parent_req_stops_at_global_cap`, `second_peer_take_due_does_not_drop_other_peers_keys` ([053](./053-parent-req-cap.md)) | | N2 | low | Wtxid follow-up requested as a txid | fixed | `wtxid_followup_is_requested_as_wtx` ([054](./054-wtxid-getdata.md)) | | H1 | high | Decoy and invalid-type packets skip the rate window | fixed | `decoy_packet_is_handed_to_the_rate_hook` ([055](./055-decoy-rate.md)) | -| H2 | high | Inbound eviction drops the longest-connected peer | open | — | -| H2-ban | high | Misbehavior disconnect is not remembered | open | — | +| H2 | high | Inbound eviction drops the longest-connected peer | fixed | `eviction_drops_the_newest_in_the_largest_netgroup` ([060](./060-evict-newest-netgroup.md)) | +| H2-ban | high | Misbehavior disconnect is not remembered | fixed | `misbehavior_disconnect_refuses_the_same_address` ([061](./061-misbehavior-remembered.md)) | | H3 | high | REST always on and shares the RPC work queue | open | — | | H4 | medium | Silent-payment unsubscribe logs the scan secret | open | — | -| H5 | high | IBD reader credits unsolicited data as progress | open | — | +| H5 | high | IBD reader credits unsolicited data as progress | fixed | `unsolicited_block_does_not_refresh_progress` ([062](./062-ibd-requested-progress.md)) | | N1 | medium | Inv getdata does not charge the send budget | fixed | `inv_getdata_charges_send_budget` ([056](./056-inv-getdata-budget.md)) | | M1 | medium | Block getdata can queue past the send budget | fixed | `getdata_stops_when_send_budget_is_already_over` ([057](./057-block-getdata-budget.md)) | | M2 | medium | Silent-payment scan span is unbounded when start is set | open | — | diff --git a/docs/external_findings/060-evict-newest-netgroup.md b/docs/external_findings/060-evict-newest-netgroup.md new file mode 100644 index 000000000..e431f1e02 --- /dev/null +++ b/docs/external_findings/060-evict-newest-netgroup.md @@ -0,0 +1,9 @@ +# 060 — Evict the newest inbound in the largest netgroup + +**Severity:** high +**Status:** fixed +**Found by:** Stephan Livera, 2026-10-02 (H2) + +Inbound eviction disconnected the oldest unprotected peer. It protects a share of the longest-connected peers, then disconnects the newest peer in the largest netgroup. The netgroup key is fixed when the peer is accepted. Asmap stays outbound-only. + +**Regression:** `rbitcoin-net` `eviction_drops_the_newest_in_the_largest_netgroup`. diff --git a/docs/external_findings/061-misbehavior-remembered.md b/docs/external_findings/061-misbehavior-remembered.md new file mode 100644 index 000000000..be7203561 --- /dev/null +++ b/docs/external_findings/061-misbehavior-remembered.md @@ -0,0 +1,9 @@ +# 061 — Misbehavior disconnect is remembered in memory + +**Severity:** high +**Status:** fixed +**Found by:** Stephan Livera, 2026-10-02 (H2-ban) + +A misbehavior disconnect was forgotten, so the same address could reconnect immediately. The address is refused from a capped in-memory set. There is no ban file and no new ban-time flag. A netgroup that just lost an inbound slot waits ten minutes. + +**Regression:** `rbitcoin-net` `misbehavior_disconnect_refuses_the_same_address`. diff --git a/docs/external_findings/062-ibd-requested-progress.md b/docs/external_findings/062-ibd-requested-progress.md new file mode 100644 index 000000000..c9ad606c5 --- /dev/null +++ b/docs/external_findings/062-ibd-requested-progress.md @@ -0,0 +1,9 @@ +# 062 — Only a requested block moves the IBD stall clock + +**Severity:** high +**Status:** fixed +**Found by:** Stephan Livera, 2026-10-02 (H5) + +During initial download, unsolicited blocks and decoys refreshed the stall clock, and unsolicited blocks were queued. Only a block this node requested moves the clock or is queued. Requested bodies are still accepted past 16 MB/s. The requested-body channel stays unbounded, and the reader does not wait on a decode permit. + +**Regression:** `rbitcoin-net` `unsolicited_block_does_not_refresh_progress`. diff --git a/docs/external_findings/README.md b/docs/external_findings/README.md index e8aafb3e8..8b1d6609d 100644 --- a/docs/external_findings/README.md +++ b/docs/external_findings/README.md @@ -64,6 +64,9 @@ rbitcoin reference, or redteam static analysis). Numbered reports live beside th | [057](./057-block-getdata-budget.md) | medium | Block serving stops when the send budget is over | fixed | `getdata_stops_when_send_budget_is_already_over` | | [058](./058-tx-inv-batch.md) | medium | Mempool announcements batch into one inv | fixed | `tx_inv_over_one_thousand_is_two_messages` | | [059](./059-rate-window-boundary.md) | low | Rate window keeps the previous second | fixed | `rate_limiter_boundary_does_not_grant_a_second_budget` | +| [060](./060-evict-newest-netgroup.md) | high | Evict the newest inbound in the largest netgroup | fixed | `eviction_drops_the_newest_in_the_largest_netgroup` | +| [061](./061-misbehavior-remembered.md) | high | Misbehavior disconnect is remembered in memory | fixed | `misbehavior_disconnect_refuses_the_same_address` | +| [062](./062-ibd-requested-progress.md) | high | Only a requested block moves the IBD stall clock | fixed | `unsolicited_block_does_not_refresh_progress` | **012–021:** fuzzamoto differential report (`rbitcoin-report.tar.gz`, baseline `8f3990f`). Report-local 001–010 are **renumbered** here. Identity/BIP30 diff --git a/scripts/core-functional/inventory.toml b/scripts/core-functional/inventory.toml index d967f9e87..6dbaad092 100644 --- a/scripts/core-functional/inventory.toml +++ b/scripts/core-functional/inventory.toml @@ -605,7 +605,9 @@ reason = "core-net-policy" [[test]] name = "p2p_eviction.py" -status = "run" +status = "skip" +reason = "core-net-policy" +analog = "21 inbound peers (4 block, 5 slow, 4 tx, 8 fast ping) are all protected, including longest-connected; the extra inbound is rejected (eviction::tests::eviction_protects_block_tx_ping_and_netgroup). Core accepts that peer and evicts one slow peer." [[test]] name = "p2p_feefilter.py"