From 22b8d0b54d771faea624a4bcf6c47b4642fa4b7e Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 11:41:51 -0700 Subject: [PATCH 1/9] net: parse BIP155 I2P destinations as NetAddr Display is Core-shaped {52}.b32.i2p:port. SOCKS still refuses I2P until SAM exists; Direct onion path is unchanged. Co-authored-by: Cursor --- crates/rbitcoin-net/src/netaddr.rs | 120 +++++++++++++++++++++++++++-- crates/rbitcoin-net/src/peers.rs | 8 +- crates/rbitcoin-net/src/socks.rs | 3 + 3 files changed, 119 insertions(+), 12 deletions(-) diff --git a/crates/rbitcoin-net/src/netaddr.rs b/crates/rbitcoin-net/src/netaddr.rs index cbdad3ab0..83f5b123f 100644 --- a/crates/rbitcoin-net/src/netaddr.rs +++ b/crates/rbitcoin-net/src/netaddr.rs @@ -10,6 +10,8 @@ use std::str::FromStr; const B32: &[u8; 32] = b"abcdefghijklmnopqrstuvwxyz234567"; const ONION_VERSION: u8 = 3; const ONION_NAME_LEN: usize = 56; +const I2P_B32_LEN: usize = 52; +const I2P_SUFFIX: &str = ".b32.i2p"; #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub enum OnlyNet { @@ -47,6 +49,7 @@ pub fn addr_allowed(addr: NetAddr, only: &[OnlyNet]) -> bool { pub enum NetAddr { Ip(SocketAddr), Onion { pk: [u8; 32], port: u16 }, + I2p { dest: [u8; 32], port: u16 }, } impl fmt::Display for NetAddr { @@ -56,6 +59,9 @@ impl fmt::Display for NetAddr { NetAddr::Onion { pk, port } => { write!(f, "{}.onion:{port}", encode_onion_name(&pk)) } + NetAddr::I2p { dest, port } => { + write!(f, "{}{I2P_SUFFIX}:{port}", encode_i2p_name(&dest)) + } } } } @@ -79,28 +85,32 @@ impl NetAddr { pk: *pk, port: msg.port, }), - AddrV2::TorV2(_) | AddrV2::I2p(_) | AddrV2::Cjdns(_) | AddrV2::Unknown(_, _) => None, + AddrV2::I2p(dest) => Some(NetAddr::I2p { + dest: *dest, + port: msg.port, + }), + AddrV2::TorV2(_) | AddrV2::Cjdns(_) | AddrV2::Unknown(_, _) => None, } } pub fn socket_addr(self) -> Option { match self { NetAddr::Ip(s) => Some(s), - NetAddr::Onion { .. } => None, + NetAddr::Onion { .. } | NetAddr::I2p { .. } => None, } } pub fn is_ipv6(self) -> bool { match self { NetAddr::Ip(s) => s.is_ipv6(), - NetAddr::Onion { .. } => false, + NetAddr::Onion { .. } | NetAddr::I2p { .. } => false, } } pub fn port(self) -> u16 { match self { NetAddr::Ip(s) => s.port(), - NetAddr::Onion { port, .. } => port, + NetAddr::Onion { port, .. } | NetAddr::I2p { port, .. } => port, } } @@ -108,6 +118,7 @@ impl NetAddr { match self { NetAddr::Ip(s) => s.ip().to_string(), NetAddr::Onion { pk, .. } => format!("{}.onion", encode_onion_name(&pk)), + NetAddr::I2p { dest, .. } => format!("{}{I2P_SUFFIX}", encode_i2p_name(&dest)), } } @@ -116,6 +127,7 @@ impl NetAddr { NetAddr::Ip(s) if s.is_ipv4() => "ipv4", NetAddr::Ip(_) => "ipv6", NetAddr::Onion { .. } => "onion", + NetAddr::I2p { .. } => "i2p", } } } @@ -124,6 +136,14 @@ fn parse_net_addr(s: &str) -> Result { let Some((host, port_s)) = s.rsplit_once(':') else { return Err(NetError::Encode(format!("bad peer address {s}"))); }; + if let Some(name) = strip_i2p_suffix(host) { + let port: u16 = port_s + .parse() + .map_err(|_| NetError::Encode(format!("bad i2p port {s}")))?; + let dest = decode_i2p_name(name) + .ok_or_else(|| NetError::Encode(format!("bad i2p address {s}")))?; + return Ok(NetAddr::I2p { dest, port }); + } if let Some(name) = strip_onion_suffix(host) { let port: u16 = port_s .parse() @@ -137,6 +157,68 @@ fn parse_net_addr(s: &str) -> Result { .map_err(|_| NetError::Encode(format!("bad peer address {s}"))) } +fn strip_i2p_suffix(host: &str) -> Option<&str> { + let b = host.as_bytes(); + let suf = I2P_SUFFIX.as_bytes(); + if b.len() <= suf.len() { + return None; + } + if !b[b.len() - suf.len()..].eq_ignore_ascii_case(suf) { + return None; + } + Some(&host[..host.len() - suf.len()]) +} + +fn encode_i2p_name(dest: &[u8; 32]) -> String { + let mut out = String::with_capacity(I2P_B32_LEN); + let mut acc = 0u32; + let mut bits = 0u32; + for &b in dest { + acc = (acc << 8) | u32::from(b); + bits += 8; + while bits >= 5 { + bits -= 5; + out.push(B32[((acc >> bits) & 31) as usize] as char); + } + } + if bits > 0 { + out.push(B32[((acc << (5 - bits)) & 31) as usize] as char); + } + out +} + +fn decode_i2p_name(name: &str) -> Option<[u8; 32]> { + if name.len() != I2P_B32_LEN { + return None; + } + let mut out = [0u8; 32]; + let mut acc = 0u32; + let mut bits = 0u32; + let mut n = 0usize; + for c in name.bytes() { + let v = match c { + b'a'..=b'z' => c - b'a', + b'A'..=b'Z' => c - b'A', + b'2'..=b'7' => 26 + (c - b'2'), + _ => return None, + }; + acc = (acc << 5) | u32::from(v); + bits += 5; + if bits >= 8 { + bits -= 8; + if n >= 32 { + return None; + } + out[n] = (acc >> bits) as u8; + n += 1; + } + } + if n != 32 { + return None; + } + Some(out) +} + fn strip_onion_suffix(host: &str) -> Option<&str> { let b = host.as_bytes(); if b.len() < 6 { @@ -249,7 +331,7 @@ mod tests { ); assert_eq!(port, 8333); } - NetAddr::Ip(_) => panic!("expected onion"), + NetAddr::Ip(_) | NetAddr::I2p { .. } => panic!("expected onion"), } assert!( "qg6mmjiyjmcrsslvykfwnntlaru7p5svn6y2ymmju6nubxndf4pscryd.onion:8333" @@ -275,4 +357,32 @@ mod tests { ); assert_eq!(a.to_string(), s); } + + #[test] + fn netaddr_i2p_addrv2_roundtrip() { + let dest = [0x11u8; 32]; + let msg = AddrV2Message { + time: 1, + services: bitcoin::p2p::ServiceFlags::NETWORK, + addr: AddrV2::I2p(dest), + port: 8333, + }; + let a = NetAddr::from_addrv2(&msg).expect("i2p addrv2"); + assert_eq!(a, NetAddr::I2p { dest, port: 8333 }); + let s = a.to_string(); + assert!(s.ends_with(".b32.i2p:8333"), "{s}"); + assert_eq!(s.parse::().unwrap(), a); + assert_eq!(a.network_label(), "i2p"); + assert_eq!(a.port(), 8333); + assert!(a.socket_addr().is_none()); + let zeros = NetAddr::I2p { + dest: [0u8; 32], + port: 1, + }; + assert_eq!( + zeros.to_string(), + "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa.b32.i2p:1" + ); + assert!("short.b32.i2p:1".parse::().is_err()); + } } diff --git a/crates/rbitcoin-net/src/peers.rs b/crates/rbitcoin-net/src/peers.rs index 1936674bc..074cd64bb 100644 --- a/crates/rbitcoin-net/src/peers.rs +++ b/crates/rbitcoin-net/src/peers.rs @@ -3140,12 +3140,6 @@ mod tests { addr: AddrV2::Ipv4(Ipv4Addr::new(1, 2, 3, 4)), port: 18444, }, - AddrV2Message { - time: 1, - services: ServiceFlags::NETWORK, - addr: AddrV2::I2p([0u8; 32]), - port: 1, - }, ]); let book = am.lock().unwrap_or_else(|e| e.into_inner()); let ents = book.entries(); @@ -3157,7 +3151,7 @@ mod tests { .iter() .any(|e| e.addr == crate::NetAddr::Ip(SocketAddr::from((Ipv4Addr::new(1, 2, 3, 4), 18444))))); - assert_eq!(ents.len(), 2, "I2P is not stored until plan 04"); + assert_eq!(ents.len(), 2); } #[test] diff --git a/crates/rbitcoin-net/src/socks.rs b/crates/rbitcoin-net/src/socks.rs index bfbfdf954..d505d299f 100644 --- a/crates/rbitcoin-net/src/socks.rs +++ b/crates/rbitcoin-net/src/socks.rs @@ -140,6 +140,9 @@ impl Dialer { } self.connect_domain(&addr.host_str(), port).await } + crate::NetAddr::I2p { .. } => { + Err(NetError::Encode("i2p dial requires SAM (--i2p-sam)".into())) + } } } } From 9e82e19658da2cfb177d862fb367ff8fb16efb22 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 11:45:00 -0700 Subject: [PATCH 2/9] net: SAM v3 STREAM CONNECT against a fake router HELLO + SESSION CREATE once per process; each STREAM CONNECT is a new SAM socket with the session id. Bad HELLO is an error. No live i2pd. Co-authored-by: Cursor --- crates/rbitcoin-net/src/i2p_sam.rs | 177 +++++++++++++++++++++++++++++ crates/rbitcoin-net/src/lib.rs | 1 + 2 files changed, 178 insertions(+) create mode 100644 crates/rbitcoin-net/src/i2p_sam.rs diff --git a/crates/rbitcoin-net/src/i2p_sam.rs b/crates/rbitcoin-net/src/i2p_sam.rs new file mode 100644 index 000000000..057f08112 --- /dev/null +++ b/crates/rbitcoin-net/src/i2p_sam.rs @@ -0,0 +1,177 @@ +//! I2P SAM v3 STREAM CONNECT (system router, not SOCKS). + +use crate::error::NetError; +use std::net::SocketAddr; +use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; +use tokio::net::TcpStream; + +pub struct I2pSam { + sam_addr: SocketAddr, + session_id: String, + _control: TcpStream, +} + +impl I2pSam { + pub async fn connect(sam_addr: SocketAddr) -> Result { + let mut control = TcpStream::connect(sam_addr) + .await + .map_err(|e| NetError::Encode(format!("i2p sam connect {sam_addr}: {e}")))?; + hello(&mut control).await?; + let session_id = fresh_session_id(); + write_line( + &mut control, + &format!("SESSION CREATE STYLE=STREAM ID={session_id} DESTINATION=TRANSIENT"), + ) + .await?; + let reply = read_line(&mut control).await?; + if !reply.to_ascii_uppercase().contains("RESULT=OK") { + return Err(NetError::Encode(format!("i2p sam session: {reply}"))); + } + Ok(Self { + sam_addr, + session_id, + _control: control, + }) + } + + pub async fn stream_connect(&self, dest_b32: &str) -> Result { + let mut s = TcpStream::connect(self.sam_addr) + .await + .map_err(|e| NetError::Encode(format!("i2p sam stream connect: {e}")))?; + hello(&mut s).await?; + write_line( + &mut s, + &format!( + "STREAM CONNECT ID={} DESTINATION={}", + self.session_id, dest_b32 + ), + ) + .await?; + let reply = read_line(&mut s).await?; + if !reply.to_ascii_uppercase().contains("RESULT=OK") { + return Err(NetError::Encode(format!("i2p sam stream: {reply}"))); + } + Ok(s) + } +} + +fn fresh_session_id() -> String { + let mut b = [0u8; 8]; + getrandom::fill(&mut b).expect("CSPRNG for SAM session id"); + let mut id = String::from("rbtc"); + for x in b { + id.push_str(&format!("{x:02x}")); + } + id +} + +async fn hello(s: &mut TcpStream) -> Result<(), NetError> { + write_line(s, "HELLO VERSION MIN=3.1 MAX=3.3").await?; + let reply = read_line(s).await?; + let up = reply.to_ascii_uppercase(); + if !up.starts_with("HELLO REPLY") || !up.contains("RESULT=OK") { + return Err(NetError::Encode(format!("i2p sam hello: {reply}"))); + } + Ok(()) +} + +async fn write_line(s: &mut TcpStream, line: &str) -> Result<(), NetError> { + s.write_all(line.as_bytes()) + .await + .map_err(|e| NetError::Encode(format!("i2p sam write: {e}")))?; + s.write_all(b"\n") + .await + .map_err(|e| NetError::Encode(format!("i2p sam write: {e}")))?; + s.flush() + .await + .map_err(|e| NetError::Encode(format!("i2p sam write: {e}"))) +} + +async fn read_line(s: &mut TcpStream) -> Result { + let mut reader = BufReader::new(s); + let mut line = String::new(); + let n = reader + .read_line(&mut line) + .await + .map_err(|e| NetError::Encode(format!("i2p sam read: {e}")))?; + if n == 0 { + return Err(NetError::Encode("i2p sam: connection closed".into())); + } + Ok(line.trim_end_matches(['\r', '\n']).to_string()) +} + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::{Arc, Mutex}; + use tokio::net::TcpListener; + + async fn fake_sam(ok_hello: bool, dest_log: Arc>>) -> SocketAddr { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { + loop { + let Ok((mut s, _)) = listener.accept().await else { + break; + }; + let log = Arc::clone(&dest_log); + let ok = ok_hello; + tokio::spawn(async move { + loop { + let line = match read_line(&mut s).await { + Ok(l) => l, + Err(_) => break, + }; + let up = line.to_ascii_uppercase(); + if up.starts_with("HELLO VERSION") { + if ok { + let _ = + write_line(&mut s, "HELLO REPLY RESULT=OK VERSION=3.1").await; + } else { + let _ = write_line(&mut s, "HELLO REPLY RESULT=NOVERSION").await; + break; + } + } else if up.starts_with("SESSION CREATE") { + let _ = write_line(&mut s, "SESSION STATUS RESULT=OK DESTINATION=fake") + .await; + } else if let Some(rest) = line.strip_prefix("STREAM CONNECT ") { + log.lock().unwrap().push(rest.to_string()); + let _ = write_line(&mut s, "STREAM STATUS RESULT=OK").await; + break; + } else { + let _ = write_line(&mut s, "PING").await; + } + } + }); + } + }); + addr + } + + #[tokio::test] + async fn i2p_sam_stream_connect_fake() { + let log = Arc::new(Mutex::new(Vec::new())); + let addr = fake_sam(true, Arc::clone(&log)).await; + let sam = I2pSam::connect(addr).await.unwrap(); + let dest = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa.b32.i2p"; + sam.stream_connect(dest).await.unwrap(); + sam.stream_connect(dest).await.unwrap(); + let got = log.lock().unwrap().clone(); + assert_eq!(got.len(), 2, "{got:?}"); + for g in &got { + assert!(g.contains(dest), "{g}"); + assert!(g.contains("ID=rbtc"), "{g}"); + } + + let bad = fake_sam(false, Arc::new(Mutex::new(Vec::new()))).await; + let err = match I2pSam::connect(bad).await { + Err(e) => e, + Ok(_) => panic!("bad HELLO must fail"), + }; + let msg = format!("{err}"); + assert!( + msg.contains("hello") || msg.contains("NOVERSION") || msg.contains("sam"), + "{msg}" + ); + } +} diff --git a/crates/rbitcoin-net/src/lib.rs b/crates/rbitcoin-net/src/lib.rs index d4b7fb1c0..b0191b7ce 100644 --- a/crates/rbitcoin-net/src/lib.rs +++ b/crates/rbitcoin-net/src/lib.rs @@ -7,6 +7,7 @@ mod codec; mod compact; mod error; mod eviction; +mod i2p_sam; mod ibd; mod most_work; mod msg_decode; From c57ee2ebcc6e52b000f4d53149474935c27fb2c5 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 11:45:00 -0700 Subject: [PATCH 3/9] net: keep I2P rows in AddrMan and peers v2 learn_addrv2 already mapped AddrV2::I2p. Persist the .b32.i2p token so a restart still has the destination when SAM is configured later. Co-authored-by: Cursor --- crates/rbitcoin-net/src/peers.rs | 25 ++++++++++++++++++++++++- crates/rbitcoin-net/src/seeds.rs | 30 ++++++++++++++++++++++++++++++ 2 files changed, 54 insertions(+), 1 deletion(-) diff --git a/crates/rbitcoin-net/src/peers.rs b/crates/rbitcoin-net/src/peers.rs index 074cd64bb..f7ec98f15 100644 --- a/crates/rbitcoin-net/src/peers.rs +++ b/crates/rbitcoin-net/src/peers.rs @@ -1289,7 +1289,7 @@ impl PeerHub { *self.addrman.lock().unwrap_or_else(|e| e.into_inner()) = Some(am); } - /// Learn IPv4/IPv6/Tor v3 rows from BIP155 `addrv2` (`p2p_addrv2_relay.py`). + /// Learn IPv4/IPv6/Tor v3/I2P rows from BIP155 `addrv2`. pub fn learn_addrv2(&self, list: &[bitcoin::p2p::address::AddrV2Message]) { let g = self.addrman.lock().unwrap_or_else(|e| e.into_inner()); let Some(am) = g.as_ref() else { @@ -3167,4 +3167,27 @@ mod tests { assert!(hint.ip().is_unspecified(), "{hint}"); assert_eq!(hint.port(), 8333); } + + #[test] + fn learn_addrv2_keeps_i2p() { + use bitcoin::p2p::address::{AddrV2, AddrV2Message}; + + let hub = PeerHub::new(); + let am = Arc::new(Mutex::new(crate::seeds::AddrMan::new())); + hub.set_addrman(am.clone()); + let dest = [0x11u8; 32]; + hub.learn_addrv2(&[AddrV2Message { + time: 1, + services: ServiceFlags::NETWORK, + addr: AddrV2::I2p(dest), + port: 8333, + }]); + let book = am.lock().unwrap_or_else(|e| e.into_inner()); + let want = crate::NetAddr::I2p { dest, port: 8333 }; + assert!( + book.entries().iter().any(|e| e.addr == want), + "I2P must stay in the book, got {:?}", + book.entries() + ); + } } diff --git a/crates/rbitcoin-net/src/seeds.rs b/crates/rbitcoin-net/src/seeds.rs index cca8fe077..900951ebe 100644 --- a/crates/rbitcoin-net/src/seeds.rs +++ b/crates/rbitcoin-net/src/seeds.rs @@ -1563,6 +1563,36 @@ mod tests { let _ = std::fs::remove_dir_all(&dir); } + #[test] + fn peers_file_roundtrip_i2p() { + let dir = std::env::temp_dir().join(format!( + "rbitcoin-peers-i2p-{}-{}", + std::process::id(), + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_nanos() + )); + let _ = std::fs::create_dir_all(&dir); + let path = dir.join("peers"); + let i2p = NetAddr::I2p { + dest: [0x11u8; 32], + port: 8333, + }; + let mut am = AddrMan::new(); + am.add_addr(i2p); + am.save(&path).unwrap(); + let loaded = AddrMan::load(&path).unwrap(); + assert!( + loaded.entries().iter().any(|e| e.addr == i2p), + "i2p must persist, got {:?}", + loaded.entries() + ); + let body = std::fs::read_to_string(&path).unwrap(); + assert!(body.contains(".b32.i2p:8333"), "{body}"); + let _ = std::fs::remove_dir_all(&dir); + } + #[test] fn peers_file_v1_ipv4_still_loads() { let dir = std::env::temp_dir().join(format!( From 4dae5d31481cfe66d623cb0aed0b71eddfda4e96 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 11:50:30 -0700 Subject: [PATCH 4/9] node: --i2p-sam and --only-net=i2p Omitted SAM ADDR is 127.0.0.1:7656. only-net=i2p without --i2p-sam is a start error. Start opens one SAM session; Direct still refuses I2P dial. Co-authored-by: Cursor --- crates/rbitcoin-net/src/i2p_sam.rs | 20 ++++++++++++++ crates/rbitcoin-net/src/lib.rs | 1 + crates/rbitcoin-net/src/netaddr.rs | 5 +++- crates/rbitcoin-net/src/seeds.rs | 20 +++++++++++++- crates/rbitcoin-node/src/cli.rs | 42 ++++++++++++++++++++++++++++-- crates/rbitcoin-node/src/config.rs | 17 ++++++++++++ crates/rbitcoin-node/src/run.rs | 9 +++++++ 7 files changed, 110 insertions(+), 4 deletions(-) diff --git a/crates/rbitcoin-net/src/i2p_sam.rs b/crates/rbitcoin-net/src/i2p_sam.rs index 057f08112..22ad040cb 100644 --- a/crates/rbitcoin-net/src/i2p_sam.rs +++ b/crates/rbitcoin-net/src/i2p_sam.rs @@ -174,4 +174,24 @@ mod tests { "{msg}" ); } + + #[tokio::test] + async fn dial_i2p_uses_sam() { + let log = Arc::new(Mutex::new(Vec::new())); + let addr = fake_sam(true, Arc::clone(&log)).await; + let sam = I2pSam::connect(addr).await.unwrap(); + let peer = crate::NetAddr::I2p { + dest: [0u8; 32], + port: 8333, + }; + sam.stream_connect(&peer.host_str()).await.unwrap(); + let got = log.lock().unwrap().clone(); + assert!(got.iter().any(|g| g.contains(&peer.host_str())), "{got:?}"); + let err = match crate::socks::Dialer::Direct.connect_net(peer).await { + Err(e) => e, + Ok(_) => panic!("Direct must not dial I2P"), + }; + let msg = format!("{err}"); + assert!(msg.contains("SAM") || msg.contains("i2p"), "{msg}"); + } } diff --git a/crates/rbitcoin-net/src/lib.rs b/crates/rbitcoin-net/src/lib.rs index b0191b7ce..605fce89d 100644 --- a/crates/rbitcoin-net/src/lib.rs +++ b/crates/rbitcoin-net/src/lib.rs @@ -36,6 +36,7 @@ pub use compact::{ }; pub use error::NetError; pub use eviction::{select_inbound_eviction, InboundEvictCandidate}; +pub use i2p_sam::I2pSam; pub use ibd::{ format_tip_perf_sizes, read_proc_rss, rehydrate_block_queue_residue, IbdConfig, ProcRss, TipPerfSizes, DEFAULT_BLOCKS_IN_TRANSIT_PER_PEER, DEFAULT_IBD_WINDOW, diff --git a/crates/rbitcoin-net/src/netaddr.rs b/crates/rbitcoin-net/src/netaddr.rs index 83f5b123f..eed17a294 100644 --- a/crates/rbitcoin-net/src/netaddr.rs +++ b/crates/rbitcoin-net/src/netaddr.rs @@ -18,6 +18,7 @@ pub enum OnlyNet { Ipv4, Ipv6, Onion, + I2p, } impl OnlyNet { @@ -26,7 +27,8 @@ impl OnlyNet { "ipv4" => Ok(Self::Ipv4), "ipv6" => Ok(Self::Ipv6), "onion" => Ok(Self::Onion), - "i2p" | "cjdns" => Err(format!("unknown network {s} (not yet implemented)")), + "i2p" => Ok(Self::I2p), + "cjdns" => Err(format!("unknown network {s} (not yet implemented)")), other => Err(format!("unknown network {other}")), } } @@ -36,6 +38,7 @@ impl OnlyNet { (Self::Ipv4, NetAddr::Ip(s)) => s.is_ipv4(), (Self::Ipv6, NetAddr::Ip(s)) => s.is_ipv6(), (Self::Onion, NetAddr::Onion { .. }) => true, + (Self::I2p, NetAddr::I2p { .. }) => true, _ => false, } } diff --git a/crates/rbitcoin-net/src/seeds.rs b/crates/rbitcoin-net/src/seeds.rs index 900951ebe..47137a8c7 100644 --- a/crates/rbitcoin-net/src/seeds.rs +++ b/crates/rbitcoin-net/src/seeds.rs @@ -360,7 +360,10 @@ impl AddrMan { if out.len() >= max { break; } - if !matches!(a, NetAddr::Onion { .. }) || !self.allowed(a) || exclude.contains(&a) { + if !matches!(a, NetAddr::Onion { .. } | NetAddr::I2p { .. }) + || !self.allowed(a) + || exclude.contains(&a) + { continue; } if !out.contains(&a) { @@ -1625,4 +1628,19 @@ mod tests { assert_eq!(got, vec![onion]); assert!(am.take_dial_candidates(8, &HashSet::new(), &[]).is_empty()); } + + #[test] + fn only_net_i2p_filters_ipv4_candidates() { + let i2p = NetAddr::I2p { + dest: [0x11u8; 32], + port: 8333, + }; + let mut am = AddrMan::new(); + am.add(addr(1)); + am.add_addr(i2p); + am.set_only_net(vec![OnlyNet::I2p]); + let got = am.take_dial_candidates_net(8, &HashSet::new(), &[]); + assert_eq!(got, vec![i2p]); + assert!(am.take_dial_candidates(8, &HashSet::new(), &[]).is_empty()); + } } diff --git a/crates/rbitcoin-node/src/cli.rs b/crates/rbitcoin-node/src/cli.rs index 2b8d7a9a7..f1b386846 100644 --- a/crates/rbitcoin-node/src/cli.rs +++ b/crates/rbitcoin-node/src/cli.rs @@ -295,6 +295,7 @@ fn operator_usage() -> String { [--signet-challenge HEX] [--signet-block-time SECS] \\\n\ [--listen ADDR] [--no-listen] [--connect ADDR]... [--seed-node HOST]... [--proxy HOST:PORT] [--onion HOST:PORT] [--proxy-randomize[=0|1]] [--only-net NET]... \\\n\ [--tor-control [HOST:PORT]] [--tor-control-cookie PATH] [--tor-control-password PASS] \\\n\ + [--i2p-sam [HOST:PORT]] \\\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\ @@ -325,6 +326,7 @@ Peers: --max-outbound (default 16 live download), --max-inbound (default 125).\n --tor-control [HOST:PORT] talks to system tor (default 127.0.0.1:9051). Cookie or password AUTH;\n\ failed AUTH is a start error. Unset: no control connection.\n\ --tor-control-cookie PATH (default /run/tor/control.authcookie). --tor-control-password PASS.\n\ + --i2p-sam [HOST:PORT] SAM v3 to system i2pd (default 127.0.0.1:7656). --only-net=i2p requires it.\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\ @@ -393,7 +395,7 @@ fn is_bool_key(key: &str) -> bool { fn is_optional_addr_key(key: &str) -> bool { matches!( key, - "rpc_listen" | "electrum_listen" | "esplora_listen" | "tor_control" + "rpc_listen" | "electrum_listen" | "esplora_listen" | "tor_control" | "i2p_sam" ) } @@ -581,6 +583,7 @@ mod tests { "--tor-control", "--tor-control-cookie", "--tor-control-password", + "--i2p-sam", ] { assert!(h.contains(flag), "help must list {flag}"); } @@ -613,6 +616,7 @@ mod tests { "--nolisten", "--nodiscover", "--torcontrol", + "--i2psam", ] { assert!(!h.contains(concat), "help must not advertise {concat}"); } @@ -847,7 +851,10 @@ mod tests { assert!(msg.contains("SOCKS") && msg.contains("only-net"), "{msg}"); c.apply_kv("proxy", "127.0.0.1:9050").unwrap(); c.validate().unwrap(); - assert!(NodeConfig::default().apply_kv("only_net", "i2p").is_err()); + let mut i2p_ok = NodeConfig::default(); + i2p_ok.apply_kv("only_net", "i2p").unwrap(); + assert_eq!(i2p_ok.listen.only_net, vec![rbitcoin_net::OnlyNet::I2p]); + assert!(NodeConfig::default().apply_kv("only_net", "cjdns").is_err()); let ok = ready_config([ "rbitcoin-node", "--only-net", @@ -894,6 +901,37 @@ mod tests { assert!(!h.contains("--torcontrol")); } + #[test] + fn i2p_sam_cli() { + let omitted = ready_config(["rbitcoin-node", "--i2p-sam"]); + assert_eq!( + omitted.listen.i2p_sam, + Some("127.0.0.1:7656".parse().unwrap()) + ); + let explicit = ready_config(["rbitcoin-node", "--i2p-sam", "127.0.0.1:7656"]); + assert_eq!( + explicit.listen.i2p_sam, + Some("127.0.0.1:7656".parse().unwrap()) + ); + let h = operator_usage(); + assert!(h.contains("--i2p-sam")); + assert!(!h.contains("--i2psam")); + } + + #[test] + fn only_net_i2p_without_sam_is_config_error() { + let mut c = NodeConfig::default(); + c.apply_kv("only_net", "i2p").unwrap(); + let err = c.validate().unwrap_err(); + let msg = format!("{err}"); + assert!(msg.contains("SAM") && msg.contains("i2p"), "{msg}"); + c.apply_kv("i2p_sam", "").unwrap(); + c.validate().unwrap(); + let ok = ready_config(["rbitcoin-node", "--only-net", "i2p", "--i2p-sam"]); + assert_eq!(ok.listen.only_net, vec![rbitcoin_net::OnlyNet::I2p]); + assert_eq!(ok.listen.i2p_sam, Some("127.0.0.1:7656".parse().unwrap())); + } + #[test] fn no_discover_conf() { let _g = OPERATOR_ENV_TEST_LOCK.lock().unwrap(); diff --git a/crates/rbitcoin-node/src/config.rs b/crates/rbitcoin-node/src/config.rs index 8327bfd08..dbbd6f873 100644 --- a/crates/rbitcoin-node/src/config.rs +++ b/crates/rbitcoin-node/src/config.rs @@ -103,6 +103,8 @@ pub struct ListenOpts { pub discover: bool, /// Empty = all networks. Repeatable `--only-net`. pub only_net: Vec, + /// SAM v3 TCP port (`--i2p-sam`). + pub i2p_sam: Option, } impl Default for ListenOpts { @@ -125,6 +127,7 @@ impl Default for ListenOpts { proxy_randomize: true, discover: true, only_net: Vec::new(), + i2p_sam: None, } } } @@ -519,6 +522,13 @@ impl NodeConfig { "only-net=onion requires SOCKS (--proxy or --onion)".into(), )); } + if self.listen.only_net.contains(&rbitcoin_net::OnlyNet::I2p) + && self.listen.i2p_sam.is_none() + { + return Err(NodeError::Config( + "only-net=i2p requires SAM (--i2p-sam)".into(), + )); + } Ok(()) } @@ -750,6 +760,13 @@ impl NodeConfig { "tor_control_password" => { self.tor.password = Some(val.to_string()); } + "i2p_sam" => { + self.listen.i2p_sam = Some(if val.is_empty() { + SocketAddr::from(([127, 0, 0, 1], 7656)) + } else { + parse_required_socket(val, "i2p_sam")? + }); + } "proxy_randomize" => { self.listen.proxy_randomize = parse_conf_bool(val) .map_err(|e| NodeError::Config(format!("conf proxy_randomize: {e}")))?; diff --git a/crates/rbitcoin-node/src/run.rs b/crates/rbitcoin-node/src/run.rs index 16b8fa513..31909c122 100644 --- a/crates/rbitcoin-node/src/run.rs +++ b/crates/rbitcoin-node/src/run.rs @@ -370,6 +370,15 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { .expect("control addr set when session exists") ); } + let _i2p_sam = if let Some(addr) = config.listen.i2p_sam { + let s = rbitcoin_net::I2pSam::connect(addr) + .await + .map_err(|e| NodeError::Init(format!("i2p sam {addr}: {e}")))?; + info!("i2p SAM session on {addr}"); + Some(s) + } else { + None + }; // One Class B appender thread. Join it at shutdown so apply does not race flush. let sh_writebehind = if config.shindex { Some(spawn_sh_writebehind( From ce1e9e64cf5064ac8a1a42a0e09487246943bb67 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 12:26:45 -0700 Subject: [PATCH 5/9] node: --i2p-accept-incoming STREAM FORWARDs to the P2P bind Persist {datadir}/i2p/p2p.priv from SESSION CREATE DESTINATION. --listen=0 is a Config error: overlay incoming still needs a loopback P2P listener. Co-authored-by: Cursor --- crates/rbitcoin-net/src/i2p_sam.rs | 198 +++++++++++++++++++++++++++-- crates/rbitcoin-node/src/cli.rs | 107 +++++++++++++++- crates/rbitcoin-node/src/config.rs | 20 +++ crates/rbitcoin-node/src/run.rs | 23 +++- 4 files changed, 330 insertions(+), 18 deletions(-) diff --git a/crates/rbitcoin-net/src/i2p_sam.rs b/crates/rbitcoin-net/src/i2p_sam.rs index 22ad040cb..aa4a292ee 100644 --- a/crates/rbitcoin-net/src/i2p_sam.rs +++ b/crates/rbitcoin-net/src/i2p_sam.rs @@ -1,7 +1,9 @@ -//! I2P SAM v3 STREAM CONNECT (system router, not SOCKS). +//! I2P SAM v3 STREAM CONNECT / FORWARD (system router, not SOCKS). use crate::error::NetError; +use std::io::Write; use std::net::SocketAddr; +use std::path::Path; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::net::TcpStream; @@ -9,29 +11,80 @@ pub struct I2pSam { sam_addr: SocketAddr, session_id: String, _control: TcpStream, + _forward: Option, } impl I2pSam { pub async fn connect(sam_addr: SocketAddr) -> Result { + Self::connect_session(sam_addr, None).await + } + + pub async fn connect_persistent( + sam_addr: SocketAddr, + dest_path: &Path, + ) -> Result { + let stored = match std::fs::read_to_string(dest_path) { + Ok(s) => { + let s = s.trim(); + if s.is_empty() { + None + } else { + Some(s.to_string()) + } + } + Err(e) if e.kind() == std::io::ErrorKind::NotFound => None, + Err(e) => { + return Err(NetError::Encode(format!( + "i2p dest {}: {e}", + dest_path.display() + ))); + } + }; + let (sam, dest) = Self::connect_session_dest(sam_addr, stored.as_deref()).await?; + if stored.is_none() { + write_dest_file(dest_path, &dest)?; + } + Ok(sam) + } + + async fn connect_session(sam_addr: SocketAddr, dest: Option<&str>) -> Result { + let (sam, _) = Self::connect_session_dest(sam_addr, dest).await?; + Ok(sam) + } + + async fn connect_session_dest( + sam_addr: SocketAddr, + dest: Option<&str>, + ) -> Result<(Self, String), NetError> { let mut control = TcpStream::connect(sam_addr) .await .map_err(|e| NetError::Encode(format!("i2p sam connect {sam_addr}: {e}")))?; hello(&mut control).await?; let session_id = fresh_session_id(); + let dest_arg = dest.unwrap_or("TRANSIENT"); write_line( &mut control, - &format!("SESSION CREATE STYLE=STREAM ID={session_id} DESTINATION=TRANSIENT"), + &format!("SESSION CREATE STYLE=STREAM ID={session_id} DESTINATION={dest_arg}"), ) .await?; let reply = read_line(&mut control).await?; if !reply.to_ascii_uppercase().contains("RESULT=OK") { return Err(NetError::Encode(format!("i2p sam session: {reply}"))); } - Ok(Self { - sam_addr, - session_id, - _control: control, - }) + let destination = sam_kv(&reply, "DESTINATION") + .ok_or_else(|| { + NetError::Encode(format!("i2p sam session missing DESTINATION: {reply}")) + })? + .to_string(); + Ok(( + Self { + sam_addr, + session_id, + _control: control, + _forward: None, + }, + destination, + )) } pub async fn stream_connect(&self, dest_b32: &str) -> Result { @@ -53,6 +106,55 @@ impl I2pSam { } Ok(s) } + + pub async fn stream_forward(&mut self, port: u16) -> Result<(), NetError> { + let mut s = TcpStream::connect(self.sam_addr) + .await + .map_err(|e| NetError::Encode(format!("i2p sam forward connect: {e}")))?; + hello(&mut s).await?; + write_line( + &mut s, + &format!("STREAM FORWARD ID={} PORT={port}", self.session_id), + ) + .await?; + let reply = read_line(&mut s).await?; + if !reply.to_ascii_uppercase().contains("RESULT=OK") { + return Err(NetError::Encode(format!("i2p sam forward: {reply}"))); + } + self._forward = Some(s); + Ok(()) + } +} + +fn sam_kv<'a>(line: &'a str, key: &str) -> Option<&'a str> { + let prefix = format!("{key}="); + line.split_whitespace().find_map(|tok| { + if tok.len() >= prefix.len() && tok[..prefix.len()].eq_ignore_ascii_case(&prefix) { + Some(&tok[prefix.len()..]) + } else { + None + } + }) +} + +fn write_dest_file(path: &Path, dest: &str) -> Result<(), NetError> { + if let Some(parent) = path.parent() { + std::fs::create_dir_all(parent) + .map_err(|e| NetError::Encode(format!("i2p dest dir {}: {e}", parent.display())))?; + } + let mut opts = std::fs::OpenOptions::new(); + opts.write(true).create(true).truncate(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + opts.mode(0o600); + } + let mut f = opts + .open(path) + .map_err(|e| NetError::Encode(format!("i2p dest {}: {e}", path.display())))?; + writeln!(f, "{dest}") + .map_err(|e| NetError::Encode(format!("i2p dest {}: {e}", path.display())))?; + Ok(()) } fn fresh_session_id() -> String { @@ -104,8 +206,11 @@ async fn read_line(s: &mut TcpStream) -> Result { mod tests { use super::*; use std::sync::{Arc, Mutex}; + use std::time::{SystemTime, UNIX_EPOCH}; use tokio::net::TcpListener; + const FAKE_DEST: &str = "fakeprivdest"; + async fn fake_sam(ok_hello: bool, dest_log: Arc>>) -> SocketAddr { let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let addr = listener.local_addr().unwrap(); @@ -132,10 +237,23 @@ mod tests { break; } } else if up.starts_with("SESSION CREATE") { - let _ = write_line(&mut s, "SESSION STATUS RESULT=OK DESTINATION=fake") - .await; - } else if let Some(rest) = line.strip_prefix("STREAM CONNECT ") { - log.lock().unwrap().push(rest.to_string()); + log.lock().unwrap().push(line.clone()); + let dest = sam_kv(&line, "DESTINATION").unwrap_or("TRANSIENT"); + let reply_dest = if dest.eq_ignore_ascii_case("TRANSIENT") { + FAKE_DEST + } else { + dest + }; + let _ = write_line( + &mut s, + &format!("SESSION STATUS RESULT=OK DESTINATION={reply_dest}"), + ) + .await; + } else if up.starts_with("STREAM FORWARD") { + log.lock().unwrap().push(line.clone()); + let _ = write_line(&mut s, "STREAM STATUS RESULT=OK").await; + } else if up.starts_with("STREAM CONNECT") { + log.lock().unwrap().push(line.clone()); let _ = write_line(&mut s, "STREAM STATUS RESULT=OK").await; break; } else { @@ -148,6 +266,15 @@ mod tests { addr } + fn stream_lines(log: &Arc>>, prefix: &str) -> Vec { + log.lock() + .unwrap() + .iter() + .filter(|g| g.to_ascii_uppercase().starts_with(prefix)) + .cloned() + .collect() + } + #[tokio::test] async fn i2p_sam_stream_connect_fake() { let log = Arc::new(Mutex::new(Vec::new())); @@ -156,7 +283,7 @@ mod tests { let dest = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa.b32.i2p"; sam.stream_connect(dest).await.unwrap(); sam.stream_connect(dest).await.unwrap(); - let got = log.lock().unwrap().clone(); + let got = stream_lines(&log, "STREAM CONNECT"); assert_eq!(got.len(), 2, "{got:?}"); for g in &got { assert!(g.contains(dest), "{g}"); @@ -185,7 +312,7 @@ mod tests { port: 8333, }; sam.stream_connect(&peer.host_str()).await.unwrap(); - let got = log.lock().unwrap().clone(); + let got = stream_lines(&log, "STREAM CONNECT"); assert!(got.iter().any(|g| g.contains(&peer.host_str())), "{got:?}"); let err = match crate::socks::Dialer::Direct.connect_net(peer).await { Err(e) => e, @@ -194,4 +321,49 @@ mod tests { let msg = format!("{err}"); assert!(msg.contains("SAM") || msg.contains("i2p"), "{msg}"); } + + #[tokio::test] + async fn i2p_accept_incoming_forwards_to_loopback() { + let log = Arc::new(Mutex::new(Vec::new())); + let addr = fake_sam(true, Arc::clone(&log)).await; + let dir = std::env::temp_dir().join(format!( + "rbtc-i2p-{}-{}", + std::process::id(), + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_nanos() + )); + let dest_path = dir.join("i2p").join("p2p.priv"); + let mut sam = I2pSam::connect_persistent(addr, &dest_path).await.unwrap(); + sam.stream_forward(18444).await.unwrap(); + let stored = std::fs::read_to_string(&dest_path).unwrap(); + assert_eq!(stored.trim(), FAKE_DEST); + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + let mode = std::fs::metadata(&dest_path).unwrap().permissions().mode() & 0o777; + assert_eq!(mode, 0o600); + } + let fw = stream_lines(&log, "STREAM FORWARD"); + assert_eq!(fw.len(), 1, "{fw:?}"); + assert!(fw[0].contains("PORT=18444"), "{}", fw[0]); + assert!(fw[0].contains("ID=rbtc"), "{}", fw[0]); + let creates = stream_lines(&log, "SESSION CREATE"); + assert!( + creates.iter().any(|c| c.contains("DESTINATION=TRANSIENT")), + "{creates:?}" + ); + + let mut sam2 = I2pSam::connect_persistent(addr, &dest_path).await.unwrap(); + sam2.stream_forward(18444).await.unwrap(); + let creates = stream_lines(&log, "SESSION CREATE"); + assert!( + creates + .iter() + .any(|c| c.contains(&format!("DESTINATION={FAKE_DEST}"))), + "{creates:?}" + ); + let _ = std::fs::remove_dir_all(&dir); + } } diff --git a/crates/rbitcoin-node/src/cli.rs b/crates/rbitcoin-node/src/cli.rs index f1b386846..3d5f414b1 100644 --- a/crates/rbitcoin-node/src/cli.rs +++ b/crates/rbitcoin-node/src/cli.rs @@ -295,7 +295,7 @@ fn operator_usage() -> String { [--signet-challenge HEX] [--signet-block-time SECS] \\\n\ [--listen ADDR] [--no-listen] [--connect ADDR]... [--seed-node HOST]... [--proxy HOST:PORT] [--onion HOST:PORT] [--proxy-randomize[=0|1]] [--only-net NET]... \\\n\ [--tor-control [HOST:PORT]] [--tor-control-cookie PATH] [--tor-control-password PASS] \\\n\ - [--i2p-sam [HOST:PORT]] \\\n\ + [--i2p-sam [HOST:PORT]] [--i2p-accept-incoming] \\\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\ @@ -327,6 +327,7 @@ Peers: --max-outbound (default 16 live download), --max-inbound (default 125).\n failed AUTH is a start error. Unset: no control connection.\n\ --tor-control-cookie PATH (default /run/tor/control.authcookie). --tor-control-password PASS.\n\ --i2p-sam [HOST:PORT] SAM v3 to system i2pd (default 127.0.0.1:7656). --only-net=i2p requires it.\n\ + --i2p-accept-incoming persist {{datadir}}/i2p/p2p.priv and STREAM FORWARD to the P2P bind. Needs --listen.\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\ @@ -384,6 +385,7 @@ fn is_bool_key(key: &str) -> bool { | "no_listen" | "no_discover" | "proxy_randomize" + | "i2p_accept_incoming" | "inhibit_suspend" | "trusted" | "always_relay" @@ -584,6 +586,7 @@ mod tests { "--tor-control-cookie", "--tor-control-password", "--i2p-sam", + "--i2p-accept-incoming", ] { assert!(h.contains(flag), "help must list {flag}"); } @@ -916,6 +919,108 @@ mod tests { let h = operator_usage(); assert!(h.contains("--i2p-sam")); assert!(!h.contains("--i2psam")); + assert!(h.contains("--i2p-accept-incoming")); + assert!(!h.contains("--i2pacceptincoming")); + } + + #[test] + fn i2p_accept_incoming_cli() { + let c = ready_config(["rbitcoin-node", "--i2p-sam", "--i2p-accept-incoming"]); + assert!(c.listen.i2p_accept_incoming); + assert_eq!(c.listen.i2p_sam, Some("127.0.0.1:7656".parse().unwrap())); + let mut conf = NodeConfig::default(); + conf.apply_kv("i2p_sam", "").unwrap(); + conf.apply_kv("i2p_accept_incoming", "1").unwrap(); + conf.validate().unwrap(); + assert!(conf.listen.i2p_accept_incoming); + } + + #[test] + fn i2p_accept_incoming_without_listen_is_config_error() { + let c = ready_config([ + "rbitcoin-node", + "--no-listen", + "--i2p-sam", + "--i2p-accept-incoming", + ]); + let err = c.validate().unwrap_err(); + let msg = format!("{err}"); + assert!(msg.contains("--listen") && msg.contains("i2p"), "{msg}"); + let mut no_sam = NodeConfig::default(); + no_sam.apply_kv("i2p_accept_incoming", "1").unwrap(); + let err = no_sam.validate().unwrap_err(); + let msg = format!("{err}"); + assert!(msg.contains("SAM") && msg.contains("i2p"), "{msg}"); + } + + #[tokio::test] + async fn i2p_accept_incoming_forwards_to_loopback() { + use std::sync::{Arc, Mutex}; + use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; + use tokio::net::{TcpListener, TcpStream}; + + async fn write_line(s: &mut TcpStream, line: &str) { + s.write_all(line.as_bytes()).await.unwrap(); + s.write_all(b"\n").await.unwrap(); + s.flush().await.unwrap(); + } + async fn read_line(s: &mut TcpStream) -> Option { + let mut reader = BufReader::new(s); + let mut line = String::new(); + let n = reader.read_line(&mut line).await.ok()?; + if n == 0 { + return None; + } + Some(line.trim_end_matches(['\r', '\n']).to_string()) + } + + let log = Arc::new(Mutex::new(Vec::new())); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let log_acc = Arc::clone(&log); + tokio::spawn(async move { + loop { + let Ok((mut s, _)) = listener.accept().await else { + break; + }; + let log = Arc::clone(&log_acc); + tokio::spawn(async move { + loop { + let Some(line) = read_line(&mut s).await else { + break; + }; + let up = line.to_ascii_uppercase(); + if up.starts_with("HELLO VERSION") { + write_line(&mut s, "HELLO REPLY RESULT=OK VERSION=3.1").await; + } else if up.starts_with("SESSION CREATE") { + write_line(&mut s, "SESSION STATUS RESULT=OK DESTINATION=fakeprivdest") + .await; + } else if up.starts_with("STREAM FORWARD") { + log.lock().unwrap().push(line); + write_line(&mut s, "STREAM STATUS RESULT=OK").await; + } else if up.starts_with("STREAM CONNECT") { + write_line(&mut s, "STREAM STATUS RESULT=OK").await; + break; + } + } + }); + } + }); + + let dir = tmp_datadir(); + let dest_path = dir.join("i2p").join("p2p.priv"); + let mut sam = rbitcoin_net::I2pSam::connect_persistent(addr, &dest_path) + .await + .unwrap(); + sam.stream_forward(18444).await.unwrap(); + assert_eq!( + std::fs::read_to_string(&dest_path).unwrap().trim(), + "fakeprivdest" + ); + let fw = log.lock().unwrap().clone(); + assert_eq!(fw.len(), 1, "{fw:?}"); + assert!(fw[0].contains("PORT=18444"), "{}", fw[0]); + let _ = std::fs::remove_dir_all(&dir); } #[test] diff --git a/crates/rbitcoin-node/src/config.rs b/crates/rbitcoin-node/src/config.rs index dbbd6f873..e139ac0e1 100644 --- a/crates/rbitcoin-node/src/config.rs +++ b/crates/rbitcoin-node/src/config.rs @@ -105,6 +105,8 @@ pub struct ListenOpts { pub only_net: Vec, /// SAM v3 TCP port (`--i2p-sam`). pub i2p_sam: Option, + /// Persistent SAM destination + STREAM FORWARD to the P2P bind (`--i2p-accept-incoming`). + pub i2p_accept_incoming: bool, } impl Default for ListenOpts { @@ -128,6 +130,7 @@ impl Default for ListenOpts { discover: true, only_net: Vec::new(), i2p_sam: None, + i2p_accept_incoming: false, } } } @@ -529,6 +532,19 @@ impl NodeConfig { "only-net=i2p requires SAM (--i2p-sam)".into(), )); } + if self.listen.i2p_accept_incoming { + if self.listen.i2p_sam.is_none() { + return Err(NodeError::Config( + "i2p-accept-incoming requires SAM (--i2p-sam)".into(), + )); + } + if matches!(self.listen.p2p, P2pListen::Off) { + return Err(NodeError::Config( + "i2p-accept-incoming needs a P2P listener (--listen); --listen=0 has no loopback to STREAM FORWARD" + .into(), + )); + } + } Ok(()) } @@ -767,6 +783,10 @@ impl NodeConfig { parse_required_socket(val, "i2p_sam")? }); } + "i2p_accept_incoming" => { + self.listen.i2p_accept_incoming = parse_conf_bool(val) + .map_err(|e| NodeError::Config(format!("conf i2p_accept_incoming: {e}")))?; + } "proxy_randomize" => { self.listen.proxy_randomize = parse_conf_bool(val) .map_err(|e| NodeError::Config(format!("conf proxy_randomize: {e}")))?; diff --git a/crates/rbitcoin-node/src/run.rs b/crates/rbitcoin-node/src/run.rs index 31909c122..5fb4f0853 100644 --- a/crates/rbitcoin-node/src/run.rs +++ b/crates/rbitcoin-node/src/run.rs @@ -370,15 +370,30 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { .expect("control addr set when session exists") ); } - let _i2p_sam = if let Some(addr) = config.listen.i2p_sam { - let s = rbitcoin_net::I2pSam::connect(addr) - .await - .map_err(|e| NodeError::Init(format!("i2p sam {addr}: {e}")))?; + let mut i2p_sam = if let Some(addr) = config.listen.i2p_sam { + let s = if config.listen.i2p_accept_incoming { + let dest = config.datadir.path().join("i2p").join("p2p.priv"); + rbitcoin_net::I2pSam::connect_persistent(addr, &dest).await + } else { + rbitcoin_net::I2pSam::connect(addr).await + } + .map_err(|e| NodeError::Init(format!("i2p sam {addr}: {e}")))?; info!("i2p SAM session on {addr}"); Some(s) } else { None }; + if config.listen.i2p_accept_incoming { + let port = node.local_addr.port(); + let sam = i2p_sam + .as_mut() + .expect("validate requires --i2p-sam with --i2p-accept-incoming"); + sam.stream_forward(port) + .await + .map_err(|e| NodeError::Init(format!("i2p STREAM FORWARD {port}: {e}")))?; + info!("i2p STREAM FORWARD to {}", node.local_addr); + } + let _i2p_sam = i2p_sam; // One Class B appender thread. Join it at shutdown so apply does not race flush. let sh_writebehind = if config.shindex { Some(spawn_sh_writebehind( From a0e14c49e4baf122d7e65e768fc5667c491e7070 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 12:34:35 -0700 Subject: [PATCH 6/9] docs+nix: --i2p-sam, --i2p-accept-incoming, After=i2pd NixOS i2p.sam / i2p.acceptIncoming are first-class. Eval pins argv. Runtime dummy i2pd.service only for After/Wants; CI does not start i2pd. Co-authored-by: Cursor --- OPERATOR.md | 18 ++++++++++--- nix/modules/rbitcoin.nix | 42 +++++++++++++++++++++++++++--- nix/tests/nixos-module-eval.nix | 16 +++++++++++- nix/tests/nixos-module-runtime.nix | 14 ++++++++++ 4 files changed, 83 insertions(+), 7 deletions(-) diff --git a/OPERATOR.md b/OPERATOR.md index 3f2f353e7..cd6a721b3 100644 --- a/OPERATOR.md +++ b/OPERATOR.md @@ -368,14 +368,16 @@ Clean smoke: | `--listen ADDR` | `listen=` | bind later default port | | `--no-listen` / `--listen=0` | `listen=0` / `no_listen=` | bind a loopback default; **off** = no P2P socket (outbound-only) | | `--no-discover` | `no_discover=` | discover **on**; flag off = no self-announce / `localaddresses` | -| `--only-net NET` | `only_net=` | all nets; repeatable `ipv4` / `ipv6` / `onion` (`i2p`/`cjdns` later) | -| `--connect ADDR` | `connect=` (repeatable) | seeds; `IP:port` or Tor v3 `.onion:port` | +| `--only-net NET` | `only_net=` | all nets; repeatable `ipv4` / `ipv6` / `onion` / `i2p` (`cjdns` later) | +| `--connect ADDR` | `connect=` (repeatable) | seeds; `IP:port`, Tor v3 `.onion:port`, or `{52}.b32.i2p:port` | | `--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) | | `--tor-control [HOST:PORT]` | `tor_control=` | unset — no control connection; omit ADDR → `127.0.0.1:9051` | | `--tor-control-cookie PATH` | `tor_control_cookie=` | `/run/tor/control.authcookie` when `--tor-control` is set and password is unset | | `--tor-control-password PASS` | `tor_control_password=` | unset — cookie AUTH unless set | +| `--i2p-sam [HOST:PORT]` | `i2p_sam=` | unset — no SAM; omit ADDR → `127.0.0.1:7656` | +| `--i2p-accept-incoming` | `i2p_accept_incoming=` | **off** — persist `{datadir}/i2p/p2p.priv` and STREAM FORWARD to the P2P bind | | `--milestone HEIGHT` | `milestone=` | network default (mainnet 840000) | | `--max-outbound N` | `max_outbound=` | 16 live download peers | | `--max-inbound N` | `max_inbound=` | 125 inbound sessions; **0** = no inbound slots (outbound-only) | @@ -441,7 +443,7 @@ mempool_size_mb=100 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. -`--only-net onion` (repeatable with `ipv4`/`ipv6`) filters dial and learn; +`--only-net onion` (repeatable with `ipv4`/`ipv6`/`i2p`) filters dial and learn; onion requires `--proxy` or `--onion`. `--connect foo.onion:8333` is a start error when the v3 checksum is invalid. The peers file is `rbitcoin-peers-v2` (v1 IPv4/IPv6 still loads). @@ -460,6 +462,16 @@ socket. With `--electrum-listen`, the node `ADD_ONION`s that TCP port to the onion (`rpc.sock` / `--rpc-listen` only). Cookie path differs by distro; pass `--tor-control-cookie` rather than globbing. +`--i2p-sam [HOST:PORT]` talks to **system i2pd** SAM v3 (not SOCKS, not Arti). +Omit ADDR for `127.0.0.1:7656`. Failed HELLO / `SESSION CREATE` is a start +error. Unset: I2P rows may still load from `peers` v2 but are not dialed. +`--only-net i2p` without `--i2p-sam` is a start error. `--i2p-accept-incoming` +creates a persistent local destination (`{datadir}/i2p/p2p.priv`, 0600) and +`STREAM FORWARD`s to the P2P bind. With `--listen=0` that is a start error +naming `--listen` (overlay incoming still needs a loopback P2P accept). NixOS: +`services.rbitcoin.i2p.sam` / `i2p.acceptIncoming`; the unit `After`/`Wants` +`i2pd.service` when SAM is set. Do not start i2pd from this module. + `--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 214f0f539..ac40cf019 100644 --- a/nix/modules/rbitcoin.nix +++ b/nix/modules/rbitcoin.nix @@ -40,6 +40,7 @@ let "${address}:${toString port}"; needTorControl = cfg.tor.control != null || cfg.electrum.hiddenService; + needI2pSam = cfg.i2p.sam != null; torControlAddr = if cfg.tor.control != null then cfg.tor.control @@ -91,6 +92,9 @@ let ++ optional (torControlAddr != null) torControlAddr ++ optional (cfg.tor.controlCookie != null) "--tor-control-cookie" ++ optional (cfg.tor.controlCookie != null) (toString cfg.tor.controlCookie) + ++ optional needI2pSam "--i2p-sam" + ++ optional needI2pSam cfg.i2p.sam + ++ optional cfg.i2p.acceptIncoming "--i2p-accept-incoming" ++ cfg.extraArgs; in { @@ -196,6 +200,21 @@ in }; }; + i2p = { + sam = mkOption { + type = types.nullOr types.str; + default = null; + example = "127.0.0.1:7656"; + description = "I2P SAM v3 HOST:PORT (system i2pd). Unset skips SAM."; + }; + + acceptIncoming = mkOption { + type = types.bool; + default = false; + description = "STREAM FORWARD to the P2P bind. Requires i2p.sam and p2p.listen. Persists {dataDir}/i2p/p2p.priv."; + }; + }; + proxy = mkOption { type = types.nullOr types.str; default = null; @@ -221,9 +240,10 @@ in "ipv4" "ipv6" "onion" + "i2p" ]); default = [ ]; - description = "Restrict P2P to these networks. onion requires proxy or onionProxy."; + description = "Restrict P2P to these networks. onion requires proxy or onionProxy; i2p requires i2p.sam."; }; p2p = { @@ -337,6 +357,14 @@ in assertion = cfg.coldDataDir == null || cfg.coldDataDir != cfg.dataDir; message = "services.rbitcoin.coldDataDir must differ from dataDir"; } + { + assertion = !cfg.i2p.acceptIncoming || cfg.i2p.sam != null; + message = "services.rbitcoin.i2p.acceptIncoming requires i2p.sam"; + } + { + assertion = !cfg.i2p.acceptIncoming || cfg.p2p.listen; + message = "services.rbitcoin.i2p.acceptIncoming requires p2p.listen (STREAM FORWARD needs a P2P bind)"; + } ]; users.groups.${cfg.group} = { }; @@ -365,8 +393,16 @@ in description = "rbitcoin full node"; documentation = [ "https://github.com/reardencode/rbitcoin/blob/master/OPERATOR.md" ]; wantedBy = [ "multi-user.target" ]; - wants = [ "network-online.target" ] ++ optional needTorControl "tor.service"; - after = [ "network-online.target" ] ++ optional needTorControl "tor.service"; + wants = [ + "network-online.target" + ] + ++ optional needTorControl "tor.service" + ++ optional needI2pSam "i2pd.service"; + after = [ + "network-online.target" + ] + ++ optional needTorControl "tor.service" + ++ optional needI2pSam "i2pd.service"; environment = cfg.environment; serviceConfig = { diff --git a/nix/tests/nixos-module-eval.nix b/nix/tests/nixos-module-eval.nix index 427e5d6af..52d366c04 100644 --- a/nix/tests/nixos-module-eval.nix +++ b/nix/tests/nixos-module-eval.nix @@ -31,11 +31,18 @@ let proxy = "127.0.0.1:9050"; onionProxy = "127.0.0.1:9050"; proxyRandomize = true; - onlyNet = [ "onion" ]; + onlyNet = [ + "onion" + "i2p" + ]; tor = { control = "127.0.0.1:9051"; controlCookie = "/run/tor/control.authcookie"; }; + i2p = { + sam = "127.0.0.1:7656"; + acceptIncoming = true; + }; p2p = { address = "127.0.0.1"; openFirewall = true; @@ -99,6 +106,8 @@ assert defaultCfg.onlyNet == [ ]; assert defaultCfg.tor.control == null; assert defaultCfg.tor.controlCookie == null; assert defaultCfg.electrum.hiddenService == false; +assert defaultCfg.i2p.sam == null; +assert defaultCfg.i2p.acceptIncoming == false; assert cfg.services.rbitcoin.p2p.port == 18444; assert cfg.services.rbitcoin.rpc.port == 18443; assert @@ -125,10 +134,15 @@ 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; assert builtins.match ".*--only-net onion.*" execStart != null; +assert builtins.match ".*--only-net i2p.*" execStart != null; assert builtins.match ".*--tor-control 127.0.0.1:9051.*" execStart != null; assert builtins.match ".*--tor-control-cookie /run/tor/control.authcookie.*" execStart != null; +assert builtins.match ".*--i2p-sam 127.0.0.1:7656.*" execStart != null; +assert builtins.match ".*--i2p-accept-incoming.*" execStart != null; assert builtins.elem "tor.service" service.after; assert builtins.elem "tor.service" service.wants; +assert builtins.elem "i2pd.service" service.after; +assert builtins.elem "i2pd.service" service.wants; assert builtins.match ".*--no-listen.*" listenOffExec != null; assert builtins.match ".*--listen .*" listenOffExec == null; assert builtins.match ".*--max-inbound 0.*" listenOffExec != null; diff --git a/nix/tests/nixos-module-runtime.nix b/nix/tests/nixos-module-runtime.nix index 614eeea32..84455e298 100644 --- a/nix/tests/nixos-module-runtime.nix +++ b/nix/tests/nixos-module-runtime.nix @@ -32,6 +32,7 @@ pkgs.testers.runNixOSTest { }; rpc.enable = true; tor.control = "127.0.0.1:9051"; + i2p.sam = "127.0.0.1:7656"; extraArgs = [ "--max-outbound" "4" @@ -47,6 +48,16 @@ pkgs.testers.runNixOSTest { ExecStart = "${pkgs.coreutils}/bin/true"; }; }; + + systemd.services.i2pd = { + description = "fake i2pd unit for After= ordering"; + wantedBy = [ "multi-user.target" ]; + serviceConfig = { + Type = "oneshot"; + RemainAfterExit = true; + ExecStart = "${pkgs.coreutils}/bin/true"; + }; + }; }; testScript = '' @@ -62,7 +73,10 @@ pkgs.testers.runNixOSTest { machine.succeed("grep -Fx -- '4' /var/lib/rbitcoin-test/args") machine.succeed("grep -Fx -- '--tor-control' /var/lib/rbitcoin-test/args") machine.succeed("grep -Fx -- '127.0.0.1:9051' /var/lib/rbitcoin-test/args") + machine.succeed("grep -Fx -- '--i2p-sam' /var/lib/rbitcoin-test/args") + machine.succeed("grep -Fx -- '127.0.0.1:7656' /var/lib/rbitcoin-test/args") machine.succeed("systemctl show -p After rbitcoin.service | grep -F tor.service") + machine.succeed("systemctl show -p After rbitcoin.service | grep -F i2pd.service") machine.succeed("systemctl stop rbitcoin.service") machine.succeed("test -e /var/lib/rbitcoin-test/stopped") ''; From 71bd62bf5e348de11466f8d9ce4bd9d682b08a4b Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 19 Sep 2026 22:02:39 -0700 Subject: [PATCH 7/9] net: route i2p outbound/follow through SAM dial targets Co-authored-by: Cursor --- crates/rbitcoin-net/src/i2p_sam.rs | 45 ++++++++++++++++++-------- crates/rbitcoin-net/src/ibd/dial.rs | 37 ++++++++++++--------- crates/rbitcoin-net/src/ibd/mod.rs | 17 ++++++---- crates/rbitcoin-net/src/ibd/peer_io.rs | 11 ++++--- crates/rbitcoin-net/src/lib.rs | 2 +- crates/rbitcoin-net/src/peers.rs | 9 ++++++ crates/rbitcoin-net/src/seeds.rs | 20 +++++++++--- crates/rbitcoin-net/src/service.rs | 19 ++++++++--- crates/rbitcoin-net/src/socks.rs | 12 ++++++- crates/rbitcoin-net/tests/ibd_smoke.rs | 11 +++++-- crates/rbitcoin-node/src/run.rs | 44 +++++++++++++++++-------- 11 files changed, 162 insertions(+), 65 deletions(-) diff --git a/crates/rbitcoin-net/src/i2p_sam.rs b/crates/rbitcoin-net/src/i2p_sam.rs index aa4a292ee..62bba2c04 100644 --- a/crates/rbitcoin-net/src/i2p_sam.rs +++ b/crates/rbitcoin-net/src/i2p_sam.rs @@ -14,6 +14,12 @@ pub struct I2pSam { _forward: Option, } +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct I2pDialer { + sam_addr: SocketAddr, + session_id: String, +} + impl I2pSam { pub async fn connect(sam_addr: SocketAddr) -> Result { Self::connect_session(sam_addr, None).await @@ -88,41 +94,54 @@ impl I2pSam { } pub async fn stream_connect(&self, dest_b32: &str) -> Result { + self.dialer().stream_connect(dest_b32).await + } + + pub fn dialer(&self) -> I2pDialer { + I2pDialer { + sam_addr: self.sam_addr, + session_id: self.session_id.clone(), + } + } + + pub async fn stream_forward(&mut self, port: u16) -> Result<(), NetError> { let mut s = TcpStream::connect(self.sam_addr) .await - .map_err(|e| NetError::Encode(format!("i2p sam stream connect: {e}")))?; + .map_err(|e| NetError::Encode(format!("i2p sam forward connect: {e}")))?; hello(&mut s).await?; write_line( &mut s, - &format!( - "STREAM CONNECT ID={} DESTINATION={}", - self.session_id, dest_b32 - ), + &format!("STREAM FORWARD ID={} PORT={port}", self.session_id), ) .await?; let reply = read_line(&mut s).await?; if !reply.to_ascii_uppercase().contains("RESULT=OK") { - return Err(NetError::Encode(format!("i2p sam stream: {reply}"))); + return Err(NetError::Encode(format!("i2p sam forward: {reply}"))); } - Ok(s) + self._forward = Some(s); + Ok(()) } +} - pub async fn stream_forward(&mut self, port: u16) -> Result<(), NetError> { +impl I2pDialer { + pub async fn stream_connect(&self, dest_b32: &str) -> Result { let mut s = TcpStream::connect(self.sam_addr) .await - .map_err(|e| NetError::Encode(format!("i2p sam forward connect: {e}")))?; + .map_err(|e| NetError::Encode(format!("i2p sam stream connect: {e}")))?; hello(&mut s).await?; write_line( &mut s, - &format!("STREAM FORWARD ID={} PORT={port}", self.session_id), + &format!( + "STREAM CONNECT ID={} DESTINATION={}", + self.session_id, dest_b32 + ), ) .await?; let reply = read_line(&mut s).await?; if !reply.to_ascii_uppercase().contains("RESULT=OK") { - return Err(NetError::Encode(format!("i2p sam forward: {reply}"))); + return Err(NetError::Encode(format!("i2p sam stream: {reply}"))); } - self._forward = Some(s); - Ok(()) + Ok(s) } } diff --git a/crates/rbitcoin-net/src/ibd/dial.rs b/crates/rbitcoin-net/src/ibd/dial.rs index 23f61f279..ffb3229a4 100644 --- a/crates/rbitcoin-net/src/ibd/dial.rs +++ b/crates/rbitcoin-net/src/ibd/dial.rs @@ -223,8 +223,8 @@ fn classify_dial_err(e: &NetError) -> DialFailKind { /// Result of a dial batch: live slots + failures for the peer book. pub(crate) struct DialBatchResult { pub slots: Vec, - pub failed: Vec<(SocketAddr, DialFailKind)>, - pub attempted: Vec, + pub failed: Vec<(crate::NetAddr, DialFailKind)>, + pub attempted: Vec, } #[allow(clippy::too_many_arguments)] // call-site args stay unbundled @@ -235,7 +235,7 @@ pub(crate) async fn dial_batch( book: &AddrMan, next_id: &AtomicUsize, count: usize, - mut already: HashSet, + mut already: HashSet, occupied: &[SocketAddr], magic: Magic, local_addr: SocketAddr, @@ -260,7 +260,7 @@ pub(crate) async fn dial_batch( .unwrap_or(false) }; - let candidates = book.take_dial_candidates(count, &already, occupied); + let candidates = book.take_dial_candidates_net(count, &already, occupied); out.attempted = candidates.clone(); let mut handles = Vec::new(); for addr in candidates { @@ -333,13 +333,13 @@ pub(crate) fn redial_want(alive: usize, target: usize) -> usize { /// Apply dial successes / failures to the peer book. pub(crate) fn apply_dial_result(book: &mut AddrMan, result: &DialBatchResult) { for &addr in &result.attempted { - book.note_attempt(addr); + book.note_attempt_addr(addr); } for s in &result.slots { book.note_connected(s.addr); } for &(addr, kind) in &result.failed { - book.note_connect_failed(addr, kind == DialFailKind::Incompatible); + book.note_connect_failed_addr(addr, kind == DialFailKind::Incompatible); } } @@ -441,11 +441,12 @@ pub(crate) fn dial_blocked_addrs( slots: &[PeerSlot], cooldown: &HashMap, now: Instant, -) -> HashSet { - let mut blocked: HashSet = slots.iter().map(|s| s.addr).collect(); +) -> HashSet { + let mut blocked: HashSet = + slots.iter().map(|s| crate::NetAddr::Ip(s.addr)).collect(); for (&addr, &until) in cooldown { if until > now { - blocked.insert(addr); + blocked.insert(crate::NetAddr::Ip(addr)); } } blocked @@ -700,9 +701,9 @@ mod tests { cooldown.insert(addr(2), now + Duration::from_secs(60)); cooldown.insert(addr(3), now - Duration::from_secs(1)); // expired let blocked = dial_blocked_addrs(&[s], &cooldown, now); - assert!(blocked.contains(&addr(1))); - assert!(blocked.contains(&addr(2))); - assert!(!blocked.contains(&addr(3))); + assert!(blocked.contains(&crate::NetAddr::Ip(addr(1)))); + assert!(blocked.contains(&crate::NetAddr::Ip(addr(2)))); + assert!(!blocked.contains(&crate::NetAddr::Ip(addr(3)))); expire_addr_cooldown(&mut cooldown, now); assert!(cooldown.contains_key(&addr(2))); @@ -724,7 +725,7 @@ mod tests { assert_eq!(book.flags(&lemon).dial_tier(), 2); assert!(cooldown.contains_key(&lemon)); let blocked = dial_blocked_addrs(&[], &cooldown, now); - assert!(blocked.contains(&lemon)); + assert!(blocked.contains(&crate::NetAddr::Ip(lemon))); let good = addr(5); book.note_connected(good); @@ -903,10 +904,14 @@ mod tests { let result = DialBatchResult { slots: vec![slot], failed: vec![ - (bad, DialFailKind::Network), - (inc, DialFailKind::Incompatible), + (crate::NetAddr::Ip(bad), DialFailKind::Network), + (crate::NetAddr::Ip(inc), DialFailKind::Incompatible), + ], + attempted: vec![ + crate::NetAddr::Ip(good), + crate::NetAddr::Ip(bad), + crate::NetAddr::Ip(inc), ], - attempted: vec![good, bad, inc], }; apply_dial_result(&mut book, &result); assert!(book.flags(&good).has_connected()); diff --git a/crates/rbitcoin-net/src/ibd/mod.rs b/crates/rbitcoin-net/src/ibd/mod.rs index 39015c5c0..41b723d8a 100644 --- a/crates/rbitcoin-net/src/ibd/mod.rs +++ b/crates/rbitcoin-net/src/ibd/mod.rs @@ -208,14 +208,16 @@ struct PeerBookSession { impl PeerBookSession { fn new( shared: Option>>, - seed_peers: &[SocketAddr], + seed_peers: &[crate::NetAddr], ) -> Self { let mut book = if let Some(ref s) = shared { s.lock().unwrap_or_else(|e| e.into_inner()).clone() } else { crate::seeds::AddrMan::new() }; - book.inject(seed_peers.iter().copied()); + for &addr in seed_peers { + book.add_addr(addr); + } Self { book, shared } } @@ -247,7 +249,7 @@ pub async fn ibd_cancellable( hub: Arc, magic: Magic, local_addr: SocketAddr, - peers: &[SocketAddr], + peers: &[crate::NetAddr], cfg: IbdConfig, cancel: Option>, ) -> Result { @@ -672,7 +674,7 @@ pub async fn ibd_cancellable( let mut n = 0usize; for s in result.slots { // Race: same addr may have connected on another path. - if blocked.contains(&s.addr) + if blocked.contains(&crate::NetAddr::Ip(s.addr)) || st.slots.iter().any(|x| x.addr == s.addr) { warn!( @@ -1104,7 +1106,10 @@ mod peer_book_and_config_tests { fn peer_book_session_injects_seeds_and_flushes_on_drop() { let shared = Arc::new(Mutex::new(AddrMan::new())); { - let mut sess = PeerBookSession::new(Some(Arc::clone(&shared)), &[sa(1), sa(2)]); + let mut sess = PeerBookSession::new( + Some(Arc::clone(&shared)), + &[crate::NetAddr::Ip(sa(1)), crate::NetAddr::Ip(sa(2))], + ); assert!(sess.book().entry(&sa(1)).is_some()); assert!(sess.book().entry(&sa(2)).is_some()); // Mutate book via book_mut. @@ -1119,7 +1124,7 @@ mod peer_book_and_config_tests { assert!(shared.lock().unwrap().entry(&sa(1)).is_some()); // No shared book — seeds only, flush is a no-op. - let sess2 = PeerBookSession::new(None, &[sa(9)]); + let sess2 = PeerBookSession::new(None, &[crate::NetAddr::Ip(sa(9))]); assert!(sess2.book().entry(&sa(9)).is_some()); sess2.flush(); } diff --git a/crates/rbitcoin-net/src/ibd/peer_io.rs b/crates/rbitcoin-net/src/ibd/peer_io.rs index 50d3746e9..a8dfebc93 100644 --- a/crates/rbitcoin-net/src/ibd/peer_io.rs +++ b/crates/rbitcoin-net/src/ibd/peer_io.rs @@ -153,14 +153,17 @@ pub(crate) fn note_block_rx(slots: &mut [PeerSlot], peer: usize, wire_bytes: usi pub(crate) async fn spawn_peer( id: usize, - addr: SocketAddr, + addr: crate::NetAddr, magic: Magic, local: SocketAddr, tip_h: Option, sinks: PeerEventSinks, dialer: crate::socks::Dialer, ) -> Result { - let stream = dialer.connect_net(crate::NetAddr::Ip(addr)).await?; + let stream = dialer.connect_net(addr).await?; + let version_socket = addr + .socket_addr() + .unwrap_or_else(|| SocketAddr::from(([0, 0, 0, 0], addr.port()))); 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( @@ -168,7 +171,7 @@ pub(crate) async fn spawn_peer( stream, magic, local, - addr, + version_socket, tip_h.map(|h| h as i32).unwrap_or(0), false, &ua, @@ -422,7 +425,7 @@ pub(crate) async fn spawn_peer( Ok(PeerSlot { id, - addr, + addr: version_socket, cmd_tx, in_flight: HashSet::new(), peer_height, diff --git a/crates/rbitcoin-net/src/lib.rs b/crates/rbitcoin-net/src/lib.rs index 605fce89d..d7f554494 100644 --- a/crates/rbitcoin-net/src/lib.rs +++ b/crates/rbitcoin-net/src/lib.rs @@ -67,7 +67,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 socks::{install_i2p_dialer, Dialer}; pub use tx_relay::{ ElectrumMempoolItem, MempoolAnnounce, MempoolHub, MempoolPerfSample, MempoolTxSnapEntry, MempoolTxSnapshot, diff --git a/crates/rbitcoin-net/src/peers.rs b/crates/rbitcoin-net/src/peers.rs index f7ec98f15..cd4269d9f 100644 --- a/crates/rbitcoin-net/src/peers.rs +++ b/crates/rbitcoin-net/src/peers.rs @@ -2019,6 +2019,15 @@ impl PeerHub { self.dial(addr, typ) } + pub fn dial_net(&self, addr: crate::NetAddr, typ: PeerConnType) -> Result<(), String> { + match addr { + crate::NetAddr::Ip(ip) => self.dial(ip, typ), + crate::NetAddr::Onion { .. } | crate::NetAddr::I2p { .. } => { + self.dial_domain(addr.host_str(), addr.port(), typ) + } + } + } + /// Outbound full-relay sessions eligible for stale-tip slot rotation. /// Empty when this hub is `noban` (functional keep-alive). pub fn outbound_full_relay_ids(&self) -> Vec { diff --git a/crates/rbitcoin-net/src/seeds.rs b/crates/rbitcoin-net/src/seeds.rs index 47137a8c7..bb777a801 100644 --- a/crates/rbitcoin-net/src/seeds.rs +++ b/crates/rbitcoin-net/src/seeds.rs @@ -553,8 +553,12 @@ impl AddrMan { /// Successful BIP324 handshake. pub fn note_connected(&mut self, addr: SocketAddr) { - self.add(addr); - if let Some(f) = self.by_addr.get_mut(&NetAddr::Ip(addr)) { + self.note_connected_addr(NetAddr::Ip(addr)); + } + + pub fn note_connected_addr(&mut self, addr: NetAddr) { + self.add_addr(addr); + if let Some(f) = self.by_addr.get_mut(&addr) { f.insert(PeerFlags::HAS_CONNECTED); f.remove(PeerFlags::FAILED_LAST_CONNECT); f.remove(PeerFlags::INCOMPATIBLE); @@ -570,6 +574,10 @@ impl AddrMan { self.last_attempt.insert(NetAddr::Ip(addr), when); } + pub fn note_attempt_addr(&mut self, addr: NetAddr) { + self.last_attempt.insert(addr, Instant::now()); + } + fn recently_attempted(&self, addr: SocketAddr, now: Instant) -> bool { self.last_attempt .get(&NetAddr::Ip(addr)) @@ -578,8 +586,12 @@ impl AddrMan { /// Dial failed. `incompatible` = no v2 / protocol reject; else network/timeout. pub fn note_connect_failed(&mut self, addr: SocketAddr, incompatible: bool) { - self.add(addr); - if let Some(f) = self.by_addr.get_mut(&NetAddr::Ip(addr)) { + self.note_connect_failed_addr(NetAddr::Ip(addr), incompatible); + } + + pub fn note_connect_failed_addr(&mut self, addr: NetAddr, incompatible: bool) { + self.add_addr(addr); + if let Some(f) = self.by_addr.get_mut(&addr) { if incompatible { f.insert(PeerFlags::INCOMPATIBLE); f.remove(PeerFlags::FAILED_LAST_CONNECT); diff --git a/crates/rbitcoin-net/src/service.rs b/crates/rbitcoin-net/src/service.rs index 0a63af95e..db6fa5ea6 100644 --- a/crates/rbitcoin-net/src/service.rs +++ b/crates/rbitcoin-net/src/service.rs @@ -285,14 +285,14 @@ impl P2PNode { /// IBD / catch-up: multi-peer densify across `peers`. /// /// This is the only history-sync path. Tip-follow is [`Self::follow_from`]. - pub async fn sync(&self, peers: &[SocketAddr], cfg: IbdConfig) -> Result { + pub async fn sync(&self, peers: &[crate::NetAddr], cfg: IbdConfig) -> Result { self.sync_cancellable(peers, cfg, None).await } /// IBD with optional cooperative cancel flag (SIGINT / SIGTERM path). pub async fn sync_cancellable( &self, - peers: &[SocketAddr], + peers: &[crate::NetAddr], cfg: IbdConfig, cancel: Option>, ) -> Result { @@ -308,7 +308,7 @@ impl P2PNode { } /// IBD with default window (1024 concurrent getdata, 16/peer). - pub async fn sync_default(&self, peers: &[SocketAddr]) -> Result { + pub async fn sync_default(&self, peers: &[crate::NetAddr]) -> Result { self.sync(peers, IbdConfig::default()).await } @@ -320,8 +320,19 @@ impl P2PNode { /// any gap (e.g. blocks mined during SH materialize) is filled actively. /// Call [`Self::sync`] first when far behind (multi-thousand height IBD). pub async fn follow_from(&mut self, peer: SocketAddr) -> Result<(), NetError> { + self.follow_from_net(crate::NetAddr::Ip(peer)).await + } + + pub async fn follow_from_net(&mut self, peer: crate::NetAddr) -> Result<(), NetError> { + let target = match peer { + crate::NetAddr::Ip(addr) => DialTarget::Socket(addr), + crate::NetAddr::Onion { .. } | crate::NetAddr::I2p { .. } => DialTarget::Domain { + host: peer.host_str(), + port: peer.port(), + }, + }; let prepared = prepare_outbound_session( - DialTarget::Socket(peer), + target, self.magic, self.local_addr, self.hub.clone(), diff --git a/crates/rbitcoin-net/src/socks.rs b/crates/rbitcoin-net/src/socks.rs index d505d299f..a4c8854c2 100644 --- a/crates/rbitcoin-net/src/socks.rs +++ b/crates/rbitcoin-net/src/socks.rs @@ -3,9 +3,16 @@ use crate::error::NetError; use std::net::SocketAddr; use std::sync::Arc; +use std::sync::OnceLock; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpStream; +static I2P_DIALER: OnceLock = OnceLock::new(); + +pub fn install_i2p_dialer(dialer: crate::i2p_sam::I2pDialer) { + let _ = I2P_DIALER.set(dialer); +} + #[derive(Clone, Debug, PartialEq, Eq)] pub(crate) struct ProxyCreds { pub(crate) username: Vec, @@ -141,7 +148,10 @@ impl Dialer { self.connect_domain(&addr.host_str(), port).await } crate::NetAddr::I2p { .. } => { - Err(NetError::Encode("i2p dial requires SAM (--i2p-sam)".into())) + let dialer = I2P_DIALER + .get() + .ok_or_else(|| NetError::Encode("i2p dial requires SAM (--i2p-sam)".into()))?; + dialer.stream_connect(&addr.host_str()).await } } } diff --git a/crates/rbitcoin-net/tests/ibd_smoke.rs b/crates/rbitcoin-net/tests/ibd_smoke.rs index 2d1295a45..4518428da 100644 --- a/crates/rbitcoin-net/tests/ibd_smoke.rs +++ b/crates/rbitcoin-net/tests/ibd_smoke.rs @@ -126,7 +126,11 @@ async fn ibd_cancellable_exits_when_flag_set() { cfg.target_peers = 1; // Cancelled IBD should return Ok (partial) or complete if race finishes first. let _ = peer - .sync_cancellable(&[seed.local_addr], cfg, Some(cancel)) + .sync_cancellable( + &[rbitcoin_net::NetAddr::Ip(seed.local_addr)], + cfg, + Some(cancel), + ) .await; seed.shutdown().await; @@ -158,7 +162,10 @@ async fn ibd_unreachable_peer_errors() { cfg.connect_timeout = Duration::from_millis(150); // TEST-NET-1 documentation address — closed / non-listening. let dead: std::net::SocketAddr = "127.0.0.1:1".parse().unwrap(); - let err = node.sync(&[dead], cfg).await.unwrap_err(); + let err = node + .sync(&[rbitcoin_net::NetAddr::Ip(dead)], cfg) + .await + .unwrap_err(); let s = err.to_string(); assert!( s.contains("no peers") || s.contains("protocol") || s.contains("connect"), diff --git a/crates/rbitcoin-node/src/run.rs b/crates/rbitcoin-node/src/run.rs index 5fb4f0853..8e8ab9248 100644 --- a/crates/rbitcoin-node/src/run.rs +++ b/crates/rbitcoin-node/src/run.rs @@ -393,6 +393,9 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { .map_err(|e| NodeError::Init(format!("i2p STREAM FORWARD {port}: {e}")))?; info!("i2p STREAM FORWARD to {}", node.local_addr); } + if let Some(sam) = i2p_sam.as_ref() { + rbitcoin_net::install_i2p_dialer(sam.dialer()); + } let _i2p_sam = i2p_sam; // One Class B appender thread. Join it at shutdown so apply does not race flush. let sh_writebehind = if config.shindex { @@ -596,7 +599,7 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { const FOLLOW_CONNECT_SECS: u64 = 8; if catch_up.dial_failed_all() { for peer in targets.iter().take(follow_n) { - if let Err(e) = node.peers.dial(*peer, PeerConnType::OutboundFullRelay) { + if let Err(e) = node.peers.dial_net(*peer, PeerConnType::OutboundFullRelay) { warn!("node: follow dial {peer}: {e}"); } } @@ -616,7 +619,7 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { } result = tokio::time::timeout( Duration::from_secs(FOLLOW_CONNECT_SECS), - node.follow_from(*peer), + node.follow_from_net(*peer), ) => { match result { Ok(Ok(())) => { @@ -946,7 +949,7 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { _ = shutdown.cancelled() => break, result = tokio::time::timeout( Duration::from_secs(8), - node.follow_from(peer), + node.follow_from_net(rbitcoin_net::NetAddr::Ip(peer)), ) => { match result { Ok(Ok(())) => { @@ -965,7 +968,7 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { 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]; + let retry_peers = [rbitcoin_net::NetAddr::Ip(peer)]; tokio::select! { biased; _ = shutdown.cancelled() => break, @@ -1195,7 +1198,7 @@ fn apply_startup_index_mode( async fn run_ibd_or_skip( node: &P2PNode, - ibd_targets: &[SocketAddr], + ibd_targets: &[rbitcoin_net::NetAddr], max_out: usize, shared_peers: &std::sync::Arc>, addrman: &mut AddrMan, @@ -1707,15 +1710,14 @@ pub(crate) fn follow_dial_targets( book: &AddrMan, max: usize, occupied: &[SocketAddr], -) -> Vec { +) -> Vec { if !connect.is_empty() { - connect - .iter() - .copied() - .filter_map(rbitcoin_net::NetAddr::socket_addr) - .collect() + connect.to_vec() } else { book.take_outbound_occupied(max, occupied) + .into_iter() + .map(rbitcoin_net::NetAddr::Ip) + .collect() } } @@ -1815,13 +1817,27 @@ mod tests { 8333, ))]; let occupied = vec![SocketAddr::new(IpAddr::V4(Ipv4Addr::new(1, 2, 0, 9)), 8333)]; - let want = vec![SocketAddr::new( + let want = vec![rbitcoin_net::NetAddr::Ip(SocketAddr::new( IpAddr::V4(Ipv4Addr::new(10, 0, 0, 1)), 8333, - )]; + ))]; assert_eq!(follow_dial_targets(&connect, &am, 8, &occupied), want); } + #[test] + fn follow_dial_targets_keeps_i2p_connect() { + use std::net::{IpAddr, Ipv4Addr, SocketAddr}; + let mut am = AddrMan::new(); + am.add(SocketAddr::new(IpAddr::V4(Ipv4Addr::new(1, 2, 0, 1)), 8333)); + let i2p = rbitcoin_net::NetAddr::I2p { + dest: [7u8; 32], + port: 8333, + }; + let connect = vec![i2p]; + let got = follow_dial_targets(&connect, &am, 8, &[]); + assert_eq!(got, vec![i2p]); + } + #[test] fn follow_dial_targets_skips_occupied_group() { use std::net::{IpAddr, Ipv4Addr, SocketAddr}; @@ -1832,7 +1848,7 @@ mod tests { am.add(other); let occupied = vec![SocketAddr::new(IpAddr::V4(Ipv4Addr::new(1, 2, 0, 9)), 8333)]; let got = follow_dial_targets(&[], &am, 1, &occupied); - assert_eq!(got, vec![other]); + assert_eq!(got, vec![rbitcoin_net::NetAddr::Ip(other)]); } #[test] From de20d483bb650f546ee7a42ef6cdda911ee12c66 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sun, 20 Sep 2026 07:30:39 -0700 Subject: [PATCH 8/9] test: sync integration_multinode with NetAddr peer inputs Co-authored-by: Cursor --- .../tests/integration_multinode.rs | 18 +++++++++++++++--- 1 file changed, 15 insertions(+), 3 deletions(-) diff --git a/crates/rbitcoin-test/tests/integration_multinode.rs b/crates/rbitcoin-test/tests/integration_multinode.rs index d3acc05b0..f9ef1fdc9 100644 --- a/crates/rbitcoin-test/tests/integration_multinode.rs +++ b/crates/rbitcoin-test/tests/integration_multinode.rs @@ -126,7 +126,7 @@ async fn seed_chain(node: &P2PNode, blocks: u32) { /// IBD from a single peer (test helper). async fn sync_ibd(node: &P2PNode, peer: SocketAddr) -> u32 { - node.sync(&[peer], IbdConfig::for_test()) + node.sync(&[rbitcoin_net::NetAddr::Ip(peer)], IbdConfig::for_test()) .await .expect("ibd sync") } @@ -1585,7 +1585,13 @@ async fn ibd_two_peers() { let client = start_node(&peer_dir).await; let n = client - .sync(&[seed.local_addr, mid.local_addr], IbdConfig::for_test()) + .sync( + &[ + rbitcoin_net::NetAddr::Ip(seed.local_addr), + rbitcoin_net::NetAddr::Ip(mid.local_addr), + ], + IbdConfig::for_test(), + ) .await .expect("ibd"); assert!(n >= 8, "accepted {n}"); @@ -1619,7 +1625,13 @@ async fn ibd_skips_dead_peer() { let peer = start_node(&peer_dir).await; let bad: SocketAddr = "127.0.0.1:1".parse().unwrap(); let n = peer - .sync(&[bad, seed.local_addr], IbdConfig::for_test()) + .sync( + &[ + rbitcoin_net::NetAddr::Ip(bad), + rbitcoin_net::NetAddr::Ip(seed.local_addr), + ], + IbdConfig::for_test(), + ) .await .expect("ibd with bad+good"); assert!(n >= 4, "downloaded {n}"); From 56ce20b130daea1f7935d3810acaa2eb6a82e14f Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sun, 20 Sep 2026 19:51:23 -0700 Subject: [PATCH 9/9] net: recreate the I2P SAM session after i2pd drops it OnceLock kept a dead session id forever. Store DESTINATION, keep the control socket on reconnect, and retry STREAM CONNECT once on INVALID_ID. Co-authored-by: Cursor --- crates/rbitcoin-net/src/i2p_sam.rs | 157 +++++++++++++++++++++++++++-- crates/rbitcoin-net/src/socks.rs | 10 +- 2 files changed, 152 insertions(+), 15 deletions(-) diff --git a/crates/rbitcoin-net/src/i2p_sam.rs b/crates/rbitcoin-net/src/i2p_sam.rs index 62bba2c04..a30489a8c 100644 --- a/crates/rbitcoin-net/src/i2p_sam.rs +++ b/crates/rbitcoin-net/src/i2p_sam.rs @@ -4,12 +4,21 @@ use crate::error::NetError; use std::io::Write; use std::net::SocketAddr; use std::path::Path; +use std::sync::Mutex; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::net::TcpStream; +static INSTALLED: Mutex> = Mutex::new(None); + +struct Installed { + dialer: I2pDialer, + _keepalive: Option, +} + pub struct I2pSam { sam_addr: SocketAddr, session_id: String, + destination: String, _control: TcpStream, _forward: Option, } @@ -18,6 +27,61 @@ pub struct I2pSam { pub struct I2pDialer { sam_addr: SocketAddr, session_id: String, + destination: String, +} + +pub fn install(dialer: I2pDialer) { + *INSTALLED.lock().unwrap_or_else(|e| e.into_inner()) = Some(Installed { + dialer, + _keepalive: None, + }); +} + +fn install_kept(dialer: I2pDialer, keepalive: TcpStream) { + *INSTALLED.lock().unwrap_or_else(|e| e.into_inner()) = Some(Installed { + dialer, + _keepalive: Some(keepalive), + }); +} + +pub async fn stream_connect_installed(dest_b32: &str) -> Result { + let first = installed()?; + match first.stream_connect(dest_b32).await { + Ok(s) => Ok(s), + Err(e) if session_dead(&e) => { + let (fresh, keepalive) = first.recreate_session().await?; + install_kept(fresh.clone(), keepalive); + fresh.stream_connect(dest_b32).await + } + Err(e) => Err(e), + } +} + +fn installed() -> Result { + INSTALLED + .lock() + .unwrap_or_else(|e| e.into_inner()) + .as_ref() + .map(|s| s.dialer.clone()) + .ok_or_else(|| NetError::Encode("i2p dial requires SAM (--i2p-sam)".into())) +} + +fn session_dead(err: &NetError) -> bool { + match err { + NetError::Io(_) | NetError::Disconnected => true, + NetError::Encode(s) => { + let up = s.to_ascii_uppercase(); + up.contains("INVALID_ID") + || up.contains("CONNECTION CLOSED") + || up.contains("STREAM CONNECT:") + } + _ => false, + } +} + +#[cfg(test)] +fn clear_installed() { + *INSTALLED.lock().unwrap_or_else(|e| e.into_inner()) = None; } impl I2pSam { @@ -86,6 +150,7 @@ impl I2pSam { Self { sam_addr, session_id, + destination: destination.clone(), _control: control, _forward: None, }, @@ -101,9 +166,15 @@ impl I2pSam { I2pDialer { sam_addr: self.sam_addr, session_id: self.session_id.clone(), + destination: self.destination.clone(), } } + fn into_keepalive(self) -> (I2pDialer, TcpStream) { + let dialer = self.dialer(); + (dialer, self._control) + } + pub async fn stream_forward(&mut self, port: u16) -> Result<(), NetError> { let mut s = TcpStream::connect(self.sam_addr) .await @@ -143,6 +214,11 @@ impl I2pDialer { } Ok(s) } + + async fn recreate_session(&self) -> Result<(Self, TcpStream), NetError> { + let (sam, _) = I2pSam::connect_session_dest(self.sam_addr, Some(&self.destination)).await?; + Ok(sam.into_keepalive()) + } } fn sam_kv<'a>(line: &'a str, key: &str) -> Option<&'a str> { @@ -224,23 +300,32 @@ async fn read_line(s: &mut TcpStream) -> Result { #[cfg(test)] mod tests { use super::*; + use std::collections::HashSet; use std::sync::{Arc, Mutex}; use std::time::{SystemTime, UNIX_EPOCH}; use tokio::net::TcpListener; const FAKE_DEST: &str = "fakeprivdest"; + static INSTALL_GATE: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); - async fn fake_sam(ok_hello: bool, dest_log: Arc>>) -> SocketAddr { + async fn fake_sam( + ok_hello: bool, + dest_log: Arc>>, + ) -> (SocketAddr, Arc>>) { let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let addr = listener.local_addr().unwrap(); + let live = Arc::new(Mutex::new(HashSet::new())); + let live_accept = Arc::clone(&live); tokio::spawn(async move { loop { let Ok((mut s, _)) = listener.accept().await else { break; }; let log = Arc::clone(&dest_log); + let live = Arc::clone(&live_accept); let ok = ok_hello; tokio::spawn(async move { + let mut created_id: Option = None; loop { let line = match read_line(&mut s).await { Ok(l) => l, @@ -257,6 +342,10 @@ mod tests { } } else if up.starts_with("SESSION CREATE") { log.lock().unwrap().push(line.clone()); + if let Some(id) = sam_kv(&line, "ID") { + live.lock().unwrap().insert(id.to_string()); + created_id = Some(id.to_string()); + } let dest = sam_kv(&line, "DESTINATION").unwrap_or("TRANSIENT"); let reply_dest = if dest.eq_ignore_ascii_case("TRANSIENT") { FAKE_DEST @@ -273,16 +362,25 @@ mod tests { let _ = write_line(&mut s, "STREAM STATUS RESULT=OK").await; } else if up.starts_with("STREAM CONNECT") { log.lock().unwrap().push(line.clone()); - let _ = write_line(&mut s, "STREAM STATUS RESULT=OK").await; + let id = sam_kv(&line, "ID").unwrap_or(""); + let known = live.lock().unwrap().contains(id); + if known { + let _ = write_line(&mut s, "STREAM STATUS RESULT=OK").await; + } else { + let _ = write_line(&mut s, "STREAM STATUS RESULT=INVALID_ID").await; + } break; } else { let _ = write_line(&mut s, "PING").await; } } + if let Some(id) = created_id { + live.lock().unwrap().remove(&id); + } }); } }); - addr + (addr, live) } fn stream_lines(log: &Arc>>, prefix: &str) -> Vec { @@ -297,7 +395,7 @@ mod tests { #[tokio::test] async fn i2p_sam_stream_connect_fake() { let log = Arc::new(Mutex::new(Vec::new())); - let addr = fake_sam(true, Arc::clone(&log)).await; + let (addr, _live) = fake_sam(true, Arc::clone(&log)).await; let sam = I2pSam::connect(addr).await.unwrap(); let dest = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa.b32.i2p"; sam.stream_connect(dest).await.unwrap(); @@ -309,7 +407,7 @@ mod tests { assert!(g.contains("ID=rbtc"), "{g}"); } - let bad = fake_sam(false, Arc::new(Mutex::new(Vec::new()))).await; + let (bad, _live) = fake_sam(false, Arc::new(Mutex::new(Vec::new()))).await; let err = match I2pSam::connect(bad).await { Err(e) => e, Ok(_) => panic!("bad HELLO must fail"), @@ -323,8 +421,10 @@ mod tests { #[tokio::test] async fn dial_i2p_uses_sam() { + let _gate = INSTALL_GATE.lock().await; + clear_installed(); let log = Arc::new(Mutex::new(Vec::new())); - let addr = fake_sam(true, Arc::clone(&log)).await; + let (addr, _live) = fake_sam(true, Arc::clone(&log)).await; let sam = I2pSam::connect(addr).await.unwrap(); let peer = crate::NetAddr::I2p { dest: [0u8; 32], @@ -344,7 +444,7 @@ mod tests { #[tokio::test] async fn i2p_accept_incoming_forwards_to_loopback() { let log = Arc::new(Mutex::new(Vec::new())); - let addr = fake_sam(true, Arc::clone(&log)).await; + let (addr, _live) = fake_sam(true, Arc::clone(&log)).await; let dir = std::env::temp_dir().join(format!( "rbtc-i2p-{}-{}", std::process::id(), @@ -385,4 +485,47 @@ mod tests { ); let _ = std::fs::remove_dir_all(&dir); } + + #[tokio::test] + async fn i2p_sam_reconnects_after_invalid_id() { + let _gate = INSTALL_GATE.lock().await; + clear_installed(); + let log = Arc::new(Mutex::new(Vec::new())); + let (addr, live) = fake_sam(true, Arc::clone(&log)).await; + let sam = I2pSam::connect(addr).await.unwrap(); + crate::socks::install_i2p_dialer(sam.dialer()); + let peer = crate::NetAddr::I2p { + dest: [0u8; 32], + port: 8333, + }; + crate::socks::Dialer::Direct + .connect_net(peer) + .await + .unwrap(); + assert_eq!(stream_lines(&log, "SESSION CREATE").len(), 1); + + live.lock().unwrap().clear(); + crate::socks::Dialer::Direct + .connect_net(peer) + .await + .expect("SAM INVALID_ID must recreate the session and retry"); + let creates = stream_lines(&log, "SESSION CREATE"); + assert_eq!(creates.len(), 2, "{creates:?}"); + assert!( + creates[1].contains(&format!("DESTINATION={FAKE_DEST}")), + "{}", + creates[1] + ); + + crate::socks::Dialer::Direct + .connect_net(peer) + .await + .unwrap(); + assert_eq!( + stream_lines(&log, "SESSION CREATE").len(), + 2, + "published dialer must keep the new session id" + ); + clear_installed(); + } } diff --git a/crates/rbitcoin-net/src/socks.rs b/crates/rbitcoin-net/src/socks.rs index a4c8854c2..448799dff 100644 --- a/crates/rbitcoin-net/src/socks.rs +++ b/crates/rbitcoin-net/src/socks.rs @@ -3,14 +3,11 @@ use crate::error::NetError; use std::net::SocketAddr; use std::sync::Arc; -use std::sync::OnceLock; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpStream; -static I2P_DIALER: OnceLock = OnceLock::new(); - pub fn install_i2p_dialer(dialer: crate::i2p_sam::I2pDialer) { - let _ = I2P_DIALER.set(dialer); + crate::i2p_sam::install(dialer); } #[derive(Clone, Debug, PartialEq, Eq)] @@ -148,10 +145,7 @@ impl Dialer { self.connect_domain(&addr.host_str(), port).await } crate::NetAddr::I2p { .. } => { - let dialer = I2P_DIALER - .get() - .ok_or_else(|| NetError::Encode("i2p dial requires SAM (--i2p-sam)".into()))?; - dialer.stream_connect(&addr.host_str()).await + crate::i2p_sam::stream_connect_installed(&addr.host_str()).await } } }