diff --git a/changelog.d/ibd-index-append.md b/changelog.d/ibd-index-append.md new file mode 100644 index 000000000..694a9f595 --- /dev/null +++ b/changelog.d/ibd-index-append.md @@ -0,0 +1,7 @@ +Changed + +- **Block filters and silent-payment tweaks are sealed during sync when + they are enabled before it.** `--block-filter-index` and `--sp-tweaks` + append each connected batch on the confirm write thread. A restart gap of + at most one write drain is sealed at startup. Turning either index on + after those blocks were connected still builds the gap after catch-up. diff --git a/crates/rbitcoin-consensus/src/confirm_run/index.rs b/crates/rbitcoin-consensus/src/confirm_run/index.rs new file mode 100644 index 000000000..1081852bd --- /dev/null +++ b/crates/rbitcoin-consensus/src/confirm_run/index.rs @@ -0,0 +1,470 @@ +//! Live filter and tweak bytes for one confirm batch. +//! +//! Assemble runs only while [`Query::index_live`] is set, from the wire block +//! and [`BatchParents`] (same-batch creates and pinned parents). No store read. +//! The flag is set at startup, after a short restart gap is sealed, so load +//! does not sample `next_height` (the previous batch can still be committing). +//! A write that finds the watermark is not this batch is corrupt. + +use std::collections::HashMap; +use std::sync::{Arc, Mutex}; +use std::time::Instant; + +use bitcoin::hashes::Hash; +use bitcoin::{Amount, Block, ScriptBuf, TxOut}; +use rbitcoin_primitives::{Fk, Height}; +use rbitcoin_query::Query; +use rbitcoin_store::StoreError; + +use crate::index_rows::{self, IndexHeightOut, IndexRows}; +use crate::silent_payments::tweak_from_tx; +use crate::ConsensusError; + +use super::{LoadedBatch, Prepared}; + +/// Which indexes this batch should assemble. Set at load from [`Query::index_live`]. +#[derive(Clone, Debug, Default)] +pub(super) struct IndexWant { + pub filters: bool, + /// Tweaks for heights at or above this origin. `None` when tweaks are off. + pub tweak_origin: Option, +} + +/// One confirm batch of filter bytes and tweak vecs. Dropped with the batch. +pub(super) type IndexSeal = IndexRows; + +struct LiveJob { + block: Arc, + hash: [u8; 32], + height: Height, + header_fk: Fk, + /// Per tx. Empty for a coinbase. Non-coinbase length matches `tx.input`. + prevouts: Vec>, + want_filter: bool, + want_tweaks: bool, + out: Arc>>, +} + +pub(super) fn index_want(query: &Query) -> IndexWant { + if !query.index_live() { + return IndexWant::default(); + } + IndexWant { + filters: query.filter_index_next().is_some(), + tweak_origin: query.tweak_index_next().map(|_| query.sptweaks_origin().0), + } +} + +pub(super) fn live_index_seal(batch: &LoadedBatch) -> Result<(IndexRows, u64), ConsensusError> { + let want = &batch.index_want; + if !want.filters && want.tweak_origin.is_none() { + return Ok((IndexSeal::default(), 0)); + } + let t0 = Instant::now(); + if batch.prepared.len() != batch.wire_blocks.len() { + return Err(corrupt("invariant: live index block count")); + } + let mut slots = Vec::with_capacity(batch.prepared.len()); + let mut jobs = Vec::new(); + for (prep, block) in batch.prepared.iter().zip(batch.wire_blocks.iter()) { + let want_tweaks = want.tweak_origin.is_some_and(|o| prep.height.0 >= o); + let out = Arc::new(Mutex::new(None)); + if want.filters || want_tweaks { + let prevouts = prevouts_for_block(prep, block, &batch.batch_parents)?; + jobs.push(LiveJob { + block: Arc::clone(block), + hash: prep.hash, + height: prep.height, + header_fk: prep.header_fk, + prevouts, + want_filter: want.filters, + want_tweaks, + out: Arc::clone(&out), + }); + } + slots.push(out); + } + index_rows::finish_wave(jobs, assemble_live_height)?; + let ns = t0.elapsed().as_nanos() as u64; + Ok((index_rows::rows_from_slots(&slots), ns)) +} + +fn assemble_live_height(job: &LiveJob) -> Result<(), ConsensusError> { + let filter = if job.want_filter { + Some(( + job.height, + basic_filter_from_wire(&job.hash, &job.block, &job.prevouts)?, + job.header_fk, + )) + } else { + None + }; + let tweaks = if job.want_tweaks { + Some(( + job.height, + job.header_fk, + tweaks_from_wire(&job.block, &job.prevouts), + )) + } else { + None + }; + *job.out.lock().unwrap_or_else(|e| e.into_inner()) = Some(IndexHeightOut { filter, tweaks }); + Ok(()) +} + +fn basic_filter_from_wire( + hash: &[u8; 32], + block: &Block, + prevouts: &[Vec], +) -> Result { + let outputs = block + .txdata + .iter() + .flat_map(|tx| tx.output.iter().map(|o| o.script_pubkey.as_bytes())); + let spent = prevouts + .iter() + .flatten() + .map(|prev| prev.script_pubkey.as_bytes()); + rbitcoin_query::basic_filter_from_scripts(hash, outputs, spent).map_err(ConsensusError::from) +} + +fn tweaks_from_wire(block: &Block, prevouts: &[Vec]) -> index_rows::HeightTweaks { + block + .txdata + .iter() + .zip(prevouts.iter()) + .map(|(tx, prev)| tweak_from_tx(tx, prev).map(|t| t.tweak)) + .collect() +} + +fn prevouts_for_block( + prep: &Prepared, + block: &Block, + parents: &rbitcoin_query::BatchParents, +) -> Result>, ConsensusError> { + let mut same_block: HashMap<[u8; 32], usize> = HashMap::new(); + for (i, tx) in block.txdata.iter().enumerate() { + same_block.insert(tx.compute_txid().to_byte_array(), i); + } + let mut by_tx: HashMap> = HashMap::new(); + for &(prev_txid, vout, spending_fk, create_fk, vin) in &prep.spends { + if spending_fk.is_null() { + return Err(corrupt("invariant: live index spend fk")); + } + let txout = if create_fk.is_null() { + let Some(&pi) = same_block.get(&prev_txid) else { + return Err(corrupt("invariant: blockfilter prevout missing")); + }; + let Some(o) = block.txdata[pi].output.get(vout as usize) else { + return Err(corrupt("invariant: blockfilter prevout missing")); + }; + o.clone() + } else { + let Some((value, script, txid)) = + parents.get_parent_txout_parts(create_fk, vout, |v, s, txid| (v, s.to_vec(), txid)) + else { + return Err(corrupt("invariant: blockfilter prevout missing")); + }; + if txid != prev_txid { + return Err(corrupt("invariant: live index prevout txid")); + } + if value < 0 { + return Err(corrupt("invariant: live index prevout value")); + } + TxOut { + value: Amount::from_sat(value as u64), + script_pubkey: ScriptBuf::from_bytes(script), + } + }; + by_tx.entry(spending_fk).or_default().push((vin, txout)); + } + let mut out = Vec::with_capacity(block.txdata.len()); + for (ti, tx) in block.txdata.iter().enumerate() { + if tx.is_coinbase() { + out.push(Vec::new()); + continue; + } + let fk = prep + .tx_fks + .get(ti) + .copied() + .ok_or_else(|| corrupt("invariant: live index tx fk"))?; + let mut rows = by_tx.remove(&fk).unwrap_or_default(); + rows.sort_by_key(|(vin, _)| *vin); + if rows.len() != tx.input.len() { + return Err(corrupt("invariant: live index prevout count")); + } + out.push(rows.into_iter().map(|(_, o)| o).collect()); + } + Ok(out) +} + +/// After `confirmed[h]` is set: commit this batch. A watermark that is not +/// the first row is corrupt. An empty seal (live append off) is a no-op. +pub(super) fn seal_live_indexes(query: &Query, seal: IndexSeal) -> Result { + if seal.filters.is_empty() && seal.tweaks.is_empty() { + return Ok(0); + } + let commit = index_rows::commit_index_rows(query, seal)?; + if commit.moved { + return Err(corrupt("invariant: live index commit")); + } + Ok(commit.put_ns) +} + +fn corrupt(msg: &'static str) -> ConsensusError { + ConsensusError::Store(StoreError::Corrupt(msg)) +} + +#[cfg(test)] +mod tests { + use super::*; + use bitcoin::absolute::LockTime; + use bitcoin::hashes::{hash160, Hash}; + use bitcoin::secp256k1::{PublicKey, Secp256k1, SecretKey}; + use bitcoin::transaction::Version as TxVersion; + use bitcoin::{Amount, OutPoint, Sequence, TxIn, TxOut, Witness}; + use rbitcoin_query::testutil::tiny_query_labeled; + use rbitcoin_query::Query; + + use crate::silent_payments::tweak_records_from_window; + use crate::{ + accept_and_connect_block, confirm_wire_run, genesis_block, mine_empty_regtest, + mine_regtest_paying, ChainParams, Milestone, + }; + + fn keys() -> (Vec, bitcoin::ScriptBuf, bitcoin::ScriptBuf) { + let secp = Secp256k1::new(); + let sk = SecretKey::from_slice(&[2u8; 32]).unwrap(); + let pk = PublicKey::from_secret_key(&secp, &sk); + let ser = pk.serialize().to_vec(); + 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()); + ( + ser, + bitcoin::ScriptBuf::from_bytes(p2wpkh), + bitcoin::ScriptBuf::from_bytes(p2tr), + ) + } + + fn enable(q: &Query) { + q.set_block_filter_index(true).unwrap(); + q.set_sptweaks_enabled(true, Height(0)).unwrap(); + } + + fn go_live(q: &Query) { + crate::prepare_live_indexes(q).unwrap(); + } + + fn assert_sealed_matches_window(q: &Query, tip: u32) { + assert_eq!(q.filter_index_next(), Some(tip + 1)); + assert_eq!(q.tweak_index_next(), Some(tip + 1)); + for h in 0..=tip { + let stored = q.basic_filter_at(h).unwrap().unwrap().0; + let heights = q.index_heights(h, h, Some(0)).unwrap(); + let window = q.read_index_window(&heights).unwrap(); + let expect = q.basic_filter_from_window(&window, 0).unwrap(); + assert_eq!(stored, expect.content, "filter h={h}"); + let recs = tweak_records_from_window(&window, 0).unwrap(); + let want: Vec<[u8; 33]> = recs.into_iter().flatten().collect(); + let got = q.load_thin_tweaks(Height(h)).unwrap().unwrap_or_default(); + let got: Vec<[u8; 33]> = got.into_iter().map(|r| r.tweak).collect(); + assert_eq!(got, want, "tweaks h={h}"); + } + } + + /// Three roles in one batch: a P2WPKH coinbase, a same-batch spend into + /// P2TR once that coinbase is mature, and the blocks between them. + #[test] + fn live_batch_matches_window_and_reaches_tip() { + let _gate = crate::script_pool::steal_test_gate(); + let (_dir, q) = tiny_query_labeled("idx-live"); + enable(&q); + go_live(&q); + let params = ChainParams::regtest(); + let ms = Milestone::height(u32::MAX); + let genesis = genesis_block(¶ms); + accept_and_connect_block(&q, ¶ms, Height(0), &genesis, ms).unwrap(); + + let (ser, p2wpkh, p2tr) = keys(); + let spend_h = 101u32; + let mut blocks = Vec::new(); + let mut prev = genesis.block_hash(); + let mut time = genesis.header.time; + let mut pay_value = Amount::ZERO; + let mut pay_txid = None; + for h in 1..=spend_h { + time += 600; + let block = if h == 1 { + mine_regtest_paying(prev, time, h, p2wpkh.clone(), Vec::new()) + } else if h == spend_h { + let txid = pay_txid.expect("paying coinbase"); + let spend = bitcoin::Transaction { + version: TxVersion::TWO, + lock_time: LockTime::ZERO, + input: vec![TxIn { + previous_output: OutPoint::new(txid, 0), + script_sig: bitcoin::ScriptBuf::new(), + sequence: Sequence::MAX, + witness: Witness::from_slice(&[vec![0u8; 64].as_slice(), ser.as_slice()]), + }], + output: vec![TxOut { + value: pay_value, + script_pubkey: p2tr.clone(), + }], + }; + let child = bitcoin::Transaction { + version: TxVersion::TWO, + lock_time: LockTime::ZERO, + input: vec![TxIn { + previous_output: OutPoint::new(spend.compute_txid(), 0), + script_sig: bitcoin::ScriptBuf::new(), + sequence: Sequence::MAX, + witness: Witness::new(), + }], + output: vec![TxOut { + value: pay_value, + script_pubkey: p2tr.clone(), + }], + }; + mine_regtest_paying(prev, time, h, p2wpkh.clone(), vec![spend, child]) + } else { + mine_empty_regtest(prev, time, h) + }; + if h == 1 { + pay_value = block.txdata[0].output[0].value; + pay_txid = Some(block.txdata[0].compute_txid()); + } + prev = block.block_hash(); + time = block.header.time; + blocks.push((Height(h), block)); + } + confirm_wire_run(&q, ¶ms, ms, &blocks).unwrap(); + assert_sealed_matches_window(&q, spend_h); + let rows = q.load_thin_tweaks(Height(spend_h)).unwrap().unwrap(); + assert_eq!( + rows.len(), + 2, + "the P2TR spend and its same-block child are eligible" + ); + } + + #[test] + fn far_watermark_does_not_advance() { + let _gate = crate::script_pool::steal_test_gate(); + let (_dir, q) = tiny_query_labeled("idx-gap"); + let params = ChainParams::regtest(); + let ms = Milestone::height(u32::MAX); + let genesis = genesis_block(¶ms); + accept_and_connect_block(&q, ¶ms, Height(0), &genesis, ms).unwrap(); + let mut prev = genesis.block_hash(); + let mut time = genesis.header.time; + for h in 1..=4 { + time += 600; + let block = mine_empty_regtest(prev, time, h); + accept_and_connect_block(&q, ¶ms, Height(h), &block, ms).unwrap(); + prev = block.block_hash(); + time = block.header.time; + } + enable(&q); + crate::index_writebehind::prepare_live_indexes_limited(&q, 2).unwrap(); + assert!(!q.index_live()); + assert_eq!(q.filter_index_next(), Some(0)); + time += 600; + let block = mine_empty_regtest(prev, time, 5); + confirm_wire_run(&q, ¶ms, ms, &[(Height(5), block)]).unwrap(); + assert_eq!(q.tip_height(), Some(Height(5))); + assert_eq!(q.filter_index_next(), Some(0)); + assert_eq!(q.tweak_index_next(), Some(0)); + } + + #[test] + fn startup_seals_a_short_gap_then_the_batch_appends() { + let _gate = crate::script_pool::steal_test_gate(); + let (_dir, q) = tiny_query_labeled("idx-repair"); + let params = ChainParams::regtest(); + let ms = Milestone::height(u32::MAX); + let genesis = genesis_block(¶ms); + accept_and_connect_block(&q, ¶ms, Height(0), &genesis, ms).unwrap(); + let b1 = mine_empty_regtest(genesis.block_hash(), genesis.header.time + 600, 1); + accept_and_connect_block(&q, ¶ms, Height(1), &b1, ms).unwrap(); + enable(&q); + go_live(&q); + assert!(q.index_live()); + assert_eq!(q.filter_index_next(), Some(2)); + let b2 = mine_empty_regtest(b1.block_hash(), b1.header.time + 600, 2); + confirm_wire_run(&q, ¶ms, ms, &[(Height(2), b2)]).unwrap(); + assert_sealed_matches_window(&q, 2); + } + + /// Filters still at genesis and tweaks at a later origin are two holes. + /// A limit that fits only the tweak hole does not turn live append on. + #[test] + fn startup_repairs_each_fitting_index() { + let _gate = crate::script_pool::steal_test_gate(); + let (_dir, q) = tiny_query_labeled("idx-split"); + let params = ChainParams::regtest(); + let ms = Milestone::height(u32::MAX); + let genesis = genesis_block(¶ms); + accept_and_connect_block(&q, ¶ms, Height(0), &genesis, ms).unwrap(); + let mut prev = genesis.block_hash(); + let mut time = genesis.header.time; + for h in 1..=3 { + time += 600; + let block = mine_empty_regtest(prev, time, h); + accept_and_connect_block(&q, ¶ms, Height(h), &block, ms).unwrap(); + prev = block.block_hash(); + time = block.header.time; + } + q.set_block_filter_index(true).unwrap(); + q.set_sptweaks_enabled(true, Height(2)).unwrap(); + crate::index_writebehind::prepare_live_indexes_limited(&q, 2).unwrap(); + assert!(!q.index_live()); + assert_eq!(q.filter_index_next(), Some(0)); + assert_eq!(q.tweak_index_next(), Some(4)); + time += 600; + let block = mine_empty_regtest(prev, time, 4); + confirm_wire_run(&q, ¶ms, ms, &[(Height(4), block)]).unwrap(); + assert_eq!(q.filter_index_next(), Some(0)); + assert_eq!(q.tweak_index_next(), Some(4)); + + let (_dir, both) = tiny_query_labeled("idx-split-both"); + accept_and_connect_block(&both, ¶ms, Height(0), &genesis, ms).unwrap(); + let b1 = mine_empty_regtest(genesis.block_hash(), genesis.header.time + 600, 1); + accept_and_connect_block(&both, ¶ms, Height(1), &b1, ms).unwrap(); + both.set_block_filter_index(true).unwrap(); + both.set_sptweaks_enabled(true, Height(0)).unwrap(); + go_live(&both); + assert!(both.index_live()); + assert_sealed_matches_window(&both, 1); + } + + #[test] + fn tip_entry_spawns_only_while_a_watermark_is_behind() { + let (_dir, q) = tiny_query_labeled("idx-tip"); + enable(&q); + go_live(&q); + let idle = crate::index_tip_entry(&q); + assert!(!idle.spawn); + assert!(!idle.advertise_filters); + + let params = ChainParams::regtest(); + let ms = Milestone::height(u32::MAX); + let genesis = genesis_block(¶ms); + accept_and_connect_block(&q, ¶ms, Height(0), &genesis, ms).unwrap(); + let caught = crate::index_tip_entry(&q); + assert!(!caught.spawn); + assert!(caught.advertise_filters); + + let (_dir, behind) = tiny_query_labeled("idx-tip-behind"); + accept_and_connect_block(&behind, ¶ms, Height(0), &genesis, ms).unwrap(); + enable(&behind); + let gap = crate::index_tip_entry(&behind); + assert!(gap.spawn); + assert!(!gap.advertise_filters); + } +} diff --git a/crates/rbitcoin-consensus/src/confirm_run/lookup.rs b/crates/rbitcoin-consensus/src/confirm_run/lookup.rs index 12780b325..1dab366ac 100644 --- a/crates/rbitcoin-consensus/src/confirm_run/lookup.rs +++ b/crates/rbitcoin-consensus/src/confirm_run/lookup.rs @@ -357,6 +357,7 @@ pub fn confirm_wire_load_from_plan( batch_parents, script_preverified: preverified.clone(), archive_plan: plan, + index_want: super::index::index_want(query), stats: query.confirm_stats_arc(), }, work_ns, diff --git a/crates/rbitcoin-consensus/src/confirm_run/mod.rs b/crates/rbitcoin-consensus/src/confirm_run/mod.rs index 4f8ea131c..3418a4ccc 100644 --- a/crates/rbitcoin-consensus/src/confirm_run/mod.rs +++ b/crates/rbitcoin-consensus/src/confirm_run/mod.rs @@ -45,6 +45,7 @@ use std::time::Instant; mod bq_resolve; mod head_drain; +mod index; mod lookup; mod phases; mod pin; @@ -170,6 +171,8 @@ pub struct LoadedBatch { script_preverified: ScriptPreverified, /// Planned Class A write from wire lookup/load (committed in write stage). pub archive_plan: Option, + /// Indexes to assemble. Copied from `index_live` at load. + index_want: index::IndexWant, stats: Arc, } @@ -181,6 +184,8 @@ pub struct ScriptOkBatch { wire_blocks: Vec>, batch_parents: rbitcoin_query::BatchParents, pub archive_plan: Option, + /// Filter bytes and tweak vecs for this batch. Empty when the indexes are off. + index_seal: index::IndexSeal, } /// Outcome of load: batch ready for scripts + pure work wall. @@ -195,6 +200,8 @@ pub struct ConfirmScriptOutcome { pub batch: ScriptOkBatch, /// Script verify only (when produced by [`confirm_scripts_phase`]). pub work_ns: u64, + /// Filter and tweak assemble on this stage. Not part of `work_ns`. + pub idx_asm_ns: u64, } /// LOAD STAGE from **raw wire blocks** (unified height-ordered pipeline). @@ -381,9 +388,18 @@ impl ScriptOkBatch { if self.archive_plan.is_some() != other.archive_plan.is_some() { return Err(other); } + let self_f = !self.index_seal.filters.is_empty(); + let other_f = !other.index_seal.filters.is_empty(); + if self_f != other_f { + return Err(other); + } self.prepared.append(&mut other.prepared); self.wire_blocks.append(&mut other.wire_blocks); self.batch_parents.extend_from(other.batch_parents); + self.index_seal + .filters + .append(&mut other.index_seal.filters); + self.index_seal.tweaks.append(&mut other.index_seal.tweaks); if let (Some(dst), Some(src)) = (self.archive_plan.as_mut(), other.archive_plan.take()) { dst.append(src); } diff --git a/crates/rbitcoin-consensus/src/confirm_run/scripts.rs b/crates/rbitcoin-consensus/src/confirm_run/scripts.rs index efd89129e..e01edabe2 100644 --- a/crates/rbitcoin-consensus/src/confirm_run/scripts.rs +++ b/crates/rbitcoin-consensus/src/confirm_run/scripts.rs @@ -44,16 +44,27 @@ fn take_script_jobs( jobs } -fn outcome_from(batch: LoadedBatch, work_ns: u64) -> ConfirmScriptOutcome { +fn outcome_from( + batch: LoadedBatch, + work_ns: u64, + seal: super::index::IndexSeal, + idx_asm_ns: u64, +) -> ConfirmScriptOutcome { rbitcoin_query::note_confirm(&batch.stats.script_ns, work_ns); + if idx_asm_ns > 0 { + rbitcoin_query::note_confirm(&batch.stats.idx_asm_ns, idx_asm_ns); + rbitcoin_query::note_confirm(&batch.stats.blockfilter_ns, idx_asm_ns); + } ConfirmScriptOutcome { batch: ScriptOkBatch { prepared: batch.prepared, wire_blocks: batch.wire_blocks, batch_parents: batch.batch_parents, archive_plan: batch.archive_plan, + index_seal: seal, }, work_ns, + idx_asm_ns, } } @@ -104,7 +115,11 @@ impl Inflight { if let Some(w) = self.wave { w.finish()?; } - Ok((outcome_from(self.batch, work_ns), self.meta)) + let (seal, idx_asm_ns) = super::index::live_index_seal(&self.batch)?; + Ok(( + outcome_from(self.batch, work_ns, seal, idx_asm_ns), + self.meta, + )) } } diff --git a/crates/rbitcoin-consensus/src/confirm_run/write.rs b/crates/rbitcoin-consensus/src/confirm_run/write.rs index 006a294a9..2111ddd30 100644 --- a/crates/rbitcoin-consensus/src/confirm_run/write.rs +++ b/crates/rbitcoin-consensus/src/confirm_run/write.rs @@ -128,7 +128,7 @@ fn apply_archive_plan( } /// COMMIT STAGE: optional Class A plan commit → structural → class_c → spend annotate → tip GC -/// → optional SP tweak index (**Tip write-through only**; Direct defers to backfill). +/// → live filter and tweak append when `index_live` assembled this batch. /// /// When `batch.archive_plan` is set (wire lookup/load path), Class A is appended in this /// same stage before structural/annotate — single ordered commit era. @@ -221,6 +221,13 @@ pub fn confirm_write_phase( ); } + let seal = std::mem::take(&mut batch.index_seal); + let idx_put_ns = super::index::seal_live_indexes(query, seal)?; + if idx_put_ns > 0 { + rbitcoin_query::note_confirm(&query.confirm_stats().idx_put_ns, idx_put_ns); + rbitcoin_query::note_confirm(&query.confirm_stats().blockfilter_ns, idx_put_ns); + } + let spend_ann_ns = post_commit(query, &slots)?; if let Some(tip) = query.tip_height() { let sync_ns = query diff --git a/crates/rbitcoin-consensus/src/confirm_run/write_idempotent_tests.rs b/crates/rbitcoin-consensus/src/confirm_run/write_idempotent_tests.rs index 7b4d345de..97eb02bc3 100644 --- a/crates/rbitcoin-consensus/src/confirm_run/write_idempotent_tests.rs +++ b/crates/rbitcoin-consensus/src/confirm_run/write_idempotent_tests.rs @@ -63,6 +63,7 @@ fn script_ok_append_contiguous_and_gap() { ))], batch_parents: rbitcoin_query::BatchParents::new(), archive_plan: None, + index_seal: super::index::IndexSeal::default(), } } let mut a = batch_one(10); @@ -85,6 +86,7 @@ fn script_ok_append_contiguous_and_gap() { wire_blocks: vec![], batch_parents: rbitcoin_query::BatchParents::new(), archive_plan: None, + index_seal: super::index::IndexSeal::default(), }) .is_ok()); assert_eq!(a.len(), 3); @@ -95,6 +97,7 @@ fn script_ok_append_contiguous_and_gap() { wire_blocks: vec![], batch_parents: rbitcoin_query::BatchParents::new(), archive_plan: None, + index_seal: super::index::IndexSeal::default(), }; assert!(empty.append_contiguous(batch_one(50)).is_ok()); assert_eq!(empty.len(), 1); @@ -174,6 +177,7 @@ fn three_stage_write_filter_and_scripts_surface() { batch_parents: rbitcoin_query::BatchParents::new(), script_preverified: ScriptPreverified::new(), archive_plan: None, + index_want: super::index::IndexWant::default(), stats: std::sync::Arc::new(rbitcoin_query::ConfirmStats::default()), }; assert!(batch.is_empty()); @@ -226,6 +230,7 @@ fn empty_loaded_batch() -> super::LoadedBatch { batch_parents: rbitcoin_query::BatchParents::new(), script_preverified: super::ScriptPreverified::new(), archive_plan: None, + index_want: super::index::IndexWant::default(), stats: std::sync::Arc::new(rbitcoin_query::ConfirmStats::default()), } } @@ -326,6 +331,7 @@ fn loaded_at( batch_parents: rbitcoin_query::BatchParents::new(), script_preverified: super::ScriptPreverified::new(), archive_plan: None, + index_want: super::index::IndexWant::default(), stats: std::sync::Arc::new(rbitcoin_query::ConfirmStats::default()), } } @@ -622,6 +628,7 @@ fn expected_bits_extending_height0_and_no_retarget() { batch_parents: rbitcoin_query::BatchParents::new(), script_preverified: ScriptPreverified::new(), archive_plan: None, + index_want: super::index::IndexWant::default(), stats: std::sync::Arc::new(rbitcoin_query::ConfirmStats::default()), }; let ok = confirm_scripts_phase(loaded).unwrap(); @@ -814,6 +821,7 @@ fn script_wave_skips_preverified_txids() { batch_parents: rbitcoin_query::BatchParents::new(), script_preverified: pre, archive_plan: None, + index_want: super::index::IndexWant::default(), stats: std::sync::Arc::new(rbitcoin_query::ConfirmStats::default()), }; confirm_scripts_phase(batch).expect("preverified skip avoids bad script fail"); diff --git a/crates/rbitcoin-consensus/src/index_rows.rs b/crates/rbitcoin-consensus/src/index_rows.rs new file mode 100644 index 000000000..8d1c19c40 --- /dev/null +++ b/crates/rbitcoin-consensus/src/index_rows.rs @@ -0,0 +1,125 @@ +//! One height-wave and one commit for filter and tweak rows. +//! +//! The confirm batch and the post-IBD window both produce [`IndexRows`]. +//! A commit that finds the watermark or `confirmed[h]` moved returns +//! [`IndexCommit::moved`]. The live path treats that as corrupt. The +//! materialize treats it as a re-plan. + +use std::sync::{Arc, Mutex}; +use std::time::Instant; + +use bitcoin::bip158::BlockFilter; +use rbitcoin_primitives::{Fk, Height}; +use rbitcoin_query::Query; + +use crate::script_pool::start_for_each_owned_chunk; +use crate::ConsensusError; + +/// Per-tx tweak (`None` = ineligible), block order. +pub(crate) type HeightTweaks = Vec>; + +/// Filter bytes and tweak vecs for a run of heights. Heights travel with the rows. +#[derive(Default)] +pub(crate) struct IndexRows { + pub filters: Vec<(Height, BlockFilter, Fk)>, + pub tweaks: Vec<(Height, Fk, HeightTweaks)>, +} + +pub(crate) struct IndexHeightOut { + pub filter: Option<(Height, BlockFilter, Fk)>, + pub tweaks: Option<(Height, Fk, HeightTweaks)>, +} + +/// Result of [`commit_index_rows`]. +pub(crate) struct IndexCommit { + pub put_ns: u64, + /// Watermark or `confirmed[h]` moved, or the rows leave a hole at `next`. + pub moved: bool, +} + +/// Claim one height per job. `f` writes its [`IndexHeightOut`] into the job. +pub(crate) fn finish_wave( + jobs: Vec, + f: fn(&T) -> Result<(), ConsensusError>, +) -> Result<(), ConsensusError> { + if let Some(wave) = start_for_each_owned_chunk(jobs, f, 1)? { + wave.finish()?; + } + Ok(()) +} + +pub(crate) fn rows_from_slots(slots: &[Arc>>]) -> IndexRows { + let mut rows = IndexRows::default(); + 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 { + rows.filters.push(filter); + } + if let Some(tweaks) = done.tweaks { + rows.tweaks.push(tweaks); + } + } + rows +} + +/// Commit each index from its watermark. Empty rows are a no-op. +pub(crate) fn commit_index_rows( + query: &Query, + rows: IndexRows, +) -> Result { + let t0 = Instant::now(); + let moved = commit_filters(query, rows.filters)? || commit_tweaks(query, rows.tweaks)?; + Ok(IndexCommit { + put_ns: t0.elapsed().as_nanos() as u64, + moved, + }) +} + +/// `true` when the rows do not start at `next`, or the put returned 0. +fn commit_filters( + query: &Query, + filters: Vec<(Height, BlockFilter, Fk)>, +) -> Result { + let Some(next) = query.filter_index_next() else { + return Ok(false); + }; + let Some(start) = filters.iter().position(|(h, _, _)| h.0 >= next) else { + return Ok(false); + }; + if filters[start..] + .iter() + .enumerate() + .any(|(i, (h, _, _))| h.0 != next.saturating_add(i as u32)) + { + return Ok(true); + } + let items: Vec<(BlockFilter, Fk)> = filters + .into_iter() + .skip(start) + .map(|(_, filter, fk)| (filter, fk)) + .collect(); + Ok(query.commit_window_filters(next, &items)? == 0) +} + +fn commit_tweaks( + query: &Query, + tweaks: Vec<(Height, Fk, HeightTweaks)>, +) -> Result { + let Some(next) = query.tweak_index_next() else { + return Ok(false); + }; + let Some(start) = tweaks.iter().position(|(h, _, _)| h.0 >= next) else { + return Ok(false); + }; + let slice = &tweaks[start..]; + if slice + .iter() + .enumerate() + .any(|(i, (h, _, _))| h.0 != next.saturating_add(i as u32)) + { + return Ok(true); + } + Ok(query.commit_window_tweaks(slice)? == 0) +} diff --git a/crates/rbitcoin-consensus/src/index_writebehind.rs b/crates/rbitcoin-consensus/src/index_writebehind.rs index 0c82edd0e..5ad6b3383 100644 --- a/crates/rbitcoin-consensus/src/index_writebehind.rs +++ b/crates/rbitcoin-consensus/src/index_writebehind.rs @@ -1,5 +1,9 @@ -//! Post-IBD block index write-behind (`rbtc-idx-wb`): BIP158 basic filters -//! and BIP-352 tweaks from one read of Class A. +//! Block index materialize (`rbtc-idx-wb`): BIP158 basic filters and BIP-352 +//! tweaks from one read of Class A when a watermark is behind the tip. +//! +//! [`prepare_live_indexes`] seals a restart gap of at most one IBD write +//! drain before confirm runs. A larger gap stays on this reader. Confirm +//! batches do not call it. //! //! The IO thread plans windows of consecutive heights from the lower index //! watermark up to the released tip and reads each through @@ -11,10 +15,12 @@ //! 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_chunk; +use crate::index_rows::{self, IndexHeightOut, IndexRows}; use crate::silent_payments::tweak_records_from_window; use crate::ConsensusError; -use rbitcoin_primitives::{Fk, Height}; +#[cfg(test)] +use rbitcoin_primitives::Fk; +use rbitcoin_primitives::Height; use rbitcoin_query::Query; use rbitcoin_store::IndexWindow; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; @@ -39,7 +45,106 @@ fn next_needed(query: &Query) -> Option { } } -/// Heights `start..` for the next window, capped at `target`. +/// Heights one IBD write drain can leave unsealed: write queue (14) times +/// the confirm run cap (144). +pub const INDEX_STARTUP_GAP_HEIGHTS: u32 = 14 * 144; + +/// Tip entry: spawn the materialize only while an enabled watermark is still +/// behind the tip and live append is off. Advertise filters immediately when +/// that watermark already covers the tip. +pub struct IndexTipEntry { + pub spawn: bool, + pub advertise_filters: bool, +} + +pub fn index_tip_entry(query: &Query) -> IndexTipEntry { + let tip = query.tip_height().map(|h| h.0); + let behind = |next: Option| next.is_some_and(|n| tip.is_some_and(|t| n <= t)); + let caught = |next: Option| next.is_some_and(|n| tip.is_some_and(|t| n > t)); + IndexTipEntry { + spawn: !query.index_live() + && (behind(query.filter_index_next()) || behind(query.tweak_index_next())), + advertise_filters: caught(query.filter_index_next()), + } +} + +/// Before confirm: seal each enabled index whose gap through the tip fits +/// [`INDEX_STARTUP_GAP_HEIGHTS`], then set [`Query::index_live`] when every +/// enabled index is contiguous with the tip. +pub fn prepare_live_indexes(query: &Query) -> Result<(), ConsensusError> { + prepare_live_indexes_limited(query, INDEX_STARTUP_GAP_HEIGHTS) +} + +pub(crate) fn prepare_live_indexes_limited( + query: &Query, + max_heights: u32, +) -> Result<(), ConsensusError> { + let tip_h = query.tip_height().map(|h| h.0); + if let Some(tip) = tip_h { + let repair_f = repair_from(query.filter_index_next(), tip, max_heights); + let repair_t = repair_from(query.tweak_index_next(), tip, max_heights); + if repair_f.is_some() || repair_t.is_some() { + seal_through(query, repair_f, repair_t, tip)?; + } + } + let tip_h = query.tip_height().map(|h| h.0); + let any_on = query.filter_index_next().is_some() || query.tweak_index_next().is_some(); + query.set_index_live( + any_on + && index_caught_up(query.filter_index_next(), tip_h) + && index_caught_up(query.tweak_index_next(), tip_h), + ); + Ok(()) +} + +fn repair_from(next: Option, tip: u32, max_heights: u32) -> Option { + let next = next?; + if next > tip { + return None; + } + let gap = tip.saturating_add(1).saturating_sub(next); + (gap <= max_heights).then_some(next) +} + +fn index_caught_up(next: Option, tip: Option) -> bool { + match (next, tip) { + (None, _) => true, + (Some(_), None) => true, + (Some(n), Some(t)) => n > t, + } +} + +/// Read and commit `start..=tip` for the indexes in `filter_from` / `tweak_from`. +/// `None` means that index is not part of this repair. +fn seal_through( + query: &Query, + filter_from: Option, + tweak_from: Option, + tip: u32, +) -> Result<(), ConsensusError> { + let mut start = match (filter_from, tweak_from) { + (Some(f), Some(t)) => f.min(t), + (Some(f), None) => f, + (None, Some(t)) => t, + (None, None) => return Ok(()), + }; + while start <= tip { + let end = plan_window(query, start, tip)?; + let window = Arc::new(read_window_from(query, start, end, tweak_from)?); + let rows = assemble_window_from(&window, filter_from, tweak_from)?; + if index_rows::commit_index_rows(query, rows)?.moved { + return Err(ConsensusError::Store(rbitcoin_store::StoreError::Corrupt( + "invariant: startup index repair", + ))); + } + if end >= tip { + break; + } + start = end + 1; + } + Ok(()) +} + fn plan_window(query: &Query, start: u32, target: u32) -> Result { let mut end = start; let mut txs = 0u64; @@ -62,7 +167,15 @@ fn plan_window(query: &Query, start: u32, target: u32) -> Result Result { - let tweaks_from = query.tweak_index_next(); + read_window_from(query, start, end, query.tweak_index_next()) +} + +fn read_window_from( + query: &Query, + start: u32, + end: u32, + tweaks_from: Option, +) -> Result { let heights = query.index_heights(start, end, tweaks_from)?; let mut window = query.read_index_window(&heights)?; for block in &mut window.blocks { @@ -71,14 +184,6 @@ fn read_window(query: &Query, start: u32, end: u32) -> Result>; - -struct IndexHeightOut { - filter: Option<(bitcoin::bip158::BlockFilter, Fk)>, - tweaks: Option<(Height, Fk, HeightTweaks)>, -} - struct IndexHeightJob { window: Arc, index: usize, @@ -91,6 +196,7 @@ fn assemble_index_height(job: &IndexHeightJob) -> Result<(), ConsensusError> { let block = &job.window.blocks[job.index]; let filter = if job.want_filter { Some(( + block.height, rbitcoin_query::basic_filter_of(&block.hash, &job.window, job.index)?, block.header_fk, )) @@ -110,16 +216,17 @@ fn assemble_index_height(job: &IndexHeightJob) -> Result<(), ConsensusError> { 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(); +fn assemble_window(query: &Query, window: &Arc) -> Result { + assemble_window_from(window, query.filter_index_next(), query.tweak_index_next()) +} + +fn assemble_window_from( + window: &Arc, + filter_from: Option, + tweak_from: Option, +) -> Result { let mut slots = Vec::with_capacity(window.blocks.len()); let mut jobs = Vec::new(); for (index, block) in window.blocks.iter().enumerate() { @@ -138,23 +245,8 @@ fn assemble_window(query: &Query, window: &Arc) -> Result, stages: &IndexStageMs, ) -> Result { - let Some(first) = window.blocks.first().map(|b| b.height.0) else { + if window.blocks.is_empty() { 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; + let committed = index_rows::commit_index_rows(query, assembled)?; + let commit_ns = committed.put_ns; stages.add_build(build_ns); stages.add_commit(commit_ns); rbitcoin_query::note_confirm( &query.confirm_stats().blockfilter_ns, build_ns.saturating_add(commit_ns), ); - committed.map(|ok| CommitTimes { - ok, + Ok(CommitTimes { + ok: !committed.moved, build_ns, commit_ns, }) @@ -626,8 +707,8 @@ mod tests { 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); + assert_eq!(got.filters[i].1.content, serial_f.content, "filter {i}"); + assert_eq!(got.filters[i].2, 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}"); } diff --git a/crates/rbitcoin-consensus/src/lib.rs b/crates/rbitcoin-consensus/src/lib.rs index 9db288cc8..8bd4de35f 100644 --- a/crates/rbitcoin-consensus/src/lib.rs +++ b/crates/rbitcoin-consensus/src/lib.rs @@ -6,6 +6,7 @@ mod confirm_run; mod convert; mod error; mod header; +mod index_rows; mod index_writebehind; mod milestone; mod params; @@ -71,7 +72,10 @@ pub use error::{block_reject_log_line, block_reject_reason, script_flag_paren, C pub use header::{ expected_next_bits, median_time_past, validate_header, validate_header_on_parent, }; -pub use index_writebehind::{build_indexes_released, spawn_index_writebehind}; +pub use index_writebehind::{ + build_indexes_released, index_tip_entry, prepare_live_indexes, spawn_index_writebehind, + IndexTipEntry, INDEX_STARTUP_GAP_HEIGHTS, +}; pub use milestone::{Milestone, MilestoneAnchor}; pub use params::{ default_milestone_height, genesis_block, mainnet_milestone_anchor, mainnet_min_chain_work_be, diff --git a/crates/rbitcoin-consensus/src/silent_payments.rs b/crates/rbitcoin-consensus/src/silent_payments.rs index 045261df9..758a41699 100644 --- a/crates/rbitcoin-consensus/src/silent_payments.rs +++ b/crates/rbitcoin-consensus/src/silent_payments.rs @@ -980,26 +980,17 @@ mod tests { let _ = std::fs::remove_dir_all(&dir); } - /// Tweaks come only from the index builder: each connect is sealed once - /// released, a disconnect truncates, and the replacement block is indexed. + /// A connect with tweaks on seals that height. Disconnect truncates, and + /// the replacement block is sealed on the next connect. #[test] - fn builder_writes_tweaks_and_reorg_truncates() { + fn live_connect_seals_tweaks_and_reorg_truncates() { let (dir, q) = tmp_store(); let params = ChainParams::regtest(); q.set_sptweaks_enabled(true, Height(0)).unwrap(); - let seal = |h: u32| { - q.release_index_writebehind(Height(h)); - crate::build_indexes_released(&q).unwrap(); - }; + crate::prepare_live_indexes(&q).unwrap(); let genesis = bitcoin::blockdata::constants::genesis_block(bitcoin::Network::Regtest); crate::accept_and_connect_block(&q, ¶ms, Height::GENESIS, &genesis, Milestone::NONE) .unwrap(); - assert_eq!( - q.sptweaks_next_height(), - Some(Height(0)), - "connect writes none" - ); - seal(0); assert_eq!(q.sptweaks_next_height(), Some(Height(1))); let thin0 = q.load_thin_tweaks(Height(0)).unwrap().expect("indexed"); assert!( @@ -1009,7 +1000,6 @@ mod tests { let b1 = crate::mine_empty_regtest(genesis.block_hash(), genesis.header.time + 600, 1); crate::accept_and_connect_block(&q, ¶ms, Height(1), &b1, Milestone::NONE).unwrap(); - seal(1); assert_eq!(q.sptweaks_next_height(), Some(Height(2))); q.disconnect_tip().unwrap(); @@ -1018,29 +1008,26 @@ mod tests { let b1b = crate::mine_empty_regtest(genesis.block_hash(), genesis.header.time + 601, 2); crate::accept_and_connect_block(&q, ¶ms, Height(1), &b1b, Milestone::NONE).unwrap(); - seal(1); assert_eq!(q.sptweaks_next_height(), Some(Height(2))); assert!(q.load_thin_tweaks(Height(1)).unwrap().is_some()); let _ = std::fs::remove_dir_all(&dir); } - /// Neither Direct nor Tip confirms write tweaks; one builder pass fills - /// origin..=tip. + /// Contiguous connects seal tweaks. The materialize does not rewrite them. #[test] - fn builder_fills_the_gap_then_follows_the_tip() { + fn live_connect_seals_tweaks_builder_leaves_them() { let (dir, q) = tmp_store(); let params = ChainParams::regtest(); - q.enter_direct_index_mode().unwrap(); q.set_sptweaks_enabled(true, Height(0)).unwrap(); + crate::prepare_live_indexes(&q).unwrap(); let genesis = bitcoin::blockdata::constants::genesis_block(bitcoin::Network::Regtest); crate::accept_and_connect_block(&q, ¶ms, Height::GENESIS, &genesis, Milestone::NONE) .unwrap(); let b1 = crate::mine_empty_regtest(genesis.block_hash(), genesis.header.time + 600, 1); crate::accept_and_connect_block(&q, ¶ms, Height(1), &b1, Milestone::NONE).unwrap(); - q.enter_tip_index_mode(); let b2 = crate::mine_empty_regtest(b1.block_hash(), b1.header.time + 600, 2); crate::accept_and_connect_block(&q, ¶ms, Height(2), &b2, Milestone::NONE).unwrap(); - assert_eq!(q.sptweaks_next_height(), Some(Height(0))); + assert_eq!(q.sptweaks_next_height(), Some(Height(3))); q.release_index_writebehind(Height(2)); crate::build_indexes_released(&q).unwrap(); diff --git a/crates/rbitcoin-net/src/ibd/confirm/mod.rs b/crates/rbitcoin-net/src/ibd/confirm/mod.rs index b1fa53835..9fd6463eb 100644 --- a/crates/rbitcoin-net/src/ibd/confirm/mod.rs +++ b/crates/rbitcoin-net/src/ibd/confirm/mod.rs @@ -1676,12 +1676,15 @@ pub(crate) fn spawn_confirm_engine( ); }, |outcome, meta| { - loop_stats_sc - .confirm_ns - .fetch_add(outcome.work_ns, Ordering::Relaxed); + loop_stats_sc.confirm_ns.fetch_add( + outcome.work_ns.saturating_add(outcome.idx_asm_ns), + Ordering::Relaxed, + ); confirm_thr_stats::add_script_work( &stats, - confirm_thr_stats::script_work_from_verify_ns(outcome.work_ns), + confirm_thr_stats::script_work_from_verify_ns( + outcome.work_ns.saturating_add(outcome.idx_asm_ns), + ), ); let script_ms = outcome.work_ns / 1_000_000; let mat_ms = meta.mat_ns / 1_000_000; diff --git a/crates/rbitcoin-net/src/ibd/confirm/tests.rs b/crates/rbitcoin-net/src/ibd/confirm/tests.rs index e70c49450..9c90bc6bb 100644 --- a/crates/rbitcoin-net/src/ibd/confirm/tests.rs +++ b/crates/rbitcoin-net/src/ibd/confirm/tests.rs @@ -2,8 +2,17 @@ use super::{ format_conf_q, format_queue_depth, format_stamp_reject_missing_prevout, - stamp_reject_operator_msg, ConfirmFeed, ConfirmQueueDepths, + stamp_reject_operator_msg, write_drain_max_parts, write_queue_cap, ConfirmFeed, + ConfirmQueueDepths, CONFIRM_RUN_MAX_BLOCKS, }; + +#[test] +fn index_startup_gap_matches_one_write_drain() { + assert_eq!( + rbitcoin_consensus::INDEX_STARTUP_GAP_HEIGHTS as usize, + write_drain_max_parts(write_queue_cap()) * CONFIRM_RUN_MAX_BLOCKS + ); +} use bitcoin::hashes::Hash; use bitcoin::BlockHash; use rbitcoin_primitives::Fk; diff --git a/crates/rbitcoin-net/src/ibd/perf_log.rs b/crates/rbitcoin-net/src/ibd/perf_log.rs index 5ecb6e955..4401b66d2 100644 --- a/crates/rbitcoin-net/src/ibd/perf_log.rs +++ b/crates/rbitcoin-net/src/ibd/perf_log.rs @@ -30,11 +30,12 @@ //! (plan=None / S0 only), clone, and post-stamp prune on a marked last load //! batch (`load_thr pack/stamp/pin/asm/prune`). //! - **script=** = `SCRIPT_NS` (publish → first `is_complete` per batch on -//! `ibd-confirm`; excludes head-of-line wait for write handoff). `thr script work` -//! is that same ns. Recv/send are wait. Publisher parks; it does not `wait_done` -//! on steal workers. +//! `ibd-confirm`; excludes head-of-line wait for write handoff and `idx_asm=`). +//! `idx_asm=` is filter and tweak assemble after that verify. `thr script work` +//! is `script=` plus `idx_asm=`. Recv/send are wait. Publisher parks; it does +//! not `wait_done` on steal workers. //! - **write** = Class A + ensure + structural + class_c + spend -//! + `pins=` / `head_sub=` / `drain_join=` / `dequeue=`. +//! + `pins=` / `head_sub=` / `drain_join=` / `dequeue=` / `idx_put=`. //! `other=` is write-thread work minus that inventory. //! //! **Inventory rule:** new work on lookup / load / scripts / write (or a sidecar @@ -97,6 +98,9 @@ pub(crate) struct WriteStageSample { /// Body-queue dequeue after confirm (`dequeue=`) pub dequeue_ms: u64, pub dequeue_ns: u64, + /// Live filter and tweak put (`idx_put=`) + pub idx_put_ms: u64, + pub idx_put_ns: u64, } type WriteInvTok = ( @@ -127,6 +131,7 @@ impl WriteStageSample { ("class_c_join", |s| s.class_c_join_ms, |s| s.class_c_join_ns), ("drain_join", |s| s.drain_join_ms, |s| s.drain_join_ns), ("dequeue", |s| s.dequeue_ms, |s| s.dequeue_ns), + ("idx_put", |s| s.idx_put_ms, |s| s.idx_put_ns), ]; /// Same inventory in nanoseconds (`format_debug` us/blk write=). @@ -169,6 +174,8 @@ pub(crate) struct IbdPerfSample { pub phase_blks: u64, pub connect_ms: u64, pub script_ms: u64, + /// Filter and tweak assemble on the scripts stage (`idx_asm=`). + pub idx_asm_ms: u64, /// Write-stage exclusive tokens (`write=` = [`WriteStageSample::stage_ms`]). pub write: WriteStageSample, /// Ensure mix: residency/pin hits vs cold denserels body loads. @@ -461,6 +468,7 @@ impl Default for IbdPerfSample { phase_blks: 0, connect_ms: 0, script_ms: 0, + idx_asm_ms: 0, write: WriteStageSample::default(), ensure_res_hit: 0, ensure_cold_n: 0, @@ -860,6 +868,8 @@ pub(crate) fn sample( let w = stats.take_window(); let connect_ns = w.connect_ns; let script_ns = w.script_ns; + let idx_asm_ns = w.idx_asm_ns; + let idx_put_ns = w.idx_put_ns; let milestone_gate_ns = w.milestone_gate_ns; let class_c_ns = w.class_c_ns; let strong_ns = w.strong_ns; @@ -949,6 +959,7 @@ pub(crate) fn sample( phase_blks, connect_ms: ns_ms(connect_ns), script_ms: ns_ms(script_ns), + idx_asm_ms: ns_ms(idx_asm_ns), write: WriteStageSample { class_a_ms: ns_ms(class_a_ns), class_a_ns, @@ -972,6 +983,8 @@ pub(crate) fn sample( drain_join_ns, dequeue_ms: ns_ms(dequeue_ns), dequeue_ns, + idx_put_ms: ns_ms(idx_put_ns), + idx_put_ns, }, ensure_res_hit, ensure_cold_n, @@ -1298,7 +1311,7 @@ pub(crate) fn format_info(s: &IbdPerfSample) -> String { let stamp_head_ms = s.stamp_batch_head_fk_ms; let stamp_pack_ms = s.thr_load_stamp_ms.saturating_sub(stamp_head_ms); out.push_str(&format!( - " | conf blks={} lookup={}ms load={}ms script={}ms(jobs={} skip={}) write={}ms \ + " | conf blks={} lookup={}ms load={}ms script={}ms(jobs={} skip={}) idx_asm={}ms write={}ms \ lookup_thr busy={}ms(claim={}ms wave={}ms(decode={}ms precompute={}ms collect={}ms head={}ms(probe={}ms io={}ms preads={}) loc={}ms) other={}ms send_w={}ms) \ load_thr busy/wait={}/{}ms(pack={}ms clone={}ms stamp={}ms(pack={}ms head={}ms) pin={}ms asm={}ms prune={}ms send_w={}ms) \ thr script={}/{}ms write={}/{}ms \ @@ -1309,6 +1322,7 @@ pub(crate) fn format_info(s: &IbdPerfSample) -> String { s.script_ms, s.script_jobs, s.script_skip, + s.idx_asm_ms, write_ms, thr_lookup_busy, s.thr_lookup_claim_ms, diff --git a/crates/rbitcoin-node/src/run.rs b/crates/rbitcoin-node/src/run.rs index 29c698401..dd2e45142 100644 --- a/crates/rbitcoin-node/src/run.rs +++ b/crates/rbitcoin-node/src/run.rs @@ -213,6 +213,8 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { ); } apply_startup_index_mode(&handle.query, &config, params.taproot_height())?; + rbitcoin_consensus::prepare_live_indexes(&handle.query) + .map_err(|e| crate::error::NodeError::Init(format!("index startup repair failed: {e}")))?; let bind = config.listen.start_p2p_bind(config.network); let start_tip = handle.query.tip_height().map(|h| h.0).unwrap_or(0); @@ -596,7 +598,13 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { tip_follow_ready = gates.tip_follow_ready; sh_tip_ready = gates.sh_tip_ready; if tip_follow_ready && !shutdown.requested() { - if config.block_filter_index || config.sptweaks { + let index_entry = rbitcoin_consensus::index_tip_entry(&node.hub.query); + if index_entry.advertise_filters { + info!("blockfilter: already at tip; advertising NODE_COMPACT_FILTERS"); + rbitcoin_net::set_compact_filters_service(true); + } + if index_entry.spawn { + let advertise_later = !index_entry.advertise_filters; index_writebehind = Some(rbitcoin_consensus::spawn_index_writebehind( Arc::clone(&node.hub.query), Arc::clone(&shutdown.flag), @@ -606,7 +614,11 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { }, // Version carries the bit once per peer: advertise only // when served filters reach the tip, not during materialize. - || rbitcoin_net::set_compact_filters_service(true), + move || { + if advertise_later { + rbitcoin_net::set_compact_filters_service(true); + } + }, )); } if relay_while_following( diff --git a/crates/rbitcoin-query/src/block_filter.rs b/crates/rbitcoin-query/src/block_filter.rs index a25935d4c..0ea4628a9 100644 --- a/crates/rbitcoin-query/src/block_filter.rs +++ b/crates/rbitcoin-query/src/block_filter.rs @@ -52,30 +52,40 @@ 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. +/// No store IO. 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); - } - } - } + let mut spent = Vec::new(); 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); + spent.push(out.script.as_slice()); } - encode_basic_filter(block_hash, elements.into_iter()) + let outputs = block + .txs + .iter() + .flat_map(|tx| tx.outs.iter().map(|o| o.script.as_slice())); + basic_filter_from_scripts(block_hash, outputs, spent) +} + +/// GCS-encode a basic filter. Non-`OP_RETURN` output scripts, then each spent +/// prevout script. `block_hash` is internal byte order. +pub fn basic_filter_from_scripts<'a>( + block_hash: &[u8; 32], + output_scripts: impl IntoIterator, + spent_scripts: impl IntoIterator, +) -> Result { + const OP_RETURN: u8 = 0x6a; + let elements = output_scripts + .into_iter() + .filter(|script| script.first() != Some(&OP_RETURN)) + .chain(spent_scripts); + encode_basic_filter(block_hash, elements) } /// GCS-encode a basic filter keyed by `block_hash` (internal byte order). @@ -307,6 +317,16 @@ impl Query { Ok(built.len() as u32) } + /// Confirm may append filter and tweak rows for batches it connects. + pub fn index_live(&self) -> bool { + self.index_live.load(Ordering::Acquire) + } + + /// The write side of [`Self::index_live`]. Startup sets it; confirm does not. + pub fn set_index_live(&self, on: bool) { + self.index_live.store(on, Ordering::Release); + } + /// Next filter height to seal (`None` when the filter index is off). pub fn filter_index_next(&self) -> Option { self.block_filter_table().map(|t| t.next_height().0) diff --git a/crates/rbitcoin-query/src/confirm_stats.rs b/crates/rbitcoin-query/src/confirm_stats.rs index dd3c9caac..cb9f3c92f 100644 --- a/crates/rbitcoin-query/src/confirm_stats.rs +++ b/crates/rbitcoin-query/src/confirm_stats.rs @@ -140,8 +140,10 @@ confirm_window! { structural_create_h_ns, structural_bip68_ns, class_c_ns, - // `rbtc-bf-wb` build + commit (off the write thread) + // Live idx assemble + put, and materialize build + commit. `tip: accept bf=` is this sum. blockfilter_ns, + idx_asm_ns, + idx_put_ns, ensure_layout_ns, write_class_c_join_ns, write_drain_join_ns, diff --git a/crates/rbitcoin-query/src/lib.rs b/crates/rbitcoin-query/src/lib.rs index e6645afc4..3ebdaf09a 100644 --- a/crates/rbitcoin-query/src/lib.rs +++ b/crates/rbitcoin-query/src/lib.rs @@ -168,7 +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 block_filter::{basic_filter_from_scripts, basic_filter_of}; pub use catchup::IndexMode; pub use chain_view::{ChainView, ChainViewKind}; pub use confirm_load::SpendEdges; @@ -309,6 +309,10 @@ pub struct Query { /// SH collect/enqueue/durable write-through entirely (tip follow independent). sh_index_enabled: std::sync::atomic::AtomicBool, block_filter_enabled: std::sync::atomic::AtomicBool, + /// Confirm appends filter and tweak rows. Set by `prepare_live_indexes` + /// when every enabled index is contiguous with the tip. Load reads it; + /// the write thread is the only writer. + index_live: std::sync::atomic::AtomicBool, /// BIP158 basic filter table, opened when the index is first turned on. block_filters: std::sync::OnceLock, bf_wb: block_filter::BlockFilterWriteBehind, @@ -462,6 +466,7 @@ impl Query { // `--shindex` off before entering Direct. sh_index_enabled: std::sync::atomic::AtomicBool::new(true), block_filter_enabled: std::sync::atomic::AtomicBool::new(false), + index_live: std::sync::atomic::AtomicBool::new(false), block_filters: std::sync::OnceLock::new(), bf_wb: block_filter::BlockFilterWriteBehind::new(), sp_tweaks: Mutex::new(sp_tweaks), diff --git a/crates/rbitcoin-test/tests/integration_multinode.rs b/crates/rbitcoin-test/tests/integration_multinode.rs index d42aadfa9..97d2b1a06 100644 --- a/crates/rbitcoin-test/tests/integration_multinode.rs +++ b/crates/rbitcoin-test/tests/integration_multinode.rs @@ -48,6 +48,10 @@ async fn start_node_inbound(dir: &TempDir, max_inbound: usize) -> P2PNode { } /// Mature regtest pad with basic filters sealed through height 1 only. +/// +/// Connecting the chain releases the index through the tip. Commit 0..=1 +/// directly and leave live append off, so a later connect does not seal +/// the rest. fn open_padded_query(dir: &TempDir) -> Query { use rbitcoin_consensus::accept_and_connect_block; use rbitcoin_test::pad_empty_from; @@ -56,11 +60,27 @@ fn open_padded_query(dir: &TempDir) -> Query { let params = ChainParams::regtest(); let genesis = regtest_genesis(); accept_and_connect_block(&q, ¶ms, Height::GENESIS, &genesis, Milestone::NONE).unwrap(); - let (tip, time) = pad_empty_from(&q, ¶ms, genesis.block_hash(), genesis.header.time, 1, 1); + pad_empty_from( + &q, + ¶ms, + genesis.block_hash(), + genesis.header.time, + 1, + params.coinbase_maturity() + 1, + ); q.set_block_filter_index(true).unwrap(); - q.release_index_writebehind(Height(1)); - rbitcoin_consensus::build_indexes_released(&q).expect("filters through height 1"); - pad_empty_from(&q, ¶ms, tip, time, 2, params.coinbase_maturity() + 1); + let heights = q.index_heights(0, 1, None).unwrap(); + let window = q.read_index_window(&heights).unwrap(); + let built: Vec<_> = (0..heights.len()) + .map(|i| { + ( + q.basic_filter_from_window(&window, i).unwrap(), + heights[i].header_fk, + ) + }) + .collect(); + assert_eq!(q.commit_window_filters(0, &built).unwrap(), 2); + assert_eq!(q.basic_filter_hwm().unwrap(), Some(1)); q } @@ -1751,8 +1771,8 @@ async fn ibd_skips_dead_peer() { } /// After IBD, seed announces a new tip; follower picks it up via inv/headers. -/// With the filter index on, IBD confirm builds no basic filters; the -/// write-behind appender materializes them, then follows the new tip. +/// Filters and tweaks enabled before sync are sealed by IBD confirm. A +/// caught-up follower seals the new tip on the confirm path. #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn tip_follow_after_ibd() { use std::sync::Arc; @@ -1769,19 +1789,20 @@ async fn tip_follow_after_ibd() { peer.query .set_sptweaks_enabled(true, Height(ChainParams::regtest().taproot_height())) .unwrap(); + rbitcoin_consensus::prepare_live_indexes(&peer.query).unwrap(); sync_ibd(&peer, seed.local_addr).await; peer.wait_height(5, Duration::from_secs(10)) .await .expect("ibd"); assert_eq!( peer.query.basic_filter_hwm().unwrap(), - None, - "IBD confirm leaves basic filters to the appender" + Some(5), + "IBD confirm seals filters when the index is on from the start" ); assert_eq!( peer.query.sptweaks_next_height(), - Some(Height(0)), - "the confirm write thread writes no tweaks" + Some(Height(6)), + "IBD confirm seals tweaks from the taproot origin" ); let pq = Arc::clone(&peer.query); // Records the filter watermark when a builder first reports caught up. @@ -1839,29 +1860,19 @@ async fn tip_follow_after_ibd() { assert_eq!(peer.query.tip_height(), Some(Height(6))); assert_eq!( indexed(), - (Some(5), Some(Height(6))), - "with no builder running, the confirm path writes no index data" + (Some(6), Some(Height(7))), + "a caught-up confirm seals the new tip" ); - let (stop, builder) = spawn_builder(); - wait_ms_until( - 5_000, - || indexed() == (Some(6), Some(Height(7))), - || format!("catch up {:?}", indexed()), - ) - .await; let h7 = mine_next(7); peer.wait_tip_hash(h7, Duration::from_secs(10)) .await .expect("tip follow"); - wait_ms_until( - 5_000, - || indexed() == (Some(7), Some(Height(8))), - || format!("follow {:?}", indexed()), - ) - .await; - stop.store(true, std::sync::atomic::Ordering::SeqCst); - builder.join().unwrap(); + assert_eq!( + indexed(), + (Some(7), Some(Height(8))), + "the next tip is sealed on the confirm path too" + ); seed.shutdown().await; peer.shutdown().await; diff --git a/docs/concurrency.md b/docs/concurrency.md index e9159537e..c6fab5c95 100644 --- a/docs/concurrency.md +++ b/docs/concurrency.md @@ -10,7 +10,7 @@ Short map of who may write which tables. **Format is unstable until 1.0.** | Confirm **lookup** | 1 OS thread | pack/hold from BQ stamped Σ inputs; emit: decode BQ raw + TipOnly `head_fk` + **`take_raw` onto loadq** (`Arc` + `TxPrecompute`). Does **not** plan_batch / structure | | Confirm **load** | 1 OS thread | `confirm_wire_lookup_stamp` (structure with **loadq `pres`**, plan_batch binds carried BQ keys, skeleton bind) + pin `txout` + assemble | | Confirm **scripts** | 1 OS thread (`ibd-confirm`) publishes waves + `rbtc-scripts` steal (lock-free claim) | **none** — pure CPU | -| Confirm **write** | 1 OS thread (`ibd-confirm-write`) + 1 process-wide `ibd-confirm-head` | **sole Class A appender** (`txout`+`seqsigwit`+`spent`; IBD encodes ins from `Arc` + SpendEdges) + structural + Class C + spend annotate on **`spent.body`** + tip GC; **`block_queue_dequeue_height`**. `tx.head` write-behind insert runs on **`ibd-confirm-head`** overlapping structural + Class C (not a per-batch spawn). Class A **never leads tip** (same commit era; no archive-ahead DONTNEED) | +| Confirm **write** | 1 OS thread (`ibd-confirm-write`) + 1 process-wide `ibd-confirm-head` | **sole Class A appender** (`txout`+`seqsigwit`+`spent`; IBD encodes ins from `Arc` + SpendEdges) + structural + Class C + spend annotate on **`spent.body`** + tip GC; **`block_queue_dequeue_height`**. `tx.head` write-behind insert runs on **`ibd-confirm-head`** overlapping structural + Class C (not a per-batch spawn). Class A **never leads tip** (same commit era; no archive-ahead DONTNEED). When `index_live` is set, this thread appends that index after `confirmed[h]` is set (one put per confirm batch) and does not read Class A for it. Startup `prepare_live_indexes` seals a restart gap of at most one write drain (14 × 144 heights) before confirm; a larger gap is left for `rbtc-idx-wb` | | IBD main loop | 1 tokio task | none (orchestration only). Event drain + confirm-offer every turn; getdata assign ≤50 ms (immediate if inflight empty); main-loop `getheaders`/`locator_hashes` ≤500 ms (empty path immediate); peer stall/relative-slow and work-path hygiene ≤1 s. Full-batch header continuation stays on the Headers event. | **IoSession TLS:** one completion session per OS thread (`with_thread_local`). Harvest / poison / drain / do-not-flatten: [`io-modality.md`](./io-modality.md). `RBITCOIN_IO=pread` disables the session. SH k-way merge submits 256 KiB ahead preads on that TLS session and waits only when promote needs a page that has not completed. @@ -21,7 +21,7 @@ Short map of who may write which tables. **Format is unstable until 1.0.** **IBD lookup resolve wave:** TipOnly `head_fk` bounded by remaining loadq slots × load pack (safety cap **1080** heights / soft **64000** inputs; include-overshoot; 256 k unique-key safety cap). Decode parks as resolved BQ rows; `taken_hi` bumps **per load-batch send**. Unsent tail stays on the BQ (no re-decode). Hard **min 8000** inputs per wave when more unresolved heights can still join — including `ready=0` / load-frontier / unknown window. Last available thin wave still emits. Each load chunk carries a `BatchParentIds` skeleton (shared wave ids/spent + per-chunk need-vouts). Same-wave creates are omitted from TipOnly need. When `ready >` half the 1-min BQ window, lookup waits for a full wave instead of minting a 1-block layer — unless the first unresolved height is within `path_lo + win/2` **and** the collect is already ≥8000 inputs (load is about to claim it; O(1) from the already-sorted unresolved list). Pack/hold uses enqueue-stamped `n_inputs` (peer CompactSize walk). `wave_intake` does **not** clone raw payloads; decode clones `raw_payload(height)` only on emit. Hold is `decode=0` / `precompute=0`. `lookup_thr wave=(decode= precompute= collect= head= loc=)`: `precompute=` is `from_tx_wire` on payload slices (`sighash` midstates only when scripts run); `head=` is TipOnly `get_fk_by_txid_batch`, not load stamp; `loc=` is create.loc fill inside that TipOnly. -**`ibd: perf`:** `load=` is pin+assemble only. Load OS-thread `stamp=` nests `pack=` (plan HashMap) vs `head=` (leftover TipOnly `prep_head_fk_ns`). IBD skeleton keeps `head=` ~0. `stamp_sub struct_txid=` must be **0** on IBD (loadq `pres`). Non-zero means load dropped lookup hashes and `from_tx`'d again. Post-stamp in-flight drop on the last load batch of a wave is `prune=` (height-index remove; IBD has no pstore Weak walk). Lookup-wave `decode=` / `precompute=` / `collect=` / `head=` / `spent=` nest under `lookup_thr wave=`. Pin names `thin=` for the vout-map prefix. `script=` is per-batch wave wall on `ibd-confirm` (`jobs=`/`skip=`). Steal claim is an `ArcSwap` snapshot + `fetch_add(32)` (no `WAVES` mutex per job; `in_wave` counts in-flight chunks). The scripts thread publishes another `scriptq` batch when steal is empty (up to 4 in-flight, matching `scriptq`) and parks until a worker unparks it on wave complete (load also unparks on `scriptq` send). `script=` stamps when the wave first reports complete, not at write-queue pop. Write snapshot-drains the whole **writeq** (cap 14) so one `flush_class_c_tip` covers the queued run. `append_contiguous` stops when `archive_plan` polarity differs (`Some` vs `None`); leftover is the next meta-batch. `ready>0` + `scriptq=1` + high `stamp=` `head=` means leftover on load, not a hungry script pool. High `stamp=` `pack=` with `head=0` is HashMap CPU on load (leave it there; lookup is already the wave). Do not retune steal or borrow `rbtc-scripts-*` for decode. Process restart leftover is empty RAM identity (horizon does not survive). +**`ibd: perf`:** `load=` is pin+assemble only. Load OS-thread `stamp=` nests `pack=` (plan HashMap) vs `head=` (leftover TipOnly `prep_head_fk_ns`). IBD skeleton keeps `head=` ~0. `stamp_sub struct_txid=` must be **0** on IBD (loadq `pres`). Non-zero means load dropped lookup hashes and `from_tx`'d again. Post-stamp in-flight drop on the last load batch of a wave is `prune=` (height-index remove; IBD has no pstore Weak walk). Lookup-wave `decode=` / `precompute=` / `collect=` / `head=` / `spent=` nest under `lookup_thr wave=`. Pin names `thin=` for the vout-map prefix. `script=` is per-batch wave wall on `ibd-confirm` (`jobs=`/`skip=`). `idx_asm=` is filter and tweak assemble on that stage after script verify (one height per claim) and is not inside `script=`. `idx_put=` is the write-thread index append and is inside `write=`. Steal claim is an `ArcSwap` snapshot + `fetch_add(32)` (no `WAVES` mutex per job; `in_wave` counts in-flight chunks). The scripts thread publishes another `scriptq` batch when steal is empty (up to 4 in-flight, matching `scriptq`) and parks until a worker unparks it on wave complete (load also unparks on `scriptq` send). `script=` stamps when the wave first reports complete, not at write-queue pop. Write snapshot-drains the whole **writeq** (cap 14) so one `flush_class_c_tip` covers the queued run. `append_contiguous` stops when `archive_plan` polarity differs (`Some` vs `None`); leftover is the next meta-batch. `ready>0` + `scriptq=1` + high `stamp=` `head=` means leftover on load, not a hungry script pool. High `stamp=` `pack=` with `head=0` is HashMap CPU on load (leave it there; lookup is already the wave). Do not retune steal or borrow `rbtc-scripts-*` for decode. Process restart leftover is empty RAM identity (horizon does not survive). **Tip follow / reorg:** peer reconstructs on the session task, then offers the body to the process-wide **`tip-accept`** OS thread (bounded queue 8, oneshot wait). That thread is the sole production owner of `connect_lock` / `accept_and_connect_block_preverified` / `confirm_wire_run_preverified` at tip (1-block, or one `accept_branch` run). Scripts still publish via `start_for_each_owned` → `rbtc-scripts-*` steal — same as IBD’s publisher, not a second interpreter. Disconnect keeps Class A archive; re-extension always supplies **wire** from the peer, not hash-only load. **IBD most-work reorg** rewinds the tip to the LCA (`rewind_to_height`) from the **IBD orchestration task only** — never from confirm lookup/load/scripts/write threads — then the confirm pipeline connects the winner as a linear extension. `accept_branch` remains for tip-follow / held-body apply (also on `tip-accept`). Selector / rewind / invalid-heavy: [`architecture.md`](./architecture.md#most-work-chain-selection-ibd--tip-follow). @@ -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). 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. | +| `rbtc-idx-wb` / `rbtc-idx-cpu` | **One** block index builder for BIP158 basic filters (`--block-filter-index`) and BIP-352 tweaks (`--sp-tweaks`), spawned at tip entry only when an enabled watermark is still behind the tip after startup repair. Live batches append on the confirm write thread; this reader does not walk a chain that was indexed as it connected. 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 | diff --git a/docs/ibd-memory.md b/docs/ibd-memory.md index 40816c1e2..3ecacc3d3 100644 --- a/docs/ibd-memory.md +++ b/docs/ibd-memory.md @@ -69,7 +69,8 @@ Not page cache. Caps on **decoded `Block` objects and live outbound sessions**: | **getheaders continuation** | full 2000-header reply locates from last hash | Next batch after that hash, not a replay from our tip. | | **headers poll** | skip if `best_known` cannot beat our tip | 120s `getheaders` only for peers that can still add work. | | **Chainwork prefix** | `Vec` `prefix[h] = work through h` (~32 B × tip; ≈28–32 MiB at 900k) | Process cache. Extend/truncate to `query.tip_height()`. Not durable. Restart rebuilds on first `chain_work`. | -| **Block index windows** | ≤ 3 windows (one being read, ≤2 queued) of ≤64 heights / ≤50k creates: block outputs, input edges, P2TR-output witnesses, spent parents (tens of MiB each at mainnet sizes) | Only while `rbtc-idx-wb` builds filters / tweaks. A committed window is dropped. | +| **Block index windows** | ≤ 3 windows (one being read, ≤2 queued) of ≤64 heights / ≤50k creates: block outputs, input edges, P2TR-output witnesses, spent parents (tens of MiB each at mainnet sizes) | Only while `rbtc-idx-wb` builds a watermark that is behind the tip. A committed window is dropped. | +| **Confirm-batch index bytes** | One confirm batch of BIP158 filter bytes and BIP-352 tweak vecs while `index_live` is set | Drop with the batch at write. Not a cache. | | **Fee history** | ≤1008 `(height, p10)` entries (~16 KiB) | Read from the chain per connected block; backfilled over the newest 1008 blocks when relay turns on (`txstat` + `spent` span reads, no bodies). Not a cache: every entry is a chain fact. | | **Mempool fee snapshot** | Published Arc (chunks + live count/vsize/total_fee) | Dirty/singleflight ≤~1 s. Admit only marks dirty. `GET /mempool` Arc-loads; no graph walk, no body clones. | | **Mempool tx-body snapshot** | Lazy; ≤ one extra live-pool of `Arc` + JSON `OnceLock` after first unix `/internal` mempool-tx page | Dirty/singleflight. Not FIFO/LRU. Operators who never hit unix `/internal` do not keep this. | diff --git a/docs/invariants.md b/docs/invariants.md index 3bd7461a1..31167fb8e 100644 --- a/docs/invariants.md +++ b/docs/invariants.md @@ -49,9 +49,13 @@ wire / body-queue → ensure abs (holes only: same-batch after Class A / missing stamp; post-condition: every spend has abs) → structural spentness (pin abs bulk pread of spent.body; multi-list protocol cold only) → Class C tip + → index append (filter and tweak bytes from the scripts stage; no store read) → abs spend annotate (put_spend_batch_by_abs_meta on spent.body only) ``` +Startup, before this pipeline, may `read_index_window` to seal a filter or tweak +gap of at most one IBD write drain. The confirm write thread does not. + IBD thread split (same IO table): lookup **thread** is decode + TipOnly + `take_raw` onto loadq. Structure + plan_batch (`confirm_wire_lookup_stamp`) run on the **load** thread and consume `LoadBatch.pres` — they do not read diff --git a/docs/operator/operations.md b/docs/operator/operations.md index cb1d68a86..c8671844e 100644 --- a/docs/operator/operations.md +++ b/docs/operator/operations.md @@ -90,11 +90,11 @@ 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 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 | +| `--block-filter-index` | `block_filter_index=` | **off** — BIP158 basic filters. Independent of `--sh-index`. IBD seals them when the flag is on from the start. A later enable still materializes after catch-up: the `rbtc-idx-wb` builder fills 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`). A caught-up chain seals new blocks on the confirm write thread; 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. | -| `--sp-tweaks` | `sp_tweaks=` | **off** — thin BIP-352 tweak index (`sp_tweaks.*`), built after catch-up by `rbtc-idx-wb` (with block filters, one read of Class A) from the taproot origin, then sealed per tip block off the confirm path. Heights not yet sealed are served by the naive walk. Refused with `--prune-seqsigwit`: tweaks read input keys from scriptSig and witness | +| `--sp-tweaks` | `sp_tweaks=` | **off** — thin BIP-352 tweak index (`sp_tweaks.*`). IBD seals it from the taproot origin when the flag is on from the start. A later enable still materializes after catch-up on `rbtc-idx-wb` (with block filters, one read of Class A). New tips are sealed on the confirm write thread once the watermark has caught the chain. Heights not yet sealed are served by the naive walk. Refused with `--prune-seqsigwit`: tweaks read input keys from scriptSig and witness | | `--sp-tweaks-dust SATS` | `sp_tweaks_dust=` | **1000** — omit served P2TR outs with `value <= SATS` (`0` = serve all; **546** matches Cake electrs) | | `--electrum-listen [ADDR]` | `electrum_listen=` | disabled; omit ADDR → `127.0.0.1:50001`. Address/scripthash methods need `--sh-index` | | `--esplora-listen [ADDR\|PATH]` | `esplora_listen=` | disabled (Esplora REST); omit ADDR → `127.0.0.1:3000`; a filesystem path is unix HTTP (mode **0660**, dummy `Host: api` is fine). Address/scripthash methods need `--sh-index` | diff --git a/docs/operator/storage.md b/docs/operator/storage.md index dfbf72604..e24d0013b 100644 --- a/docs/operator/storage.md +++ b/docs/operator/storage.md @@ -264,10 +264,12 @@ with `value <= SATS` and drops txs that then have none. Default **1000**. not Core dust: P2TR at 1 sat/vB is about **330** sats; 546 is the P2PKH figure Cake’s server used. The index is unchanged — only the Electrum JSON. -The index is built after catch-up by the block index builder (`rbtc-idx-wb`, -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 +IBD seals the index from the taproot origin (709632 on mainnet) when +`--sp-tweaks` is on from the start, on the confirm write thread, shared with +`--block-filter-index`. A restart gap of at most one write drain is sealed +at startup; a later enable still materializes after catch-up +(`rbtc-idx-wb`). Once the watermark covers the tip, new blocks are sealed on +the confirm write thread. That materialize is one IO thread reading windows of heights on one completion session (`seqsigwit` and parent txids only for P2TR-output txs); one CPU thread (`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 + diff --git a/docs/rpc.md b/docs/rpc.md index 83e521bd5..58706c463 100644 --- a/docs/rpc.md +++ b/docs/rpc.md @@ -27,7 +27,7 @@ history. | `--rpc-listen [ADDR]` / conf `rpc_listen=` | **off** | TCP JSON-RPC; omit ADDR → `127.0.0.1` and Core-matching port (8332 / 18332 / 38332 / 18443). Implies `--rpc`. | | `--rpc-token-file PATH` | `{datadir}/rpc.token` | CSPRNG hex token; TCP `Authorization: Bearer` | | `--sh-index` | **off** | Class B scripthash (Electrum/Esplora only; RPC by height/hash/txid does not need it) | -| `--block-filter-index` | **off** | BIP158 basic, built after catch-up by a write-behind appender (not during IBD). `NODE_COMPACT_FILTERS` is advertised once filters first reach the tip (`getnetworkinfo` lists `COMPACT_FILTERS`), then for the life of the process. `getblockfilter` and `/rest/blockfilter/` serve heights the watermark already covers. Independent of `--sh-index` | +| `--block-filter-index` | **off** | BIP158 basic. IBD seals them when the flag is on from the start. A later enable still materializes after catch-up. `NODE_COMPACT_FILTERS` is advertised once filters first reach the tip (`getnetworkinfo` lists `COMPACT_FILTERS`), then for the life of the process. `getblockfilter` and `/rest/blockfilter/` serve heights the watermark already covers. Independent of `--sh-index` | | `--rpc-work-queue N` | **16** | In-flight HTTP RPC (Core `-rpcworkqueue`). One POST is one slot (a JSON-RPC array is still one slot). Full permit is HTTP **503** `Work queue depth exceeded`. **0** is the default queue of 16. | TLS is external (reverse proxy). Unix socket needs no HTTP header. TCP is