From c00d6b053ae26f96cdb09a07de1045a6a87abda2 Mon Sep 17 00:00:00 2001 From: "rbitcoin-grok[bot]" Date: Sat, 29 Aug 2026 07:35:33 -0700 Subject: [PATCH 1/4] net: prompt disconnectnode tear-down for peer getpeerinfo disconnect_id now clears the writer channel, aborts the writer task (TCP FIN to the peer), and unregisters immediately so getpeerinfo drops the row without waiting for the session loop. Fixes flaky mempool_reorg disconnect_nodes (5s far-side wait). Also treat ConnectionReset/BrokenPipe/Aborted as clean peer close. --- crates/rbitcoin-net/src/peer.rs | 13 ++++++- crates/rbitcoin-net/src/peers.rs | 58 ++++++++++++++++++++++++++------ 2 files changed, 59 insertions(+), 12 deletions(-) diff --git a/crates/rbitcoin-net/src/peer.rs b/crates/rbitcoin-net/src/peer.rs index d19db172..59eca567 100644 --- a/crates/rbitcoin-net/src/peer.rs +++ b/crates/rbitcoin-net/src/peer.rs @@ -669,6 +669,9 @@ pub async fn peer_session_with( } } }); + if let Some(s) = meta.session.as_ref() { + s.set_writer_abort(writer_task.abort_handle()); + } if let Some(s) = meta.session.as_ref() { let _ = maybe_queue_addrfetch_getaddr(&out_tx, s); @@ -937,7 +940,15 @@ pub async fn peer_session_with( frame = read_v2_frame(&mut reader, magic) => { let frame = match frame { Ok(f) => f, - Err(NetError::Io(e)) if e.kind() == std::io::ErrorKind::UnexpectedEof => { + Err(NetError::Io(e)) + if matches!( + e.kind(), + std::io::ErrorKind::UnexpectedEof + | std::io::ErrorKind::ConnectionReset + | std::io::ErrorKind::BrokenPipe + | std::io::ErrorKind::ConnectionAborted + ) => + { return Ok(()); } Err(NetError::MessageTooLarge(n)) => { diff --git a/crates/rbitcoin-net/src/peers.rs b/crates/rbitcoin-net/src/peers.rs index 1d694cf0..5b1135ce 100644 --- a/crates/rbitcoin-net/src/peers.rs +++ b/crates/rbitcoin-net/src/peers.rs @@ -129,6 +129,9 @@ pub struct LivePeer { connected_at: AtomicU64, /// Skip INV for mempool txs with `accept_gen < floor` (post-verack privacy). inv_gen_floor: AtomicU64, + /// Writer-task abort — `disconnectnode` aborts the write half so the peer + /// sees FIN without waiting for our read loop (`mempool_reorg` disconnect_nodes). + writer_abort: Mutex>, } impl LivePeer { @@ -196,6 +199,21 @@ impl LivePeer { self.stop.store(true, Ordering::SeqCst); } + pub fn set_writer_abort(&self, handle: tokio::task::AbortHandle) { + *self.writer_abort.lock().unwrap_or_else(|e| e.into_inner()) = Some(handle); + } + + fn take_writer_abort(&self) -> Option { + self.writer_abort + .lock() + .unwrap_or_else(|e| e.into_inner()) + .take() + } + + fn clear_out_tx(&self) { + *self.out_tx.lock().unwrap_or_else(|e| e.into_inner()) = None; + } + pub fn note_failed_cmpct(&self, hash: BlockHash) { self.failed_cmpct .lock() @@ -1015,6 +1033,7 @@ impl PeerHub { out_tx: Mutex::new(None), connected_at: AtomicU64::new(connected_at), inv_gen_floor: AtomicU64::new(0), + writer_abort: Mutex::new(None), }); // Handshake already exchanged version + verack (+ maybe ping). peer.note_recv("version", 100); @@ -1170,12 +1189,20 @@ impl PeerHub { } pub fn disconnect_id(&self, id: u64) -> bool { - if let Some(p) = self.get(id) { - p.request_disconnect(); - true - } else { - false + let Some(p) = self.get(id) else { + return false; + }; + p.request_disconnect(); + // Drop our writer-channel sender and abort the writer task so TCP FIN + // goes out even if the read loop is between ticks. + p.clear_out_tx(); + if let Some(h) = p.take_writer_abort() { + h.abort(); } + // Drop from getpeerinfo before the session task finishes tearing down + // so disconnect_nodes' far-side wait sees us gone promptly. + self.unregister(id); + true } /// Core `AttemptToEvictConnection`: disconnect one unprotected inbound. @@ -1207,11 +1234,16 @@ impl PeerHub { } pub fn disconnect_addr(&self, addr: SocketAddr) -> bool { - let g = self.live.read().unwrap_or_else(|e| e.into_inner()); + let ids: Vec = { + let g = self.live.read().unwrap_or_else(|e| e.into_inner()); + g.values() + .filter(|p| p.addr == addr) + .map(|p| p.id) + .collect() + }; let mut n = 0usize; - for p in g.values() { - if p.addr == addr { - p.request_disconnect(); + for id in ids { + if self.disconnect_id(id) { n += 1; } } @@ -1306,8 +1338,12 @@ mod tests { assert!(snap[0].bytesrecv_per_msg.get("pong").copied().unwrap() >= 29); assert!(hub.disconnect_id(0)); assert!(p.stop.load(Ordering::SeqCst)); - hub.unregister(0); - assert!(hub.snapshot().is_empty()); + // disconnectnode must clear getpeerinfo immediately (mempool_reorg + // disconnect_nodes waits ≤5s on the far side seeing us gone). + assert!( + hub.snapshot().is_empty(), + "disconnect_id must unregister before the session task exits" + ); } #[test] From de797b385f22dd23b88e65d7749638ba6b503a7b Mon Sep 17 00:00:00 2001 From: "rbitcoin-grok[bot]" Date: Sat, 29 Aug 2026 08:44:30 -0700 Subject: [PATCH 2/4] net: abort whole session on disconnectnode Writer abort alone was not enough for mempool_reorg's far-side getpeerinfo wait. Attach a session AbortHandle and abort it on disconnect_id so reader+writer drop immediately (TCP RST/FIN). --- crates/rbitcoin-net/src/peers.rs | 27 +++++++++++++++++++++------ crates/rbitcoin-net/src/service.rs | 23 ++++++++++++++++++++--- 2 files changed, 41 insertions(+), 9 deletions(-) diff --git a/crates/rbitcoin-net/src/peers.rs b/crates/rbitcoin-net/src/peers.rs index 5b1135ce..f9cdbf8b 100644 --- a/crates/rbitcoin-net/src/peers.rs +++ b/crates/rbitcoin-net/src/peers.rs @@ -129,9 +129,10 @@ pub struct LivePeer { connected_at: AtomicU64, /// Skip INV for mempool txs with `accept_gen < floor` (post-verack privacy). inv_gen_floor: AtomicU64, - /// Writer-task abort — `disconnectnode` aborts the write half so the peer - /// sees FIN without waiting for our read loop (`mempool_reorg` disconnect_nodes). + /// Writer-task abort — FIN via dropping the write half. writer_abort: Mutex>, + /// Whole session-task abort — drops reader+writer if the loop is stuck. + session_abort: Mutex>, } impl LivePeer { @@ -203,6 +204,10 @@ impl LivePeer { *self.writer_abort.lock().unwrap_or_else(|e| e.into_inner()) = Some(handle); } + pub fn set_session_abort(&self, handle: tokio::task::AbortHandle) { + *self.session_abort.lock().unwrap_or_else(|e| e.into_inner()) = Some(handle); + } + fn take_writer_abort(&self) -> Option { self.writer_abort .lock() @@ -210,6 +215,13 @@ impl LivePeer { .take() } + fn take_session_abort(&self) -> Option { + self.session_abort + .lock() + .unwrap_or_else(|e| e.into_inner()) + .take() + } + fn clear_out_tx(&self) { *self.out_tx.lock().unwrap_or_else(|e| e.into_inner()) = None; } @@ -1034,6 +1046,7 @@ impl PeerHub { connected_at: AtomicU64::new(connected_at), inv_gen_floor: AtomicU64::new(0), writer_abort: Mutex::new(None), + session_abort: Mutex::new(None), }); // Handshake already exchanged version + verack (+ maybe ping). peer.note_recv("version", 100); @@ -1193,14 +1206,16 @@ impl PeerHub { return false; }; p.request_disconnect(); - // Drop our writer-channel sender and abort the writer task so TCP FIN - // goes out even if the read loop is between ticks. + // Drop writer channel + abort writer/session so both TCP halves close + // even if the read loop is mid-frame (`mempool_reorg` disconnect_nodes). p.clear_out_tx(); if let Some(h) = p.take_writer_abort() { h.abort(); } - // Drop from getpeerinfo before the session task finishes tearing down - // so disconnect_nodes' far-side wait sees us gone promptly. + if let Some(h) = p.take_session_abort() { + h.abort(); + } + // Drop from getpeerinfo before teardown finishes. self.unregister(id); true } diff --git a/crates/rbitcoin-net/src/service.rs b/crates/rbitcoin-net/src/service.rs index d1a2d1fa..0538a87f 100644 --- a/crates/rbitcoin-net/src/service.rs +++ b/crates/rbitcoin-net/src/service.rs @@ -150,12 +150,19 @@ impl P2PNode { let peers = dial_peers.clone(); let ua = dial_ua.clone(); let live = dial_live.clone(); + let (sess_tx, sess_rx) = tokio::sync::oneshot::channel::>(); let h = tokio::spawn(async move { - let _ = run_outbound_session( - req.addr, magic, local_addr, hub, peers, ua, live, req.typ, + let _ = run_outbound_session_with_sess_hook( + req.addr, magic, local_addr, hub, peers, ua, live, req.typ, sess_tx, ) .await; }); + let ah = h.abort_handle(); + tokio::spawn(async move { + if let Ok(sess) = sess_rx.await { + sess.set_session_abort(ah); + } + }); push_session_task(&sessions_dial, h); } }); @@ -406,6 +413,7 @@ fn spawn_inbound_accept( Err(_) => our, }; let sessions = session_tasks.clone(); + let (sess_tx, sess_rx) = tokio::sync::oneshot::channel::>(); let h = tokio::spawn(async move { let _session_slot = permit; let (ver, reader, writer, wire) = match inbound_connect_and_handshake( @@ -438,6 +446,7 @@ fn spawn_inbound_accept( } sess.attach_wire(wire); let id = sess.id; + let _ = sess_tx.send(Arc::clone(&sess)); let meta = FollowSessionMeta { peer: Some(peer_addr), live: None, @@ -446,6 +455,12 @@ fn spawn_inbound_accept( let _ = peer_session_with(reader, writer, magic, hub, tip_rx, meta).await; peers.unregister(id); }); + let ah = h.abort_handle(); + tokio::spawn(async move { + if let Ok(sess) = sess_rx.await { + sess.set_session_abort(ah); + } + }); push_session_task(&sessions, h); } Ok(Err(_)) => break, @@ -554,7 +569,7 @@ async fn run_prepared_outbound(prepared: PreparedOutbound) -> Result<(), NetErro out } -async fn run_outbound_session( +async fn run_outbound_session_with_sess_hook( peer: SocketAddr, magic: Magic, local: SocketAddr, @@ -563,6 +578,7 @@ async fn run_outbound_session( user_agent: String, follow_live: Arc, typ: PeerConnType, + sess_tx: tokio::sync::oneshot::Sender>, ) -> Result<(), NetError> { if typ == PeerConnType::Feeler { let stream = TcpStream::connect(peer).await?; @@ -572,6 +588,7 @@ async fn run_outbound_session( let prepared = prepare_outbound_session(peer, magic, local, hub, peers, user_agent, follow_live, typ) .await?; + let _ = sess_tx.send(Arc::clone(&prepared.sess)); run_prepared_outbound(prepared).await } From 11a9b2d2537a48b221108c83364e9c7873edf36b Mon Sep 17 00:00:00 2001 From: "rbitcoin-grok[bot]" Date: Sat, 29 Aug 2026 09:01:46 -0700 Subject: [PATCH 3/4] net: wire session AbortHandle without detached spawn ast-grep detached-tokio-spawn rejected the oneshot helper tasks that installed set_session_abort. Pass the AbortHandle into the session task instead so disconnectnode still aborts the whole session. --- crates/rbitcoin-net/src/service.rs | 35 +++++++++++++----------------- 1 file changed, 15 insertions(+), 20 deletions(-) diff --git a/crates/rbitcoin-net/src/service.rs b/crates/rbitcoin-net/src/service.rs index 0538a87f..9dacf2ee 100644 --- a/crates/rbitcoin-net/src/service.rs +++ b/crates/rbitcoin-net/src/service.rs @@ -150,19 +150,14 @@ impl P2PNode { let peers = dial_peers.clone(); let ua = dial_ua.clone(); let live = dial_live.clone(); - let (sess_tx, sess_rx) = tokio::sync::oneshot::channel::>(); + let (ah_tx, ah_rx) = tokio::sync::oneshot::channel::(); let h = tokio::spawn(async move { - let _ = run_outbound_session_with_sess_hook( - req.addr, magic, local_addr, hub, peers, ua, live, req.typ, sess_tx, + let _ = run_outbound_session_with_abort( + req.addr, magic, local_addr, hub, peers, ua, live, req.typ, ah_rx, ) .await; }); - let ah = h.abort_handle(); - tokio::spawn(async move { - if let Ok(sess) = sess_rx.await { - sess.set_session_abort(ah); - } - }); + let _ = ah_tx.send(h.abort_handle()); push_session_task(&sessions_dial, h); } }); @@ -413,7 +408,8 @@ fn spawn_inbound_accept( Err(_) => our, }; let sessions = session_tasks.clone(); - let (sess_tx, sess_rx) = tokio::sync::oneshot::channel::>(); + let (ah_tx, ah_rx) = + tokio::sync::oneshot::channel::(); let h = tokio::spawn(async move { let _session_slot = permit; let (ver, reader, writer, wire) = match inbound_connect_and_handshake( @@ -441,12 +437,14 @@ fn spawn_inbound_accept( }; let sess = peers.register(peer_addr, bind, &ver, true, PeerConnType::Inbound); + if let Ok(ah) = ah_rx.await { + sess.set_session_abort(ah); + } if let Some(mp) = hub.mempool() { sess.set_inv_gen_floor(mp.next_accept_gen()); } sess.attach_wire(wire); let id = sess.id; - let _ = sess_tx.send(Arc::clone(&sess)); let meta = FollowSessionMeta { peer: Some(peer_addr), live: None, @@ -455,12 +453,7 @@ fn spawn_inbound_accept( let _ = peer_session_with(reader, writer, magic, hub, tip_rx, meta).await; peers.unregister(id); }); - let ah = h.abort_handle(); - tokio::spawn(async move { - if let Ok(sess) = sess_rx.await { - sess.set_session_abort(ah); - } - }); + let _ = ah_tx.send(h.abort_handle()); push_session_task(&sessions, h); } Ok(Err(_)) => break, @@ -569,7 +562,7 @@ async fn run_prepared_outbound(prepared: PreparedOutbound) -> Result<(), NetErro out } -async fn run_outbound_session_with_sess_hook( +async fn run_outbound_session_with_abort( peer: SocketAddr, magic: Magic, local: SocketAddr, @@ -578,7 +571,7 @@ async fn run_outbound_session_with_sess_hook( user_agent: String, follow_live: Arc, typ: PeerConnType, - sess_tx: tokio::sync::oneshot::Sender>, + ah_rx: tokio::sync::oneshot::Receiver, ) -> Result<(), NetError> { if typ == PeerConnType::Feeler { let stream = TcpStream::connect(peer).await?; @@ -588,7 +581,9 @@ async fn run_outbound_session_with_sess_hook( let prepared = prepare_outbound_session(peer, magic, local, hub, peers, user_agent, follow_live, typ) .await?; - let _ = sess_tx.send(Arc::clone(&prepared.sess)); + if let Ok(ah) = ah_rx.await { + prepared.sess.set_session_abort(ah); + } run_prepared_outbound(prepared).await } From b0a08835409addfdc1345f5567ffcb06f1823fba Mon Sep 17 00:00:00 2001 From: "rbitcoin-grok[bot]" Date: Sat, 29 Aug 2026 10:04:07 -0700 Subject: [PATCH 4/4] net: shutdown TCP on disconnectnode for far-side EOF mempool_reorg disconnect_nodes still flaked waiting for the far side's getpeerinfo: aborting our session was not enough when the peer was mid-frame. Clone the std TCP fd at open_v2 and Shutdown::Both on disconnect_id so the far read sees EOF inside the 5s wait. --- crates/rbitcoin-net/src/ibd/peer_io.rs | 2 +- crates/rbitcoin-net/src/peer.rs | 39 +++++++++-- crates/rbitcoin-net/src/peer_tests.rs | 94 ++++++++++++++++++++++++++ crates/rbitcoin-net/src/peers.rs | 74 +++++++++++++++++++- crates/rbitcoin-net/src/service.rs | 51 +++++++------- crates/rbitcoin-net/src/v2.rs | 16 ++++- 6 files changed, 241 insertions(+), 35 deletions(-) diff --git a/crates/rbitcoin-net/src/ibd/peer_io.rs b/crates/rbitcoin-net/src/ibd/peer_io.rs index 25c419d8..1193092e 100644 --- a/crates/rbitcoin-net/src/ibd/peer_io.rs +++ b/crates/rbitcoin-net/src/ibd/peer_io.rs @@ -177,7 +177,7 @@ pub(crate) async fn spawn_peer( let stream = TcpStream::connect(addr).await?; let ua = rbitcoin_primitives::rbitcoin_subversion(env!("CARGO_PKG_VERSION"), &[] as &[&str]) .unwrap_or_else(|_| format!("/rbitcoin:{}/", env!("CARGO_PKG_VERSION"))); - let (ver, reader, writer, _wire) = connect_and_handshake_timed( + let (ver, reader, writer, _wire, _tcp_shutdown) = connect_and_handshake_timed( HANDSHAKE_TIMEOUT, stream, magic, diff --git a/crates/rbitcoin-net/src/peer.rs b/crates/rbitcoin-net/src/peer.rs index 59eca567..c1581673 100644 --- a/crates/rbitcoin-net/src/peer.rs +++ b/crates/rbitcoin-net/src/peer.rs @@ -287,8 +287,17 @@ pub async fn connect_and_handshake( inbound: bool, user_agent: &str, policy: HandshakePolicy<'_>, -) -> Result<(VersionMessage, V2Reader, V2Writer, crate::v2::WireBytes), NetError> { - let (mut reader, mut writer, wire) = open_v2(stream, magic, inbound).await?; +) -> Result< + ( + VersionMessage, + V2Reader, + V2Writer, + crate::v2::WireBytes, + std::net::TcpStream, + ), + NetError, +> { + let (mut reader, mut writer, wire, tcp_shutdown) = open_v2(stream, magic, inbound).await?; let their_version = application_handshake( &mut reader, &mut writer, @@ -301,7 +310,7 @@ pub async fn connect_and_handshake( policy, ) .await?; - Ok((their_version, reader, writer, wire)) + Ok((their_version, reader, writer, wire, tcp_shutdown)) } /// Core VERSION/VERACK bound: 60s from TCP connect/accept. Timeout drops the stream. @@ -315,7 +324,16 @@ pub(crate) async fn inbound_connect_and_handshake( start_height: i32, user_agent: &str, policy: HandshakePolicy<'_>, -) -> Result<(VersionMessage, V2Reader, V2Writer, crate::v2::WireBytes), NetError> { +) -> Result< + ( + VersionMessage, + V2Reader, + V2Writer, + crate::v2::WireBytes, + std::net::TcpStream, + ), + NetError, +> { connect_and_handshake_timed( HANDSHAKE_TIMEOUT, stream, @@ -340,7 +358,16 @@ pub(crate) async fn connect_and_handshake_timed( inbound: bool, user_agent: &str, policy: HandshakePolicy<'_>, -) -> Result<(VersionMessage, V2Reader, V2Writer, crate::v2::WireBytes), NetError> { +) -> Result< + ( + VersionMessage, + V2Reader, + V2Writer, + crate::v2::WireBytes, + std::net::TcpStream, + ), + NetError, +> { tokio::time::timeout( limit, connect_and_handshake( @@ -411,7 +438,7 @@ async fn run_feeler_inner( start_height: i32, user_agent: &str, ) -> Result<(), NetError> { - let (mut reader, mut writer, _wire) = open_v2(stream, magic, false).await?; + let (mut reader, mut writer, _wire, _tcp_shutdown) = open_v2(stream, magic, false).await?; let services = local_service_flags(); let now = SystemTime::now() .duration_since(UNIX_EPOCH) diff --git a/crates/rbitcoin-net/src/peer_tests.rs b/crates/rbitcoin-net/src/peer_tests.rs index 3feb215e..b203e338 100644 --- a/crates/rbitcoin-net/src/peer_tests.rs +++ b/crates/rbitcoin-net/src/peer_tests.rs @@ -5573,3 +5573,97 @@ async fn tip_burst_past_broadcast_capacity_still_syncs_peer() { nb.shutdown().await; let _ = std::fs::remove_dir_all(dir); } + +/// Core `disconnect_nodes` waits ≤5s for the far side's `getpeerinfo` to drop us. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn disconnect_clears_far_side_getpeerinfo_within_5s() { + use crate::P2PNode; + use std::time::Duration; + + let n = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_nanos(); + let dir = std::env::temp_dir().join(format!("rbitcoin-disc-far-{n}")); + std::fs::create_dir_all(dir.join("a")).unwrap(); + std::fs::create_dir_all(dir.join("b")).unwrap(); + let qa = Query::open_or_create(dir.join("a/store")).unwrap(); + let qb = Query::open_or_create(dir.join("b/store")).unwrap(); + let params = ChainParams::regtest(); + let mut na = P2PNode::start_with_agent( + "127.0.0.1:0".parse().unwrap(), + qa, + params.clone(), + Milestone::NONE, + "/rbitcoin:0.1.0(testnode0)/".into(), + crate::DEFAULT_MAX_INBOUND, + ) + .await + .unwrap(); + let nb = P2PNode::start_with_agent( + "127.0.0.1:0".parse().unwrap(), + qb, + params, + Milestone::NONE, + "/rbitcoin:0.1.0(testnode1)/".into(), + crate::DEFAULT_MAX_INBOUND, + ) + .await + .unwrap(); + + na.follow_from(nb.local_addr).await.unwrap(); + let mut linked = false; + for _ in 0..100 { + let a_sees = na + .peers + .snapshot() + .iter() + .any(|p| p.subver.contains("testnode1")); + let b_sees = nb + .peers + .snapshot() + .iter() + .any(|p| p.subver.contains("testnode0")); + if a_sees && b_sees { + linked = true; + break; + } + tokio::time::sleep(Duration::from_millis(20)).await; + } + assert!(linked, "both sides must list each other before disconnect"); + + let peer_id = na + .peers + .snapshot() + .into_iter() + .find(|p| p.subver.contains("testnode1")) + .map(|p| p.id) + .expect("outbound peer id"); + assert!(na.peers.disconnect_id(peer_id)); + assert!( + na.peers.snapshot().is_empty(), + "local getpeerinfo clears immediately" + ); + + let mut far_clear = false; + for _ in 0..100 { + if !nb + .peers + .snapshot() + .iter() + .any(|p| p.subver.contains("testnode0")) + { + far_clear = true; + break; + } + tokio::time::sleep(Duration::from_millis(50)).await; + } + assert!( + far_clear, + "far side getpeerinfo must drop us within 5s (Core disconnect_nodes)" + ); + + na.shutdown().await; + nb.shutdown().await; + let _ = std::fs::remove_dir_all(dir); +} diff --git a/crates/rbitcoin-net/src/peers.rs b/crates/rbitcoin-net/src/peers.rs index f9cdbf8b..e7aeecd3 100644 --- a/crates/rbitcoin-net/src/peers.rs +++ b/crates/rbitcoin-net/src/peers.rs @@ -133,6 +133,9 @@ pub struct LivePeer { writer_abort: Mutex>, /// Whole session-task abort — drops reader+writer if the loop is stuck. session_abort: Mutex>, + /// Cloned std TCP fd for `Shutdown::Both` on `disconnectnode` so the far + /// side sees EOF even if our session task is mid-frame. + tcp_shutdown: Mutex>, } impl LivePeer { @@ -208,6 +211,10 @@ impl LivePeer { *self.session_abort.lock().unwrap_or_else(|e| e.into_inner()) = Some(handle); } + pub fn attach_tcp_shutdown(&self, stream: std::net::TcpStream) { + *self.tcp_shutdown.lock().unwrap_or_else(|e| e.into_inner()) = Some(stream); + } + fn take_writer_abort(&self) -> Option { self.writer_abort .lock() @@ -222,6 +229,13 @@ impl LivePeer { .take() } + fn take_tcp_shutdown(&self) -> Option { + self.tcp_shutdown + .lock() + .unwrap_or_else(|e| e.into_inner()) + .take() + } + fn clear_out_tx(&self) { *self.out_tx.lock().unwrap_or_else(|e| e.into_inner()) = None; } @@ -1047,6 +1061,7 @@ impl PeerHub { inv_gen_floor: AtomicU64::new(0), writer_abort: Mutex::new(None), session_abort: Mutex::new(None), + tcp_shutdown: Mutex::new(None), }); // Handshake already exchanged version + verack (+ maybe ping). peer.note_recv("version", 100); @@ -1206,8 +1221,13 @@ impl PeerHub { return false; }; p.request_disconnect(); - // Drop writer channel + abort writer/session so both TCP halves close - // even if the read loop is mid-frame (`mempool_reorg` disconnect_nodes). + // Hard-close TCP first so the far side's read sees EOF inside the + // Core `disconnect_nodes` 5s wait (`mempool_reorg`). + if let Some(s) = p.take_tcp_shutdown() { + let _ = s.shutdown(std::net::Shutdown::Both); + } + // Drop writer channel + abort writer/session so local halves tear down + // even if the read loop is mid-frame. p.clear_out_tx(); if let Some(h) = p.take_writer_abort() { h.abort(); @@ -1361,6 +1381,56 @@ mod tests { ); } + #[test] + fn disconnect_id_shuts_down_tcp_so_far_side_sees_eof() { + use std::io::{Read, Write}; + use std::net::TcpListener; + use std::time::{Duration, Instant}; + + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let addr = listener.local_addr().unwrap(); + let mut local = std::net::TcpStream::connect(addr).unwrap(); + let (mut far, _) = listener.accept().unwrap(); + far.set_read_timeout(Some(Duration::from_millis(200))) + .unwrap(); + + let hub = PeerHub::new(); + let a = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 18444); + let b = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 18445); + let p = hub.register( + a, + b, + &ver("/rbitcoin:0.1.0(testnode0)/"), + false, + PeerConnType::OutboundFullRelay, + ); + let killer = local.try_clone().unwrap(); + p.attach_tcp_shutdown(killer); + + assert!(hub.disconnect_id(0)); + + let start = Instant::now(); + let mut buf = [0u8; 1]; + let n = far.read(&mut buf); + assert!( + start.elapsed() < Duration::from_secs(1), + "far side must observe close promptly, took {:?}", + start.elapsed() + ); + match n { + Ok(0) => {} + Err(e) + if matches!( + e.kind(), + std::io::ErrorKind::ConnectionReset + | std::io::ErrorKind::ConnectionAborted + | std::io::ErrorKind::UnexpectedEof + ) => {} + other => panic!("expected EOF/reset after disconnect_id, got {other:?}"), + } + let _ = local.write(&[1]); + } + #[test] fn outbound_nonce_detects_self_connect() { let hub = PeerHub::new(); diff --git a/crates/rbitcoin-net/src/service.rs b/crates/rbitcoin-net/src/service.rs index 9dacf2ee..0886cbaf 100644 --- a/crates/rbitcoin-net/src/service.rs +++ b/crates/rbitcoin-net/src/service.rs @@ -412,31 +412,33 @@ fn spawn_inbound_accept( tokio::sync::oneshot::channel::(); let h = tokio::spawn(async move { let _session_slot = permit; - let (ver, reader, writer, wire) = match inbound_connect_and_handshake( - stream, - magic, - our, - peer_addr, - height, - &ua, - HandshakePolicy { - hub: Some(hub.as_ref()), - peers: Some(peers.as_ref()), - conn_type: PeerConnType::Inbound, - }, - ) - .await - { - Ok(x) => x, - Err(e) => { - rbitcoin_log::debug!( - "p2p: inbound handshake {peer_addr} failed: {e}" - ); - return; - } - }; + let (ver, reader, writer, wire, tcp_shutdown) = + match inbound_connect_and_handshake( + stream, + magic, + our, + peer_addr, + height, + &ua, + HandshakePolicy { + hub: Some(hub.as_ref()), + peers: Some(peers.as_ref()), + conn_type: PeerConnType::Inbound, + }, + ) + .await + { + Ok(x) => x, + Err(e) => { + rbitcoin_log::debug!( + "p2p: inbound handshake {peer_addr} failed: {e}" + ); + return; + } + }; let sess = peers.register(peer_addr, bind, &ver, true, PeerConnType::Inbound); + sess.attach_tcp_shutdown(tcp_shutdown); if let Ok(ah) = ah_rx.await { sess.set_session_abort(ah); } @@ -514,7 +516,7 @@ async fn prepare_outbound_session( }, ) .await; - let (ver, reader, writer, wire) = match handshake { + let (ver, reader, writer, wire, tcp_shutdown) = match handshake { Ok(x) => x, Err(e) => { peers.unregister(provisional_id); @@ -523,6 +525,7 @@ async fn prepare_outbound_session( }; peers.unregister(provisional_id); let sess = peers.register_with_id(provisional_id, peer, bind, &ver, false, typ); + sess.attach_tcp_shutdown(tcp_shutdown); if let Some(mp) = hub.mempool() { sess.set_inv_gen_floor(mp.next_accept_gen()); } diff --git a/crates/rbitcoin-net/src/v2.rs b/crates/rbitcoin-net/src/v2.rs index 0dae93e9..e9e94d6b 100644 --- a/crates/rbitcoin-net/src/v2.rs +++ b/crates/rbitcoin-net/src/v2.rs @@ -448,14 +448,21 @@ fn map_protocol_error(e: ProtocolError) -> NetError { /// Complete BIP324 handshake on a connected TCP stream; return split encrypted halves. /// +/// The fourth value is a cloned std TCP handle for [`std::net::TcpStream::shutdown`] +/// on `disconnectnode` (far-side EOF without waiting on our session task). +/// /// Not cancellation-safe (BIP324 handshake). Callers should not wrap this in /// `select!` without a dedicated task. pub async fn open_v2( stream: TcpStream, magic: Magic, inbound: bool, -) -> Result<(V2Reader, V2Writer, WireBytes), NetError> { +) -> Result<(V2Reader, V2Writer, WireBytes, std::net::TcpStream), NetError> { let _ = stream.set_nodelay(true); + let std = stream.into_std().map_err(NetError::Io)?; + std.set_nonblocking(true).map_err(NetError::Io)?; + let tcp_shutdown = std.try_clone().map_err(NetError::Io)?; + let stream = TcpStream::from_std(std).map_err(NetError::Io)?; let role = if inbound { Role::Responder } else { @@ -477,7 +484,12 @@ pub async fn open_v2( .await .map_err(map_protocol_error)?; let (r, w) = protocol.into_split(); - Ok((V2SessionReader::from_protocol_reader(r), w, wire)) + Ok(( + V2SessionReader::from_protocol_reader(r), + w, + wire, + tcp_shutdown, + )) } /// Encrypt and send one application message.