From 0167daf060ef0dbde5251f27340450635011171f Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 08:19:11 -0700 Subject: [PATCH 01/10] net: SOCKS5 IPv4 CONNECT through a fake proxy Hand-rolled client so P2P outbound can tunnel via system Tor. Test records ATYP=1 and echoes a byte through the tunnel. No live Tor. Co-authored-by: Cursor --- crates/rbitcoin-net/src/lib.rs | 1 + crates/rbitcoin-net/src/socks.rs | 130 +++++++++++++++++++++++++++++++ 2 files changed, 131 insertions(+) create mode 100644 crates/rbitcoin-net/src/socks.rs diff --git a/crates/rbitcoin-net/src/lib.rs b/crates/rbitcoin-net/src/lib.rs index 3003860c8..cf5bbf9eb 100644 --- a/crates/rbitcoin-net/src/lib.rs +++ b/crates/rbitcoin-net/src/lib.rs @@ -19,6 +19,7 @@ mod reactor; mod seeds; mod serve_perf; mod service; +mod socks; mod tip_accept; mod tx_relay; mod v2; diff --git a/crates/rbitcoin-net/src/socks.rs b/crates/rbitcoin-net/src/socks.rs new file mode 100644 index 000000000..030af9dc8 --- /dev/null +++ b/crates/rbitcoin-net/src/socks.rs @@ -0,0 +1,130 @@ +//! SOCKS5 CONNECT client for P2P outbound (system Tor / generic proxy). + +use crate::error::NetError; +use std::net::{Ipv4Addr, SocketAddr}; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::TcpStream; + +pub(crate) async fn socks5_connect( + proxy: SocketAddr, + target: SocketAddr, + creds: Option<&[u8]>, +) -> Result { + let mut s = TcpStream::connect(proxy).await?; + greet(&mut s, creds).await?; + connect_ipv4(&mut s, target).await?; + Ok(s) +} + +async fn greet(s: &mut TcpStream, creds: Option<&[u8]>) -> Result<(), NetError> { + if creds.is_some() { + return Err(NetError::Protocol( + "socks username/password not implemented", + )); + } + s.write_all(&[5, 1, 0x00]).await?; + let mut sel = [0u8; 2]; + s.read_exact(&mut sel).await?; + if sel[0] != 5 || sel[1] != 0x00 { + return Err(NetError::Protocol("socks method rejected")); + } + Ok(()) +} + +async fn connect_ipv4(s: &mut TcpStream, target: SocketAddr) -> Result<(), NetError> { + let SocketAddr::V4(v4) = target else { + return Err(NetError::Protocol("socks CONNECT needs IPv4")); + }; + let ip: Ipv4Addr = *v4.ip(); + let port = v4.port().to_be_bytes(); + let mut req = [0u8; 10]; + req[0] = 5; + req[1] = 1; + req[3] = 1; + req[4..8].copy_from_slice(&ip.octets()); + req[8..10].copy_from_slice(&port); + s.write_all(&req).await?; + read_connect_reply(s).await +} + +async fn read_connect_reply(s: &mut TcpStream) -> Result<(), NetError> { + let mut hdr = [0u8; 4]; + s.read_exact(&mut hdr).await?; + if hdr[0] != 5 { + return Err(NetError::Protocol("socks reply version")); + } + if hdr[1] != 0 { + return Err(NetError::Protocol("socks CONNECT refused")); + } + match hdr[3] { + 1 => { + let mut rest = [0u8; 6]; + s.read_exact(&mut rest).await?; + } + 4 => { + let mut rest = [0u8; 18]; + s.read_exact(&mut rest).await?; + } + 3 => { + let mut n = [0u8; 1]; + s.read_exact(&mut n).await?; + let mut rest = vec![0u8; n[0] as usize + 2]; + s.read_exact(&mut rest).await?; + } + _ => return Err(NetError::Protocol("socks reply ATYP")), + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::socks5_connect; + use std::net::{Ipv4Addr, SocketAddr}; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::net::TcpListener; + + #[tokio::test] + async fn socks5_connect_ipv4_against_fake_proxy() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let proxy = listener.local_addr().unwrap(); + let target = SocketAddr::from((Ipv4Addr::new(203, 0, 113, 7), 8333)); + + let server = tokio::spawn(async move { + let (mut s, _) = listener.accept().await.unwrap(); + let mut ver_n = [0u8; 2]; + s.read_exact(&mut ver_n).await.unwrap(); + assert_eq!(ver_n[0], 5, "SOCKS version"); + let nmethods = ver_n[1] as usize; + let mut methods = vec![0u8; nmethods]; + s.read_exact(&mut methods).await.unwrap(); + assert!(methods.contains(&0x00), "NOAUTH offered, got {methods:?}"); + s.write_all(&[5, 0x00]).await.unwrap(); + + let mut hdr = [0u8; 4]; + s.read_exact(&mut hdr).await.unwrap(); + assert_eq!(hdr[0], 5); + assert_eq!(hdr[1], 1, "CONNECT"); + assert_eq!(hdr[2], 0); + assert_eq!(hdr[3], 1, "ATYP IPv4"); + let mut addr = [0u8; 4]; + s.read_exact(&mut addr).await.unwrap(); + let mut port = [0u8; 2]; + s.read_exact(&mut port).await.unwrap(); + assert_eq!(addr, [203, 0, 113, 7]); + assert_eq!(u16::from_be_bytes(port), 8333); + + s.write_all(&[5, 0, 0, 1, 0, 0, 0, 0, 0, 0]).await.unwrap(); + + let mut b = [0u8; 1]; + s.read_exact(&mut b).await.unwrap(); + s.write_all(&b).await.unwrap(); + }); + + let mut stream = socks5_connect(proxy, target, None).await.unwrap(); + stream.write_all(&[0xab]).await.unwrap(); + let mut echo = [0u8; 1]; + stream.read_exact(&mut echo).await.unwrap(); + assert_eq!(echo, [0xab]); + server.await.unwrap(); + } +} From fc3a0f1b78de47f2377f08c93db806a4ae6c5427 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 08:20:57 -0700 Subject: [PATCH 02/10] net: SOCKS5 domain CONNECT without local DNS ATYP=3 sends the hostname bytes so seed lookup can use remote DNS through the proxy. One CONNECT encoder covers IPv4, IPv6, and domain. Co-authored-by: Cursor --- crates/rbitcoin-net/src/socks.rs | 108 ++++++++++++++++++++++++++----- 1 file changed, 92 insertions(+), 16 deletions(-) diff --git a/crates/rbitcoin-net/src/socks.rs b/crates/rbitcoin-net/src/socks.rs index 030af9dc8..bd280b9a9 100644 --- a/crates/rbitcoin-net/src/socks.rs +++ b/crates/rbitcoin-net/src/socks.rs @@ -1,18 +1,41 @@ //! SOCKS5 CONNECT client for P2P outbound (system Tor / generic proxy). use crate::error::NetError; -use std::net::{Ipv4Addr, SocketAddr}; +use std::net::SocketAddr; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpStream; +enum SocksDest<'a> { + Socket(SocketAddr), + Domain { host: &'a str, port: u16 }, +} + pub(crate) async fn socks5_connect( proxy: SocketAddr, target: SocketAddr, creds: Option<&[u8]>, +) -> Result { + socks5_connect_dest(proxy, SocksDest::Socket(target), creds).await +} + +pub(crate) async fn socks5_connect_domain( + proxy: SocketAddr, + host: &str, + port: u16, + creds: Option<&[u8]>, +) -> Result { + socks5_connect_dest(proxy, SocksDest::Domain { host, port }, creds).await +} + +async fn socks5_connect_dest( + proxy: SocketAddr, + dest: SocksDest<'_>, + creds: Option<&[u8]>, ) -> Result { let mut s = TcpStream::connect(proxy).await?; greet(&mut s, creds).await?; - connect_ipv4(&mut s, target).await?; + write_connect(&mut s, dest).await?; + read_connect_reply(&mut s).await?; Ok(s) } @@ -31,20 +54,33 @@ async fn greet(s: &mut TcpStream, creds: Option<&[u8]>) -> Result<(), NetError> Ok(()) } -async fn connect_ipv4(s: &mut TcpStream, target: SocketAddr) -> Result<(), NetError> { - let SocketAddr::V4(v4) = target else { - return Err(NetError::Protocol("socks CONNECT needs IPv4")); - }; - let ip: Ipv4Addr = *v4.ip(); - let port = v4.port().to_be_bytes(); - let mut req = [0u8; 10]; - req[0] = 5; - req[1] = 1; - req[3] = 1; - req[4..8].copy_from_slice(&ip.octets()); - req[8..10].copy_from_slice(&port); +async fn write_connect(s: &mut TcpStream, dest: SocksDest<'_>) -> Result<(), NetError> { + let mut req = Vec::with_capacity(22); + req.extend_from_slice(&[5, 1, 0]); + match dest { + SocksDest::Socket(SocketAddr::V4(v4)) => { + req.push(1); + req.extend_from_slice(&v4.ip().octets()); + req.extend_from_slice(&v4.port().to_be_bytes()); + } + SocksDest::Socket(SocketAddr::V6(v6)) => { + req.push(4); + req.extend_from_slice(&v6.ip().octets()); + req.extend_from_slice(&v6.port().to_be_bytes()); + } + SocksDest::Domain { host, port } => { + let bytes = host.as_bytes(); + if bytes.is_empty() || bytes.len() > 255 { + return Err(NetError::Protocol("socks domain length")); + } + req.push(3); + req.push(bytes.len() as u8); + req.extend_from_slice(bytes); + req.extend_from_slice(&port.to_be_bytes()); + } + } s.write_all(&req).await?; - read_connect_reply(s).await + Ok(()) } async fn read_connect_reply(s: &mut TcpStream) -> Result<(), NetError> { @@ -78,7 +114,7 @@ async fn read_connect_reply(s: &mut TcpStream) -> Result<(), NetError> { #[cfg(test)] mod tests { - use super::socks5_connect; + use super::{socks5_connect, socks5_connect_domain}; use std::net::{Ipv4Addr, SocketAddr}; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpListener; @@ -127,4 +163,44 @@ mod tests { assert_eq!(echo, [0xab]); server.await.unwrap(); } + + #[tokio::test] + async fn socks5_connect_domain_does_not_resolve_locally() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let proxy = listener.local_addr().unwrap(); + let host = "seed.example"; + let port = 8333u16; + + let server = tokio::spawn(async move { + let (mut s, _) = listener.accept().await.unwrap(); + let mut ver_n = [0u8; 2]; + s.read_exact(&mut ver_n).await.unwrap(); + let nmethods = ver_n[1] as usize; + let mut methods = vec![0u8; nmethods]; + s.read_exact(&mut methods).await.unwrap(); + s.write_all(&[5, 0x00]).await.unwrap(); + + let mut hdr = [0u8; 4]; + s.read_exact(&mut hdr).await.unwrap(); + assert_eq!(hdr[0], 5); + assert_eq!(hdr[1], 1, "CONNECT"); + assert_eq!(hdr[2], 0); + assert_eq!(hdr[3], 3, "ATYP domain; must not resolve locally"); + let mut n = [0u8; 1]; + s.read_exact(&mut n).await.unwrap(); + let mut name = vec![0u8; n[0] as usize]; + s.read_exact(&mut name).await.unwrap(); + let mut p = [0u8; 2]; + s.read_exact(&mut p).await.unwrap(); + assert_eq!(name, b"seed.example"); + assert_eq!(u16::from_be_bytes(p), 8333); + + s.write_all(&[5, 0, 0, 1, 0, 0, 0, 0, 0, 0]).await.unwrap(); + }); + + socks5_connect_domain(proxy, host, port, None) + .await + .unwrap(); + server.await.unwrap(); + } } From 7a8b87da024fb9435de89de805aa00053d8ee2af Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 08:25:04 -0700 Subject: [PATCH 03/10] net: SOCKS username/password isolation creds ProxyCreds::fresh and dial_isolated mint a new RFC 1929 username per call so Tor can isolate circuits. Fake proxy records the bytes. Co-authored-by: Cursor --- Cargo.lock | 1 + crates/rbitcoin-net/Cargo.toml | 1 + crates/rbitcoin-net/src/socks.rs | 163 ++++++++++++++++++++++++++++--- 3 files changed, 149 insertions(+), 16 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 602abab81..e0ef107ad 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -803,6 +803,7 @@ dependencies = [ "arc-swap", "bip324", "bitcoin", + "getrandom 0.4.3", "libc", "rbitcoin-consensus", "rbitcoin-log", diff --git a/crates/rbitcoin-net/Cargo.toml b/crates/rbitcoin-net/Cargo.toml index c8e87c05d..5b8d6e432 100644 --- a/crates/rbitcoin-net/Cargo.toml +++ b/crates/rbitcoin-net/Cargo.toml @@ -18,6 +18,7 @@ bitcoin = { workspace = true } bip324 = { workspace = true } tokio = { workspace = true } arc-swap = { workspace = true } +getrandom = "0.4" [target.'cfg(target_os = "linux")'.dependencies] # POLLRDHUP: peer FIN with unread bytes still in the receive buffer. diff --git a/crates/rbitcoin-net/src/socks.rs b/crates/rbitcoin-net/src/socks.rs index bd280b9a9..4c163d6cc 100644 --- a/crates/rbitcoin-net/src/socks.rs +++ b/crates/rbitcoin-net/src/socks.rs @@ -5,6 +5,21 @@ use std::net::SocketAddr; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpStream; +pub(crate) struct ProxyCreds { + pub(crate) username: Vec, + pub(crate) password: Vec, +} + +impl ProxyCreds { + pub(crate) fn fresh() -> Self { + let mut username = vec![0u8; 16]; + let mut password = vec![0u8; 16]; + getrandom::fill(&mut username).expect("CSPRNG for SOCKS creds"); + getrandom::fill(&mut password).expect("CSPRNG for SOCKS creds"); + Self { username, password } + } +} + enum SocksDest<'a> { Socket(SocketAddr), Domain { host: &'a str, port: u16 }, @@ -13,7 +28,7 @@ enum SocksDest<'a> { pub(crate) async fn socks5_connect( proxy: SocketAddr, target: SocketAddr, - creds: Option<&[u8]>, + creds: Option<&ProxyCreds>, ) -> Result { socks5_connect_dest(proxy, SocksDest::Socket(target), creds).await } @@ -22,15 +37,23 @@ pub(crate) async fn socks5_connect_domain( proxy: SocketAddr, host: &str, port: u16, - creds: Option<&[u8]>, + creds: Option<&ProxyCreds>, ) -> Result { socks5_connect_dest(proxy, SocksDest::Domain { host, port }, creds).await } +pub(crate) async fn dial_isolated( + proxy: SocketAddr, + target: SocketAddr, +) -> Result { + let creds = ProxyCreds::fresh(); + socks5_connect(proxy, target, Some(&creds)).await +} + async fn socks5_connect_dest( proxy: SocketAddr, dest: SocksDest<'_>, - creds: Option<&[u8]>, + creds: Option<&ProxyCreds>, ) -> Result { let mut s = TcpStream::connect(proxy).await?; greet(&mut s, creds).await?; @@ -39,19 +62,46 @@ async fn socks5_connect_dest( Ok(s) } -async fn greet(s: &mut TcpStream, creds: Option<&[u8]>) -> Result<(), NetError> { - if creds.is_some() { - return Err(NetError::Protocol( - "socks username/password not implemented", - )); - } - s.write_all(&[5, 1, 0x00]).await?; - let mut sel = [0u8; 2]; - s.read_exact(&mut sel).await?; - if sel[0] != 5 || sel[1] != 0x00 { - return Err(NetError::Protocol("socks method rejected")); +async fn greet(s: &mut TcpStream, creds: Option<&ProxyCreds>) -> Result<(), NetError> { + match creds { + None => { + s.write_all(&[5, 1, 0x00]).await?; + let mut sel = [0u8; 2]; + s.read_exact(&mut sel).await?; + if sel[0] != 5 || sel[1] != 0x00 { + return Err(NetError::Protocol("socks method rejected")); + } + Ok(()) + } + Some(c) => { + if c.username.is_empty() + || c.username.len() > 255 + || c.password.is_empty() + || c.password.len() > 255 + { + return Err(NetError::Protocol("socks username/password length")); + } + s.write_all(&[5, 1, 0x02]).await?; + let mut sel = [0u8; 2]; + s.read_exact(&mut sel).await?; + if sel[0] != 5 || sel[1] != 0x02 { + return Err(NetError::Protocol("socks method rejected")); + } + let mut auth = Vec::with_capacity(3 + c.username.len() + c.password.len()); + auth.push(1); + auth.push(c.username.len() as u8); + auth.extend_from_slice(&c.username); + auth.push(c.password.len() as u8); + auth.extend_from_slice(&c.password); + s.write_all(&auth).await?; + let mut st = [0u8; 2]; + s.read_exact(&mut st).await?; + if st[0] != 1 || st[1] != 0 { + return Err(NetError::Protocol("socks username/password rejected")); + } + Ok(()) + } } - Ok(()) } async fn write_connect(s: &mut TcpStream, dest: SocksDest<'_>) -> Result<(), NetError> { @@ -114,7 +164,7 @@ async fn read_connect_reply(s: &mut TcpStream) -> Result<(), NetError> { #[cfg(test)] mod tests { - use super::{socks5_connect, socks5_connect_domain}; + use super::{dial_isolated, socks5_connect, socks5_connect_domain, ProxyCreds}; use std::net::{Ipv4Addr, SocketAddr}; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpListener; @@ -203,4 +253,85 @@ mod tests { .unwrap(); server.await.unwrap(); } + + async fn serve_userpass_ipv4( + s: &mut tokio::net::TcpStream, + want_ip: [u8; 4], + want_port: u16, + ) -> (Vec, Vec) { + let mut ver_n = [0u8; 2]; + s.read_exact(&mut ver_n).await.unwrap(); + assert_eq!(ver_n[0], 5); + let nmethods = ver_n[1] as usize; + let mut methods = vec![0u8; nmethods]; + s.read_exact(&mut methods).await.unwrap(); + assert!(methods.contains(&0x02), "USERPASS offered, got {methods:?}"); + s.write_all(&[5, 0x02]).await.unwrap(); + + let mut ver = [0u8; 1]; + s.read_exact(&mut ver).await.unwrap(); + assert_eq!(ver[0], 1); + let mut ulen = [0u8; 1]; + s.read_exact(&mut ulen).await.unwrap(); + let mut user = vec![0u8; ulen[0] as usize]; + s.read_exact(&mut user).await.unwrap(); + let mut plen = [0u8; 1]; + s.read_exact(&mut plen).await.unwrap(); + let mut pass = vec![0u8; plen[0] as usize]; + s.read_exact(&mut pass).await.unwrap(); + s.write_all(&[1, 0]).await.unwrap(); + + let mut hdr = [0u8; 4]; + s.read_exact(&mut hdr).await.unwrap(); + assert_eq!(hdr[3], 1); + let mut addr = [0u8; 4]; + s.read_exact(&mut addr).await.unwrap(); + let mut p = [0u8; 2]; + s.read_exact(&mut p).await.unwrap(); + assert_eq!(addr, want_ip); + assert_eq!(u16::from_be_bytes(p), want_port); + s.write_all(&[5, 0, 0, 1, 0, 0, 0, 0, 0, 0]).await.unwrap(); + (user, pass) + } + + #[tokio::test] + async fn socks5_username_password_seen_by_fake_proxy() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let proxy = listener.local_addr().unwrap(); + let target = SocketAddr::from((Ipv4Addr::new(198, 51, 100, 1), 8333)); + let creds = ProxyCreds { + username: b"alice".to_vec(), + password: b"secret".to_vec(), + }; + let server = tokio::spawn(async move { + let (mut s, _) = listener.accept().await.unwrap(); + serve_userpass_ipv4(&mut s, [198, 51, 100, 1], 8333).await + }); + socks5_connect(proxy, target, Some(&creds)).await.unwrap(); + let (user, pass) = server.await.unwrap(); + assert_eq!(user, b"alice"); + assert_eq!(pass, b"secret"); + } + + #[tokio::test] + async fn dial_isolated_uses_new_creds_each_call() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let proxy = listener.local_addr().unwrap(); + let target = SocketAddr::from((Ipv4Addr::new(198, 51, 100, 2), 8333)); + let (tx, mut rx) = tokio::sync::mpsc::channel(2); + let server = tokio::spawn(async move { + for _ in 0..2 { + let (mut s, _) = listener.accept().await.unwrap(); + let creds = serve_userpass_ipv4(&mut s, [198, 51, 100, 2], 8333).await; + tx.send(creds).await.unwrap(); + } + }); + dial_isolated(proxy, target).await.unwrap(); + dial_isolated(proxy, target).await.unwrap(); + let (u1, _) = rx.recv().await.unwrap(); + let (u2, _) = rx.recv().await.unwrap(); + server.await.unwrap(); + assert_ne!(u1, u2, "each isolated dial must use fresh SOCKS creds"); + assert!(!u1.is_empty() && !u2.is_empty()); + } } From d260991d37207bca05be355873cb1c7f60259028 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 08:30:45 -0700 Subject: [PATCH 04/10] net: Dialer routes all P2P outbound through SOCKS when set spawn_peer, tip-follow, and feelers share one Dialer. A fake SOCKS splice plus BIP324 handshake proves CONNECT is used; Direct is unchanged. Co-authored-by: Cursor --- crates/rbitcoin-net/src/ibd/dial.rs | 6 +- crates/rbitcoin-net/src/ibd/mod.rs | 8 +- crates/rbitcoin-net/src/ibd/peer_io.rs | 4 +- crates/rbitcoin-net/src/lib.rs | 1 + crates/rbitcoin-net/src/service.rs | 52 ++++++- crates/rbitcoin-net/src/socks.rs | 188 ++++++++++++++++++++++++- 6 files changed, 246 insertions(+), 13 deletions(-) diff --git a/crates/rbitcoin-net/src/ibd/dial.rs b/crates/rbitcoin-net/src/ibd/dial.rs index 8c0211c84..23f61f279 100644 --- a/crates/rbitcoin-net/src/ibd/dial.rs +++ b/crates/rbitcoin-net/src/ibd/dial.rs @@ -243,6 +243,7 @@ pub(crate) async fn dial_batch( sinks: PeerEventSinks, connect_timeout: Duration, cancel: Option>, + dialer: crate::socks::Dialer, ) -> DialBatchResult { let mut out = DialBatchResult { slots: Vec::new(), @@ -278,8 +279,9 @@ pub(crate) async fn dial_batch( "{}", trying_connection_log(PeerConnType::OutboundFullRelay, addr) ); + let dialer = dialer.clone(); handles.push(tokio::spawn(async move { - let fut = spawn_peer(id, addr, magic, local_addr, tip_h, sinks); + let fut = spawn_peer(id, addr, magic, local_addr, tip_h, sinks, dialer); match tokio::time::timeout(connect_timeout, fut).await { Ok(Ok(slot)) => Ok(slot), Ok(Err(e)) => { @@ -968,6 +970,7 @@ mod tests { sinks.clone(), Duration::from_millis(50), None, + crate::socks::Dialer::Direct, )); assert!(r.slots.is_empty() && r.failed.is_empty()); let r2 = rt.block_on(dial_batch( @@ -982,6 +985,7 @@ mod tests { sinks, Duration::from_millis(50), None, + crate::socks::Dialer::Direct, )); assert!(r2.slots.is_empty() && r2.failed.is_empty()); } diff --git a/crates/rbitcoin-net/src/ibd/mod.rs b/crates/rbitcoin-net/src/ibd/mod.rs index 417076355..39015c5c0 100644 --- a/crates/rbitcoin-net/src/ibd/mod.rs +++ b/crates/rbitcoin-net/src/ibd/mod.rs @@ -164,6 +164,8 @@ pub struct IbdConfig { /// Optional shared peer book (discovered addrs + flags). Seeded at start and /// written back on IBD exit so the node can persist across runs. pub peers: Option>>, + /// Outbound TCP: direct or SOCKS5. + pub dialer: crate::socks::Dialer, } impl Default for IbdConfig { @@ -176,6 +178,7 @@ impl Default for IbdConfig { headers_batch: MAX_HEADERS_RESULTS, connect_timeout: Duration::from_secs(8), peers: None, + dialer: crate::socks::Dialer::Direct, } } } @@ -191,6 +194,7 @@ impl IbdConfig { headers_batch: MAX_HEADERS_RESULTS, connect_timeout: Duration::from_millis(400), peers: None, + dialer: crate::socks::Dialer::Direct, } } } @@ -301,6 +305,7 @@ pub async fn ibd_cancellable( sinks.clone(), cfg.connect_timeout, cancel.as_ref().map(Arc::clone), + cfg.dialer.clone(), ) .await; apply_dial_result(peer_sess.book_mut(), &initial); @@ -736,10 +741,11 @@ pub async fn ibd_cancellable( let sinks_r = sinks.clone(); let cto = cfg.connect_timeout; let cancel_c = cancel.as_ref().map(Arc::clone); + let dialer_c = cfg.dialer.clone(); redial_handle = Some(tokio::spawn(async move { dial_batch( &book, &next_id, want, already, &occupied, magic, local_addr, tip_h, sinks_r, - cto, cancel_c, + cto, cancel_c, dialer_c, ) .await })); diff --git a/crates/rbitcoin-net/src/ibd/peer_io.rs b/crates/rbitcoin-net/src/ibd/peer_io.rs index c8e0b6501..c992a1b11 100644 --- a/crates/rbitcoin-net/src/ibd/peer_io.rs +++ b/crates/rbitcoin-net/src/ibd/peer_io.rs @@ -21,7 +21,6 @@ use std::net::SocketAddr; use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::Arc; use std::time::Instant; -use tokio::net::TcpStream; use tokio::sync::mpsc; use tokio::task::JoinHandle; @@ -159,8 +158,9 @@ pub(crate) async fn spawn_peer( local: SocketAddr, tip_h: Option, sinks: PeerEventSinks, + dialer: crate::socks::Dialer, ) -> Result { - let stream = TcpStream::connect(addr).await?; + let stream = dialer.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, _tcp_shutdown) = connect_and_handshake_timed( diff --git a/crates/rbitcoin-net/src/lib.rs b/crates/rbitcoin-net/src/lib.rs index cf5bbf9eb..d4c7db730 100644 --- a/crates/rbitcoin-net/src/lib.rs +++ b/crates/rbitcoin-net/src/lib.rs @@ -62,6 +62,7 @@ pub use seeds::{ }; pub use serve_perf::{format_serve_perf, sample_reset_serve_perf, ServePerfSample}; pub use service::P2PNode; +pub use socks::Dialer; pub use tx_relay::{ ElectrumMempoolItem, MempoolAnnounce, MempoolHub, MempoolPerfSample, MempoolTxSnapEntry, MempoolTxSnapshot, diff --git a/crates/rbitcoin-net/src/service.rs b/crates/rbitcoin-net/src/service.rs index 141474c1d..7f4e5de74 100644 --- a/crates/rbitcoin-net/src/service.rs +++ b/crates/rbitcoin-net/src/service.rs @@ -20,7 +20,7 @@ use std::net::SocketAddr; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; use std::time::Duration; -use tokio::net::{TcpListener, TcpStream}; +use tokio::net::TcpListener; use tokio::task::JoinHandle; /// Running P2P node handle (listen + optional outbound sync / tip follow). @@ -46,6 +46,7 @@ pub struct P2PNode { pub max_inbound: usize, /// Shared inbound slots across all listen sockets. inbound_sem: Arc, + dialer: crate::socks::Dialer, } impl P2PNode { @@ -76,6 +77,28 @@ impl P2PNode { milestone: Milestone, user_agent: String, max_inbound: usize, + ) -> Result { + Self::start_with_dialer( + listen, + query, + params, + milestone, + user_agent, + max_inbound, + crate::socks::Dialer::Direct, + ) + .await + } + + /// Like [`Self::start_with_agent`] with an outbound [`crate::Dialer`]. + pub async fn start_with_dialer( + listen: SocketAddr, + query: Query, + params: ChainParams, + milestone: Milestone, + user_agent: String, + max_inbound: usize, + dialer: crate::socks::Dialer, ) -> Result { let magic = magic_for_params(¶ms); let hub = Arc::new(ChainHub::new(query, params, milestone)); @@ -113,6 +136,7 @@ impl P2PNode { let dial_live = follow_live.clone(); let dial_shutdown = shutdown.clone(); let sessions_dial = session_tasks.clone(); + let dialer_task = dialer.clone(); let dial_task = tokio::spawn(async move { while let Some(req) = dial_rx.recv().await { if dial_shutdown.load(Ordering::SeqCst) { @@ -122,10 +146,11 @@ impl P2PNode { let peers = dial_peers.clone(); let ua = dial_ua.clone(); let live = dial_live.clone(); + let d = dialer_task.clone(); let (ah_tx, ah_rx) = tokio::sync::oneshot::channel::(); let h = tokio::spawn(async move { let _ = run_outbound_session_with_abort( - req.addr, magic, local_addr, hub, peers, ua, live, req.typ, ah_rx, + req.addr, magic, local_addr, hub, peers, ua, live, req.typ, ah_rx, d, ) .await; }); @@ -148,6 +173,7 @@ impl P2PNode { user_agent, max_inbound, inbound_sem, + dialer, }) } @@ -241,6 +267,7 @@ impl P2PNode { self.user_agent.clone(), self.follow_live.clone(), PeerConnType::OutboundFullRelay, + self.dialer.clone(), ) .await?; let handle = tokio::spawn(async move { @@ -454,9 +481,10 @@ async fn prepare_outbound_session( user_agent: String, follow_live: Arc, typ: PeerConnType, + dialer: crate::socks::Dialer, ) -> Result { rbitcoin_log::debug!("{}", crate::peers::trying_connection_log(typ, peer)); - let stream = TcpStream::connect(peer).await?; + let stream = dialer.connect(peer).await?; let bind = stream.local_addr().unwrap_or(local); let height = hub.tip_height().map(|h| h as i32).unwrap_or(0); // Core adds CNode before VERSION. Provisional row so getpeerinfo is non-empty @@ -551,15 +579,25 @@ async fn run_outbound_session_with_abort( follow_live: Arc, typ: PeerConnType, ah_rx: tokio::sync::oneshot::Receiver, + dialer: crate::socks::Dialer, ) -> Result<(), NetError> { if typ == PeerConnType::Feeler { - let stream = TcpStream::connect(peer).await?; + let stream = dialer.connect(peer).await?; let height = hub.tip_height().map(|h| h as i32).unwrap_or(0); return crate::peer::run_feeler(stream, magic, local, peer, height, &user_agent).await; } - let prepared = - prepare_outbound_session(peer, magic, local, hub, peers, user_agent, follow_live, typ) - .await?; + let prepared = prepare_outbound_session( + peer, + magic, + local, + hub, + peers, + user_agent, + follow_live, + typ, + dialer, + ) + .await?; if let Ok(ah) = ah_rx.await { prepared.sess.set_session_abort(ah); } diff --git a/crates/rbitcoin-net/src/socks.rs b/crates/rbitcoin-net/src/socks.rs index 4c163d6cc..e19b7575e 100644 --- a/crates/rbitcoin-net/src/socks.rs +++ b/crates/rbitcoin-net/src/socks.rs @@ -50,6 +50,57 @@ pub(crate) async fn dial_isolated( socks5_connect(proxy, target, Some(&creds)).await } +#[derive(Clone, Debug, Default, PartialEq, Eq)] +pub enum Dialer { + #[default] + Direct, + Socks { + proxy: SocketAddr, + randomize: bool, + }, +} + +impl Dialer { + pub async fn connect(&self, target: SocketAddr) -> Result { + match self { + Dialer::Direct => Ok(TcpStream::connect(target).await?), + Dialer::Socks { proxy, randomize } => { + if *randomize { + let creds = ProxyCreds::fresh(); + socks5_connect(*proxy, target, Some(&creds)).await + } else { + socks5_connect(*proxy, target, None).await + } + } + } + } + + pub async fn connect_domain(&self, host: &str, port: u16) -> Result { + match self { + Dialer::Direct => { + let mut addrs = tokio::net::lookup_host((host, port)).await?; + let addr = addrs.next().ok_or(NetError::Protocol("dns lookup empty"))?; + self.connect(addr).await + } + Dialer::Socks { proxy, randomize } => { + if *randomize { + let creds = ProxyCreds::fresh(); + socks5_connect_domain(*proxy, host, port, Some(&creds)).await + } else { + socks5_connect_domain(*proxy, host, port, None).await + } + } + } + } + + pub async fn connect_isolated(&self, target: SocketAddr) -> Result { + match self { + Dialer::Direct => self.connect(target).await, + Dialer::Socks { proxy, .. } => dial_isolated(*proxy, target).await, + } + } +} + async fn socks5_connect_dest( proxy: SocketAddr, dest: SocksDest<'_>, @@ -164,10 +215,10 @@ async fn read_connect_reply(s: &mut TcpStream) -> Result<(), NetError> { #[cfg(test)] mod tests { - use super::{dial_isolated, socks5_connect, socks5_connect_domain, ProxyCreds}; + use super::{dial_isolated, socks5_connect, socks5_connect_domain, Dialer, ProxyCreds}; use std::net::{Ipv4Addr, SocketAddr}; use tokio::io::{AsyncReadExt, AsyncWriteExt}; - use tokio::net::TcpListener; + use tokio::net::{TcpListener, TcpStream}; #[tokio::test] async fn socks5_connect_ipv4_against_fake_proxy() { @@ -334,4 +385,137 @@ mod tests { assert_ne!(u1, u2, "each isolated dial must use fresh SOCKS creds"); assert!(!u1.is_empty() && !u2.is_empty()); } + + async fn splice_one_socks( + listener: TcpListener, + saw: tokio::sync::oneshot::Sender, + ) { + let (mut c, _) = listener.accept().await.unwrap(); + let mut ver_n = [0u8; 2]; + c.read_exact(&mut ver_n).await.unwrap(); + let nmethods = ver_n[1] as usize; + let mut methods = vec![0u8; nmethods]; + c.read_exact(&mut methods).await.unwrap(); + if methods.contains(&0x02) { + c.write_all(&[5, 0x02]).await.unwrap(); + let mut ver = [0u8; 1]; + c.read_exact(&mut ver).await.unwrap(); + let mut ulen = [0u8; 1]; + c.read_exact(&mut ulen).await.unwrap(); + let mut user = vec![0u8; ulen[0] as usize]; + c.read_exact(&mut user).await.unwrap(); + let mut plen = [0u8; 1]; + c.read_exact(&mut plen).await.unwrap(); + let mut pass = vec![0u8; plen[0] as usize]; + c.read_exact(&mut pass).await.unwrap(); + c.write_all(&[1, 0]).await.unwrap(); + } else { + c.write_all(&[5, 0x00]).await.unwrap(); + } + let mut hdr = [0u8; 4]; + c.read_exact(&mut hdr).await.unwrap(); + assert_eq!(hdr[1], 1); + let dest = match hdr[3] { + 1 => { + let mut a = [0u8; 4]; + c.read_exact(&mut a).await.unwrap(); + let mut p = [0u8; 2]; + c.read_exact(&mut p).await.unwrap(); + SocketAddr::from((Ipv4Addr::new(a[0], a[1], a[2], a[3]), u16::from_be_bytes(p))) + } + _ => panic!("test splice expects IPv4 CONNECT"), + }; + let mut peer = TcpStream::connect(dest).await.unwrap(); + c.write_all(&[5, 0, 0, 1, 0, 0, 0, 0, 0, 0]).await.unwrap(); + let _ = saw.send(dest); + let _ = tokio::io::copy_bidirectional(&mut c, &mut peer).await; + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn outbound_dial_uses_proxy_when_set() { + use crate::peer::{connect_and_handshake_timed, HandshakePolicy, HANDSHAKE_TIMEOUT}; + use bitcoin::p2p::Magic; + use std::time::Duration; + + let bitcoin_l = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let peer_addr = bitcoin_l.local_addr().unwrap(); + let inbound = tokio::spawn(async move { + let (stream, from) = bitcoin_l.accept().await.unwrap(); + connect_and_handshake_timed( + Duration::from_secs(5), + stream, + Magic::REGTEST, + peer_addr, + from, + 0, + true, + "/rbitcoin:test/", + HandshakePolicy::plain(), + ) + .await + }); + + let socks_l = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let proxy = socks_l.local_addr().unwrap(); + let (saw_tx, saw_rx) = tokio::sync::oneshot::channel(); + let splice = tokio::spawn(splice_one_socks(socks_l, saw_tx)); + + let stream = Dialer::Socks { + proxy, + randomize: true, + } + .connect(peer_addr) + .await + .unwrap(); + assert_eq!(saw_rx.await.unwrap(), peer_addr); + + connect_and_handshake_timed( + Duration::from_secs(5), + stream, + Magic::REGTEST, + proxy, + peer_addr, + 0, + false, + "/rbitcoin:test/", + HandshakePolicy::plain(), + ) + .await + .unwrap(); + inbound.await.unwrap().unwrap(); + splice.abort(); + + let direct_l = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let direct_addr = direct_l.local_addr().unwrap(); + let inbound = tokio::spawn(async move { + let (stream, from) = direct_l.accept().await.unwrap(); + connect_and_handshake_timed( + HANDSHAKE_TIMEOUT, + stream, + Magic::REGTEST, + direct_addr, + from, + 0, + true, + "/rbitcoin:test/", + HandshakePolicy::plain(), + ) + .await + }); + let stream = Dialer::Direct.connect(direct_addr).await.unwrap(); + connect_and_handshake_timed( + Duration::from_secs(5), + stream, + Magic::REGTEST, + direct_addr, + direct_addr, + 0, + false, + "/rbitcoin:test/", + HandshakePolicy::plain(), + ) + .await + .unwrap(); + inbound.await.unwrap().unwrap(); + } } From f7e9c721fb17aa79451589ca518e9d34e183e16f Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 08:36:04 -0700 Subject: [PATCH 05/10] node: --proxy/--onion CLI and conf set SOCKS endpoints Operator SOCKS is a ListenOpts dialer passed into P2P start and IBD, so outbound never bypasses the proxy when one is configured. Co-authored-by: Cursor --- crates/rbitcoin-node/src/cli.rs | 71 +++++++++++++++++++++++++++++- crates/rbitcoin-node/src/config.rs | 39 ++++++++++++++++ crates/rbitcoin-node/src/run.rs | 23 +++++++--- 3 files changed, 126 insertions(+), 7 deletions(-) diff --git a/crates/rbitcoin-node/src/cli.rs b/crates/rbitcoin-node/src/cli.rs index 61414d1fa..3827f5804 100644 --- a/crates/rbitcoin-node/src/cli.rs +++ b/crates/rbitcoin-node/src/cli.rs @@ -288,7 +288,8 @@ fn operator_usage() -> String { "rbitcoin-node {} — usage:\n\ rbitcoin-node [--conf FILE] [--datadir PATH] [--datadir-cold PATH] [--network NET] \\\n\ [--signet-challenge HEX] [--signet-block-time SECS] \\\n\ - [--listen ADDR] [--connect ADDR]... [--seed-node HOST]... [--electrum-listen ADDR] [--esplora-listen ADDR] \\\n\ + [--listen ADDR] [--connect ADDR]... [--seed-node HOST]... [--proxy HOST:PORT] [--onion HOST:PORT] [--proxy-randomize[=0|1]] \\\n\ + [--electrum-listen ADDR] [--esplora-listen ADDR] \\\n\ [--sh-index] [--sp-tweaks] [--sp-tweaks-dust SATS] [--max-sh-creates N] [--esplora-block-template] \\\n\ [--rpc] [--rpc-listen [ADDR]] [--rpc-token-file PATH] [--rpc-work-queue N] \\\n\ [--milestone HEIGHT] \\\n\ @@ -313,6 +314,8 @@ Milestone: skip script/sig checks at/below HEIGHT.\n\ Check-blocks: --check-blocks N revalidates the last N confirmed heights on open (default 6; 0 = all).\n\ Mempool: --mempool-size-mb (default ~300 MiB weight budget).\n\ Peers: --max-outbound (default 16 live download), --max-inbound (default 125).\n\ + --proxy HOST:PORT SOCKS5 for all P2P outbound; --onion HOST:PORT SOCKS for onion (02).\n\ + --proxy-randomize (default on) uses a fresh SOCKS username per peer (Tor circuit isolation).\n\ --trusted / --always-relay / --relay are inbound permission knobs.\n\ --net-permission / --net-permission-bind are CIDR or bind grants (noban, relay, …; IPv4 and IPv6).\n\ --net-permission-relay (default on) / --net-permission-force-relay (default off) are implicit bits on a bare CIDR grant.\n\ @@ -367,6 +370,7 @@ fn is_bool_key(key: &str) -> bool { | "net_permission_relay" | "net_permission_force_relay" | "no_seeds" + | "proxy_randomize" | "inhibit_suspend" | "trusted" | "always_relay" @@ -555,6 +559,9 @@ mod tests { "--rpc", "--rpc-listen", "--rpc-token-file", + "--proxy", + "--onion", + "--proxy-randomize", ] { assert!(h.contains(flag), "help must list {flag}"); } @@ -614,6 +621,68 @@ mod tests { ); } + #[test] + fn proxy_conf_and_cli() { + let _g = OPERATOR_ENV_TEST_LOCK.lock().unwrap(); + let cfg = ready_config(["rbitcoin-node", "--proxy", "127.0.0.1:9050"]); + assert_eq!(cfg.listen.proxy, Some("127.0.0.1:9050".parse().unwrap())); + assert!(cfg.listen.onion.is_none()); + assert_eq!( + cfg.listen.dialer(), + rbitcoin_net::Dialer::Socks { + proxy: "127.0.0.1:9050".parse().unwrap(), + randomize: true, + } + ); + assert_eq!( + NodeConfig::default().listen.dialer(), + rbitcoin_net::Dialer::Direct + ); + + let mut from_conf = NodeConfig::default(); + assert_eq!( + from_conf.apply_kv("proxy", "127.0.0.1:9050").unwrap(), + ConfApply::Applied + ); + assert_eq!( + from_conf.listen.proxy, + Some("127.0.0.1:9050".parse().unwrap()) + ); + + let empty = NodeConfig::default().apply_kv("proxy", "").unwrap_err(); + let empty_msg = format!("{empty}"); + assert!( + empty_msg.contains("proxy"), + "empty proxy must be a start error: {empty_msg}" + ); + + let bad = NodeConfig::default() + .apply_kv("proxy", "not-an-addr") + .unwrap_err(); + let bad_msg = format!("{bad}"); + assert!( + bad_msg.contains("proxy"), + "invalid proxy must be a start error: {bad_msg}" + ); + + let split = ready_config([ + "rbitcoin-node", + "--proxy", + "127.0.0.1:9050", + "--onion", + "127.0.0.1:9051", + ]); + assert_eq!(split.listen.proxy, Some("127.0.0.1:9050".parse().unwrap())); + assert_eq!(split.listen.onion, Some("127.0.0.1:9051".parse().unwrap())); + assert!(split.listen.proxy.is_some()); + assert!(split.listen.onion.is_some()); + assert_ne!(split.listen.proxy, split.listen.onion); + + let h = operator_usage(); + assert!(h.contains("--proxy"), "help must list kebab --proxy"); + assert!(h.contains("--onion"), "help must list kebab --onion"); + } + #[test] fn kebab_seed_node_and_min_relay_tx_fee_parse() { let seeds = ready_config(["rbitcoin-node", "--seed-node", "127.0.0.1:8333"]); diff --git a/crates/rbitcoin-node/src/config.rs b/crates/rbitcoin-node/src/config.rs index 33d678b37..3361cfdc0 100644 --- a/crates/rbitcoin-node/src/config.rs +++ b/crates/rbitcoin-node/src/config.rs @@ -84,6 +84,12 @@ pub struct ListenOpts { pub max_inbound_explicit: bool, pub external_ips: Vec, pub peer_timeout_secs: Option, + /// SOCKS5 for all P2P outbound (`--proxy`). + pub proxy: Option, + /// SOCKS5 for onion destinations (`--onion`); stored until plan 02. + pub onion: Option, + /// Fresh SOCKS USERPASS per peer (Core `-proxyrandomize`; default on). + pub proxy_randomize: bool, } impl Default for ListenOpts { @@ -101,6 +107,21 @@ impl Default for ListenOpts { max_inbound_explicit: false, external_ips: Vec::new(), peer_timeout_secs: None, + proxy: None, + onion: None, + proxy_randomize: true, + } + } +} + +impl ListenOpts { + pub fn dialer(&self) -> rbitcoin_net::Dialer { + match self.proxy { + None => rbitcoin_net::Dialer::Direct, + Some(proxy) => rbitcoin_net::Dialer::Socks { + proxy, + randomize: self.proxy_randomize, + }, } } } @@ -635,6 +656,16 @@ impl NodeConfig { .map_err(|e| NodeError::Config(format!("conf connect: {e}")))?, ); } + "proxy" => { + self.listen.proxy = Some(parse_required_socket(val, "proxy")?); + } + "onion" => { + self.listen.onion = Some(parse_required_socket(val, "onion")?); + } + "proxy_randomize" => { + self.listen.proxy_randomize = parse_conf_bool(val) + .map_err(|e| NodeError::Config(format!("conf proxy_randomize: {e}")))?; + } "seed_node" => { if !val.is_empty() { self.listen.seednodes.push(val.to_string()); @@ -963,6 +994,14 @@ pub(crate) fn parse_signet_challenge(value: &str) -> Result { .map_err(|e| format!("must be hexadecimal: {e}")) } +fn parse_required_socket(val: &str, key: &str) -> Result { + if val.is_empty() { + return Err(NodeError::Config(format!("conf {key}: empty"))); + } + val.parse() + .map_err(|e| NodeError::Config(format!("conf {key}: {e}"))) +} + fn is_conf_true(val: &str) -> bool { matches!( val.to_ascii_lowercase().as_str(), diff --git a/crates/rbitcoin-node/src/run.rs b/crates/rbitcoin-node/src/run.rs index 96832212e..1e6bec20c 100644 --- a/crates/rbitcoin-node/src/run.rs +++ b/crates/rbitcoin-node/src/run.rs @@ -7,8 +7,8 @@ use rbitcoin_esplora::{run_esplora, BlockTemplateFn, EsploraConfig, EsploraHandl use rbitcoin_log::{debug, enabled, info, warn, Level}; use rbitcoin_net::{ default_port, format_serve_perf, format_tip_perf_sizes, netgroup, read_proc_rss, - sample_reset_serve_perf, AddrMan, AsMap, BlockingRegion, ChainHub, IbdConfig, MempoolHub, - P2PNode, PeerConnType, TipEvent, TipPerfSizes, + sample_reset_serve_perf, AddrMan, AsMap, BlockingRegion, ChainHub, Dialer, IbdConfig, + MempoolHub, P2PNode, PeerConnType, TipEvent, TipPerfSizes, }; use rbitcoin_primitives::Network; use rbitcoin_query::{spawn_sh_writebehind, Query}; @@ -208,13 +208,14 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { let p2p_ua = rbitcoin_primitives::rbitcoin_subversion(env!("CARGO_PKG_VERSION"), &config.uacomments) .unwrap_or_else(|_| format!("/rbitcoin:{}/", env!("CARGO_PKG_VERSION"))); - let mut node = P2PNode::start_with_agent( + let mut node = P2PNode::start_with_dialer( listen, query, params.clone(), milestone, p2p_ua, config.listen.max_inbound as usize, + config.listen.dialer(), ) .await .map_err(|e| NodeError::Config(format!("p2p start: {e}")))?; @@ -413,6 +414,7 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { &mut addrman, &peers_path, &shutdown, + config.listen.dialer(), ) .await; @@ -885,7 +887,10 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { } } else { info!("ibd: retry catch-up from {peer} (tip stagnant, catch-up incomplete)"); - let retry_cfg = catch_up_retry_config(std::sync::Arc::clone(&shared_peers)); + let retry_cfg = catch_up_retry_config( + std::sync::Arc::clone(&shared_peers), + config.listen.dialer(), + ); let cancel = Some(Arc::clone(&shutdown.flag)); let retry_peers = [peer]; tokio::select! { @@ -1109,6 +1114,7 @@ async fn run_ibd_or_skip( addrman: &mut AddrMan, peers_path: &std::path::Path, shutdown: &Shutdown, + dialer: Dialer, ) -> CatchUp { if ibd_targets.is_empty() { info!("ibd: no outbound peers; serving only (use --connect or seeds)"); @@ -1126,6 +1132,7 @@ async fn run_ibd_or_skip( // could deliver mid-chain blocks). Default 30s is enough. stall: std::time::Duration::from_secs(30), peers: Some(std::sync::Arc::clone(shared_peers)), + dialer, ..IbdConfig::default() }; info!( @@ -1510,10 +1517,14 @@ pub(crate) fn enter_tip_mode( /// /// Uses [`IbdConfig::default`] (window 1024, stall 30s, connect 8s, …) — not /// [`IbdConfig::for_test`], which is only for unit/integration test harnesses. -fn catch_up_retry_config(peers: std::sync::Arc>) -> IbdConfig { +fn catch_up_retry_config( + peers: std::sync::Arc>, + dialer: Dialer, +) -> IbdConfig { IbdConfig { target_peers: 1, peers: Some(peers), + dialer, ..IbdConfig::default() } } @@ -1882,7 +1893,7 @@ mod tests { #[test] fn catch_up_retry_config_uses_production_not_for_test() { let peers = std::sync::Arc::new(std::sync::Mutex::new(rbitcoin_net::AddrMan::new())); - let cfg = catch_up_retry_config(std::sync::Arc::clone(&peers)); + let cfg = catch_up_retry_config(std::sync::Arc::clone(&peers), Dialer::Direct); let prod = IbdConfig::default(); let test = IbdConfig::for_test(); From 6632a5f6acedc00bddc2106ab8506ae8497b0781 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 08:38:37 -0700 Subject: [PATCH 06/10] node: skip local DNS seeds when --proxy is set SOCKS does remote DNS; ToSocketAddrs on seed hostnames would leak queries off-proxy. --proxy-randomize stays on unless explicitly off. Co-authored-by: Cursor --- crates/rbitcoin-net/src/lib.rs | 3 ++- crates/rbitcoin-net/src/seeds.rs | 33 ++++++++++++++++++++++++++++++++ crates/rbitcoin-node/src/cli.rs | 27 ++++++++++++++++++++++++++ crates/rbitcoin-node/src/run.rs | 23 +++++++++++++++++++--- 4 files changed, 82 insertions(+), 4 deletions(-) diff --git a/crates/rbitcoin-net/src/lib.rs b/crates/rbitcoin-net/src/lib.rs index d4c7db730..910ec0066 100644 --- a/crates/rbitcoin-net/src/lib.rs +++ b/crates/rbitcoin-net/src/lib.rs @@ -58,7 +58,8 @@ pub(crate) use rbitcoin_mempool::MempoolGraphStats; pub use reactor::BlockingRegion; pub use seeds::{ default_port, default_rpc_port, dns_seeds, fixed_seed_hosts, resolve_all_seeds, - resolve_dns_seeds, resolve_fixed_seeds, AddrMan, PeerEntry, PeerFlags, MAX_ADDR_MAN, + resolve_dns_seeds, resolve_fixed_seeds, socks_dns_seed_dests, AddrMan, PeerEntry, PeerFlags, + MAX_ADDR_MAN, }; pub use serve_perf::{format_serve_perf, sample_reset_serve_perf, ServePerfSample}; pub use service::P2PNode; diff --git a/crates/rbitcoin-net/src/seeds.rs b/crates/rbitcoin-net/src/seeds.rs index 73454365d..de932d014 100644 --- a/crates/rbitcoin-net/src/seeds.rs +++ b/crates/rbitcoin-net/src/seeds.rs @@ -155,6 +155,17 @@ pub fn resolve_all_seeds(network: Network) -> Vec { out } +/// DNS seed hostnames with the network default port for SOCKS domain CONNECT. +/// +/// Does not call [`ToSocketAddrs`] — the proxy performs remote DNS. +pub fn socks_dns_seed_dests(network: Network) -> Vec<(String, u16)> { + let port = default_port(network); + dns_seeds(network) + .iter() + .map(|host| ((*host).to_string(), port)) + .collect() +} + /// Informational peer flags packed into one byte (more bits reserved for later). /// /// | bit | name | meaning | @@ -1233,6 +1244,28 @@ mod tests { assert!(resolve_all_seeds(Network::Regtest).is_empty()); } + #[test] + fn dns_seeds_not_resolved_locally_when_proxy() { + for net in [ + Network::Mainnet, + Network::Testnet, + Network::Signet, + Network::Regtest, + ] { + let dests = socks_dns_seed_dests(net); + let names = dns_seeds(net); + assert_eq!(dests.len(), names.len()); + for ((host, port), want) in dests.iter().zip(names.iter()) { + assert_eq!(host, want); + assert_eq!(*port, default_port(net)); + assert!( + host.parse::().is_err(), + "SOCKS seed dest must stay a hostname, got {host}" + ); + } + } + } + #[test] fn peer_flags_set_remove_and_mid_range_speed() { let mut f = PeerFlags::empty(); diff --git a/crates/rbitcoin-node/src/cli.rs b/crates/rbitcoin-node/src/cli.rs index 3827f5804..072b3b5ed 100644 --- a/crates/rbitcoin-node/src/cli.rs +++ b/crates/rbitcoin-node/src/cli.rs @@ -681,6 +681,33 @@ mod tests { let h = operator_usage(); assert!(h.contains("--proxy"), "help must list kebab --proxy"); assert!(h.contains("--onion"), "help must list kebab --onion"); + assert!( + h.contains("--proxy-randomize"), + "help must list kebab --proxy-randomize" + ); + } + + #[test] + fn proxy_randomize_defaults_on() { + let _g = OPERATOR_ENV_TEST_LOCK.lock().unwrap(); + assert!(NodeConfig::default().listen.proxy_randomize); + let off = ready_config([ + "rbitcoin-node", + "--proxy", + "127.0.0.1:9050", + "--proxy-randomize=0", + ]); + assert!(!off.listen.proxy_randomize); + match off.listen.dialer() { + rbitcoin_net::Dialer::Socks { randomize, .. } => assert!(!randomize), + other => panic!("expected socks dialer, got {other:?}"), + } + let on = ready_config(["rbitcoin-node", "--proxy", "127.0.0.1:9050"]); + assert!(on.listen.proxy_randomize); + match on.listen.dialer() { + rbitcoin_net::Dialer::Socks { randomize, .. } => assert!(randomize), + other => panic!("expected socks dialer, got {other:?}"), + } } #[test] diff --git a/crates/rbitcoin-node/src/run.rs b/crates/rbitcoin-node/src/run.rs index 1e6bec20c..9ba8acd7a 100644 --- a/crates/rbitcoin-node/src/run.rs +++ b/crates/rbitcoin-node/src/run.rs @@ -7,8 +7,8 @@ use rbitcoin_esplora::{run_esplora, BlockTemplateFn, EsploraConfig, EsploraHandl use rbitcoin_log::{debug, enabled, info, warn, Level}; use rbitcoin_net::{ default_port, format_serve_perf, format_tip_perf_sizes, netgroup, read_proc_rss, - sample_reset_serve_perf, AddrMan, AsMap, BlockingRegion, ChainHub, Dialer, IbdConfig, - MempoolHub, P2PNode, PeerConnType, TipEvent, TipPerfSizes, + sample_reset_serve_perf, socks_dns_seed_dests, AddrMan, AsMap, BlockingRegion, ChainHub, + Dialer, IbdConfig, MempoolHub, P2PNode, PeerConnType, TipEvent, TipPerfSizes, }; use rbitcoin_primitives::Network; use rbitcoin_query::{spawn_sh_writebehind, Query}; @@ -392,6 +392,9 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { addrman.len().saturating_sub(n_before), addrman.len() ); + } else if config.listen.proxy.is_some() && config.listen.use_seeds { + let n = socks_dns_seed_dests(config.network).len(); + info!("ibd: SOCKS proxy set — skipping local DNS for {n} seed hostnames (use --connect)"); } else if config.signet_challenge.is_some() && config.listen.connect.is_empty() && addrman.is_empty() @@ -990,7 +993,10 @@ fn tip_meets_min_work(config: &NodeConfig, hub: &rbitcoin_net::ChainHub) -> bool } fn should_resolve_default_seeds(config: &NodeConfig) -> bool { - config.listen.use_seeds && config.listen.connect.is_empty() && config.signet_challenge.is_none() + config.listen.use_seeds + && config.listen.connect.is_empty() + && config.signet_challenge.is_none() + && config.listen.proxy.is_none() } /// One walker per process: SH-warm start and post-IBD `enter_tip_mode` both call this. @@ -1920,6 +1926,17 @@ mod tests { assert!(!should_resolve_default_seeds(&cfg)); } + #[test] + fn dns_seeds_not_resolved_locally_when_proxy() { + let mut cfg = NodeConfig::default(); + assert!(should_resolve_default_seeds(&cfg)); + cfg.listen.proxy = Some("127.0.0.1:9050".parse().unwrap()); + assert!( + !should_resolve_default_seeds(&cfg), + "proxy path must not ToSocketAddrs DNS/fixed seeds" + ); + } + fn coinbase_block( h: u32, prev: rbitcoin_primitives::Fk, From ead3ea621abbfb49b93ef789b16a1d08f0cf0049 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 08:39:46 -0700 Subject: [PATCH 07/10] docs+nix: first-class SOCKS --proxy options Home-node Tor SOCKS is a module option and OPERATOR flag, not extraArgs. proxyRandomize stays on so Tor isolates circuits unless the operator turns it off. Co-authored-by: Cursor --- OPERATOR.md | 11 +++++++++++ nix/modules/rbitcoin.nix | 25 +++++++++++++++++++++++++ nix/tests/nixos-module-eval.nix | 8 ++++++++ 3 files changed, 44 insertions(+) diff --git a/OPERATOR.md b/OPERATOR.md index 5f37a8334..c104f2f04 100644 --- a/OPERATOR.md +++ b/OPERATOR.md @@ -367,6 +367,9 @@ Clean smoke: | `--signet-block-time SECS` | `signet_block_time=` | 600; requires a custom challenge | | `--listen ADDR` | `listen=` | bind later default port | | `--connect ADDR` | `connect=` (repeatable) | seeds | +| `--proxy HOST:PORT` | `proxy=` | unset — SOCKS5 for all P2P outbound | +| `--onion HOST:PORT` | `onion=` | unset — SOCKS5 for onion destinations | +| `--proxy-randomize[=0\|1]` | `proxy_randomize=` | **on** — fresh SOCKS username per peer (Tor circuit isolation) | | `--milestone HEIGHT` | `milestone=` | network default (mainnet 840000) | | `--max-outbound N` | `max_outbound=` | 16 live download peers | | `--max-inbound N` | `max_inbound=` | 125 inbound sessions | @@ -425,6 +428,14 @@ max_inbound=64 mempool_size_mb=100 ``` +### P2P via system Tor SOCKS + +`--proxy 127.0.0.1:9050` sends every P2P outbound through SOCKS5 CONNECT +(system `tor`, not Arti). DNS seeds are not resolved locally on that path — +pass `--connect ADDR` (or reuse a `peers` file). `--proxy-randomize` (default +on) uses a fresh SOCKS username per peer so Tor isolates circuits. +`--onion HOST:PORT` stores a separate SOCKS endpoint for onion destinations. + `--datadir` holds the node root (`store/`, `mempool/`, `peers`, `rpc.token`, `rpc.sock`). Omit `--datadir-cold` and cold files live there too. Set it to put the large rarely-read Class A **inwit** stem (`inwit.body` + `inwit.loc`, ~486 GiB + loc diff --git a/nix/modules/rbitcoin.nix b/nix/modules/rbitcoin.nix index 6101a0d76..356e724b2 100644 --- a/nix/modules/rbitcoin.nix +++ b/nix/modules/rbitcoin.nix @@ -60,6 +60,11 @@ let ++ optional cfg.esplora.enable (socket cfg.esplora.address cfg.esplora.port) ++ optional (cfg.scripthashIndex || cfg.electrum.enable || cfg.esplora.enable) "--shindex" ++ optional cfg.silentPaymentIndex "--sptweaks" + ++ optional (cfg.proxy != null) "--proxy" + ++ optional (cfg.proxy != null) cfg.proxy + ++ optional (cfg.onionProxy != null) "--onion" + ++ optional (cfg.onionProxy != null) cfg.onionProxy + ++ optional (!cfg.proxyRandomize) "--proxy-randomize=0" ++ cfg.extraArgs; in { @@ -149,6 +154,26 @@ in description = "Additional command-line arguments appended after module-managed arguments."; }; + proxy = mkOption { + type = types.nullOr types.str; + default = null; + example = "127.0.0.1:9050"; + description = "SOCKS5 proxy HOST:PORT for all P2P outbound."; + }; + + onionProxy = mkOption { + type = types.nullOr types.str; + default = null; + example = "127.0.0.1:9050"; + description = "SOCKS5 proxy HOST:PORT for onion destinations."; + }; + + proxyRandomize = mkOption { + type = types.bool; + default = true; + description = "Fresh SOCKS username per peer so Tor isolates circuits."; + }; + p2p = { address = mkOption { type = types.str; diff --git a/nix/tests/nixos-module-eval.nix b/nix/tests/nixos-module-eval.nix index bd4720737..f8ff099e5 100644 --- a/nix/tests/nixos-module-eval.nix +++ b/nix/tests/nixos-module-eval.nix @@ -28,6 +28,9 @@ let "--max-outbound" "8" ]; + proxy = "127.0.0.1:9050"; + onionProxy = "127.0.0.1:9050"; + proxyRandomize = true; p2p = { address = "127.0.0.1"; openFirewall = true; @@ -61,6 +64,9 @@ assert defaultCfg.package == expectedPackage; assert defaultCfg.network == "mainnet"; assert defaultCfg.p2p.port == 8333; assert defaultCfg.rpc.port == 8332; +assert defaultCfg.proxy == null; +assert defaultCfg.onionProxy == null; +assert defaultCfg.proxyRandomize == true; assert cfg.services.rbitcoin.p2p.port == 18444; assert cfg.services.rbitcoin.rpc.port == 18443; assert @@ -84,4 +90,6 @@ assert builtins.match ".*--shindex.*" execStart != null; assert builtins.match ".*--sptweaks.*" execStart != null; assert builtins.match ".*--log-level debug.*" execStart != null; assert builtins.match ".*--max-outbound 8.*" execStart != null; +assert builtins.match ".*--proxy 127.0.0.1:9050.*" execStart != null; +assert builtins.match ".*--onion 127.0.0.1:9050.*" execStart != null; pkgs.runCommand "rbitcoin-nixos-module-eval" { } "touch $out" From 1f8a637ede3df9ffb5c8f6d1d4c096b53d6125e9 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 08:41:41 -0700 Subject: [PATCH 08/10] net: IBD reuses the node's Dialer Catch-up and retry already start through P2PNode; cloning that dialer keeps SOCKS on the IBD path without growing run_ibd_or_skip's arity. Co-authored-by: Cursor --- crates/rbitcoin-net/src/service.rs | 5 +++++ crates/rbitcoin-node/src/run.rs | 10 +++------- 2 files changed, 8 insertions(+), 7 deletions(-) diff --git a/crates/rbitcoin-net/src/service.rs b/crates/rbitcoin-net/src/service.rs index 7f4e5de74..a19944036 100644 --- a/crates/rbitcoin-net/src/service.rs +++ b/crates/rbitcoin-net/src/service.rs @@ -177,6 +177,11 @@ impl P2PNode { }) } + /// Outbound TCP path (direct or SOCKS). + pub fn dialer(&self) -> crate::socks::Dialer { + self.dialer.clone() + } + /// Bind an additional listen socket (Core multi-`-bind`). pub async fn add_listen(&mut self, listen: SocketAddr) -> Result { let listener = TcpListener::bind(listen).await?; diff --git a/crates/rbitcoin-node/src/run.rs b/crates/rbitcoin-node/src/run.rs index 9ba8acd7a..164a427c4 100644 --- a/crates/rbitcoin-node/src/run.rs +++ b/crates/rbitcoin-node/src/run.rs @@ -417,7 +417,6 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { &mut addrman, &peers_path, &shutdown, - config.listen.dialer(), ) .await; @@ -890,10 +889,8 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { } } else { info!("ibd: retry catch-up from {peer} (tip stagnant, catch-up incomplete)"); - let retry_cfg = catch_up_retry_config( - std::sync::Arc::clone(&shared_peers), - config.listen.dialer(), - ); + let retry_cfg = + catch_up_retry_config(std::sync::Arc::clone(&shared_peers), node.dialer()); let cancel = Some(Arc::clone(&shutdown.flag)); let retry_peers = [peer]; tokio::select! { @@ -1120,7 +1117,6 @@ async fn run_ibd_or_skip( addrman: &mut AddrMan, peers_path: &std::path::Path, shutdown: &Shutdown, - dialer: Dialer, ) -> CatchUp { if ibd_targets.is_empty() { info!("ibd: no outbound peers; serving only (use --connect or seeds)"); @@ -1138,7 +1134,7 @@ async fn run_ibd_or_skip( // could deliver mid-chain blocks). Default 30s is enough. stall: std::time::Duration::from_secs(30), peers: Some(std::sync::Arc::clone(shared_peers)), - dialer, + dialer: node.dialer(), ..IbdConfig::default() }; info!( From 372720268a71ad3bc641952adcf69b9cd34f91db Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 09:16:09 -0700 Subject: [PATCH 09/10] test: exercise Dialer::connect_domain SOCKS and direct Coverage CRAP flagged the unused Dialer domain path. The SOCKS encoder was already tested; this drives the production Dialer match arms. Co-authored-by: Cursor --- crates/rbitcoin-net/src/socks.rs | 89 ++++++++++++++++++++++++++++++++ 1 file changed, 89 insertions(+) diff --git a/crates/rbitcoin-net/src/socks.rs b/crates/rbitcoin-net/src/socks.rs index e19b7575e..33667eed0 100644 --- a/crates/rbitcoin-net/src/socks.rs +++ b/crates/rbitcoin-net/src/socks.rs @@ -305,6 +305,95 @@ mod tests { server.await.unwrap(); } + async fn accept_domain_connect( + listener: TcpListener, + want_host: &'static [u8], + want_port: u16, + ) { + let (mut s, _) = listener.accept().await.unwrap(); + let mut ver_n = [0u8; 2]; + s.read_exact(&mut ver_n).await.unwrap(); + let nmethods = ver_n[1] as usize; + let mut methods = vec![0u8; nmethods]; + s.read_exact(&mut methods).await.unwrap(); + if methods.contains(&0x02) { + s.write_all(&[5, 0x02]).await.unwrap(); + let mut ver = [0u8; 1]; + s.read_exact(&mut ver).await.unwrap(); + let mut ulen = [0u8; 1]; + s.read_exact(&mut ulen).await.unwrap(); + let mut user = vec![0u8; ulen[0] as usize]; + s.read_exact(&mut user).await.unwrap(); + let mut plen = [0u8; 1]; + s.read_exact(&mut plen).await.unwrap(); + let mut pass = vec![0u8; plen[0] as usize]; + s.read_exact(&mut pass).await.unwrap(); + s.write_all(&[1, 0]).await.unwrap(); + } else { + s.write_all(&[5, 0x00]).await.unwrap(); + } + let mut hdr = [0u8; 4]; + s.read_exact(&mut hdr).await.unwrap(); + assert_eq!(hdr[3], 3, "ATYP domain"); + let mut n = [0u8; 1]; + s.read_exact(&mut n).await.unwrap(); + let mut name = vec![0u8; n[0] as usize]; + s.read_exact(&mut name).await.unwrap(); + let mut p = [0u8; 2]; + s.read_exact(&mut p).await.unwrap(); + assert_eq!(name, want_host); + assert_eq!(u16::from_be_bytes(p), want_port); + s.write_all(&[5, 0, 0, 1, 0, 0, 0, 0, 0, 0]).await.unwrap(); + } + + #[tokio::test] + async fn dialer_connect_domain_covers_socks_and_direct() { + let host = "seed.example"; + let port = 8333u16; + + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let proxy = listener.local_addr().unwrap(); + let server = tokio::spawn(accept_domain_connect(listener, b"seed.example", port)); + Dialer::Socks { + proxy, + randomize: true, + } + .connect_domain(host, port) + .await + .unwrap(); + server.await.unwrap(); + + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let proxy = listener.local_addr().unwrap(); + let server = tokio::spawn(accept_domain_connect(listener, b"seed.example", port)); + Dialer::Socks { + proxy, + randomize: false, + } + .connect_domain(host, port) + .await + .unwrap(); + server.await.unwrap(); + + let echo = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let echo_addr = echo.local_addr().unwrap(); + let server = tokio::spawn(async move { + let (mut s, _) = echo.accept().await.unwrap(); + let mut b = [0u8; 1]; + s.read_exact(&mut b).await.unwrap(); + s.write_all(&b).await.unwrap(); + }); + let mut stream = Dialer::Direct + .connect_domain("127.0.0.1", echo_addr.port()) + .await + .unwrap(); + stream.write_all(&[0x42]).await.unwrap(); + let mut got = [0u8; 1]; + stream.read_exact(&mut got).await.unwrap(); + assert_eq!(got, [0x42]); + server.await.unwrap(); + } + async fn serve_userpass_ipv4( s: &mut tokio::net::TcpStream, want_ip: [u8; 4], From 1b0244faa6c6c23c47626cf7cff65aa4fff7f3db Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 21:47:59 -0700 Subject: [PATCH 10/10] net: stabilize non-randomized SOCKS creds and proxy seed bootstrap Co-authored-by: Cursor --- crates/rbitcoin-net/src/lib.rs | 4 +- crates/rbitcoin-net/src/peers.rs | 51 +++++++++- crates/rbitcoin-net/src/service.rs | 40 +++++--- crates/rbitcoin-net/src/socks.rs | 150 ++++++++++++++++++++++------- crates/rbitcoin-node/src/cli.rs | 12 ++- crates/rbitcoin-node/src/config.rs | 5 +- crates/rbitcoin-node/src/run.rs | 33 ++++++- 7 files changed, 227 insertions(+), 68 deletions(-) diff --git a/crates/rbitcoin-net/src/lib.rs b/crates/rbitcoin-net/src/lib.rs index 910ec0066..f17c57557 100644 --- a/crates/rbitcoin-net/src/lib.rs +++ b/crates/rbitcoin-net/src/lib.rs @@ -50,8 +50,8 @@ pub use peer::{ }; pub use peer_dos::DEFAULT_MAX_INBOUND; pub use peers::{ - parse_peer_addr, pick_stale_follow_evict, DialRequest, LivePeer, PeerConnType, PeerHub, - PeerInfo, PeerOut, PingAction, + parse_peer_addr, pick_stale_follow_evict, DialRequest, DialTarget, LivePeer, PeerConnType, + PeerHub, PeerInfo, PeerOut, PingAction, }; pub use rbitcoin_mempool::AcceptError; pub(crate) use rbitcoin_mempool::MempoolGraphStats; diff --git a/crates/rbitcoin-net/src/peers.rs b/crates/rbitcoin-net/src/peers.rs index d10123e10..9305b5d8f 100644 --- a/crates/rbitcoin-net/src/peers.rs +++ b/crates/rbitcoin-net/src/peers.rs @@ -123,10 +123,34 @@ pub(crate) fn trying_connection_log(typ: PeerConnType, addr: impl std::fmt::Disp format!("p2p: trying connection ({}) to {addr}", typ.as_str()) } +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum DialTarget { + Socket(SocketAddr), + Domain { host: String, port: u16 }, +} + +impl DialTarget { + pub fn peer_hint(&self) -> SocketAddr { + match self { + Self::Socket(addr) => *addr, + Self::Domain { port, .. } => SocketAddr::from(([0, 0, 0, 0], *port)), + } + } +} + +impl std::fmt::Display for DialTarget { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Socket(addr) => write!(f, "{addr}"), + Self::Domain { host, port } => write!(f, "{host}:{port}"), + } + } +} + /// Request that the node dial `addr` as `typ`. #[derive(Clone, Debug)] pub struct DialRequest { - pub addr: SocketAddr, + pub target: DialTarget, pub typ: PeerConnType, } @@ -1980,8 +2004,29 @@ impl PeerHub { pub fn dial(&self, addr: SocketAddr, typ: PeerConnType) -> Result<(), String> { let g = self.dial_tx.lock().unwrap_or_else(|e| e.into_inner()); let tx = g.as_ref().ok_or("no dialer attached")?; - tx.send(DialRequest { addr, typ }) - .map_err(|_| "dialer closed".to_string()) + tx.send(DialRequest { + target: DialTarget::Socket(addr), + typ, + }) + .map_err(|_| "dialer closed".to_string()) + } + + pub fn dial_domain( + &self, + host: impl Into, + port: u16, + typ: PeerConnType, + ) -> Result<(), String> { + let g = self.dial_tx.lock().unwrap_or_else(|e| e.into_inner()); + let tx = g.as_ref().ok_or("no dialer attached")?; + tx.send(DialRequest { + target: DialTarget::Domain { + host: host.into(), + port, + }, + typ, + }) + .map_err(|_| "dialer closed".to_string()) } } diff --git a/crates/rbitcoin-net/src/service.rs b/crates/rbitcoin-net/src/service.rs index a19944036..f003f83a4 100644 --- a/crates/rbitcoin-net/src/service.rs +++ b/crates/rbitcoin-net/src/service.rs @@ -9,7 +9,7 @@ use crate::peer::{ FollowSessionMeta, HandshakePolicy, HANDSHAKE_TIMEOUT, }; use crate::peer_dos::{inbound_semaphore, DEFAULT_MAX_INBOUND}; -use crate::peers::{DialRequest, LivePeer, PeerConnType, PeerHub}; +use crate::peers::{DialRequest, DialTarget, LivePeer, PeerConnType, PeerHub}; use crate::v2::{V2Reader, V2Writer}; use bitcoin::p2p::Magic; use bitcoin::Block; @@ -150,7 +150,7 @@ impl P2PNode { let (ah_tx, ah_rx) = tokio::sync::oneshot::channel::(); let h = tokio::spawn(async move { let _ = run_outbound_session_with_abort( - req.addr, magic, local_addr, hub, peers, ua, live, req.typ, ah_rx, d, + req.target, magic, local_addr, hub, peers, ua, live, req.typ, ah_rx, d, ) .await; }); @@ -264,7 +264,7 @@ impl P2PNode { /// Call [`Self::sync`] first when far behind (multi-thousand height IBD). pub async fn follow_from(&mut self, peer: SocketAddr) -> Result<(), NetError> { let prepared = prepare_outbound_session( - peer, + DialTarget::Socket(peer), self.magic, self.local_addr, self.hub.clone(), @@ -465,7 +465,7 @@ fn default_user_agent() -> String { } struct PreparedOutbound { - peer: SocketAddr, + peer_hint: SocketAddr, magic: Magic, hub: Arc, peers: Arc, @@ -478,7 +478,7 @@ struct PreparedOutbound { #[allow(clippy::too_many_arguments)] // call-site args stay unbundled async fn prepare_outbound_session( - peer: SocketAddr, + peer: DialTarget, magic: Magic, local: SocketAddr, hub: Arc, @@ -488,20 +488,24 @@ async fn prepare_outbound_session( typ: PeerConnType, dialer: crate::socks::Dialer, ) -> Result { - rbitcoin_log::debug!("{}", crate::peers::trying_connection_log(typ, peer)); - let stream = dialer.connect(peer).await?; + rbitcoin_log::debug!("{}", crate::peers::trying_connection_log(typ, &peer)); + let stream = match &peer { + DialTarget::Socket(addr) => dialer.connect(*addr).await?, + DialTarget::Domain { host, port } => dialer.connect_domain(host, *port).await?, + }; + let peer_hint = peer.peer_hint(); let bind = stream.local_addr().unwrap_or(local); let height = hub.tip_height().map(|h| h as i32).unwrap_or(0); // Core adds CNode before VERSION. Provisional row so getpeerinfo is non-empty // during handshake (p2p_handshake self-connect wait_until + assert_debug_log). - let provisional = peers.register_connecting(peer, bind, false, typ); + let provisional = peers.register_connecting(peer_hint, bind, false, typ); let provisional_id = provisional.id; let handshake = connect_and_handshake_timed( HANDSHAKE_TIMEOUT, stream, magic, local, - peer, + peer_hint, height, false, &user_agent, @@ -523,7 +527,7 @@ async fn prepare_outbound_session( let wants_addrv2 = provisional.wants_addrv2(); let wtxid_relay = provisional.wtxid_relay(); peers.unregister(provisional_id); - let sess = peers.register_with_id(provisional_id, peer, bind, &ver, false, typ); + let sess = peers.register_with_id(provisional_id, peer_hint, bind, &ver, false, typ); sess.mark_handshake_complete(); if wants_addrv2 { sess.set_wants_addrv2(); @@ -541,7 +545,7 @@ async fn prepare_outbound_session( let id = sess.id; follow_live.fetch_add(1, Ordering::SeqCst); Ok(PreparedOutbound { - peer, + peer_hint, magic, hub, peers, @@ -556,7 +560,7 @@ async fn prepare_outbound_session( async fn run_prepared_outbound(prepared: PreparedOutbound) -> Result<(), NetError> { let tip_rx = prepared.hub.subscribe_tips(); let meta = FollowSessionMeta { - peer: Some(prepared.peer), + peer: Some(prepared.peer_hint), live: Some(prepared.follow_live), session: Some(prepared.sess), }; @@ -575,7 +579,7 @@ async fn run_prepared_outbound(prepared: PreparedOutbound) -> Result<(), NetErro #[allow(clippy::too_many_arguments)] // call-site args stay unbundled async fn run_outbound_session_with_abort( - peer: SocketAddr, + peer: DialTarget, magic: Magic, local: SocketAddr, hub: Arc, @@ -587,9 +591,15 @@ async fn run_outbound_session_with_abort( dialer: crate::socks::Dialer, ) -> Result<(), NetError> { if typ == PeerConnType::Feeler { - let stream = dialer.connect(peer).await?; + let peer_addr = match peer { + DialTarget::Socket(addr) => addr, + DialTarget::Domain { .. } => { + return Err(NetError::Encode("feeler requires ip:port target".into())) + } + }; + let stream = dialer.connect(peer_addr).await?; let height = hub.tip_height().map(|h| h as i32).unwrap_or(0); - return crate::peer::run_feeler(stream, magic, local, peer, height, &user_agent).await; + return crate::peer::run_feeler(stream, magic, local, peer_addr, height, &user_agent).await; } let prepared = prepare_outbound_session( peer, diff --git a/crates/rbitcoin-net/src/socks.rs b/crates/rbitcoin-net/src/socks.rs index 33667eed0..ae06a288c 100644 --- a/crates/rbitcoin-net/src/socks.rs +++ b/crates/rbitcoin-net/src/socks.rs @@ -2,9 +2,11 @@ use crate::error::NetError; use std::net::SocketAddr; +use std::sync::Arc; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpStream; +#[derive(Clone, Debug, PartialEq, Eq)] pub(crate) struct ProxyCreds { pub(crate) username: Vec, pub(crate) password: Vec, @@ -57,20 +59,47 @@ pub enum Dialer { Socks { proxy: SocketAddr, randomize: bool, + shared_creds: Option, Vec)>>, }, } impl Dialer { + pub fn socks(proxy: SocketAddr, randomize: bool) -> Self { + let shared_creds = (!randomize).then(|| { + let creds = ProxyCreds::fresh(); + Arc::new((creds.username, creds.password)) + }); + Self::Socks { + proxy, + randomize, + shared_creds, + } + } + + fn dial_creds( + randomize: bool, + shared_creds: Option<&Arc<(Vec, Vec)>>, + ) -> Option { + if randomize { + Some(ProxyCreds::fresh()) + } else { + shared_creds.map(|c| ProxyCreds { + username: c.0.clone(), + password: c.1.clone(), + }) + } + } + pub async fn connect(&self, target: SocketAddr) -> Result { match self { Dialer::Direct => Ok(TcpStream::connect(target).await?), - Dialer::Socks { proxy, randomize } => { - if *randomize { - let creds = ProxyCreds::fresh(); - socks5_connect(*proxy, target, Some(&creds)).await - } else { - socks5_connect(*proxy, target, None).await - } + Dialer::Socks { + proxy, + randomize, + shared_creds, + } => { + let creds = Self::dial_creds(*randomize, shared_creds.as_ref()); + socks5_connect(*proxy, target, creds.as_ref()).await } } } @@ -82,13 +111,13 @@ impl Dialer { let addr = addrs.next().ok_or(NetError::Protocol("dns lookup empty"))?; self.connect(addr).await } - Dialer::Socks { proxy, randomize } => { - if *randomize { - let creds = ProxyCreds::fresh(); - socks5_connect_domain(*proxy, host, port, Some(&creds)).await - } else { - socks5_connect_domain(*proxy, host, port, None).await - } + Dialer::Socks { + proxy, + randomize, + shared_creds, + } => { + let creds = Self::dial_creds(*randomize, shared_creds.as_ref()); + socks5_connect_domain(*proxy, host, port, creds.as_ref()).await } } } @@ -132,12 +161,17 @@ async fn greet(s: &mut TcpStream, creds: Option<&ProxyCreds>) -> Result<(), NetE { return Err(NetError::Protocol("socks username/password length")); } - s.write_all(&[5, 1, 0x02]).await?; + s.write_all(&[5, 2, 0x00, 0x02]).await?; let mut sel = [0u8; 2]; s.read_exact(&mut sel).await?; - if sel[0] != 5 || sel[1] != 0x02 { + if sel[0] != 5 { return Err(NetError::Protocol("socks method rejected")); } + match sel[1] { + 0x00 => return Ok(()), + 0x02 => {} + _ => return Err(NetError::Protocol("socks method rejected")), + } let mut auth = Vec::with_capacity(3 + c.username.len() + c.password.len()); auth.push(1); auth.push(c.username.len() as u8); @@ -305,6 +339,37 @@ mod tests { server.await.unwrap(); } + #[tokio::test] + async fn randomize_off_with_creds_still_connects_on_noauth_proxy() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let proxy = listener.local_addr().unwrap(); + let target = SocketAddr::from((Ipv4Addr::new(203, 0, 113, 7), 8333)); + let server = tokio::spawn(async move { + let (mut s, _) = listener.accept().await.unwrap(); + let mut ver_n = [0u8; 2]; + s.read_exact(&mut ver_n).await.unwrap(); + let nmethods = ver_n[1] as usize; + let mut methods = vec![0u8; nmethods]; + s.read_exact(&mut methods).await.unwrap(); + assert!( + methods.contains(&0x00) && methods.contains(&0x02), + "{methods:?}" + ); + s.write_all(&[5, 0x00]).await.unwrap(); + let mut hdr = [0u8; 4]; + s.read_exact(&mut hdr).await.unwrap(); + assert_eq!(hdr[3], 1); + let mut addr = [0u8; 4]; + s.read_exact(&mut addr).await.unwrap(); + let mut port = [0u8; 2]; + s.read_exact(&mut port).await.unwrap(); + s.write_all(&[5, 0, 0, 1, 0, 0, 0, 0, 0, 0]).await.unwrap(); + }); + let d = Dialer::socks(proxy, false); + d.connect(target).await.unwrap(); + server.await.unwrap(); + } + async fn accept_domain_connect( listener: TcpListener, want_host: &'static [u8], @@ -354,25 +419,19 @@ mod tests { let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let proxy = listener.local_addr().unwrap(); let server = tokio::spawn(accept_domain_connect(listener, b"seed.example", port)); - Dialer::Socks { - proxy, - randomize: true, - } - .connect_domain(host, port) - .await - .unwrap(); + Dialer::socks(proxy, true) + .connect_domain(host, port) + .await + .unwrap(); server.await.unwrap(); let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let proxy = listener.local_addr().unwrap(); let server = tokio::spawn(accept_domain_connect(listener, b"seed.example", port)); - Dialer::Socks { - proxy, - randomize: false, - } - .connect_domain(host, port) - .await - .unwrap(); + Dialer::socks(proxy, false) + .connect_domain(host, port) + .await + .unwrap(); server.await.unwrap(); let echo = TcpListener::bind("127.0.0.1:0").await.unwrap(); @@ -475,6 +534,29 @@ mod tests { assert!(!u1.is_empty() && !u2.is_empty()); } + #[tokio::test] + async fn randomize_off_reuses_stable_userpass_per_dialer() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let proxy = listener.local_addr().unwrap(); + let target = SocketAddr::from((Ipv4Addr::new(198, 51, 100, 3), 8333)); + let (tx, mut rx) = tokio::sync::mpsc::channel(2); + let server = tokio::spawn(async move { + for _ in 0..2 { + let (mut s, _) = listener.accept().await.unwrap(); + let creds = serve_userpass_ipv4(&mut s, [198, 51, 100, 3], 8333).await; + tx.send(creds).await.unwrap(); + } + }); + let d = Dialer::socks(proxy, false); + d.connect(target).await.unwrap(); + d.connect(target).await.unwrap(); + let (u1, p1) = rx.recv().await.unwrap(); + let (u2, p2) = rx.recv().await.unwrap(); + server.await.unwrap(); + assert_eq!(u1, u2, "randomize=0 must keep one SOCKS username"); + assert_eq!(p1, p2, "randomize=0 must keep one SOCKS password"); + } + async fn splice_one_socks( listener: TcpListener, saw: tokio::sync::oneshot::Sender, @@ -549,13 +631,7 @@ mod tests { let (saw_tx, saw_rx) = tokio::sync::oneshot::channel(); let splice = tokio::spawn(splice_one_socks(socks_l, saw_tx)); - let stream = Dialer::Socks { - proxy, - randomize: true, - } - .connect(peer_addr) - .await - .unwrap(); + let stream = Dialer::socks(proxy, true).connect(peer_addr).await.unwrap(); assert_eq!(saw_rx.await.unwrap(), peer_addr); connect_and_handshake_timed( diff --git a/crates/rbitcoin-node/src/cli.rs b/crates/rbitcoin-node/src/cli.rs index 072b3b5ed..4ca737670 100644 --- a/crates/rbitcoin-node/src/cli.rs +++ b/crates/rbitcoin-node/src/cli.rs @@ -627,13 +627,15 @@ mod tests { let cfg = ready_config(["rbitcoin-node", "--proxy", "127.0.0.1:9050"]); assert_eq!(cfg.listen.proxy, Some("127.0.0.1:9050".parse().unwrap())); assert!(cfg.listen.onion.is_none()); - assert_eq!( - cfg.listen.dialer(), + match cfg.listen.dialer() { rbitcoin_net::Dialer::Socks { - proxy: "127.0.0.1:9050".parse().unwrap(), - randomize: true, + proxy, randomize, .. + } => { + assert_eq!(proxy, "127.0.0.1:9050".parse().unwrap()); + assert!(randomize); } - ); + other => panic!("expected socks dialer, got {other:?}"), + } assert_eq!( NodeConfig::default().listen.dialer(), rbitcoin_net::Dialer::Direct diff --git a/crates/rbitcoin-node/src/config.rs b/crates/rbitcoin-node/src/config.rs index 3361cfdc0..4fab7d61d 100644 --- a/crates/rbitcoin-node/src/config.rs +++ b/crates/rbitcoin-node/src/config.rs @@ -118,10 +118,7 @@ impl ListenOpts { pub fn dialer(&self) -> rbitcoin_net::Dialer { match self.proxy { None => rbitcoin_net::Dialer::Direct, - Some(proxy) => rbitcoin_net::Dialer::Socks { - proxy, - randomize: self.proxy_randomize, - }, + Some(proxy) => rbitcoin_net::Dialer::socks(proxy, self.proxy_randomize), } } } diff --git a/crates/rbitcoin-node/src/run.rs b/crates/rbitcoin-node/src/run.rs index 164a427c4..71cefeb2b 100644 --- a/crates/rbitcoin-node/src/run.rs +++ b/crates/rbitcoin-node/src/run.rs @@ -393,8 +393,8 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { addrman.len() ); } else if config.listen.proxy.is_some() && config.listen.use_seeds { - let n = socks_dns_seed_dests(config.network).len(); - info!("ibd: SOCKS proxy set — skipping local DNS for {n} seed hostnames (use --connect)"); + let n = queue_proxy_seed_addrfetch(&node.peers, config.network); + info!("ibd: SOCKS proxy set — queued {n} seed hostnames via SOCKS addrfetch"); } else if config.signet_challenge.is_some() && config.listen.connect.is_empty() && addrman.is_empty() @@ -996,6 +996,17 @@ fn should_resolve_default_seeds(config: &NodeConfig) -> bool { && config.listen.proxy.is_none() } +fn queue_proxy_seed_addrfetch(peers: &Arc, network: Network) -> usize { + let mut n = 0usize; + for (host, port) in socks_dns_seed_dests(network) { + match peers.dial_domain(host.clone(), port, PeerConnType::AddrFetch) { + Ok(()) => n += 1, + Err(e) => warn!("seednode dial {host}:{port}: {e}"), + } + } + n +} + /// One walker per process: SH-warm start and post-IBD `enter_tip_mode` both call this. fn spawn_sptweaks_backfill( query: Arc, @@ -1933,6 +1944,24 @@ mod tests { ); } + #[test] + fn proxy_seed_bootstrap_queues_domain_addrfetch() { + let peers = rbitcoin_net::PeerHub::new(); + let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::(); + peers.set_dialer(tx); + let n = queue_proxy_seed_addrfetch(&peers, rbitcoin_primitives::Network::Signet); + assert_eq!(n, 1, "signet has one default DNS seed"); + let req = rx.try_recv().expect("queued dial request"); + assert_eq!(req.typ, PeerConnType::AddrFetch); + match req.target { + rbitcoin_net::DialTarget::Domain { host, port } => { + assert_eq!(host, "seed.signet.bitcoin.sprovoost.nl"); + assert_eq!(port, 38333); + } + other => panic!("expected domain target, got {other:?}"), + } + } + fn coinbase_block( h: u32, prev: rbitcoin_primitives::Fk,