diff --git a/crates/rbitcoin-consensus/src/index_writebehind.rs b/crates/rbitcoin-consensus/src/index_writebehind.rs index 96c44b39d..0c35bd0af 100644 --- a/crates/rbitcoin-consensus/src/index_writebehind.rs +++ b/crates/rbitcoin-consensus/src/index_writebehind.rs @@ -4,21 +4,22 @@ //! 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}; 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::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; /// Heights per window. Bounds how long a disconnect waits and how far the @@ -63,48 +64,222 @@ 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) } -/// Build and commit what each index still needs from `window`. Returns -/// false when a commit found the watermark or a `confirmed[h]` moved. -fn commit_window(query: &Query, window: &IndexWindow) -> Result { - let t0 = Instant::now(); - let Some(first) = window.blocks.first().map(|b| b.height.0) else { - return Ok(true); +/// 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 }; - let mut ok = true; - 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>>()?; - if !built.is_empty() { - ok &= query.commit_window_filters(first + skip as u32, &built)? > 0; + *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. +#[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); } } - 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>>()?; - if !items.is_empty() { - ok &= query.commit_window_tweaks(&items)? > 0; + + 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: &Arc, + stages: &IndexStageMs, +) -> Result { + let Some(first) = window.blocks.first().map(|b| b.height.0) else { + return Ok(CommitTimes { + ok: true, + build_ns: 0, + commit_ns: 0, + }); + }; + let t_build = Instant::now(); + 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 !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 !assembled.tweaks.is_empty() { + ok &= query.commit_window_tweaks(&assembled.tweaks)? > 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 +289,51 @@ 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(); + let window = Arc::new(read_window(query, start, end)?); + if !commit_window(query, &window, &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 window = Arc::new(ready.window); + let times = commit_window(query, &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 +371,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 +395,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 +452,239 @@ 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)); + } + + /// 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`.