diff --git a/changelog.d/parent-same-peer-wtxid.md b/changelog.d/parent-same-peer-wtxid.md new file mode 100644 index 000000000..5591326c4 --- /dev/null +++ b/changelog.d/parent-same-peer-wtxid.md @@ -0,0 +1,6 @@ +Fixed + +- **A txid parent is requested after the same peer's wtxid window ends.** + Expiry matched the first announcement for that peer. When that row was + the waiting txid parent, the in-flight wtxid window stayed indexed and + no getdata followed. diff --git a/changelog.d/peer-resources.md b/changelog.d/peer-resources.md new file mode 100644 index 000000000..e23cdb120 --- /dev/null +++ b/changelog.d/peer-resources.md @@ -0,0 +1,19 @@ +Security + +- Cap the parent-request tracker per peer and process-wide. A full + process-wide table skips the new announcement and does not disconnect + the peer. Only a peer at its own cap is disconnected. A wtxid + announcement is re-requested as a wtxid and does not change another + peer's txid parent. While that request is in flight, the same hash is + not asked again as a txid. +- Charge outbound getdata and tx announcements against the per-peer send + budget, and stop serving blocks once that budget is already over. +- The per-peer rate window keeps the previous second so a boundary does + not grant a second full budget. +- Count v2 decoy packets and unknown message types in the per-peer rate + window on tip-follow and IBD. One decoy that does not fit adds the + rate-limit score. The peer is disconnected at the same threshold as + other frames. +- Batch mempool transaction announcements into one inv per thousand, + still charged against the per-peer send budget. +- Findings write-ups: 053, 054, 055, 056, 057, 058, 059. diff --git a/crates/rbitcoin-net/src/codec.rs b/crates/rbitcoin-net/src/codec.rs index aa5efce63..947ca8475 100644 --- a/crates/rbitcoin-net/src/codec.rs +++ b/crates/rbitcoin-net/src/codec.rs @@ -22,6 +22,10 @@ pub const MAX_PROTOCOL_MESSAGE_LENGTH: usize = 4_000_000; /// Bitcoin Core `MAX_INV_SZ` — max inventory items in inv/getdata/notfound. pub const MAX_INV_SIZE: usize = 50_000; +/// Tx announcements per `inv`. Separate from [`MAX_INV_SIZE`]: one message +/// stays at a thousand even though the wire cap is fifty thousand. +pub const TX_INV_BATCH: usize = 1_000; + /// Bitcoin Core `MAX_HEADERS_RESULTS` — max headers in a `headers` message. pub const MAX_HEADERS_RESULTS: usize = 2_000; diff --git a/crates/rbitcoin-net/src/ibd/peer_io.rs b/crates/rbitcoin-net/src/ibd/peer_io.rs index 2da0f1951..c2fd8dd34 100644 --- a/crates/rbitcoin-net/src/ibd/peer_io.rs +++ b/crates/rbitcoin-net/src/ibd/peer_io.rs @@ -5,11 +5,14 @@ //! - writer: encode offloaded for heavy payloads; then encrypt + write use super::rate::PeerRate; -use crate::codec::MAX_INV_SIZE; +use crate::codec::{FramedMessage, MAX_INV_SIZE}; use crate::error::NetError; use crate::msg_decode::spawn_decode_then_with_err; -use crate::peer::{connect_and_handshake_timed, HandshakePolicy, HANDSHAKE_TIMEOUT}; -use crate::v2::{read_v2_frame_with_progress, write_v2_msg_offload}; +use crate::peer::{ + connect_and_handshake_timed, HandshakePolicy, BAN_SCORE_THRESHOLD, HANDSHAKE_TIMEOUT, +}; +use crate::peer_dos::{PeerRateLimiter, RATE_LIMIT_BAN_SCORE}; +use crate::v2::{read_v2_frame_with_progress, write_v2_msg_offload, V2Reader, V2Writer}; use bitcoin::block::Header; use bitcoin::hashes::Hash; use bitcoin::p2p::message::NetworkMessage; @@ -167,6 +170,316 @@ pub(crate) fn note_block_rx(slots: &mut [PeerSlot], peer: usize, wire_bytes: usi } } +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")) + } +} + +fn relay_headers(id: usize, sinks: &PeerEventSinks, headers: Vec
) { + sinks.send_ctrl(PeerEvent::Headers { peer: id, headers }); +} + +fn relay_notfound(id: usize, sinks: &PeerEventSinks, inv: &[Inventory]) { + let hashes = block_inventory_hashes(inv); + if !hashes.is_empty() { + sinks.send_body(PeerEvent::NotFound { peer: id, hashes }); + } +} + +fn relay_addr(id: usize, sinks: &PeerEventSinks, list: &[(u32, bitcoin::p2p::address::Address)]) { + let addrs = net_addrs_from_addr(list); + if !addrs.is_empty() { + sinks.send_ctrl(PeerEvent::Addrs { peer: id, addrs }); + } +} + +fn relay_addrv2(id: usize, sinks: &PeerEventSinks, list: &[bitcoin::p2p::address::AddrV2Message]) { + let addrs = net_addrs_from_addrv2(list); + if !addrs.is_empty() { + sinks.send_ctrl(PeerEvent::Addrs { peer: id, addrs }); + } +} + +fn relay_blocks_inv(id: usize, sinks: &PeerEventSinks, inv: &[Inventory]) { + let hashes = block_inventory_hashes(inv); + if !hashes.is_empty() { + sinks.send_ctrl(PeerEvent::BlocksInv { peer: id, hashes }); + } +} + +fn apply_decoded_message( + id: usize, + sinks: &PeerEventSinks, + msg: bitcoin::p2p::message::RawNetworkMessage, +) { + match msg.into_payload() { + NetworkMessage::Headers(h) => relay_headers(id, sinks, h), + NetworkMessage::NotFound(inv) => relay_notfound(id, sinks, &inv), + 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(), + }); +} + +fn is_eof_or_reset(err: &std::io::Error) -> bool { + err.kind() == std::io::ErrorKind::UnexpectedEof + || err.kind() == std::io::ErrorKind::ConnectionReset +} + +/// Reader state for one IBD peer. The socket loop only pulls frames. +struct IbdReadCtx { + id: usize, + out_tx: mpsc::UnboundedSender, + sinks: PeerEventSinks, + /// 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, + logged_invalid_v2: bool, + ban_score: u32, +} + +impl IbdReadCtx { + /// `true` stops the reader. + fn handle(&mut self, frame: Result) -> bool { + match frame { + Ok(frame) => self.dispatch_frame(frame), + Err(err) => self.dispatch_err(err), + } + } + + fn dispatch_frame(&mut self, frame: FramedMessage) -> bool { + if frame.is_ping() { + self.on_ping(&frame); + return false; + } + if frame.is_block() { + self.on_block(frame); + return false; + } + self.on_heavy(frame); + false + } + + fn on_ping(&self, frame: &FramedMessage) { + if let Some(n) = frame.ping_nonce() { + let _ = self.out_tx.send(NetworkMessage::Pong(n)); + } + } + + fn on_block(&self, frame: FramedMessage) { + 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, + }); + } + // 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 + // event variant remains the reader's response if the check splits. + Some(hash) => { + self.sinks.send_body(PeerEvent::BlockDecodeFailed { + peer: self.id, + hash, + }); + } + None => { + rbitcoin_log::debug!( + "ibd: peer[{}] block frame without usable header hash", + self.id + ); + } + } + } + + /// Never await a decode permit on the reader (stalls TCP). + fn on_heavy(&self, frame: FramedMessage) { + let id = self.id; + let sinks_ok = self.sinks.clone(); + let sinks_err = self.sinks.clone(); + spawn_decode_then_with_err( + frame, + move |msg| apply_decoded_message(id, &sinks_ok, msg), + move |err| on_heavy_err(id, &sinks_err, err), + ); + } + + fn dispatch_err(&mut self, err: NetError) -> bool { + match err { + NetError::InvalidV2Type { contents_len } => self.on_invalid_v2(contents_len), + NetError::Io(err) if is_eof_or_reset(&err) => { + self.sinks.send_body(PeerEvent::Dead { + peer: self.id, + reason: format!("eof: {err}"), + }); + true + } + other => { + self.sinks.send_body(PeerEvent::Dead { + peer: self.id, + reason: other.to_string(), + }); + true + } + } + } + + fn on_invalid_v2(&mut self, contents_len: usize) -> bool { + if !self.logged_invalid_v2 { + 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(), + }); + true + } else { + false + } + } +} + +async fn read_ibd_peer( + id: usize, + magic: Magic, + mut reader: V2Reader, + out_tx: mpsc::UnboundedSender, + sinks_r: PeerEventSinks, + bytes_io: Arc, +) { + let mut prog_mark = 0usize; + let mut ctx = IbdReadCtx { + id, + out_tx, + sinks: sinks_r, + rate: PeerRateLimiter::default_limits(), + logged_invalid_v2: false, + ban_score: 0, + }; + loop { + let frame = read_v2_frame_with_progress( + &mut reader, + magic, + |buffered| note_read_progress(&bytes_io, &mut prog_mark, buffered), + |n| decoy_hook(&mut ctx.rate, &mut ctx.ban_score, n), + ) + .await; + // The progress mark is per read. The next frame starts at zero. + prog_mark = 0; + if ctx.handle(frame) { + break; + } + } +} + +async fn write_ibd_peer( + id: usize, + mut cmd_rx: mpsc::UnboundedReceiver, + mut out_rx: mpsc::UnboundedReceiver, + mut writer: V2Writer, + sinks_w: PeerEventSinks, +) { + loop { + tokio::select! { + cmd = cmd_rx.recv() => { + match cmd { + Some(PeerCmd::GetHeaders { locator }) => { + let locator = if locator.len() > crate::codec::MAX_LOCATOR_SZ { + locator[..crate::codec::MAX_LOCATOR_SZ].to_vec() + } else { + locator + }; + let gh = GetHeadersMessage::new( + locator, + BlockHash::from_byte_array([0u8; 32]), + ); + if write_v2_msg_offload( + &mut writer, + NetworkMessage::GetHeaders(gh), + ) + .await + .is_err() + { + sinks_w.send_body(PeerEvent::Dead { + peer: id, + reason: "write getheaders failed".into(), + }); + break; + } + } + Some(PeerCmd::GetData { hashes }) => { + for chunk in hashes.chunks(MAX_INV_SIZE) { + let inv: Vec<_> = chunk + .iter() + .copied() + .map(Inventory::WitnessBlock) + .collect(); + if inv.is_empty() { + continue; + } + if write_v2_msg_offload( + &mut writer, + NetworkMessage::GetData(inv), + ) + .await + .is_err() + { + sinks_w.send_body(PeerEvent::Dead { + peer: id, + reason: "write getdata failed".into(), + }); + return; + } + } + } + Some(PeerCmd::Shutdown) | None => break, + } + } + msg = out_rx.recv() => { + match msg { + Some(payload) => { + if write_v2_msg_offload(&mut writer, payload).await.is_err() { + sinks_w.send_body(PeerEvent::Dead { + peer: id, + reason: "write outbound failed".into(), + }); + break; + } + } + None => break, + } + } + } + } +} + pub(crate) async fn spawn_peer( id: usize, addr: crate::NetAddr, @@ -196,10 +509,10 @@ pub(crate) async fn spawn_peer( .await?; let peer_height = u32::try_from(ver.start_height).unwrap_or(0); - let (cmd_tx, mut cmd_rx) = mpsc::unbounded_channel::(); + let (cmd_tx, cmd_rx) = mpsc::unbounded_channel::(); // 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 (out_tx, out_rx) = mpsc::unbounded_channel::(); let bytes_rx_total = Arc::new(AtomicU64::new(0)); let bytes_io = Arc::clone(&bytes_rx_total); @@ -223,213 +536,14 @@ pub(crate) async fn spawn_peer( // July 18 cold-start worked with **no** post-handshake getaddr/sendaddrv2 // before getheaders; those writes raced Core's pipeline and peers closed // (ordered=0 / inflight=0 / never archive). - let mut reader = reader; let sinks_r = sinks.clone(); - let reader_task = tokio::spawn(async move { - let mut prog_mark = 0usize; - loop { - let frame = read_v2_frame_with_progress(&mut reader, magic, |buffered| { - let delta = buffered.saturating_sub(prog_mark); - note_stream_bytes(&bytes_io, delta as u64); - prog_mark = buffered; - }) - .await; - prog_mark = 0; - match frame { - Ok(frame) => { - if frame.is_ping() { - if let Some(n) = frame.ping_nonce() { - let _ = out_tx.send(NetworkMessage::Pong(n)); - } - continue; - } - - if frame.is_block() { - match frame.block_hash_from_header() { - Some(hash) if frame.payload.len() >= 80 => { - sinks_r.send_body(PeerEvent::BlockFramed { - peer: id, - hash, - payload: frame.payload, - }); - } - Some(hash) => { - sinks_r - .send_body(PeerEvent::BlockDecodeFailed { peer: id, hash }); - } - None => { - rbitcoin_log::debug!( - "ibd: peer[{id}] block frame without usable header hash" - ); - } - } - continue; - } - - 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. - spawn_decode_then_with_err( - frame, - move |msg| { - match msg.into_payload() { - NetworkMessage::Headers(h) => { - sinks_d.send_ctrl(PeerEvent::Headers { - peer: id, - headers: h, - }); - } - NetworkMessage::NotFound(inv) => { - let hashes = block_inventory_hashes(&inv); - if !hashes.is_empty() { - sinks_d.send_body(PeerEvent::NotFound { - peer: id, - hashes, - }); - } - } - NetworkMessage::Addr(list) => { - let addrs = net_addrs_from_addr(&list); - if !addrs.is_empty() { - sinks_d.send_ctrl(PeerEvent::Addrs { peer: id, addrs }); - } - } - NetworkMessage::AddrV2(list) => { - let addrs = net_addrs_from_addrv2(&list); - if !addrs.is_empty() { - sinks_d.send_ctrl(PeerEvent::Addrs { peer: id, addrs }); - } - } - NetworkMessage::SendAddrV2 => {} - // Blocks must not reach decode (handled above). - NetworkMessage::Block(_) => {} - NetworkMessage::Inv(inv) => { - let hashes = block_inventory_hashes(&inv); - if !hashes.is_empty() { - sinks_d.send_ctrl(PeerEvent::BlocksInv { - peer: id, - hashes, - }); - } - } - _other => {} - } - }, - { - let sinks_e = sinks_r.clone(); - move |e| { - sinks_e.send_body(PeerEvent::Dead { - peer: id, - reason: e.to_string(), - }); - } - }, - ); - } - Err(NetError::InvalidV2Type { .. }) => { - // Core logs and stays connected (same as tip-follow). - continue; - } - Err(NetError::Io(e)) - if e.kind() == std::io::ErrorKind::UnexpectedEof - || e.kind() == std::io::ErrorKind::ConnectionReset => - { - sinks_r.send_body(PeerEvent::Dead { - peer: id, - reason: format!("eof: {e}"), - }); - break; - } - Err(e) => { - sinks_r.send_body(PeerEvent::Dead { - peer: id, - reason: e.to_string(), - }); - break; - } - } - } - }); + let reader_task = tokio::spawn(read_ibd_peer(id, magic, reader, out_tx, sinks_r, bytes_io)); // Let the reader poll once before we accept write work (getheaders). tokio::task::yield_now().await; - let mut writer = writer; let sinks_w = sinks; - let writer_task = tokio::spawn(async move { - loop { - tokio::select! { - cmd = cmd_rx.recv() => { - match cmd { - Some(PeerCmd::GetHeaders { locator }) => { - let locator = if locator.len() > crate::codec::MAX_LOCATOR_SZ { - locator[..crate::codec::MAX_LOCATOR_SZ].to_vec() - } else { - locator - }; - let gh = GetHeadersMessage::new( - locator, - BlockHash::from_byte_array([0u8; 32]), - ); - if write_v2_msg_offload( - &mut writer, - NetworkMessage::GetHeaders(gh), - ) - .await - .is_err() - { - sinks_w.send_body(PeerEvent::Dead { - peer: id, - reason: "write getheaders failed".into(), - }); - break; - } - } - Some(PeerCmd::GetData { hashes }) => { - for chunk in hashes.chunks(MAX_INV_SIZE) { - let inv: Vec<_> = chunk - .iter() - .copied() - .map(Inventory::WitnessBlock) - .collect(); - if inv.is_empty() { - continue; - } - if write_v2_msg_offload( - &mut writer, - NetworkMessage::GetData(inv), - ) - .await - .is_err() - { - sinks_w.send_body(PeerEvent::Dead { - peer: id, - reason: "write getdata failed".into(), - }); - return; - } - } - } - Some(PeerCmd::Shutdown) | None => break, - } - } - msg = out_rx.recv() => { - match msg { - Some(payload) => { - if write_v2_msg_offload(&mut writer, payload).await.is_err() { - sinks_w.send_body(PeerEvent::Dead { - peer: id, - reason: "write outbound failed".into(), - }); - break; - } - } - None => break, - } - } - } - } - }); + let writer_task = tokio::spawn(write_ibd_peer(id, cmd_rx, out_rx, writer, sinks_w)); let mut guard = PeerIoTasks { reader: reader_task, @@ -552,6 +666,272 @@ mod tests { } } + fn framed(command: [u8; 12], payload: Vec) -> FramedMessage { + FramedMessage { + magic: Magic::from(bitcoin::Network::Regtest), + command, + payload, + } + } + + fn test_ctx() -> ( + IbdReadCtx, + mpsc::UnboundedReceiver, + mpsc::UnboundedReceiver, + mpsc::UnboundedReceiver, + ) { + let (out_tx, out_rx) = mpsc::unbounded_channel(); + let (body_tx, body_rx) = mpsc::unbounded_channel(); + let (ctrl_tx, ctrl_rx) = mpsc::unbounded_channel(); + let ctx = IbdReadCtx { + id: 3, + out_tx, + sinks: PeerEventSinks { + body: body_tx, + ctrl: ctrl_tx, + }, + rate: PeerRateLimiter::default_limits(), + logged_invalid_v2: false, + ban_score: 0, + }; + (ctx, out_rx, body_rx, ctrl_rx) + } + + fn one_header() -> Header { + use bitcoin::block::Version; + use bitcoin::CompactTarget; + Header { + version: Version::from_consensus(4), + prev_blockhash: BlockHash::from_byte_array([0u8; 32]), + merkle_root: bitcoin::TxMerkleNode::from_byte_array([0u8; 32]), + time: 1, + bits: CompactTarget::from_consensus(0x207f_ffff), + nonce: 1, + } + } + + #[test] + fn read_progress_counts_the_delta_from_the_mark() { + let bytes = AtomicU64::new(0); + let mut mark = 0usize; + note_read_progress(&bytes, &mut mark, 10); + note_read_progress(&bytes, &mut mark, 25); + assert_eq!(bytes.load(Ordering::Relaxed), 25); + assert_eq!(mark, 25); + note_read_progress(&bytes, &mut mark, 4); + assert_eq!(bytes.load(Ordering::Relaxed), 25); + assert_eq!(mark, 4); + } + + #[test] + fn second_decoy_past_the_window_is_the_threshold() { + let mut rate = PeerRateLimiter::default_limits(); + let mut score = 0u32; + let n = (crate::peer_dos::DEFAULT_MAX_BYTES_PER_SEC as usize) + 1; + assert!(decoy_hook(&mut rate, &mut score, n).is_ok()); + assert_eq!(score, RATE_LIMIT_BAN_SCORE); + let err = decoy_hook(&mut rate, &mut score, n).unwrap_err(); + assert!(matches!( + err, + NetError::Protocol("peer misbehavior threshold") + )); + } + + #[test] + fn ping_and_block_frames_stay_on_their_channels() { + let (mut ctx, mut out_rx, mut body_rx, mut ctrl_rx) = test_ctx(); + 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(), + )))); + 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])))); + assert!(out_rx.try_recv().is_err()); + + 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"); + assert!(!ctx.handle(Ok(block))); + match body_rx.try_recv() { + Ok(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"), + } + assert!(ctrl_rx.try_recv().is_err()); + 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 invalid_v2_and_io_errors_stop_on_the_shipped_reasons() { + let (mut ctx, _out, mut body_rx, _ctrl) = test_ctx(); + 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"), + _ => panic!("second invalid type must stop"), + } + + let (mut ctx, _out, mut body_rx, _ctrl) = test_ctx(); + 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}"), + _ => panic!("eof must be a dead peer"), + } + + let (mut ctx, _out, mut body_rx, _ctrl) = test_ctx(); + 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}"), + _ => panic!("reset must be a dead peer"), + } + + let (mut ctx, _out, mut body_rx, _ctrl) = test_ctx(); + assert!(ctx.handle(Err(NetError::Protocol("bye")))); + match body_rx.try_recv() { + Ok(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; + 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(); + 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("ctrl open"); + match ev { + PeerEvent::BlocksInv { peer, hashes } => { + assert_eq!(peer, 3); + assert_eq!(hashes, vec![hash]); + } + _ => panic!("block inv must reach ctrl"), + } + + 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 ev = tokio::time::timeout(std::time::Duration::from_secs(2), ctrl_rx.recv()) + .await + .expect("headers") + .expect("ctrl open"); + match ev { + PeerEvent::Headers { headers, .. } => assert_eq!(headers.len(), 1), + _ => panic!("headers must reach ctrl"), + } + + let mut too_many = Vec::new(); + bitcoin::consensus::encode::VarInt((crate::codec::MAX_HEADERS_RESULTS as u64) + 1) + .consensus_encode(&mut too_many) + .unwrap(); + assert!(!ctx.handle(Ok(framed(*b"headers\0\0\0\0\0", too_many)))); + let ev = tokio::time::timeout(std::time::Duration::from_secs(2), body_rx.recv()) + .await + .expect("oversize headers") + .expect("body open"); + match ev { + PeerEvent::Dead { reason, .. } => { + assert!(reason.contains("message too large"), "{reason}") + } + _ => panic!("oversize headers must die in the decoder"), + } + }); + } + + fn event_sinks() -> ( + PeerEventSinks, + mpsc::UnboundedReceiver, + mpsc::UnboundedReceiver, + ) { + let (body_tx, body_rx) = mpsc::unbounded_channel(); + let (ctrl_tx, ctrl_rx) = mpsc::unbounded_channel(); + ( + PeerEventSinks { + body: body_tx, + ctrl: ctrl_tx, + }, + body_rx, + ctrl_rx, + ) + } + + #[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()])); + assert!(matches!( + ctrl_rx.try_recv(), + Ok(PeerEvent::Headers { headers, .. }) if headers.len() == 1 + )); + apply(NetworkMessage::NotFound(vec![])); + assert!(body_rx.try_recv().is_err()); + let block = BlockHash::from_byte_array([1u8; 32]); + apply(NetworkMessage::NotFound(vec![Inventory::Block(block)])); + assert!(matches!( + body_rx.try_recv(), + Ok(PeerEvent::NotFound { hashes, .. }) if hashes == vec![block] + )); + apply(NetworkMessage::Inv(vec![Inventory::Block(block)])); + 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))])); + 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, + }])); + 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![])); + 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)); diff --git a/crates/rbitcoin-net/src/parent_req.rs b/crates/rbitcoin-net/src/parent_req.rs new file mode 100644 index 000000000..d144a404a --- /dev/null +++ b/crates/rbitcoin-net/src/parent_req.rs @@ -0,0 +1,562 @@ +//! Missing-parent GETDATA tracker. +//! +//! One mutex owns the hash map and a per-peer time index. A heartbeat walks +//! only that peer's due entries. Caps are announcements, not a new knob. + +use std::collections::{BTreeMap, HashMap, HashSet}; + +use super::{GETDATA_TX_INTERVAL_SECS, MAX_PARENTS_PER_PARK}; + +/// Announcements one peer may add. Past this the map does not grow. +pub(super) const MAX_PARENT_ANN_PER_PEER: usize = 5_000; +/// Process-wide announcements. A few hundred thousand, not tens of millions. +pub(super) const MAX_PARENT_ANN_GLOBAL: usize = 200_000; + +/// One parent GETDATA the sweep should send. +pub(crate) struct DueParent { + pub hash: [u8; 32], + pub wtxid: bool, +} + +/// Why `note_inv` did or did not record an announcement. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum ParentNote { + Accepted, + /// This peer's own counter is full. + PeerCapped, + /// The process-wide table is full. The announcement is skipped. + GlobalFull, +} + +struct ParentAnn { + peer: u64, + preferred: bool, + reqtime: u64, + /// Bucket in `due_by_peer` while this ann is waiting to be selected. + due_at: Option, + requested_until: Option, + failed: bool, + /// This announcement's getdata type. A wtxid inv does not retarget + /// another peer's txid row. + wtxid: bool, +} + +struct ParentSlot { + anns: Vec, +} + +/// Hash map plus per-peer due and in-flight indexes. The indexes are one +/// entry per announcement and stop at the same caps (extra RAM, not a scan +/// of every peer on the heartbeat). +pub(super) struct ParentTracker { + by_hash: HashMap<[u8; 32], ParentSlot>, + due_by_peer: HashMap>>, + inflight_by_peer: HashMap>>, + hashes_by_peer: HashMap>, + ann_per_peer: HashMap, + announcements: usize, +} + +impl ParentTracker { + pub(super) fn new() -> Self { + Self { + by_hash: HashMap::new(), + due_by_peer: HashMap::new(), + inflight_by_peer: HashMap::new(), + hashes_by_peer: HashMap::new(), + ann_per_peer: HashMap::new(), + announcements: 0, + } + } + + fn peer_at_cap(&self, peer: u64) -> bool { + self.ann_per_peer.get(&peer).copied().unwrap_or(0) >= MAX_PARENT_ANN_PER_PEER + } + + fn global_full(&self) -> bool { + self.announcements >= MAX_PARENT_ANN_GLOBAL + } + + pub(super) fn note_inv( + &mut self, + peer: u64, + hash: [u8; 32], + inbound: bool, + now: u64, + wtxid: bool, + ) -> ParentNote { + if self.rearm_existing(peer, hash, now, wtxid) { + return ParentNote::Accepted; + } + if self.peer_at_cap(peer) { + return ParentNote::PeerCapped; + } + if self.global_full() { + return ParentNote::GlobalFull; + } + let exp = now.saturating_add(GETDATA_TX_INTERVAL_SECS); + // Same-kind in-flight stays the one request. This announcer waits + // until that window ends and keeps its own getdata type. + let (due_at, requested_until) = match self.kind_inflight_until(&hash, wtxid) { + Some(until) => (Some(until), None), + None => (None, Some(exp)), + }; + self.add_ann( + hash, + ParentAnn { + peer, + preferred: !inbound, + reqtime: now, + due_at, + requested_until, + failed: false, + wtxid, + }, + ); + ParentNote::Accepted + } + + pub(super) fn schedule(&mut self, hash: [u8; 32], peer: u64, preferred: bool, reqtime: u64) { + if self + .by_hash + .get(&hash) + .is_some_and(|slot| slot.anns.iter().any(|a| a.peer == peer && !a.wtxid)) + { + return; + } + if self.peer_at_cap(peer) || self.global_full() { + return; + } + self.add_ann( + hash, + ParentAnn { + peer, + preferred, + reqtime, + due_at: Some(reqtime), + requested_until: None, + failed: false, + wtxid: false, + }, + ); + } + + /// True when this peer already had this kind of announcement for `hash`. + fn rearm_existing(&mut self, peer: u64, hash: [u8; 32], now: u64, wtxid: bool) -> bool { + let Some(slot) = self.by_hash.get_mut(&hash) else { + return false; + }; + let Some(pos) = slot + .anns + .iter() + .position(|a| a.peer == peer && a.wtxid == wtxid) + else { + return false; + }; + let due_at = { + if slot.anns[pos].requested_until.is_some() || slot.anns[pos].failed { + return true; + } + let exp = now.saturating_add(GETDATA_TX_INTERVAL_SECS); + let due_at = slot.anns[pos].due_at; + slot.anns[pos].requested_until = Some(exp); + slot.anns[pos].due_at = None; + (due_at, exp) + }; + if let Some(t) = due_at.0 { + unindex(&mut self.due_by_peer, peer, t, &hash); + } + index_at(&mut self.inflight_by_peer, peer, due_at.1, hash); + true + } + + fn add_ann(&mut self, hash: [u8; 32], ann: ParentAnn) { + let peer = ann.peer; + let due_at = ann.due_at; + let exp = ann.requested_until; + let slot = self + .by_hash + .entry(hash) + .or_insert_with(|| ParentSlot { anns: Vec::new() }); + slot.anns.push(ann); + *self.ann_per_peer.entry(peer).or_default() += 1; + self.announcements += 1; + self.hashes_by_peer.entry(peer).or_default().insert(hash); + if let Some(exp) = exp { + index_at(&mut self.inflight_by_peer, peer, exp, hash); + } else if let Some(t) = due_at { + index_at(&mut self.due_by_peer, peer, t, hash); + } + } + + pub(super) fn forget_peer(&mut self, peer: u64) { + self.due_by_peer.remove(&peer); + self.inflight_by_peer.remove(&peer); + let Some(hashes) = self.hashes_by_peer.remove(&peer) else { + return; + }; + for hash in hashes { + self.remove_peer_ann(&hash, peer); + } + } + + pub(super) fn announcer_peers(&self, hashes: [[u8; 32]; 2]) -> Vec { + let mut peers = Vec::new(); + for hash in hashes { + let Some(slot) = self.by_hash.get(&hash) else { + continue; + }; + for a in &slot.anns { + if !a.failed && !peers.contains(&a.peer) { + peers.push(a.peer); + } + } + } + peers + } + + pub(super) fn resolve(&mut self, hashes: [[u8; 32]; 2], admitted: bool) { + for hash in hashes { + if admitted { + self.remove_hash(&hash); + continue; + } + let inflight = self.inflight_peers(&hash); + for (peer, exp) in inflight { + self.fail_inflight(&hash, peer, exp); + } + if self.slot_all_failed(&hash) { + self.remove_hash(&hash); + } + } + } + + pub(super) fn take_due( + &mut self, + peer: u64, + now: u64, + mut already_have: impl FnMut(&[u8; 32], bool) -> bool, + ) -> Vec { + self.expire_peer_inflight(peer, now); + let due_hashes = self.due_hashes(peer, now); + let mut out = Vec::new(); + let mut seen = HashSet::new(); + for hash in due_hashes { + if out.len() >= MAX_PARENTS_PER_PARK { + break; + } + if !seen.insert(hash) { + continue; + } + self.expire_hash_inflight(&hash, now); + if !self.by_hash.contains_key(&hash) { + continue; + } + if self.slot_all_failed(&hash) { + self.remove_hash(&hash); + continue; + } + let Some(wtxid) = self.select(peer, now, &hash) else { + continue; + }; + if already_have(&hash, wtxid) { + self.drop_kind(&hash, wtxid); + continue; + } + out.push(DueParent { hash, wtxid }); + } + out + } + + fn due_hashes(&self, peer: u64, now: u64) -> Vec<[u8; 32]> { + let Some(tree) = self.due_by_peer.get(&peer) else { + return Vec::new(); + }; + tree.range(..=now) + .flat_map(|(_, v)| v.iter().copied()) + .collect() + } + + /// When this kind is already in flight, the time that request ends. + fn kind_inflight_until(&self, hash: &[u8; 32], wtxid: bool) -> Option { + self.by_hash.get(hash).and_then(|slot| { + slot.anns + .iter() + .filter(|a| a.wtxid == wtxid) + .filter_map(|a| a.requested_until) + .min() + }) + } + + fn select(&mut self, peer: u64, now: u64, hash: &[u8; 32]) -> Option { + let exp = now.saturating_add(GETDATA_TX_INTERVAL_SECS); + let (due_at, wtxid) = { + let slot = self.by_hash.get_mut(hash)?; + let has_pref = slot + .anns + .iter() + .any(|a| a.preferred && !a.failed && a.reqtime <= now); + let pos = slot.anns.iter().position(|a| { + !a.failed && a.reqtime <= now && a.peer == peer && (!has_pref || a.preferred) + })?; + let wtxid = slot.anns[pos].wtxid; + // Any in-flight request for these bytes blocks a second getdata. + // A non-segwit wtxid inv uses the txid, so the orphan parent is + // already in flight. The kind we would send stays this ann's own. + if slot.anns.iter().any(|a| a.requested_until.is_some()) { + return None; + } + let ann = &mut slot.anns[pos]; + let due_at = ann.due_at.unwrap_or(ann.reqtime); + ann.requested_until = Some(exp); + ann.due_at = None; + (due_at, wtxid) + }; + unindex(&mut self.due_by_peer, peer, due_at, hash); + index_at(&mut self.inflight_by_peer, peer, exp, *hash); + Some(wtxid) + } + + fn expire_peer_inflight(&mut self, peer: u64, now: u64) { + let expired = self + .inflight_by_peer + .get(&peer) + .map(|tree| { + tree.range(..=now) + .flat_map(|(t, v)| v.iter().map(|h| (*t, *h))) + .collect::>() + }) + .unwrap_or_default(); + for (exp, hash) in expired { + self.fail_inflight(&hash, peer, exp); + if self.slot_all_failed(&hash) { + self.remove_hash(&hash); + } + } + } + + fn expire_hash_inflight(&mut self, hash: &[u8; 32], now: u64) { + let expired = self.inflight_peers(hash); + for (peer, exp) in expired { + if exp <= now { + self.fail_inflight(hash, peer, exp); + } + } + if self.slot_all_failed(hash) { + self.remove_hash(hash); + } + } + + fn inflight_peers(&self, hash: &[u8; 32]) -> Vec<(u64, u64)> { + self.by_hash + .get(hash) + .map(|slot| { + slot.anns + .iter() + .filter_map(|a| a.requested_until.map(|exp| (a.peer, exp))) + .collect() + }) + .unwrap_or_default() + } + + fn fail_inflight(&mut self, hash: &[u8; 32], peer: u64, exp: u64) { + // Expire the ann whose window is `exp`. A same-peer txid parent is a + // different row and must stay eligible once this window ends. + if let Some(slot) = self.by_hash.get_mut(hash) { + if let Some(ann) = slot + .anns + .iter_mut() + .find(|a| a.peer == peer && a.requested_until == Some(exp)) + { + ann.requested_until = None; + ann.failed = true; + } + } + unindex(&mut self.inflight_by_peer, peer, exp, hash); + } + + fn slot_all_failed(&self, hash: &[u8; 32]) -> bool { + self.by_hash + .get(hash) + .is_some_and(|slot| !slot.anns.is_empty() && slot.anns.iter().all(|a| a.failed)) + } + + fn remove_hash(&mut self, hash: &[u8; 32]) { + let Some(slot) = self.by_hash.remove(hash) else { + return; + }; + for ann in slot.anns { + self.note_removed_ann(&ann, hash); + } + } + + fn drop_kind(&mut self, hash: &[u8; 32], wtxid: bool) { + let peers: Vec = self + .by_hash + .get(hash) + .map(|slot| { + slot.anns + .iter() + .filter(|a| a.wtxid == wtxid) + .map(|a| a.peer) + .collect() + }) + .unwrap_or_default(); + for peer in peers { + self.remove_peer_ann_kind(hash, peer, Some(wtxid)); + } + } + + fn remove_peer_ann(&mut self, hash: &[u8; 32], peer: u64) { + self.remove_peer_ann_kind(hash, peer, None); + } + + fn remove_peer_ann_kind(&mut self, hash: &[u8; 32], peer: u64, kind: Option) { + loop { + let Some(ann) = self.by_hash.get_mut(hash).and_then(|slot| { + let pos = slot + .anns + .iter() + .position(|a| a.peer == peer && kind.is_none_or(|k| a.wtxid == k))?; + Some(slot.anns.swap_remove(pos)) + }) else { + return; + }; + if self + .by_hash + .get(hash) + .is_some_and(|slot| slot.anns.is_empty()) + { + self.by_hash.remove(hash); + } + self.note_removed_ann(&ann, hash); + if kind.is_some() { + return; + } + } + } + + fn note_removed_ann(&mut self, ann: &ParentAnn, hash: &[u8; 32]) { + self.announcements = self.announcements.saturating_sub(1); + if let Some(c) = self.ann_per_peer.get_mut(&ann.peer) { + *c = c.saturating_sub(1); + if *c == 0 { + self.ann_per_peer.remove(&ann.peer); + } + } + if let Some(set) = self.hashes_by_peer.get_mut(&ann.peer) { + let still = self + .by_hash + .get(hash) + .is_some_and(|slot| slot.anns.iter().any(|a| a.peer == ann.peer)); + if !still { + set.remove(hash); + if set.is_empty() { + self.hashes_by_peer.remove(&ann.peer); + } + } + } + if let Some(exp) = ann.requested_until { + unindex(&mut self.inflight_by_peer, ann.peer, exp, hash); + } + if let Some(t) = ann.due_at { + unindex(&mut self.due_by_peer, ann.peer, t, hash); + } + } +} + +fn index_at( + tree: &mut HashMap>>, + peer: u64, + time: u64, + hash: [u8; 32], +) { + tree.entry(peer) + .or_default() + .entry(time) + .or_default() + .push(hash); +} + +fn unindex( + tree: &mut HashMap>>, + peer: u64, + time: u64, + hash: &[u8; 32], +) { + let Some(by_time) = tree.get_mut(&peer) else { + return; + }; + let Some(v) = by_time.get_mut(&time) else { + return; + }; + if let Some(i) = v.iter().position(|h| h == hash) { + v.swap_remove(i); + } + if v.is_empty() { + by_time.remove(&time); + } + if by_time.is_empty() { + tree.remove(&peer); + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn txid_parent_already_have_does_not_follow_another_peers_wtxid() { + let mut t = ParentTracker::new(); + let hash = [0xcd; 32]; + t.schedule(hash, 2, true, 1_000); + assert_eq!( + t.note_inv(1, hash, false, 1_000, true), + ParentNote::Accepted + ); + assert_eq!( + t.note_inv(3, hash, false, 1_000, true), + ParentNote::Accepted + ); + let now = 1_000 + GETDATA_TX_INTERVAL_SECS; + let mut saw_txid = false; + let dropped = t.take_due(2, now, |h, wtxid| { + assert_eq!(*h, hash); + assert!(!wtxid, "txid parent is checked as a txid"); + saw_txid = true; + true + }); + assert!(saw_txid); + assert!( + dropped.is_empty(), + "a txid parent already in the mempool is not requested" + ); + let follow = t.take_due(3, now, |_, wtxid| { + assert!(wtxid, "wtxid follow-up is not checked as a txid"); + false + }); + assert_eq!(follow.len(), 1); + assert!(follow[0].wtxid); + } + + #[test] + fn inflight_wtxid_of_the_same_hash_is_not_requested_again_as_txid() { + let mut t = ParentTracker::new(); + let inflight = [0x11; 32]; + let missing = [0x22; 32]; + assert_eq!( + t.note_inv(1, inflight, true, 1_000, true), + ParentNote::Accepted + ); + t.schedule(inflight, 2, false, 1_000); + t.schedule(missing, 2, false, 1_000); + let due = t.take_due(2, 1_000, |_, _| false); + assert_eq!( + due.len(), + 1, + "the in-flight hash must not join this getdata" + ); + assert_eq!(due[0].hash, missing); + assert!(!due[0].wtxid, "the other parent stays a txid request"); + } +} diff --git a/crates/rbitcoin-net/src/peer.rs b/crates/rbitcoin-net/src/peer.rs index 07d494ded..0858a1a5b 100644 --- a/crates/rbitcoin-net/src/peer.rs +++ b/crates/rbitcoin-net/src/peer.rs @@ -6,14 +6,16 @@ use crate::chain::{ received_getdata_wtx_log, received_tx_log, synchronizing_blockheaders_log, AcceptOutcome, ChainHub, }; -use crate::codec::{FramedMessage, MAX_HEADERS_RESULTS, MAX_INV_SIZE, MAX_LOCATOR_SZ}; +use crate::codec::{ + FramedMessage, MAX_HEADERS_RESULTS, MAX_INV_SIZE, MAX_LOCATOR_SZ, TX_INV_BATCH, +}; use crate::error::NetError; use crate::msg_decode::decode_framed_offload; use crate::peer_dos::{PeerRateLimiter, OVERSIZE_BAN_SCORE, RATE_LIMIT_BAN_SCORE}; use crate::peers::{CappedSet, PeerOut, PingAction}; use crate::v2::{ - open_v2, open_v2_with_wire, read_v2_contents, read_v2_frame, write_v2_contents, write_v2_msg, - write_v2_msg_offload, V2Reader, V2Writer, + open_v2, open_v2_with_wire, read_v2_contents, read_v2_frame, read_v2_frame_with_progress, + write_v2_contents, write_v2_msg, write_v2_msg_offload, V2Reader, V2Writer, }; use bitcoin::bip152::{BlockTransactions, BlockTransactionsRequest, HeaderAndShortIds}; use bitcoin::hashes::Hash; @@ -28,7 +30,7 @@ use rbitcoin_primitives::Height; use rbitcoin_query::Query; use std::collections::{BTreeSet, HashMap, HashSet, VecDeque}; use std::net::SocketAddr; -use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicU32, AtomicUsize, Ordering}; use std::sync::Arc; use std::time::{Duration, SystemTime, UNIX_EPOCH}; use tokio::net::TcpStream; @@ -1566,6 +1568,8 @@ pub async fn peer_session_with( } let mut requested_since: Option = None; let mut rate = PeerRateLimiter::default_limits(); + let decoy_score = AtomicU32::new(0); + let mut logged_invalid_v2 = false; let mut tx_announce_rx = hub.mempool().map(|m| m.subscribe_announces()); let mut inv_flush_rx = hub.mempool().map(|m| m.subscribe_inv_flush()); let mut headers_poll = tokio::time::interval(Duration::from_secs(HEADERS_POLL_SECS)); @@ -1587,6 +1591,7 @@ pub async fn peer_session_with( } } let hb_wait = SESSION_HEARTBEAT.saturating_sub(last_hb.elapsed()); + decoy_score.store(follow.ban_score, Ordering::Relaxed); tokio::select! { biased; // Peer half-close / write failure: tear down so getpeerinfo @@ -1597,7 +1602,22 @@ pub async fn peer_session_with( } // Inbound before local tip announce so GetData during a // generate burst is not queued behind hundreds of cmpctblocks. - frame = read_v2_frame(&mut reader, magic) => { + frame = read_v2_frame_with_progress(&mut reader, magic, |_| {}, |n| { + let mut score = decoy_score.load(Ordering::Relaxed); + let keep = crate::peer_dos::decoy_stays( + &mut rate, + &mut score, + n, + BAN_SCORE_THRESHOLD, + ); + decoy_score.store(score, Ordering::Relaxed); + if keep { + Ok(()) + } else { + Err(NetError::Protocol("peer misbehavior threshold")) + } + }) => { + follow.ban_score = decoy_score.load(Ordering::Relaxed); let frame = match frame { Ok(f) => f, // Any socket Io means the peer is gone — exit cleanly so @@ -1616,12 +1636,27 @@ pub async fn peer_session_with( } Err(NetError::InvalidV2Type { contents_len }) => { // Core stays connected; counts raw v2 size as `*other*`. + if !logged_invalid_v2 { + logged_invalid_v2 = true; + rbitcoin_log::debug!("{}", crate::v2::v2_invalid_message_type_log()); + } if let Some(ref sess) = session { sess.note_recv_raw( "*other*", crate::v2::v2_other_recv_bytes(contents_len), ); } + if !rate.note(contents_len) { + follow.ban_score = + follow.ban_score.saturating_add(RATE_LIMIT_BAN_SCORE); + rbitcoin_log::warn!( + "p2p: {peer_s} rate limit exceeded misbehavior={}", + follow.ban_score + ); + if follow.ban_score >= BAN_SCORE_THRESHOLD { + return Err(NetError::Protocol("peer misbehavior threshold")); + } + } continue; } Err(e) => return Err(e), @@ -1978,7 +2013,14 @@ fn maybe_expire_pending_cmpct( } for hash in stale { drop_pending_cmpct(follow, session, hash); - queue_block_getdata(hub, out, &mut follow.requested_blocks, &[hash], false)?; + queue_block_getdata( + hub, + out, + session, + &mut follow.requested_blocks, + &[hash], + false, + )?; } Ok(true) } @@ -1986,6 +2028,7 @@ fn maybe_expire_pending_cmpct( fn queue_block_getdata( hub: &ChainHub, out: &mpsc::UnboundedSender, + session: Option<&crate::peers::LivePeer>, requested_blocks: &mut HashSet, want: &[BlockHash], compact: bool, @@ -2008,7 +2051,7 @@ fn queue_block_getdata( hub.note_asked_block(*h); } for chunk in inv.chunks(MAX_INV_SIZE.min(500)) { - queue_out(out, NetworkMessage::GetData(chunk.to_vec()))?; + queue_accounted(session, out, NetworkMessage::GetData(chunk.to_vec()))?; } Ok(()) } @@ -2157,7 +2200,7 @@ pub fn force_announce_txid(hub: &ChainHub, peers: &crate::peers::PeerHub, txid: continue; }; s.note_announced_wtx(w); - let _ = queue_out(&out, NetworkMessage::Inv(vec![Inventory::WTx(w)])); + queue_tx_inv_batches(&s, &out, vec![Inventory::WTx(w)]); if let Some(seq) = mp.relay_seq_of(&w) { s.note_tx_inv_seq(s.last_inv_sequence().max(seq.saturating_add(1))); } @@ -2287,6 +2330,7 @@ fn queue_due_tx_invs( let Some(live_wtx) = mp.try_list_live_wtxids() else { return; }; + let mut due = Vec::new(); for (txid, w) in live_wtx { if !tx_inv_candidate_ok( mp, @@ -2300,12 +2344,13 @@ fn queue_due_tx_invs( continue; } session.note_announced_wtx(w); - let _ = queue_out(out_tx, NetworkMessage::Inv(vec![Inventory::WTx(w)])); + due.push(Inventory::WTx(w)); n += 1; if let Some(seq) = mp.relay_seq_of(&w) { max_ann = max_ann.max(seq.saturating_add(1)); } } + queue_tx_inv_batches(session, out_tx, due); if let Some((due, gen)) = mp.try_age_inv_watermark(mp_now) { session.note_age_inv_seen(due, gen); } @@ -2314,6 +2359,7 @@ fn queue_due_tx_invs( return; }; session.note_age_inv_seen(last.0, last.1); + let mut due = Vec::new(); for (txid, w) in due_wtx { if !tx_inv_candidate_ok( mp, @@ -2327,12 +2373,13 @@ fn queue_due_tx_invs( continue; } session.note_announced_wtx(w); - let _ = queue_out(out_tx, NetworkMessage::Inv(vec![Inventory::WTx(w)])); + due.push(Inventory::WTx(w)); n += 1; if let Some(seq) = mp.relay_seq_of(&w) { max_ann = max_ann.max(seq.saturating_add(1)); } } + queue_tx_inv_batches(session, out_tx, due); } if n > 0 { // Only INV txs that existed when this INV was built. @@ -2487,11 +2534,12 @@ fn on_compact_filters( payload: &NetworkMessage, hub: &ChainHub, out_tx: &mpsc::UnboundedSender, + session: Option<&crate::peers::LivePeer>, ) -> Result<(), NetError> { match payload { - NetworkMessage::GetCFilters(m) => on_getcfilters(hub, out_tx, m), - NetworkMessage::GetCFHeaders(m) => on_getcfheaders(hub, out_tx, m), - NetworkMessage::GetCFCheckpt(m) => on_getcfcheckpt(hub, out_tx, m), + NetworkMessage::GetCFilters(m) => on_getcfilters(hub, out_tx, session, m), + NetworkMessage::GetCFHeaders(m) => on_getcfheaders(hub, out_tx, session, m), + NetworkMessage::GetCFCheckpt(m) => on_getcfcheckpt(hub, out_tx, session, m), _ => Ok(()), } } @@ -2533,6 +2581,7 @@ fn within_filter_watermark(hub: &ChainHub, stop: u32) -> Result fn on_getcfilters( hub: &ChainHub, out_tx: &mpsc::UnboundedSender, + session: Option<&crate::peers::LivePeer>, m: &bitcoin::p2p::message_filter::GetCFilters, ) -> Result<(), NetError> { use bitcoin::hashes::Hash; @@ -2551,7 +2600,8 @@ fn on_getcfilters( else { return Ok(()); }; - queue_out( + queue_accounted( + session, out_tx, NetworkMessage::CFilter(bitcoin::p2p::message_filter::CFilter { filter_type: 0, @@ -2566,6 +2616,7 @@ fn on_getcfilters( fn on_getcfheaders( hub: &ChainHub, out_tx: &mpsc::UnboundedSender, + session: Option<&crate::peers::LivePeer>, m: &bitcoin::p2p::message_filter::GetCFHeaders, ) -> Result<(), NetError> { use bitcoin::bip158::FilterHeader; @@ -2587,7 +2638,8 @@ fn on_getcfheaders( } else { (rows[0].1, &rows[1..]) }; - queue_out( + queue_accounted( + session, out_tx, NetworkMessage::CFHeaders(bitcoin::p2p::message_filter::CFHeaders { filter_type: 0, @@ -2602,6 +2654,7 @@ fn on_getcfheaders( fn on_getcfcheckpt( hub: &ChainHub, out_tx: &mpsc::UnboundedSender, + session: Option<&crate::peers::LivePeer>, m: &bitcoin::p2p::message_filter::GetCFCheckpt, ) -> Result<(), NetError> { use bitcoin::hashes::Hash; @@ -2620,7 +2673,8 @@ fn on_getcfcheckpt( filter_headers.push(row[0].1); h = h.saturating_add(1000); } - queue_out( + queue_accounted( + session, out_tx, NetworkMessage::CFCheckpt(bitcoin::p2p::message_filter::CFCheckpt { filter_type: 0, @@ -2660,7 +2714,7 @@ fn handle_peer_inventory_msg( NetworkMessage::GetAddr => on_getaddr(hub, out_tx, session)?, NetworkMessage::GetCFilters(_) | NetworkMessage::GetCFHeaders(_) - | NetworkMessage::GetCFCheckpt(_) => on_compact_filters(payload, hub, out_tx)?, + | NetworkMessage::GetCFCheckpt(_) => on_compact_filters(payload, hub, out_tx, session)?, NetworkMessage::Unknown { .. } | NetworkMessage::GetData(_) | NetworkMessage::Block(_) @@ -2915,6 +2969,9 @@ async fn serve_getdata( let inflight = session.map(|s| &s.serve_inflight); let mut notfound: Vec = Vec::new(); for item in inv { + if session.is_some_and(|s| s.send_over_budget()) { + break; + } match item { // Unknown hash: Core ProcessGetData answers notfound. Silence // holds the requester's getdata until the stall floor. @@ -3129,6 +3186,7 @@ fn on_inv( } let mut want = Vec::new(); let mut inv_tx_n = 0u64; + let mut parent_capped = false; let mut need_headers = false; let mut tx_inv_hex: Option = None; let relay = !hub.in_ibd() @@ -3149,7 +3207,7 @@ fn on_inv( if tx_inv_hex.is_none() { tx_inv_hex = Some(txid.to_string()); } - if let Some(inv) = on_inv_txid(hub, session, relay, txid) { + if let Some(inv) = on_inv_txid(hub, session, relay, txid, &mut parent_capped) { want.push(inv); inv_tx_n = inv_tx_n.saturating_add(1); } @@ -3158,7 +3216,7 @@ fn on_inv( if tx_inv_hex.is_none() { tx_inv_hex = Some(wtxid.to_string()); } - if let Some(inv) = on_inv_wtxid(hub, session, relay, wtxid) { + if let Some(inv) = on_inv_wtxid(hub, session, relay, wtxid, &mut parent_capped) { want.push(inv); inv_tx_n = inv_tx_n.saturating_add(1); } @@ -3181,6 +3239,10 @@ fn on_inv( .count() as u64; mp.note_getdata_tx(gd_tx); } + if parent_capped { + punish_disconnect(&mut follow.ban_score, session); + return Ok(()); + } if let Some(hx) = tx_inv_hex { if reject_unsolicited_tx(hub, session) { rbitcoin_log::info!( @@ -3194,7 +3256,7 @@ fn on_inv( let _ = queue_getheaders(out_tx, hub, session, true, None); } if !want.is_empty() { - queue_out(out_tx, NetworkMessage::GetData(want))?; + queue_accounted(session, out_tx, NetworkMessage::GetData(want))?; } Ok(()) } @@ -3222,6 +3284,7 @@ fn on_inv_txid( session: Option<&crate::peers::LivePeer>, relay: bool, txid: &bitcoin::Txid, + parent_capped: &mut bool, ) -> Option { if !relay { return None; @@ -3235,7 +3298,15 @@ fn on_inv_txid( return None; } if let Some(s) = session { - mp.note_inv_tx_requested(s.id, txid.to_byte_array(), s.inbound, s.clock_now()); + match mp.note_inv_tx_requested(s.id, txid.to_byte_array(), s.inbound, s.clock_now(), false) + { + crate::tx_relay::ParentNote::Accepted => {} + crate::tx_relay::ParentNote::PeerCapped => { + *parent_capped = true; + return None; + } + crate::tx_relay::ParentNote::GlobalFull => return None, + } } Some(Inventory::WitnessTransaction(*txid)) } @@ -3245,6 +3316,7 @@ fn on_inv_wtxid( session: Option<&crate::peers::LivePeer>, relay: bool, wtxid: &bitcoin::Wtxid, + parent_capped: &mut bool, ) -> Option { if !relay { return None; @@ -3252,7 +3324,20 @@ fn on_inv_wtxid( let mp = hub.mempool()?; if !mp.try_contains_wtxid(wtxid) { if let Some(s) = session { - mp.note_inv_tx_requested(s.id, wtxid.to_byte_array(), s.inbound, s.clock_now()); + match mp.note_inv_tx_requested( + s.id, + wtxid.to_byte_array(), + s.inbound, + s.clock_now(), + true, + ) { + crate::tx_relay::ParentNote::Accepted => {} + crate::tx_relay::ParentNote::PeerCapped => { + *parent_capped = true; + return None; + } + crate::tx_relay::ParentNote::GlobalFull => return None, + } } return Some(Inventory::WTx(*wtxid)); } @@ -3346,6 +3431,7 @@ fn on_headers( queue_block_getdata( hub, out_tx, + session, &mut follow.requested_blocks, &want, getdata_use_compact(hub, follow.cmpct_version), @@ -3552,7 +3638,7 @@ async fn on_cmpctblock( } let keep = keep_pending_connecting_paths(hub, &follow.pending_headers); release_asks_off_path(hub, &mut follow.requested_blocks, &keep); - on_cmpctblock_queue_ancestors(hub, out_tx, follow, hash)?; + on_cmpctblock_queue_ancestors(hub, out_tx, follow, session, hash)?; if compact_header_low_work(hub, &hsi.header) && !follow.requested_blocks.contains(&hash) { let id = session.map(|s| s.id).unwrap_or(0); rbitcoin_log::info!("p2p: ignore low-work compact block from peer {id}"); @@ -3610,6 +3696,7 @@ fn on_cmpctblock_queue_ancestors( hub: &ChainHub, out_tx: &mpsc::UnboundedSender, follow: &mut PeerFollowState, + session: Option<&crate::peers::LivePeer>, hash: BlockHash, ) -> Result<(), NetError> { let mut ancestors: Vec = fetchable_header_path_bodies( @@ -3626,6 +3713,7 @@ fn on_cmpctblock_queue_ancestors( queue_block_getdata( hub, out_tx, + session, &mut follow.requested_blocks, &ancestors, getdata_use_compact(hub, follow.cmpct_version), @@ -3985,14 +4073,17 @@ fn queue_due_parent_getdata( return; } mp.note_getdata_tx(want.len() as u64); - let _ = queue_out( - out_tx, - NetworkMessage::GetData( - want.into_iter() - .map(Inventory::WitnessTransaction) - .collect(), - ), - ); + let inv = want + .into_iter() + .map(|due| { + if due.wtxid { + Inventory::WTx(bitcoin::Wtxid::from_byte_array(due.hash)) + } else { + Inventory::WitnessTransaction(bitcoin::Txid::from_byte_array(due.hash)) + } + }) + .collect(); + let _ = queue_accounted(Some(session), out_tx, NetworkMessage::GetData(inv)); } async fn on_tx( @@ -4722,6 +4813,20 @@ fn queue_out(out: &mpsc::UnboundedSender, msg: NetworkMessage) -> Resul queue_accounted(None, out, msg) } +/// Send filtered tx announcements in [`TX_INV_BATCH`] chunks. +/// +/// The caller already owns the filtered inventories. Chunking does not clone +/// the live wtxid set and does not hold the mempool lock across the send. +fn queue_tx_inv_batches( + session: &crate::peers::LivePeer, + out_tx: &mpsc::UnboundedSender, + inv: Vec, +) { + for chunk in inv.chunks(TX_INV_BATCH) { + let _ = queue_accounted(Some(session), out_tx, NetworkMessage::Inv(chunk.to_vec())); + } +} + fn queue_accounted( session: Option<&crate::peers::LivePeer>, out: &mpsc::UnboundedSender, @@ -4905,7 +5010,7 @@ async fn drain_pending( missing.retain(|h| !requested_blocks.contains(h)); let room = MAX_SERVE_BLOCKS.saturating_sub(requested_blocks.len()); missing.truncate(room); - queue_block_getdata(hub, out, requested_blocks, &missing, compact)?; + queue_block_getdata(hub, out, session, requested_blocks, &missing, compact)?; Ok(()) } diff --git a/crates/rbitcoin-net/src/peer_dos.rs b/crates/rbitcoin-net/src/peer_dos.rs index 8ccf9a91f..674fc3384 100644 --- a/crates/rbitcoin-net/src/peer_dos.rs +++ b/crates/rbitcoin-net/src/peer_dos.rs @@ -9,7 +9,7 @@ use tokio::sync::Semaphore; /// Default max concurrent **inbound** P2P sessions (post-handshake work). pub const DEFAULT_MAX_INBOUND: usize = 125; -/// Sliding window for rate accounting. +/// One-second rate window (current second plus the previous second). pub const RATE_WINDOW: Duration = Duration::from_secs(1); /// Max application messages per peer per window (after decrypt/frame). /// @@ -31,12 +31,19 @@ pub fn inbound_semaphore(max: usize) -> Arc { Arc::new(Semaphore::new(max.max(1))) } -/// Per-session sliding-window message and byte counters. +/// Per-session message and byte counters. +/// +/// Two buckets (the current second and the previous one). A note weighs the +/// previous bucket by how much of it is still inside the one-second window. +/// No per-message allocation. A tumbling reset granted a second full budget +/// at the boundary. #[derive(Debug, Clone)] pub struct PeerRateLimiter { - window_start: Instant, - msgs: u32, - bytes: u64, + current_start: Instant, + cur_msgs: u32, + cur_bytes: u64, + prev_msgs: u32, + prev_bytes: u64, max_msgs: u32, max_bytes: u64, } @@ -44,9 +51,11 @@ pub struct PeerRateLimiter { impl PeerRateLimiter { pub fn new(max_msgs: u32, max_bytes: u64) -> Self { Self { - window_start: Instant::now(), - msgs: 0, - bytes: 0, + current_start: Instant::now(), + cur_msgs: 0, + cur_bytes: 0, + prev_msgs: 0, + prev_bytes: 0, max_msgs: max_msgs.max(1), max_bytes: max_bytes.max(1), } @@ -59,21 +68,66 @@ impl PeerRateLimiter { /// Record one framed message of `payload_len` bytes. /// Returns `false` if this message would exceed the window budget. pub fn note(&mut self, payload_len: usize) -> bool { - let now = Instant::now(); - if now.duration_since(self.window_start) >= RATE_WINDOW { - self.window_start = now; - self.msgs = 0; - self.bytes = 0; - } - let next_msgs = self.msgs.saturating_add(1); - let next_bytes = self.bytes.saturating_add(payload_len as u64); - if next_msgs > self.max_msgs || next_bytes > self.max_bytes { + self.note_at(payload_len, Instant::now()) + } + + /// Same as [`Self::note`] at a chosen instant. Tests pin the window edge. + pub fn note_at(&mut self, payload_len: usize, now: Instant) -> bool { + self.roll(now); + let elapsed_ms = now + .saturating_duration_since(self.current_start) + .as_millis() + .min(1_000) as u64; + let prev_weight = 1_000 - elapsed_ms; + let eff_msgs = u64::from(self.cur_msgs) + (u64::from(self.prev_msgs) * prev_weight) / 1_000; + let eff_bytes = self.cur_bytes + (self.prev_bytes * prev_weight) / 1_000; + let next_msgs = eff_msgs.saturating_add(1); + let next_bytes = eff_bytes.saturating_add(payload_len as u64); + if next_msgs > u64::from(self.max_msgs) || next_bytes > self.max_bytes { return false; } - self.msgs = next_msgs; - self.bytes = next_bytes; + self.cur_msgs = self.cur_msgs.saturating_add(1); + self.cur_bytes = self.cur_bytes.saturating_add(payload_len as u64); true } + + fn roll(&mut self, now: Instant) { + let elapsed = now.saturating_duration_since(self.current_start); + if elapsed < RATE_WINDOW { + return; + } + if elapsed >= RATE_WINDOW + RATE_WINDOW { + self.prev_msgs = 0; + self.prev_bytes = 0; + self.cur_msgs = 0; + self.cur_bytes = 0; + self.current_start = now; + return; + } + self.prev_msgs = self.cur_msgs; + self.prev_bytes = self.cur_bytes; + self.cur_msgs = 0; + self.cur_bytes = 0; + self.current_start += RATE_WINDOW; + } +} + +/// Count one decoy (or other unsolicited frame) in `rate`. +/// +/// A frame that does not fit adds [`RATE_LIMIT_BAN_SCORE`] and stays +/// connected until `disconnect_at`. Callers use the same threshold as an +/// unknown message type. +pub fn decoy_stays( + rate: &mut PeerRateLimiter, + ban_score: &mut u32, + payload_len: usize, + disconnect_at: u32, +) -> bool { + if rate.note(payload_len) { + return true; + } + *ban_score = ban_score.saturating_add(RATE_LIMIT_BAN_SCORE); + *ban_score < disconnect_at } #[cfg(test)] @@ -96,15 +150,33 @@ mod tests { assert!(!r.note(1)); } + #[test] + fn rate_limiter_boundary_does_not_grant_a_second_budget() { + let mut r = PeerRateLimiter::new(2, 10_000); + let t0 = Instant::now(); + assert!(r.note_at(1, t0)); + assert!(r.note_at(1, t0)); + assert!(!r.note_at(1, t0)); + let boundary = t0 + RATE_WINDOW; + assert!( + !r.note_at(1, boundary), + "a full previous second must not grant another budget at the boundary" + ); + let cleared = t0 + RATE_WINDOW + RATE_WINDOW + Duration::from_millis(1); + assert!(r.note_at(1, cleared)); + assert!(r.note_at(1, cleared)); + assert!(!r.note_at(1, cleared)); + } + #[test] fn rate_limiter_window_resets() { let mut r = PeerRateLimiter::new(2, 10_000); - assert!(r.note(1)); - assert!(r.note(1)); - assert!(!r.note(1)); - // Simulate window expiry. - r.window_start = Instant::now() - RATE_WINDOW - Duration::from_millis(1); - assert!(r.note(1)); + let t0 = Instant::now(); + assert!(r.note_at(1, t0)); + assert!(r.note_at(1, t0)); + assert!(!r.note_at(1, t0)); + let cleared = t0 + RATE_WINDOW + RATE_WINDOW + Duration::from_millis(1); + assert!(r.note_at(1, cleared)); } #[test] @@ -116,4 +188,23 @@ mod tests { assert_eq!(RATE_LIMIT_BAN_SCORE, 50); assert_eq!(OVERSIZE_BAN_SCORE, 100); } + + #[test] + fn one_overflow_scores_and_the_second_disconnects() { + let mut rate = PeerRateLimiter::new(1, 10_000); + let mut score = 0u32; + let disconnect_at = RATE_LIMIT_BAN_SCORE.saturating_mul(2); + assert!(decoy_stays(&mut rate, &mut score, 1, disconnect_at)); + assert_eq!(score, 0, "a frame that fits does not score"); + assert!( + decoy_stays(&mut rate, &mut score, 1, disconnect_at), + "one frame over the window stays connected" + ); + assert_eq!(score, RATE_LIMIT_BAN_SCORE); + assert!( + !decoy_stays(&mut rate, &mut score, 1, disconnect_at), + "the next overflow reaches the disconnect threshold" + ); + assert_eq!(score, disconnect_at); + } } diff --git a/crates/rbitcoin-net/src/peer_tests.rs b/crates/rbitcoin-net/src/peer_tests.rs index cc096eae9..711d66b22 100644 --- a/crates/rbitcoin-net/src/peer_tests.rs +++ b/crates/rbitcoin-net/src/peer_tests.rs @@ -2958,6 +2958,270 @@ async fn over_budget_reader_waits_until_one_byte_is_written() { .unwrap(); } +fn inbound_peer( + peers: &std::sync::Arc, +) -> std::sync::Arc { + let addr = std::net::SocketAddr::from(([127, 0, 0, 1], 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: 1, + user_agent: "/rbitcoin:test/".into(), + start_height: 0, + relay: true, + }; + peers.register(addr, addr, &ver, true, crate::peers::PeerConnType::Inbound) +} + +#[tokio::test] +async fn inv_getdata_charges_send_budget() { + use bitcoin::hashes::Hash; + use bitcoin::Txid; + + let (dir, hub) = crate::chain::tiny_regtest_hub_labeled("inv-budget"); + hub.ensure_genesis().unwrap(); + let t = hub.tip_header().unwrap().time; + hub.clock.set_mock(i64::from(t) + 1); + assert!(!hub.in_ibd(), "tx inv getdata pin is not IBD"); + let mp = crate::tx_relay::MempoolHub::open(dir.join("mp"), std::sync::Arc::clone(&hub.query)) + .unwrap(); + mp.set_relay_enabled(true); + assert!(hub.attach_mempool(mp).is_ok()); + let peers = crate::peers::PeerHub::new(); + let peer = inbound_peer(&peers); + let (out_tx, _rx) = mpsc::unbounded_channel(); + let mut follow = PeerFollowState::new(); + let tx = Inventory::WitnessTransaction(Txid::from_byte_array([0x42; 32])); + on_inv(&hub, &out_tx, &mut follow, Some(&peer), &[tx]).unwrap(); + assert!( + peer.send_queued() > 0, + "inv getdata must charge the send budget" + ); + let _ = std::fs::remove_dir_all(dir); +} + +/// A full process-wide parent table skips a new announcement. The peer stays +/// connected, and getdata / getheaders already collected in that inv still go out. +#[tokio::test] +async fn full_parent_table_keeps_headers_and_collected_getdata() { + use bitcoin::hashes::Hash; + use bitcoin::Txid; + + let (dir, hub) = crate::chain::tiny_regtest_hub_labeled("parent-global"); + hub.ensure_genesis().unwrap(); + let t = hub.tip_header().unwrap().time; + hub.clock.set_mock(i64::from(t) + 1); + assert!(!hub.in_ibd(), "parent-cap inv pin is not IBD"); + let mp = crate::tx_relay::MempoolHub::open(dir.join("mp"), std::sync::Arc::clone(&hub.query)) + .unwrap(); + mp.set_relay_enabled(true); + assert!(hub.attach_mempool(mp).is_ok()); + let mp = hub.mempool().unwrap(); + let peers = crate::peers::PeerHub::new(); + let peer = inbound_peer(&peers); + let keep = [0x42; 32]; + assert_eq!( + mp.note_inv_tx_requested(peer.id, keep, true, 1_000, false), + crate::tx_relay::ParentNote::Accepted + ); + let mut n = 0u32; + let mut filler = peer.id.saturating_add(1); + let mut per_peer = 0u32; + loop { + let mut hash = [0u8; 32]; + hash[..4].copy_from_slice(&n.to_le_bytes()); + match mp.note_inv_tx_requested(filler, hash, false, 1_000, false) { + crate::tx_relay::ParentNote::Accepted => { + n += 1; + per_peer += 1; + if per_peer == 5_000 { + filler += 1; + per_peer = 0; + } + } + crate::tx_relay::ParentNote::GlobalFull => break, + crate::tx_relay::ParentNote::PeerCapped => { + panic!("filler peer {filler} hit its own cap before the table filled") + } + } + } + let (out_tx, mut rx) = mpsc::unbounded_channel(); + let mut follow = PeerFollowState::new(); + let block = Inventory::Block(BlockHash::from_byte_array([0x91; 32])); + let kept = Inventory::WitnessTransaction(Txid::from_byte_array(keep)); + let fresh = Inventory::WitnessTransaction(Txid::from_byte_array([0x77; 32])); + on_inv( + &hub, + &out_tx, + &mut follow, + Some(&peer), + &[block, kept, fresh], + ) + .unwrap(); + assert_eq!( + follow.ban_score, 0, + "a full process-wide table is not this peer's misbehavior" + ); + let mut saw_headers = false; + let mut getdata = Vec::new(); + while let Ok(msg) = rx.try_recv() { + match msg.expect_msg() { + NetworkMessage::GetHeaders(_) => saw_headers = true, + NetworkMessage::GetData(v) => getdata = v, + other => panic!("unexpected inv follow-up: {other:?}"), + } + } + assert!(saw_headers, "block inv still requests headers"); + assert_eq!( + getdata, + vec![kept], + "getdata already collected is sent, and the refused announcement is not" + ); + let _ = std::fs::remove_dir_all(dir); +} + +#[tokio::test] +async fn getdata_stops_when_send_budget_is_already_over() { + let (dir, hub) = crate::chain::tiny_regtest_hub_labeled("gd-budget"); + hub.ensure_genesis().unwrap(); + let peers = crate::peers::PeerHub::new(); + let peer = inbound_peer(&peers); + peer.note_send_queued(crate::peers::PEER_SEND_BUDGET + 1); + let (out_tx, mut rx) = mpsc::unbounded_channel(); + let mut follow = PeerFollowState::new(); + let genesis = hub.tip_hash().unwrap(); + serve_getdata( + &hub, + &out_tx, + &mut follow, + Some(&peer), + &[Inventory::WitnessBlock(genesis)], + ) + .await + .unwrap(); + assert!( + rx.try_recv().is_err(), + "a getdata must not queue another block once the send budget is over" + ); + assert_eq!(peer.send_queued(), crate::peers::PEER_SEND_BUDGET + 1); + let _ = std::fs::remove_dir_all(dir); +} + +/// 1001 due wtxids are two inv messages, and both charge the send budget. +#[tokio::test(flavor = "current_thread")] +async fn tx_inv_over_one_thousand_is_two_messages() { + use bitcoin::absolute::LockTime; + use bitcoin::transaction::Version as TxVersion; + use bitcoin::{Amount, OutPoint, Sequence, TxIn, TxOut, Witness}; + + let (dir, hub) = crate::chain::tiny_regtest_hub_labeled("inv-batch"); + hub.ensure_genesis().unwrap(); + hub.generate_to_script(101, op_true(), vec![]).unwrap(); + let t = hub.tip_header().unwrap().time; + hub.clock.set_mock(i64::from(t) + 1); + assert!(!hub.in_ibd(), "tx inv batch pin is not IBD"); + let coinbase = hub + .query + .reconstruct_block_at_height(rbitcoin_primitives::Height(1)) + .unwrap() + .txdata[0] + .clone(); + let each = coinbase.output[0].value.to_sat() / 2 / 1001; + assert!(each > 2_000, "coinbase must fund 1001 spends"); + let fanout = bitcoin::Transaction { + version: TxVersion::TWO, + lock_time: LockTime::ZERO, + input: vec![TxIn { + previous_output: OutPoint { + txid: coinbase.compute_txid(), + vout: 0, + }, + script_sig: ScriptBuf::new(), + sequence: Sequence::ENABLE_RBF_NO_LOCKTIME, + witness: Witness::new(), + }], + output: (0..1001) + .map(|_| TxOut { + value: Amount::from_sat(each), + script_pubkey: op_true(), + }) + .collect(), + }; + hub.generate_to_script(1, op_true(), vec![fanout.clone()]) + .unwrap(); + let mp = crate::tx_relay::MempoolHub::open(dir.join("mp"), std::sync::Arc::clone(&hub.query)) + .unwrap(); + mp.set_relay_enabled(true); + assert!(hub.attach_mempool(mp).is_ok()); + let mempool = hub.mempool().unwrap(); + let parent = fanout.compute_txid(); + for vout in 0..1001u32 { + let child = bitcoin::Transaction { + version: TxVersion::TWO, + lock_time: LockTime::ZERO, + input: vec![TxIn { + previous_output: OutPoint { txid: parent, vout }, + script_sig: ScriptBuf::new(), + sequence: Sequence::ENABLE_RBF_NO_LOCKTIME, + witness: Witness::new(), + }], + output: vec![TxOut { + value: Amount::from_sat(each - 1_000), + script_pubkey: op_true(), + }], + }; + mempool + .accept_tx(&child) + .unwrap_or_else(|e| panic!("accept vout {vout}: {e}")); + } + let peers = crate::peers::PeerHub::new(); + let addr = std::net::SocketAddr::from(([127, 0, 0, 1], 9)); + 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, + false, + crate::peers::PeerConnType::OutboundFullRelay, + ); + peer.request_tx_inv(); + let (out_tx, mut rx) = mpsc::unbounded_channel(); + let before = peer.send_queued(); + queue_due_tx_invs(&hub, peer.as_ref(), &CappedSet::new(), &out_tx); + let mut msgs = 0usize; + let mut items = 0usize; + while let Ok(msg) = rx.try_recv() { + match msg.expect_msg() { + NetworkMessage::Inv(v) => { + msgs += 1; + assert!(v.len() <= 1000, "one inv holds at most 1000"); + items += v.len(); + } + other => panic!("expected inv, got {other:?}"), + } + } + assert_eq!(items, 1001, "every accepted tx is announced"); + assert_eq!(msgs, 2, "1001 tx invs are two messages"); + assert!( + peer.send_queued() > before, + "batched inv must charge the send budget" + ); + let _ = std::fs::remove_dir_all(dir); +} + /// FNV-1a of the address bytes. Duplicated here so a broken mixer in /// `addr_relay_key` cannot satisfy the assertion by changing both sides. fn addr_key_oracle(msg: &bitcoin::p2p::address::AddrV2Message) -> u64 { diff --git a/crates/rbitcoin-net/src/tx_relay.rs b/crates/rbitcoin-net/src/tx_relay.rs index fe3d08e31..23e559d20 100644 --- a/crates/rbitcoin-net/src/tx_relay.rs +++ b/crates/rbitcoin-net/src/tx_relay.rs @@ -29,6 +29,10 @@ use tokio::sync::broadcast; use crate::fee_history::{FeeHistory, HistoricalFeeBlock}; use crate::fee_history_file; +#[path = "parent_req.rs"] +mod parent_req; +pub(crate) use parent_req::{DueParent, ParentNote}; + /// Max age of a published fee snapshot before refresh (request path is still Arc-load only /// after a concurrent refresh has finished; see [`MempoolHub::maybe_refresh_fee_snapshot`]). const FEE_SNAPSHOT_MAX_AGE: Duration = Duration::from_secs(1); @@ -469,13 +473,6 @@ pub(crate) const GETDATA_TX_INTERVAL_SECS: u64 = 60; /// Core `MAX_PEER_TX_REQUEST_IN_FLIGHT`. const MAX_PARENTS_PER_PARK: usize = 100; -struct ParentAnn { - peer: u64, - preferred: bool, - reqtime: u64, - requested_until: Option, - failed: bool, -} /// INV AlreadyHave for recently confirmed txid/wtxid (Core rolling bloom is ~100k). const RECENT_CONFIRMED_CAP: usize = 65_536; @@ -640,8 +637,8 @@ pub struct MempoolHub { min_live_accept_at: AtomicU64, /// Cached tip MTP for accept (`{header_fk, ctx}`). tip_ctx: Mutex>, - /// TxRequestTracker-shaped missing-parent GETDATA (txid hash → announcers). - parent_req: Mutex>>, + /// TxRequestTracker-shaped missing-parent GETDATA. + parent_req: Mutex, } impl MempoolHub { @@ -772,7 +769,7 @@ impl MempoolHub { age_inv: Mutex::new(BTreeMap::new()), min_live_accept_at: AtomicU64::new(u64::MAX), tip_ctx: Mutex::new(None), - parent_req: Mutex::new(HashMap::new()), + parent_req: Mutex::new(parent_req::ParentTracker::new()), }; // Schema ≤ 2 records carry no sigop cost: fill (or evict) before serving. hub.lock_write() @@ -2529,130 +2526,55 @@ impl MempoolHub { if self.parent_already_have(p) { continue; } - let anns = g.entry(p.to_byte_array()).or_default(); - if anns.iter().any(|a| a.peer == peer) { - continue; - } - anns.push(ParentAnn { - peer, - preferred, - reqtime, - requested_until: None, - failed: false, - }); + g.schedule(p.to_byte_array(), peer, preferred, reqtime); } } - pub(crate) fn note_inv_tx_requested(&self, peer: u64, hash: [u8; 32], inbound: bool, now: u64) { - let mut g = self.parent_req.lock().unwrap(); - let anns = g.entry(hash).or_default(); - if let Some(a) = anns.iter_mut().find(|a| a.peer == peer) { - if a.requested_until.is_none() && !a.failed { - a.requested_until = Some(now.saturating_add(GETDATA_TX_INTERVAL_SECS)); - } - return; - } - anns.push(ParentAnn { - peer, - preferred: !inbound, - reqtime: now, - requested_until: Some(now.saturating_add(GETDATA_TX_INTERVAL_SECS)), - failed: false, - }); + /// `wtxid` records that this announcement is a wtxid inv. + /// [`ParentNote::PeerCapped`] is this peer's own cap. [`ParentNote::GlobalFull`] + /// skips the row and is not misbehavior. + pub(crate) fn note_inv_tx_requested( + &self, + peer: u64, + hash: [u8; 32], + inbound: bool, + now: u64, + wtxid: bool, + ) -> ParentNote { + self.parent_req + .lock() + .unwrap() + .note_inv(peer, hash, inbound, now, wtxid) } pub(crate) fn forget_parent_anns_for_peer(&self, peer: u64) { - let mut g = self.parent_req.lock().unwrap(); - g.retain(|_, anns| { - anns.retain(|a| a.peer != peer); - !anns.is_empty() - }); + self.parent_req.lock().unwrap().forget_peer(peer); } pub(crate) fn announcer_peers_for(&self, txid: &Txid, wtxid: &Wtxid) -> Vec { - let g = self.parent_req.lock().unwrap(); - let mut peers = Vec::new(); - for hash in [txid.to_byte_array(), wtxid.to_byte_array()] { - if let Some(anns) = g.get(&hash) { - for a in anns { - if !a.failed && !peers.contains(&a.peer) { - peers.push(a.peer); - } - } - } - } - peers + self.parent_req + .lock() + .unwrap() + .announcer_peers([txid.to_byte_array(), wtxid.to_byte_array()]) } pub(crate) fn resolve_tx_request(&self, txid: &Txid, wtxid: &Wtxid, admitted: bool) { - let mut g = self.parent_req.lock().unwrap(); - for hash in [txid.to_byte_array(), wtxid.to_byte_array()] { - if admitted { - g.remove(&hash); - continue; - } - let Some(anns) = g.get_mut(&hash) else { - continue; - }; - for a in anns.iter_mut() { - if a.requested_until.is_some() { - a.requested_until = None; - a.failed = true; - } - } - if anns.iter().all(|a| a.failed) { - g.remove(&hash); - } - } + self.parent_req + .lock() + .unwrap() + .resolve([txid.to_byte_array(), wtxid.to_byte_array()], admitted); } - pub(crate) fn take_due_parent_getdata(&self, peer: u64, now: u64) -> Vec { + pub(crate) fn take_due_parent_getdata(&self, peer: u64, now: u64) -> Vec { let mut g = self.parent_req.lock().unwrap(); - let hashes: Vec<[u8; 32]> = g.keys().copied().collect(); - let mut out = Vec::new(); - let mut drop_keys = Vec::new(); - for hash in hashes { - let Some(anns) = g.get_mut(&hash) else { - continue; - }; - for a in anns.iter_mut() { - if let Some(exp) = a.requested_until { - if exp <= now { - a.requested_until = None; - a.failed = true; - } - } - } - if anns.iter().all(|a| a.failed) { - drop_keys.push(hash); - continue; - } - let txid = Txid::from_byte_array(hash); - if self.parent_already_have(&txid) { - drop_keys.push(hash); - continue; - } - if anns.iter().any(|a| a.requested_until.is_some()) { - continue; - } - let has_pref = anns - .iter() - .any(|a| a.preferred && !a.failed && a.reqtime <= now); - let chosen = anns.iter_mut().find(|a| { - !a.failed && a.reqtime <= now && a.peer == peer && (!has_pref || a.preferred) - }); - if let Some(a) = chosen { - a.requested_until = Some(now.saturating_add(GETDATA_TX_INTERVAL_SECS)); - out.push(txid); - if out.len() >= MAX_PARENTS_PER_PARK { - break; - } + g.take_due(peer, now, |hash, wtxid| { + if wtxid { + let w = Wtxid::from_byte_array(*hash); + self.try_contains_wtxid(&w) || self.try_recent_reject(&w) + } else { + self.parent_already_have(&Txid::from_byte_array(*hash)) } - } - for k in drop_keys { - g.remove(&k); - } - out + }) } /// Re-admit txs after reorg disconnect (best-effort). @@ -5045,6 +4967,168 @@ mod tests { let _ = std::fs::remove_dir_all(&store_dir); } + fn parent_hash(i: u32) -> [u8; 32] { + let mut h = [0u8; 32]; + h[..4].copy_from_slice(&i.to_le_bytes()); + h + } + + #[test] + fn parent_req_stops_at_per_peer_cap() { + let dir = tmp(); + let store_dir = tmp(); + let q = Query::open_or_create_tiny(&store_dir).unwrap(); + let hub = MempoolHub::open(&dir, Arc::new(q)).unwrap(); + let cap = parent_req::MAX_PARENT_ANN_PER_PEER; + for i in 0..(cap as u32 + 10) { + let accepted = hub.note_inv_tx_requested(7, parent_hash(i), false, 1_000, false); + if i < cap as u32 { + assert_eq!( + accepted, + ParentNote::Accepted, + "announcement {i} under the cap" + ); + } else { + assert_eq!( + accepted, + ParentNote::PeerCapped, + "announcement {i} past this peer's cap" + ); + } + } + let _ = std::fs::remove_dir_all(&dir); + let _ = std::fs::remove_dir_all(&store_dir); + } + + #[test] + fn parent_req_stops_at_global_cap() { + let dir = tmp(); + let store_dir = tmp(); + let q = Query::open_or_create_tiny(&store_dir).unwrap(); + let hub = MempoolHub::open(&dir, Arc::new(q)).unwrap(); + let per_peer = parent_req::MAX_PARENT_ANN_PER_PEER as u32; + let global = parent_req::MAX_PARENT_ANN_GLOBAL as u32; + let mut n = 0u32; + let mut peer = 1u64; + while n < global { + assert_eq!( + hub.note_inv_tx_requested(peer, parent_hash(n), false, 1_000, false), + ParentNote::Accepted + ); + n += 1; + if n.is_multiple_of(per_peer) { + peer += 1; + } + } + assert_eq!( + hub.note_inv_tx_requested(peer, parent_hash(n), false, 1_000, false), + ParentNote::GlobalFull, + "a peer under its own cap is not misbehavior when the table is full" + ); + let _ = std::fs::remove_dir_all(&dir); + let _ = std::fs::remove_dir_all(&store_dir); + } + + #[test] + fn second_peer_take_due_does_not_drop_other_peers_keys() { + let dir = tmp(); + let store_dir = tmp(); + let q = Query::open_or_create_tiny(&store_dir).unwrap(); + let hub = MempoolHub::open(&dir, Arc::new(q)).unwrap(); + let hash = [0x91; 32]; + let txid = Txid::from_byte_array(hash); + let wtxid = Wtxid::from_byte_array(hash); + assert_eq!( + hub.note_inv_tx_requested(1, hash, false, 1_000, false), + ParentNote::Accepted + ); + hub.note_recent_reject(wtxid); + assert_eq!(hub.announcer_peers_for(&txid, &wtxid), vec![1]); + assert!(hub.take_due_parent_getdata(2, 2_000).is_empty()); + assert_eq!( + hub.announcer_peers_for(&txid, &wtxid), + vec![1], + "peer 2's heartbeat must not walk peer 1's keys" + ); + let _ = std::fs::remove_dir_all(&dir); + let _ = std::fs::remove_dir_all(&store_dir); + } + + #[test] + fn wtxid_followup_is_requested_as_wtx() { + let dir = tmp(); + let store_dir = tmp(); + let q = Query::open_or_create_tiny(&store_dir).unwrap(); + let hub = MempoolHub::open(&dir, Arc::new(q)).unwrap(); + let hash = [0xab; 32]; + assert_eq!( + hub.note_inv_tx_requested(1, hash, false, 1_000, true), + ParentNote::Accepted + ); + assert_eq!( + hub.note_inv_tx_requested(3, hash, false, 1_000, true), + ParentNote::Accepted, + "the wtxid follow-up is its own row" + ); + let missing = BTreeSet::from([Txid::from_byte_array(hash)]); + hub.schedule_orphan_parents(&missing, 2, false, 1_000); + let expired = 1_000 + GETDATA_TX_INTERVAL_SECS; + assert!(hub.take_due_parent_getdata(1, expired).is_empty()); + let follow = hub.take_due_parent_getdata(3, expired); + assert_eq!(follow.len(), 1); + assert_eq!(follow[0].hash, hash); + assert!(follow[0].wtxid, "wtxid follow-up must stay WTx"); + assert!( + hub.take_due_parent_getdata(2, expired).is_empty(), + "a txid parent waits while those bytes are in flight" + ); + let parent = hub.take_due_parent_getdata(2, expired + GETDATA_TX_INTERVAL_SECS); + assert_eq!(parent.len(), 1); + assert_eq!(parent[0].hash, hash); + assert!( + !parent[0].wtxid, + "a txid parent is not fetched as a wtxid because another peer announced these bytes" + ); + let _ = std::fs::remove_dir_all(&dir); + let _ = std::fs::remove_dir_all(&store_dir); + } + + #[test] + fn same_peer_wtxid_inflight_expires_back_to_txid_parent() { + let dir = tmp(); + let store_dir = tmp(); + let q = Query::open_or_create_tiny(&store_dir).unwrap(); + let hub = MempoolHub::open(&dir, Arc::new(q)).unwrap(); + let hash = [0xbc; 32]; + let missing = BTreeSet::from([Txid::from_byte_array(hash)]); + let t0 = 1_000u64; + hub.schedule_orphan_parents(&missing, 2, false, t0); + assert_eq!( + hub.note_inv_tx_requested(2, hash, false, t0, true), + ParentNote::Accepted, + "the same peer's wtxid inv is a second announcement" + ); + let due = t0 + TXID_RELAY_DELAY_SECS; + assert!( + hub.take_due_parent_getdata(2, due).is_empty(), + "the txid parent waits while this peer's wtxid request is in flight" + ); + let expired = t0 + GETDATA_TX_INTERVAL_SECS; + let parent = hub.take_due_parent_getdata(2, expired); + assert_eq!( + parent.len(), + 1, + "the txid parent is requested once the wtxid window ends" + ); + assert_eq!(parent[0].hash, hash); + assert!( + !parent[0].wtxid, + "after the window the parent is requested as a txid" + ); + let _ = std::fs::remove_dir_all(&dir); + let _ = std::fs::remove_dir_all(&store_dir); + } + #[test] fn take_due_parent_getdata_waits_inbound_txid_delay() { let dir = tmp(); @@ -5093,7 +5177,7 @@ mod tests { let hash = [0x44; 32]; let txid = Txid::from_byte_array(hash); let wtxid = Wtxid::from_byte_array(hash); - hub.note_inv_tx_requested(1, hash, true, 1_000); + hub.note_inv_tx_requested(1, hash, true, 1_000, false); let missing = BTreeSet::from([txid]); hub.schedule_orphan_parents(&missing, 2, true, 1_000); let due = 1_000 + NONPREF_PEER_TX_DELAY_SECS + TXID_RELAY_DELAY_SECS; @@ -5105,7 +5189,9 @@ mod tests { hub.schedule_orphan_parents(&missing, 2, true, due); let got = hub .take_due_parent_getdata(2, due + NONPREF_PEER_TX_DELAY_SECS + TXID_RELAY_DELAY_SECS); - assert_eq!(got, vec![txid]); + assert_eq!(got.len(), 1); + assert_eq!(got[0].hash, txid.to_byte_array()); + assert!(!got[0].wtxid); let _ = std::fs::remove_dir_all(&dir); let _ = std::fs::remove_dir_all(&store_dir); } diff --git a/crates/rbitcoin-net/src/v2.rs b/crates/rbitcoin-net/src/v2.rs index fae36ffdd..52ddc4bec 100644 --- a/crates/rbitcoin-net/src/v2.rs +++ b/crates/rbitcoin-net/src/v2.rs @@ -208,13 +208,19 @@ impl V2SessionReader { } /// Next genuine application contents (skips decoys). Checks the decrypted /// length prefix against [`MAX_V2_CONTENTS_LEN`] before reading the body. - async fn read_genuine_contents(&mut self, mut on_progress: F) -> Result, NetError> + async fn read_genuine_contents( + &mut self, + mut on_progress: F, + mut on_decoy: D, + ) -> Result, NetError> where F: FnMut(usize), + D: FnMut(usize) -> Result<(), NetError>, { loop { let (packet_type, plaintext) = self.read_packet(&mut on_progress).await?; if packet_type == PacketType::Decoy { + on_decoy(plaintext.len())?; continue; } // plaintext = 1-byte ignore header + application contents @@ -440,7 +446,6 @@ pub fn parse_v2_contents(magic: Magic, contents: &[u8]) -> Result (command_to_12(name), contents[1..].to_vec()), None => { - rbitcoin_log::info!("{}", v2_invalid_message_type_log()); return Err(NetError::InvalidV2Type { contents_len: contents.len(), }); @@ -453,7 +458,6 @@ pub fn parse_v2_contents(magic: Magic, contents: &[u8]) -> Result(reader: &mut V2SessionReader) -> Result( where R: AsyncRead + Unpin + Send, { - read_v2_frame_with_progress(reader, magic, |_| {}).await + read_v2_frame_with_progress(reader, magic, |_| {}, |_| Ok(())).await } /// Read the next genuine frame; `on_progress` is invoked as ciphertext body /// bytes arrive, then again with decrypted content length. -pub async fn read_v2_frame_with_progress( +pub async fn read_v2_frame_with_progress( reader: &mut V2SessionReader, magic: Magic, mut on_progress: F, + on_decoy: D, ) -> Result where R: AsyncRead + Unpin + Send, F: FnMut(usize), + D: FnMut(usize) -> Result<(), NetError>, { - let contents = reader.read_genuine_contents(&mut on_progress).await?; + let contents = reader + .read_genuine_contents(&mut on_progress, on_decoy) + .await?; on_progress(contents.len()); parse_v2_contents(magic, &contents) } @@ -1054,6 +1062,138 @@ mod tests { } } + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn decoy_packet_is_handed_to_the_rate_hook() { + use bip324::OutboundCipher; + use tokio::io::AsyncWriteExt; + + let magic = signet_magic(); + let magic_b = magic.to_bytes(); + let (client, server) = duplex(64 * 1024); + let server_task = tokio::spawn(async move { + let (rh, wh) = tokio::io::split(server); + let reader = BufReader::new(rh); + let protocol = Protocol::new(magic_b, Role::Responder, None, None, reader, wh) + .await + .expect("server handshake"); + let (r, _w) = protocol.into_split(); + let mut reader = V2SessionReader::from_protocol_reader(r); + let mut seen = Vec::new(); + let frame = read_v2_frame_with_progress( + &mut reader, + magic, + |_| {}, + |n| { + seen.push(n); + Ok(()) + }, + ) + .await + .expect("genuine frame after decoy"); + (seen, frame.is_ping()) + }); + + let (rh, wh) = tokio::io::split(client); + let reader = BufReader::new(rh); + let protocol = Protocol::new(magic_b, Role::Initiator, None, None, reader, wh) + .await + .expect("client handshake"); + let (_r, w) = protocol.into_split(); + let (mut cipher, mut raw_w) = w.into_inner(); + let decoy_plain = vec![0u8; 32]; + let mut packet = vec![0u8; OutboundCipher::encryption_buffer_len(decoy_plain.len())]; + cipher + .encrypt(&decoy_plain, &mut packet, PacketType::Decoy, None) + .expect("encrypt decoy"); + raw_w.write_all(&packet).await.unwrap(); + let ping = encode_v2_contents(NetworkMessage::Ping(7)).unwrap(); + let mut packet = vec![0u8; OutboundCipher::encryption_buffer_len(ping.len())]; + cipher + .encrypt(&ping, &mut packet, PacketType::Genuine, None) + .expect("encrypt ping"); + raw_w.write_all(&packet).await.unwrap(); + raw_w.flush().await.unwrap(); + + let (seen, is_ping) = tokio::time::timeout(std::time::Duration::from_secs(5), server_task) + .await + .expect("decoy read timed out") + .expect("server task join"); + assert!(!seen.is_empty(), "a decoy must reach the rate hook"); + assert!(is_ping, "the genuine packet after the decoy is the ping"); + } + + /// Tip-follow and the IBD reader both call [`crate::peer_dos::decoy_stays`]. + /// One decoy that does not fit scores and the following frame is still read. + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn decoy_over_the_window_scores_like_other_frames() { + use crate::peer::BAN_SCORE_THRESHOLD; + use crate::peer_dos::{decoy_stays, PeerRateLimiter, RATE_LIMIT_BAN_SCORE}; + use bip324::OutboundCipher; + use tokio::io::AsyncWriteExt; + + let magic = signet_magic(); + let magic_b = magic.to_bytes(); + let (client, server) = duplex(64 * 1024); + let server_task = tokio::spawn(async move { + let (rh, wh) = tokio::io::split(server); + let reader = BufReader::new(rh); + let protocol = Protocol::new(magic_b, Role::Responder, None, None, reader, wh) + .await + .expect("server handshake"); + let (r, _w) = protocol.into_split(); + let mut reader = V2SessionReader::from_protocol_reader(r); + let mut rate = PeerRateLimiter::new(1, 10_000); + let mut score = 0u32; + assert!(rate.note(1), "the window is already full"); + let frame = read_v2_frame_with_progress( + &mut reader, + magic, + |_| {}, + |n| { + if decoy_stays(&mut rate, &mut score, n, BAN_SCORE_THRESHOLD) { + Ok(()) + } else { + Err(NetError::Protocol("peer misbehavior threshold")) + } + }, + ) + .await + .expect("one decoy over the window must not abort the read"); + (score, frame.is_ping()) + }); + + let (rh, wh) = tokio::io::split(client); + let reader = BufReader::new(rh); + let protocol = Protocol::new(magic_b, Role::Initiator, None, None, reader, wh) + .await + .expect("client handshake"); + let (_r, w) = protocol.into_split(); + let (mut cipher, mut raw_w) = w.into_inner(); + let decoy_plain = vec![0u8; 32]; + let mut packet = vec![0u8; OutboundCipher::encryption_buffer_len(decoy_plain.len())]; + cipher + .encrypt(&decoy_plain, &mut packet, PacketType::Decoy, None) + .expect("encrypt decoy"); + raw_w.write_all(&packet).await.unwrap(); + let ping = encode_v2_contents(NetworkMessage::Ping(7)).unwrap(); + let mut packet = vec![0u8; OutboundCipher::encryption_buffer_len(ping.len())]; + cipher + .encrypt(&ping, &mut packet, PacketType::Genuine, None) + .expect("encrypt ping"); + raw_w.write_all(&packet).await.unwrap(); + raw_w.flush().await.unwrap(); + + let (score, is_ping) = tokio::time::timeout(std::time::Duration::from_secs(5), server_task) + .await + .expect("decoy overflow read timed out") + .expect("server task join"); + assert_eq!(score, RATE_LIMIT_BAN_SCORE); + assert!( + is_ping, + "the genuine packet after the over-window decoy is the ping" + ); + } + #[test] fn truncated_v2_long_command_is_protocol() { assert!(matches!( diff --git a/docs/external_findings/052-livera-review-index.md b/docs/external_findings/052-livera-review-index.md index 643f170da..bde016c61 100644 --- a/docs/external_findings/052-livera-review-index.md +++ b/docs/external_findings/052-livera-review-index.md @@ -6,22 +6,22 @@ steps. | Id | Severity | Topic | Status | Regression | |----|----------|--------|--------|------------| -| C1 | critical | Unbounded parent-request tracker | open | — | -| N2 | low | Wtxid follow-up requested as a txid | open | — | -| H1 | high | Decoy and invalid-type packets skip the rate window | open | — | +| 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 | — | | 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 | — | -| N1 | medium | Inv getdata does not charge the send budget | open | — | -| M1 | medium | Block getdata can queue past the send budget | open | — | +| 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 | — | | M3 | medium | RPC listener has no accept timeout; long-poll holds a permit | open | — | | M4 | medium | fuse8 segment length need not be a power of two | open | — | | M5 | medium | Class A bulk read ignores the published end | open | — | | M6 | medium | Testnet milestone is height-only | open | — | -| M7 | medium | Mempool inv is one message per transaction | open | — | +| M7 | medium | Mempool inv is one message per transaction | fixed | `tx_inv_over_one_thousand_is_two_messages` ([058](./058-tx-inv-batch.md)) | | M8 | medium | Pending blocks and orphans are count-capped only | open | — | | M9 | high | io_uring drop can free a buffer the kernel still owns | open | — | | M10 | low | Secret types derive Debug; create-then-chmod | open | — | @@ -34,7 +34,7 @@ steps. | L10 | low | Datadir lock follows a symlink | open | — | | L11 | low | Conf parse errors echo the raw line | open | — | | L12 | low | Invalid-hash set grows without a cap | open | — | -| L13 | low | Rate window grants two budgets at the boundary | open | — | +| L13 | low | Rate window grants two budgets at the boundary | fixed | `rate_limiter_boundary_does_not_grant_a_second_budget` ([059](./059-rate-window-boundary.md)) | | L14 | low | Mempool expiry runs only on admission | open | — | | L15 | low | Write jobs form a mutable slice over a shared buffer | open | — | | L2 | — | P2PKH fast path skips FindAndDelete | rejected | The fast path is the 25-byte template. A DER signature does not fit in that scriptCode, so FindAndDelete cannot change it. | diff --git a/docs/external_findings/053-parent-req-cap.md b/docs/external_findings/053-parent-req-cap.md new file mode 100644 index 000000000..8dcae505d --- /dev/null +++ b/docs/external_findings/053-parent-req-cap.md @@ -0,0 +1,9 @@ +# 053 — Parent-request tracker cap + +**Severity:** critical +**Status:** fixed +**Found by:** Stephan Livera, 2026-10-02 (C1) + +The parent-request tracker grew with every in-flight parent on every peer. It now stops at a per-peer cap and a process-wide cap. Taking a due key on one peer does not drop another peer's keys. + +**Regression:** `rbitcoin-net` `parent_req_stops_at_per_peer_cap`, `parent_req_stops_at_global_cap`, `second_peer_take_due_does_not_drop_other_peers_keys`. diff --git a/docs/external_findings/054-wtxid-getdata.md b/docs/external_findings/054-wtxid-getdata.md new file mode 100644 index 000000000..7d88894de --- /dev/null +++ b/docs/external_findings/054-wtxid-getdata.md @@ -0,0 +1,9 @@ +# 054 — Wtxid follow-up is a wtxid getdata + +**Severity:** low +**Status:** fixed +**Found by:** Stephan Livera, 2026-10-02 (N2) + +A wtxid announcement was re-requested as a txid. The follow-up getdata uses a wtxid. + +**Regression:** `rbitcoin-net` `wtxid_followup_is_requested_as_wtx`. diff --git a/docs/external_findings/055-decoy-rate.md b/docs/external_findings/055-decoy-rate.md new file mode 100644 index 000000000..aa7c182d1 --- /dev/null +++ b/docs/external_findings/055-decoy-rate.md @@ -0,0 +1,9 @@ +# 055 — Decoy packets count toward the rate window + +**Severity:** high +**Status:** fixed +**Found by:** Stephan Livera, 2026-10-02 (H1) + +Decoy and invalid-type packets skipped the per-peer rate window and were logged at info. They are handed to the rate hook on tip-follow and during initial download. A full window disconnects. The window does not allocate per message. + +**Regression:** `rbitcoin-net` `decoy_packet_is_handed_to_the_rate_hook`. diff --git a/docs/external_findings/056-inv-getdata-budget.md b/docs/external_findings/056-inv-getdata-budget.md new file mode 100644 index 000000000..c7b5c4d9c --- /dev/null +++ b/docs/external_findings/056-inv-getdata-budget.md @@ -0,0 +1,9 @@ +# 056 — Inv getdata charges the send budget + +**Severity:** medium +**Status:** fixed +**Found by:** Stephan Livera, 2026-10-02 (N1) + +Getdata driven by an inv was not charged against the per-peer send budget. Those bytes are charged. + +**Regression:** `rbitcoin-net` `inv_getdata_charges_send_budget`. diff --git a/docs/external_findings/057-block-getdata-budget.md b/docs/external_findings/057-block-getdata-budget.md new file mode 100644 index 000000000..4ec540382 --- /dev/null +++ b/docs/external_findings/057-block-getdata-budget.md @@ -0,0 +1,9 @@ +# 057 — Block serving stops when the send budget is over + +**Severity:** medium +**Status:** fixed +**Found by:** Stephan Livera, 2026-10-02 (M1) + +Block getdata could queue more data after the per-peer send budget was already over. Serving stops once that budget is over. + +**Regression:** `rbitcoin-net` `getdata_stops_when_send_budget_is_already_over`. diff --git a/docs/external_findings/058-tx-inv-batch.md b/docs/external_findings/058-tx-inv-batch.md new file mode 100644 index 000000000..42844a6d8 --- /dev/null +++ b/docs/external_findings/058-tx-inv-batch.md @@ -0,0 +1,9 @@ +# 058 — Mempool announcements batch into one inv + +**Severity:** medium +**Status:** fixed +**Found by:** Stephan Livera, 2026-10-02 (M7) + +Each mempool transaction was its own inv. Announcements are batched up to 1000 per message and still charged against the per-peer send budget. + +**Regression:** `rbitcoin-net` `tx_inv_over_one_thousand_is_two_messages`. diff --git a/docs/external_findings/059-rate-window-boundary.md b/docs/external_findings/059-rate-window-boundary.md new file mode 100644 index 000000000..f044e2d60 --- /dev/null +++ b/docs/external_findings/059-rate-window-boundary.md @@ -0,0 +1,9 @@ +# 059 — Rate window keeps the previous second + +**Severity:** low +**Status:** fixed +**Found by:** Stephan Livera, 2026-10-02 (L13) + +At a one-second boundary the rate window dropped the previous second and granted a second full budget. It keeps that second. The window does not allocate per message. + +**Regression:** `rbitcoin-net` `rate_limiter_boundary_does_not_grant_a_second_budget`. diff --git a/docs/external_findings/README.md b/docs/external_findings/README.md index 349e7cb92..e8aafb3e8 100644 --- a/docs/external_findings/README.md +++ b/docs/external_findings/README.md @@ -57,6 +57,13 @@ rbitcoin reference, or redteam static analysis). Numbered reports live beside th | [043](./043-peer-send-buffer.md) | high | Unbounded per-peer outbound queue | fixed | `hostile_peer_session` | | [040](./040-corrupt-bounds.md) | medium | Corrupt uleb128, seqsigwit lengths, and BDZ modulus | fixed | `read_packed_zero_modulus_or_vertices_is_corrupt` | | [052](./052-livera-review-index.md) | — | Livera review index | index | Q-72 | +| [053](./053-parent-req-cap.md) | critical | Parent-request tracker cap | fixed | `parent_req_stops_at_per_peer_cap` | +| [054](./054-wtxid-getdata.md) | low | Wtxid follow-up is a wtxid getdata | fixed | `wtxid_followup_is_requested_as_wtx` | +| [055](./055-decoy-rate.md) | high | Decoy packets count toward the rate window | fixed | `decoy_packet_is_handed_to_the_rate_hook` | +| [056](./056-inv-getdata-budget.md) | medium | Inv getdata charges the send budget | fixed | `inv_getdata_charges_send_budget` | +| [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` | **012–021:** fuzzamoto differential report (`rbitcoin-report.tar.gz`, baseline `8f3990f`). Report-local 001–010 are **renumbered** here. Identity/BIP30