diff --git a/OPERATOR.md b/OPERATOR.md index 836cb7652..67df7062e 100644 --- a/OPERATOR.md +++ b/OPERATOR.md @@ -443,6 +443,14 @@ mempool_size_mb=100 (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. +`--proxy` or `--onion` also turns on **isolated local-tx broadcast**: +`sendrawtransaction`, Electrum `transaction.broadcast`, and Esplora +`POST /tx` (and packages) are not INV'd on standing peers. After mempool +accept the node opens a short-lived SOCKS circuit (fresh isolation +credentials), BIP324-handshakes one or two AddrMan peers (onion first), +sends `tx`, and disconnects. This is **not** Dandelion++. If that +one-shot fails, the tx stays in the mempool and is still not INV'd; +confirmation can still arrive in a block. `--onion HOST:PORT` stores a separate SOCKS endpoint for onion destinations. `--only-net onion` (repeatable with `ipv4`/`ipv6`/`i2p`) filters dial and learn; onion requires `--proxy` or `--onion`. `--connect foo.onion:8333` is a start diff --git a/crates/rbitcoin-electrum/src/server.rs b/crates/rbitcoin-electrum/src/server.rs index 215a5ea45..02260d3dc 100644 --- a/crates/rbitcoin-electrum/src/server.rs +++ b/crates/rbitcoin-electrum/src/server.rs @@ -1438,6 +1438,89 @@ fn features_hosts_json(config: &ElectrumConfig) -> Value { } } +fn broadcast_raw_tx( + params: &Value, + config: &ElectrumConfig, + mempool: Option<&MempoolHub>, +) -> Result { + let raw_hex = param_str(params, 0)?; + if raw_hex.len() > config.max_broadcast_hex { + return Err(format!( + "transaction hex too large (max {} chars)", + config.max_broadcast_hex + )); + } + let raw = rbitcoin_primitives::hex_decode(raw_hex).map_err(|e| e.to_string())?; + if raw.len() > 4_000_000 { + return Err("transaction too large".into()); + } + let tx: bitcoin::Transaction = + bitcoin::consensus::deserialize(&raw).map_err(|e| e.to_string())?; + let mp = mempool.ok_or_else(|| "mempool not available".to_string())?; + let r = mp + .accept_tx(&tx) + .map_err(|e| format!("broadcast reject: {e}"))?; + mp.mark_local_origin(r.txid); + Ok(json!(format!("{}", r.txid))) +} + +fn broadcast_package( + params: &Value, + config: &ElectrumConfig, + mempool: Option<&MempoolHub>, +) -> Result { + let arr = params + .as_array() + .and_then(|a| a.first()) + .and_then(|v| v.as_array()) + .ok_or_else(|| "broadcast_package expected array of hex txs".to_string())?; + let verbose = params + .as_array() + .and_then(|a| a.get(1)) + .and_then(|v| v.as_bool()) + .unwrap_or(false); + let mut txs = Vec::with_capacity(arr.len()); + let mut total_hex = 0usize; + for h in arr { + let raw_hex = h + .as_str() + .ok_or_else(|| "broadcast_package tx must be hex".to_string())?; + total_hex = total_hex.saturating_add(raw_hex.len()); + if total_hex > config.max_broadcast_hex { + return Err("package hex too large".into()); + } + let raw = rbitcoin_primitives::hex_decode(raw_hex).map_err(|e| e.to_string())?; + if raw.len() > 4_000_000 { + return Err("transaction too large".into()); + } + let tx: bitcoin::Transaction = + bitcoin::consensus::deserialize(&raw).map_err(|e| e.to_string())?; + txs.push(tx); + } + let mp = mempool.ok_or_else(|| "mempool not available".to_string())?; + let accepted = mp + .accept_package(&txs) + .map_err(|e| format!("broadcast_package reject: {e}"))?; + for r in &accepted { + mp.mark_local_origin(r.txid); + } + if verbose { + let mut tx_results = serde_json::Map::new(); + for r in &accepted { + tx_results.insert( + r.txid.to_string(), + json!({"txid": r.txid.to_string(), "allowed": true}), + ); + } + Ok(json!({ + "package_msg": "success", + "tx-results": tx_results, + })) + } else { + Ok(json!("success")) + } +} + #[allow(clippy::too_many_arguments)] // call-site args stay unbundled fn dispatch_pinned( method: &str, @@ -1727,77 +1810,8 @@ fn dispatch_pinned( "pos": proof.pos, })) } - "blockchain.transaction.broadcast" => { - let raw_hex = param_str(params, 0)?; - if raw_hex.len() > config.max_broadcast_hex { - return Err(format!( - "transaction hex too large (max {} chars)", - config.max_broadcast_hex - )); - } - let raw = rbitcoin_primitives::hex_decode(raw_hex).map_err(|e| e.to_string())?; - // Consensus max block weight is 4M; reject absurd raw sizes early. - if raw.len() > 4_000_000 { - return Err("transaction too large".into()); - } - let tx: bitcoin::Transaction = - bitcoin::consensus::deserialize(&raw).map_err(|e| e.to_string())?; - let mp = mempool.ok_or_else(|| "mempool not available".to_string())?; - let r = mp - .accept_tx(&tx) - .map_err(|e| format!("broadcast reject: {e}"))?; - let _ = chain.network; - Ok(json!(format!("{}", r.txid))) - } - "blockchain.transaction.broadcast_package" => { - let arr = params - .as_array() - .and_then(|a| a.first()) - .and_then(|v| v.as_array()) - .ok_or_else(|| "broadcast_package expected array of hex txs".to_string())?; - let verbose = params - .as_array() - .and_then(|a| a.get(1)) - .and_then(|v| v.as_bool()) - .unwrap_or(false); - let mut txs = Vec::with_capacity(arr.len()); - let mut total_hex = 0usize; - for h in arr { - let raw_hex = h - .as_str() - .ok_or_else(|| "broadcast_package tx must be hex".to_string())?; - total_hex = total_hex.saturating_add(raw_hex.len()); - if total_hex > config.max_broadcast_hex { - return Err("package hex too large".into()); - } - let raw = rbitcoin_primitives::hex_decode(raw_hex).map_err(|e| e.to_string())?; - if raw.len() > 4_000_000 { - return Err("transaction too large".into()); - } - let tx: bitcoin::Transaction = - bitcoin::consensus::deserialize(&raw).map_err(|e| e.to_string())?; - txs.push(tx); - } - let mp = mempool.ok_or_else(|| "mempool not available".to_string())?; - let accepted = mp - .accept_package(&txs) - .map_err(|e| format!("broadcast_package reject: {e}"))?; - if verbose { - let mut tx_results = serde_json::Map::new(); - for r in &accepted { - tx_results.insert( - r.txid.to_string(), - json!({"txid": r.txid.to_string(), "allowed": true}), - ); - } - Ok(json!({ - "package_msg": "success", - "tx-results": tx_results, - })) - } else { - Ok(json!("success")) - } - } + "blockchain.transaction.broadcast" => broadcast_raw_tx(params, config, mempool), + "blockchain.transaction.broadcast_package" => broadcast_package(params, config, mempool), "mempool.get_info" => { let min = MempoolHub::relay_fee_btc_per_kb(); let unbroadcast = mempool.map(|m| m.unbroadcast_count()).unwrap_or(0); diff --git a/crates/rbitcoin-esplora/src/handlers.rs b/crates/rbitcoin-esplora/src/handlers.rs index 2533081e2..c0c4480cb 100644 --- a/crates/rbitcoin-esplora/src/handlers.rs +++ b/crates/rbitcoin-esplora/src/handlers.rs @@ -1723,6 +1723,7 @@ async fn admit_broadcast(st: AppState, tx: bitcoin::Transaction) -> Response { }; match mp.accept_tx_async(tx).await { Ok(r) => { + mp.mark_local_origin(r.txid); let tid = r.txid.to_byte_array(); plain_ok(block_hash_hex(&tid)) } @@ -1959,6 +1960,9 @@ pub async fn post_tx_package(State(st): State, body: Bytes) -> Respons } match mp.accept_package_async(txs).await { Ok(results) => { + for r in &results { + mp.mark_local_origin(r.txid); + } let txids: Vec = results .iter() .map(|r| block_hash_hex(&r.txid.to_byte_array())) diff --git a/crates/rbitcoin-net/src/ephemeral.rs b/crates/rbitcoin-net/src/ephemeral.rs new file mode 100644 index 000000000..0c628a60d --- /dev/null +++ b/crates/rbitcoin-net/src/ephemeral.rs @@ -0,0 +1,493 @@ +//! One-shot isolated SOCKS broadcast of locally submitted transactions. + +use crate::error::NetError; +use crate::peer::{connect_and_handshake_timed, HandshakePolicy, HANDSHAKE_TIMEOUT}; +use crate::seeds::AddrMan; +use crate::socks::Dialer; +use crate::tx_relay::MempoolHub; +use crate::v2::write_v2_msg; +use crate::NetAddr; +use bitcoin::p2p::message::NetworkMessage; +use bitcoin::p2p::Magic; +use bitcoin::Transaction; +use std::net::SocketAddr; +use std::sync::{Arc, Mutex}; +use std::time::Duration; +use tokio::task::JoinHandle; + +pub(crate) const ISOLATED_BROADCAST_PEERS: usize = 2; + +pub(crate) fn isolated_broadcast_targets(am: &AddrMan, max: usize) -> Vec { + let mut onions = Vec::new(); + let mut ips = Vec::new(); + for e in am.entries() { + match e.addr { + NetAddr::Onion { .. } => onions.push(e.addr), + NetAddr::Ip(_) => ips.push(e.addr), + NetAddr::I2p { .. } => {} + } + } + let mut out = Vec::new(); + for a in onions.into_iter().chain(ips) { + if out.len() >= max { + break; + } + out.push(a); + } + out +} + +pub(crate) async fn send_tx_isolated( + dialer: &Dialer, + target: NetAddr, + magic: Magic, + tx: Transaction, + user_agent: &str, +) -> Result<(), NetError> { + send_tx_isolated_timed(dialer, target, magic, tx, user_agent, HANDSHAKE_TIMEOUT).await +} + +pub(crate) async fn send_tx_isolated_timed( + dialer: &Dialer, + target: NetAddr, + magic: Magic, + tx: Transaction, + user_agent: &str, + limit: Duration, +) -> Result<(), NetError> { + let stream = dialer.connect_isolated_net(target).await?; + let our = stream + .local_addr() + .unwrap_or_else(|_| SocketAddr::from(([127, 0, 0, 1], 0))); + let their = target + .socket_addr() + .unwrap_or_else(|| SocketAddr::from(([0, 0, 0, 0], target.port()))); + let (_ver, _reader, mut writer, _wire, tcp_shutdown) = connect_and_handshake_timed( + limit, + stream, + magic, + our, + their, + 0, + false, + user_agent, + HandshakePolicy::plain(), + ) + .await?; + write_v2_msg(&mut writer, NetworkMessage::Tx(tx)).await?; + let _ = tcp_shutdown.shutdown(std::net::Shutdown::Both); + Ok(()) +} + +async fn isolated_broadcast_known_tx( + dialer: &Dialer, + addrman: &Mutex, + magic: Magic, + user_agent: &str, + txid: bitcoin::Txid, + tx: Transaction, +) { + let targets = { + let am = addrman.lock().unwrap_or_else(|e| e.into_inner()); + isolated_broadcast_targets(&am, ISOLATED_BROADCAST_PEERS) + }; + if targets.is_empty() { + rbitcoin_log::warn!("isolated broadcast {txid}: no AddrMan targets"); + return; + } + for t in targets { + if let Err(e) = send_tx_isolated(dialer, t, magic, tx.clone(), user_agent).await { + rbitcoin_log::warn!("isolated broadcast {txid} to {t}: {e}"); + } + } +} + +pub fn spawn_isolated_broadcast_loop( + mp: Arc, + dialer: Dialer, + addrman: Arc>, + magic: Magic, + user_agent: String, +) -> JoinHandle<()> { + if !mp.isolated_broadcast() { + return tokio::spawn(async {}); + } + let mut rx = mp.subscribe_isolated(); + tokio::spawn(async move { + while let Ok(txid) = rx.recv().await { + let Some(tx) = mp.get_tx(&txid) else { + continue; + }; + isolated_broadcast_known_tx(&dialer, &addrman, magic, &user_agent, txid, tx).await; + } + }) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::msg_decode::decode_framed_offload; + use crate::peer::connect_and_handshake_timed; + use crate::v2::read_v2_frame; + use bitcoin::absolute::LockTime; + use bitcoin::p2p::message::NetworkMessage; + use bitcoin::p2p::Magic; + use bitcoin::transaction::Version as TxVersion; + use bitcoin::{Amount, Transaction, TxOut}; + use bitcoin::{ScriptBuf, Sequence, TxIn, Witness}; + use std::net::{IpAddr, Ipv4Addr, SocketAddr}; + use std::sync::{Arc, Mutex}; + use std::time::Duration; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::net::{TcpListener, TcpStream}; + + fn dummy_tx() -> Transaction { + Transaction { + version: TxVersion::TWO, + lock_time: LockTime::ZERO, + input: vec![TxIn { + previous_output: bitcoin::OutPoint::null(), + script_sig: ScriptBuf::new(), + sequence: Sequence::MAX, + witness: Witness::new(), + }], + output: vec![TxOut { + value: Amount::from_sat(50), + script_pubkey: ScriptBuf::new(), + }], + } + } + + async fn socks_userpass(s: &mut TcpStream) -> Vec { + 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(&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(); + let mut hdr = [0u8; 4]; + s.read_exact(&mut hdr).await.unwrap(); + match hdr[3] { + 1 => { + let mut a = [0u8; 4]; + s.read_exact(&mut a).await.unwrap(); + let mut p = [0u8; 2]; + s.read_exact(&mut p).await.unwrap(); + } + 3 => { + let mut n = [0u8; 1]; + s.read_exact(&mut n).await.unwrap(); + let mut host = vec![0u8; n[0] as usize]; + s.read_exact(&mut host).await.unwrap(); + let mut p = [0u8; 2]; + s.read_exact(&mut p).await.unwrap(); + } + _ => panic!("unexpected ATYP {}", hdr[3]), + } + s.write_all(&[5, 0, 0, 1, 0, 0, 0, 0, 0, 0]).await.unwrap(); + user + } + + async fn splice_after_userpass(mut c: TcpStream, dest: SocketAddr) -> Vec { + let user = socks_userpass(&mut c).await; + let mut peer = TcpStream::connect(dest).await.unwrap(); + let _ = tokio::io::copy_bidirectional(&mut c, &mut peer).await; + user + } + + #[test] + fn isolated_broadcast_targets_prefers_onion() { + let mut am = AddrMan::new(); + am.add(SocketAddr::from((Ipv4Addr::new(1, 2, 3, 4), 8333))); + let onion: NetAddr = "pg6mmjiyjmcrsslvykfwnntlaru7p5svn6y2ymmju6nubxndf4pscryd.onion:8333" + .parse() + .unwrap(); + am.add_addr(onion); + let got = isolated_broadcast_targets(&am, 1); + assert_eq!(got, vec![onion]); + let got = isolated_broadcast_targets(&am, 2); + assert_eq!(got[0], onion); + assert!(matches!(got[1], NetAddr::Ip(_))); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn ephemeral_broadcast_one_shot_tx_and_new_socks_creds() { + 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(); + let (_ver, mut reader, _writer, _wire, _tcp) = connect_and_handshake_timed( + Duration::from_secs(5), + stream, + Magic::REGTEST, + peer_addr, + from, + 0, + true, + "/rbitcoin:test/", + HandshakePolicy::plain(), + ) + .await + .unwrap(); + let frame = tokio::time::timeout( + Duration::from_secs(5), + read_v2_frame(&mut reader, Magic::REGTEST), + ) + .await + .expect("tx frame") + .expect("tx decrypt"); + let msg = decode_framed_offload(frame).await.unwrap(); + assert!( + matches!(msg.payload(), NetworkMessage::Tx(_)), + "expected tx, got {:?}", + msg.payload() + ); + let eof = tokio::time::timeout( + Duration::from_secs(5), + read_v2_frame(&mut reader, Magic::REGTEST), + ) + .await + .ok() + .and_then(Result::ok) + .is_none(); + eof + }); + + let socks_l = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let proxy = socks_l.local_addr().unwrap(); + let (cred_tx, mut cred_rx) = tokio::sync::mpsc::channel::>(2); + let splice = tokio::spawn(async move { + let (mut s, _) = socks_l.accept().await.unwrap(); + let u = socks_userpass(&mut s).await; + cred_tx.send(u).await.unwrap(); + let (s, _) = socks_l.accept().await.unwrap(); + let u = splice_after_userpass(s, peer_addr).await; + cred_tx.send(u).await.unwrap(); + }); + + let dialer = Dialer::socks(proxy, true); + let _standing = dialer.connect(peer_addr).await.unwrap(); + let standing_user = cred_rx.recv().await.unwrap(); + + let tx = dummy_tx(); + let want = tx.compute_txid(); + send_tx_isolated_timed( + &dialer, + NetAddr::Ip(peer_addr), + Magic::REGTEST, + tx, + "/rbitcoin:test/", + Duration::from_secs(5), + ) + .await + .unwrap(); + let iso_user = cred_rx.recv().await.unwrap(); + assert_ne!( + standing_user, iso_user, + "isolated dial must use fresh SOCKS creds" + ); + assert!(!standing_user.is_empty() && !iso_user.is_empty()); + + let eof = inbound.await.unwrap(); + assert!(eof, "one-shot must disconnect after tx ({want})"); + splice.abort(); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn ephemeral_broadcast_fail_does_not_inv_standing() { + use bitcoin::p2p::address::Address; + use bitcoin::p2p::message_network::VersionMessage; + use bitcoin::p2p::ServiceFlags; + use bitcoin::OutPoint; + use rbitcoin_primitives::Height; + use tokio::sync::mpsc; + + let (dir, hub) = crate::chain::tiny_regtest_hub_labeled("iso-fail-inv"); + hub.ensure_genesis().unwrap(); + hub.generate_to_script(102, ScriptBuf::from_bytes(vec![0x51]), vec![]) + .expect("pad"); + let mp = MempoolHub::open(dir.join("mp"), Arc::clone(&hub.query)).unwrap(); + mp.set_relay_enabled(true); + mp.set_isolated_broadcast(true); + assert!(hub.attach_mempool(mp).is_ok()); + let cb = hub + .query + .reconstruct_block_at_height(Height(1)) + .unwrap() + .txdata[0] + .compute_txid(); + let local = Transaction { + version: TxVersion::TWO, + lock_time: LockTime::ZERO, + input: vec![TxIn { + previous_output: OutPoint { txid: cb, vout: 0 }, + script_sig: ScriptBuf::new(), + sequence: Sequence::ENABLE_RBF_NO_LOCKTIME, + witness: Witness::new(), + }], + output: vec![TxOut { + value: Amount::from_sat(49_9999_0000), + script_pubkey: ScriptBuf::from_bytes(vec![0x51]), + }], + }; + hub.mempool().unwrap().accept_tx(&local).expect("local"); + hub.mempool() + .unwrap() + .mark_local_origin(local.compute_txid()); + + let closed = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let dest = closed.local_addr().unwrap(); + drop(closed); + let err = send_tx_isolated_timed( + &Dialer::Direct, + NetAddr::Ip(dest), + Magic::REGTEST, + local.clone(), + "/rbitcoin:test/", + Duration::from_millis(200), + ) + .await; + assert!(err.is_err(), "closed listener must fail the one-shot"); + + let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 18444); + let ver = VersionMessage { + version: 70016, + services: ServiceFlags::NETWORK, + timestamp: 0, + receiver: Address::new(&addr, ServiceFlags::NONE), + sender: Address::new(&addr, ServiceFlags::NONE), + nonce: 1, + user_agent: "/rbitcoin:test/".into(), + start_height: 0, + relay: true, + }; + let (tx, mut rx) = mpsc::unbounded_channel(); + let peers = crate::peers::PeerHub::new(); + let sess = peers.register( + addr, + addr, + &ver, + false, + crate::peers::PeerConnType::OutboundFullRelay, + ); + sess.attach_out(tx); + crate::force_announce_txid(&hub, &peers, local.compute_txid()); + assert!( + rx.try_recv().is_err(), + "failed isolated send must not fall back to standing INV" + ); + let _ = std::fs::remove_dir_all(dir); + } + + #[tokio::test] + async fn ephemeral_broadcast_loop_idle_when_not_isolated() { + let (dir, hub) = crate::chain::tiny_regtest_hub_labeled("iso-idle"); + let mp = MempoolHub::open(dir.join("mp"), Arc::clone(&hub.query)).unwrap(); + let am = Arc::new(Mutex::new(AddrMan::new())); + spawn_isolated_broadcast_loop( + mp, + Dialer::Direct, + am, + Magic::REGTEST, + "/rbitcoin:test/".into(), + ) + .await + .unwrap(); + let _ = std::fs::remove_dir_all(dir); + } + + #[tokio::test] + async fn ephemeral_broadcast_loop_skips_unknown_txid() { + let (dir, hub) = crate::chain::tiny_regtest_hub_labeled("iso-skip"); + let mp = MempoolHub::open(dir.join("mp"), Arc::clone(&hub.query)).unwrap(); + mp.set_isolated_broadcast(true); + let am = Arc::new(Mutex::new(AddrMan::new())); + let h = spawn_isolated_broadcast_loop( + mp.clone(), + Dialer::Direct, + am, + Magic::REGTEST, + "/rbitcoin:test/".into(), + ); + tokio::time::sleep(Duration::from_millis(20)).await; + mp.mark_local_origin(dummy_tx().compute_txid()); + tokio::time::sleep(Duration::from_millis(20)).await; + h.abort(); + let _ = std::fs::remove_dir_all(dir); + } + + #[tokio::test] + async fn ephemeral_broadcast_known_tx_no_addrman_targets() { + let tx = dummy_tx(); + isolated_broadcast_known_tx( + &Dialer::Direct, + &Mutex::new(AddrMan::new()), + Magic::REGTEST, + "/rbitcoin:test/", + tx.compute_txid(), + tx, + ) + .await; + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn ephemeral_broadcast_known_tx_one_shot() { + 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(); + let (_ver, mut reader, _writer, _wire, _tcp) = connect_and_handshake_timed( + Duration::from_secs(5), + stream, + Magic::REGTEST, + peer_addr, + from, + 0, + true, + "/rbitcoin:test/", + HandshakePolicy::plain(), + ) + .await + .unwrap(); + let frame = tokio::time::timeout( + Duration::from_secs(5), + read_v2_frame(&mut reader, Magic::REGTEST), + ) + .await + .expect("tx frame") + .expect("tx decrypt"); + let msg = decode_framed_offload(frame).await.unwrap(); + assert!( + matches!(msg.payload(), NetworkMessage::Tx(_)), + "expected tx, got {:?}", + msg.payload() + ); + true + }); + let mut am = AddrMan::new(); + am.add(peer_addr); + let tx = dummy_tx(); + isolated_broadcast_known_tx( + &Dialer::Direct, + &Mutex::new(am), + Magic::REGTEST, + "/rbitcoin:test/", + tx.compute_txid(), + tx, + ) + .await; + assert!(inbound.await.unwrap(), "one-shot reached the peer"); + } +} diff --git a/crates/rbitcoin-net/src/lib.rs b/crates/rbitcoin-net/src/lib.rs index d7f554494..37281926a 100644 --- a/crates/rbitcoin-net/src/lib.rs +++ b/crates/rbitcoin-net/src/lib.rs @@ -5,6 +5,7 @@ mod cache; mod chain; mod codec; mod compact; +mod ephemeral; mod error; mod eviction; mod i2p_sam; @@ -34,6 +35,7 @@ pub use compact::{ classify_v2_cmpct_peer, prefilled_indexes_ok, shortid_map_from_txs, try_reconstruct, CmpctPeerFrame, }; +pub use ephemeral::spawn_isolated_broadcast_loop; pub use error::NetError; pub use eviction::{select_inbound_eviction, InboundEvictCandidate}; pub use i2p_sam::I2pSam; diff --git a/crates/rbitcoin-net/src/peer.rs b/crates/rbitcoin-net/src/peer.rs index b51f70196..55c74178f 100644 --- a/crates/rbitcoin-net/src/peer.rs +++ b/crates/rbitcoin-net/src/peer.rs @@ -1311,7 +1311,8 @@ fn tx_announce_should_queue( mp: &crate::tx_relay::MempoolHub, txid: &bitcoin::Txid, ) -> bool { - tx_announce_peer_ok(session, mp, txid) + !mp.skip_standing_inv(txid) + && tx_announce_peer_ok(session, mp, txid) && mp.try_contains(txid) && (mp.relay_enabled() || mp.is_unbroadcast(txid)) && !tx_announce_below_feefilter(session, mp, txid) @@ -2010,6 +2011,9 @@ pub fn force_announce_txid(hub: &ChainHub, peers: &crate::peers::PeerHub, txid: let Some(mp) = hub.mempool() else { return; }; + if mp.skip_standing_inv(&txid) { + return; + } let Some(tx) = mp.try_get_tx(&txid) else { return; }; @@ -2092,6 +2096,9 @@ fn tx_inv_candidate_ok( clock_due: bool, inbound_age_gate: bool, ) -> bool { + if mp.skip_standing_inv(&txid) { + return false; + } if from_this_peer.contains_key(&txid) { return false; } diff --git a/crates/rbitcoin-net/src/peer_tests.rs b/crates/rbitcoin-net/src/peer_tests.rs index cbcfc12aa..1dfb7c17f 100644 --- a/crates/rbitcoin-net/src/peer_tests.rs +++ b/crates/rbitcoin-net/src/peer_tests.rs @@ -1413,6 +1413,98 @@ fn force_announce_txid_skips_then_invs_full_relay() { let _ = std::fs::remove_dir_all(dir); } +#[test] +fn local_origin_not_inv_on_standing_peer() { + use bitcoin::absolute::LockTime; + use bitcoin::p2p::address::Address; + use bitcoin::p2p::message::NetworkMessage; + use bitcoin::p2p::message_network::VersionMessage; + use bitcoin::p2p::ServiceFlags; + use bitcoin::script::ScriptBuf; + use bitcoin::transaction::Version as TxVersion; + use bitcoin::{Amount, OutPoint, Sequence, Transaction, TxIn, TxOut, Witness}; + use rbitcoin_primitives::Height; + use std::net::{IpAddr, Ipv4Addr, SocketAddr}; + + let (dir, hub) = crate::chain::tiny_regtest_hub_labeled("local-origin-inv"); + hub.ensure_genesis().unwrap(); + hub.generate_to_script(102, ScriptBuf::from_bytes(vec![0x51]), vec![]) + .expect("pad"); + let mp = crate::tx_relay::MempoolHub::open(dir.join("mp"), Arc::clone(&hub.query)).unwrap(); + mp.set_relay_enabled(true); + mp.set_isolated_broadcast(true); + assert!(hub.attach_mempool(mp).is_ok()); + let cb0 = hub + .query + .reconstruct_block_at_height(Height(1)) + .unwrap() + .txdata[0] + .compute_txid(); + let cb1 = hub + .query + .reconstruct_block_at_height(Height(2)) + .unwrap() + .txdata[0] + .compute_txid(); + let spend = |cb: bitcoin::Txid, fee: u64| Transaction { + version: TxVersion::TWO, + lock_time: LockTime::ZERO, + input: vec![TxIn { + previous_output: OutPoint { txid: cb, vout: 0 }, + script_sig: ScriptBuf::new(), + sequence: Sequence::ENABLE_RBF_NO_LOCKTIME, + witness: Witness::new(), + }], + output: vec![TxOut { + value: Amount::from_sat(49_9999_0000 - fee), + script_pubkey: ScriptBuf::from_bytes(vec![0x51]), + }], + }; + let local = spend(cb0, 0); + let p2p = spend(cb1, 1_000); + hub.mempool().unwrap().accept_tx(&local).expect("local"); + hub.mempool() + .unwrap() + .mark_local_origin(local.compute_txid()); + hub.mempool().unwrap().accept_tx(&p2p).expect("p2p"); + + let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 18444); + let ver = VersionMessage { + version: 70016, + services: ServiceFlags::NETWORK, + timestamp: 0, + receiver: Address::new(&addr, ServiceFlags::NONE), + sender: Address::new(&addr, ServiceFlags::NONE), + nonce: 1, + user_agent: "/rbitcoin:test/".into(), + start_height: 0, + relay: true, + }; + let (tx, mut rx) = mpsc::unbounded_channel(); + let peers = crate::peers::PeerHub::new(); + let sess = peers.register( + addr, + addr, + &ver, + false, + crate::peers::PeerConnType::OutboundFullRelay, + ); + sess.attach_out(tx); + crate::force_announce_txid(&hub, &peers, local.compute_txid()); + assert!( + rx.try_recv().is_err(), + "isolated local-origin must not INV standing peers" + ); + crate::force_announce_txid(&hub, &peers, p2p.compute_txid()); + match rx.try_recv().expect("p2p origin INV").expect_msg() { + NetworkMessage::Inv(v) => { + assert_eq!(v, vec![Inventory::WTx(p2p.compute_wtxid())]); + } + other => panic!("expected WTx inv, got {other:?}"), + } + let _ = std::fs::remove_dir_all(dir); +} + /// When relay is on, unbroadcast must not skip the inbound 30s INV gate /// (`mempool_reorg.py:71`). #[test] diff --git a/crates/rbitcoin-net/src/service.rs b/crates/rbitcoin-net/src/service.rs index db6fa5ea6..5064fdeec 100644 --- a/crates/rbitcoin-net/src/service.rs +++ b/crates/rbitcoin-net/src/service.rs @@ -239,6 +239,14 @@ impl P2PNode { self.dialer.clone() } + pub fn magic(&self) -> Magic { + self.magic + } + + pub fn user_agent(&self) -> &str { + &self.user_agent + } + /// 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-net/src/socks.rs b/crates/rbitcoin-net/src/socks.rs index 448799dff..b3b623785 100644 --- a/crates/rbitcoin-net/src/socks.rs +++ b/crates/rbitcoin-net/src/socks.rs @@ -56,6 +56,15 @@ pub(crate) async fn dial_isolated( socks5_connect(proxy, target, Some(&creds)).await } +pub(crate) async fn dial_isolated_domain( + proxy: SocketAddr, + host: &str, + port: u16, +) -> Result { + let creds = ProxyCreds::fresh(); + socks5_connect_domain(proxy, host, port, Some(&creds)).await +} + #[derive(Clone, Debug, Default, PartialEq, Eq)] pub enum Dialer { #[default] @@ -133,6 +142,26 @@ impl Dialer { } } + pub(crate) async fn connect_isolated_net( + &self, + addr: crate::NetAddr, + ) -> Result { + match addr { + crate::NetAddr::Ip(s) => self.connect_isolated(s).await, + crate::NetAddr::Onion { port, .. } => match self { + Dialer::Direct => Err(NetError::Encode( + "onion dial requires SOCKS (--proxy or --onion)".into(), + )), + Dialer::Socks { proxy, .. } => { + dial_isolated_domain(*proxy, &addr.host_str(), port).await + } + }, + crate::NetAddr::I2p { .. } => { + Err(NetError::Encode("i2p dial requires SAM (--i2p-sam)".into())) + } + } + } + pub async fn connect_net(&self, addr: crate::NetAddr) -> Result { match addr { crate::NetAddr::Ip(s) => self.connect(s).await, diff --git a/crates/rbitcoin-net/src/tx_relay.rs b/crates/rbitcoin-net/src/tx_relay.rs index d4b51481a..469cbfd01 100644 --- a/crates/rbitcoin-net/src/tx_relay.rs +++ b/crates/rbitcoin-net/src/tx_relay.rs @@ -562,6 +562,9 @@ pub struct MempoolHub { dir: PathBuf, /// Locally submitted txids not yet requested by a peer (`getmempoolinfo.unbroadcastcount`). unbroadcast: Mutex>, + local_origin: Mutex>, + isolated_broadcast: AtomicBool, + isolated_kick: broadcast::Sender, /// Wtxids re-admitted from a disconnected block. Core serves these /// even if this peer has not been INV'd yet (`mempool_reorg`). reorg_servable: Mutex>, @@ -635,6 +638,7 @@ impl MempoolHub { .map_err(|e| format!("mempool open: {e}"))?; let (announce, _) = broadcast::channel(256); let (inv_flush, _) = broadcast::channel(16); + let (isolated_kick, _) = broadcast::channel(32); let unbroadcast = if persist { load_unbroadcast_file(&dir_buf) } else { @@ -685,6 +689,9 @@ impl MempoolHub { meter_get_coin_create_mtp: AtomicU64::new(0), sh_index: Mutex::new(MempoolShIndex::new()), unbroadcast: Mutex::new(unbroadcast), + local_origin: Mutex::new(HashSet::new()), + isolated_broadcast: AtomicBool::new(false), + isolated_kick, reorg_servable: Mutex::new(HashSet::new()), relay_seq: Mutex::new(HashMap::new()), wtxid_by_txid: Mutex::new(HashMap::new()), @@ -3147,6 +3154,33 @@ impl MempoolHub { persist_unbroadcast_file(&self.dir, &u); } + pub fn mark_local_origin(&self, txid: Txid) { + self.local_origin.lock().unwrap().insert(txid); + if self.isolated_broadcast() { + let _ = self.isolated_kick.send(txid); + } + } + + pub(crate) fn subscribe_isolated(&self) -> broadcast::Receiver { + self.isolated_kick.subscribe() + } + + pub fn is_local_origin(&self, txid: &Txid) -> bool { + self.local_origin.lock().unwrap().contains(txid) + } + + pub fn set_isolated_broadcast(&self, on: bool) { + self.isolated_broadcast.store(on, Ordering::Relaxed); + } + + pub fn isolated_broadcast(&self) -> bool { + self.isolated_broadcast.load(Ordering::Relaxed) + } + + pub fn skip_standing_inv(&self, txid: &Txid) -> bool { + self.isolated_broadcast() && self.is_local_origin(txid) + } + /// Peer getdata served this txid — it is no longer unbroadcast. pub fn mark_broadcast(&self, txid: &Txid) { let mut u = self.unbroadcast.lock().unwrap(); @@ -4187,6 +4221,29 @@ mod tests { ); } + #[test] + fn local_origin_rpc_not_p2p() { + let dir = tmp(); + let store = tmp(); + let q = Arc::new(Query::open_or_create_tiny(&store).unwrap()); + let hub = MempoolHub::open(&dir, q).unwrap(); + let local = Txid::from_byte_array([0x11; 32]); + let wire = Txid::from_byte_array([0x22; 32]); + assert!(!hub.is_local_origin(&local)); + hub.mark_local_origin(local); + assert!(hub.is_local_origin(&local)); + assert!(!hub.is_local_origin(&wire)); + assert!( + !hub.skip_standing_inv(&local), + "without --proxy/--onion, local-origin still uses standing INV" + ); + hub.set_isolated_broadcast(true); + assert!(hub.skip_standing_inv(&local)); + assert!(!hub.skip_standing_inv(&wire)); + let _ = std::fs::remove_dir_all(&dir); + let _ = std::fs::remove_dir_all(&store); + } + #[test] fn hub_accept_remove_with_map_utxo_path() { // MempoolHub needs Query; use open empty store + MapUtxo via direct ActiveMempool diff --git a/crates/rbitcoin-node/src/config.rs b/crates/rbitcoin-node/src/config.rs index 105797e69..898fc0e0e 100644 --- a/crates/rbitcoin-node/src/config.rs +++ b/crates/rbitcoin-node/src/config.rs @@ -143,6 +143,13 @@ impl ListenOpts { } } + pub fn isolated_dialer(&self) -> rbitcoin_net::Dialer { + match self.proxy.or(self.onion) { + Some(proxy) => rbitcoin_net::Dialer::socks(proxy, true), + None => rbitcoin_net::Dialer::Direct, + } + } + pub fn p2p_bind_addr(&self, network: Network) -> Option { match self.p2p { P2pListen::Off => None, diff --git a/crates/rbitcoin-node/src/run.rs b/crates/rbitcoin-node/src/run.rs index 3eb512c8a..09767a0b5 100644 --- a/crates/rbitcoin-node/src/run.rs +++ b/crates/rbitcoin-node/src/run.rs @@ -312,6 +312,9 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { .map_err(|e| NodeError::Config(format!("mempool open join: {e}")))? .map_err(NodeError::Config)?; node.peers.attach_mempool(&mempool); + if config.listen.proxy.is_some() || config.listen.onion.is_some() { + mempool.set_isolated_broadcast(true); + } node.peers.set_net_perms(table.clone()); if let Some(secs) = config.listen.peer_timeout_secs { node.peers.set_peer_timeout_secs(secs); @@ -468,6 +471,15 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { } let shared_peers = std::sync::Arc::new(std::sync::Mutex::new(addrman.clone())); node.peers.set_addrman(std::sync::Arc::clone(&shared_peers)); + if mempool.isolated_broadcast() { + let _iso = rbitcoin_net::spawn_isolated_broadcast_loop( + Arc::clone(&mempool), + config.listen.isolated_dialer(), + std::sync::Arc::clone(&shared_peers), + node.magic(), + node.user_agent().to_string(), + ); + } let max_out = config.listen.max_outbound.max(1) as usize; let candidate_n = max_out.saturating_mul(2).clamp(16, 48); diff --git a/crates/rbitcoin-rpc/src/methods/mempool.rs b/crates/rbitcoin-rpc/src/methods/mempool.rs index 90164d885..93d074747 100644 --- a/crates/rbitcoin-rpc/src/methods/mempool.rs +++ b/crates/rbitcoin-rpc/src/methods/mempool.rs @@ -538,6 +538,7 @@ pub(crate) fn sendrawtransaction(ctx: &RpcContext, params: &RpcParams) -> Result match mp.accept_tx(&tx) { Ok(r) => { mp.note_unbroadcast(r.txid); + mp.mark_local_origin(r.txid); // `-blocksonly`: accept-time announce is skipped (relay off + not // yet unbroadcast). Re-announce after noting so inbound peers INV // (`p2p_blocksonly.py:48`). When relay is on, leave the 30s inbound @@ -1335,6 +1336,7 @@ fn submitpackage_admit( match res { Ok(ok) => { mp.note_unbroadcast(ok.txid); + mp.mark_local_origin(ok.txid); for old in &ok.replaced { replaced.push(hash_hex_display(&old.to_byte_array())); } diff --git a/crates/rbitcoin-rpc/src/methods_tests.rs b/crates/rbitcoin-rpc/src/methods_tests.rs index 8daa91177..aa7b81efd 100644 --- a/crates/rbitcoin-rpc/src/methods_tests.rs +++ b/crates/rbitcoin-rpc/src/methods_tests.rs @@ -1422,6 +1422,8 @@ fn mempool_graph_fields_follow_cluster_and_unbroadcast() { let sent = dispatch(&ctx, "sendrawtransaction", vec![json!(local_hex_tx)]).unwrap(); let local_hex = hash_hex_display(&local.compute_txid().to_byte_array()); assert_eq!(sent, json!(local_hex.clone())); + assert!(mp.is_local_origin(&local.compute_txid())); + assert!(!mp.is_local_origin(&parent.compute_txid())); let info = dispatch(&ctx, "getmempoolinfo", vec![]).unwrap(); assert_eq!(info["unbroadcastcount"], 1); diff --git a/nix/modules/rbitcoin.nix b/nix/modules/rbitcoin.nix index 265af5353..c58656a5c 100644 --- a/nix/modules/rbitcoin.nix +++ b/nix/modules/rbitcoin.nix @@ -220,7 +220,7 @@ in type = types.nullOr types.str; default = null; example = "127.0.0.1:9050"; - description = "SOCKS5 proxy HOST:PORT for all P2P outbound."; + description = "SOCKS5 proxy HOST:PORT for all P2P outbound. Also enables isolated local-tx broadcast (new SOCKS circuit after sendraw / Electrum / Esplora submit; not Dandelion++). Standing peers do not INV those txs."; }; onionProxy = mkOption {