From 76ea7294e10614d55e8134cce1371ae9f34eaa89 Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Sat, 12 Sep 2026 20:54:03 -0700 Subject: [PATCH 1/3] test: finish P2P Category B 5.1 on fewer default journeys One mature-pad journey covers HB compact, missing-tx getblocktxn connect, and orphan child-then-parent accept. The timeout pad now pins GetAddr 1000/cache, headers-sync stall replace, and self-connect. Fold overlapping compact/orphan guts and collapse the four GetAddr cache units onto one bind-key/TTL table. Keep feeler silence, eviction ranking, park-not-reject logs, and tokio-worker lock. Co-authored-by: Cursor --- crates/rbitcoin-net/src/peer_tests.rs | 227 +------ crates/rbitcoin-net/src/peers.rs | 51 +- .../tests/integration_multinode.rs | 607 +++++++++++------- 3 files changed, 378 insertions(+), 507 deletions(-) diff --git a/crates/rbitcoin-net/src/peer_tests.rs b/crates/rbitcoin-net/src/peer_tests.rs index 6ba06955..5f3d8bf3 100644 --- a/crates/rbitcoin-net/src/peer_tests.rs +++ b/crates/rbitcoin-net/src/peer_tests.rs @@ -272,23 +272,6 @@ fn tip_announce_headers_and_inv() { let _ = std::fs::remove_dir_all(dir); } -#[test] -fn cmpct_announce_uses_generated_tip_body() { - let (dir, hub) = crate::chain::tiny_regtest_hub_labeled("cmpct-announce"); - hub.ensure_genesis().unwrap(); - let hashes = hub - .generate_to_script(1, bitcoin::script::ScriptBuf::new(), vec![]) - .unwrap(); - let hash = hashes[0]; - match cmpct_announce_msg(&hub, &hash, 2) { - Some(NetworkMessage::CmpctBlock(c)) => { - assert_eq!(c.compact_block.header.block_hash(), hash); - } - other => panic!("expected CmpctBlock, got {other:?}"), - } - let _ = std::fs::remove_dir_all(dir); -} - #[test] fn header_getdata_is_compact_after_sendcmpct() { use bitcoin::block::Header; @@ -1839,19 +1822,6 @@ fn blocksonly_relay_perm_tx_invs_other_inbound() { fn cmpct_helpers_without_mempool_and_queue_out_closed() { let (dir, hub) = crate::chain::tiny_regtest_hub_labeled("cmpct-none"); hub.ensure_genesis().unwrap(); - let gen = hub - .query - .reconstruct_block_by_hash(&hub.tip_hash().unwrap().to_byte_array()) - .unwrap() - .unwrap(); - let hsi = HeaderAndShortIds::from_block(&gen, 0xabc, 2, &[]).unwrap(); - assert!( - matches!( - try_reconstruct_cmpct(&hub, &hsi, 2), - Some(CmpctReconstruct::Block(_)) - ), - "coinbase-only compact fills from prefilled txs without a mempool" - ); assert!(hub.mempool().is_none()); // Closed channel → Protocol error. @@ -2884,105 +2854,6 @@ fn parked_orphan_tx_is_not_logged_as_reject() { }); } -/// INV of an already-parked orphan must not re-GETDATA it (AlreadyHave includes -/// the orphanage, matching Core TxDownloadManager). -#[test] -fn inv_of_parked_orphan_does_not_getdata() { - use bitcoin::absolute::LockTime; - use bitcoin::consensus::encode::serialize; - use bitcoin::script::ScriptBuf; - use bitcoin::transaction::Version as TxVersion; - use bitcoin::{Amount, Network, OutPoint, Sequence, Transaction, TxIn, TxOut, Witness}; - use tokio::runtime::Builder; - - fn frame_for(msg: NetworkMessage) -> FramedMessage { - use bitcoin::p2p::message::RawNetworkMessage; - let magic = Magic::from(Network::Regtest); - let raw = RawNetworkMessage::new(magic, msg); - let full = serialize(&raw); - let command: [u8; 12] = full[4..16].try_into().unwrap(); - let payload = full[24..].to_vec(); - FramedMessage { - magic, - command, - payload, - } - } - - let rt = Builder::new_current_thread().enable_all().build().unwrap(); - rt.block_on(async { - let (dir, hub) = crate::chain::tiny_regtest_hub_labeled("orphan-inv"); - hub.ensure_genesis().unwrap(); - let t = hub.tip_header().unwrap().time; - hub.clock.set_mock(i64::from(t) + 1); - let mp = crate::tx_relay::MempoolHub::open(dir.join("mp"), Arc::clone(&hub.query)).unwrap(); - mp.set_relay_enabled(true); - assert!(hub.attach_mempool(mp).is_ok()); - - let orphan = Transaction { - version: TxVersion::TWO, - lock_time: LockTime::ZERO, - input: vec![TxIn { - previous_output: OutPoint { - txid: bitcoin::Txid::from_byte_array([0x22; 32]), - vout: 0, - }, - script_sig: ScriptBuf::new(), - sequence: Sequence::ENABLE_RBF_NO_LOCKTIME, - witness: Witness::new(), - }], - output: vec![TxOut { - value: Amount::from_sat(1000), - script_pubkey: ScriptBuf::from_bytes(vec![0x51]), - }], - }; - let txid = orphan.compute_txid(); - let wtxid = orphan.compute_wtxid(); - let (out_tx, mut out_rx) = mpsc::unbounded_channel(); - let mut follow = PeerFollowState { - wants_headers: false, - wtxid_relay: true, - send_cmpct: false, - cmpct_version: 2u32, - pending_headers: HashMap::new(), - pending_blocks: PendingBlocks::new(), - pending_cmpct: HashMap::new(), - from_this_peer: CappedSet::new(), - requested_blocks: HashSet::new(), - ban_score: 0u32, - }; - handle_peer_frame( - frame_for(NetworkMessage::Tx(orphan)), - &hub, - &out_tx, - &mut follow, - None, - ) - .await - .unwrap(); - assert_eq!(hub.mempool().unwrap().orphan_count(), 1); - while out_rx.try_recv().is_ok() {} - - handle_peer_frame( - frame_for(NetworkMessage::Inv(vec![ - Inventory::WitnessTransaction(txid), - Inventory::WTx(wtxid), - ])), - &hub, - &out_tx, - &mut follow, - None, - ) - .await - .unwrap(); - assert!( - out_rx.try_recv().is_err(), - "parked orphan INV must not GetData" - ); - let _ = std::fs::remove_dir_all(dir); - }); -} - /// Production P2P runs on `tokio-rt-worker`. Parking must not take the mempool /// inner lock on that thread (reader panics, ping/block-sync stall). #[test] @@ -3388,11 +3259,10 @@ fn invalid_getdata_type0_still_serves_tip_block() { }); } -/// Compact helpers with a live mempool hub attached (fill/missing/blocktxn). +/// Compact fill with a live mempool hub must not `list_live` every body. #[test] -fn cmpct_helpers_with_mempool_live_and_blocktxn() { +fn cmpct_helpers_with_mempool_skip_list_live() { use bitcoin::absolute::LockTime; - use bitcoin::bip152::BlockTransactions; use bitcoin::block::{Header, Version}; use bitcoin::script::ScriptBuf; use bitcoin::transaction::Version as TxVersion; @@ -3469,18 +3339,6 @@ fn cmpct_helpers_with_mempool_live_and_blocktxn() { fill.list_live ); - let pc = PendingCmpct { - hsi: hsi.clone(), - missing: missing.clone(), - version: 2, - }; - let bt = BlockTransactions { - block_hash: block.block_hash(), - transactions: vec![spend], - }; - let recon = apply_cmpct_blocktxn(&hub, &pc, &bt).expect("blocktxn fill"); - assert_eq!(recon.txdata.len(), 2); - let _ = std::fs::remove_dir_all(dir); } @@ -4109,9 +3967,6 @@ fn expect_services_from_conn_matches_core() { #[test] fn handshake_disconnect_log_needles() { - let line = connected_to_self_log("127.0.0.1:18444"); - assert!(line.contains("connected to self")); - assert!(line.contains("disconnecting")); assert_eq!( crate::peer::ping_prior_to_verack_log(0), "Unsupported message \"ping\" prior to verack from peer=0" @@ -5679,84 +5534,6 @@ fn compact_tip_announce_must_not_consume_serve_slots() { }); } -/// Coinbase-only compact must reconstruct from prefilled txs; a missing -/// mempool hub must not force a full-getdata fallback that then never -/// arrives (`p2p_compactblocks_hb` 1-block relay). -#[test] -fn coinbase_compact_fills_without_mempool() { - use bitcoin::consensus::encode::serialize; - use bitcoin::Network; - use tokio::runtime::Builder; - - fn frame_for(msg: NetworkMessage) -> FramedMessage { - use bitcoin::p2p::message::RawNetworkMessage; - let magic = Magic::from(Network::Regtest); - let raw = RawNetworkMessage::new(magic, msg); - let full = serialize(&raw); - let command: [u8; 12] = full[4..16].try_into().unwrap(); - FramedMessage { - magic, - command, - payload: full[24..].to_vec(), - } - } - - let rt = Builder::new_current_thread().enable_all().build().unwrap(); - rt.block_on(async { - let (src_dir, src) = crate::chain::tiny_regtest_hub_labeled("cmpct-nomp-src"); - src.ensure_genesis().unwrap(); - let hashes = src - .generate_to_script(1, bitcoin::ScriptBuf::from_bytes(vec![0x51]), vec![]) - .unwrap(); - let hash = hashes[0]; - let block = src - .query - .reconstruct_archived_block(&hash.to_byte_array()) - .unwrap() - .expect("src body"); - let hsi = HeaderAndShortIds::from_block(&block, 1, 2, &[0]).unwrap(); - - let (dir, hub) = crate::chain::tiny_regtest_hub_labeled("cmpct-nomp-dst"); - hub.ensure_genesis().unwrap(); - assert!(hub.mempool().is_none()); - let (out_tx, mut out_rx) = mpsc::unbounded_channel(); - let mut follow = PeerFollowState { - wants_headers: false, - wtxid_relay: false, - send_cmpct: true, - cmpct_version: 2u32, - pending_headers: HashMap::new(), - pending_blocks: PendingBlocks::new(), - pending_cmpct: HashMap::new(), - from_this_peer: CappedSet::new(), - requested_blocks: HashSet::new(), - ban_score: 0u32, - }; - handle_peer_frame( - frame_for(NetworkMessage::CmpctBlock(CmpctBlock { - compact_block: hsi, - })), - &hub, - &out_tx, - &mut follow, - None, - ) - .await - .unwrap(); - assert!( - hub.has_block(&hash), - "coinbase compact must connect without a mempool hub" - ); - while let Ok(msg) = out_rx.try_recv().map(PeerOut::expect_msg) { - if matches!(msg, NetworkMessage::GetData(_)) { - panic!("coinbase compact must not fall back to getdata, got {msg:?}"); - } - } - let _ = std::fs::remove_dir_all(src_dir); - let _ = std::fs::remove_dir_all(dir); - }); -} - #[test] fn tip_event_for_announce_on_lagged_uses_current_hub_tip() { use bitcoin::ScriptBuf; diff --git a/crates/rbitcoin-net/src/peers.rs b/crates/rbitcoin-net/src/peers.rs index 0cc76261..98869c76 100644 --- a/crates/rbitcoin-net/src/peers.rs +++ b/crates/rbitcoin-net/src/peers.rs @@ -2574,56 +2574,31 @@ mod tests { } #[test] - fn getaddr_cache_repeats_same_bind() { - let hub = PeerHub::new(); - hub.set_mock_now(1_700_000_000); - hub.set_addrman(Arc::new(Mutex::new(fill_addrman(5_000)))); - let bind = SocketAddr::from(([127, 0, 0, 1], 18444)); - let a = addr_ips(&hub.addr_response_for_bind(bind)); - let b = addr_ips(&hub.addr_response_for_bind(bind)); - assert_eq!(a.len(), 1000); - assert_eq!(a, b); - } - - #[test] - fn getaddr_cache_distinct_listens_differ() { + fn getaddr_cache_bind_key_and_ttl() { let hub = PeerHub::new(); hub.set_mock_now(1_700_000_000); hub.set_addrman(Arc::new(Mutex::new(fill_addrman(5_000)))); let a = addr_ips(&hub.addr_response_for_bind(SocketAddr::from(([127, 0, 0, 1], 18444)))); let b = addr_ips(&hub.addr_response_for_bind(SocketAddr::from(([127, 0, 0, 1], 18445)))); let c = addr_ips(&hub.addr_response_for_bind(SocketAddr::from(([127, 0, 0, 1], 18446)))); + let mapped = addr_ips(&hub.addr_response_for_bind(SocketAddr::from(( + Ipv4Addr::new(127, 0, 0, 1).to_ipv6_mapped(), + 18444, + )))); assert_eq!(a.len(), 1000); assert_eq!(b.len(), 1000); assert_eq!(c.len(), 1000); + assert_eq!( + a, mapped, + "IPv4-mapped IPv6 must share the clearnet cache key" + ); assert_ne!(a, b); assert_ne!(a, c); assert_ne!(b, c); - } - - #[test] - fn getaddr_cache_ipv4_mapped_shares_clearnet_key() { - let hub = PeerHub::new(); - hub.set_mock_now(1_700_000_000); - hub.set_addrman(Arc::new(Mutex::new(fill_addrman(5_000)))); - let v4 = SocketAddr::from(([127, 0, 0, 1], 18444)); - let v6 = SocketAddr::from((Ipv4Addr::new(127, 0, 0, 1).to_ipv6_mapped(), 18444)); - let a = addr_ips(&hub.addr_response_for_bind(v4)); - let b = addr_ips(&hub.addr_response_for_bind(v6)); - assert_eq!(a.len(), 1000); - assert_eq!(a, b); - } - - #[test] - fn getaddr_cache_expires_after_24h() { - let hub = PeerHub::new(); - hub.set_mock_now(1_700_000_000); - hub.set_addrman(Arc::new(Mutex::new(fill_addrman(5_000)))); - let bind = SocketAddr::from(([127, 0, 0, 1], 18444)); - let first = addr_ips(&hub.addr_response_for_bind(bind)); hub.set_mock_now(1_700_000_000 + 24 * 60 * 60); - let second = addr_ips(&hub.addr_response_for_bind(bind)); - assert_eq!(first.len(), 1000); - assert_ne!(first, second); + let expired = + addr_ips(&hub.addr_response_for_bind(SocketAddr::from(([127, 0, 0, 1], 18444)))); + assert_eq!(expired.len(), 1000); + assert_ne!(a, expired); } } diff --git a/crates/rbitcoin-test/tests/integration_multinode.rs b/crates/rbitcoin-test/tests/integration_multinode.rs index 0897f7b0..a9d66d52 100644 --- a/crates/rbitcoin-test/tests/integration_multinode.rs +++ b/crates/rbitcoin-test/tests/integration_multinode.rs @@ -2,11 +2,10 @@ //! //! **Tier A (default + CI `multinode` job):** single-hop IBD (8 blocks), cold //! reconstruct serve (10 blocks). Hard wall timeouts; hang-free on CI-class hosts. -//! **Tier B (default suite):** handshake timeout / GetAddr / keepalive ping, -//! HB compact tip-follow, compact `getblocktxn` for a missing extra tx, -//! mempool orphan child GetData of parent, outbound feeler complete-and-close, -//! inbound-full reject, hub reorg (including leftover/BadPrev orphan that must not -//! blacklist). +//! **Tier B (default suite):** handshake timeout / GetAddr cache / keepalive ping, +//! compact HB + missing-tx `getblocktxn` + orphan child→parent on one mature pad, +//! outbound feeler complete-and-close + inbound-full reject, hub reorg +//! (leftover/BadPrev orphan that must not blacklist). //! **Tier C (`#[ignore]`):** multi-hop, tip-follow, 48-block dual seeder, mesh — //! `scripts/integration.sh` or `-- --ignored` only. @@ -47,6 +46,55 @@ async fn start_node_inbound(dir: &TempDir, max_inbound: usize) -> P2PNode { .expect("listen") } +fn open_padded_query(dir: &TempDir) -> Query { + use rbitcoin_consensus::accept_and_connect_block; + use rbitcoin_test::pad_empty_from; + + let q = Query::open_or_create_tiny(dir.path().join("store")).unwrap(); + let params = ChainParams::regtest(); + let genesis = regtest_genesis(); + accept_and_connect_block(&q, ¶ms, Height::GENESIS, &genesis, Milestone::NONE).unwrap(); + let last = params.coinbase_maturity() + 1; + pad_empty_from( + &q, + ¶ms, + genesis.block_hash(), + genesis.header.time, + 1, + last, + ); + q +} + +async fn start_padded(dir: &TempDir) -> P2PNode { + let q = open_padded_query(dir); + P2PNode::start( + "127.0.0.1:0".parse().unwrap(), + q, + ChainParams::regtest(), + Milestone::NONE, + ) + .await + .expect("listen") +} + +async fn wait_ms_until( + ms: u64, + mut pred: impl FnMut() -> bool, + on_timeout: impl FnOnce() -> String, +) { + let deadline = tokio::time::Instant::now() + Duration::from_millis(ms); + loop { + if pred() { + return; + } + if tokio::time::Instant::now() >= deadline { + panic!("{}", on_timeout()); + } + tokio::time::sleep(Duration::from_millis(20)).await; + } +} + async fn seed_chain(node: &P2PNode, blocks: u32) { let genesis = regtest_genesis(); node.ingest_block(0, genesis.clone()).unwrap(); @@ -104,9 +152,11 @@ async fn two_node_header_and_block_sync() { } /// In-tree P2P client (no Core functional): peertimeout of a v1-magic inbound, -/// AddrFetch GetAddr, and one post-verack keepalive ping/pong. +/// full-relay GetAddr cache (1000 / 23%), AddrFetch GetAddr (no getheaders), +/// one post-verack keepalive ping/pong, and headers-sync stall replace. #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn p2p_timeout_getaddr_and_keepalive_ping() { + use bitcoin::p2p::message::NetworkMessage; use rbitcoin_net::{AddrMan, PeerConnType}; use std::sync::{Arc, Mutex}; use tokio::io::AsyncWriteExt; @@ -114,14 +164,21 @@ async fn p2p_timeout_getaddr_and_keepalive_ping() { let fut = async { let seed_dir = TempDir::new().unwrap(); let peer_dir = TempDir::new().unwrap(); + let dummy_dir = TempDir::new().unwrap(); let seed = start_node(&seed_dir).await; seed.peers.set_peer_timeout_secs(1); let mut book = AddrMan::new(); - book.add(std::net::SocketAddr::from(([1, 2, 3, 4], 8333))); + for i in 0..5_000u32 { + book.add(std::net::SocketAddr::from(( + [(i >> 8) as u8, (i & 0xff) as u8, 1, 1], + 8333, + ))); + } seed.peers.set_addrman(Arc::new(Mutex::new(book))); let mut peer = start_node(&peer_dir).await; + let dummy = start_node(&dummy_dir).await; tokio::time::timeout(Duration::from_secs(5), peer.follow_from(seed.local_addr)) .await .expect("follow_from must return after handshake") @@ -130,6 +187,26 @@ async fn p2p_timeout_getaddr_and_keepalive_ping() { peer.follow_live_count() >= 1, "outbound session must stay live after follow_from" ); + seed.peers + .addconnection(dummy.local_addr, PeerConnType::OutboundFullRelay) + .expect("preferred outbound for stall"); + wait_ms_until( + 3_000, + || { + seed.peers.snapshot().into_iter().any(|p| { + !p.inbound + && p.conn_type == PeerConnType::OutboundFullRelay + && !p.subver.is_empty() + }) + }, + || { + format!( + "seed outbound to dummy must complete (seed={:?})", + seed.peers.snapshot() + ) + }, + ) + .await; tokio::time::sleep(Duration::from_millis(250)).await; let inbound = seed @@ -157,6 +234,95 @@ async fn p2p_timeout_getaddr_and_keepalive_ping() { outbound.pingwait ); + wait_ms_until( + 3_000, + || { + peer.peers.live_peers().into_iter().any(|p| { + !p.inbound && p.handshake_complete() && p.queue_msg(NetworkMessage::GetAddr) + }) + }, + || { + format!( + "follower outbound must take GetAddr (peer={:?})", + peer.peers.snapshot() + ) + }, + ) + .await; + wait_ms_until( + 5_000, + || { + peer.peers.snapshot().into_iter().any(|p| { + !p.inbound && p.bytesrecv_per_msg.get("addrv2").copied().unwrap_or(0) > 0 + }) + }, + || { + format!( + "full-relay GetAddr must return addrv2 (peer={:?})", + peer.peers.snapshot() + ) + }, + ) + .await; + let bind = seed + .peers + .snapshot() + .into_iter() + .find(|p| p.inbound && !p.subver.is_empty()) + .map(|p| p.addrbind) + .unwrap_or(seed.local_addr); + let cached = seed.peers.addr_response_for_bind(bind); + assert_eq!( + cached.len(), + 1000, + "GetAddr must cap at MAX_ADDR_TO_SEND (23% of 5000 is 1150)" + ); + assert_eq!( + cached, + seed.peers.addr_response_for_bind(bind), + "same listen bind must reuse the 24h GetAddr cache" + ); + + let now = seed.peers.now_secs(); + seed.peers.set_mock_now(now + 40 * 60); + wait_ms_until( + 3_000, + || { + !seed + .peers + .snapshot() + .into_iter() + .any(|p| p.inbound && !p.subver.is_empty()) + }, + || { + format!( + "stalling headers-sync inbound must drop when a preferred outbound exists \ + (seed={:?})", + seed.peers.snapshot() + ) + }, + ) + .await; + seed.peers.set_mock_now(0); + + seed.peers + .addconnection(seed.local_addr, PeerConnType::OutboundFullRelay) + .expect("self-connect dial"); + tokio::time::sleep(Duration::from_millis(400)).await; + assert!( + !seed + .peers + .snapshot() + .into_iter() + .any(|p| p.addr == seed.local_addr && !p.subver.is_empty()), + "self-connect must not complete handshake: {:?}", + seed.peers.snapshot() + ); + + let mut one = AddrMan::new(); + one.add(std::net::SocketAddr::from(([1, 2, 3, 4], 8333))); + seed.peers.set_addrman(Arc::new(Mutex::new(one))); + let magic = bitcoin::p2p::Magic::from(bitcoin::Network::Regtest).to_bytes(); let mut raw = tokio::net::TcpStream::connect(seed.local_addr) .await @@ -212,6 +378,8 @@ async fn p2p_timeout_getaddr_and_keepalive_ping() { tokio::time::sleep(Duration::from_millis(20)).await; } + let t = seed.peers.now_secs(); + seed.peers.set_mock_now(t + 24 * 60 * 60 + 1); peer.peers .addconnection(seed.local_addr, PeerConnType::AddrFetch) .expect("addrfetch dial"); @@ -266,6 +434,7 @@ async fn p2p_timeout_getaddr_and_keepalive_ping() { seed.shutdown().await; peer.shutdown().await; + dummy.shutdown().await; }; tokio::time::timeout(Duration::from_secs(20), fut) .await @@ -281,15 +450,22 @@ fn mine_on(node: &P2PNode, height: u32) -> BlockHash { h } -/// Genesis-only follow: first new tip via headers/inv, then HB `cmpctblock`. -/// Coinbase-only compact reconstructs without `getblocktxn`. +/// Mature-pad follow: HB coinbase compact, 2-tx compact → getblocktxn + connect, +/// then orphan child GetData + parent accept (INV of parked child is AlreadyHave). #[tokio::test(flavor = "multi_thread", worker_threads = 4)] -async fn p2p_hb_compact_tip_follow() { +async fn p2p_compact_hb_getblocktxn_and_orphan() { + use bitcoin::p2p::message::NetworkMessage; + use bitcoin::p2p::message_blockdata::Inventory; + use bitcoin::Amount; + use rbitcoin_test::mine::spend_anyone_can_spend; + let fut = async { let seed_dir = TempDir::new().unwrap(); let peer_dir = TempDir::new().unwrap(); - let seed = start_node(&seed_dir).await; - let mut peer = start_node(&peer_dir).await; + let seed = start_padded(&seed_dir).await; + let mut peer = start_padded(&peer_dir).await; + attach_relay_mempool(&peer, &peer_dir); + let pad_h = seed.query.tip_height().expect("pad tip").0; tokio::time::timeout(Duration::from_secs(5), peer.follow_from(seed.local_addr)) .await .expect("follow_from handshake") @@ -299,71 +475,202 @@ async fn p2p_hb_compact_tip_follow() { "outbound follow must stay live" ); - let h1 = mine_on(&seed, 1); - peer.wait_tip_hash(h1, Duration::from_secs(5)) + let h_empty = mine_on(&seed, pad_h + 1); + peer.wait_tip_hash(h_empty, Duration::from_secs(5)) .await .expect("first tip via headers/inv"); - - let hb_deadline = tokio::time::Instant::now() + Duration::from_secs(3); - loop { - let seed_hb_from = seed - .peers - .snapshot() - .into_iter() - .any(|p| p.inbound && p.bip152_hb_from); - if seed_hb_from { - break; - } - if tokio::time::Instant::now() >= hb_deadline { - panic!( + wait_ms_until( + 3_000, + || { + seed.peers + .snapshot() + .into_iter() + .any(|p| p.inbound && p.bip152_hb_from) + }, + || { + format!( "seed inbound must see sendcmpct(1) after first tip \ (seed={:?} peer={:?})", seed.peers.snapshot(), peer.peers.snapshot() - ); - } - tokio::time::sleep(Duration::from_millis(20)).await; - } + ) + }, + ) + .await; - let h2 = mine_on(&seed, 2); - peer.wait_tip_hash(h2, Duration::from_secs(5)) + let cb1 = seed + .query + .reconstruct_block_at_height(Height(1)) + .unwrap() + .txdata[0] + .compute_txid(); + let extra = spend_anyone_can_spend(cb1, 0, Amount::from_sat(49_0000_0000)); + let tip = seed.hub.tip_hash().expect("tip"); + let tip_time = seed.hub.tip_header().expect("tip time").time; + let with_extra = mine_regtest_block(tip, tip_time + 600, pad_h + 2, vec![extra]); + assert_eq!(with_extra.txdata.len(), 2, "coinbase + extra"); + let h_extra = with_extra.block_hash(); + seed.ingest_block(pad_h + 2, with_extra).unwrap(); + wait_ms_until( + 5_000, + || { + seed.peers.snapshot().into_iter().any(|p| { + p.inbound && p.bytesrecv_per_msg.get("getblocktxn").copied().unwrap_or(0) > 0 + }) + }, + || { + format!( + "follower must GetBlockTxn the missing extra tx \ + (seed={:?} peer={:?})", + seed.peers.snapshot(), + peer.peers.snapshot() + ) + }, + ) + .await; + peer.wait_tip_hash(h_extra, Duration::from_secs(5)) .await - .expect("second tip via compact"); - assert_eq!(peer.query.tip_height(), Some(Height(2))); + .expect("2-tx compact via getblocktxn"); + assert_eq!(peer.query.tip_height(), Some(Height(pad_h + 2))); - let peer_out = peer + let cb2 = seed + .query + .reconstruct_block_at_height(Height(2)) + .unwrap() + .txdata[0] + .compute_txid(); + let parent = spend_anyone_can_spend(cb2, 0, Amount::from_sat(49_0000_0000)); + let child = + spend_anyone_can_spend(parent.compute_txid(), 0, Amount::from_sat(48_0000_0000)); + let child_txid = child.compute_txid(); + let child_wtxid = child.compute_wtxid(); + + wait_ms_until( + 3_000, + || { + seed.peers.live_peers().into_iter().any(|p| { + p.inbound + && p.handshake_complete() + && p.queue_msg(NetworkMessage::Tx(child.clone())) + }) + }, + || { + format!( + "seed inbound writer must take the child tx (seed={:?} peer={:?})", + seed.peers.snapshot(), + peer.peers.snapshot() + ) + }, + ) + .await; + wait_ms_until( + 5_000, + || { + let parked = peer.hub.mempool().map(|m| m.orphan_count()).unwrap_or(0); + let getdata = seed + .peers + .snapshot() + .into_iter() + .find(|p| p.inbound) + .map(|p| p.bytesrecv_per_msg.get("getdata").copied().unwrap_or(0)) + .unwrap_or(0); + parked == 1 && getdata > 0 + }, + || { + format!( + "peer must park the child and GetData the parent \ + (seed={:?} peer={:?})", + seed.peers.snapshot(), + peer.peers.snapshot() + ) + }, + ) + .await; + let getdata_parked = seed .peers .snapshot() .into_iter() - .find(|p| !p.inbound) - .expect("peer outbound"); - let cmpct = peer_out - .bytesrecv_per_msg - .get("cmpctblock") - .copied() + .find(|p| p.inbound) + .map(|p| p.bytesrecv_per_msg.get("getdata").copied().unwrap_or(0)) + .unwrap_or(0); + wait_ms_until( + 3_000, + || { + seed.peers.live_peers().into_iter().any(|p| { + p.inbound + && p.handshake_complete() + && p.queue_msg(NetworkMessage::Inv(vec![ + Inventory::WitnessTransaction(child_txid), + Inventory::WTx(child_wtxid), + ])) + }) + }, + || { + format!( + "seed inbound writer must take orphan INV (seed={:?})", + seed.peers.snapshot() + ) + }, + ) + .await; + tokio::time::sleep(Duration::from_millis(200)).await; + let getdata_after_inv = seed + .peers + .snapshot() + .into_iter() + .find(|p| p.inbound) + .map(|p| p.bytesrecv_per_msg.get("getdata").copied().unwrap_or(0)) .unwrap_or(0); - assert!( - cmpct > 0, - "follower must receive cmpctblock for the HB tip: {:?}", - peer_out.bytesrecv_per_msg - ); assert_eq!( - peer_out - .bytessent_per_msg - .get("getblocktxn") - .copied() - .unwrap_or(0), - 0, - "coinbase-only compact must reconstruct without getblocktxn: {:?}", - peer_out.bytessent_per_msg + getdata_after_inv, getdata_parked, + "parked orphan INV must not GetData" ); + wait_ms_until( + 3_000, + || { + seed.peers.live_peers().into_iter().any(|p| { + p.inbound + && p.handshake_complete() + && p.queue_msg(NetworkMessage::Tx(parent.clone())) + }) + }, + || { + format!( + "seed inbound writer must take the parent tx (seed={:?} peer={:?})", + seed.peers.snapshot(), + peer.peers.snapshot() + ) + }, + ) + .await; + wait_ms_until( + 5_000, + || { + let Some(mp) = peer.hub.mempool() else { + return false; + }; + mp.orphan_count() == 0 + && mp.contains(&parent.compute_txid()) + && mp.contains(&child_txid) + }, + || { + format!( + "parent accept must promote the parked child \ + (seed={:?} peer={:?})", + seed.peers.snapshot(), + peer.peers.snapshot() + ) + }, + ) + .await; + seed.shutdown().await; peer.shutdown().await; }; tokio::time::timeout(Duration::from_secs(20), fut) .await - .expect("p2p_hb_compact_tip_follow wall timeout (20s)"); + .expect("p2p_compact_hb_getblocktxn_and_orphan wall timeout (20s)"); } fn attach_relay_mempool(node: &P2PNode, dir: &TempDir) { @@ -380,194 +687,6 @@ fn attach_relay_mempool(node: &P2PNode, dir: &TempDir) { ); } -/// Child with a missing parent over a live follow session: park + GetData parent. -#[tokio::test(flavor = "multi_thread", worker_threads = 4)] -async fn p2p_orphan_child_getdatas_parent() { - use bitcoin::absolute::LockTime; - use bitcoin::p2p::message::NetworkMessage; - use bitcoin::script::ScriptBuf; - use bitcoin::transaction::Version as TxVersion; - use bitcoin::{Amount, OutPoint, Sequence, Transaction, TxIn, TxOut, Witness}; - - let fut = async { - let seed_dir = TempDir::new().unwrap(); - let peer_dir = TempDir::new().unwrap(); - let seed = start_node(&seed_dir).await; - let mut peer = start_node(&peer_dir).await; - attach_relay_mempool(&peer, &peer_dir); - tokio::time::timeout(Duration::from_secs(5), peer.follow_from(seed.local_addr)) - .await - .expect("follow_from handshake") - .expect("follow"); - assert!( - peer.follow_live_count() >= 1, - "outbound follow must stay live" - ); - - let parent_txid = bitcoin::Txid::from_byte_array([0x33; 32]); - let orphan = Transaction { - version: TxVersion::TWO, - lock_time: LockTime::ZERO, - input: vec![TxIn { - previous_output: OutPoint { - txid: parent_txid, - vout: 0, - }, - script_sig: ScriptBuf::new(), - sequence: Sequence::ENABLE_RBF_NO_LOCKTIME, - witness: Witness::new(), - }], - output: vec![TxOut { - value: Amount::from_sat(1000), - script_pubkey: ScriptBuf::from_bytes(vec![0x51]), - }], - }; - - let writer_deadline = tokio::time::Instant::now() + Duration::from_secs(3); - loop { - let queued = seed.peers.live_peers().into_iter().any(|p| { - p.inbound - && p.handshake_complete() - && p.queue_msg(NetworkMessage::Tx(orphan.clone())) - }); - if queued { - break; - } - if tokio::time::Instant::now() >= writer_deadline { - panic!( - "seed inbound writer must take the child tx (seed={:?} peer={:?})", - seed.peers.snapshot(), - peer.peers.snapshot() - ); - } - tokio::time::sleep(Duration::from_millis(20)).await; - } - - let deadline = tokio::time::Instant::now() + Duration::from_secs(5); - loop { - let parked = peer.hub.mempool().map(|m| m.orphan_count()).unwrap_or(0); - let getdata = seed - .peers - .snapshot() - .into_iter() - .find(|p| p.inbound) - .map(|p| p.bytesrecv_per_msg.get("getdata").copied().unwrap_or(0)) - .unwrap_or(0); - if parked == 1 && getdata > 0 { - break; - } - if tokio::time::Instant::now() >= deadline { - panic!( - "peer must park the child and GetData the parent \ - (parked={parked} getdata={getdata} seed={:?} peer={:?})", - seed.peers.snapshot(), - peer.peers.snapshot() - ); - } - tokio::time::sleep(Duration::from_millis(20)).await; - } - - seed.shutdown().await; - peer.shutdown().await; - }; - tokio::time::timeout(Duration::from_secs(20), fut) - .await - .expect("p2p_orphan_child_getdatas_parent wall timeout (20s)"); -} - -/// Compact of a 2-tx tip: coinbase prefilled, extra tx absent from mempool → getblocktxn. -#[tokio::test(flavor = "multi_thread", worker_threads = 4)] -async fn p2p_compact_getblocktxn_missing_extra_tx() { - use bitcoin::bip152::HeaderAndShortIds; - use bitcoin::p2p::message::NetworkMessage; - use bitcoin::p2p::message_compact_blocks::CmpctBlock; - use bitcoin::Amount; - use rbitcoin_test::mine::spend_anyone_can_spend; - - let fut = async { - let seed_dir = TempDir::new().unwrap(); - let peer_dir = TempDir::new().unwrap(); - let seed = start_node(&seed_dir).await; - let mut peer = start_node(&peer_dir).await; - attach_relay_mempool(&peer, &peer_dir); - tokio::time::timeout(Duration::from_secs(5), peer.follow_from(seed.local_addr)) - .await - .expect("follow_from handshake") - .expect("follow"); - assert!( - peer.follow_live_count() >= 1, - "outbound follow must stay live" - ); - - let extra = spend_anyone_can_spend( - bitcoin::Txid::from_byte_array([0x33; 32]), - 0, - Amount::from_sat(1000), - ); - let tip = seed.hub.tip_hash().expect("genesis"); - let tip_time = seed.hub.tip_header().expect("genesis header").time; - let block = mine_regtest_block(tip, tip_time + 600, 1, vec![extra]); - assert_eq!(block.txdata.len(), 2, "coinbase + extra"); - let hsi = HeaderAndShortIds::from_block(&block, 1, 2, &[0]).expect("compact hsi"); - assert_eq!( - hsi.short_ids.len(), - 1, - "coinbase prefilled; extra is a short-id" - ); - - let writer_deadline = tokio::time::Instant::now() + Duration::from_secs(3); - loop { - let queued = seed.peers.live_peers().into_iter().any(|p| { - p.inbound - && p.handshake_complete() - && p.queue_msg(NetworkMessage::CmpctBlock(CmpctBlock { - compact_block: hsi.clone(), - })) - }); - if queued { - break; - } - if tokio::time::Instant::now() >= writer_deadline { - panic!( - "seed inbound writer must take cmpctblock (seed={:?} peer={:?})", - seed.peers.snapshot(), - peer.peers.snapshot() - ); - } - tokio::time::sleep(Duration::from_millis(20)).await; - } - - let deadline = tokio::time::Instant::now() + Duration::from_secs(5); - loop { - let getblocktxn = seed - .peers - .snapshot() - .into_iter() - .find(|p| p.inbound) - .map(|p| p.bytesrecv_per_msg.get("getblocktxn").copied().unwrap_or(0)) - .unwrap_or(0); - if getblocktxn > 0 { - break; - } - if tokio::time::Instant::now() >= deadline { - panic!( - "follower must GetBlockTxn the missing extra tx \ - (seed={:?} peer={:?})", - seed.peers.snapshot(), - peer.peers.snapshot() - ); - } - tokio::time::sleep(Duration::from_millis(20)).await; - } - - seed.shutdown().await; - peer.shutdown().await; - }; - tokio::time::timeout(Duration::from_secs(20), fut) - .await - .expect("p2p_compact_getblocktxn_missing_extra_tx wall timeout (20s)"); -} - /// Outbound feeler: VERSION completes, then the session closes (no live follow). #[tokio::test] async fn p2p_feeler_completes_and_closes() { From e7eba6a7199fb7b2d5d695f59b3413d16d3df00f Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Sat, 12 Sep 2026 20:54:08 -0700 Subject: [PATCH 2/3] docs: catalog merged compact/orphan and GetAddr-cache timeout pad Co-authored-by: Cursor --- TESTING.md | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/TESTING.md b/TESTING.md index 8103e522..cf11b699 100644 --- a/TESTING.md +++ b/TESTING.md @@ -106,7 +106,7 @@ reads it; rustup users export it). Override coverage dir: | Remining 100-block maturity pads with `confirm_wire_run` | `pad_empty_from` / `build_mature_regtest_with_spend` **once per binary journey** (not once per skinny test) | | Wall-time multi-round microbenches in default suite | Deterministic structure / chunk-load asserts; demote wall arms to `#[ignore]` | -**Tier A timeouts:** `two_node_header_and_block_sync` 60s wall (default + job). `p2p_timeout_getaddr_and_keepalive_ping`, `p2p_hb_compact_tip_follow`, `p2p_orphan_child_getdatas_parent`, `p2p_compact_getblocktxn_missing_extra_tx`, `p2p_feeler_completes_and_closes`, and `p2p_inbound_full_rejects_extra` 20s wall (default). Reconstruct / dead-peer are **multinode job only** (`#[ignore]`; job passes `--ignored`). `coverage.sh` also `--skip`s those names plus `two_node`. Heavier topology stays `#[ignore]` (`scripts/integration.sh`). +**Tier A timeouts:** `two_node_header_and_block_sync` 60s wall (default + job). `p2p_timeout_getaddr_and_keepalive_ping`, `p2p_compact_hb_getblocktxn_and_orphan`, `p2p_feeler_completes_and_closes`, and `p2p_inbound_full_rejects_extra` 20s wall (default). Reconstruct / dead-peer are **multinode job only** (`#[ignore]`; job passes `--ignored`). `coverage.sh` also `--skip`s those names plus `two_node`. Heavier topology stays `#[ignore]` (`scripts/integration.sh`). **Speed / reliability (default suite):** prefer `pad_empty_from` / `build_mature_regtest_with_spend` **once per journey** (tx_relay live hub, Electrum protocol, core_analogs assumevalid+mempool) over remine pads; SH run-builder sleeps are 1 ms under `cfg(test)` (40 ms in production). `pin_compose_multi_pack_timed` keeps functional + layout/covered short-circuit gates (multi-ms floor); sticky vs cold assemble is log-only (not a hard timing assert). Schema-13 wire rebuild must stamp create identity from `txid.body` — zero batch identity is treated as missing (regression covered by `reconstruct_and_connect_error_arms` + multi-vout confirm scenarios). Coverage vs speed: prefer **one** scenario at the real entry over N micro-opens that only paint lines; when adding coverage for reduce/materialize, use a **tiny** target, not production stream depth. @@ -236,10 +236,8 @@ Prefer **one high-level scenario** per behavior cluster. Delete lower-level test | `electrum_idle_timeout_disconnects_quiet_client` | Electrum | Idle timeout closes a quiet socket | | `esplora_broadcast_visible_in_rpc_and_electrum` | Node + Electrum + Esplora + RPC | One `run_p2p` datadir: Esplora `POST /tx` parent and mempool child appear in `getrawmempool` and Electrum mempool/history (`fee` on unconfirmed, including child `height = -1`) | | `two_node_header_and_block_sync` | P2P (**default + multinode CI**) | Seeder → peer 8-block IBD; peer `last_write` meter. **Not** re-run under `coverage.sh`. | -| `p2p_timeout_getaddr_and_keepalive_ping` | P2P (**default**) | One pad: v1-magic inbound drops at `peertimeout=1`, AddrFetch `getaddr`/`addrv2` (no `getheaders`), one keepalive ping/pong | -| `p2p_hb_compact_tip_follow` | P2P (**default**) | Genesis follow: first tip via headers/inv, then HB `sendcmpct` + `cmpctblock` reconstruct (coinbase-only, no `getblocktxn`) | -| `p2p_compact_getblocktxn_missing_extra_tx` | P2P (**default**) | Follow session: 2-tx compact (coinbase prefilled, extra not in mempool) → `getblocktxn`. Does **not** pin depth-10 full-block serve | -| `p2p_orphan_child_getdatas_parent` | P2P (**default**) | Follow session: missing-parent child parks; peer `GetData`s the parent. Does **not** pin log-not-reject or tokio-worker lock | +| `p2p_timeout_getaddr_and_keepalive_ping` | P2P (**default**) | One pad: v1-magic inbound drops at `peertimeout=1`, full-relay GetAddr cache 1000, headers-sync stall replace, self-connect refuses, AddrFetch `getaddr`/`addrv2` (no `getheaders`), one keepalive ping/pong | +| `p2p_compact_hb_getblocktxn_and_orphan` | P2P (**default**) | One mature pad: HB coinbase `cmpctblock`, 2-tx compact → `getblocktxn` + connect, orphan child GetData then parent accept (INV AlreadyHave). Does **not** pin depth-10 full-block serve, tokio-worker lock, or park-not-reject logs | | `p2p_feeler_completes_and_closes` | P2P (**default**) | Outbound feeler: VERSION then close (`feeler connection completed`). No live follow; dummy has no completed inbound. Does **not** pin feeler silence timeout | | `p2p_inbound_full_rejects_extra` | P2P (**default**) | `max_inbound=1`: second follow is refused; first inbound stays. Does **not** pin SelectNodeToEvict ranking | | `badprev_orphan_does_not_blacklist_then_reorg_reconstructs` | P2P/chain (default) | Orphan whose prev is not on the tip is held (not `BLOCK_FAILED`); winner branch reconstructs | @@ -259,7 +257,7 @@ Removed (covered by the rows above): `confirm_cross_block_prevout_without_tx_hea ### Integration / multi-node -Default `cargo test` runs `two_node_header_and_block_sync` (8-block), `p2p_timeout_getaddr_and_keepalive_ping`, `p2p_hb_compact_tip_follow`, `p2p_orphan_child_getdatas_parent`, `p2p_compact_getblocktxn_missing_extra_tx`, `p2p_feeler_completes_and_closes`, `p2p_inbound_full_rejects_extra`, and `badprev_orphan_does_not_blacklist_then_reorg_reconstructs`. The required **multinode** job also runs reconstruct and slim dead-peer (`--ignored` filters in `ci.yml`). +Default `cargo test` runs `two_node_header_and_block_sync` (8-block), `p2p_timeout_getaddr_and_keepalive_ping`, `p2p_compact_hb_getblocktxn_and_orphan`, `p2p_feeler_completes_and_closes`, `p2p_inbound_full_rejects_extra`, and `badprev_orphan_does_not_blacklist_then_reorg_reconstructs`. The required **multinode** job also runs reconstruct and slim dead-peer (`--ignored` filters in `ci.yml`). Heavy topology (3-hop, 48-block, mesh, `run_p2p`) stays `#[ignore]` for `scripts/integration.sh`: ```bash From f5572a116fab8f3d3beedcdea949f4ef3e8bb8ef Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Sat, 12 Sep 2026 21:15:40 -0700 Subject: [PATCH 3/3] test: fold remaining Category B 5.1 PeerHub twins into journeys The timeout pad already disconnects a stalling inbound, so drop that heartbeat unit. Fold pre-verack peertimeout math into the V2 timeout log case, and run inbound/outbound/feeler/plain silence timeouts as one test. Sole-preferred stall KEEP stays. Co-authored-by: Cursor --- TESTING.md | 4 +- crates/rbitcoin-net/src/peer_tests.rs | 207 ++++++++++++-------------- crates/rbitcoin-net/src/peers.rs | 70 +-------- 3 files changed, 99 insertions(+), 182 deletions(-) diff --git a/TESTING.md b/TESTING.md index cf11b699..ba04eb2d 100644 --- a/TESTING.md +++ b/TESTING.md @@ -236,9 +236,9 @@ Prefer **one high-level scenario** per behavior cluster. Delete lower-level test | `electrum_idle_timeout_disconnects_quiet_client` | Electrum | Idle timeout closes a quiet socket | | `esplora_broadcast_visible_in_rpc_and_electrum` | Node + Electrum + Esplora + RPC | One `run_p2p` datadir: Esplora `POST /tx` parent and mempool child appear in `getrawmempool` and Electrum mempool/history (`fee` on unconfirmed, including child `height = -1`) | | `two_node_header_and_block_sync` | P2P (**default + multinode CI**) | Seeder → peer 8-block IBD; peer `last_write` meter. **Not** re-run under `coverage.sh`. | -| `p2p_timeout_getaddr_and_keepalive_ping` | P2P (**default**) | One pad: v1-magic inbound drops at `peertimeout=1`, full-relay GetAddr cache 1000, headers-sync stall replace, self-connect refuses, AddrFetch `getaddr`/`addrv2` (no `getheaders`), one keepalive ping/pong | +| `p2p_timeout_getaddr_and_keepalive_ping` | P2P (**default**) | One pad: v1-magic inbound drops at `peertimeout=1`, full-relay GetAddr cache 1000, headers-sync stall replace, self-connect refuses, AddrFetch `getaddr`/`addrv2` (no `getheaders`), one keepalive ping/pong. Sole-preferred stall KEEP stays a PeerHub unit. | | `p2p_compact_hb_getblocktxn_and_orphan` | P2P (**default**) | One mature pad: HB coinbase `cmpctblock`, 2-tx compact → `getblocktxn` + connect, orphan child GetData then parent accept (INV AlreadyHave). Does **not** pin depth-10 full-block serve, tokio-worker lock, or park-not-reject logs | -| `p2p_feeler_completes_and_closes` | P2P (**default**) | Outbound feeler: VERSION then close (`feeler connection completed`). No live follow; dummy has no completed inbound. Does **not** pin feeler silence timeout | +| `p2p_feeler_completes_and_closes` | P2P (**default**) | Outbound feeler: VERSION then close (`feeler connection completed`). No live follow; dummy has no completed inbound. Does **not** pin feeler silence timeout (`handshake_timeout_after_silence`) | | `p2p_inbound_full_rejects_extra` | P2P (**default**) | `max_inbound=1`: second follow is refused; first inbound stays. Does **not** pin SelectNodeToEvict ranking | | `badprev_orphan_does_not_blacklist_then_reorg_reconstructs` | P2P/chain (default) | Orphan whose prev is not on the tip is held (not `BLOCK_FAILED`); winner branch reconstructs | | `serve_after_restart_via_reconstruct` | P2P (**multinode job only**) | Cold serve via reconstruct | diff --git a/crates/rbitcoin-net/src/peer_tests.rs b/crates/rbitcoin-net/src/peer_tests.rs index 5f3d8bf3..09092bea 100644 --- a/crates/rbitcoin-net/src/peer_tests.rs +++ b/crates/rbitcoin-net/src/peer_tests.rs @@ -5212,93 +5212,107 @@ fn pending_header_walk_is_ram_then_one_store_lookup() { } #[tokio::test] -async fn inbound_handshake_timeout_after_silence() { +async fn handshake_timeout_after_silence() { + use std::net::SocketAddr; use tokio::net::{TcpListener, TcpStream}; - let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); - let addr = listener.local_addr().unwrap(); - let _silent = TcpStream::connect(addr).await.unwrap(); - let (stream, peer) = listener.accept().await.unwrap(); - - let handle = tokio::spawn(async move { - connect_and_handshake_timed( - Duration::from_millis(50), - stream, - Magic::REGTEST, - addr, - peer, - 0, - true, - "/rbitcoin:test/", - HandshakePolicy::plain(), - ) - .await - }); - tokio::time::sleep(Duration::from_millis(10)).await; - assert!(!handle.is_finished(), "must still wait during handshake"); - tokio::time::sleep(Duration::from_millis(200)).await; - assert!( - handle.is_finished(), - "silence past the bound must end handshake" - ); - match handle.await.unwrap() { - Err(NetError::Timeout) => {} - Err(e) => panic!("expected Timeout, got {e}"), - Ok(_) => panic!("handshake succeeded on a silent peer"), - } -} - -#[test] -fn inbound_handshake_timeout_is_core_60s() { assert_eq!(HANDSHAKE_TIMEOUT, Duration::from_secs(60)); -} - -#[tokio::test] -async fn outbound_handshake_timeout_after_silence() { - use tokio::net::{TcpListener, TcpStream}; - let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); - let addr = listener.local_addr().unwrap(); - let stream = TcpStream::connect(addr).await.unwrap(); - let _accepted = listener.accept().await.unwrap(); + async fn bind_pair() -> (TcpStream, TcpStream, SocketAddr, SocketAddr) { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let client = TcpStream::connect(addr).await.unwrap(); + let (server, peer) = listener.accept().await.unwrap(); + (client, server, addr, peer) + } - let handle = tokio::spawn(async move { - connect_and_handshake_timed( - Duration::from_millis(50), - stream, - Magic::REGTEST, - addr, - addr, - 0, - false, - "/rbitcoin:test/", - HandshakePolicy::plain(), - ) - .await - }); - tokio::time::sleep(Duration::from_millis(10)).await; - assert!(!handle.is_finished(), "must still wait during handshake"); - tokio::time::sleep(Duration::from_millis(200)).await; - assert!( - handle.is_finished(), - "silence past the bound must end handshake" - ); - match handle.await.unwrap() { - Err(NetError::Timeout) => {} - Err(e) => panic!("expected Timeout, got {e}"), - Ok(_) => panic!("handshake succeeded on a silent peer"), + async fn join_timeout( + handle: tokio::task::JoinHandle>, + still: &str, + done: &str, + succeeded: &str, + ) { + tokio::time::sleep(Duration::from_millis(10)).await; + assert!(!handle.is_finished(), "{still}"); + tokio::time::sleep(Duration::from_millis(200)).await; + assert!(handle.is_finished(), "{done}"); + match handle.await.unwrap() { + Err(NetError::Timeout) => {} + Err(e) => panic!("expected Timeout, got {e}"), + Ok(_) => panic!("{succeeded}"), + } } -} -#[tokio::test] -async fn outbound_regtest_plain_session_timeout_after_silence() { - use tokio::net::{TcpListener, TcpStream}; + let (client, server, addr, peer) = bind_pair().await; + let _silent_in = client; + join_timeout( + tokio::spawn(async move { + connect_and_handshake_timed( + Duration::from_millis(50), + server, + Magic::REGTEST, + addr, + peer, + 0, + true, + "/rbitcoin:test/", + HandshakePolicy::plain(), + ) + .await + }), + "must still wait during inbound handshake", + "silence past the bound must end inbound handshake", + "inbound handshake succeeded on a silent peer", + ) + .await; + + let (client, server, addr, _) = bind_pair().await; + let _silent_out = server; + join_timeout( + tokio::spawn(async move { + connect_and_handshake_timed( + Duration::from_millis(50), + client, + Magic::REGTEST, + addr, + addr, + 0, + false, + "/rbitcoin:test/", + HandshakePolicy::plain(), + ) + .await + }), + "must still wait during outbound handshake", + "silence past the bound must end outbound handshake", + "outbound handshake succeeded on a silent peer", + ) + .await; + + let (client, server, addr, _) = bind_pair().await; + let _silent_feeler = server; + join_timeout( + tokio::spawn(async move { + run_feeler_timed( + Duration::from_millis(50), + client, + Magic::REGTEST, + addr, + addr, + 0, + "/rbitcoin:test/", + ) + .await + }), + "must still wait during feeler", + "silence past the bound must end feeler", + "feeler succeeded on a silent peer", + ) + .await; - let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); - let addr = listener.local_addr().unwrap(); - let stream = TcpStream::connect(addr).await.unwrap(); - let _accepted = listener.accept().await.unwrap(); - match V2PlainSession::outbound_regtest(stream, "/rbitcoin:test/", Duration::from_millis(50)) + let (client, server, _, _) = bind_pair().await; + let _silent_plain = server; + match V2PlainSession::outbound_regtest(client, "/rbitcoin:test/", Duration::from_millis(50)) .await { Err(NetError::Timeout) => {} @@ -5307,41 +5321,6 @@ async fn outbound_regtest_plain_session_timeout_after_silence() { } } -#[tokio::test] -async fn feeler_handshake_timeout_after_silence() { - use tokio::net::{TcpListener, TcpStream}; - - let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); - let addr = listener.local_addr().unwrap(); - let stream = TcpStream::connect(addr).await.unwrap(); - let _accepted = listener.accept().await.unwrap(); - - let handle = tokio::spawn(async move { - run_feeler_timed( - Duration::from_millis(50), - stream, - Magic::REGTEST, - addr, - addr, - 0, - "/rbitcoin:test/", - ) - .await - }); - tokio::time::sleep(Duration::from_millis(10)).await; - assert!(!handle.is_finished(), "must still wait during feeler"); - tokio::time::sleep(Duration::from_millis(200)).await; - assert!( - handle.is_finished(), - "silence past the bound must end feeler" - ); - match handle.await.unwrap() { - Err(NetError::Timeout) => {} - Err(e) => panic!("expected Timeout, got {e}"), - Ok(_) => panic!("feeler succeeded on a silent peer"), - } -} - /// Writer used to `fetch_sub` every `CmpctBlock`, including tip announces that /// never `fetch_add`. That wrapped `serve_inflight` to `usize::MAX` and skipped /// every later reconstruct (`sync_blocks` 60s on long-lived node-to-node). diff --git a/crates/rbitcoin-net/src/peers.rs b/crates/rbitcoin-net/src/peers.rs index 98869c76..a581a4ed 100644 --- a/crates/rbitcoin-net/src/peers.rs +++ b/crates/rbitcoin-net/src/peers.rs @@ -2064,7 +2064,8 @@ mod tests { } #[test] - fn connecting_peer_times_out_at_peertimeout() { + fn connecting_peer_v2_timeout_log_before_transport() { + rbitcoin_log::capture_logs(true); let hub = PeerHub::new(); hub.set_peer_timeout_secs(3); hub.set_mock_now(1_700_000_000); @@ -2072,10 +2073,10 @@ mod tests { let p = hub.register_connecting(a, a, true, PeerConnType::Inbound); assert!(!p.handshake_complete()); hub.set_mock_now(1_700_000_002); - hub.on_session_heartbeat(); assert!(!p.stop.load(Ordering::SeqCst), "still inside peertimeout"); hub.set_mock_now(1_700_000_003); - hub.on_session_heartbeat(); + let logs = rbitcoin_log::take_logs(); + rbitcoin_log::capture_logs(false); assert!( p.stop.load(Ordering::SeqCst), "peertimeout must disconnect pre-verack" @@ -2084,21 +2085,6 @@ mod tests { hub.get(p.id).is_none(), "timed-out connecting peer is dropped" ); - } - - #[test] - fn connecting_peer_v2_timeout_log_before_transport() { - rbitcoin_log::capture_logs(true); - let hub = PeerHub::new(); - hub.set_peer_timeout_secs(3); - hub.set_mock_now(1_700_000_000); - let a = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 1); - let p = hub.register_connecting(a, a, true, PeerConnType::Inbound); - hub.set_mock_now(1_700_000_003); - hub.on_session_heartbeat(); - let logs = rbitcoin_log::take_logs(); - rbitcoin_log::capture_logs(false); - assert!(p.stop.load(Ordering::SeqCst)); assert!( logs.iter() .any(|(_, m)| m.contains("V2 handshake timeout, disconnecting peer=0")), @@ -2106,26 +2092,6 @@ mod tests { ); } - #[test] - fn version_handshake_timeout_log_after_v2_ready() { - rbitcoin_log::capture_logs(true); - let hub = PeerHub::new(); - hub.set_peer_timeout_secs(3); - hub.set_mock_now(1_700_000_000); - let a = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 1); - let p = hub.register_connecting(a, a, true, PeerConnType::Inbound); - p.mark_v2_transport_ready(); - hub.set_mock_now(1_700_000_003); - hub.on_session_heartbeat(); - let logs = rbitcoin_log::take_logs(); - rbitcoin_log::capture_logs(false); - assert!( - logs.iter() - .any(|(_, m)| m.contains("version handshake timeout, disconnecting peer=0")), - "expected version handshake timeout after v2 ready, got {logs:?}" - ); - } - #[test] fn p2p_timeouts_v2_logs_all_three_connecting_peer_ids() { rbitcoin_log::capture_logs(true); @@ -2339,34 +2305,6 @@ mod tests { assert!(hub.try_start_headers_sync(&p2, now, 1_231_006_505)); } - #[test] - fn session_heartbeat_disconnects_stalling_headers_sync() { - let hub = PeerHub::new(); - let a = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 1); - let b = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 2); - let inbound = hub.register(a, a, &ver("/rbitcoin:0.1.0/"), true, PeerConnType::Inbound); - let _outbound = hub.register( - b, - b, - &ver("/rbitcoin:0.1.0/"), - false, - PeerConnType::OutboundFullRelay, - ); - // Deadline is computed from the `now` passed to try_start, not the hub - // clock. Start 16 minutes behind wall time with a tip of the same age - // so variable timeout is 0: deadline = wall − 60s. - let wall = hub.now_secs(); - let start = wall.saturating_sub(16 * 60); - assert!(hub.try_start_headers_sync(&inbound, start, start)); - assert!(inbound.is_sync_started()); - assert!(!inbound.stop.load(Ordering::SeqCst)); - hub.on_session_heartbeat(); - assert!( - inbound.stop.load(Ordering::SeqCst), - "stalling inbound must disconnect when another preferred peer exists" - ); - } - #[test] fn session_heartbeat_keeps_sole_preferred_headers_sync_peer() { let hub = PeerHub::new();