From eca656bc8278b306f0288a3621df1ac5dda47c7c Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Fri, 28 Aug 2026 20:23:31 -0700 Subject: [PATCH 1/6] store: size tx.head rebuild workers at 1 GiB peak BDZ peel scratch + keys + g at the default 2^25 seal peaks around 1 GiB per worker, but the free-RAM gate still used 750 MiB. Two workers on a 1.5 GiB-free host would oversubscribe. Match the gate (and owner docs) to the actual peak. Co-authored-by: Cursor --- CHANGELOG.md | 2 +- SCHEMA.md | 2 +- crates/rbitcoin-store/src/sorted_run.rs | 8 ++++---- crates/rbitcoin-store/src/tx_table/mod.rs | 7 ++++--- crates/rbitcoin-store/src/tx_table/tests.rs | 11 ++++++----- docs/env-knobs.md | 2 +- docs/ibd-memory.md | 2 +- 7 files changed, 18 insertions(+), 16 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 12e3fb84..dc657103 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -98,7 +98,7 @@ before 1.0). unchanged). No `sp_tweaks` schema change. - **Parallel `tx.head` wipe-rebuild:** ranges seal concurrently, - min(CPUs, free RAM / **750 MiB**, range count). Distinct from SH + min(CPUs, free RAM / **1 GiB**, range count). Distinct from SH materialize (1.5 GiB/worker). `RBITCOIN_TX_HEAD_REBUILD_WORKERS` overrides (`1` = serial). diff --git a/SCHEMA.md b/SCHEMA.md index 7f8c2744..4401956e 100644 --- a/SCHEMA.md +++ b/SCHEMA.md @@ -481,7 +481,7 @@ seals it. Two unsealed non-tails is **Corrupt**. **Capacity @ 0.80 load (25-bit):** ≈ **26.8 M creates/segment**, ~29 MiB fuse8 when sealed (~6.1 B total sealed storage per create including head slots). **Wipe / empty-head rebuild:** writes MPHF+fuse8 **directly** from `txid.body` -(no historical OA). Ranges seal in parallel: min(CPUs, free RAM / 750 MiB, +(no historical OA). Ranges seal in parallel: min(CPUs, free RAM / 1 GiB, range count); `RBITCOIN_TX_HEAD_REBUILD_WORKERS` overrides (`1` = serial). Default range **2²⁵ keys** (`RBITCOIN_TX_HEAD_REBUILD_SEAL_BITS=25`); **26** is wider. Remainder is sealed; an empty open tail is created. Live IBD rolls OA diff --git a/crates/rbitcoin-store/src/sorted_run.rs b/crates/rbitcoin-store/src/sorted_run.rs index 9f61b909..94704868 100644 --- a/crates/rbitcoin-store/src/sorted_run.rs +++ b/crates/rbitcoin-store/src/sorted_run.rs @@ -914,16 +914,16 @@ mod tests { } #[test] - fn workers_for_free_ram_750mib_is_not_sh_1_5gib() { - const MIB: u64 = 1024 * 1024; - const HEAD: u64 = 750 * MIB; + fn workers_for_free_ram_1gib_head_is_not_sh_1_5gib() { + const HEAD: u64 = crate::tx_table::TX_HEAD_REBUILD_WORKER_FREE_RAM_BYTES; assert_eq!(workers_for_free_ram(8, 0, HEAD), 1); assert_eq!(workers_for_free_ram(8, HEAD.saturating_sub(1), HEAD), 1); assert_eq!(workers_for_free_ram(8, HEAD, HEAD), 1); assert_eq!(workers_for_free_ram(8, 2 * HEAD, HEAD), 2); assert_eq!(workers_for_free_ram(8, 3 * HEAD, HEAD), 3); assert_eq!(workers_for_free_ram(4, 20 * HEAD, HEAD), 4); - assert_eq!(workers_for_free_ram(8, 3 * (1 << 30) / 2, HEAD), 2); + assert_eq!(workers_for_free_ram(8, 3 * (1 << 30) / 2, HEAD), 1); + assert_eq!(workers_for_free_ram(8, 2 * (1 << 30), HEAD), 2); } #[test] diff --git a/crates/rbitcoin-store/src/tx_table/mod.rs b/crates/rbitcoin-store/src/tx_table/mod.rs index e1786848..7b4f6a0e 100644 --- a/crates/rbitcoin-store/src/tx_table/mod.rs +++ b/crates/rbitcoin-store/src/tx_table/mod.rs @@ -18,7 +18,8 @@ thread_local! { } /// Host RAM budget per parallel `tx.head` rebuild worker (not SH pack's 2 GiB). -pub const TX_HEAD_REBUILD_WORKER_FREE_RAM_BYTES: u64 = 750 * 1024 * 1024; +/// BDZ peel scratch + keys + g at the default 2²⁵ seal is ≈1 GiB peak. +pub const TX_HEAD_REBUILD_WORKER_FREE_RAM_BYTES: u64 = 1024 * 1024 * 1024; /// Max coalesced `txout.body` pread for SH Class A collect. Sequential libc /// pread (not TLS uring); 16 MiB matches Class A locality. @@ -1553,7 +1554,7 @@ impl TxTable { /// /// Range width is [`Self::rebuild_seal_keys`] (default 2²⁵). Remainder is /// sealed too; an empty open tail is created for later inserts. - /// Workers: [`Self::rebuild_workers`] (min of CPUs, free RAM / 750 MiB, + /// Workers: [`Self::rebuild_workers`] (min of CPUs, free RAM / 1 GiB, /// and range count). Distinct from SH pack's 2 GiB cap. pub fn rebuild_head_from_bodies( &self, @@ -1649,7 +1650,7 @@ impl TxTable { } /// Parallel wipe-rebuild workers. Env `RBITCOIN_TX_HEAD_REBUILD_WORKERS`; - /// default min(CPUs, free RAM / 750 MiB). Not SH pack's 2 GiB cap. + /// default min(CPUs, free RAM / 1 GiB). Not SH pack's 2 GiB cap. pub fn rebuild_workers() -> usize { #[cfg(test)] if let Some(n) = TEST_REBUILD_WORKERS.with(std::cell::Cell::get) { diff --git a/crates/rbitcoin-store/src/tx_table/tests.rs b/crates/rbitcoin-store/src/tx_table/tests.rs index d02c178b..3aa8f404 100644 --- a/crates/rbitcoin-store/src/tx_table/tests.rs +++ b/crates/rbitcoin-store/src/tx_table/tests.rs @@ -2846,18 +2846,19 @@ fn parse_rebuild_seal_bits_default_25() { } #[test] -fn parse_rebuild_workers_and_750mib_cap() { +fn parse_rebuild_workers_and_1gib_cap() { assert_eq!(parse_rebuild_workers(None), None); assert_eq!(parse_rebuild_workers(Some("foo")), None); assert_eq!(parse_rebuild_workers(Some("1")), Some(1)); assert_eq!(parse_rebuild_workers(Some("0")), Some(1)); assert_eq!(parse_rebuild_workers(Some("8")), Some(8)); assert_eq!(parse_rebuild_workers(Some("999")), Some(256)); - const MIB: u64 = 1024 * 1024; + const GIB: u64 = 1024 * 1024 * 1024; + assert_eq!(TX_HEAD_REBUILD_WORKER_FREE_RAM_BYTES, GIB); assert_eq!(tx_head_rebuild_workers_for_free_ram(8, 0), 1); - assert_eq!(tx_head_rebuild_workers_for_free_ram(8, 750 * MIB), 1); - assert_eq!(tx_head_rebuild_workers_for_free_ram(8, 1500 * MIB), 2); - assert_eq!(tx_head_rebuild_workers_for_free_ram(4, 20 * 750 * MIB), 4); + assert_eq!(tx_head_rebuild_workers_for_free_ram(8, GIB), 1); + assert_eq!(tx_head_rebuild_workers_for_free_ram(8, 2 * GIB), 2); + assert_eq!(tx_head_rebuild_workers_for_free_ram(4, 20 * GIB), 4); } #[test] diff --git a/docs/env-knobs.md b/docs/env-knobs.md index 4fda2d46..aa1cb2e0 100644 --- a/docs/env-knobs.md +++ b/docs/env-knobs.md @@ -27,7 +27,7 @@ for signet/mainnet sync. **Not** CLI. | `RBITCOIN_CLASS_C_INRAM_MAX_MB` | 256 | L2 cap for `confirmed` / `header_txs_*`; over → fd L0. `strong_tx` always L2 | | `RBITCOIN_TX_HEAD_BITS` | scale default | `tx.head` bits (dangerous on a live datadir) | | `RBITCOIN_TX_HEAD_REBUILD_SEAL_BITS` | 25 | Wipe/empty-head MPHF range `2^bits` (26 wider; clamp 6..=26) | -| `RBITCOIN_TX_HEAD_REBUILD_WORKERS` | min(n-cpu, free-RAM/750 MiB) | Wipe/empty-head MPHF parallelism (`1` = serial). Unset = auto. **Not** SH pack's 2 GiB cap | +| `RBITCOIN_TX_HEAD_REBUILD_WORKERS` | min(n-cpu, free-RAM/1 GiB) | Wipe/empty-head MPHF parallelism (`1` = serial). Unset = auto. **Not** SH pack's 2 GiB cap | | `RBITCOIN_TX_IDX_SOFT_SPAN` | 16 GiB | Per-stem idx soft rollover (do not set above 32 GiB hard span). Does **not** cut `tx.head`. | | `RBITCOIN_HEAD_SLOTS_HEADER` | scale default | Header hash-head initial slots (power of two) | | `RBITCOIN_SH_UNIQUE_HINT` | off | SH unique-hint probe | diff --git a/docs/ibd-memory.md b/docs/ibd-memory.md index 1ee7df59..b9eba5ba 100644 --- a/docs/ibd-memory.md +++ b/docs/ibd-memory.md @@ -39,7 +39,7 @@ pres and **not** the raw bytes. Reorg gather that wants wire re-encodes. | **Confirm plans / headers** | offer-ahead window | `ConfirmParentCache::advance_tip` from write `post_commit` | | **SH catalog runs** | leftover `scripthash.runs` discarded at tip (unsorted collect does not write them) | write-behind / discard; not during Direct confirm | | **SH unsorted collect / pack** | Collect: nCPU (no env / RAM cap; 1 MiB grow-on-demand write buffers; per-shard mutex so pwrite issues in offset order; 64 MiB fallocate on full flushes). Pack: min(CPUs, free RAM / 2 GiB); `RBITCOIN_SH_MERGE_WORKERS` override. One anonymous file image per pack worker (in-place unique-sort; MPHF has no HashSet). Class A collect spans 1 MiB | Tip finalize. Unsorted files under `scripthash.unsorted/` | -| **`tx.head` wipe-rebuild workers** | min(CPUs, host free RAM / 750 MiB, range count); floor 1. Same free-RAM probe as SH. **Not** the SH pack 2 GiB cap. Env `RBITCOIN_TX_HEAD_REBUILD_WORKERS` override (`1` = serial) | Empty/wipe `tx.head` rebuild from `txid.body`. Logs `workers=` `free_GiB=` | +| **`tx.head` wipe-rebuild workers** | min(CPUs, host free RAM / 1 GiB, range count); floor 1. Same free-RAM probe as SH. Matches BDZ peel+keys+g peak at default 2²⁵. **Not** the SH pack 2 GiB cap. Env `RBITCOIN_TX_HEAD_REBUILD_WORKERS` override (`1` = serial) | Empty/wipe `tx.head` rebuild from `txid.body`. Logs `workers=` `free_GiB=` | | **Ordered work path** | `MAX_ORDERED_HEADERS` | `IbdWorkState::hygiene` | Tests that need a clean process must call these **same** entry points (or drop the From 0e9bdccb44425c7694b00994a8cfe432374ab706 Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Fri, 28 Aug 2026 20:25:22 -0700 Subject: [PATCH 2/6] query: meter in-flight bytes from live creates/outs Same-txid overwrite saturating-added another 40-byte creates listing while the HashMap still held one slot, so iflight= could exceed live Arc bytes until the older pack pruned. Charge creates occupancy on vacant insert only, and subtract only when remove_keys actually drops the map entry (or the CreatePin). Co-authored-by: Cursor --- crates/rbitcoin-query/src/in_flight.rs | 50 +++++++++++++++++++------- 1 file changed, 37 insertions(+), 13 deletions(-) diff --git a/crates/rbitcoin-query/src/in_flight.rs b/crates/rbitcoin-query/src/in_flight.rs index de715394..28370f40 100644 --- a/crates/rbitcoin-query/src/in_flight.rs +++ b/crates/rbitcoin-query/src/in_flight.rs @@ -70,14 +70,13 @@ impl InFlight { for (fk, pin) in pins { self.note_create_fk(pin.0.txid, fk); keys.txids.push((pin.0.txid, fk)); - keys.approx_bytes = keys.approx_bytes.saturating_add(40); if let Some(id) = fk.get() { self.outs.insert(id, Arc::clone(pin)); keys.out_ids.push(id); - keys.approx_bytes = keys - .approx_bytes - .saturating_add(40) + let pin_bytes = 40u64 .saturating_add(crate::archive::create_pin_approx_bytes(pin) as u64); + keys.approx_bytes = keys.approx_bytes.saturating_add(pin_bytes); + self.approx_bytes = self.approx_bytes.saturating_add(pin_bytes); } } self.commit_keys(keys, height); @@ -93,7 +92,6 @@ impl InFlight { for (txid, fk) in pairs { self.note_create_fk(txid, fk); keys.txids.push((txid, fk)); - keys.approx_bytes = keys.approx_bytes.saturating_add(40); } self.commit_keys(keys, height); } @@ -102,7 +100,6 @@ impl InFlight { if keys.is_empty() { return; } - self.approx_bytes = self.approx_bytes.saturating_add(keys.approx_bytes); let slot = match height { Some(h) => self.by_height.entry(h).or_default(), None => &mut self.untagged, @@ -113,10 +110,14 @@ impl InFlight { } fn note_create_fk(&mut self, txid: [u8; 32], fk: Fk) { - if let Some(old) = self.creates.insert(txid, fk) { - if old != fk { + match self.creates.insert(txid, fk) { + None => { + self.approx_bytes = self.approx_bytes.saturating_add(40); + } + Some(old) if old != fk => { *self.evictions.entry(txid).or_insert(0) += 1; } + Some(_) => {} } } @@ -134,23 +135,28 @@ impl InFlight { fn remove_keys(&mut self, keys: HeightKeys) { if self.evictions.is_empty() { for (tid, _) in &keys.txids { - self.creates.remove(tid); + if self.creates.remove(tid).is_some() { + self.approx_bytes = self.approx_bytes.saturating_sub(40); + } } } else { for (tid, fk) in &keys.txids { if self.creates.get(tid) == Some(fk) { self.creates.remove(tid); + self.approx_bytes = self.approx_bytes.saturating_sub(40); } else { - // This pack's entry was clobbered (or the clobberer was - // dropped first); retire one eviction credit instead. self.consume_eviction(tid); } } } for id in &keys.out_ids { - self.outs.remove(id); + if let Some(pin) = self.outs.remove(id) { + self.approx_bytes = self + .approx_bytes + .saturating_sub(40) + .saturating_sub(crate::archive::create_pin_approx_bytes(&pin) as u64); + } } - self.approx_bytes = self.approx_bytes.saturating_sub(keys.approx_bytes); } /// Drop tagged packs at or above a disconnected height. Untagged stay. @@ -452,6 +458,24 @@ mod tests { assert_eq!(m.get_create_fk(&txid), Some(Fk(2))); } + #[test] + fn same_txid_overwrite_does_not_double_count_creates_bytes() { + let mut m = InFlight::new(); + let mut txid = [0u8; 32]; + txid[0] = 0xe3; + m.note_creates([(txid, Fk(1))], Some(91722)); + let (_, _, once) = m.size_snapshot(); + m.note_creates([(txid, Fk(2))], Some(91880)); + let (_, _, twice) = m.size_snapshot(); + assert_eq!( + twice, once, + "creates map still holds one slot; iflight= must not grow on clobber" + ); + m.prune_below_height(Some(91750)); + let (_, _, after) = m.size_snapshot(); + assert_eq!(after, once); + } + #[test] fn empty_note_is_noop() { let mut m = InFlight::new(); From acc05d226b2bc9529948bda8cb4ac289139481b3 Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Fri, 28 Aug 2026 20:31:19 -0700 Subject: [PATCH 3/6] net: evict pending blocks FIFO at the 128 cap MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit stash_pending_block dropped pending.keys().next() — HashMap iteration order — so a burst could evict the body on the current header path. Track insert order and pop the oldest hash, matching Q-60's pending FIFO. Re-stash of the same hash keeps its place. Co-authored-by: Cursor --- crates/rbitcoin-net/src/peer.rs | 73 +++++++++++++++++++++------ crates/rbitcoin-net/src/peer_tests.rs | 69 +++++++++++++------------ docs/ibd-memory.md | 2 +- 3 files changed, 96 insertions(+), 48 deletions(-) diff --git a/crates/rbitcoin-net/src/peer.rs b/crates/rbitcoin-net/src/peer.rs index a22181b0..c008ddce 100644 --- a/crates/rbitcoin-net/src/peer.rs +++ b/crates/rbitcoin-net/src/peer.rs @@ -22,7 +22,7 @@ use bitcoin::p2p::{Magic, ServiceFlags, PROTOCOL_VERSION}; use bitcoin::{Block, BlockHash, Transaction}; use rbitcoin_primitives::Height; use rbitcoin_query::Query; -use std::collections::{HashMap, HashSet}; +use std::collections::{HashMap, HashSet, VecDeque}; use std::net::SocketAddr; use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::Arc; @@ -123,17 +123,58 @@ pub(crate) fn insert_capped_txid( #[cfg(test)] pub(crate) const MAX_PENDING_BLOCKS_FOR_TEST: usize = MAX_PENDING_BLOCKS; -fn stash_pending_block( - pending: &mut HashMap, - hash: BlockHash, - block: bitcoin::Block, -) { - if pending.len() >= MAX_PENDING_BLOCKS && !pending.contains_key(&hash) { - if let Some(k) = pending.keys().next().copied() { - pending.remove(&k); +/// Tip-follow decoded bodies waiting for a connectable parent. Cap 128; +/// insert evicts the oldest hash (FIFO), not `HashMap::keys().next()`. +#[derive(Default)] +pub(crate) struct PendingBlocks { + map: HashMap, + fifo: VecDeque, +} + +impl PendingBlocks { + pub(crate) fn new() -> Self { + Self::default() + } + + pub(crate) fn len(&self) -> usize { + self.map.len() + } + + pub(crate) fn contains_key(&self, hash: &BlockHash) -> bool { + self.map.contains_key(hash) + } + + pub(crate) fn values(&self) -> std::collections::hash_map::Values<'_, BlockHash, bitcoin::Block> { + self.map.values() + } + + pub(crate) fn keys(&self) -> std::collections::hash_map::Keys<'_, BlockHash, bitcoin::Block> { + self.map.keys() + } + + pub(crate) fn insert(&mut self, hash: BlockHash, block: bitcoin::Block) { + stash_pending_block(self, hash, block); + } + + pub(crate) fn remove(&mut self, hash: &BlockHash) -> Option { + let b = self.map.remove(hash)?; + if let Some(i) = self.fifo.iter().position(|h| h == hash) { + self.fifo.remove(i); } + Some(b) + } +} + +fn stash_pending_block(pending: &mut PendingBlocks, hash: BlockHash, block: bitcoin::Block) { + if pending.map.len() >= MAX_PENDING_BLOCKS && !pending.map.contains_key(&hash) { + if let Some(k) = pending.fifo.pop_front() { + pending.map.remove(&k); + } + } + if !pending.map.contains_key(&hash) { + pending.fifo.push_back(hash); } - pending.insert(hash, block); + pending.map.insert(hash, block); } /// Services we advertise once store-backed reconstruct serve is available. @@ -645,7 +686,7 @@ pub async fn peer_session_with( // relay peer getdata CMPCT and broke tests that only serve `msg_block`. let mut peer_cmpct_version: u32 = 0; let mut pending_headers: HashMap = HashMap::new(); - let mut pending_blocks: HashMap = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct: HashMap = HashMap::new(); let mut from_this_peer: HashMap = HashMap::new(); let mut requested_blocks: HashSet = HashSet::new(); @@ -1374,7 +1415,7 @@ async fn handle_peer_frame( peer_send_cmpct: &mut bool, peer_cmpct_version: &mut u32, pending_headers: &mut HashMap, - pending_blocks: &mut HashMap, + pending_blocks: &mut PendingBlocks, pending_cmpct: &mut HashMap, from_this_peer: &mut HashMap, requested_blocks: &mut HashSet, @@ -2555,7 +2596,7 @@ fn fetchable_header_path_bodies( hub: &ChainHub, pending: &HashMap, tip: BlockHash, - pending_blocks: &HashMap, + pending_blocks: &PendingBlocks, requested: &HashSet, ) -> Vec { if !header_path_meets_minwork(hub, pending, tip) { @@ -2575,7 +2616,7 @@ fn missing_blocks_on_header_path( hub: &ChainHub, pending: &HashMap, tip: BlockHash, - pending_blocks: &HashMap, + pending_blocks: &PendingBlocks, requested: &HashSet, ) -> Vec { let mut path = Vec::new(); @@ -2687,7 +2728,7 @@ fn pending_header_leaves(pending: &HashMap) - fn drain_pending( hub: &ChainHub, out: &mpsc::UnboundedSender, - pending_blocks: &mut HashMap, + pending_blocks: &mut PendingBlocks, pending_headers: &mut HashMap, requested_blocks: &mut HashSet, compact: bool, @@ -2737,7 +2778,7 @@ fn drain_pending( /// download window, not a second most-work assembler. fn drain_pending_once( hub: &ChainHub, - pending_blocks: &mut HashMap, + pending_blocks: &mut PendingBlocks, pending_headers: &mut HashMap, ) -> Result<(), NetError> { let mut progress = true; diff --git a/crates/rbitcoin-net/src/peer_tests.rs b/crates/rbitcoin-net/src/peer_tests.rs index e8cf6f79..a17f3e7a 100644 --- a/crates/rbitcoin-net/src/peer_tests.rs +++ b/crates/rbitcoin-net/src/peer_tests.rs @@ -284,7 +284,7 @@ fn header_getdata_is_compact_after_sendcmpct() { hub.ensure_genesis().unwrap(); let (out_tx, mut out_rx) = mpsc::unbounded_channel(); let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -368,7 +368,7 @@ fn submitheader_parent_p2p_child_header_getdatas_body() { let (out_tx, mut out_rx) = mpsc::unbounded_channel(); let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -655,7 +655,7 @@ fn minchainwork_does_not_getdata_below_floor() { let (out_tx, mut out_rx) = mpsc::unbounded_channel(); let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -793,7 +793,7 @@ fn minchainwork_one_header_announces_ignore_height_14() { let (out_tx, mut out_rx) = mpsc::unbounded_channel(); let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -897,7 +897,7 @@ fn blocksonly_tx_and_inv_raise_ban() { let (out_tx, _out_rx) = mpsc::unbounded_channel(); let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -1123,7 +1123,7 @@ fn blocksonly_sendraw_invs_unbroadcast_to_inbound() { let mut send_cmpct = false; let mut cmpct_ver = 2u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut ban = 0u32; @@ -1553,7 +1553,7 @@ fn mocktime_jump_does_not_inv_or_serve_new_sendraw() { let mut send_cmpct = false; let mut cmpct_ver = 2u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut ban = 0u32; @@ -1679,7 +1679,7 @@ fn blocksonly_relay_perm_tx_invs_other_inbound() { let mut send_cmpct = false; let mut cmpct_ver = 2u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_first = HashMap::new(); let mut ban = 0u32; @@ -1753,7 +1753,7 @@ fn cmpct_helpers_without_mempool_and_queue_out_closed() { assert!(hdrs.is_empty() || !hdrs.is_empty()); // drain_pending empty is a no-op. - let mut pb = HashMap::new(); + let mut pb = PendingBlocks::new(); let mut ph = HashMap::new(); let (tx, _rx) = mpsc::unbounded_channel(); drain_pending(&hub, &tx, &mut pb, &mut ph, &mut HashSet::new(), false).unwrap(); @@ -1870,7 +1870,7 @@ fn compact_child_of_invalid_disconnects_cached_same_stays() { let mut send_cmpct = false; let mut cmpct_ver = 2u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -2006,7 +2006,7 @@ fn handle_peer_frame_control_and_inv_paths() { let mut send_cmpct = false; let mut cmpct_ver = 2u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut ban = 0u32; @@ -2607,7 +2607,7 @@ fn handle_peer_frame_mempool_tx_and_inv_paths() { let mut send_cmpct = false; let mut cmpct_ver = 2u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut ban = 0u32; @@ -2885,7 +2885,7 @@ fn getdata_tx_notfound_unless_announced_or_reorg() { let mut send_cmpct = false; let mut cmpct_ver = 2u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut ban = 0u32; @@ -3070,7 +3070,7 @@ fn invalid_getdata_type0_still_serves_tip_block() { let mut send_cmpct = false; let mut cmpct_ver = 0u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut ban = 0u32; @@ -3444,7 +3444,7 @@ fn drain_requests_missing_parent_of_pending_branch() { } } let (tx, mut rx) = mpsc::unbounded_channel(); - let mut pb = HashMap::new(); + let mut pb = PendingBlocks::new(); pb.insert(orphan.block_hash(), orphan); let mut ph = HashMap::new(); drain_pending(&hub, &tx, &mut pb, &mut ph, &mut HashSet::new(), false).unwrap(); @@ -3526,7 +3526,7 @@ fn drain_connects_pending_child_of_new_tip_after_reorg() { let b1 = mine(gen, 1_300_001_000, 1); let b2 = mine(b1.block_hash(), 1_300_001_600, 2); - let mut pb = HashMap::new(); + let mut pb = PendingBlocks::new(); pb.insert(b1.block_hash(), b1.clone()); pb.insert(b2.block_hash(), b2.clone()); let mut ph = HashMap::new(); @@ -3595,7 +3595,7 @@ fn inv_of_already_asked_block_does_not_getdata() { let (out_tx, mut out_rx) = mpsc::unbounded_channel(); let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -3685,7 +3685,7 @@ fn bloom_disabled_messages_request_disconnect() { hub.ensure_genesis().unwrap(); let (out_tx, _out_rx) = mpsc::unbounded_channel(); let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -3763,7 +3763,7 @@ fn oversize_locator_request_disconnect() { hub.ensure_genesis().unwrap(); let (out_tx, _out_rx) = mpsc::unbounded_channel(); let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -3931,7 +3931,7 @@ fn redundant_verack_is_ignored_and_logged() { let mut send_cmpct = false; let mut cmpct_ver = 0u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut ban = 0u32; @@ -4081,7 +4081,7 @@ fn addrfetch_multi_addr_disconnects() { let mut send_cmpct = false; let mut cmpct_ver = 0u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut ban = 0u32; @@ -4224,8 +4224,9 @@ fn addrfetch_times_out_after_300s() { #[test] fn pending_blocks_insert_evicts_at_cap() { - let mut pending = HashMap::new(); + let mut pending = PendingBlocks::new(); let bits = bitcoin::CompactTarget::from_consensus(0x207f_ffff); + let mut hashes = Vec::new(); for i in 0u32..(MAX_PENDING_BLOCKS_FOR_TEST as u32 + 1) { let mut b = bitcoin::blockdata::constants::genesis_block(bitcoin::Network::Regtest); b.header.merkle_root = bitcoin::TxMerkleNode::from_byte_array({ @@ -4238,9 +4239,15 @@ fn pending_blocks_insert_evicts_at_cap() { }); b.header.bits = bits; let h = b.block_hash(); + hashes.push(h); stash_pending_block(&mut pending, h, b); } assert_eq!(pending.len(), MAX_PENDING_BLOCKS_FOR_TEST); + assert!( + !pending.contains_key(&hashes[0]), + "cap eviction must drop the oldest insert, not HashMap::keys().next()" + ); + assert!(pending.contains_key(&hashes[MAX_PENDING_BLOCKS_FOR_TEST])); } #[test] @@ -4318,7 +4325,7 @@ fn shorter_higher_work_fork_is_not_hopeless() { !announced_tip_is_hopeless(hub.tip_height().unwrap(), announced_h, work_cmp), "shorter higher-work path must not be hopeless" ); - let want = fetchable_header_path_bodies(&hub, &pending, tip, &HashMap::new(), &HashSet::new()); + let want = fetchable_header_path_bodies(&hub, &pending, tip, &PendingBlocks::new(), &HashSet::new()); assert!( !want.is_empty(), "must not skip bodies on a shorter higher-work path" @@ -4413,7 +4420,7 @@ fn connecting_ancient_weaker_headers_request_disconnect() { let mut send_cmpct = false; let mut cmpct_ver = 2u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -4540,7 +4547,7 @@ fn getdata_skips_reconstruct_when_serve_inflight_at_cap() { let mut send_cmpct = false; let mut cmpct_ver = 2u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -4643,7 +4650,7 @@ fn catchup_headers_getdata_stays_in_serve_window() { hub.ensure_genesis().unwrap(); let (out_tx, mut out_rx) = mpsc::unbounded_channel(); let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -4788,7 +4795,7 @@ fn catchup_compact_getdata_clears_requested_for_next_window() { assert!(hub.attach_mempool(mp).is_ok()); let (out_tx, mut out_rx) = mpsc::unbounded_channel(); let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -4961,7 +4968,7 @@ fn full_headers_batch_continues_from_last_header() { let mut send_cmpct = false; let mut cmpct_ver = 2u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -5238,7 +5245,7 @@ fn compact_tip_announce_must_not_wrap_serve_inflight() { let mut send_cmpct = true; let mut cmpct_ver = 2u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -5344,7 +5351,7 @@ fn compact_tip_announce_must_not_consume_serve_slots() { let mut send_cmpct = true; let mut cmpct_ver = 2u32; let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); @@ -5423,7 +5430,7 @@ fn coinbase_compact_fills_without_mempool() { assert!(hub.mempool().is_none()); let (out_tx, mut out_rx) = mpsc::unbounded_channel(); let mut pending_headers = HashMap::new(); - let mut pending_blocks = HashMap::new(); + let mut pending_blocks = PendingBlocks::new(); let mut pending_cmpct = HashMap::new(); let mut from_peer = HashMap::new(); let mut requested = HashSet::new(); diff --git a/docs/ibd-memory.md b/docs/ibd-memory.md index b9eba5ba..088307d2 100644 --- a/docs/ibd-memory.md +++ b/docs/ibd-memory.md @@ -57,7 +57,7 @@ Not page cache. Caps on **decoded `Block` objects and live outbound sessions**: | **GetData serve inflight** | **16** full `Block`/`CmpctBlock` per session writer | Writer saturating-decrements after send so unpaired compact tip announce cannot wrap to `usize::MAX`. Announce is not counted on this cap (a burst would starve reconstruct). Extra inv hashes not reconstructed. | | **Tip-follow catch-up getdata** | **16** (`MAX_SERVE_BLOCKS`) hashes per ask | `requested` tracks inflight; after those bodies connect, drain asks the next window. Asking the whole header path left hashes stuck while the peer served only 16. | | **`from_this_peer`** | **50_000** txids / session (INV origin skip) | Clear and restart on overflow (`announced_wtx` pattern). | -| **`pending_blocks`** | **128** decoded bodies / session | Insert evicts one existing hash at cap. Unsolicited BIP130 window still 16. | +| **`pending_blocks`** | **128** decoded bodies / session | Insert evicts the **oldest** hash (FIFO). Unsolicited BIP130 window still 16. | | **Hub `BlockCache` bodies** | **16** decoded (`DEFAULT_BODY_DEPTH`) | Hashes kept for locators. Compact/`getblocktxn` serve is depth 5/10; IBD does not `push_best`. Reconstruct from store past the window. | | **Hub `held_bodies`** | **320** count + **288** height window | Side-branch hold for most-work apply. After a successful tip connect, drop bodies whose connected height is more than **288** below tip (Core unrequested window). Do not hold `IgnoredWeaker` already that far behind. Count cap still evicts an arbitrary hash at 320. | | **Query `sh_heads`** | **65_536** process-local SH body heads | Evict arbitrary key at cap (`keys().next()`). Miss path `locate_head`s. Catch-up and tip SH apply share this map. | From 694cfe54997404a94038a7825869ff761ffc4298 Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Fri, 28 Aug 2026 20:44:55 -0700 Subject: [PATCH 4/6] consensus: decrement tapscript weight before empty-pubkey fail Core EvalChecksigTapscript subtracts VALIDATION_WEIGHT_PER_SIGOP for a non-empty sig, then rejects an empty pubkey. We rejected empty pubkey first, so a CHECKSIG that was also over the weight budget reported \"empty pubkey\" instead of \"validation weight\". Same reject either way; match Core's error class. Co-authored-by: Cursor --- .../src/script/interpreter.rs | 20 ++++++++++++++++--- 1 file changed, 17 insertions(+), 3 deletions(-) diff --git a/crates/rbitcoin-consensus/src/script/interpreter.rs b/crates/rbitcoin-consensus/src/script/interpreter.rs index dc7c04a9..d824cf9d 100644 --- a/crates/rbitcoin-consensus/src/script/interpreter.rs +++ b/crates/rbitcoin-consensus/src/script/interpreter.rs @@ -1000,9 +1000,6 @@ fn tapscript_sig_result( pubkey: &[u8], ctx: &EvalContext<'_>, ) -> Result { - if pubkey.is_empty() { - return Err(ConsensusError::Script("tapscript empty pubkey".into())); - } if !sig.is_empty() { let left = ctx.validation_weight_left.get() - TAPSCRIPT_VALIDATION_WEIGHT_PER_SIGOP; ctx.validation_weight_left.set(left); @@ -1010,6 +1007,9 @@ fn tapscript_sig_result( return Err(ConsensusError::Script("tapscript validation weight".into())); } } + if pubkey.is_empty() { + return Err(ConsensusError::Script("tapscript empty pubkey".into())); + } // Unknown public key type (not 32 bytes): treat signature as valid (soft-fork hook). if pubkey.len() != 32 { if sig.is_empty() { @@ -1794,6 +1794,20 @@ mod success_and_disabled_tests { assert!(format!("{err}").contains("empty pubkey")); } + #[test] + fn tapscript_empty_pubkey_reports_weight_when_budget_exhausted() { + // OP_1 OP_1 CHECKSIG then OP_1 OP_0 CHECKSIG. + // Dummy witness init weight is 51; first non-empty sig burns 50. + // Core decrements weight before the empty-pubkey fail, so the second + // CHECKSIG is TAPSCRIPT_VALIDATION_WEIGHT, not empty pubkey. + let script = vec![0x51, 0x51, 0xac, 0x51, 0x00, 0xac]; + let err = eval(&script, SigVersion::TapScript).unwrap_err(); + assert!( + format!("{err}").contains("validation weight"), + "Core order: weight before empty pubkey, got {err}" + ); + } + #[test] fn op_1sub_and_unary_arith() { // OP_3 OP_1SUB → 2; OP_2 EQUAL From 735481891cb9807c88b67290d4298d406ae12377 Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Fri, 28 Aug 2026 21:40:46 -0700 Subject: [PATCH 5/6] query: batch Electrum cut-through spentness per wave Cake historicalMode=false walked spent.idx + spent.body once per eligible tx (up to 16384 serial reads). One spent_range_batch plus one spent-body walk per create, matching serial unspent_create_vouts. Co-authored-by: Cursor --- CHANGELOG.md | 4 +- crates/rbitcoin-net/src/peer.rs | 4 +- crates/rbitcoin-net/src/peer_tests.rs | 3 +- crates/rbitcoin-query/src/in_flight.rs | 4 +- crates/rbitcoin-query/src/lib.rs | 8 +++ crates/rbitcoin-query/src/sp_tweaks.rs | 35 +++++++++---- crates/rbitcoin-store/src/store.rs | 56 +++++++++++++++++++++ crates/rbitcoin-store/src/tx_table/mod.rs | 21 +++++--- crates/rbitcoin-store/src/tx_table/tests.rs | 4 +- 9 files changed, 118 insertions(+), 21 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index dc657103..5f4f2d96 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -14,7 +14,9 @@ before 1.0). - **Electrum `blockchain.tweaks.subscribe`:** pre-taproot empty heights go out as **one notify** with ≤1024 keys (Cake last-key progress), not one line per height. Cake `historicalMode=false` (param `[2]`) **cut-through**: omit - confirmed-spent P2TR outs (and txs with none left). `true` keeps spent outs + confirmed-spent P2TR outs (and txs with none left). Spentness is one + `spent.idx` batch plus one spent-body walk per create, not a serial + idx+body per eligible tx. `true` keeps spent outs for restore. Probe `[0,1,false]` is still `{"0": {}}`. `{"message":"done"}` ends a **chunk** (60s wall at a wave boundary, or the requested `count` if sooner) so Cake resubscribes; it is not “`count` through tip”. diff --git a/crates/rbitcoin-net/src/peer.rs b/crates/rbitcoin-net/src/peer.rs index c008ddce..0f2b9da9 100644 --- a/crates/rbitcoin-net/src/peer.rs +++ b/crates/rbitcoin-net/src/peer.rs @@ -144,7 +144,9 @@ impl PendingBlocks { self.map.contains_key(hash) } - pub(crate) fn values(&self) -> std::collections::hash_map::Values<'_, BlockHash, bitcoin::Block> { + pub(crate) fn values( + &self, + ) -> std::collections::hash_map::Values<'_, BlockHash, bitcoin::Block> { self.map.values() } diff --git a/crates/rbitcoin-net/src/peer_tests.rs b/crates/rbitcoin-net/src/peer_tests.rs index a17f3e7a..4a12304d 100644 --- a/crates/rbitcoin-net/src/peer_tests.rs +++ b/crates/rbitcoin-net/src/peer_tests.rs @@ -4325,7 +4325,8 @@ fn shorter_higher_work_fork_is_not_hopeless() { !announced_tip_is_hopeless(hub.tip_height().unwrap(), announced_h, work_cmp), "shorter higher-work path must not be hopeless" ); - let want = fetchable_header_path_bodies(&hub, &pending, tip, &PendingBlocks::new(), &HashSet::new()); + let want = + fetchable_header_path_bodies(&hub, &pending, tip, &PendingBlocks::new(), &HashSet::new()); assert!( !want.is_empty(), "must not skip bodies on a shorter higher-work path" diff --git a/crates/rbitcoin-query/src/in_flight.rs b/crates/rbitcoin-query/src/in_flight.rs index 28370f40..3b142354 100644 --- a/crates/rbitcoin-query/src/in_flight.rs +++ b/crates/rbitcoin-query/src/in_flight.rs @@ -73,8 +73,8 @@ impl InFlight { if let Some(id) = fk.get() { self.outs.insert(id, Arc::clone(pin)); keys.out_ids.push(id); - let pin_bytes = 40u64 - .saturating_add(crate::archive::create_pin_approx_bytes(pin) as u64); + let pin_bytes = + 40u64.saturating_add(crate::archive::create_pin_approx_bytes(pin) as u64); keys.approx_bytes = keys.approx_bytes.saturating_add(pin_bytes); self.approx_bytes = self.approx_bytes.saturating_add(pin_bytes); } diff --git a/crates/rbitcoin-query/src/lib.rs b/crates/rbitcoin-query/src/lib.rs index efe45b13..5269ab17 100644 --- a/crates/rbitcoin-query/src/lib.rs +++ b/crates/rbitcoin-query/src/lib.rs @@ -1473,6 +1473,14 @@ impl Query { Ok(self.store.unspent_create_vouts(create_fk, vouts, None)?) } + /// Batch [`Self::unspent_create_vouts`]: one `spent.idx` walk across creates. + pub fn unspent_create_vouts_batch( + &self, + items: &[(Fk, Vec)], + ) -> Result>, QueryError> { + Ok(self.store.unspent_create_vouts_batch(items)?) + } + /// Enable/disable txid hash-head inserts on archive (default on). Off under /// milestone IBD; Class A bodies remain complete via header_txs fk lists. pub fn set_tx_index(&self, enabled: bool) { diff --git a/crates/rbitcoin-query/src/sp_tweaks.rs b/crates/rbitcoin-query/src/sp_tweaks.rs index d221a7be..1eaf35bd 100644 --- a/crates/rbitcoin-query/src/sp_tweaks.rs +++ b/crates/rbitcoin-query/src/sp_tweaks.rs @@ -238,7 +238,8 @@ impl Query { /// span** from first..=last eligible fk in the wave (ineligible txout in /// the hole is included; `inwit` is not). `sp_tweaks` mutex is not held /// during Class A IO. `limits.cut_through` drops confirmed-spent P2TR - /// outs after the join (txs with none left are omitted; the height remains). + /// outs after the join (one `spent.idx` batch, then one spent-body walk + /// per create; txs with none left are omitted; the height remains). pub fn load_thin_tweaks_range( &self, start: Height, @@ -362,14 +363,31 @@ impl Query { } if limits.cut_through && !elig_fks.is_empty() { - let mut k = 0usize; + let n_rows: usize = out_rows.iter().map(Vec::len).sum(); + if n_rows != elig_fks.len() { + return Err(StoreError::Corrupt( + "invariant: thin cut_through row/fk count", + )); + } + let items: Vec<(Fk, Vec)> = out_rows + .iter() + .flatten() + .zip(elig_fks.iter().copied()) + .map(|(row, fk)| (fk, row.p2tr.iter().map(|p| p.0).collect())) + .collect(); + let live_rows = self.unspent_create_vouts_batch(&items)?; + if live_rows.len() != items.len() { + return Err(StoreError::Corrupt( + "invariant: unspent_create_vouts_batch length", + )); + } + let mut live_iter = live_rows.into_iter(); for rows in &mut out_rows { - let n = rows.len(); - let fks = &elig_fks[k..k + n]; - let mut kept = Vec::with_capacity(n); - for (mut row, fk) in rows.drain(..).zip(fks.iter()) { - let vouts: Vec = row.p2tr.iter().map(|p| p.0).collect(); - let live = self.store.unspent_create_vouts(*fk, &vouts, None)?; + let mut kept = Vec::with_capacity(rows.len()); + for mut row in rows.drain(..) { + let live = live_iter + .next() + .ok_or(StoreError::Corrupt("invariant: thin cut_through live rows"))?; if live.len() != row.p2tr.len() { row.p2tr.retain(|(v, _, _)| live.iter().any(|u| u == v)); } @@ -378,7 +396,6 @@ impl Query { } } *rows = kept; - k = k.saturating_add(n); } } diff --git a/crates/rbitcoin-store/src/store.rs b/crates/rbitcoin-store/src/store.rs index 90ee6725..40215053 100644 --- a/crates/rbitcoin-store/src/store.rs +++ b/crates/rbitcoin-store/src/store.rs @@ -1022,6 +1022,27 @@ impl Store { Ok(unspent) } + /// Batch [`Self::unspent_create_vouts`]: one `spent.idx` walk, then one + /// spent-body read per create that has a range. + pub fn unspent_create_vouts_batch( + &self, + items: &[(Fk, Vec)], + ) -> Result>, StoreError> { + if items.is_empty() { + return Ok(Vec::new()); + } + let fks: Vec = items.iter().map(|(fk, _)| *fk).collect(); + let ranges = self.txs.spent_range_batch(&fks)?; + if ranges.len() != items.len() { + return Err(StoreError::Corrupt("invariant: spent_range_batch length")); + } + let mut out = Vec::with_capacity(items.len()); + for ((fk, vouts), range) in items.iter().zip(ranges) { + out.push(self.unspent_create_vouts(*fk, vouts, range)?); + } + Ok(out) + } + /// Multi-list node count only (sole spends do not allocate body rows). pub fn spender_list_count(&self) -> u64 { self.spenders.count() @@ -2982,6 +3003,11 @@ mod tests { assert_eq!(u, vec![0]); // empty vouts assert!(s.unspent_create_vouts(fk, &[], None).unwrap().is_empty()); + let batch = s + .unspent_create_vouts_batch(&[(fk, vec![0u32]), (fk, vec![])]) + .unwrap(); + assert_eq!(batch[0], vec![0]); + assert!(batch[1].is_empty()); // has_confirmed without range, no spender assert!(!s.has_confirmed_strong_spender_create(fk, 0, None).unwrap()); assert!(!s.has_confirmed_strong_spender(&[20u8; 32], 0).unwrap()); @@ -2993,6 +3019,36 @@ mod tests { let _ = std::fs::remove_dir_all(&dir); } + #[test] + fn unspent_create_vouts_batch_matches_serial() { + let dir = tmp(); + let s = Store::create(&dir).unwrap(); + let mut fks = Vec::new(); + for i in 1u8..=4 { + let item = coinbase_item( + [i; 32], + vec![ + OutputRecord::unspent(10, vec![0x51]), + OutputRecord::unspent(11, vec![0x51]), + ], + ); + fks.push(s.put_tx_full_batch_indexed(&[item], true).unwrap()[0]); + } + s.flush().unwrap(); + let items: Vec<(Fk, Vec)> = fks.iter().map(|fk| (*fk, vec![0, 1])).collect(); + let batch = s.unspent_create_vouts_batch(&items).unwrap(); + assert_eq!(batch.len(), 4); + for (i, fk) in fks.iter().enumerate() { + assert_eq!( + batch[i], + s.unspent_create_vouts(*fk, &[0, 1], None).unwrap(), + "fk {fk:?}" + ); + } + assert!(s.unspent_create_vouts_batch(&[]).unwrap().is_empty()); + let _ = std::fs::remove_dir_all(&dir); + } + #[test] fn resolve_txid_prefers_connected_over_newer_unconnected() { let dir = tmp(); diff --git a/crates/rbitcoin-store/src/tx_table/mod.rs b/crates/rbitcoin-store/src/tx_table/mod.rs index 7b4f6a0e..0877f335 100644 --- a/crates/rbitcoin-store/src/tx_table/mod.rs +++ b/crates/rbitcoin-store/src/tx_table/mod.rs @@ -1162,13 +1162,22 @@ impl TxTable { if vouts.is_empty() { return Ok(Vec::new()); } - let mut out = Vec::with_capacity(vouts.len()); - for &v in vouts { - if let Ok((multi, field)) = self.get_output_spender_meta_at(body_off, body_len, v) { - out.push((v, multi, field)); + let slot = OutputRecord::SPENT_SLOT_LEN; + self.spent.with_bytes_at(body_off, body_len, |raw| { + let mut out = Vec::with_capacity(vouts.len()); + for &v in vouts { + let start = (v as usize).saturating_mul(slot); + let end = start.saturating_add(slot); + if end > raw.len() { + continue; + } + let Ok((flags, field)) = decode_spent_slot_v17(&raw[start..end]) else { + continue; + }; + out.push((v, flags & output_flags::MULTI_SPENDER != 0, field)); } - } - Ok(out) + Ok(out) + }) } /// Patch multi + spender_field on create tx output (packed Class A body). diff --git a/crates/rbitcoin-store/src/tx_table/tests.rs b/crates/rbitcoin-store/src/tx_table/tests.rs index 3aa8f404..c86cb748 100644 --- a/crates/rbitcoin-store/src/tx_table/tests.rs +++ b/crates/rbitcoin-store/src/tx_table/tests.rs @@ -1362,7 +1362,9 @@ fn get_output_spender_metas_at_one_walk() { let s1 = Fk(10); t.put_spends_on_create_at(&spenders, off, len, &[(0, s1), (2, Fk(20))]) .unwrap(); - let metas = t.get_output_spender_metas_at(off, len, &[0, 1, 2]).unwrap(); + let metas = t + .get_output_spender_metas_at(off, len, &[0, 1, 2, 99]) + .unwrap(); assert_eq!(metas.len(), 3); assert!(!metas[0].1 && metas[0].2 == s1); assert!(!metas[1].1 && metas[1].2.is_null()); From abe350a0f7d0d8220887f4c9bee8719b64990466 Mon Sep 17 00:00:00 2001 From: rbitcoin-grok Date: Fri, 28 Aug 2026 21:48:47 -0700 Subject: [PATCH 6/6] net: insert pending blocks through the FIFO API PendingBlocks::insert and len were only reached from tests, so -D warnings failed the non-test lib build. Tip-follow now inserts through the wrapper; the cap pin uses keys().len(). Co-authored-by: Cursor --- crates/rbitcoin-net/src/peer.rs | 6 +----- crates/rbitcoin-net/src/peer_tests.rs | 4 ++-- 2 files changed, 3 insertions(+), 7 deletions(-) diff --git a/crates/rbitcoin-net/src/peer.rs b/crates/rbitcoin-net/src/peer.rs index 0f2b9da9..06f6ba70 100644 --- a/crates/rbitcoin-net/src/peer.rs +++ b/crates/rbitcoin-net/src/peer.rs @@ -136,10 +136,6 @@ impl PendingBlocks { Self::default() } - pub(crate) fn len(&self) -> usize { - self.map.len() - } - pub(crate) fn contains_key(&self, hash: &BlockHash) -> bool { self.map.contains_key(hash) } @@ -1926,7 +1922,7 @@ async fn handle_peer_frame( requested_blocks.remove(&hash); pending_headers.entry(hash).or_insert(block.header); if !any_header_path_meets_minwork(hub, pending_headers, hash) { - stash_pending_block(pending_blocks, hash, block.clone()); + pending_blocks.insert(hash, block.clone()); return Ok(()); } match hub.accept_received_block(block.clone()) { diff --git a/crates/rbitcoin-net/src/peer_tests.rs b/crates/rbitcoin-net/src/peer_tests.rs index 4a12304d..3feb215e 100644 --- a/crates/rbitcoin-net/src/peer_tests.rs +++ b/crates/rbitcoin-net/src/peer_tests.rs @@ -4240,9 +4240,9 @@ fn pending_blocks_insert_evicts_at_cap() { b.header.bits = bits; let h = b.block_hash(); hashes.push(h); - stash_pending_block(&mut pending, h, b); + pending.insert(h, b); } - assert_eq!(pending.len(), MAX_PENDING_BLOCKS_FOR_TEST); + assert_eq!(pending.keys().len(), MAX_PENDING_BLOCKS_FOR_TEST); assert!( !pending.contains_key(&hashes[0]), "cap eviction must drop the oldest insert, not HashMap::keys().next()"