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
4 changes: 2 additions & 2 deletions crates/rbitcoin-consensus/src/index_writebehind.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -138,7 +138,7 @@ fn assemble_window(query: &Query, window: &Arc<IndexWindow>) -> Result<Assembled
}
slots.push(out);
}
if let Some(wave) = start_for_each_owned(jobs, assemble_index_height)? {
if let Some(wave) = start_for_each_owned_chunk(jobs, assemble_index_height, 1)? {
wave.finish()?;
}
let mut filters = Vec::new();
Expand Down
50 changes: 47 additions & 3 deletions crates/rbitcoin-consensus/src/script_pool.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,9 @@ unsafe impl Sync for Apply {}

struct Wave {
n: usize,
/// Jobs taken per successful steal. Script verify uses [`STEAL_CHUNK`].
/// An index window claims one height.
chunk: usize,
next: AtomicUsize,
in_wave: AtomicUsize,
failed: AtomicBool,
Expand Down Expand Up @@ -81,7 +84,8 @@ impl Wave {
self.notify_if_complete();
return None;
}
let i = self.next.fetch_add(STEAL_CHUNK, Ordering::Relaxed);
let chunk = self.chunk.max(1);
let i = self.next.fetch_add(chunk, Ordering::Relaxed);
if i >= self.n {
self.in_wave.fetch_sub(1, Ordering::AcqRel);
self.notify_if_complete();
Expand All @@ -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 {
Expand Down Expand Up @@ -263,10 +267,20 @@ fn unpublish_fg(wave: &Arc<Wave>) {
}

/// 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<T: Sync>(
items: Vec<T>,
f: fn(&T) -> Result<(), ConsensusError>,
) -> Result<Option<OwnedWave<T>>, 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<T: Sync>(
items: Vec<T>,
f: fn(&T) -> Result<(), ConsensusError>,
chunk: usize,
) -> Result<Option<OwnedWave<T>>, ConsensusError> {
if on_steal_worker() {
return Err(ConsensusError::BadBlock(
Expand All @@ -291,6 +305,7 @@ pub(crate) fn start_for_each_owned<T: Sync>(
}
let wave = Arc::new(Wave {
n: items.len(),
chunk,
next: AtomicUsize::new(0),
in_wave: AtomicUsize::new(0),
failed: AtomicBool::new(false),
Expand Down Expand Up @@ -815,6 +830,35 @@ mod tests {
);
}

fn run_owned_chunk(items: Vec<u32>, 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
Expand Down
2 changes: 1 addition & 1 deletion docs/concurrency.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |

Expand Down
Loading