From b5c430f9fa644995a97c95cdecbc2595d87de5ea Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Sun, 13 Sep 2026 13:20:39 -0700 Subject: [PATCH] mempool: admit package members sequentially for 1p1c Prepare-all left a child orphaned because the parent was not yet in the graph. Commit each member before the next admit. Pin Esplora POST /txs/package success, process gettxout/getchaintips, live /mempool*, and serving-only submitpackage refusing while relay is off. Co-authored-by: Cursor --- CHANGELOG.md | 12 ++ TESTING.md | 2 +- crates/rbitcoin-net/src/tx_relay.rs | 121 +++++++++---------- crates/rbitcoin-test/tests/cross_surface.rs | 124 ++++++++++++++++++-- 4 files changed, 182 insertions(+), 77 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index e14a3a02b..55f6f4f36 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,13 @@ before 1.0). ### Added +- **Process package + live mempool HTTP:** `esplora_broadcast_visible_in_rpc_and_electrum` + pins Esplora `POST /txs/package` 1p1c success, process `gettxout` (confirmed, + mempool create, mempool-spent hide), `getchaintips`, live `GET /mempool` / + `/mempool/txids` / `/mempool/recent` / `/fee-estimates`, and serving-only + `submitpackage` refusing while relay is off (`sendraw` still admits). + Package JSON errors and dispatch `gettxout` stay units. + - **P2P getblocks / feefilter / bloom:** live follower `getblocks` is answered with `inv`, inbound BIP133 `feefilter` is recorded, and `filterload` disconnects (bloom off). Oversize locator and MemPool/`filteradd`/`filterclear` @@ -74,6 +81,11 @@ before 1.0). ### Fixed +- **Esplora `POST /txs/package` 1p1c:** `MempoolHub::accept_package` commits each + member before the next admit so a child can spend an in-package parent + (same sequential shape as `ActiveMempool::accept_package`). Prepare-all + left the child orphaned. + - **Coverage `integration_multinode` SIGABRT:** live `P2PNode` tests serialize on a process mutex so overlapping `shutdown` abort cannot smash the shared `rbtc-scripts` pool under llvm-cov. diff --git a/TESTING.md b/TESTING.md index 9d9e666a4..903385027 100644 --- a/TESTING.md +++ b/TESTING.md @@ -282,7 +282,7 @@ Prefer **one high-level scenario** per behavior cluster. Delete lower-level test | `electrum_tweaks_subscribe_streams_then_done` | Electrum | Cake `tweaks.subscribe`: one-height result, per-height notifies, `done`, and BIP352 hash-bind on a P2WPKH→P2TR spend. Keep zero-chunk / pre-taproot units | | `electrum_max_connections_rejects_extra_client` | Electrum | TCP cap drops the extra client | | `electrum_idle_timeout_disconnects_quiet_client` | Electrum | Idle timeout closes a quiet socket | -| `esplora_broadcast_visible_in_rpc_and_electrum` | Node + Electrum + Esplora + RPC | One `run_p2p` datadir: HTTP `sendrawtransaction` / `testmempoolaccept` (allowed, missing-or-spent, min-relay, RBF too-low reject + replacement), Esplora `POST /tx` parent and mempool child appear in `getrawmempool` and Electrum mempool/history (`fee` on unconfirmed, including child `height = -1`); `generate` includes those txs; immature coinbase sendraw rejects. Keep `accept.rs` reject units and RPC dry-run orphan-count | +| `esplora_broadcast_visible_in_rpc_and_electrum` | Node + Electrum + Esplora + RPC | One `run_p2p` datadir: HTTP `sendrawtransaction` / `testmempoolaccept` (allowed, missing-or-spent, min-relay, RBF too-low reject + replacement), Esplora `POST /tx` parent and mempool child appear in `getrawmempool` and Electrum mempool/history (`fee` on unconfirmed, including child `height = -1`); process `gettxout` / `getchaintips`; Esplora `POST /txs/package` 1p1c; serving-only `submitpackage` refuses (relay off); live `GET /mempool` / `/mempool/txids` / `/mempool/recent` / `/fee-estimates`; `generate` includes those txs; immature coinbase sendraw rejects. Keep `accept.rs` reject units, package JSON errors, RPC dry-run orphan-count, and dispatch `gettxout` | | `two_node_header_and_block_sync` | P2P (**default**) | Seeder → peer 8-block IBD; peer `last_write` meter | | `p2p_timeout_getaddr_and_keepalive_ping` | P2P (**default**) | One pad: v1-magic inbound drops at `peertimeout=1`, obsolete VERSION and pre-verack ping close the peer, full-relay GetAddr cache 1000, headers-sync stall replace, self-connect refuses, AddrFetch `getaddr`/`addrv2` (no `getheaders`), one keepalive ping/pong. Handshake **format** needles stay. Sole-preferred stall KEEP stays a PeerHub unit. | | `p2p_compact_hb_getblocktxn_and_orphan` | P2P (**default**) | One mature pad: HB coinbase `cmpctblock`, 2-tx compact → `getblocktxn` + connect, orphan child GetData then parent accept (INV AlreadyHave), then live `getblocks` → `inv`, inbound `feefilter`, `filterload` disconnect. Oversize locator and MemPool/`filteradd`/`filterclear` stay PeerHub units. Does **not** pin depth-10 full-block serve, tokio-worker lock, or park-not-reject logs | diff --git a/crates/rbitcoin-net/src/tx_relay.rs b/crates/rbitcoin-net/src/tx_relay.rs index 54bd49697..aac464132 100644 --- a/crates/rbitcoin-net/src/tx_relay.rs +++ b/crates/rbitcoin-net/src/tx_relay.rs @@ -1588,7 +1588,8 @@ impl MempoolHub { let utxo = self.utxo_provider(); let mut stages = rbitcoin_mempool::AcceptStageUs::default(); let mut lock_us = 0u64; - let mut preps = Vec::with_capacity(txs.len()); + let mut accepted: Vec = Vec::with_capacity(txs.len()); + let mut prevouts: Vec> = Vec::with_capacity(txs.len()); for tx in txs { utxo.note_spender(tx); let delta = self.fee_delta(&tx.compute_txid()); @@ -1601,85 +1602,71 @@ impl MempoolHub { let prep = match self.admit_staged(tx, &utxo, spec, &mut stages, &mut lock_us) { Ok(p) => p, Err(e) => { + if !accepted.is_empty() { + let mut g = self.lock_write(); + for r in accepted.iter().rev() { + let _ = g.remove_txid(&r.txid); + } + } let us = t0.elapsed().as_micros() as u64; self.meter_accept_stages(lock_us, stages); return Err(self.finish_accept_err(us, e).unwrap_err()); } }; - preps.push(prep); - } - let prevouts: Vec> = preps.iter().map(|p| p.prevouts.clone()).collect(); - let result = { + let prev = prep.prevouts.clone(); let t_lock = Instant::now(); - let mut g = self.lock_write(); - g.last_accept_stages = stages; - let mut accepted: Vec = Vec::with_capacity(txs.len()); - let mut err = None; - for (tx, prep) in txs.iter().zip(preps) { - match g.commit_after_script(tx, prep) { - Ok(r) => accepted.push(r), - Err(e) => { - for r in accepted.iter().rev() { - let _ = g.remove_txid(&r.txid); - } - err = Some(e); - break; + let commit = { + let mut g = self.lock_write(); + g.last_accept_stages = stages; + let r = g.commit_after_script(tx, prep); + stages = g.last_accept_stages; + r + }; + lock_us = lock_us.saturating_add(t_lock.elapsed().as_micros() as u64); + match commit { + Ok(r) => { + prevouts.push(prev); + accepted.push(r); + } + Err(e) => { + let mut g = self.lock_write(); + for r in accepted.iter().rev() { + let _ = g.remove_txid(&r.txid); } + let us = t0.elapsed().as_micros() as u64; + self.meter_accept_stages(lock_us, stages); + return Err(self.finish_accept_err(us, e).unwrap_err()); } } - stages = g.last_accept_stages; - lock_us = lock_us.saturating_add(t_lock.elapsed().as_micros() as u64); - match err { - Some(e) => Err(e), - None => Ok(accepted), - } - }; + } let us = t0.elapsed().as_micros() as u64; self.meter_accept_stages(lock_us, stages); - match result { - Ok(res) => { - let per = us / (res.len().max(1) as u64); - for (i, (tx, r)) in txs.iter().zip(res.iter()).enumerate() { - self.meter_accept_wall(per, true); - self.note_fee_flow_admit(r.weight, r.fee_sat); - self.push_recent(tx, r); - for old in &r.replaced { - self.unindex_txid(old); - } - self.index_txid( - r.txid, - tx, - prevouts.get(i).map(Vec::as_slice).unwrap_or(&[]), - ); - let shs = self - .sh_index - .lock() - .unwrap() - .by_tx - .get(&r.txid) - .cloned() - .unwrap_or_default(); - self.publish_announce(r, shs); - self.promote_orphans_staged(r.txid, &utxo); - } - self.note_template_update(); - Ok(res) - } - Err(e) => { - let hard = !matches!( - e, - AcceptError::Duplicate(_) - | AcceptError::Orphaned { .. } - | AcceptError::Policy("mempool full") - ); - if hard { - self.meter_accept_wall(us, false); - } else { - self.meter_accept_us.fetch_add(us, Ordering::Relaxed); - } - Err(e) + let per = us / (accepted.len().max(1) as u64); + for (i, (tx, r)) in txs.iter().zip(accepted.iter()).enumerate() { + self.meter_accept_wall(per, true); + self.note_fee_flow_admit(r.weight, r.fee_sat); + self.push_recent(tx, r); + for old in &r.replaced { + self.unindex_txid(old); } + self.index_txid( + r.txid, + tx, + prevouts.get(i).map(Vec::as_slice).unwrap_or(&[]), + ); + let shs = self + .sh_index + .lock() + .unwrap() + .by_tx + .get(&r.txid) + .cloned() + .unwrap_or_default(); + self.publish_announce(r, shs); + self.promote_orphans_staged(r.txid, &utxo); } + self.note_template_update(); + Ok(accepted) } /// Remove confirmed txids (tip connect / archive confirm) and re-try orphans diff --git a/crates/rbitcoin-test/tests/cross_surface.rs b/crates/rbitcoin-test/tests/cross_surface.rs index fc00ea7e6..af3673e5a 100644 --- a/crates/rbitcoin-test/tests/cross_surface.rs +++ b/crates/rbitcoin-test/tests/cross_surface.rs @@ -152,7 +152,7 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { let td = TestDatadir::new().unwrap(); let params = ChainParams::regtest(); let genesis = bitcoin::blockdata::constants::genesis_block(bitcoin::Network::Regtest); - let (coinbase_txid, rpc_cb) = { + let (coinbase_txid, rpc_cb, pkg_cb) = { let q = Query::open_or_create_tiny(td.store_path()).unwrap(); accept_and_connect_block(&q, ¶ms, Height::GENESIS, &genesis, Milestone::NONE).unwrap(); let (_tip, _time, cbs) = pad_empty_from( @@ -162,10 +162,10 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { genesis.header.time, 1, 102, - 2, + 3, ); q.flush().unwrap(); - (cbs[0], cbs[1]) + (cbs[0], cbs[1], cbs[2]) }; let electrum_addr = ephemeral_addr(); @@ -195,6 +195,13 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { assert_eq!(height, "102"); let count = jsonrpc(rpc_addr, "getblockcount", json!([])).await; assert_eq!(count["result"], 102, "{count}"); + let tips = jsonrpc(rpc_addr, "getchaintips", json!([])).await; + assert_eq!(tips["result"][0]["height"], 102, "{tips}"); + assert_eq!(tips["result"][0]["status"], "active", "{tips}"); + let cb_hex = coinbase_txid.to_string(); + let utxo = jsonrpc(rpc_addr, "gettxout", json!([cb_hex.clone(), 0])).await; + assert_eq!(utxo["result"]["coinbase"], true, "{utxo}"); + assert_eq!(utxo["result"]["confirmations"], 102, "{utxo}"); let rpc_spk = ScriptBuf::from_bytes(vec![0x54]); let rpc_spend = acs_spend(rpc_cb, 50_0000_0000, 1_000, rpc_spk); @@ -213,6 +220,9 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { assert_eq!(tma["result"][0]["allowed"], true, "{tma}"); let sent = jsonrpc(rpc_addr, "sendrawtransaction", json!([rpc_hex])).await; assert_eq!(sent["result"], rpc_txid, "{sent}"); + let mem_utxo = jsonrpc(rpc_addr, "gettxout", json!([rpc_txid.clone(), 0])).await; + assert_eq!(mem_utxo["result"]["confirmations"], 0, "{mem_utxo}"); + assert_eq!(mem_utxo["result"]["coinbase"], false, "{mem_utxo}"); let mem = jsonrpc(rpc_addr, "getrawmempool", json!([])).await; assert!( mempool_has(&mem, &rpc_txid), @@ -253,6 +263,13 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { let (st, body) = http_post(esplora_addr, "/tx", &hex).await; assert_eq!(st, 200, "POST /tx: {body}"); assert_eq!(body, txid_hex); + let hidden = jsonrpc(rpc_addr, "gettxout", json!([cb_hex.clone(), 0])).await; + assert!( + hidden["result"].is_null(), + "default include_mempool hides mempool-spent coinbase: {hidden}" + ); + let shown = jsonrpc(rpc_addr, "gettxout", json!([cb_hex, 0, false])).await; + assert_eq!(shown["result"]["coinbase"], true, "{shown}"); let dup = jsonrpc(rpc_addr, "sendrawtransaction", json!([hex.clone()])).await; assert_eq!(dup["result"], txid_hex, "sendraw of live mempool tx: {dup}"); @@ -408,6 +425,87 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { "replaced tx must leave mempool: {mem}" ); + let pkg_parent = acs_spend( + pkg_cb, + 50_0000_0000, + 1_000, + ScriptBuf::from_bytes(vec![0x57]), + ); + let pkg_child = acs_spend( + pkg_parent.compute_txid(), + 50_0000_0000 - 1_000, + 1_000, + ScriptBuf::from_bytes(vec![0x58]), + ); + let pkg_parent_txid = pkg_parent.compute_txid().to_string(); + let pkg_child_txid = pkg_child.compute_txid().to_string(); + let pkg_body = json!([encode_tx(&pkg_parent), encode_tx(&pkg_child)]).to_string(); + let (st, body) = http_post(esplora_addr, "/txs/package", &pkg_body).await; + assert_eq!(st, 200, "POST /txs/package: {body}"); + let pkg_v: Value = serde_json::from_str(&body).unwrap_or_else(|e| { + panic!("POST /txs/package json: {e} body={body}"); + }); + assert_eq!( + pkg_v["txids"], + json!([pkg_parent_txid.clone(), pkg_child_txid.clone()]), + "{pkg_v}" + ); + + let submitted = jsonrpc( + rpc_addr, + "submitpackage", + json!([[encode_tx(&pkg_parent), encode_tx(&pkg_child)]]), + ) + .await; + assert_eq!(submitted["error"]["code"], -1, "{submitted}"); + assert_eq!( + submitted["error"]["message"], "mempool relay disabled (still in IBD or tip not ready)", + "{submitted}" + ); + + let mem = jsonrpc(rpc_addr, "getrawmempool", json!([])).await; + for tid in [ + &high_txid, + &txid_hex, + &child_txid, + &pkg_parent_txid, + &pkg_child_txid, + ] { + assert!(mempool_has(&mem, tid), "getrawmempool missing {tid}: {mem}"); + } + + let (st, body) = http_get(esplora_addr, "/mempool").await; + assert_eq!(st, 200, "GET /mempool: {body}"); + let mem_info: Value = serde_json::from_str(&body).unwrap(); + assert!( + mem_info["count"].as_u64().unwrap_or(0) >= 5, + "live mempool count: {mem_info}" + ); + let (st, body) = http_get(esplora_addr, "/mempool/txids").await; + assert_eq!(st, 200, "GET /mempool/txids: {body}"); + let txids_v: Value = serde_json::from_str(&body).unwrap(); + let txid_list = txids_v.as_array().expect("mempool/txids array"); + assert!( + txid_list + .iter() + .any(|v| v.as_str() == Some(pkg_parent_txid.as_str())), + "mempool/txids missing package parent: {body}" + ); + let (st, body) = http_get(esplora_addr, "/mempool/recent").await; + assert_eq!(st, 200, "GET /mempool/recent: {body}"); + let recent: Value = serde_json::from_str(&body).unwrap(); + let recent_rows = recent.as_array().expect("mempool/recent array"); + assert!( + recent_rows + .iter() + .any(|r| r["txid"] == pkg_child_txid || r["txid"] == pkg_parent_txid), + "mempool/recent missing package tx: {body}" + ); + let (st, body) = http_get(esplora_addr, "/fee-estimates").await; + assert_eq!(st, 200, "GET /fee-estimates: {body}"); + let fees: Value = serde_json::from_str(&body).unwrap(); + assert!(fees.get("1").is_some(), "{fees}"); + let mined = jsonrpc(rpc_addr, "generate", json!([1])).await; assert_eq!( mined["result"].as_array().map(|a| a.len()), @@ -422,13 +520,21 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { let blk = jsonrpc(rpc_addr, "getblock", json!([tip["result"].clone(), 2])).await; let txs = blk["result"]["tx"].as_array().expect("mined tx array"); assert!( - txs.len() >= 4, - "coinbase + RBF replacement + esplora parent + child: {blk}" - ); - assert!( - txs.iter().any(|t| t["txid"] == high_txid), - "generate must include RBF replacement: {blk}" + txs.len() >= 6, + "coinbase + RBF + esplora parent/child + package: {blk}" ); + for tid in [ + &high_txid, + &txid_hex, + &child_txid, + &pkg_parent_txid, + &pkg_child_txid, + ] { + assert!( + txs.iter().any(|t| t["txid"] == *tid), + "generate must include {tid}: {blk}" + ); + } let cb_txid = txs[0]["txid"].as_str().expect("coinbase txid").to_string(); let cb_val = (txs[0]["vout"][0]["value"].as_f64().unwrap() * 100_000_000.0).round() as u64; let immature = Transaction {