From 55175a5a6ec65efc1609857f7002b9e1e8dad557 Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Sun, 30 Aug 2026 07:25:39 -0700 Subject: [PATCH 1/9] ibd: add PeerRate EWMA core with stall predicate Integer per-peer rate accrued only while getdata is in flight. Idle rebaseline freezes the EWMA; stall is now - max(rx, work start). Co-authored-by: Cursor --- crates/rbitcoin-net/src/ibd/mod.rs | 1 + crates/rbitcoin-net/src/ibd/rate.rs | 173 ++++++++++++++++++++++++++++ 2 files changed, 174 insertions(+) create mode 100644 crates/rbitcoin-net/src/ibd/rate.rs diff --git a/crates/rbitcoin-net/src/ibd/mod.rs b/crates/rbitcoin-net/src/ibd/mod.rs index b591bc04..add76795 100644 --- a/crates/rbitcoin-net/src/ibd/mod.rs +++ b/crates/rbitcoin-net/src/ibd/mod.rs @@ -23,6 +23,7 @@ mod path; mod peer_io; mod perf_log; mod progress; +mod rate; mod reorg; mod state; mod status; diff --git a/crates/rbitcoin-net/src/ibd/rate.rs b/crates/rbitcoin-net/src/ibd/rate.rs new file mode 100644 index 00000000..8551579e --- /dev/null +++ b/crates/rbitcoin-net/src/ibd/rate.rs @@ -0,0 +1,173 @@ +//! Per-peer IBD byte-rate EWMA, accrued only while block getdata is in flight. + +/// EWMA time constant (ms). +pub(crate) const TAU_MS: u64 = 15_000; +/// Minimum dt to fold a sample (ms). +pub(crate) const MIN_SAMPLE_MS: u64 = 250; +/// `bps()` becomes `Some` after this much inflight-accrued time (ms). +pub(crate) const RANK_MATURE_MS: u64 = 5_000; +/// Relative-slow sample maturity (ms of inflight-accrued EWMA). +pub(crate) const RELSLOW_ACTIVE_MS: u64 = 30_000; +/// Qualifying stream delta for rx progress. +pub(crate) const PROGRESS_STEP: u64 = 64 * 1024; + +/// Integer EWMA of received bytes, sampled on the IBD main thread. +/// +/// Accrues only while the peer has block getdata in flight. Idle time rebaselines +/// the byte cursor so it does not dilute the rate. `progress_ms` is last qualifying +/// rx (stream ≥ [`PROGRESS_STEP`] or event-path `note_rx`); `work_started_ms` is the +/// last empty→nonempty getdata. Stall is `now - max(progress, work_started) > stall` +/// while inflight. +#[derive(Clone, Copy, Debug, Default)] +pub(crate) struct PeerRate { + ewma: u64, + pub(crate) active_ms: u64, + pub(crate) progress_ms: u64, + pub(crate) work_started_ms: u64, + last_bytes: u64, + last_ms: Option, +} + +impl PeerRate { + /// Fold `bytes_total` at `now_ms`. `inflight` false freezes the EWMA. + pub(crate) fn sample(&mut self, now_ms: u64, bytes_total: u64, inflight: bool) { + if !inflight { + self.last_bytes = bytes_total; + self.last_ms = Some(now_ms); + return; + } + let Some(prev_ms) = self.last_ms else { + self.last_bytes = bytes_total; + self.last_ms = Some(now_ms); + return; + }; + if now_ms < prev_ms { + self.last_bytes = bytes_total; + self.last_ms = Some(now_ms); + return; + } + let dt = now_ms.saturating_sub(prev_ms); + if dt < MIN_SAMPLE_MS { + return; + } + let delta = bytes_total.saturating_sub(self.last_bytes); + let inst = delta.saturating_mul(1000) / dt; + let den = TAU_MS.saturating_add(dt); + self.ewma = self + .ewma + .saturating_mul(TAU_MS) + .saturating_add(inst.saturating_mul(dt)) + / den.max(1); + self.active_ms = self.active_ms.saturating_add(dt); + if delta >= PROGRESS_STEP { + self.progress_ms = now_ms; + } + self.last_bytes = bytes_total; + self.last_ms = Some(now_ms); + } + + pub(crate) fn note_work_started(&mut self, now_ms: u64) { + self.work_started_ms = now_ms; + } + + pub(crate) fn note_rx(&mut self, now_ms: u64) { + self.progress_ms = now_ms; + } + + pub(crate) fn bps(&self) -> Option { + if self.active_ms >= RANK_MATURE_MS { + Some(self.ewma) + } else { + None + } + } + + pub(crate) fn stalled(&self, now_ms: u64, stall_ms: u64, inflight: bool) -> bool { + if !inflight { + return false; + } + let last = self.progress_ms.max(self.work_started_ms); + now_ms.saturating_sub(last) > stall_ms + } + +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn relslow_active_is_longer_than_rank() { + let mut r = PeerRate::default(); + r.sample(0, 0, true); + r.sample(RANK_MATURE_MS, RANK_MATURE_MS * 1_000, true); + assert!(r.bps().is_some()); + assert!(r.active_ms < RELSLOW_ACTIVE_MS); + } + + + #[test] + fn idle_rebaseline_freezes_ewma() { + let mut r = PeerRate::default(); + r.sample(0, 0, true); + r.sample(5_000, 5_000_000, true); + let frozen = r.bps().expect("mature after 5s"); + r.sample(15_000, 5_000_000 + 50_000_000, false); + assert_eq!(r.bps(), Some(frozen)); + } + + #[test] + fn active_zero_delta_decays_ewma() { + let mut r = PeerRate::default(); + r.sample(0, 0, true); + r.sample(5_000, 5_000_000, true); + let high = r.bps().expect("mature"); + r.sample(10_000, 5_000_000, true); + r.sample(15_000, 5_000_000, true); + assert!(r.bps().unwrap() < high); + } + + #[test] + fn bps_none_until_rank_mature() { + let mut r = PeerRate::default(); + r.sample(0, 0, true); + r.sample(4_000, 4_000_000, true); + assert!(r.bps().is_none()); + r.sample(5_000, 5_000_000, true); + assert!(r.bps().is_some()); + } + + #[test] + fn progress_marks_on_64k_not_on_small_delta() { + let mut r = PeerRate::default(); + r.sample(0, 0, true); + r.sample(1_000, 1_000, true); + assert_eq!(r.progress_ms, 0); + r.sample(2_000, 1_000 + PROGRESS_STEP, true); + assert_eq!(r.progress_ms, 2_000); + } + + #[test] + fn stalled_uses_work_started_when_no_rx() { + let mut r = PeerRate::default(); + r.note_work_started(0); + assert!(!r.stalled(29_999, 30_000, true)); + assert!(r.stalled(30_001, 30_000, true)); + } + + #[test] + fn stalled_false_when_idle_inflight_empty() { + let mut r = PeerRate::default(); + r.note_work_started(0); + assert!(!r.stalled(40_000, 30_000, false)); + } + + #[test] + fn stalled_false_when_rx_within_grace() { + let mut r = PeerRate::default(); + r.note_work_started(0); + r.note_rx(20_000); + assert!(!r.stalled(45_000, 30_000, true)); + assert!(r.stalled(50_001, 30_000, true)); + } +} From 320c57693ca1dc701cb71b456c77158c72958164 Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Sun, 30 Aug 2026 07:27:46 -0700 Subject: [PATCH 2/9] ibd: count all streamed peer bytes in the reader Stall clocks stay on the event path; the reader only fetch_adds ciphertext/plaintext progress into bytes_rx_total for the EWMA. Co-authored-by: Cursor --- crates/rbitcoin-net/src/ibd/assign.rs | 1 + crates/rbitcoin-net/src/ibd/dial.rs | 1 + .../src/ibd/events/confirm_reject_tests.rs | 7 ++++ crates/rbitcoin-net/src/ibd/peer_io.rs | 34 ++++++++++++------- 4 files changed, 31 insertions(+), 12 deletions(-) diff --git a/crates/rbitcoin-net/src/ibd/assign.rs b/crates/rbitcoin-net/src/ibd/assign.rs index a2874bb2..047bace9 100644 --- a/crates/rbitcoin-net/src/ibd/assign.rs +++ b/crates/rbitcoin-net/src/ibd/assign.rs @@ -763,6 +763,7 @@ mod tests { connected_ms: 1, first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), + bytes_rx_total: Arc::new(AtomicU64::new(0)), alive: true, task, } diff --git a/crates/rbitcoin-net/src/ibd/dial.rs b/crates/rbitcoin-net/src/ibd/dial.rs index 0f3230fe..5b81f610 100644 --- a/crates/rbitcoin-net/src/ibd/dial.rs +++ b/crates/rbitcoin-net/src/ibd/dial.rs @@ -544,6 +544,7 @@ mod tests { connected_ms: 0, first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), + bytes_rx_total: Arc::new(AtomicU64::new(0)), alive, task, } 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 095769a2..1a4f3402 100644 --- a/crates/rbitcoin-net/src/ibd/events/confirm_reject_tests.rs +++ b/crates/rbitcoin-net/src/ibd/events/confirm_reject_tests.rs @@ -886,6 +886,7 @@ fn confirmed_height_mids_blocked_while_densify_ahead_leaves_tip_hole() { connected_ms: 1, first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), + bytes_rx_total: Arc::new(AtomicU64::new(0)), alive: true, task, }; @@ -1083,6 +1084,7 @@ fn zombie_pending_mid_at_confirmed_height_never_reget() { connected_ms: 1, first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), + bytes_rx_total: Arc::new(AtomicU64::new(0)), alive: true, task, }; @@ -1656,6 +1658,7 @@ fn apply_peer_event_body_and_control_surface() { connected_ms: 1, first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), + bytes_rx_total: Arc::new(AtomicU64::new(0)), alive: true, task, } @@ -1907,6 +1910,7 @@ fn apply_peer_event_block_framed_bq_horizon_and_headers_done() { connected_ms: 1, first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), + bytes_rx_total: Arc::new(AtomicU64::new(0)), alive: true, task, } @@ -2173,6 +2177,7 @@ fn block_framed_raw_offers_body_queue_with_confirm_feed() { connected_ms: 1, first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), + bytes_rx_total: Arc::new(AtomicU64::new(0)), alive: true, task, } @@ -2320,6 +2325,7 @@ fn known_headers_re_admit_to_ordered_after_tip_drain() { connected_ms: 1, first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), + bytes_rx_total: Arc::new(AtomicU64::new(0)), alive: true, task, } @@ -2480,6 +2486,7 @@ fn path_slot_first_wins_chained_via_headers() { connected_ms: 1, first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), + bytes_rx_total: Arc::new(AtomicU64::new(0)), alive: true, task, } diff --git a/crates/rbitcoin-net/src/ibd/peer_io.rs b/crates/rbitcoin-net/src/ibd/peer_io.rs index 1193092e..fc751477 100644 --- a/crates/rbitcoin-net/src/ibd/peer_io.rs +++ b/crates/rbitcoin-net/src/ibd/peer_io.rs @@ -102,6 +102,8 @@ pub(crate) struct PeerSlot { pub first_data_ms: AtomicU64, /// Cumulative block payload bytes (speed sample). pub bytes_rx: AtomicU64, + /// All streamed wire bytes (EWMA input). Reader-only `fetch_add`. + pub bytes_rx_total: Arc, pub alive: bool, pub task: JoinHandle<()>, } @@ -153,6 +155,13 @@ pub(crate) fn touch_block_progress(ms: &AtomicU64) { ms.store(ibd_mono_ms(), Ordering::Relaxed); } +pub(crate) fn note_stream_bytes(counter: &AtomicU64, n: u64) { + if n == 0 { + return; + } + counter.fetch_add(n, Ordering::Relaxed); +} + pub(crate) fn note_block_progress(slots: &mut [PeerSlot], peer: usize) { if let Some(s) = slots.iter_mut().find(|s| s.id == peer) { touch_block_progress(&s.block_progress_ms); @@ -196,7 +205,8 @@ pub(crate) async fn spawn_peer( // stall the receive half and look like a peer stall). let (out_tx, mut out_rx) = mpsc::unbounded_channel::(); let block_progress_ms = Arc::new(AtomicU64::new(ibd_mono_ms())); - let progress_io = Arc::clone(&block_progress_ms); + let bytes_rx_total = Arc::new(AtomicU64::new(0)); + let bytes_io = Arc::clone(&bytes_rx_total); // Parent owns concurrent read + write tasks. Aborting the parent (PeerSlot // Drop / stall disconnect) must abort both children — plain JoinHandle drop @@ -224,11 +234,9 @@ pub(crate) async fn spawn_peer( let mut prog_mark = 0usize; loop { let frame = read_v2_frame_with_progress(&mut reader, magic, |buffered| { - const STEP: usize = 64 * 1024; - if buffered >= prog_mark + STEP || buffered <= STEP { - prog_mark = buffered; - touch_block_progress(&progress_io); - } + let delta = buffered.saturating_sub(prog_mark); + note_stream_bytes(&bytes_io, delta as u64); + prog_mark = buffered; }) .await; prog_mark = 0; @@ -241,10 +249,6 @@ pub(crate) async fn spawn_peer( continue; } - if frame.is_block() || frame.is_notfound() { - touch_block_progress(&progress_io); - } - if frame.is_block() { match frame.block_hash_from_header() { Some(hash) if frame.payload.len() >= 80 => { @@ -267,7 +271,6 @@ pub(crate) async fn spawn_peer( continue; } - let progress = Arc::clone(&progress_io); let sinks_d = sinks_r.clone(); // Non-block: decode off-thread. Never await a decode permit // on the reader (stalls TCP). Soft budgets gate *requests* only. @@ -282,7 +285,6 @@ pub(crate) async fn spawn_peer( }); } NetworkMessage::NotFound(inv) => { - touch_block_progress(&progress); let hashes: Vec = inv .iter() .filter_map(|i| match i { @@ -452,6 +454,7 @@ pub(crate) async fn spawn_peer( connected_ms: ibd_mono_ms(), first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), + bytes_rx_total, alive: true, task, }) @@ -533,6 +536,7 @@ mod tests { connected_ms: 1, first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), + bytes_rx_total: Arc::new(AtomicU64::new(0)), alive: true, task, } @@ -583,6 +587,12 @@ mod tests { let sample = s.speed_sample().expect("bps after ≥64KiB"); assert!(sample.1 > 0); + 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); + touch_block_progress(&s.block_progress_ms); assert!(s.block_progress_ms.load(Ordering::Relaxed) > 0); note_block_progress(std::slice::from_mut(&mut s), 7); From 154aa607942dec7ca98584a4626b396d08613fca Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Sun, 30 Aug 2026 07:31:40 -0700 Subject: [PATCH 3/9] ibd: sample PeerRate in the loop and rank by EWMA Tip-hole peer order and AddrMan note_speed on death use the inflight EWMA instead of lifetime complete-block speed_sample. Co-authored-by: Cursor --- crates/rbitcoin-net/src/ibd/assign.rs | 24 +++++++------------ crates/rbitcoin-net/src/ibd/dial.rs | 1 + .../src/ibd/events/confirm_reject_tests.rs | 7 ++++++ crates/rbitcoin-net/src/ibd/events/mod.rs | 4 +++- crates/rbitcoin-net/src/ibd/mod.rs | 3 ++- crates/rbitcoin-net/src/ibd/peer_io.rs | 19 +++++++++++++++ crates/rbitcoin-net/src/ibd/rate.rs | 2 -- 7 files changed, 41 insertions(+), 19 deletions(-) diff --git a/crates/rbitcoin-net/src/ibd/assign.rs b/crates/rbitcoin-net/src/ibd/assign.rs index 047bace9..6c30340c 100644 --- a/crates/rbitcoin-net/src/ibd/assign.rs +++ b/crates/rbitcoin-net/src/ibd/assign.rs @@ -570,7 +570,7 @@ fn demote_zombie_pending_for_fetch( const TIP_HOLE_INFLIGHT_STALE: Duration = Duration::from_secs(6); /// Rank alive peer ids for tip-hole getdata: prefer peers not in `avoid`, then -/// higher live `speed_sample` bps, then lower id. Unsampled peers sort last +/// higher live EWMA bps, then lower id. Unsampled peers sort last /// among non-avoided (bps=0). pub(crate) fn rank_tip_hole_peers( slots: &[PeerSlot], @@ -586,8 +586,7 @@ pub(crate) fn rank_tip_hole_peers( slots .iter() .find(|s| s.id == pid && s.alive) - .and_then(|s| s.speed_sample()) - .map(|(_, bps)| bps) + .and_then(|s| s.rate.bps()) .unwrap_or(0) }; bps(b).cmp(&bps(a)).then_with(|| a.cmp(&b)) @@ -764,6 +763,7 @@ mod tests { first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), bytes_rx_total: Arc::new(AtomicU64::new(0)), + rate: Default::default(), alive: true, task, } @@ -1137,10 +1137,9 @@ mod tests { let _ = std::fs::remove_dir_all(dir); } - /// Tip-hole cover prefers higher speed_sample bps peers first. + /// Tip-hole cover prefers higher EWMA bps peers first. #[test] fn cover_tip_holes_prefers_fast_peers() { - use super::super::peer_io::ibd_mono_ms; let (dir, hub) = tmp_hub(); hub.ensure_genesis().unwrap(); let mut st = IbdWorkState::new( @@ -1155,17 +1154,12 @@ mod tests { st.height_to_hash.insert(ht, hole); st.body.mark_missing(hole); - // Peer 0 slow, peer 1 fast, peer 2 medium — inject lifetime samples. - let now = ibd_mono_ms().max(2_000); - for (i, bps_equiv_bytes) in [(0usize, 100_000u64), (1, 10_000_000u64), (2, 1_000_000u64)] { - st.slots[i].connected_ms = now.saturating_sub(2_000); + // Peer 0 slow, peer 1 fast, peer 2 medium — inject mature EWMA samples. + for (i, bytes_per_sec) in [(0usize, 100_000u64), (1, 10_000_000u64), (2, 1_000_000u64)] { + st.slots[i].rate.sample(0, 0, true); st.slots[i] - .first_data_ms - .store(now.saturating_sub(1_000), Ordering::Relaxed); - // bytes such that bps ≈ bytes*1000/1000ms = bytes for ~1s elapsed - st.slots[i] - .bytes_rx - .store(bps_equiv_bytes, Ordering::Relaxed); + .rate + .sample(5_000, bytes_per_sec.saturating_mul(5), true); } // Cap want to 2 so only the top two speeds get work if ranking works. // TIP_HOLE_MAX_PEERS is 4 but we only have 3 peers — all may get work. diff --git a/crates/rbitcoin-net/src/ibd/dial.rs b/crates/rbitcoin-net/src/ibd/dial.rs index 5b81f610..4731ff74 100644 --- a/crates/rbitcoin-net/src/ibd/dial.rs +++ b/crates/rbitcoin-net/src/ibd/dial.rs @@ -545,6 +545,7 @@ mod tests { first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), bytes_rx_total: Arc::new(AtomicU64::new(0)), + rate: Default::default(), alive, task, } 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 1a4f3402..36dbf460 100644 --- a/crates/rbitcoin-net/src/ibd/events/confirm_reject_tests.rs +++ b/crates/rbitcoin-net/src/ibd/events/confirm_reject_tests.rs @@ -887,6 +887,7 @@ fn confirmed_height_mids_blocked_while_densify_ahead_leaves_tip_hole() { first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), bytes_rx_total: Arc::new(AtomicU64::new(0)), + rate: Default::default(), alive: true, task, }; @@ -1085,6 +1086,7 @@ fn zombie_pending_mid_at_confirmed_height_never_reget() { first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), bytes_rx_total: Arc::new(AtomicU64::new(0)), + rate: Default::default(), alive: true, task, }; @@ -1659,6 +1661,7 @@ fn apply_peer_event_body_and_control_surface() { first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), bytes_rx_total: Arc::new(AtomicU64::new(0)), + rate: Default::default(), alive: true, task, } @@ -1911,6 +1914,7 @@ fn apply_peer_event_block_framed_bq_horizon_and_headers_done() { first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), bytes_rx_total: Arc::new(AtomicU64::new(0)), + rate: Default::default(), alive: true, task, } @@ -2178,6 +2182,7 @@ fn block_framed_raw_offers_body_queue_with_confirm_feed() { first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), bytes_rx_total: Arc::new(AtomicU64::new(0)), + rate: Default::default(), alive: true, task, } @@ -2326,6 +2331,7 @@ fn known_headers_re_admit_to_ordered_after_tip_drain() { first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), bytes_rx_total: Arc::new(AtomicU64::new(0)), + rate: Default::default(), alive: true, task, } @@ -2487,6 +2493,7 @@ fn path_slot_first_wins_chained_via_headers() { first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), bytes_rx_total: Arc::new(AtomicU64::new(0)), + rate: Default::default(), alive: true, task, } diff --git a/crates/rbitcoin-net/src/ibd/events/mod.rs b/crates/rbitcoin-net/src/ibd/events/mod.rs index 2a8ee78f..0d46bb21 100644 --- a/crates/rbitcoin-net/src/ibd/events/mod.rs +++ b/crates/rbitcoin-net/src/ibd/events/mod.rs @@ -432,7 +432,9 @@ pub(crate) fn apply_peer_event( PeerEvent::Dead { peer, reason } => { warn!("ibd: peer[{peer}] dead: {reason}"); if let Some(s) = st.slots.iter().find(|s| s.id == peer) { - if let Some((lat, bps)) = s.speed_sample() { + if let Some(bps) = s.rate.bps() { + let first = s.first_data_ms.load(Ordering::Relaxed); + let lat = first.saturating_sub(s.connected_ms); peer_book.note_speed(s.addr, lat, bps); } } diff --git a/crates/rbitcoin-net/src/ibd/mod.rs b/crates/rbitcoin-net/src/ibd/mod.rs index add76795..9366560c 100644 --- a/crates/rbitcoin-net/src/ibd/mod.rs +++ b/crates/rbitcoin-net/src/ibd/mod.rs @@ -48,7 +48,7 @@ use exit::{ ibd_caught_up, path_drained, should_unlatch_headers_done, AllPeersDead, }; use path::{path_hashes_above_tip, seed_work_path_from_store, work_path_tips}; -use peer_io::{PeerCmd, PeerEvent, PeerEventSinks}; +use peer_io::{sample_peer_rates, PeerCmd, PeerEvent, PeerEventSinks}; use progress::{ claim_ready, format_progress_line, ibd_pct, work_chain_progress, ProgressLineInput, TipRateTracker, @@ -664,6 +664,7 @@ pub async fn ibd_cancellable( last_progress = Instant::now(); } + sample_peer_rates(&mut st.slots, peer_io::ibd_mono_ms()); let now = Instant::now(); disconnect_stalled_block_peers( &mut st.slots, diff --git a/crates/rbitcoin-net/src/ibd/peer_io.rs b/crates/rbitcoin-net/src/ibd/peer_io.rs index fc751477..dcbffc56 100644 --- a/crates/rbitcoin-net/src/ibd/peer_io.rs +++ b/crates/rbitcoin-net/src/ibd/peer_io.rs @@ -4,6 +4,7 @@ //! - reader: decrypt frame + cheap ping handling; heavy decode off-thread //! - writer: encode offloaded for heavy payloads; then encrypt + write +use super::rate::PeerRate; use crate::codec::MAX_INV_SIZE; use crate::error::NetError; use crate::msg_decode::spawn_decode_then_with_err; @@ -104,6 +105,7 @@ pub(crate) struct PeerSlot { pub bytes_rx: AtomicU64, /// All streamed wire bytes (EWMA input). Reader-only `fetch_add`. pub bytes_rx_total: Arc, + pub rate: PeerRate, pub alive: bool, pub task: JoinHandle<()>, } @@ -162,6 +164,16 @@ pub(crate) fn note_stream_bytes(counter: &AtomicU64, n: u64) { 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 note_block_progress(slots: &mut [PeerSlot], peer: usize) { if let Some(s) = slots.iter_mut().find(|s| s.id == peer) { touch_block_progress(&s.block_progress_ms); @@ -455,6 +467,7 @@ pub(crate) async fn spawn_peer( first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), bytes_rx_total, + rate: PeerRate::default(), alive: true, task, }) @@ -537,6 +550,7 @@ mod tests { first_data_ms: AtomicU64::new(0), bytes_rx: AtomicU64::new(0), bytes_rx_total: Arc::new(AtomicU64::new(0)), + rate: Default::default(), alive: true, task, } @@ -593,6 +607,11 @@ mod tests { 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()); + touch_block_progress(&s.block_progress_ms); assert!(s.block_progress_ms.load(Ordering::Relaxed) > 0); note_block_progress(std::slice::from_mut(&mut s), 7); diff --git a/crates/rbitcoin-net/src/ibd/rate.rs b/crates/rbitcoin-net/src/ibd/rate.rs index 8551579e..37e3935f 100644 --- a/crates/rbitcoin-net/src/ibd/rate.rs +++ b/crates/rbitcoin-net/src/ibd/rate.rs @@ -89,7 +89,6 @@ impl PeerRate { let last = self.progress_ms.max(self.work_started_ms); now_ms.saturating_sub(last) > stall_ms } - } #[cfg(test)] @@ -105,7 +104,6 @@ mod tests { assert!(r.active_ms < RELSLOW_ACTIVE_MS); } - #[test] fn idle_rebaseline_freezes_ewma() { let mut r = PeerRate::default(); From 6379acd26734c2951f588568e1c3b3648e25470f Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Sun, 30 Aug 2026 07:33:39 -0700 Subject: [PATCH 4/9] ibd: stall on rx + work start, not getdata issue Issuing getdata only starts the 30s grace clock. Qualifying stream or block/notfound events count as progress; hung sockets disconnect. Co-authored-by: Cursor --- crates/rbitcoin-net/src/ibd/assign.rs | 21 +++++++- crates/rbitcoin-net/src/ibd/dial.rs | 75 +++++++++++++++++--------- crates/rbitcoin-net/src/ibd/peer_io.rs | 4 +- 3 files changed, 71 insertions(+), 29 deletions(-) diff --git a/crates/rbitcoin-net/src/ibd/assign.rs b/crates/rbitcoin-net/src/ibd/assign.rs index 6c30340c..3f79daaa 100644 --- a/crates/rbitcoin-net/src/ibd/assign.rs +++ b/crates/rbitcoin-net/src/ibd/assign.rs @@ -17,7 +17,7 @@ //! - One body-queue copy per height (receive path drops duplicates). use super::assign_plan::far_slots_per_peer; -use super::peer_io::{touch_block_progress, PeerCmd, PeerSlot}; +use super::peer_io::{ibd_mono_ms, PeerCmd, PeerSlot}; use super::state::{self, IbdWorkState}; use super::status::LoopStats; use super::{ @@ -447,7 +447,7 @@ pub(crate) fn issue_batch( st.slots[idx].in_flight.insert(h); } if empty { - touch_block_progress(&st.slots[idx].block_progress_ms); + st.slots[idx].rate.note_work_started(ibd_mono_ms()); } let _ = st.slots[idx].cmd_tx.send(PeerCmd::GetData { hashes: batch.clone(), @@ -836,6 +836,23 @@ mod tests { let _ = std::fs::remove_dir_all(dir); } + #[test] + fn issue_batch_does_not_count_as_rx() { + let (dir, _hub) = tmp_hub(); + let mut st = IbdWorkState::new(vec![dummy_slot(0)], None, Some(0)); + st.slots[0].rate.progress_ms = 42; + st.slots[0].rate.work_started_ms = 7; + let mut room = 10usize; + let mut issued = 0u64; + let t0 = super::super::peer_io::ibd_mono_ms(); + assert!(issue_one(&mut st, 0, h(30), &mut room, &mut issued)); + let t1 = super::super::peer_io::ibd_mono_ms(); + assert_eq!(st.slots[0].rate.progress_ms, 42); + assert!(st.slots[0].rate.work_started_ms >= t0); + assert!(st.slots[0].rate.work_started_ms <= t1); + let _ = std::fs::remove_dir_all(dir); + } + /// Off-path getdata (mainnet 08:16:23: ordered empty, h2h=0, inflight=7) /// must not occupy slots; tip+1 and live awaiting-reorg need stay. /// Speculative explore-need at an empty remainder is leftover — drop it. diff --git a/crates/rbitcoin-net/src/ibd/dial.rs b/crates/rbitcoin-net/src/ibd/dial.rs index 4731ff74..beeac881 100644 --- a/crates/rbitcoin-net/src/ibd/dial.rs +++ b/crates/rbitcoin-net/src/ibd/dial.rs @@ -434,14 +434,24 @@ pub(crate) fn disconnect_stalled_block_peers( addr_cooldown: &mut HashMap, now: Instant, stall: Duration, +) { + disconnect_stalled_block_peers_at(slots, inflight, addr_cooldown, now, stall, ibd_mono_ms()); +} + +pub(crate) fn disconnect_stalled_block_peers_at( + slots: &mut [PeerSlot], + inflight: &mut HashMap, + addr_cooldown: &mut HashMap, + now: Instant, + stall: Duration, + now_ms: u64, ) { let stall = stall.max(Duration::from_secs(30)); let stall_ms = stall.as_millis() as u64; - let now_ms = ibd_mono_ms(); let stalled_peers: Vec<(usize, usize, SocketAddr)> = slots .iter() .filter(|s| s.alive && !s.in_flight.is_empty()) - .filter(|s| now_ms.saturating_sub(s.block_progress_ms.load(Ordering::Relaxed)) > stall_ms) + .filter(|s| s.rate.stalled(now_ms, stall_ms, true)) .map(|s| (s.id, s.in_flight.len(), s.addr)) .collect(); for (id, n_work, addr) in stalled_peers { @@ -796,38 +806,54 @@ mod tests { } #[test] - fn disconnect_stalled_releases_and_cools_addr() { - use std::sync::atomic::Ordering as AtOrd; + fn disconnect_stalled_after_30s_without_rx() { let a = addr(11); let mut slot = dummy_slot(5, a, true); let h = BlockHash::from_byte_array([0xee; 32]); slot.in_flight.insert(h); - // Old progress → stalled relative to mono clock. - slot.block_progress_ms.store(0, AtOrd::Relaxed); - // Ensure mono has advanced past stall window. - while ibd_mono_ms() < 50 { - std::thread::sleep(Duration::from_millis(5)); - } + slot.rate.note_work_started(0); let mut inflight = HashMap::new(); inflight.insert(h, super::super::state::InflightReq::new(5)); let mut cooldown = HashMap::new(); - let now = Instant::now(); - disconnect_stalled_block_peers( - &mut [slot], + disconnect_stalled_block_peers_at( + std::slice::from_mut(&mut slot), &mut inflight, &mut cooldown, - now, - Duration::from_millis(1), // clamped to 30s internally + Instant::now(), + Duration::from_secs(30), + 30_001, + ); + assert!(cooldown.contains_key(&a)); + assert!(inflight.is_empty()); + } + + #[test] + fn disconnect_stalled_not_when_rx_recent() { + let a = addr(12); + let mut slot = dummy_slot(6, a, true); + let h = BlockHash::from_byte_array([0xee; 32]); + slot.in_flight.insert(h); + slot.rate.note_work_started(0); + slot.rate.note_rx(25_000); + let mut inflight = HashMap::new(); + inflight.insert(h, super::super::state::InflightReq::new(6)); + let mut cooldown = HashMap::new(); + disconnect_stalled_block_peers_at( + std::slice::from_mut(&mut slot), + &mut inflight, + &mut cooldown, + Instant::now(), + Duration::from_secs(30), + 45_000, ); - // With stall floor 30s, may not disconnect if mono elapsed < 30s. - // Drive with a very old progress and long stall requirement by faking: - // store progress far in the past relative to mono. - let mut slot2 = dummy_slot(6, addr(12), true); - slot2.in_flight.insert(h); - slot2.block_progress_ms.store(0, AtOrd::Relaxed); - // Advance mono if needed is limited — instead assert cooldown path when - // we force-release after a simulated stall disconnect (release already tested). - // Call with empty inflight peer (no work) → no-op. + assert!(cooldown.get(&a).is_none()); + assert!(inflight.contains_key(&h)); + } + + #[test] + fn disconnect_stalled_releases_and_cools_addr() { + let now = Instant::now(); + let mut cooldown = HashMap::new(); disconnect_stalled_block_peers( &mut [dummy_slot(7, addr(13), true)], &mut HashMap::new(), @@ -835,7 +861,6 @@ mod tests { now, Duration::from_secs(30), ); - // Alive peer with no in_flight is never stalled. assert!(cooldown.get(&addr(13)).is_none()); } diff --git a/crates/rbitcoin-net/src/ibd/peer_io.rs b/crates/rbitcoin-net/src/ibd/peer_io.rs index dcbffc56..e4ad43a5 100644 --- a/crates/rbitcoin-net/src/ibd/peer_io.rs +++ b/crates/rbitcoin-net/src/ibd/peer_io.rs @@ -176,13 +176,13 @@ pub(crate) fn sample_peer_rates(slots: &mut [PeerSlot], now_ms: u64) { pub(crate) fn note_block_progress(slots: &mut [PeerSlot], peer: usize) { if let Some(s) = slots.iter_mut().find(|s| s.id == peer) { - touch_block_progress(&s.block_progress_ms); + 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) { - touch_block_progress(&s.block_progress_ms); + s.rate.note_rx(ibd_mono_ms()); s.note_rx_bytes(wire_bytes as u64); } } From 0af7b60b8a971c4216170d07bc7e8d37f7b6d111 Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Sun, 30 Aug 2026 07:34:50 -0700 Subject: [PATCH 5/9] ibd: relative-slow maturity is EWMA active time Drop the 2MiB payload floor and 60s-since-first-byte gate. Pack warmup is still 60s since first block payload. Co-authored-by: Cursor --- crates/rbitcoin-net/src/ibd/dial.rs | 86 ++++++++++++----------------- 1 file changed, 35 insertions(+), 51 deletions(-) diff --git a/crates/rbitcoin-net/src/ibd/dial.rs b/crates/rbitcoin-net/src/ibd/dial.rs index beeac881..cdce4523 100644 --- a/crates/rbitcoin-net/src/ibd/dial.rs +++ b/crates/rbitcoin-net/src/ibd/dial.rs @@ -1,6 +1,7 @@ //! Peer dial, header request, stall disconnect / cooldown. use super::peer_io::{ibd_mono_ms, spawn_peer, PeerCmd, PeerEventSinks, PeerSlot}; +use super::rate::RELSLOW_ACTIVE_MS; use crate::chain::ChainHub; use crate::error::NetError; use crate::peers::{trying_connection_log, PeerConnType}; @@ -22,10 +23,6 @@ pub(crate) const STALL_ADDR_COOLDOWN: Duration = Duration::from_secs(10 * 60); pub(crate) const RELATIVE_SLOW_CLUSTER_SPREAD: u64 = 2; /// Disconnect only if peer bps ≤ median / this (default 2 → half median). pub(crate) const RELATIVE_SLOW_OUTLIER_RATIO: u64 = 2; -/// Per-peer mature sample age (ms since first block byte). -pub(crate) const RELATIVE_SLOW_MIN_AGE_MS: u64 = 60_000; -/// Per-peer mature payload floor for disconnect decisions (tip-rank may use less). -pub(crate) const RELATIVE_SLOW_MIN_BYTES: u64 = 2 * 1024 * 1024; /// Target mature peers before relative rule runs (full IBD peer set). pub(crate) const RELATIVE_SLOW_MIN_SAMPLES: usize = 8; /// Floor when fewer than 16 alive peers. @@ -131,26 +128,17 @@ pub(crate) fn relative_slow_with_hysteresis( } } -/// Build mature relative-slow samples from live slots (age + bytes floors). +/// Build mature relative-slow samples from live slots (`active_ms` floor). pub(crate) fn mature_relative_slow_samples(slots: &[PeerSlot]) -> Vec { - let now = ibd_mono_ms(); let mut out = Vec::new(); for s in slots { if !s.alive { continue; } - let first = s.first_data_ms.load(Ordering::Relaxed); - if first == 0 { - continue; - } - if now.saturating_sub(first) < RELATIVE_SLOW_MIN_AGE_MS { + if s.rate.active_ms < RELSLOW_ACTIVE_MS { continue; } - let bytes = s.bytes_rx.load(Ordering::Relaxed); - if bytes < RELATIVE_SLOW_MIN_BYTES { - continue; - } - let Some((_, bps)) = s.speed_sample() else { + let Some(bps) = s.rate.bps() else { continue; }; out.push(RelativeSlowSample { @@ -928,47 +916,21 @@ mod tests { ]; assert_eq!(relative_slow_pick(&tie, 8), Some(2)); - // Mature filters: dead / no first_data / young / under bytes / no sample. + // Mature filters: dead / young active_ms / no sample. let dead = { - let s = dummy_slot(0, addr(20), false); - s.first_data_ms - .store(1, std::sync::atomic::Ordering::Relaxed); - s.bytes_rx.store( - RELATIVE_SLOW_MIN_BYTES, - std::sync::atomic::Ordering::Relaxed, - ); - s - }; - let no_first = { - let s = dummy_slot(1, addr(21), true); - s.bytes_rx.store( - RELATIVE_SLOW_MIN_BYTES * 2, - std::sync::atomic::Ordering::Relaxed, - ); + let mut s = dummy_slot(0, addr(20), false); + s.rate.sample(0, 0, true); + s.rate.sample(RELSLOW_ACTIVE_MS, RELSLOW_ACTIVE_MS * 1_000, true); s }; let young = { - let s = dummy_slot(2, addr(22), true); - // first_data near "now" → age too small for RELATIVE_SLOW_MIN_AGE_MS. - s.first_data_ms - .store(ibd_mono_ms().max(1), std::sync::atomic::Ordering::Relaxed); - s.bytes_rx.store( - RELATIVE_SLOW_MIN_BYTES * 2, - std::sync::atomic::Ordering::Relaxed, - ); - s - }; - let thin = { - let s = dummy_slot(3, addr(23), true); - s.first_data_ms - .store(1, std::sync::atomic::Ordering::Relaxed); - s.bytes_rx.store(1024, std::sync::atomic::Ordering::Relaxed); + let mut s = dummy_slot(2, addr(22), true); + s.rate.sample(0, 0, true); + s.rate.sample(5_000, 50_000_000, true); s }; - let samples = mature_relative_slow_samples(&[dead, no_first, young, thin]); - // Fresh process: none mature (age/bytes). Long-lived llvm-cov may mature - // a slot with first=1 + enough mono — still never includes dead/no_first. - assert!(samples.iter().all(|s| s.peer_id != 0 && s.peer_id != 1)); + let samples = mature_relative_slow_samples(&[dead, young]); + assert!(samples.iter().all(|s| s.peer_id != 0 && s.peer_id != 2)); assert_eq!(global_first_block_ms(&[]), 0); let a_empty = dummy_slot(10, addr(30), true); @@ -1021,4 +983,26 @@ mod tests { assert!(suspect.is_none()); assert!(cooldown.is_empty()); } + + #[test] + fn mature_relative_slow_samples_uses_ewma_active_ms() { + let h = BlockHash::from_byte_array([9u8; 32]); + let mut young = dummy_slot(0, addr(50), true); + young.rate.sample(0, 0, true); + young.rate.sample(5_000, 50_000_000, true); + young.in_flight.insert(h); + assert!(young.rate.bps().is_some()); + assert!(young.rate.active_ms < RELSLOW_ACTIVE_MS); + + let mut mature = dummy_slot(1, addr(51), true); + mature.rate.sample(0, 0, true); + mature + .rate + .sample(RELSLOW_ACTIVE_MS, RELSLOW_ACTIVE_MS * 10_000, true); + mature.in_flight.insert(h); + + let samples = mature_relative_slow_samples(&[young, mature]); + assert!(samples.iter().all(|s| s.peer_id != 0)); + assert!(samples.iter().any(|s| s.peer_id == 1)); + } } From a413481bdc20406a44683ff14896d43bdd4c9df2 Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Sun, 30 Aug 2026 07:51:22 -0700 Subject: [PATCH 6/9] ibd: port tip-hole drop and densify policy onto PeerRate MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Drop one no-rx or relative-slow hole owner instead of clearing the whole race set at 6s. Steal hung densify only when a faster peer has a slot; issue new densify by EWMA rank; skip the band walk at cap; give 16 slots only to 2×-median outliers. Co-authored-by: Cursor --- crates/rbitcoin-net/src/ibd/assign.rs | 888 ++++++++++++++++++--- crates/rbitcoin-net/src/ibd/assign_plan.rs | 68 ++ crates/rbitcoin-net/src/ibd/dial.rs | 3 +- crates/rbitcoin-net/src/ibd/rate.rs | 6 + 4 files changed, 841 insertions(+), 124 deletions(-) diff --git a/crates/rbitcoin-net/src/ibd/assign.rs b/crates/rbitcoin-net/src/ibd/assign.rs index 3f79daaa..d799ceb8 100644 --- a/crates/rbitcoin-net/src/ibd/assign.rs +++ b/crates/rbitcoin-net/src/ibd/assign.rs @@ -16,7 +16,10 @@ //! - Never request beyond densify horizon; events refuse far bodies too. //! - One body-queue copy per height (receive path drops duplicates). -use super::assign_plan::far_slots_per_peer; +use super::assign_plan::densify_slots_for_peer; +use super::dial::{ + median_u64, relative_slow_pick, RelativeSlowSample, RELATIVE_SLOW_CLUSTER_SPREAD, +}; use super::peer_io::{ibd_mono_ms, PeerCmd, PeerSlot}; use super::state::{self, IbdWorkState}; use super::status::LoopStats; @@ -210,15 +213,32 @@ pub(crate) fn assign_work_ordered( return; } + let tip_hole = !tip_holes.is_empty(); + let (pack_median, pack_tight) = pack_ewma_bps(&st.slots, &alive); + let caps: HashMap = alive + .iter() + .map(|&pid| { + ( + pid, + densify_cap_for( + &st.slots, + pid, + cfg.per_peer, + tip_hole, + pack_median, + pack_tight, + ), + ) + }) + .collect(); + issued += steal_hung_densify(st, hub, &alive, tip_batch_hi, &caps); + let mut room = cfg.window.saturating_sub(st.inflight.len()); if room == 0 { finish_assign(loop_stats, t0, issued); return; } - // Leave per-peer headroom for tip races while hole>0. - let densify_per_peer = far_slots_per_peer(cfg.per_peer, !tip_holes.is_empty()); - let densify_hi = path_lo.saturating_add(CONTIG_DENSIFY_AHEAD); let depth_bytes = hub.query.block_queue_stats().1; let fetched_hi = hub @@ -241,6 +261,13 @@ pub(crate) fn assign_work_ordered( } st.assign_path_lo = path_lo; st.densify_scan_lo = st.densify_scan_lo.max(path_lo); + if !alive + .iter() + .any(|&pid| peer_has_slot(st, pid, caps.get(&pid).copied().unwrap_or(1))) + { + finish_assign(loop_stats, t0, issued); + return; + } let densify_lo = path_lo.max(st.densify_scan_lo); let densify = collect_height_band(st, hub, densify_lo, band_hi, room.max(1)); if densify.is_empty() { @@ -248,30 +275,24 @@ pub(crate) fn assign_work_ordered( return; } - let mut peer_i = st.assign_rot; - st.assign_rot = st.assign_rot.wrapping_add(1); + let ranked = rank_peers_by_speed(&st.slots, &alive, &HashSet::new()); let mut densify_q = densify; - while room > 0 && !densify_q.is_empty() { - let mut any = false; - for _ in 0..alive.len() { - if room == 0 || densify_q.is_empty() { + for &pid in &ranked { + if room == 0 || densify_q.is_empty() { + break; + } + let cap = caps.get(&pid).copied().unwrap_or(1); + while room > 0 && !densify_q.is_empty() { + if !peer_has_slot(st, pid, cap) { break; } - let pid = alive[peer_i % alive.len()]; - peer_i += 1; - if !peer_has_slot(st, pid, densify_per_peer) { - continue; - } let Some(h) = pop_need(&mut densify_q, st, hub) else { break; }; - if issue_one(st, pid, h, &mut room, &mut issued) { - any = true; + if !issue_one(st, pid, h, &mut room, &mut issued) { + break; } } - if !any { - break; - } } finish_assign(loop_stats, t0, issued); @@ -562,17 +583,210 @@ fn demote_zombie_pending_for_fetch( body.mark_missing(hash); } -/// Tip-hole getdata older than this with no claimable wire is cleared and re-issued -/// (mainnet freeze: inflight stuck, soft frozen, hole=1 forever). +/// Stream rx older than this is not “recent” for tip-hole owner eviction. +/// Matches the absolute stall floor so slow-but-steady 64 KiB ticks stay live. +const TIP_HOLE_RX_STALE: Duration = Duration::from_secs(30); + +fn peer_has_recent_rx(slot: &PeerSlot, now_ms: u64) -> bool { + slot.rate + .has_recent_rx(now_ms, TIP_HOLE_RX_STALE.as_millis() as u64) +} + +/// Which current owner of a tip-hole hash to drop from **this hash** (not disconnect). /// -/// Short on purpose: confirm claim waits ~5s per tick while tip is blocked; 20s -/// left the same slow race set holding hole=1 while densify progressed. -const TIP_HOLE_INFLIGHT_STALE: Duration = Duration::from_secs(6); +/// - No owner has recent rx → none (too early / first 64 KiB still in flight). +/// - Some have recent rx, some do not → drop a no-rx owner (quick dead-racer). +/// - All have recent rx → [`relative_slow_pick`] among those owners (`min_samples` = +/// owner count). Tight cluster → none. +/// - Solo owner: drop only when no recent rx and inflight age ≥ [`TIP_HOLE_RX_STALE`]. +pub(crate) fn tip_hole_owner_to_drop( + owners: &[usize], + slots: &[PeerSlot], + started_at: Instant, +) -> Option { + if owners.is_empty() { + return None; + } + let now_ms = ibd_mono_ms(); + let mut recent = Vec::new(); + let mut stale = Vec::new(); + for &id in owners { + let Some(slot) = slots.iter().find(|s| s.id == id && s.alive) else { + stale.push(id); + continue; + }; + if peer_has_recent_rx(slot, now_ms) { + recent.push(id); + } else { + stale.push(id); + } + } + if owners.len() == 1 { + if !recent.is_empty() { + return None; + } + if Instant::now().duration_since(started_at) >= TIP_HOLE_RX_STALE { + return Some(owners[0]); + } + return None; + } + if recent.is_empty() { + return None; + } + if let Some(&id) = stale.iter().min() { + return Some(id); + } + let samples: Vec = owners + .iter() + .filter_map(|&id| { + let s = slots.iter().find(|s| s.id == id && s.alive)?; + Some(RelativeSlowSample { + peer_id: id, + bps: s.rate.bps().unwrap_or(0), + has_inflight: true, + }) + }) + .collect(); + relative_slow_pick(&samples, samples.len()) +} + +fn drop_hash_owner(st: &mut IbdWorkState, hash: BlockHash, pid: usize) { + if let Some(s) = st.slots.iter_mut().find(|s| s.id == pid) { + s.in_flight.remove(&hash); + } + if let Some(req) = st.inflight.get_mut(&hash) { + if req.remove_peer(pid) { + st.inflight.remove(&hash); + } + } +} + +fn peer_bps(slots: &[PeerSlot], pid: usize) -> u64 { + slots + .iter() + .find(|s| s.id == pid && s.alive) + .and_then(|s| s.rate.bps()) + .unwrap_or(0) +} + +fn pack_ewma_bps(slots: &[PeerSlot], alive: &[usize]) -> (Option, bool) { + let mut samples: Vec = alive + .iter() + .filter_map(|&pid| { + slots + .iter() + .find(|s| s.id == pid && s.alive) + .and_then(|s| s.rate.bps()) + }) + .collect(); + if samples.is_empty() { + return (None, true); + } + samples.sort_unstable(); + let lo = samples[0]; + let hi = samples[samples.len() - 1]; + let tight = if lo == 0 { + hi == 0 + } else { + hi <= lo.saturating_mul(RELATIVE_SLOW_CLUSTER_SPREAD) + }; + (Some(median_u64(&samples)), tight) +} + +fn densify_cap_for( + slots: &[PeerSlot], + pid: usize, + per_peer: usize, + tip_hole: bool, + pack_median: Option, + pack_tight: bool, +) -> usize { + let bps = slots + .iter() + .find(|s| s.id == pid && s.alive) + .and_then(|s| s.rate.bps()); + densify_slots_for_peer(per_peer, tip_hole, bps, pack_median, pack_tight) +} + +/// Move hung single-peer densify getdata to a faster peer with a free slot. +fn steal_hung_densify( + st: &mut IbdWorkState, + hub: &ChainHub, + alive: &[usize], + tip_batch_hi: u32, + densify_caps: &HashMap, +) -> u64 { + let now = Instant::now(); + let now_ms = ibd_mono_ms(); + let candidates: Vec<(BlockHash, u32, usize, Instant)> = st + .inflight + .iter() + .filter_map(|(h, req)| { + if req.len() != 1 { + return None; + } + let &ht = st.hash_height.get(h)?; + if ht <= tip_batch_hi { + return None; + } + let &pid = req.peers.iter().next()?; + Some((*h, ht, pid, req.started_at)) + }) + .collect(); + let hung: Vec = candidates + .into_iter() + .filter(|(h, ht, pid, started)| { + if super::progress::claim_ready(hub, &mut st.body, *ht, h) { + return false; + } + let Some(slot) = st.slots.iter().find(|s| s.id == *pid && s.alive) else { + return false; + }; + if peer_has_recent_rx(slot, now_ms) { + return false; + } + now.duration_since(*started) >= TIP_HOLE_RX_STALE + }) + .map(|(h, _, _, _)| h) + .collect(); + let mut issued = 0u64; + for h in hung { + let Some(owner) = st + .inflight + .get(&h) + .and_then(|req| req.peers.iter().copied().next()) + else { + continue; + }; + let owner_bps = peer_bps(&st.slots, owner); + let mut faster: Vec = alive + .iter() + .copied() + .filter(|&pid| pid != owner && peer_bps(&st.slots, pid) > owner_bps) + .collect(); + if faster.is_empty() { + continue; + } + faster.sort_by(|&a, &b| peer_bps(&st.slots, b).cmp(&peer_bps(&st.slots, a))); + let dest = faster + .into_iter() + .find(|&pid| peer_has_slot(st, pid, densify_caps.get(&pid).copied().unwrap_or(1))); + let ht = st.hash_height.get(&h).copied(); + drop_hash_owner(st, h, owner); + if let Some(pid) = dest { + let mut room = 1usize; + let _ = issue_one(st, pid, h, &mut room, &mut issued); + } else if let Some(ht) = ht { + st.densify_scan_lo = st.densify_scan_lo.min(ht); + } + } + issued +} -/// Rank alive peer ids for tip-hole getdata: prefer peers not in `avoid`, then +/// Rank alive peer ids for getdata: prefer peers not in `avoid`, then /// higher live EWMA bps, then lower id. Unsampled peers sort last /// among non-avoided (bps=0). -pub(crate) fn rank_tip_hole_peers( +pub(crate) fn rank_peers_by_speed( slots: &[PeerSlot], alive: &[usize], avoid: &std::collections::HashSet, @@ -582,23 +796,26 @@ pub(crate) fn rank_tip_hole_peers( let avoided_a = avoid.contains(&a) as u8; let avoided_b = avoid.contains(&b) as u8; avoided_a.cmp(&avoided_b).then_with(|| { - let bps = |pid: usize| -> u64 { - slots - .iter() - .find(|s| s.id == pid && s.alive) - .and_then(|s| s.rate.bps()) - .unwrap_or(0) - }; + let bps = |pid: usize| -> u64 { peer_bps(slots, pid) }; bps(b).cmp(&bps(a)).then_with(|| a.cmp(&b)) }) }); ranked } +pub(crate) fn rank_tip_hole_peers( + slots: &[PeerSlot], + alive: &[usize], + avoid: &std::collections::HashSet, +) -> Vec { + rank_peers_by_speed(slots, alive, avoid) +} + /// Cover each tip-hole hash with multi-peer getdata, preferring faster peers. /// -/// Stale tip-batch inflight is cleared after [`TIP_HOLE_INFLIGHT_STALE`] and -/// re-raced, preferring peers that were **not** in the cleared set. +/// While the hole is open, at most one current owner of **this hash** is dropped +/// per call when a sibling is pulling or that owner is a relative-slow outlier +/// among owners. The whole race set is never cleared on request age. pub(crate) fn cover_tip_holes( st: &mut IbdWorkState, hub: &ChainHub, @@ -613,13 +830,6 @@ pub(crate) fn cover_tip_holes( let now = Instant::now(); for &h in holes { - // Skip only when **claim-ready** (confirmed or body-queue wire present). - // Class A alone and **zombie pending without BQ** are not enough — sole - // confirm intake is BQ wire. Resume seed marks Class A as known so densify - // does not re-walk the whole band; tip-hole race must still re-get. - // - // Exception: tip+1 held for incomplete reorg mid gather — re-get would - // re-BadPrev forever while mids starve. if st.reorg.is_awaiting_held_tip(&h) { continue; } @@ -635,10 +845,10 @@ pub(crate) fn cover_tip_holes( demote_zombie_pending_for_fetch(&mut st.body, hub, h, ht); let mut avoid: HashSet = HashSet::new(); if let Some(req) = st.inflight.get(&h) { - if now.duration_since(req.started_at) >= TIP_HOLE_INFLIGHT_STALE { - avoid = req.peers.clone(); - clear_hash_inflight(&mut st.slots, &mut st.inflight, h); - st.body.mark_missing(h); + let owners: Vec = req.peers.iter().copied().collect(); + if let Some(pid) = tip_hole_owner_to_drop(&owners, &st.slots, req.started_at) { + drop_hash_owner(st, h, pid); + avoid.insert(pid); } } let (already, second_at) = st @@ -652,11 +862,14 @@ pub(crate) fn cover_tip_holes( } let mut need = want - already; let mut placed_any = false; - let ranked = rank_tip_hole_peers(&st.slots, alive, &avoid); + let ranked = rank_peers_by_speed(&st.slots, alive, &avoid); for &pid in &ranked { if need == 0 { break; } + if avoid.contains(&pid) { + continue; + } let Some(idx) = st.slots.iter().position(|s| s.id == pid && s.alive) else { continue; }; @@ -769,6 +982,34 @@ mod tests { } } + fn plant_work_path(st: &mut IbdWorkState, lo: u32, hi: u32) { + for ht in lo..=hi { + let hash = h(ht); + st.record_height(hash, ht); + st.height_to_hash.insert(ht, hash); + st.ordered_set.insert(hash); + st.ordered.push_back(hash); + st.max_ordered_height = ht; + st.body.mark_missing(hash); + } + } + + fn seed_ewma(slot: &mut PeerSlot, bytes_per_sec: u64) { + slot.rate.sample(0, 0, true); + slot.rate + .sample(5_000, bytes_per_sec.saturating_mul(5), true); + } + + fn mark_tip_batch_ready(st: &mut IbdWorkState, hub: &ChainHub, path_lo: u32) { + use bitcoin::hashes::Hash as _; + for ht in path_lo..=32 { + hub.query + .block_queue_offer(ht, h(ht).to_byte_array(), 1, &[0u8; 80]) + .unwrap(); + st.body.mark_pending(h(ht)); + } + } + fn tmp_hub() -> (std::path::PathBuf, ChainHub) { let dir = std::env::temp_dir().join(format!( "rbitcoin-assign-{}-{}", @@ -942,6 +1183,41 @@ mod tests { assert_eq!(far_slots_per_peer(16, false), 8); } + #[test] + fn tip_hole_owner_to_drop_too_early_dead_racer_and_solo() { + use super::super::peer_io::ibd_mono_ms; + let mut slots = vec![dummy_slot(0), dummy_slot(1)]; + let started = Instant::now(); + assert_eq!( + tip_hole_owner_to_drop(&[0, 1], &slots, started), + None, + "no rx yet is too early" + ); + slots[1].rate.note_rx(ibd_mono_ms().max(1)); + assert_eq!( + tip_hole_owner_to_drop(&[0, 1], &slots, started), + Some(0), + "silent owner drops when sibling has rx" + ); + slots[0].rate.note_rx(ibd_mono_ms().max(1)); + assert_eq!( + tip_hole_owner_to_drop(&[0], &slots, Instant::now() - Duration::from_secs(7)), + None, + "solo with live rx is kept" + ); + assert_eq!( + tip_hole_owner_to_drop(&[0], &slots, Instant::now() - Duration::from_secs(31)), + None, + "solo with live rx kept even if started_at is old" + ); + slots[0].rate.progress_ms = 0; + assert_eq!( + tip_hole_owner_to_drop(&[0], &slots, Instant::now() - Duration::from_secs(31)), + Some(0), + "solo hung with no rx after 30s is replaced" + ); + } + /// Wrong first-wins body at tip+1 is not claim-ready; cover must dequeue and /// re-get the work-path hash (general hole=1 with bq soft growing ahead). #[test] @@ -1088,10 +1364,10 @@ mod tests { let _ = std::fs::remove_dir_all(dir); } - /// Stale tip-hole inflight (≥6s) with no claimable wire must clear and re-race. - /// Mainnet freeze: hole=1, inflight stuck forever, soft frozen, conf=0. + /// Request age does not clear the whole tip-hole race set. #[test] - fn cover_tip_holes_re_races_stale_inflight() { + fn cover_tip_holes_does_not_clear_whole_set_on_started_at() { + use super::super::peer_io::ibd_mono_ms; use super::super::state::InflightReq; let (dir, hub) = tmp_hub(); hub.ensure_genesis().unwrap(); @@ -1106,50 +1382,126 @@ mod tests { st.record_height(hole, ht); st.height_to_hash.insert(ht, hole); st.body.mark_missing(hole); - // Fresh inflight (<6s) must not re-race yet. - let mut fresh = InflightReq::new(0); - fresh.started_at = Instant::now() - Duration::from_secs(3); - st.inflight.insert(hole, fresh); + let mut frozen = InflightReq::new(0); + frozen.add_peer(1); + frozen.started_at = Instant::now() - Duration::from_secs(7); + st.inflight.insert(hole, frozen); st.slots[0].in_flight.insert(hole); - let holes = contiguous_tip_holes(&mut st, &hub, 8); - assert_eq!(holes, vec![hole]); + st.slots[1].in_flight.insert(hole); + let now = ibd_mono_ms().max(1); + st.slots[0].rate.note_rx(now); + st.slots[1].rate.note_rx(now); + seed_ewma(&mut st.slots[0], 1_000_000); + seed_ewma(&mut st.slots[1], 1_100_000); + seed_ewma(&mut st.slots[2], 1_050_000); let cfg = IbdConfig::for_test(); let alive: Vec = st.slots.iter().filter(|s| s.alive).map(|s| s.id).collect(); - let issued_fresh = cover_tip_holes(&mut st, &hub, &cfg, &alive, &holes); - // Still at want peers (race fills), but started_at not cleared as stale. - let age_fresh = Instant::now().duration_since(st.inflight[&hole].started_at); + let _ = cover_tip_holes(&mut st, &hub, &cfg, &alive, &[hole]); + let peers = &st.inflight[&hole].peers; assert!( - age_fresh >= Duration::from_secs(2), - "fresh inflight must keep original started_at; age={age_fresh:?} issued={issued_fresh}" + peers.contains(&0) && peers.contains(&1), + "tight pack with live rx must keep owners; peers={peers:?}" ); + let _ = std::fs::remove_dir_all(dir); + } - // Frozen inflight from a prior race that never delivered wire (≥6s). - st.inflight.clear(); - for s in st.slots.iter_mut() { - s.in_flight.clear(); - } - let mut frozen = InflightReq::new(0); - frozen.started_at = Instant::now() - Duration::from_secs(7); - st.inflight.insert(hole, frozen); + #[test] + fn cover_tip_holes_drops_owner_with_no_rx_when_sibling_progresses() { + use super::super::peer_io::ibd_mono_ms; + use super::super::state::InflightReq; + let (dir, hub) = tmp_hub(); + hub.ensure_genesis().unwrap(); + let mut st = IbdWorkState::new( + vec![dummy_slot(0), dummy_slot(1), dummy_slot(2)], + hub.tip_hash(), + hub.tip_height(), + ); + let hole = h(0x52); + let tip = hub.tip_height().unwrap_or(0); + let ht = tip.saturating_add(1); + st.record_height(hole, ht); + st.height_to_hash.insert(ht, hole); + st.body.mark_missing(hole); + let mut req = InflightReq::new(0); + req.add_peer(1); + st.inflight.insert(hole, req); st.slots[0].in_flight.insert(hole); + st.slots[1].in_flight.insert(hole); + st.slots[1].rate.note_rx(ibd_mono_ms().max(1)); + let cfg = IbdConfig::for_test(); + let alive: Vec = st.slots.iter().filter(|s| s.alive).map(|s| s.id).collect(); + let _ = cover_tip_holes(&mut st, &hub, &cfg, &alive, &[hole]); + let peers = &st.inflight[&hole].peers; assert!( - !super::super::progress::claim_ready(&hub, &mut st.body, ht, &hole), - "no wire → not claim-ready" + !peers.contains(&0), + "silent owner must leave this hash; peers={peers:?}" ); - let issued = cover_tip_holes(&mut st, &hub, &cfg, &alive, &holes); - assert!( - issued >= 1, - "stale inflight must re-race getdata; issued={issued}" + let _ = std::fs::remove_dir_all(dir); + } + + #[test] + fn cover_tip_holes_drops_relative_slow_owner_among_progressing() { + use super::super::peer_io::ibd_mono_ms; + use super::super::state::InflightReq; + let (dir, hub) = tmp_hub(); + hub.ensure_genesis().unwrap(); + let mut st = IbdWorkState::new( + vec![dummy_slot(0), dummy_slot(1), dummy_slot(2)], + hub.tip_hash(), + hub.tip_height(), ); + let hole = h(0x53); + let tip = hub.tip_height().unwrap_or(0); + let ht = tip.saturating_add(1); + st.record_height(hole, ht); + st.height_to_hash.insert(ht, hole); + st.body.mark_missing(hole); + let mut req = InflightReq::new(0); + req.add_peer(1); + req.add_peer(2); + st.inflight.insert(hole, req); + for i in 0..3 { + st.slots[i].in_flight.insert(hole); + st.slots[i].rate.note_rx(ibd_mono_ms().max(1)); + } + seed_ewma(&mut st.slots[0], 400_000); + seed_ewma(&mut st.slots[1], 2_000_000); + seed_ewma(&mut st.slots[2], 1_900_000); + let cfg = IbdConfig::for_test(); + let alive: Vec = st.slots.iter().filter(|s| s.alive).map(|s| s.id).collect(); + let _ = cover_tip_holes(&mut st, &hub, &cfg, &alive, &[hole]); + let peers = &st.inflight[&hole].peers; assert!( - st.inflight.contains_key(&hole), - "hash remains inflight after re-race" + !peers.contains(&0), + "half-median owner among progressing racers drops; peers={peers:?}" ); - // Fresh started_at (not still the 7s-old stamp). - let age = Instant::now().duration_since(st.inflight[&hole].started_at); + let _ = std::fs::remove_dir_all(dir); + } + + #[test] + fn cover_tip_holes_solo_slow_but_rx_live_kept() { + use super::super::peer_io::ibd_mono_ms; + use super::super::state::InflightReq; + let (dir, hub) = tmp_hub(); + hub.ensure_genesis().unwrap(); + let mut st = IbdWorkState::new(vec![dummy_slot(0)], hub.tip_hash(), hub.tip_height()); + let hole = h(0x54); + let tip = hub.tip_height().unwrap_or(0); + let ht = tip.saturating_add(1); + st.record_height(hole, ht); + st.height_to_hash.insert(ht, hole); + st.body.mark_missing(hole); + let mut req = InflightReq::new(0); + req.started_at = Instant::now() - Duration::from_secs(31); + st.inflight.insert(hole, req); + st.slots[0].in_flight.insert(hole); + st.slots[0].rate.note_rx(ibd_mono_ms().max(1)); + let cfg = IbdConfig::for_test(); + let alive: Vec = vec![0]; + let _ = cover_tip_holes(&mut st, &hub, &cfg, &alive, &[hole]); assert!( - age < Duration::from_secs(2), - "re-race must reset started_at; age={age:?}" + st.inflight[&hole].contains_peer(0), + "solo slow-but-steady download stays" ); let _ = std::fs::remove_dir_all(dir); } @@ -1200,42 +1552,316 @@ mod tests { let _ = std::fs::remove_dir_all(dir); } - /// After stale clear, re-race prefers peers not in the cleared set. #[test] - fn cover_tip_holes_rerace_avoids_prior_peers() { + fn densify_hung_owner_stolen_to_faster_peer() { use super::super::state::InflightReq; + use bitcoin::hashes::Hash as _; let (dir, hub) = tmp_hub(); hub.ensure_genesis().unwrap(); let mut st = IbdWorkState::new( - (0..5).map(dummy_slot).collect(), + vec![dummy_slot(0), dummy_slot(1)], hub.tip_hash(), hub.tip_height(), ); - let hole = h(0x71); - let tip = hub.tip_height().unwrap_or(0); - let ht = tip.saturating_add(1); - st.record_height(hole, ht); - st.height_to_hash.insert(ht, hole); - st.body.mark_missing(hole); + let stats = LoopStats::default(); + let mut cfg = IbdConfig::for_test(); + cfg.window = 64; + cfg.per_peer = 16; + let path_lo = hub.tip_height().unwrap_or(0).saturating_add(1); + plant_work_path(&mut st, path_lo, 40); + hub.query + .block_queue_offer(path_lo, h(path_lo).to_byte_array(), 1, &[0u8; 80]) + .unwrap(); + st.body.mark_pending(h(path_lo)); + let hung = h(40); + let mut req = InflightReq::new(0); + req.started_at = Instant::now() - Duration::from_secs(31); + st.inflight.insert(hung, req); + st.slots[0].in_flight.insert(hung); + seed_ewma(&mut st.slots[1], 1_000_000); + assign_work_ordered( + &mut st, + &hub, + &cfg, + &stats, + path_lo, + AssignDepth::Full, + None, + ); + let peers = &st.inflight[&hung].peers; + assert!( + peers.contains(&1) && !peers.contains(&0), + "hung densify must move to faster peer; peers={peers:?}" + ); + let _ = std::fs::remove_dir_all(dir); + } - let mut frozen = InflightReq::new(0); - frozen.add_peer(1); - frozen.started_at = Instant::now() - Duration::from_secs(7); - st.inflight.insert(hole, frozen); - st.slots[0].in_flight.insert(hole); - st.slots[1].in_flight.insert(hole); + #[test] + fn densify_slow_but_rx_live_not_stolen() { + use super::super::peer_io::ibd_mono_ms; + use super::super::state::InflightReq; + use bitcoin::hashes::Hash as _; + let (dir, hub) = tmp_hub(); + hub.ensure_genesis().unwrap(); + let mut st = IbdWorkState::new( + vec![dummy_slot(0), dummy_slot(1)], + hub.tip_hash(), + hub.tip_height(), + ); + let stats = LoopStats::default(); + let mut cfg = IbdConfig::for_test(); + cfg.window = 64; + cfg.per_peer = 16; + let path_lo = hub.tip_height().unwrap_or(0).saturating_add(1); + plant_work_path(&mut st, path_lo, 40); + hub.query + .block_queue_offer(path_lo, h(path_lo).to_byte_array(), 1, &[0u8; 80]) + .unwrap(); + st.body.mark_pending(h(path_lo)); + let hung = h(40); + let mut req = InflightReq::new(0); + req.started_at = Instant::now() - Duration::from_secs(31); + st.inflight.insert(hung, req); + st.slots[0].in_flight.insert(hung); + st.slots[0].rate.note_rx(ibd_mono_ms().max(1)); + seed_ewma(&mut st.slots[1], 1_000_000); + assign_work_ordered( + &mut st, + &hub, + &cfg, + &stats, + path_lo, + AssignDepth::Full, + None, + ); + assert!( + st.inflight[&hung].contains_peer(0), + "live rx must not be stolen; peers={:?}", + st.inflight[&hung].peers + ); + let _ = std::fs::remove_dir_all(dir); + } - let cfg = IbdConfig::for_test(); - let alive: Vec = st.slots.iter().filter(|s| s.alive).map(|s| s.id).collect(); - let holes = vec![hole]; - let issued = cover_tip_holes(&mut st, &hub, &cfg, &alive, &holes); - assert!(issued >= 1, "issued={issued}"); - let peers = &st.inflight[&hole].peers; - // Prefer 2,3,4 over reusing 0,1 first — at least one new peer in race. - let new = peers.iter().any(|&p| p >= 2); + #[test] + fn densify_hung_no_faster_peer_does_not_steal() { + use super::super::state::InflightReq; + use bitcoin::hashes::Hash as _; + let (dir, hub) = tmp_hub(); + hub.ensure_genesis().unwrap(); + let mut st = IbdWorkState::new(vec![dummy_slot(0)], hub.tip_hash(), hub.tip_height()); + let stats = LoopStats::default(); + let mut cfg = IbdConfig::for_test(); + cfg.window = 64; + cfg.per_peer = 16; + let path_lo = hub.tip_height().unwrap_or(0).saturating_add(1); + plant_work_path(&mut st, path_lo, 40); + hub.query + .block_queue_offer(path_lo, h(path_lo).to_byte_array(), 1, &[0u8; 80]) + .unwrap(); + st.body.mark_pending(h(path_lo)); + let hung = h(40); + let mut req = InflightReq::new(0); + req.started_at = Instant::now() - Duration::from_secs(31); + st.inflight.insert(hung, req); + st.slots[0].in_flight.insert(hung); + assign_work_ordered( + &mut st, + &hub, + &cfg, + &stats, + path_lo, + AssignDepth::Full, + None, + ); + assert!( + st.inflight[&hung].contains_peer(0), + "solo hung densify has no faster peer to steal to" + ); + let _ = std::fs::remove_dir_all(dir); + } + + #[test] + fn densify_hung_no_slot_rewinds_scan_lo() { + use super::super::peer_io::ibd_mono_ms; + use super::super::state::InflightReq; + use bitcoin::hashes::Hash as _; + let (dir, hub) = tmp_hub(); + hub.ensure_genesis().unwrap(); + let mut st = IbdWorkState::new( + vec![dummy_slot(0), dummy_slot(1)], + hub.tip_hash(), + hub.tip_height(), + ); + let stats = LoopStats::default(); + let mut cfg = IbdConfig::for_test(); + cfg.window = 1; + cfg.per_peer = 2; + let path_lo = hub.tip_height().unwrap_or(0).saturating_add(1); + plant_work_path(&mut st, path_lo, 41); + hub.query + .block_queue_offer(path_lo, h(path_lo).to_byte_array(), 1, &[0u8; 80]) + .unwrap(); + st.body.mark_pending(h(path_lo)); + let hung = h(40); + let other = h(41); + let mut req = InflightReq::new(0); + req.started_at = Instant::now() - Duration::from_secs(31); + st.inflight.insert(hung, req); + st.slots[0].in_flight.insert(hung); + st.inflight.insert(other, InflightReq::new(1)); + st.slots[1].in_flight.insert(other); + st.slots[1].rate.note_rx(ibd_mono_ms().max(1)); + seed_ewma(&mut st.slots[1], 1_000_000); + st.densify_scan_lo = 90; + assign_work_ordered( + &mut st, + &hub, + &cfg, + &stats, + path_lo, + AssignDepth::Full, + None, + ); + assert!( + !st.inflight.contains_key(&hung), + "hung hash cleared when faster peer has no slot" + ); + assert!( + st.densify_scan_lo <= 40, + "scan_lo must rewind to hung height; scan_lo={}", + st.densify_scan_lo + ); + let _ = std::fs::remove_dir_all(dir); + } + + #[test] + fn densify_issues_to_fastest_peer_first() { + let _env = lock_default_assign_stop(); + let (dir, hub) = tmp_hub(); + hub.ensure_genesis().unwrap(); + let mut st = IbdWorkState::new( + vec![dummy_slot(0), dummy_slot(1), dummy_slot(2)], + hub.tip_hash(), + hub.tip_height(), + ); + let stats = LoopStats::default(); + let mut cfg = IbdConfig::for_test(); + cfg.window = 128; + cfg.per_peer = 16; + let path_lo = hub.tip_height().unwrap_or(0).saturating_add(1); + plant_work_path(&mut st, path_lo, 40); + mark_tip_batch_ready(&mut st, &hub, path_lo); + seed_ewma(&mut st.slots[0], 100_000); + seed_ewma(&mut st.slots[1], 10_000_000); + seed_ewma(&mut st.slots[2], 1_000_000); + st.densify_scan_lo = 40; + assign_work_ordered( + &mut st, + &hub, + &cfg, + &stats, + path_lo, + AssignDepth::Full, + None, + ); + let want = h(40); assert!( - new, - "re-race should include peers outside cleared set; peers={peers:?}" + st.inflight.get(&want).is_some_and(|r| r.contains_peer(1)), + "first densify hash must go to fastest peer; inflight={:?}", + st.inflight.get(&want).map(|r| &r.peers) + ); + let _ = std::fs::remove_dir_all(dir); + } + + #[test] + fn densify_skips_band_walk_when_peers_at_cap() { + use super::super::state::InflightReq; + let (dir, hub) = tmp_hub(); + hub.ensure_genesis().unwrap(); + let mut st = IbdWorkState::new( + vec![dummy_slot(0), dummy_slot(1)], + hub.tip_hash(), + hub.tip_height(), + ); + let stats = LoopStats::default(); + let mut cfg = IbdConfig::for_test(); + cfg.window = 128; + cfg.per_peer = 16; + let path_lo = hub.tip_height().unwrap_or(0).saturating_add(1); + plant_work_path(&mut st, path_lo, 70); + mark_tip_batch_ready(&mut st, &hub, path_lo); + for i in 0..16u32 { + let ht = 40 + i; + let hash = h(ht); + let pid = (i % 2) as usize; + st.inflight.insert(hash, InflightReq::new(pid)); + st.slots[pid].in_flight.insert(hash); + } + st.densify_scan_lo = 40; + let before_keys: HashSet<_> = st.inflight.keys().copied().collect(); + assign_work_ordered( + &mut st, + &hub, + &cfg, + &stats, + path_lo, + AssignDepth::Full, + None, + ); + let after_keys: HashSet<_> = st.inflight.keys().copied().collect(); + assert_eq!(after_keys, before_keys, "no new densify when peers at cap"); + assert_eq!( + st.densify_scan_lo, 40, + "band walk must not advance scan_lo; scan_lo={}", + st.densify_scan_lo + ); + let _ = std::fs::remove_dir_all(dir); + } + + #[test] + fn densify_fast_peer_receives_more_than_eight() { + let _env = lock_default_assign_stop(); + let (dir, hub) = tmp_hub(); + hub.ensure_genesis().unwrap(); + let mut st = IbdWorkState::new( + vec![dummy_slot(0), dummy_slot(1), dummy_slot(2)], + hub.tip_hash(), + hub.tip_height(), + ); + let stats = LoopStats::default(); + let mut cfg = IbdConfig::for_test(); + cfg.window = 128; + cfg.per_peer = 16; + let path_lo = hub.tip_height().unwrap_or(0).saturating_add(1); + plant_work_path(&mut st, path_lo, 52); + mark_tip_batch_ready(&mut st, &hub, path_lo); + seed_ewma(&mut st.slots[0], 5_000_000); + seed_ewma(&mut st.slots[1], 15_000_000); + seed_ewma(&mut st.slots[2], 5_000_000); + st.densify_scan_lo = 33; + assign_work_ordered( + &mut st, + &hub, + &cfg, + &stats, + path_lo, + AssignDepth::Full, + None, + ); + assert_eq!( + st.slots[1].in_flight.len(), + 16, + "2×-median outlier must get full densify cap" + ); + assert!( + st.slots[0].in_flight.len() <= 8, + "non-outlier stays at half cap; n={}", + st.slots[0].in_flight.len() + ); + assert!( + st.slots[2].in_flight.len() <= 8, + "non-outlier stays at half cap; n={}", + st.slots[2].in_flight.len() ); let _ = std::fs::remove_dir_all(dir); } @@ -1460,6 +2086,7 @@ mod tests { #[test] fn densify_requests_beyond_legacy_2k_when_soft_allows() { + let _env = lock_default_assign_stop(); let (dir, hub) = tmp_hub(); hub.ensure_genesis().unwrap(); let mut st = IbdWorkState::new(vec![dummy_slot(0), dummy_slot(1)], None, Some(0)); @@ -1674,28 +2301,43 @@ mod tests { /// Serialize env mutators — parallel suite races `bq_assign_stop_bytes`. static BQ_ASSIGN_STOP_ENV_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(()); + struct AssignStopEnvRestore(Option, Option); + impl Drop for AssignStopEnvRestore { + fn drop(&mut self) { + match self.0.take() { + Some(v) => std::env::set_var("RBITCOIN_BLOCK_QUEUE_BYTES", v), + None => std::env::remove_var("RBITCOIN_BLOCK_QUEUE_BYTES"), + } + match self.1.take() { + Some(v) => std::env::set_var("RBITCOIN_BLOCK_QUEUE_GB", v), + None => std::env::remove_var("RBITCOIN_BLOCK_QUEUE_GB"), + } + } + } + + fn lock_default_assign_stop() -> (std::sync::MutexGuard<'static, ()>, AssignStopEnvRestore) { + let g = BQ_ASSIGN_STOP_ENV_LOCK + .lock() + .unwrap_or_else(|e| e.into_inner()); + let restore = AssignStopEnvRestore( + std::env::var_os("RBITCOIN_BLOCK_QUEUE_BYTES"), + std::env::var_os("RBITCOIN_BLOCK_QUEUE_GB"), + ); + std::env::remove_var("RBITCOIN_BLOCK_QUEUE_BYTES"); + std::env::remove_var("RBITCOIN_BLOCK_QUEUE_GB"); + (g, restore) + } + /// Over assign-stop: densify within confirm window ∩ fetched; not past window. #[test] fn densify_over_assign_stop_clamps_window_and_fetched() { let _g = BQ_ASSIGN_STOP_ENV_LOCK .lock() .unwrap_or_else(|e| e.into_inner()); - let prev_b = std::env::var_os("RBITCOIN_BLOCK_QUEUE_BYTES"); - let prev_g = std::env::var_os("RBITCOIN_BLOCK_QUEUE_GB"); - struct Restore(Option, Option); - impl Drop for Restore { - fn drop(&mut self) { - match self.0.take() { - Some(v) => std::env::set_var("RBITCOIN_BLOCK_QUEUE_BYTES", v), - None => std::env::remove_var("RBITCOIN_BLOCK_QUEUE_BYTES"), - } - match self.1.take() { - Some(v) => std::env::set_var("RBITCOIN_BLOCK_QUEUE_GB", v), - None => std::env::remove_var("RBITCOIN_BLOCK_QUEUE_GB"), - } - } - } - let _restore = Restore(prev_b, prev_g); + let _restore = AssignStopEnvRestore( + std::env::var_os("RBITCOIN_BLOCK_QUEUE_BYTES"), + std::env::var_os("RBITCOIN_BLOCK_QUEUE_GB"), + ); std::env::remove_var("RBITCOIN_BLOCK_QUEUE_GB"); std::env::set_var("RBITCOIN_BLOCK_QUEUE_BYTES", "2048"); diff --git a/crates/rbitcoin-net/src/ibd/assign_plan.rs b/crates/rbitcoin-net/src/ibd/assign_plan.rs index ccb572c4..060532a3 100644 --- a/crates/rbitcoin-net/src/ibd/assign_plan.rs +++ b/crates/rbitcoin-net/src/ibd/assign_plan.rs @@ -63,6 +63,33 @@ pub(crate) fn far_slots_per_peer(per_peer: usize, tip_hole: bool) -> usize { } } +/// Per-peer densify cap: drip 2 while a tip hole is open; otherwise half of +/// `per_peer`, or `per_peer` when this peer's EWMA bps is ≥ 2× pack median +/// and the pack is not a tight cluster. +pub(crate) fn densify_slots_for_peer( + per_peer: usize, + tip_hole: bool, + peer_bps: Option, + pack_median: Option, + pack_tight: bool, +) -> usize { + let base = far_slots_per_peer(per_peer, tip_hole); + if tip_hole || pack_tight { + return base; + } + let Some(bps) = peer_bps else { + return base; + }; + let Some(med) = pack_median else { + return base; + }; + if med > 0 && bps >= med.saturating_mul(2) { + per_peer + } else { + base + } +} + /// Whether to request more headers past the soft cap. /// /// Only when the ordered path is **mostly claim-ready** (dense body-queue / @@ -138,6 +165,47 @@ mod tests { assert!(!want_headers_beyond_soft_cap(0, 0, 0, 2048)); } + #[test] + fn densify_slots_tip_hole_is_two_for_all() { + assert_eq!( + densify_slots_for_peer(16, true, Some(10_000_000), Some(1_000_000), false), + 2 + ); + assert_eq!(densify_slots_for_peer(16, true, None, None, true), 2); + } + + #[test] + fn densify_slots_tight_pack_stays_half() { + assert_eq!( + densify_slots_for_peer(16, false, Some(1_500_000), Some(1_000_000), true), + 8 + ); + assert_eq!( + densify_slots_for_peer(16, false, Some(2_000_000), Some(1_000_000), true), + 8 + ); + } + + #[test] + fn densify_slots_fast_outlier_gets_full() { + assert_eq!( + densify_slots_for_peer(16, false, Some(2_000_000), Some(1_000_000), false), + 16 + ); + assert_eq!( + densify_slots_for_peer(16, false, Some(1_500_000), Some(1_000_000), false), + 8 + ); + assert_eq!( + densify_slots_for_peer(16, false, None, Some(1_000_000), false), + 8 + ); + assert_eq!( + densify_slots_for_peer(16, false, Some(2_000_000), None, false), + 8 + ); + } + /// 292k re-admit vs tip-chatter skip. #[test] fn enqueue_header_readmit_skips_inflight_pending_and_past_tip() { diff --git a/crates/rbitcoin-net/src/ibd/dial.rs b/crates/rbitcoin-net/src/ibd/dial.rs index cdce4523..8aa3dac6 100644 --- a/crates/rbitcoin-net/src/ibd/dial.rs +++ b/crates/rbitcoin-net/src/ibd/dial.rs @@ -920,7 +920,8 @@ mod tests { let dead = { let mut s = dummy_slot(0, addr(20), false); s.rate.sample(0, 0, true); - s.rate.sample(RELSLOW_ACTIVE_MS, RELSLOW_ACTIVE_MS * 1_000, true); + s.rate + .sample(RELSLOW_ACTIVE_MS, RELSLOW_ACTIVE_MS * 1_000, true); s }; let young = { diff --git a/crates/rbitcoin-net/src/ibd/rate.rs b/crates/rbitcoin-net/src/ibd/rate.rs index 37e3935f..732f8be8 100644 --- a/crates/rbitcoin-net/src/ibd/rate.rs +++ b/crates/rbitcoin-net/src/ibd/rate.rs @@ -89,6 +89,10 @@ impl PeerRate { let last = self.progress_ms.max(self.work_started_ms); now_ms.saturating_sub(last) > stall_ms } + + pub(crate) fn has_recent_rx(&self, now_ms: u64, stale_ms: u64) -> bool { + self.progress_ms != 0 && now_ms.saturating_sub(self.progress_ms) <= stale_ms + } } #[cfg(test)] @@ -167,5 +171,7 @@ mod tests { r.note_rx(20_000); assert!(!r.stalled(45_000, 30_000, true)); assert!(r.stalled(50_001, 30_000, true)); + assert!(r.has_recent_rx(45_000, 30_000)); + assert!(!r.has_recent_rx(50_001, 30_000)); } } From 00c347659f103b3dd9a9c5d2dc6490736fae1ef2 Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Sun, 30 Aug 2026 07:59:13 -0700 Subject: [PATCH 7/9] ibd: delete lifetime speed_sample and block_progress_ms PeerRate is the only rate and stall clock. first_data_ms stays a plain u64 for AddrMan latency and relative-slow pack warmup. Co-authored-by: Cursor --- crates/rbitcoin-net/src/ibd/assign.rs | 4 +- crates/rbitcoin-net/src/ibd/dial.rs | 24 ++---- .../src/ibd/events/confirm_reject_tests.rs | 28 ++----- crates/rbitcoin-net/src/ibd/events/mod.rs | 2 +- crates/rbitcoin-net/src/ibd/peer_io.rs | 83 ++++--------------- 5 files changed, 35 insertions(+), 106 deletions(-) diff --git a/crates/rbitcoin-net/src/ibd/assign.rs b/crates/rbitcoin-net/src/ibd/assign.rs index d799ceb8..c85fb825 100644 --- a/crates/rbitcoin-net/src/ibd/assign.rs +++ b/crates/rbitcoin-net/src/ibd/assign.rs @@ -970,11 +970,9 @@ mod tests { addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 18444 + id as u16), cmd_tx, in_flight: HashSet::new(), - block_progress_ms: Arc::new(AtomicU64::new(0)), peer_height: 100, connected_ms: 1, - first_data_ms: AtomicU64::new(0), - bytes_rx: AtomicU64::new(0), + first_data_ms: 0, bytes_rx_total: Arc::new(AtomicU64::new(0)), rate: Default::default(), alive: true, diff --git a/crates/rbitcoin-net/src/ibd/dial.rs b/crates/rbitcoin-net/src/ibd/dial.rs index 8aa3dac6..c3279730 100644 --- a/crates/rbitcoin-net/src/ibd/dial.rs +++ b/crates/rbitcoin-net/src/ibd/dial.rs @@ -157,7 +157,7 @@ pub(crate) fn global_first_block_ms(slots: &[PeerSlot]) -> u64 { if !s.alive { continue; } - let first = s.first_data_ms.load(Ordering::Relaxed); + let first = s.first_data_ms; if first == 0 { continue; } @@ -537,11 +537,9 @@ mod tests { addr: a, cmd_tx, in_flight: HashSet::new(), - block_progress_ms: Arc::new(AtomicU64::new(0)), peer_height: 0, connected_ms: 0, - first_data_ms: AtomicU64::new(0), - bytes_rx: AtomicU64::new(0), + first_data_ms: 0, bytes_rx_total: Arc::new(AtomicU64::new(0)), rate: Default::default(), alive, @@ -939,29 +937,25 @@ mod tests { let b_dead = { let mut s = dummy_slot(11, addr(31), true); s.alive = false; - s.first_data_ms - .store(5, std::sync::atomic::Ordering::Relaxed); + s.first_data_ms = 5; s }; assert_eq!(global_first_block_ms(std::slice::from_ref(&b_dead)), 0); let a = { - let s = dummy_slot(10, addr(30), true); - s.first_data_ms - .store(42, std::sync::atomic::Ordering::Relaxed); + let mut s = dummy_slot(10, addr(30), true); + s.first_data_ms = 42; s }; let c = { - let s = dummy_slot(12, addr(32), true); - s.first_data_ms - .store(10, std::sync::atomic::Ordering::Relaxed); + let mut s = dummy_slot(12, addr(32), true); + s.first_data_ms = 10; s }; assert_eq!(global_first_block_ms(&[a, c]), 10); // Warmup fails when age since global first < 60s (typical unit-test process). let a2 = { - let s = dummy_slot(10, addr(30), true); - s.first_data_ms - .store(10, std::sync::atomic::Ordering::Relaxed); + let mut s = dummy_slot(10, addr(30), true); + s.first_data_ms = 10; s }; let warm = relative_slow_global_warmup_ok(std::slice::from_ref(&a2)); 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 36dbf460..9fb19417 100644 --- a/crates/rbitcoin-net/src/ibd/events/confirm_reject_tests.rs +++ b/crates/rbitcoin-net/src/ibd/events/confirm_reject_tests.rs @@ -881,11 +881,9 @@ fn confirmed_height_mids_blocked_while_densify_ahead_leaves_tip_hole() { addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 18444), cmd_tx, in_flight: HashSet::new(), - block_progress_ms: Arc::new(AtomicU64::new(0)), peer_height: 100, connected_ms: 1, - first_data_ms: AtomicU64::new(0), - bytes_rx: AtomicU64::new(0), + first_data_ms: 0, bytes_rx_total: Arc::new(AtomicU64::new(0)), rate: Default::default(), alive: true, @@ -1080,11 +1078,9 @@ fn zombie_pending_mid_at_confirmed_height_never_reget() { addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 18445), cmd_tx, in_flight: HashSet::new(), - block_progress_ms: Arc::new(AtomicU64::new(0)), peer_height: 100, connected_ms: 1, - first_data_ms: AtomicU64::new(0), - bytes_rx: AtomicU64::new(0), + first_data_ms: 0, bytes_rx_total: Arc::new(AtomicU64::new(0)), rate: Default::default(), alive: true, @@ -1655,11 +1651,9 @@ fn apply_peer_event_body_and_control_surface() { addr: a, cmd_tx, in_flight: HashSet::new(), - block_progress_ms: Arc::new(AtomicU64::new(0)), peer_height: 10, connected_ms: 1, - first_data_ms: AtomicU64::new(0), - bytes_rx: AtomicU64::new(0), + first_data_ms: 0, bytes_rx_total: Arc::new(AtomicU64::new(0)), rate: Default::default(), alive: true, @@ -1908,11 +1902,9 @@ fn apply_peer_event_block_framed_bq_horizon_and_headers_done() { addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1)), 18444), cmd_tx, in_flight: HashSet::new(), - block_progress_ms: Arc::new(AtomicU64::new(0)), peer_height: 5, connected_ms: 1, - first_data_ms: AtomicU64::new(0), - bytes_rx: AtomicU64::new(0), + first_data_ms: 0, bytes_rx_total: Arc::new(AtomicU64::new(0)), rate: Default::default(), alive: true, @@ -2176,11 +2168,9 @@ fn block_framed_raw_offers_body_queue_with_confirm_feed() { addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1)), 18444), cmd_tx, in_flight: HashSet::new(), - block_progress_ms: Arc::new(AtomicU64::new(0)), peer_height: 5, connected_ms: 1, - first_data_ms: AtomicU64::new(0), - bytes_rx: AtomicU64::new(0), + first_data_ms: 0, bytes_rx_total: Arc::new(AtomicU64::new(0)), rate: Default::default(), alive: true, @@ -2325,11 +2315,9 @@ fn known_headers_re_admit_to_ordered_after_tip_drain() { addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1)), 18444), cmd_tx, in_flight: HashSet::new(), - block_progress_ms: Arc::new(AtomicU64::new(0)), peer_height: 5, connected_ms: 1, - first_data_ms: AtomicU64::new(0), - bytes_rx: AtomicU64::new(0), + first_data_ms: 0, bytes_rx_total: Arc::new(AtomicU64::new(0)), rate: Default::default(), alive: true, @@ -2487,11 +2475,9 @@ fn path_slot_first_wins_chained_via_headers() { addr: a, cmd_tx, in_flight: HashSet::new(), - block_progress_ms: Arc::new(AtomicU64::new(0)), peer_height: 10, connected_ms: 1, - first_data_ms: AtomicU64::new(0), - bytes_rx: AtomicU64::new(0), + first_data_ms: 0, bytes_rx_total: Arc::new(AtomicU64::new(0)), rate: Default::default(), alive: true, diff --git a/crates/rbitcoin-net/src/ibd/events/mod.rs b/crates/rbitcoin-net/src/ibd/events/mod.rs index 0d46bb21..d7923250 100644 --- a/crates/rbitcoin-net/src/ibd/events/mod.rs +++ b/crates/rbitcoin-net/src/ibd/events/mod.rs @@ -433,7 +433,7 @@ pub(crate) fn apply_peer_event( warn!("ibd: peer[{peer}] dead: {reason}"); if let Some(s) = st.slots.iter().find(|s| s.id == peer) { if let Some(bps) = s.rate.bps() { - let first = s.first_data_ms.load(Ordering::Relaxed); + let first = s.first_data_ms; let lat = first.saturating_sub(s.connected_ms); peer_book.note_speed(s.addr, lat, bps); } diff --git a/crates/rbitcoin-net/src/ibd/peer_io.rs b/crates/rbitcoin-net/src/ibd/peer_io.rs index e4ad43a5..c8e0b650 100644 --- a/crates/rbitcoin-net/src/ibd/peer_io.rs +++ b/crates/rbitcoin-net/src/ibd/peer_io.rs @@ -93,16 +93,12 @@ pub(crate) struct PeerSlot { pub cmd_tx: mpsc::UnboundedSender, /// Hashes currently requested from this peer. pub in_flight: HashSet, - /// Last block-download progress as [`ibd_mono_ms`]. - pub block_progress_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). - pub first_data_ms: AtomicU64, - /// Cumulative block payload bytes (speed sample). - pub bytes_rx: AtomicU64, + /// 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`. pub bytes_rx_total: Arc, pub rate: PeerRate, @@ -110,36 +106,6 @@ pub(crate) struct PeerSlot { pub task: JoinHandle<()>, } -impl PeerSlot { - /// Record received block payload bytes for FAST/SLOW classification. - pub fn note_rx_bytes(&self, n: u64) { - if n == 0 { - return; - } - let now = ibd_mono_ms(); - let _ = self - .first_data_ms - .compare_exchange(0, now, Ordering::Relaxed, Ordering::Relaxed); - self.bytes_rx.fetch_add(n, Ordering::Relaxed); - } - - /// `(latency_ms, bytes_per_sec)` once we have ≥64 KiB of block data. - pub fn speed_sample(&self) -> Option<(u64, u64)> { - let first = self.first_data_ms.load(Ordering::Relaxed); - if first == 0 { - return None; - } - let bytes = self.bytes_rx.load(Ordering::Relaxed); - if bytes < 64 * 1024 { - return None; - } - let latency_ms = first.saturating_sub(self.connected_ms); - let elapsed_ms = ibd_mono_ms().saturating_sub(first).max(1); - let bps = bytes.saturating_mul(1000) / elapsed_ms; - Some((latency_ms, bps)) - } -} - impl Drop for PeerSlot { fn drop(&mut self) { let _ = self.cmd_tx.send(PeerCmd::Shutdown); @@ -153,10 +119,6 @@ pub(crate) fn ibd_mono_ms() -> u64 { T0.get_or_init(Instant::now).elapsed().as_millis() as u64 } -pub(crate) fn touch_block_progress(ms: &AtomicU64) { - ms.store(ibd_mono_ms(), Ordering::Relaxed); -} - pub(crate) fn note_stream_bytes(counter: &AtomicU64, n: u64) { if n == 0 { return; @@ -182,8 +144,11 @@ pub(crate) fn note_block_progress(slots: &mut [PeerSlot], peer: usize) { 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) { - s.rate.note_rx(ibd_mono_ms()); - s.note_rx_bytes(wire_bytes as u64); + let now = ibd_mono_ms(); + s.rate.note_rx(now); + if wire_bytes > 0 && s.first_data_ms == 0 { + s.first_data_ms = now; + } } } @@ -216,7 +181,6 @@ pub(crate) async fn spawn_peer( // Reader → writer for pongs (must not write on the read task — that would // stall the receive half and look like a peer stall). let (out_tx, mut out_rx) = mpsc::unbounded_channel::(); - let block_progress_ms = Arc::new(AtomicU64::new(ibd_mono_ms())); let bytes_rx_total = Arc::new(AtomicU64::new(0)); let bytes_io = Arc::clone(&bytes_rx_total); @@ -461,11 +425,9 @@ pub(crate) async fn spawn_peer( addr, cmd_tx, in_flight: HashSet::new(), - block_progress_ms, peer_height, connected_ms: ibd_mono_ms(), - first_data_ms: AtomicU64::new(0), - bytes_rx: AtomicU64::new(0), + first_data_ms: 0, bytes_rx_total, rate: PeerRate::default(), alive: true, @@ -544,11 +506,9 @@ mod tests { addr: SocketAddr::new(IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1)), 18444), cmd_tx, in_flight: HashSet::new(), - block_progress_ms: Arc::new(AtomicU64::new(0)), peer_height: 100, connected_ms: 1, - first_data_ms: AtomicU64::new(0), - bytes_rx: AtomicU64::new(0), + first_data_ms: 0, bytes_rx_total: Arc::new(AtomicU64::new(0)), rate: Default::default(), alive: true, @@ -586,21 +546,8 @@ mod tests { } #[test] - fn speed_sample_and_progress_helpers() { + fn stream_bytes_sample_and_first_data() { let mut s = dummy_slot(7); - assert!(s.speed_sample().is_none()); - s.note_rx_bytes(0); // no-op - // first_data_ms==0 is treated as "no sample yet"; wait so mono ms > 0. - while ibd_mono_ms() == 0 { - std::thread::sleep(std::time::Duration::from_millis(1)); - } - s.note_rx_bytes(32 * 1024); - assert!(s.speed_sample().is_none()); // need ≥64 KiB - s.note_rx_bytes(64 * 1024); - assert!(s.first_data_ms.load(Ordering::Relaxed) > 0); - let sample = s.speed_sample().expect("bps after ≥64KiB"); - assert!(sample.1 > 0); - 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); @@ -612,11 +559,15 @@ mod tests { sample_peer_rates(std::slice::from_mut(&mut s), 5_000); assert!(s.rate.bps().is_some()); - touch_block_progress(&s.block_progress_ms); - assert!(s.block_progress_ms.load(Ordering::Relaxed) > 0); + 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); - note_block_progress(std::slice::from_mut(&mut s), 99); // missing peer + 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); } From 6f94bbe7594f1fa997c6364ee511387509b5ca25 Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Sun, 30 Aug 2026 07:59:56 -0700 Subject: [PATCH 8/9] docs: describe IBD peer-rate EWMA hygiene Operator stall/relative-slow text and ibd-memory densify slots now match the single inflight EWMA; PeerRate rustdoc names the consumers. Co-authored-by: Cursor --- OPERATOR.md | 16 +++++++++++----- crates/rbitcoin-net/src/ibd/rate.rs | 3 ++- docs/ibd-memory.md | 3 ++- 3 files changed, 15 insertions(+), 7 deletions(-) diff --git a/OPERATOR.md b/OPERATOR.md index 9283e56e..e1f06c7e 100644 --- a/OPERATOR.md +++ b/OPERATOR.md @@ -223,12 +223,18 @@ Default INFO is `ibd: progress` only. `--log-level debug` adds perf / sizes / pe **Tip hole / peer hygiene:** `hole=` on the progress line is the fetch gap from tip+1 to the next in-hand body (confirmed, still on the BQ, or already taken -onto loadq). Tip-batch getdata races up to 4 peers -(preferring faster live rates) and re-races after ~6s without wire. WARN -`ibd: peer[…] stalled` is absolute zero block progress (~30s). WARN +onto loadq). Peer speed is one EWMA of all received bytes while that peer has +block getdata in flight. Tip-batch getdata races up to 4 peers (preferring +higher EWMA). A hole owner with no qualifying rx is dropped from that hash +when a sibling is pulling; the whole race set is not cleared on getdata age. +Densify default is 8 in-flight hashes per peer (2 while a tip hole is open); +16 only for an EWMA outlier at ≥ 2× pack median. WARN +`ibd: peer[…] stalled` is 30s without qualifying rx (≥64 KiB stream or a +block / decode-fail / NotFound event) after work start. WARN `ibd: peer[…] relative-slow (bps= med= spread=…)` disconnects a clear -half-median outlier only after ~60s warm-up and only when the peer pack is not -tight (max/min bps > 2×); good-but-slightly-slower peers are kept. +half-median outlier only after ~60s pack warm-up and 30s of inflight EWMA, +and only when the pack is not tight (max/min bps > 2×); a uniformly slow +pipe is kept. **Create pins:** pipeline-local only (`batch_pin` / `BatchParents`). No process pin FIFO. Header plans via ConfirmParentCache. Just-confirmed **identity + full create outs** stay on in-flight until a later lookup wave snapshots drain+fence past the pack height and load finishes that wave's last in-flight read. Not a coins cache. diff --git a/crates/rbitcoin-net/src/ibd/rate.rs b/crates/rbitcoin-net/src/ibd/rate.rs index 732f8be8..a5a71afa 100644 --- a/crates/rbitcoin-net/src/ibd/rate.rs +++ b/crates/rbitcoin-net/src/ibd/rate.rs @@ -17,7 +17,8 @@ pub(crate) const PROGRESS_STEP: u64 = 64 * 1024; /// the byte cursor so it does not dilute the rate. `progress_ms` is last qualifying /// rx (stream ≥ [`PROGRESS_STEP`] or event-path `note_rx`); `work_started_ms` is the /// last empty→nonempty getdata. Stall is `now - max(progress, work_started) > stall` -/// while inflight. +/// while inflight. Ranking, relative-slow, densify caps, and AddrMan FAST/SLOW +/// all read [`Self::bps`]. #[derive(Clone, Copy, Debug, Default)] pub(crate) struct PeerRate { ewma: u64, diff --git a/docs/ibd-memory.md b/docs/ibd-memory.md index 024f3c92..825dc869 100644 --- a/docs/ibd-memory.md +++ b/docs/ibd-memory.md @@ -82,7 +82,8 @@ is over target. window, outstanding requests remain finite (per-peer in-flight window). Enqueueing those bodies cannot create a truly unbounded leak; the backlog drains as confirm dequeues. Bound queue size by **not requesting**, not by -**not reading**. +**not reading**. A tight slow pack keeps eight densify getdata per peer and +is allowed to finish them. Historical regression (do not reintroduce): bounded arch_job Full-drop and reader-side decode-permit wait before the next frame made peers look dead while From 1d3d742e350e14ccd76546d13dfe2978af769ce0 Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Sun, 30 Aug 2026 08:06:32 -0700 Subject: [PATCH 9/9] ibd: drop unused rank_tip_hole_peers alias Clippy denies dead_code on lib targets; ranking is rank_peers_by_speed for both tip-hole cover and densify issue. Co-authored-by: Cursor --- crates/rbitcoin-net/src/ibd/assign.rs | 10 +--------- 1 file changed, 1 insertion(+), 9 deletions(-) diff --git a/crates/rbitcoin-net/src/ibd/assign.rs b/crates/rbitcoin-net/src/ibd/assign.rs index c85fb825..26228a5b 100644 --- a/crates/rbitcoin-net/src/ibd/assign.rs +++ b/crates/rbitcoin-net/src/ibd/assign.rs @@ -803,14 +803,6 @@ pub(crate) fn rank_peers_by_speed( ranked } -pub(crate) fn rank_tip_hole_peers( - slots: &[PeerSlot], - alive: &[usize], - avoid: &std::collections::HashSet, -) -> Vec { - rank_peers_by_speed(slots, alive, avoid) -} - /// Cover each tip-hole hash with multi-peer getdata, preferring faster peers. /// /// While the hole is open, at most one current owner of **this hash** is dropped @@ -1534,7 +1526,7 @@ mod tests { let cfg = IbdConfig::for_test(); let alive: Vec = st.slots.iter().filter(|s| s.alive).map(|s| s.id).collect(); let avoid = HashSet::new(); - let ranked = rank_tip_hole_peers(&st.slots, &alive, &avoid); + let ranked = rank_peers_by_speed(&st.slots, &alive, &avoid); assert_eq!(ranked[0], 1, "fastest peer first: ranked={ranked:?}"); assert_eq!(ranked[1], 2, "medium second: ranked={ranked:?}"); assert_eq!(ranked[2], 0, "slow last: ranked={ranked:?}");