From cea035798c35b67124329b02f9dbdd92c0a5d234 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 26 Sep 2026 17:09:18 -0700 Subject: [PATCH 1/2] index: report read, build, and commit time on the catch-up line The 10s progress line only printed a cumulative rate, so a slow tail could not be told apart from parent reads, tweak math, or fsync. Co-authored-by: Cursor --- .../src/index_writebehind.rs | 278 +++++++++++++++--- 1 file changed, 241 insertions(+), 37 deletions(-) diff --git a/crates/rbitcoin-consensus/src/index_writebehind.rs b/crates/rbitcoin-consensus/src/index_writebehind.rs index 96c44b39d..ea6244691 100644 --- a/crates/rbitcoin-consensus/src/index_writebehind.rs +++ b/crates/rbitcoin-consensus/src/index_writebehind.rs @@ -16,7 +16,7 @@ use crate::ConsensusError; use rbitcoin_primitives::{Fk, Height}; use rbitcoin_query::Query; use rbitcoin_store::IndexWindow; -use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::mpsc::{sync_channel, Receiver}; use std::sync::Arc; use std::time::{Duration, Instant}; @@ -66,15 +66,104 @@ fn read_window(query: &Query, start: u32, end: u32) -> Result Result { - let t0 = Instant::now(); +/// Milliseconds of one interval, or of one tip window. +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +struct IndexStageSample { + read_ms: u64, + build_ms: u64, + commit_ms: u64, +} + +/// Interval sums for `index: build` / `index: apply`. Not an `ibd: perf` token. +#[derive(Debug, Default)] +struct IndexStageMs { + read_ns: AtomicU64, + build_ns: AtomicU64, + commit_ns: AtomicU64, +} + +impl IndexStageMs { + fn add_ns(slot: &AtomicU64, ns: u64) { + if ns > 0 { + slot.fetch_add(ns, Ordering::Relaxed); + } + } + + fn add_read(&self, ns: u64) { + Self::add_ns(&self.read_ns, ns); + } + + fn add_build(&self, ns: u64) { + Self::add_ns(&self.build_ns, ns); + } + + fn add_commit(&self, ns: u64) { + Self::add_ns(&self.commit_ns, ns); + } + + /// Take the interval and reset. Sub-millisecond leftovers truncate. + fn take_ms(&self) -> IndexStageSample { + IndexStageSample { + read_ms: self.read_ns.swap(0, Ordering::Relaxed) / 1_000_000, + build_ms: self.build_ns.swap(0, Ordering::Relaxed) / 1_000_000, + commit_ms: self.commit_ns.swap(0, Ordering::Relaxed) / 1_000_000, + } + } +} + +struct IndexBuildProgress { + next: u32, + tip: u32, + from: u32, + elapsed: Duration, + stages: IndexStageSample, +} + +fn format_index_build_progress(p: &IndexBuildProgress) -> String { + let secs = p.elapsed.as_secs_f64(); + let rate = f64::from(p.next.saturating_sub(p.from)) / secs.max(0.001); + format!( + "index: build next={} tip={} rate={:.0}/s remain={} elapsed={:.0}s read={}ms build={}ms commit={}ms", + p.next, + p.tip, + rate, + p.tip.saturating_add(1).saturating_sub(p.next), + secs, + p.stages.read_ms, + p.stages.build_ms, + p.stages.commit_ms, + ) +} + +fn format_index_apply(start: u32, end: u32, stages: IndexStageSample) -> String { + format!( + "index: apply h={start}..={end} read={}ms build={}ms commit={}ms", + stages.read_ms, stages.build_ms, stages.commit_ms + ) +} + +struct CommitTimes { + ok: bool, + build_ns: u64, + commit_ns: u64, +} + +/// Build and commit what each index still needs from `window`. `ok` is false +/// when a commit found the watermark or a `confirmed[h]` moved. +fn commit_window( + query: &Query, + window: &IndexWindow, + stages: &IndexStageMs, +) -> Result { let Some(first) = window.blocks.first().map(|b| b.height.0) else { - return Ok(true); + return Ok(CommitTimes { + ok: true, + build_ns: 0, + commit_ns: 0, + }); }; - let mut ok = true; - if let Some(next) = query.filter_index_next() { + let t_build = Instant::now(); + let filters = if let Some(next) = query.filter_index_next() { let skip = next.saturating_sub(first) as usize; let built = window .blocks @@ -83,11 +172,11 @@ fn commit_window(query: &Query, window: &IndexWindow) -> Result, ConsensusError>>()?; - if !built.is_empty() { - ok &= query.commit_window_filters(first + skip as u32, &built)? > 0; - } - } - if let Some(next) = query.tweak_index_next() { + Some((first + skip as u32, built)) + } else { + None + }; + let tweaks = if let Some(next) = query.tweak_index_next() { let skip = next.saturating_sub(first) as usize; let items = window .blocks @@ -96,15 +185,38 @@ fn commit_window(query: &Query, window: &IndexWindow) -> Result, ConsensusError>>()?; - if !items.is_empty() { - ok &= query.commit_window_tweaks(&items)? > 0; + Some(items) + } else { + None + }; + let build_ns = t_build.elapsed().as_nanos() as u64; + let t_commit = Instant::now(); + let committed = (|| { + let mut ok = true; + if let Some((start, built)) = &filters { + if !built.is_empty() { + ok &= query.commit_window_filters(*start, built)? > 0; + } } - } + if let Some(items) = &tweaks { + if !items.is_empty() { + ok &= query.commit_window_tweaks(items)? > 0; + } + } + Ok(ok) + })(); + let commit_ns = t_commit.elapsed().as_nanos() as u64; + stages.add_build(build_ns); + stages.add_commit(commit_ns); rbitcoin_query::note_confirm( &query.confirm_stats().blockfilter_ns, - t0.elapsed().as_nanos() as u64, + build_ns.saturating_add(commit_ns), ); - Ok(ok) + committed.map(|ok| CommitTimes { + ok, + build_ns, + commit_ns, + }) } /// Seal every released height on this thread (regtest `generate`, tests). @@ -114,23 +226,49 @@ pub fn build_indexes_released(query: &Query) -> Result<(), ConsensusError> { break; } let end = plan_window(query, start, target)?; - if !commit_window(query, &read_window(query, start, end)?)? { + let stages = IndexStageMs::default(); + if !commit_window(query, &read_window(query, start, end)?, &stages)?.ok { break; } } Ok(()) } +struct ReadyWindow { + window: IndexWindow, + start: u32, + end: u32, + read_ns: u64, + /// Tip path: one height, logged after commit with this window's times. + log_apply: bool, +} + fn cpu_worker( query: &Query, - rx: Receiver, + rx: Receiver, resync: &AtomicBool, + stages: &IndexStageMs, ) -> Result<(), ConsensusError> { - for window in rx { + for ready in rx { if resync.load(Ordering::Acquire) { continue; } - if !commit_window(query, &window)? { + let times = commit_window(query, &ready.window, stages)?; + if ready.log_apply { + rbitcoin_log::info!( + "{}", + format_index_apply( + ready.start, + ready.end, + IndexStageSample { + read_ms: ready.read_ns / 1_000_000, + build_ms: times.build_ns / 1_000_000, + commit_ms: times.commit_ns / 1_000_000, + }, + ) + ); + } + if !times.ok { resync.store(true, Ordering::Release); } } @@ -168,13 +306,14 @@ fn run_writebehind( query.release_index_writebehind(tip); } let resync = AtomicBool::new(false); - let (tx, rx) = sync_channel::(2); + let stages = IndexStageMs::default(); + let (tx, rx) = sync_channel::(2); std::thread::scope(|scope| { let worker = std::thread::Builder::new() .name("rbtc-idx-cpu".into()) - .spawn_scoped(scope, || cpu_worker(query, rx, &resync)) + .spawn_scoped(scope, || cpu_worker(query, rx, &resync, &stages)) .expect("spawn index cpu worker"); - let io = io_loop(query, stop, &resync, &tx, on_filters_caught_up); + let io = io_loop(query, stop, &resync, &stages, &tx, on_filters_caught_up); drop(tx); let cpu = worker.join().expect("index cpu worker"); io.and(cpu) @@ -191,7 +330,8 @@ fn io_loop( query: &Query, stop: &AtomicBool, resync: &AtomicBool, - tx: &std::sync::mpsc::SyncSender, + stages: &IndexStageMs, + tx: &std::sync::mpsc::SyncSender, on_filters_caught_up: impl FnOnce(), ) -> Result<(), ConsensusError> { let mut on_filters_caught_up = Some(on_filters_caught_up); @@ -247,27 +387,91 @@ fn io_loop( } if let Some(p) = pass.as_mut() { if p.last_log.elapsed() >= PROGRESS_EVERY { - let secs = p.started.elapsed().as_secs_f64(); rbitcoin_log::info!( - "index: build next={start} tip={target} rate={:.0}/s remain={} elapsed={secs:.0}s", - f64::from(start - p.from) / secs.max(0.001), - target + 1 - start + "{}", + format_index_build_progress(&IndexBuildProgress { + next: start, + tip: target, + from: p.from, + elapsed: p.started.elapsed(), + stages: stages.take_ms(), + }) ); p.last_log = now; } } let end = plan_window(query, start, target)?; + let t_read = Instant::now(); let window = read_window(query, start, end)?; - if pass.is_none() { - rbitcoin_log::info!( - "index: apply h={start}..={end} read={}ms", - now.elapsed().as_millis() - ); - } - if tx.send(window).is_err() { + let read_ns = t_read.elapsed().as_nanos() as u64; + stages.add_read(read_ns); + if tx + .send(ReadyWindow { + window, + start, + end, + read_ns, + log_apply: pass.is_none(), + }) + .is_err() + { break; } cursor = Some(end + 1); } Ok(()) } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn index_build_progress_names_stage_ms() { + let line = format_index_build_progress(&IndexBuildProgress { + next: 800_000, + tip: 900_000, + from: 700_000, + elapsed: Duration::from_secs(100), + stages: IndexStageSample { + read_ms: 4_000, + build_ms: 5_500, + commit_ms: 12, + }, + }); + assert!(line.contains("next=800000"), "{line}"); + assert!(line.contains("tip=900000"), "{line}"); + assert!(line.contains("rate=1000/s"), "{line}"); + assert!(line.contains("remain=100001"), "{line}"); + assert!(line.contains("elapsed=100s"), "{line}"); + assert!(line.contains("read=4000ms"), "{line}"); + assert!(line.contains("build=5500ms"), "{line}"); + assert!(line.contains("commit=12ms"), "{line}"); + } + + #[test] + fn index_apply_names_stage_ms() { + let line = format_index_apply( + 5, + 5, + IndexStageSample { + read_ms: 20, + build_ms: 30, + commit_ms: 4, + }, + ); + assert_eq!(line, "index: apply h=5..=5 read=20ms build=30ms commit=4ms"); + } + + #[test] + fn stage_sample_take_resets() { + let s = IndexStageMs::default(); + s.add_read(2_500_000); + s.add_build(1_500_000); + s.add_commit(500_000); + let a = s.take_ms(); + assert_eq!((a.read_ms, a.build_ms, a.commit_ms), (2, 1, 0)); + let b = s.take_ms(); + assert_eq!((b.read_ms, b.build_ms, b.commit_ms), (0, 0, 0)); + } +} From 6e6c8855581d497ee1bd0ff922e2342b46f70fd4 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Sat, 26 Sep 2026 17:37:29 -0700 Subject: [PATCH 2/2] index: assemble each catch-up height on the script pool Filter GCS and tweak EC ran on one thread, which is the slow tail once blocks are large. The CPU thread still commits the window. Co-authored-by: Cursor --- .../src/index_writebehind.rs | 295 +++++++++++++++--- crates/rbitcoin-consensus/src/script_pool.rs | 7 + crates/rbitcoin-query/src/block_filter.rs | 50 +-- crates/rbitcoin-query/src/lib.rs | 1 + .../rbitcoin-store/src/index_build_uring.rs | 4 + docs/concurrency.md | 2 +- docs/operator/operations.md | 2 +- docs/operator/storage.md | 11 +- 8 files changed, 305 insertions(+), 67 deletions(-) diff --git a/crates/rbitcoin-consensus/src/index_writebehind.rs b/crates/rbitcoin-consensus/src/index_writebehind.rs index ea6244691..0c35bd0af 100644 --- a/crates/rbitcoin-consensus/src/index_writebehind.rs +++ b/crates/rbitcoin-consensus/src/index_writebehind.rs @@ -4,13 +4,14 @@ //! The IO thread plans windows of consecutive heights from the lower index //! watermark up to the released tip and reads each through //! [`rbitcoin_store::read_index_window`] (one completion session). One CPU -//! worker builds each index for the heights that index still needs (tweak -//! EC math included) and commits each index once per window +//! worker 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 lock. A commit that finds its watermark or a //! `confirmed[h]` moved (a reorg) returns 0, and the IO thread re-plans from //! 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::silent_payments::tweak_records_from_window; use crate::ConsensusError; use rbitcoin_primitives::{Fk, Height}; @@ -18,7 +19,7 @@ use rbitcoin_query::Query; use rbitcoin_store::IndexWindow; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::mpsc::{sync_channel, Receiver}; -use std::sync::Arc; +use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; /// Heights per window. Bounds how long a disconnect waits and how far the @@ -63,7 +64,97 @@ fn plan_window(query: &Query, start: u32, target: u32) -> Result Result { let tweaks_from = query.tweak_index_next(); let heights = query.index_heights(start, end, tweaks_from)?; - Ok(query.read_index_window(&heights)?) + let mut window = query.read_index_window(&heights)?; + for block in &mut window.blocks { + block.hash = query.store().get_header(block.header_fk)?.hash; + } + Ok(window) +} + +/// One tx: `None` when it has no tweak (coinbase, no P2TR output, ineligible). +type HeightTweaks = Vec>; + +struct IndexHeightOut { + filter: Option<(bitcoin::bip158::BlockFilter, Fk)>, + tweaks: Option<(Height, Fk, HeightTweaks)>, +} + +struct IndexHeightJob { + window: Arc, + index: usize, + want_filter: bool, + want_tweaks: bool, + out: Arc>>, +} + +fn assemble_index_height(job: &IndexHeightJob) -> Result<(), ConsensusError> { + let block = &job.window.blocks[job.index]; + let filter = if job.want_filter { + Some(( + rbitcoin_query::basic_filter_of(&block.hash, &job.window, job.index)?, + block.header_fk, + )) + } else { + None + }; + let tweaks = if job.want_tweaks { + Some(( + block.height, + block.header_fk, + tweak_records_from_window(&job.window, job.index)?, + )) + } else { + None + }; + *job.out.lock().unwrap_or_else(|e| e.into_inner()) = Some(IndexHeightOut { filter, tweaks }); + Ok(()) +} + +struct Assembled { + filters: Vec<(bitcoin::bip158::BlockFilter, Fk)>, + tweaks: Vec<(Height, Fk, HeightTweaks)>, +} + +/// One job per height that still needs a filter or tweaks. The wave joins +/// before return, so the `Arc` is only shared for that call. +fn assemble_window(query: &Query, window: &Arc) -> Result { + let filter_from = query.filter_index_next(); + let tweak_from = query.tweak_index_next(); + let mut slots = Vec::with_capacity(window.blocks.len()); + let mut jobs = Vec::new(); + for (index, block) in window.blocks.iter().enumerate() { + let h = block.height.0; + let want_filter = filter_from.is_some_and(|n| h >= n); + let want_tweaks = tweak_from.is_some_and(|n| h >= n); + let out = Arc::new(Mutex::new(None)); + if want_filter || want_tweaks { + jobs.push(IndexHeightJob { + window: Arc::clone(window), + index, + want_filter, + want_tweaks, + out: Arc::clone(&out), + }); + } + slots.push(out); + } + if let Some(wave) = start_for_each_owned(jobs, assemble_index_height)? { + wave.finish()?; + } + let mut filters = Vec::new(); + let mut tweaks = Vec::new(); + for slot in &slots { + let Some(done) = slot.lock().unwrap_or_else(|e| e.into_inner()).take() else { + continue; + }; + if let Some(filter) = done.filter { + filters.push(filter); + } + if let Some(tweak) = done.tweaks { + tweaks.push(tweak); + } + } + Ok(Assembled { filters, tweaks }) } /// Milliseconds of one interval, or of one tip window. @@ -152,7 +243,7 @@ struct CommitTimes { /// when a commit found the watermark or a `confirmed[h]` moved. fn commit_window( query: &Query, - window: &IndexWindow, + window: &Arc, stages: &IndexStageMs, ) -> Result { let Some(first) = window.blocks.first().map(|b| b.height.0) else { @@ -163,45 +254,17 @@ fn commit_window( }); }; let t_build = Instant::now(); - let filters = if let Some(next) = query.filter_index_next() { - let skip = next.saturating_sub(first) as usize; - let built = window - .blocks - .iter() - .enumerate() - .skip(skip) - .map(|(i, b)| Ok((query.basic_filter_from_window(window, i)?, b.header_fk))) - .collect::, ConsensusError>>()?; - Some((first + skip as u32, built)) - } else { - None - }; - let tweaks = if let Some(next) = query.tweak_index_next() { - let skip = next.saturating_sub(first) as usize; - let items = window - .blocks - .iter() - .enumerate() - .skip(skip) - .map(|(i, b)| Ok((b.height, b.header_fk, tweak_records_from_window(window, i)?))) - .collect::, ConsensusError>>()?; - Some(items) - } else { - None - }; + let assembled = assemble_window(query, window)?; let build_ns = t_build.elapsed().as_nanos() as u64; let t_commit = Instant::now(); let committed = (|| { let mut ok = true; - if let Some((start, built)) = &filters { - if !built.is_empty() { - ok &= query.commit_window_filters(*start, built)? > 0; - } + if !assembled.filters.is_empty() { + let start = query.filter_index_next().unwrap_or(first).max(first); + ok &= query.commit_window_filters(start, &assembled.filters)? > 0; } - if let Some(items) = &tweaks { - if !items.is_empty() { - ok &= query.commit_window_tweaks(items)? > 0; - } + if !assembled.tweaks.is_empty() { + ok &= query.commit_window_tweaks(&assembled.tweaks)? > 0; } Ok(ok) })(); @@ -227,7 +290,8 @@ pub fn build_indexes_released(query: &Query) -> Result<(), ConsensusError> { } let end = plan_window(query, start, target)?; let stages = IndexStageMs::default(); - if !commit_window(query, &read_window(query, start, end)?, &stages)?.ok { + let window = Arc::new(read_window(query, start, end)?); + if !commit_window(query, &window, &stages)?.ok { break; } } @@ -253,7 +317,8 @@ fn cpu_worker( if resync.load(Ordering::Acquire) { continue; } - let times = commit_window(query, &ready.window, stages)?; + let window = Arc::new(ready.window); + let times = commit_window(query, &window, stages)?; if ready.log_apply { rbitcoin_log::info!( "{}", @@ -474,4 +539,152 @@ mod tests { let b = s.take_ms(); assert_eq!((b.read_ms, b.build_ms, b.commit_ms), (0, 0, 0)); } + + /// Three heights: coinbase (ineligible), a same-window spend into P2TR, + /// and a later coinbase that is also ineligible. Pooled assemble matches + /// the serial window walk. + #[test] + fn per_height_jobs_match_serial_window() { + use bitcoin::hashes::{hash160, Hash}; + use bitcoin::secp256k1::{PublicKey, Secp256k1, SecretKey}; + use rbitcoin_query::testutil::FixtureChain; + use rbitcoin_store::{InputRecord, OutputRecord}; + + let _gate = crate::script_pool::steal_test_gate(); + let (dir, q) = rbitcoin_query::testutil::tiny_query_labeled("idx-jobs"); + q.set_block_filter_index(true).unwrap(); + q.set_sptweaks_enabled(true, Height(0)).unwrap(); + + let secp = Secp256k1::new(); + let sk = SecretKey::from_slice(&[2u8; 32]).unwrap(); + let pk = PublicKey::from_secret_key(&secp, &sk); + let ser = pk.serialize(); + let h160 = hash160::Hash::hash(&ser); + let mut p2wpkh = vec![0x00, 0x14]; + p2wpkh.extend_from_slice(h160.as_ref()); + let (xonly, _) = pk.x_only_public_key(); + let mut p2tr = vec![0x51, 0x20]; + p2tr.extend_from_slice(&xonly.serialize()); + + let mut genesis_txid = [0u8; 32]; + genesis_txid[31] = 0xcb; + let h0 = header_rec(0, Fk::NULL, None); + let fk0 = q + .connect_block( + Height(0), + &h0, + &[tx_apply( + genesis_txid, + vec![InputRecord::coinbase(u32::MAX, vec![0x00], vec![])], + vec![OutputRecord::unspent(50_0000_0000, p2wpkh.clone())], + )], + ) + .unwrap(); + let create_fk = q.block_tx_fks(Height(0)).unwrap()[0]; + + let mut spend_txid = [0u8; 32]; + spend_txid[0] = 0x11; + spend_txid[31] = 0xcd; + let h1 = header_rec(1, fk0, Some(h0.hash)); + let fk1 = q + .connect_block( + Height(1), + &h1, + &[tx_apply( + spend_txid, + vec![InputRecord { + prev_txid: genesis_txid, + create_fk, + prev_index: 0, + sequence: u32::MAX, + script_sig: vec![], + witness: vec![vec![0u8; 64], ser.to_vec()], + }], + vec![OutputRecord::unspent(49_0000_0000, p2tr)], + )], + ) + .unwrap(); + + let mut later_txid = [0u8; 32]; + later_txid[31] = 0xee; + let h2 = header_rec(2, fk1, Some(h1.hash)); + q.connect_block( + Height(2), + &h2, + &[tx_apply( + later_txid, + vec![InputRecord::coinbase(u32::MAX, vec![0x01], vec![])], + vec![OutputRecord::unspent(50_0000_0000, p2wpkh)], + )], + ) + .unwrap(); + + let window = Arc::new(read_window(&q, 0, 2).unwrap()); + assert_eq!(window.blocks.len(), 3); + let got = assemble_window(&q, &window).unwrap(); + assert_eq!(got.filters.len(), 3); + assert_eq!(got.tweaks.len(), 3); + for (i, block) in window.blocks.iter().enumerate() { + let serial_f = q.basic_filter_from_window(&window, i).unwrap(); + assert_eq!(got.filters[i].0.content, serial_f.content, "filter {i}"); + assert_eq!(got.filters[i].1, block.header_fk); + let serial_t = crate::silent_payments::tweak_records_from_window(&window, i).unwrap(); + assert_eq!(got.tweaks[i].2, serial_t, "tweaks {i}"); + } + assert!(got.tweaks[0].2.iter().all(|t| t.is_none()), "coinbase"); + assert!( + got.tweaks[1].2.iter().any(|t| t.is_some()), + "same-window P2TR spend" + ); + assert!( + got.tweaks[2].2.iter().all(|t| t.is_none()), + "no P2TR output" + ); + let _ = std::fs::remove_dir_all(dir.path()); + } + + fn header_rec( + h: u32, + prev_fk: Fk, + prev_hash: Option<[u8; 32]>, + ) -> rbitcoin_store::HeaderRecord { + let mut merkle = [0u8; 32]; + merkle[0..4].copy_from_slice(&h.to_le_bytes()); + merkle[5] = 0xec; + let hash = match prev_hash { + None => merkle, + Some(ph) => rbitcoin_store::block_header_hash(1, &ph, &merkle, h + 1, 0x207f_ffff, h), + }; + rbitcoin_store::HeaderRecord { + prev_fk, + version: 1, + timestamp: h + 1, + bits: 0x207f_ffff, + nonce: h, + merkle_root: merkle, + hash, + size: 0, + weight: 0, + } + } + + fn tx_apply( + txid: [u8; 32], + inputs: Vec, + outputs: Vec, + ) -> rbitcoin_query::TxApply { + rbitcoin_query::TxApply { + tx: rbitcoin_store::TxRecord { + txid, + version: 2, + locktime: 0, + input_start_fk: Fk::NULL, + input_count: inputs.len() as u32, + output_start_fk: Fk::NULL, + output_count: outputs.len() as u32, + }, + inputs, + outputs, + } + } } diff --git a/crates/rbitcoin-consensus/src/script_pool.rs b/crates/rbitcoin-consensus/src/script_pool.rs index 499a68411..1e1f593ac 100644 --- a/crates/rbitcoin-consensus/src/script_pool.rs +++ b/crates/rbitcoin-consensus/src/script_pool.rs @@ -162,6 +162,13 @@ static STEAL_CLAIMS_ON: AtomicBool = AtomicBool::new(false); #[cfg(test)] static STEAL_TEST: Mutex<()> = Mutex::new(()); +/// Hold across a test that publishes a steal wave, so pool tests do not +/// interleave waves. +#[cfg(test)] +pub(crate) fn steal_test_gate() -> std::sync::MutexGuard<'static, ()> { + STEAL_TEST.lock().unwrap_or_else(|p| p.into_inner()) +} + fn waves_snap() -> &'static ArcSwap>> { WAVES_SNAP.get_or_init(|| ArcSwap::from_pointee(Vec::new())) } diff --git a/crates/rbitcoin-query/src/block_filter.rs b/crates/rbitcoin-query/src/block_filter.rs index 261beaf04..a25935d4c 100644 --- a/crates/rbitcoin-query/src/block_filter.rs +++ b/crates/rbitcoin-query/src/block_filter.rs @@ -50,6 +50,34 @@ impl BlockFilterWriteBehind { } } +/// Basic filter of `window.blocks[i]` using a hash the caller already loaded. +/// +/// No store IO. Output scripts other than `OP_RETURN`, then each spent prevout +/// script. A missing prevout is corrupt. +pub fn basic_filter_of( + block_hash: &[u8; 32], + window: &IndexWindow, + i: usize, +) -> Result { + const OP_RETURN: u8 = 0x6a; + let block = &window.blocks[i]; + let mut elements: Vec<&[u8]> = Vec::new(); + for tx in &block.txs { + for o in &tx.outs { + if o.script.first() != Some(&OP_RETURN) { + elements.push(&o.script); + } + } + } + for e in block.edges.iter().flatten().filter(|e| !e.parent.is_null()) { + let out = window.prevout(e.parent, e.vout).ok_or(StoreError::Corrupt( + "invariant: blockfilter prevout missing", + ))?; + elements.push(&out.script); + } + encode_basic_filter(block_hash, elements.into_iter()) +} + /// GCS-encode a basic filter keyed by `block_hash` (internal byte order). fn encode_basic_filter<'a>( block_hash: &[u8; 32], @@ -105,30 +133,14 @@ impl Query { read_index_window(&self.store.txs, heights) } - /// Basic filter of `window.blocks[i]`. + /// Basic filter of `window.blocks[i]`. Reads the block hash from the header. pub fn basic_filter_from_window( &self, window: &IndexWindow, i: usize, ) -> Result { - const OP_RETURN: u8 = 0x6a; - let block = &window.blocks[i]; - let hash = self.store.get_header(block.header_fk)?.hash; - let mut elements: Vec<&[u8]> = Vec::new(); - for tx in &block.txs { - for o in &tx.outs { - if o.script.first() != Some(&OP_RETURN) { - elements.push(&o.script); - } - } - } - for e in block.edges.iter().flatten().filter(|e| !e.parent.is_null()) { - let out = window.prevout(e.parent, e.vout).ok_or(StoreError::Corrupt( - "invariant: blockfilter prevout missing", - ))?; - elements.push(&out.script); - } - encode_basic_filter(&hash, elements.into_iter()) + let hash = self.store.get_header(window.blocks[i].header_fk)?.hash; + basic_filter_of(&hash, window, i) } pub fn block_filter_enabled(&self) -> bool { diff --git a/crates/rbitcoin-query/src/lib.rs b/crates/rbitcoin-query/src/lib.rs index d0ff34857..e6645afc4 100644 --- a/crates/rbitcoin-query/src/lib.rs +++ b/crates/rbitcoin-query/src/lib.rs @@ -168,6 +168,7 @@ pub(crate) use batch_parents::FkSet; pub use batch_parents::{ layout_covers_need, sparse_spender_rels, BatchParents, FkMap, U32Map, U64Map, U64Set, }; +pub use block_filter::basic_filter_of; pub use catchup::IndexMode; pub use chain_view::{ChainView, ChainViewKind}; pub use confirm_load::SpendEdges; diff --git a/crates/rbitcoin-store/src/index_build_uring.rs b/crates/rbitcoin-store/src/index_build_uring.rs index b5b07aab9..86638cbd0 100644 --- a/crates/rbitcoin-store/src/index_build_uring.rs +++ b/crates/rbitcoin-store/src/index_build_uring.rs @@ -50,6 +50,9 @@ pub struct IndexHeight { pub struct IndexBlock { pub height: Height, pub header_fk: Fk, + /// Block hash, internal byte order. Zero until the index IO thread fills + /// it from the header; assemble reads this and does not touch the store. + pub hash: [u8; 32], /// `inputs` is `Some` only for tweak-height txs with a P2TR output. pub txs: Vec, pub edges: Vec>, @@ -446,6 +449,7 @@ fn decode_block( Ok(IndexBlock { height: h.height, header_fk: h.header_fk, + hash: [0u8; 32], txs, edges, }) diff --git a/docs/concurrency.md b/docs/concurrency.md index a044d2ed3..7c019a0aa 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, builds filters and tweak records (tweak EC math on that thread, not `rbtc-scripts-*`) 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) 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. | | 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 | diff --git a/docs/operator/operations.md b/docs/operator/operations.md index 8cd3e9814..cb1d68a86 100644 --- a/docs/operator/operations.md +++ b/docs/operator/operations.md @@ -90,7 +90,7 @@ Clean smoke: | `--asmap PATH` | `asmap=` | unset — try `{datadir}/ip_asn.dat` if present; else prefix groups | | `--no-seeds` | `no_seeds=` | seeds on | | `--sh-index` | `sh_index=` | **off** — Class B scripthash (address/history; Electrum/Esplora start without it) | -| `--block-filter-index` | `block_filter_index=` | **off** — BIP158 basic filters. Independent of `--sh-index`. IBD does not build them. After catch-up the `rbtc-idx-wb` builder materializes the gap from Class A in the background (`index: build from=… to=…`, progress every 10 s, `index: build done`), then seals each new tip (`index: apply h=…`); follow, relay, and Electrum do not wait for it. With `--sp-tweaks` the same pass builds both indexes. Works with `--prune-seqsigwit`. `tip: accept` shows `bf=` and `bf_lag=`. `NODE_COMPACT_FILTERS` is advertised once filters first reach the tip (`blockfilter: caught up …`); peers that connected earlier do not learn the bit. `getblockfilter`, `/rest/blockfilter/`, and P2P serve heights the watermark covers; a stop past the watermark is silence | +| `--block-filter-index` | `block_filter_index=` | **off** — BIP158 basic filters. Independent of `--sh-index`. IBD does not build them. After catch-up the `rbtc-idx-wb` builder materializes the gap from Class A in the background (`index: build from=… to=…`, progress every 10 s with `read=` `build=` `commit=` ms for that interval, `index: build done`), then seals each new tip (`index: apply h=… read=` `build=` `commit=`); follow, relay, and Electrum do not wait for it. With `--sp-tweaks` the same pass builds both indexes. Works with `--prune-seqsigwit`. `tip: accept` shows `bf=` and `bf_lag=`. `NODE_COMPACT_FILTERS` is advertised once filters first reach the tip (`blockfilter: caught up …`); peers that connected earlier do not learn the bit. `getblockfilter`, `/rest/blockfilter/`, and P2P serve heights the watermark covers; a stop past the watermark is silence | | `--prune-seqsigwit` | `prune_seqsigwit=` | **off** — unpruned reads `seqsigwit.body`. On: refuse wire reconstruct below tip−288 **heights**, advertise `NETWORK_LIMITED`, and keep those heights as `store/seqsigwit.window/{height}.bin` plus a RAM cache. Refused with `--sp-tweaks`, and once pruned a datadir serves no tweaks | | `--prune-seqsigwit-ram-threshold-bytes N` | `prune_seqsigwit_ram_threshold_bytes=` | `268435456` (256 MiB). `0` keeps nothing in RAM: every height, including tiny IBD blocks, is read from its file | | `--max-sh-creates N` | `max_sh_creates=` | **10000** — unpaged SH join above N is refused (503 / JSON-RPC error). **0** is unlimited. A request that names a page still returns that page. | diff --git a/docs/operator/storage.md b/docs/operator/storage.md index 68ab09f9d..dfbf72604 100644 --- a/docs/operator/storage.md +++ b/docs/operator/storage.md @@ -269,11 +269,12 @@ shared with `--block-filter-index`) from the taproot origin (709632 on mainnet), then sealed per released tip block; the confirm write thread writes no tweaks. One IO thread reads windows of heights on one completion session (`seqsigwit` and parent txids only for P2TR-output txs); one CPU thread -(`rbtc-idx-cpu`) computes the tweaks, secp included, and commits one batched -height-blob + idx write per window. It does not borrow `rbtc-scripts-*`, so -block scripts and mempool accept never share workers with it. Reorg -truncates with tip. Kill-safe: `next_height` is the last complete put. INFO -every 10 s: `index: build next=… tip=… rate=…/s remain=…`. +(`rbtc-idx-cpu`) publishes one job per height to `rbtc-scripts-*` (tweak EC, +and filter GCS when that index is on) and commits one batched height-blob + +idx write per window. A catch-up wave shares that pool with block scripts for +at most one window. Reorg truncates with tip. Kill-safe: `next_height` is the +last complete put. INFO every 10 s: `index: build next=… tip=… rate=…/s +remain=… read=…ms build=…ms commit=…ms`. Cake Wallet’s scan isolate may still hardcode `electrs.cakewallet.com` even after a successful probe — see `COMPAT.md`.