From 0dfb75902099a8153d1836ab9c40c45f7ee5e898 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sun, 27 Sep 2026 09:11:53 -0700 Subject: [PATCH] index: claim one height per steal on the catch-up wave A late window is shorter than the 32-job script-verify chunk, so one worker was still running the whole window. Co-authored-by: Cursor --- .../src/index_writebehind.rs | 4 +- crates/rbitcoin-consensus/src/script_pool.rs | 50 +++++++++++++++++-- docs/concurrency.md | 2 +- 3 files changed, 50 insertions(+), 6 deletions(-) diff --git a/crates/rbitcoin-consensus/src/index_writebehind.rs b/crates/rbitcoin-consensus/src/index_writebehind.rs index 0c35bd0af..fc2367ee6 100644 --- a/crates/rbitcoin-consensus/src/index_writebehind.rs +++ b/crates/rbitcoin-consensus/src/index_writebehind.rs @@ -11,7 +11,7 @@ //! the watermarks. A wide gap is the materialize; at the tip each release is //! a one-height window. -use crate::script_pool::start_for_each_owned; +use crate::script_pool::start_for_each_owned_chunk; use crate::silent_payments::tweak_records_from_window; use crate::ConsensusError; use rbitcoin_primitives::{Fk, Height}; @@ -138,7 +138,7 @@ fn assemble_window(query: &Query, window: &Arc) -> Result= self.n { self.in_wave.fetch_sub(1, Ordering::AcqRel); self.notify_if_complete(); @@ -91,7 +95,7 @@ impl Wave { if STEAL_CLAIMS_ON.load(Ordering::Relaxed) { STEAL_CLAIMS.fetch_add(1, Ordering::Relaxed); } - Some(i..self.n.min(i.saturating_add(STEAL_CHUNK))) + Some(i..self.n.min(i.saturating_add(chunk))) } fn is_complete(&self) -> bool { @@ -263,10 +267,20 @@ fn unpublish_fg(wave: &Arc) { } /// Publish `items` for steal workers without waiting. `None` = already done -/// (empty or single-item ran inline). +/// (empty or single-item ran inline). Claims [`STEAL_CHUNK`] jobs at a time. pub(crate) fn start_for_each_owned( items: Vec, f: fn(&T) -> Result<(), ConsensusError>, +) -> Result>, ConsensusError> { + start_for_each_owned_chunk(items, f, STEAL_CHUNK) +} + +/// Same as [`start_for_each_owned`] with an explicit claim size. `chunk` of 0 +/// is treated as 1. Empty and single-item lists still run inline. +pub(crate) fn start_for_each_owned_chunk( + items: Vec, + f: fn(&T) -> Result<(), ConsensusError>, + chunk: usize, ) -> Result>, ConsensusError> { if on_steal_worker() { return Err(ConsensusError::BadBlock( @@ -291,6 +305,7 @@ pub(crate) fn start_for_each_owned( } let wave = Arc::new(Wave { n: items.len(), + chunk, next: AtomicUsize::new(0), in_wave: AtomicUsize::new(0), failed: AtomicBool::new(false), @@ -815,6 +830,35 @@ mod tests { ); } + fn run_owned_chunk(items: Vec, chunk: usize) -> Result<(), ConsensusError> { + match start_for_each_owned_chunk(items, ok_u32, chunk)? { + Some(w) => w.finish(), + None => Ok(()), + } + } + + /// Index windows are at most 64 heights and often under 32, so a claim of + /// 32 assigns the whole window to one worker. Eight jobs at size 1 are + /// eight claims; the script-verify default still covers those eight in one. + #[test] + fn index_wave_claims_one_job() { + let _gate = STEAL_TEST.lock().unwrap_or_else(|p| p.into_inner()); + workers(); + STEAL_CLAIMS.store(0, Ordering::Relaxed); + STEAL_CLAIMS_ON.store(true, Ordering::Relaxed); + run_owned_chunk((0..8).collect(), 1).unwrap(); + let one = STEAL_CLAIMS.load(Ordering::Relaxed); + STEAL_CLAIMS.store(0, Ordering::Relaxed); + run_owned((0..8).collect(), ok_u32).unwrap(); + let wide = STEAL_CLAIMS.load(Ordering::Relaxed); + STEAL_CLAIMS_ON.store(false, Ordering::Relaxed); + assert_eq!(one, 8, "claim size 1 is one claim per height"); + assert_eq!( + wide, 1, + "script-verify chunk of 32 covers 8 jobs in one claim" + ); + } + #[test] fn steal_index_does_not_lock_waves_per_job() { // Claim must not take WAVES: a 256-job wave is tens of thousands of diff --git a/docs/concurrency.md b/docs/concurrency.md index 7c019a0aa..e9159537e 100644 --- a/docs/concurrency.md +++ b/docs/concurrency.md @@ -52,7 +52,7 @@ Three thread kinds only. **Tokio workers must not wait on a `std` mutex/rwlock, | `peer_session` (split read/write) | Serve + reconstruct compact/body. Offers reconstructed blocks to **`tip-accept`** (does **not** take `connect_lock` or run confirm on the tokio worker). A reconstructed/received body whose header has valid PoW and **extends the current tip** is announced as `cmpctblock` to other HB peers **before** connect (Core `NewPoWValidBlock`). That is not a connected tip. P2P `tx` is `accept_tx_async` (blocking pool). INV / getdata / compact **first-pass** fill use mempool **`try_read` only** (busy write → skip that item; never park). A pending compact owns those clones; `blocktxn` apply does not `try_read` again. Handshake `FeeFilter` reads **atomic** `-minrelaytxfee` (no `inner`). 50 ms tick: age-INV is a due-log cursor (newly due only); `clock_due` / unbroadcast may full-walk ≤1/30 s. `any_tx_inv_due` is min live `accept_at`, not a map walk. `PeerHub::on_session_heartbeat` (headers-sync stall timeout). Inbound accept is `inbound_connect_and_handshake` (60 s VERSION/VERACK); timeout drops the `max_inbound` permit. | | `tip-accept` | **One** process-wide OS thread. Queue depth 8. Sole production thread for `accept_block` / `accept_branch` / `accept_received_block` / `generate_to_script` / `connect_lock` at tip. Confirm is **`confirm_wire_run_preverified`** (lookup stamp, load pin/assemble, `confirm_scripts_phase` → `rbtc-scripts-*`, write + `ibd-confirm-head` drain). TLS uring is this thread’s `with_thread_local` session. SIGINT stays on tokio; the current job finishes, then the session sees shutdown. Dropping the session’s join future **detaches** (does not condvar-wait on the worker). | | `rbtc-sh-wb` | **One** Class B scripthash appender. Used for tip follow **and** short catch-up when a durable SH head already exists. Confirm enqueues RAM records; `connect_at` / `note_confirmed_tip` **release** after `tip_tx`. This thread `put_create_batch_append` only for released heights, then advances `sh_indexed_through`. Apply errors re-queue and halt. Post-IBD Class A collect uses pack sessions (not a second appender); it runs while still Direct so this thread no-ops. | -| `rbtc-idx-wb` / `rbtc-idx-cpu` | **One** block index builder for BIP158 basic filters (`--block-filter-index`) and BIP-352 tweaks (`--sp-tweaks`), spawned when tip mode is entered; the confirm write thread does no index work. Waits on the same release gate as `rbtc-sh-wb` (`index_released_through`). The IO thread plans windows (≤64 heights / ≤50k creates) from the lower index watermark to the released tip and reads each with `read_index_window` on one completion session (locators, block spans, parent locators + P2TR-output `seqsigwit` + parent txids, parent records). One CPU worker, at most two windows behind, publishes one job per height to `rbtc-scripts-*` (filter GCS and tweak EC) and commits each index once per window under the index write-behind mutex, only if its watermark and every `confirmed[h]` are unchanged; `disconnect_tip_with` takes that mutex before truncating both. A failed commit (reorg) makes the IO thread re-plan from the watermarks. Apply errors halt the node. | +| `rbtc-idx-wb` / `rbtc-idx-cpu` | **One** block index builder for BIP158 basic filters (`--block-filter-index`) and BIP-352 tweaks (`--sp-tweaks`), spawned when tip mode is entered; the confirm write thread does no index work. Waits on the same release gate as `rbtc-sh-wb` (`index_released_through`). The IO thread plans windows (≤64 heights / ≤50k creates) from the lower index watermark to the released tip and reads each with `read_index_window` on one completion session (locators, block spans, parent locators + P2TR-output `seqsigwit` + parent txids, parent records). One CPU worker, at most two windows behind, publishes one job per height to `rbtc-scripts-*` (filter GCS and tweak EC). That wave claims one height, so a short window still spreads across the pool. The CPU worker commits each index once per window under the index write-behind mutex, only if its watermark and every `confirmed[h]` are unchanged; `disconnect_tip_with` takes that mutex before truncating both. A failed commit (reorg) makes the IO thread re-plan from the watermarks. Apply errors halt the node. | | Electrum / Esplora | Confirmed SH reads join durable index **plus a RAM SH head** (pending jobs keyed by scripthash) and pin that visible height (live tip while jobs sit, never above published tip). A tx is in mempool overlay **or** SH (pending/durable), not both and not neither. Reorg reaccepts then drops pending. Headers subscribe is live tip. | | Epoch finalize | Single-threaded control path; flushes table maps / fd durability |