Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`
Expand Down Expand Up @@ -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.
Expand Down
2 changes: 1 addition & 1 deletion TESTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
121 changes: 54 additions & 67 deletions crates/rbitcoin-net/src/tx_relay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<AcceptResult> = Vec::with_capacity(txs.len());
let mut prevouts: Vec<Vec<TxOut>> = Vec::with_capacity(txs.len());
for tx in txs {
utxo.note_spender(tx);
let delta = self.fee_delta(&tx.compute_txid());
Expand All @@ -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<Vec<TxOut>> = 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<AcceptResult> = 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
Expand Down
124 changes: 115 additions & 9 deletions crates/rbitcoin-test/tests/cross_surface.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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, &params, Height::GENESIS, &genesis, Milestone::NONE).unwrap();
let (_tip, _time, cbs) = pad_empty_from(
Expand All @@ -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();
Expand Down Expand Up @@ -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);
Expand All @@ -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),
Expand Down Expand Up @@ -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}");

Expand Down Expand Up @@ -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()),
Expand All @@ -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 {
Expand Down