diff --git a/CHANGELOG.md b/CHANGELOG.md index d234ad35..fe406270 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -54,10 +54,10 @@ before 1.0). ### Changed -- **IBD Class A ins at write:** `confirm_wire_lookup_stamp` drops packed - `InputRecord`s after plan (SpendEdges + CreatePin remain). Write encodes - ins from `Arc` + those edges. CreatePin outs stay stamp-time for - in-flight. +- **IBD Class A ins at write:** `confirm_wire_lookup_stamp` plans from + wire (`archive_plan_batch_from_wire`) without `TxApply`. Packed ins stay + empty; SpendEdges + CreatePin remain. Write encodes ins from `Arc` + + those edges. CreatePin outs stay stamp-time for in-flight. - **Algo-review open list:** [`docs/algo-review.md`](docs/algo-review.md) is remaining work only. Landed items from that pass live here and in diff --git a/crates/rbitcoin-consensus/src/confirm_run/lookup.rs b/crates/rbitcoin-consensus/src/confirm_run/lookup.rs index 7855725d..8f9c8460 100644 --- a/crates/rbitcoin-consensus/src/confirm_run/lookup.rs +++ b/crates/rbitcoin-consensus/src/confirm_run/lookup.rs @@ -109,12 +109,6 @@ pub fn confirm_wire_lookup_stamp( Some(p) => ParentPinStamp::take_from_plan(p), None => stamp_parent_pin_archived(query, params, &metas, &wire_blocks, ifo)?, }; - if let Some(ref mut p) = plan { - for (_, ins) in &mut p.packed { - ins.clear(); - ins.shrink_to_fit(); - } - } lookup_stage_stats::BLOCKS.fetch_add(blocks.len() as u64, Ordering::Relaxed); lookup_stage_stats::HEAD_NS.fetch_add(plan_ns, Ordering::Relaxed); let work_ns = t0.elapsed().as_nanos() as u64; @@ -310,11 +304,6 @@ pub(super) fn wire_lookup_phase( } } - let mut with_fk: Vec<( - rbitcoin_primitives::Fk, - rbitcoin_store::HeaderRecord, - Vec, - )> = Vec::with_capacity(blocks.len()); let mut wire_blocks: Vec> = Vec::with_capacity(blocks.len()); let mut metas: Vec = Vec::with_capacity(blocks.len()); @@ -377,8 +366,20 @@ pub(super) fn wire_lookup_phase( struct_ns = struct_ns.saturating_add(t_struct.elapsed().as_nanos() as u64); let t_prep = Instant::now(); - let (header_rec, txs) = - crate::prepare_block_for_archive_with_txids(query, block.as_ref(), &txids)?; + let prev_fk = if i == 0 { + if block.header.prev_blockhash.to_byte_array() == [0u8; 32] { + rbitcoin_primitives::Fk::NULL + } else { + query + .get_header_by_hash(block.header.prev_blockhash.as_byte_array()) + .map_err(ConsensusError::from)? + .map(|(fk, _)| fk) + .ok_or(ConsensusError::BadPrev)? + } + } else { + metas[i - 1].header_fk + }; + let header_rec = crate::header_to_record(prev_fk, &block.header); let header_fk = if let Some((fk, _)) = query .get_header_by_hash(&header_rec.hash) .map_err(ConsensusError::from)? @@ -391,7 +392,6 @@ pub(super) fn wire_lookup_phase( .map_err(ConsensusError::from)? }; prepare_ns = prepare_ns.saturating_add(t_prep.elapsed().as_nanos() as u64); - with_fk.push((header_fk, header_rec.clone(), txs)); wire_blocks.push(block); metas.push(BodyMeta { height: *height, @@ -405,12 +405,13 @@ pub(super) fn wire_lookup_phase( } let t_filter = Instant::now(); - let (_header_fks, mut need) = query - .archive_filter_need_bodies(&mut with_fk) + let header_fks: Vec = metas.iter().map(|m| m.header_fk).collect(); + let need_fks = query + .archive_filter_need_header_fks(&header_fks) .map_err(ConsensusError::from)?; let filter_ns = t_filter.elapsed().as_nanos() as u64; let t_batch = Instant::now(); - let plan = if need.is_empty() { + let plan = if need_fks.is_empty() { for (i, m) in metas.iter_mut().enumerate() { if let Some(list) = query .store() @@ -432,18 +433,28 @@ pub(super) fn wire_lookup_phase( } None } else { + let mut need = Vec::with_capacity(need_fks.len()); + for fk in &need_fks { + let i = metas + .iter() + .position(|m| m.header_fk == *fk) + .ok_or(ConsensusError::Store(rbitcoin_store::StoreError::Corrupt( + "invariant: need-body header_fk not in batch", + )))?; + need.push((*fk, wire_blocks[i].as_ref(), metas[i].txids.as_slice())); + } let plan = match pipeline { Some(p) => query - .archive_plan_batch_from_store( - &mut need, + .archive_plan_batch_from_wire( + &need, p.next_tx_start.max(1), &p.in_flight, Some(p.published.as_ref()), ) .map_err(ConsensusError::from)?, None => query - .archive_plan_batch_from_store( - &mut need, + .archive_plan_batch_from_wire( + &need, query.tx_body_count().saturating_add(1).max(1), &rbitcoin_query::InFlightView::empty(), None, @@ -891,6 +902,16 @@ mod tests { .expect("coinbase-only stamp"); assert_eq!(stamped.metas[0].pres.len(), pres.len()); assert_eq!(stamped.metas[0].pres[0].txid, pres[0].txid); + let plan = stamped.plan.as_ref().expect("new body plans"); + assert!( + plan.packed.iter().all(|(_, ins)| ins.is_empty()), + "wire planner packed ins stay empty" + ); + let want = items[0].1.txdata[0].output[0].script_pubkey.to_bytes(); + assert_eq!( + plan.batch_pin[0].1[0].script, want, + "CreatePin outs from wire script_pubkey" + ); let _ = std::fs::remove_dir_all(&path); } } diff --git a/crates/rbitcoin-consensus/src/confirm_run/mod.rs b/crates/rbitcoin-consensus/src/confirm_run/mod.rs index 7c2c73b8..b5ddcd11 100644 --- a/crates/rbitcoin-consensus/src/confirm_run/mod.rs +++ b/crates/rbitcoin-consensus/src/confirm_run/mod.rs @@ -141,7 +141,7 @@ pub struct WireLoadPipeline { pub path_lo: u32, /// Parent of `path_lo` when ahead of store tip (last wire hash of prior loaded batch). pub parent_hash: Option<[u8; 32]>, - /// Inclusive create-fk start for [`Query::archive_plan_batch_from_store`]. + /// Inclusive create-fk start for [`Query::archive_plan_batch_from_wire`]. pub next_tx_start: u64, /// Prior uncommitted packs: immutable layer snapshot (no shared mutable map). /// @@ -248,11 +248,6 @@ pub fn confirm_wire_load_phase_pipelined( let mut ns_header = 0u64; let mut ns_prepare = 0u64; - let mut with_fk: Vec<( - rbitcoin_primitives::Fk, - rbitcoin_store::HeaderRecord, - Vec, - )> = Vec::with_capacity(blocks.len()); let mut wire_blocks: Vec> = Vec::with_capacity(blocks.len()); let mut metas: Vec = Vec::with_capacity(blocks.len()); @@ -319,8 +314,20 @@ pub fn confirm_wire_load_phase_pipelined( ns_header = ns_header.saturating_add(t.elapsed().as_nanos() as u64); let t = Instant::now(); - let (header_rec, txs) = - crate::prepare_block_for_archive_with_txids(query, block.as_ref(), &txids)?; + let prev_fk = if i == 0 { + if block.header.prev_blockhash.to_byte_array() == [0u8; 32] { + rbitcoin_primitives::Fk::NULL + } else { + query + .get_header_by_hash(block.header.prev_blockhash.as_byte_array()) + .map_err(ConsensusError::from)? + .map(|(fk, _)| fk) + .ok_or(ConsensusError::BadPrev)? + } + } else { + metas[i - 1].header_fk + }; + let header_rec = crate::header_to_record(prev_fk, &block.header); ns_prepare = ns_prepare.saturating_add(t.elapsed().as_nanos() as u64); let t = Instant::now(); let header_fk = if let Some((fk, _)) = query @@ -335,7 +342,6 @@ pub fn confirm_wire_load_phase_pipelined( .map_err(ConsensusError::from)? }; ns_header = ns_header.saturating_add(t.elapsed().as_nanos() as u64); - with_fk.push((header_fk, header_rec.clone(), txs)); wire_blocks.push(block); metas.push(BodyMeta { height: *height, @@ -349,10 +355,11 @@ pub fn confirm_wire_load_phase_pipelined( } let t_fp = Instant::now(); - let (_header_fks, mut need) = query - .archive_filter_need_bodies(&mut with_fk) + let header_fks: Vec = metas.iter().map(|m| m.header_fk).collect(); + let need_fks = query + .archive_filter_need_header_fks(&header_fks) .map_err(ConsensusError::from)?; - let mut plan = if need.is_empty() { + let mut plan = if need_fks.is_empty() { for (i, m) in metas.iter_mut().enumerate() { if let Some(list) = query .store() @@ -374,18 +381,28 @@ pub fn confirm_wire_load_phase_pipelined( } None } else { + let mut need = Vec::with_capacity(need_fks.len()); + for fk in &need_fks { + let i = metas + .iter() + .position(|m| m.header_fk == *fk) + .ok_or(ConsensusError::Store(rbitcoin_store::StoreError::Corrupt( + "invariant: need-body header_fk not in batch", + )))?; + need.push((*fk, wire_blocks[i].as_ref(), metas[i].txids.as_slice())); + } let plan = match pipeline { Some(p) => query - .archive_plan_batch_from_store( - &mut need, + .archive_plan_batch_from_wire( + &need, p.next_tx_start.max(1), &p.in_flight, Some(p.published.as_ref()), ) .map_err(ConsensusError::from)?, None => query - .archive_plan_batch_from_store( - &mut need, + .archive_plan_batch_from_wire( + &need, query.tx_body_count().saturating_add(1).max(1), &rbitcoin_query::InFlightView::empty(), None, diff --git a/crates/rbitcoin-consensus/src/confirm_run/pin.rs b/crates/rbitcoin-consensus/src/confirm_run/pin.rs index 92715b5b..843c3b95 100644 --- a/crates/rbitcoin-consensus/src/confirm_run/pin.rs +++ b/crates/rbitcoin-consensus/src/confirm_run/pin.rs @@ -56,55 +56,23 @@ pub(super) fn pin_for_wire_batch( } } let fill_vouts = parent_pin.parent_vouts.is_empty(); - if !plan.edges.is_empty() { - spend_edges = plan.edges.clone(); - if fill_vouts { - for eds in plan.edges.values() { - for e in eds { - let Some(pid) = e.create_fk.get() else { - continue; - }; - if plan.create_in_spend_header(e.spend_fk, pid) { - continue; - } - parent_vouts.entry(pid).or_default().push(e.vout); - } - } - } - } else { - for ((_pin, ins), fk) in plan.packed.iter().zip(plan.planned_fks.iter()) { - let Some(sid) = fk.get() else { continue }; - let mut edges = Vec::with_capacity(ins.len()); - for inp in ins { - if inp.is_coinbase() || inp.prev_index == u32::MAX { - edges.push(rbitcoin_query::SpendEdge { - prev_txid: [0u8; 32], - vout: u32::MAX, - spend_fk: rbitcoin_primitives::Fk(sid), - create_fk: rbitcoin_primitives::Fk::NULL, - }); + if plan.edges.is_empty() && !plan.planned_fks.is_empty() { + return Err(ConsensusError::Store(StoreError::Corrupt( + "invariant: plan spend edges empty", + ))); + } + spend_edges = plan.edges.clone(); + if fill_vouts { + for eds in plan.edges.values() { + for e in eds { + let Some(pid) = e.create_fk.get() else { + continue; + }; + if plan.create_in_spend_header(e.spend_fk, pid) { continue; } - if let Some(pid) = inp.create_fk.get() { - edges.push(rbitcoin_query::SpendEdge { - prev_txid: inp.prev_txid, - vout: inp.prev_index, - spend_fk: rbitcoin_primitives::Fk(sid), - create_fk: inp.create_fk, - }); - if fill_vouts && !plan.create_in_spend_header(*fk, pid) { - parent_vouts.entry(pid).or_default().push(inp.prev_index); - } - } else { - edges.push(rbitcoin_query::SpendEdge { - prev_txid: inp.prev_txid, - vout: inp.prev_index, - spend_fk: rbitcoin_primitives::Fk(sid), - create_fk: rbitcoin_primitives::Fk::NULL, - }); - } + parent_vouts.entry(pid).or_default().push(e.vout); } - spend_edges.insert(sid, edges); } } if !fill_vouts { 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 c13685f0..5fdea53f 100644 --- a/crates/rbitcoin-consensus/src/confirm_run/write_idempotent_tests.rs +++ b/crates/rbitcoin-consensus/src/confirm_run/write_idempotent_tests.rs @@ -1063,6 +1063,36 @@ fn tiny_query() -> (std::path::PathBuf, rbitcoin_query::Query) { (path, q) } +fn fill_edges_from_packed(plan: &mut rbitcoin_query::ArchiveWritePlan) { + use rbitcoin_primitives::Fk; + use rbitcoin_query::SpendEdge; + if !plan.edges.is_empty() { + return; + } + for ((_, ins), fk) in plan.packed.iter().zip(plan.planned_fks.iter()) { + let Some(sid) = fk.get() else { continue }; + let mut edges = Vec::with_capacity(ins.len()); + for inp in ins { + if inp.is_coinbase() || inp.prev_index == u32::MAX { + edges.push(SpendEdge { + prev_txid: [0u8; 32], + vout: u32::MAX, + spend_fk: *fk, + create_fk: Fk::NULL, + }); + } else { + edges.push(SpendEdge { + prev_txid: inp.prev_txid, + vout: inp.prev_index, + spend_fk: *fk, + create_fk: inp.create_fk, + }); + } + } + plan.edges.insert(sid, edges); + } +} + fn rec_tx(b: u8, n_out: u32) -> rbitcoin_store::TxRecord { use rbitcoin_primitives::Fk; rbitcoin_store::TxRecord { @@ -1103,6 +1133,7 @@ fn pin_and_ensure_journey() { )]; plan.planned_fks = vec![Fk(1)]; let mut stamp = ParentPinStamp::take_from_plan(&mut plan); + fill_edges_from_packed(&mut plan); let err = pin_for_wire_batch(&q, Some(&plan), &mut stamp, &[], &[], None, None) .expect_err("missing parent must hard-fail pin"); let msg = format!("{err}"); @@ -1170,6 +1201,7 @@ fn pin_and_ensure_journey() { rbitcoin_query::ParentIdent::with_body(parent_tx.txid, range), ); let mut stamp = ParentPinStamp::take_from_plan(&mut plan); + fill_edges_from_packed(&mut plan); let (parents, _, _) = pin_for_wire_batch(&q, Some(&plan), &mut stamp, &[], &[], None, None) .expect("pin via stamped range"); assert!(parents.contains(pfk)); @@ -1191,6 +1223,7 @@ fn pin_and_ensure_journey() { rbitcoin_query::ParentIdent::with_body(parent_tx.txid, range), ); let mut empty_stamp = ParentPinStamp::default(); + fill_edges_from_packed(&mut plan2); let err = pin_for_wire_batch(&q, Some(&plan2), &mut empty_stamp, &[], &[], None, None) .expect_err("plan maps must not backfill an empty stamp"); assert!(err.to_string().contains("lookup stage miss"), "got: {err}"); @@ -1242,6 +1275,7 @@ fn pin_and_ensure_journey() { }, ); let mut stamp3 = ParentPinStamp::take_from_plan(&mut plan3); + fill_edges_from_packed(&mut plan3); let (mut parents3, _, _) = pin_for_wire_batch(&q, Some(&plan3), &mut stamp3, &[], &[], None, None).unwrap(); assert!( @@ -1275,6 +1309,7 @@ fn pin_and_ensure_journey() { ]; plan4.planned_fks = vec![Fk(2), Fk(3)]; let mut stamp4 = ParentPinStamp::take_from_plan(&mut plan4); + fill_edges_from_packed(&mut plan4); let (parents4, _, _) = pin_for_wire_batch(&q, Some(&plan4), &mut stamp4, &[], &[], None, None).unwrap(); assert!( @@ -1366,6 +1401,7 @@ fn pin_for_wire_cold_range_then_adopt_skips_body_io() { let store = Arc::new(PipelineParentStore::new()); let mut plan = stamp_plan(); let mut parent_pin = ParentPinStamp::take_from_plan(&mut plan); + fill_edges_from_packed(&mut plan); let (parents, _thin, _warm) = pin_for_wire_batch( &q, Some(&plan), @@ -1384,6 +1420,7 @@ fn pin_for_wire_cold_range_then_adopt_skips_body_io() { let mut plan2 = stamp_plan(); let mut parent_pin2 = ParentPinStamp::take_from_plan(&mut plan2); + fill_edges_from_packed(&mut plan2); let (parents2, _thin2, _warm2) = pin_for_wire_batch( &q, Some(&plan2), @@ -1486,6 +1523,7 @@ fn pin_for_wire_incomplete_outs_is_invariant_error() { let ifo = log.snapshot(); let mut parent_pin = ParentPinStamp::take_from_plan(&mut plan); + fill_edges_from_packed(&mut plan); let err = pin_for_wire_batch(&q, Some(&plan), &mut parent_pin, &[], &[], Some(&ifo), None) .expect_err("incomplete outs must hard-fail pin"); let msg = format!("{err}"); @@ -1618,6 +1656,7 @@ fn pin_takes_stamp_parent_vouts() { stamp.parent_vouts.get(&parent_id).map(|v| v.as_slice()), Some(&[0u32][..]) ); + fill_edges_from_packed(&mut plan); let (parents, _, _) = pin_for_wire_batch(&q, Some(&plan), &mut stamp, &[], &[], None, None) .expect("pin via taken vouts"); assert!(stamp.parent_vouts.is_empty(), "pin must take stamp vouts"); @@ -1695,6 +1734,7 @@ fn pin_for_wire_create_pin_shares_script_bytes() { plan.per_header_ranges = vec![(Fk(10), Fk(1), 1), (Fk(11), Fk(2), 1)]; plan.external_parent_vouts.insert(1, vec![0]); let mut stamp = ParentPinStamp::take_from_plan(&mut plan); + fill_edges_from_packed(&mut plan); let (parents, edges, _) = pin_for_wire_batch(&q, Some(&plan), &mut stamp, &[], &[], None, None) .expect("cross-height CreatePin pin"); let child_edges = edges.get(&2).expect("child spend edges"); @@ -1787,6 +1827,7 @@ fn pin_plan_edges_without_packed_ins() { }], ); let mut stamp = ParentPinStamp::take_from_plan(&mut plan); + fill_edges_from_packed(&mut plan); let (parents, edges, _) = pin_for_wire_batch(&q, Some(&plan), &mut stamp, &[], &[], None, None) .expect("pin from plan.edges with empty packed ins"); let child_edges = edges.get(&2).expect("child spend edges"); @@ -1796,6 +1837,60 @@ fn pin_plan_edges_without_packed_ins() { let _ = std::fs::remove_dir_all(&path); } +#[test] +fn pin_plan_empty_edges_is_invariant() { + use super::{pin_for_wire_batch, ParentPinStamp}; + use rbitcoin_primitives::Fk; + use rbitcoin_query::{ArchiveWritePlan, Query}; + use rbitcoin_store::{OutputRecord, TxRecord}; + use std::sync::Arc; + use std::sync::Once; + + static ONCE: Once = Once::new(); + ONCE.call_once(|| { + if std::env::var_os("RBITCOIN_HEAD_SCALE").is_none() { + std::env::set_var("RBITCOIN_HEAD_SCALE", "tiny"); + } + }); + let path = std::env::temp_dir().join(format!( + "rbitcoin-pin-empty-edges-{}-{}", + std::process::id(), + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_nanos()) + .unwrap_or(0) + )); + let _ = std::fs::remove_dir_all(&path); + std::fs::create_dir_all(&path).unwrap(); + let q = Query::open_or_create(&path).unwrap(); + q.enter_direct_index_mode().unwrap(); + let pin = Arc::new(( + TxRecord { + txid: [0x11u8; 32], + version: 1, + locktime: 0, + input_start_fk: Fk::NULL, + input_count: 1, + output_start_fk: Fk::NULL, + output_count: 1, + }, + vec![OutputRecord::unspent(50, vec![0x51])], + )); + let mut plan = ArchiveWritePlan::empty(); + plan.packed = vec![(Arc::clone(&pin), vec![])]; + plan.planned_fks = vec![Fk(1)]; + plan.batch_pin = vec![Arc::clone(&pin)]; + let mut stamp = ParentPinStamp::take_from_plan(&mut plan); + let err = pin_for_wire_batch(&q, Some(&plan), &mut stamp, &[], &[], None, None) + .expect_err("empty edges with planned fks must not skip spends"); + let msg = format!("{err}"); + assert!( + msg.contains("invariant") && msg.contains("spend edges empty"), + "unexpected err: {msg}" + ); + let _ = std::fs::remove_dir_all(&path); +} + /// Need a high vout from a multi-out parent (need-vouts only, not full n_out). #[test] fn pin_sparse_need_high_vout_only() { @@ -1886,6 +1981,7 @@ fn pin_sparse_need_high_vout_only() { body_est: 0, }; let mut parent_pin = ParentPinStamp::take_from_plan(&mut plan); + fill_edges_from_packed(&mut plan); let (parents, _thin, _warm) = pin_for_wire_batch(&q, Some(&plan), &mut parent_pin, &[], &[], None, None) .expect("pin high vout"); @@ -2007,6 +2103,7 @@ fn pin_range_fill_does_not_count_as_cache_hit() { } let mut parent_pin = ParentPinStamp::take_from_plan(&mut plan); + fill_edges_from_packed(&mut plan); let (_parents, _thin, warm) = pin_for_wire_batch( &q, Some(&plan), @@ -2109,6 +2206,7 @@ fn pin_recent_outs_is_cache_not_new() { "pin must use stamp-carried CreatePin" ); q.recent_creates().drop_from(0); + fill_edges_from_packed(&mut plan); let (_parents, _thin, warm) = pin_for_wire_batch(&q, Some(&plan), &mut parent_pin, &[], &[], None, None) .expect("stamp-carried outs must cover without a live RecentCreates ring"); @@ -2200,6 +2298,7 @@ fn pin_recent_identity_without_outs_still_range_fills() { } let mut parent_pin = ParentPinStamp::take_from_plan(&mut plan); + fill_edges_from_packed(&mut plan); let (_parents, _thin, warm) = pin_for_wire_batch(&q, Some(&plan), &mut parent_pin, &[], &[], None, None) .expect("identity-only recent still range-fills"); diff --git a/crates/rbitcoin-query/src/archive.rs b/crates/rbitcoin-query/src/archive.rs index e1c6881f..7f6808e1 100644 --- a/crates/rbitcoin-query/src/archive.rs +++ b/crates/rbitcoin-query/src/archive.rs @@ -1,7 +1,8 @@ //! Class A archive write path. //! //! Split for IBD dual-thread (prep/write may overlap with a small plan queue): -//! - **Plan** ([`Query::archive_plan_batch_from_store`]): +//! - **Plan** ([`Query::archive_plan_batch_from_wire`] IBD; +//! [`Query::archive_plan_batch_from_store`] TxApply tests): //! store **reads** — assign create fks (optionally from a reserved HWM), //! in-flight planned creates + `tx.head` resolve, stamp inputs. //! Head-miss parents use **fk-only** head resolve (no denserels on plan stamp); @@ -43,7 +44,7 @@ pub fn create_pin_approx_bytes(pin: &CreatePin) -> usize { #[derive(Debug)] pub struct ArchiveWritePlan { /// Body-append rows: shared [`CreatePin`] (tx + outs) + inputs. - /// IBD stamp clears ins after plan; write fills from wire + [`Self::edges`]. + /// IBD wire planner leaves ins empty; write fills from wire + [`Self::edges`]. /// Outs live once in the pin Arc (not duplicated alongside inputs). pub packed: Vec<(CreatePin, Vec)>, pub planned_fks: Vec, @@ -175,7 +176,7 @@ impl ArchiveWritePlan { .planned_fks .iter() .position(|f| *f == first) - .unwrap_or(0); + .ok_or(StoreError::Corrupt("invariant: retain first fk missing"))?; let end = start.saturating_add(n as usize).min(self.planned_fks.len()); for f in &self.planned_fks[start..end] { if let Some(id) = f.get() { @@ -253,6 +254,72 @@ impl ArchiveWritePlan { } } +struct PlanIn { + prev_txid: [u8; 32], + prev_index: u32, + is_coinbase: bool, +} + +struct PlanRow { + tx_fk: Fk, + tx: TxRecord, + ins: Vec, + outs: Vec, + packed_ins: Vec, + ins_est: u64, +} + +fn plan_in_from_record(inp: &InputRecord) -> PlanIn { + PlanIn { + prev_txid: inp.prev_txid, + prev_index: inp.prev_index, + is_coinbase: inp.is_coinbase(), + } +} + +fn plan_in_from_txin(inp: &bitcoin::TxIn) -> PlanIn { + use bitcoin::hashes::Hash; + let is_coinbase = inp.previous_output.is_null() + || (inp.previous_output.txid.to_byte_array() == [0u8; 32] + && inp.previous_output.vout == u32::MAX); + PlanIn { + prev_txid: inp.previous_output.txid.to_byte_array(), + prev_index: if is_coinbase { + u32::MAX + } else { + inp.previous_output.vout + }, + is_coinbase, + } +} + +fn wire_ins_est(tx: &bitcoin::Transaction) -> u64 { + tx.input + .iter() + .map(|inp| { + (1 + 8 + + 9 + + 4 + + 9 + + inp.script_sig.len() + + 9 + + inp.witness.iter().map(|w| 9 + w.len()).sum::()) as u64 + }) + .sum() +} + +fn tx_record_from_wire(tx: &bitcoin::Transaction, txid: [u8; 32]) -> TxRecord { + TxRecord { + txid, + version: tx.version.0, + locktime: tx.lock_time.to_consensus_u32(), + input_start_fk: Fk::NULL, + input_count: tx.input.len() as u32, + output_start_fk: Fk::NULL, + output_count: tx.output.len() as u32, + } +} + impl Query { /// Class A only (header + bodies + `tx.head` / `header_txs`). Does **not** /// set tip / fence / strong. @@ -366,6 +433,22 @@ impl Query { Ok((header_fks, need)) } + /// Header-only need-body filter (IBD wire planner). No [`TxApply`]. + pub fn archive_filter_need_header_fks(&self, header_fks: &[Fk]) -> Result, QueryError> { + let mut need = Vec::with_capacity(header_fks.len()); + let mut seen_headers = crate::FkSet::default(); + for &fk in header_fks { + if !seen_headers.insert(fk) { + continue; + } + if self.store.header_txs.has_body(fk)? { + continue; + } + need.push(fk); + } + Ok(need) + } + /// **Prep / read path:** assign create fks, identity resolve, stamp /// inputs. No Class A body/head writes (those are [`Self::archive_commit_plan`]). /// @@ -399,11 +482,8 @@ impl Query { let t_assign = Instant::now(); let mut batch_map: HashMap<[u8; 32], Fk> = HashMap::new(); - let mut work: Vec<(Fk, TxRecord, Vec, Vec)> = Vec::new(); + let mut work: Vec = Vec::new(); let mut per_header_ranges: Vec<(Fk, Fk, u32)> = Vec::with_capacity(need.len()); - let mut spends: Vec<([u8; 32], u32, Fk, u32)> = Vec::new(); - let archive_spends = self.spend_index_enabled() && self.index_mode().is_tip(); - let index_tx = self.tx_index_enabled(); for (header_fk, txs) in need.iter_mut() { if txs.is_empty() { @@ -433,17 +513,139 @@ impl Query { tx.output_count = n_out; batch_map.insert(tx.txid, tx_fk); - work.push((tx_fk, tx, ta.inputs, ta.outputs)); + let ins: Vec = ta.inputs.iter().map(plan_in_from_record).collect(); + let ins_est = ta.inputs.iter().map(|x| x.encoded_len() as u64).sum(); + work.push(PlanRow { + tx_fk, + tx, + ins, + outs: ta.outputs, + packed_ins: ta.inputs, + ins_est, + }); + } + per_header_ranges.push((*header_fk, first_tx_fk, n_txs)); + } + let assign_ns = t_assign.elapsed().as_nanos() as u64; + self.finish_archive_plan( + work, + batch_map, + per_header_ranges, + n_headers, + assign_ns, + in_flight, + published, + ) + } + + /// IBD stamp: CreatePin + SpendEdges from wire txs. Packed ins stay empty. + /// + /// Does not build [`TxApply`] / clone `script_sig` / witness. Write fills + /// packed ins from `Arc` + edges. `body_est` uses wire compact sizes. + pub fn archive_plan_batch_from_wire( + &self, + need: &[(Fk, &bitcoin::Block, &[[u8; 32]])], + next_tx_start: u64, + in_flight: &crate::InFlightView, + published: Option<&crate::PublishedIds>, + ) -> Result { + use std::collections::{HashMap, HashSet}; + use std::time::Instant; + + self.on_load_pack()?; + if need.is_empty() { + return Ok(ArchiveWritePlan::empty()); + } + + let mut next_tx = next_tx_start.max(1); + let n_headers = need.iter().filter(|(_, b, _)| !b.txdata.is_empty()).count() as u64; + + let t_assign = Instant::now(); + let mut batch_map: HashMap<[u8; 32], Fk> = HashMap::new(); + let mut work: Vec = Vec::new(); + let mut per_header_ranges: Vec<(Fk, Fk, u32)> = Vec::with_capacity(need.len()); + + for (header_fk, block, txids) in need { + if block.txdata.is_empty() { + continue; + } + if block.txdata.len() != txids.len() { + return Err(StoreError::Corrupt("txid count mismatch").into()); + } + let first_tx_fk = Fk(next_tx); + let n_txs = block.txdata.len() as u32; + let mut seen_in_block: HashSet<[u8; 32]> = HashSet::with_capacity(block.txdata.len()); + for (tx, txid) in block.txdata.iter().zip(txids.iter()) { + if !seen_in_block.insert(*txid) { + return Err(StoreError::Corrupt( + "duplicate txid in block body (consensus violation)", + ) + .into()); + } + let tx_fk = Fk(next_tx); + next_tx += 1; + let rec = tx_record_from_wire(tx, *txid); + batch_map.insert(*txid, tx_fk); + let ins: Vec = tx.input.iter().map(plan_in_from_txin).collect(); + let outs: Vec = tx + .output + .iter() + .map(|o| { + OutputRecord::unspent(o.value.to_sat() as i64, o.script_pubkey.to_bytes()) + }) + .collect(); + work.push(PlanRow { + tx_fk, + tx: rec, + ins, + outs, + packed_ins: Vec::new(), + ins_est: wire_ins_est(tx), + }); } per_header_ranges.push((*header_fk, first_tx_fk, n_txs)); } let assign_ns = t_assign.elapsed().as_nanos() as u64; + self.finish_archive_plan( + work, + batch_map, + per_header_ranges, + n_headers, + assign_ns, + in_flight, + published, + ) + } + + fn finish_archive_plan( + &self, + work: Vec, + batch_map: std::collections::HashMap<[u8; 32], Fk>, + per_header_ranges: Vec<(Fk, Fk, u32)>, + n_headers: u64, + assign_ns: u64, + in_flight: &crate::InFlightView, + published: Option<&crate::PublishedIds>, + ) -> Result { + use std::collections::HashSet; + use std::time::Instant; + + let mut spends: Vec<([u8; 32], u32, Fk, u32)> = Vec::new(); + let archive_spends = self.spend_index_enabled() && self.index_mode().is_tip(); + let index_tx = self.tx_index_enabled(); let t_collect = Instant::now(); let mut need_external: HashSet<[u8; 32]> = HashSet::new(); - for (_sfk, _tx, inputs, _) in &work { - for inp in inputs { - if inp.is_coinbase() || !inp.create_fk.is_null() { + for row in &work { + for (i, inp) in row.ins.iter().enumerate() { + if inp.is_coinbase { + continue; + } + if row + .packed_ins + .get(i) + .is_some_and(|r| !r.create_fk.is_null()) + { continue; } if batch_map.contains_key(&inp.prev_txid) { @@ -484,6 +686,7 @@ impl Query { let mut packed: Vec<(CreatePin, Vec)> = Vec::with_capacity(work.len()); let mut batch_pin: Vec = Vec::with_capacity(work.len()); let mut planned_fks: Vec = Vec::with_capacity(work.len()); + let mut body_est = 0u64; let mut edges: crate::SpendEdges = crate::SpendEdges::default(); let mut external_parent_vouts: crate::U64Map> = crate::U64Map::default(); let mut batch_stamp = 0u64; @@ -496,12 +699,22 @@ impl Query { } } let mut prestamp_parents = false; - for (tx_fk, tx, mut inputs, outputs) in work { - let mut tx_edges: Vec = Vec::with_capacity(inputs.len()); - for (i, inp) in inputs.iter_mut().enumerate() { - if inp.is_coinbase() { - inp.create_fk = Fk::NULL; - inp.prev_index = u32::MAX; + for row in work { + let PlanRow { + tx_fk, + tx, + ins, + outs, + mut packed_ins, + ins_est, + } = row; + let mut tx_edges: Vec = Vec::with_capacity(ins.len()); + for (i, inp) in ins.iter().enumerate() { + if inp.is_coinbase { + if let Some(rec) = packed_ins.get_mut(i) { + rec.create_fk = Fk::NULL; + rec.prev_index = u32::MAX; + } tx_edges.push(crate::SpendEdge { prev_txid: [0u8; 32], vout: u32::MAX, @@ -510,12 +723,13 @@ impl Query { }); continue; } - if inp.create_fk.is_null() { + let mut create_fk = packed_ins.get(i).map(|r| r.create_fk).unwrap_or(Fk::NULL); + if create_fk.is_null() { if let Some(&cfk) = batch_map.get(&inp.prev_txid) { - inp.create_fk = cfk; + create_fk = cfk; batch_stamp = batch_stamp.saturating_add(1); } else if let Some(&cfk) = resolved.get(&inp.prev_txid) { - inp.create_fk = cfk; + create_fk = cfk; resolved_stamp = resolved_stamp.saturating_add(1); } else { return Err(StoreError::Corrupt( @@ -523,7 +737,10 @@ impl Query { )); } } - if let Some(pid) = inp.create_fk.get() { + if let Some(rec) = packed_ins.get_mut(i) { + rec.create_fk = create_fk; + } + if let Some(pid) = create_fk.get() { if !ArchiveWritePlan::create_in_header_ranges(&per_header_ranges, tx_fk, pid) { external_parent_vouts .entry(pid) @@ -541,7 +758,7 @@ impl Query { if archive_spends { spends.push((inp.prev_txid, inp.prev_index, tx_fk, i as u32)); } - if inp.is_coinbase() || inp.prev_index == u32::MAX { + if inp.prev_index == u32::MAX { tx_edges.push(crate::SpendEdge { prev_txid: [0u8; 32], vout: u32::MAX, @@ -553,7 +770,7 @@ impl Query { prev_txid: inp.prev_txid, vout: inp.prev_index, spend_fk: tx_fk, - create_fk: inp.create_fk, + create_fk, }); } } @@ -561,9 +778,18 @@ impl Query { edges.insert(sid, tx_edges); } planned_fks.push(tx_fk); - let pin = std::sync::Arc::new((tx, outputs)); + let ins_bytes = if packed_ins.is_empty() { + ins_est + } else { + packed_ins.iter().map(|x| x.encoded_len() as u64).sum() + }; + let pin = std::sync::Arc::new((tx, outs)); + body_est = body_est + .saturating_add((1 + TxRecord::ENCODED_LEN) as u64) + .saturating_add(ins_bytes) + .saturating_add(pin.1.iter().map(|x| x.encoded_len() as u64).sum::()); batch_pin.push(std::sync::Arc::clone(&pin)); - packed.push((pin, inputs)); + packed.push((pin, packed_ins)); } let stamp_ns = t_stamp.elapsed().as_nanos() as u64; for vouts in external_parent_vouts.values_mut() { @@ -575,16 +801,6 @@ impl Query { } let t_finish = Instant::now(); - let body_est: u64 = packed - .iter() - .map(|(pin, ins)| { - let (_tx, outs) = pin.as_ref(); - (1 + TxRecord::ENCODED_LEN) as u64 - + ins.iter().map(|x| x.encoded_len() as u64).sum::() - + outs.iter().map(|x| x.encoded_len() as u64).sum::() - }) - .sum(); - let batch_creates: Vec<([u8; 32], Fk)> = packed .iter() .zip(planned_fks.iter()) @@ -1636,6 +1852,29 @@ mod tests { let _ = std::fs::remove_dir_all(&dir); } + #[test] + fn archive_filter_need_header_fks_drops_archived() { + use rbitcoin_store::HeaderRecord; + let (dir, q) = temp_query("filter-header-fks"); + let header = HeaderRecord { + prev_fk: Fk::NULL, + version: 1, + timestamp: 1, + bits: 1, + nonce: 1, + merkle_root: [1u8; 32], + hash: [2u8; 32], + }; + let hfk = q + .commit_class_a_only(&header, &[coinbase_apply(1)]) + .unwrap(); + let need = q + .archive_filter_need_header_fks(&[hfk, hfk, Fk(99)]) + .unwrap(); + assert_eq!(need, vec![Fk(99)], "archived + dup dropped; missing kept"); + let _ = std::fs::remove_dir_all(&dir); + } + /// Stamp emits SpendEdges (create_fk) so pin/write need not walk packed ins. #[test] fn plan_batch_emits_spend_edges() { @@ -1662,6 +1901,107 @@ mod tests { let _ = std::fs::remove_dir_all(&dir); } + fn wire_parent_child_big_script_sig() -> (bitcoin::Block, Vec<[u8; 32]>, Vec) { + use bitcoin::absolute::LockTime; + use bitcoin::block::{Header, Version}; + use bitcoin::hashes::Hash; + use bitcoin::transaction::Version as TxVersion; + use bitcoin::{ + Amount, Block, BlockHash, CompactTarget, OutPoint, ScriptBuf, Sequence, Transaction, + TxIn, TxMerkleNode, TxOut, Witness, + }; + let parent = Transaction { + version: TxVersion::ONE, + lock_time: LockTime::ZERO, + input: vec![TxIn { + previous_output: OutPoint::null(), + script_sig: ScriptBuf::from_bytes(vec![0x01]), + sequence: Sequence::MAX, + witness: Witness::new(), + }], + output: vec![TxOut { + value: Amount::from_sat(50_0000_0000), + script_pubkey: ScriptBuf::from_bytes(vec![0x51, 0xaa]), + }], + }; + let parent_txid = parent.compute_txid(); + let script_sig = vec![0xab; 10_000]; + let child = Transaction { + version: TxVersion::ONE, + lock_time: LockTime::ZERO, + input: vec![TxIn { + previous_output: OutPoint { + txid: parent_txid, + vout: 0, + }, + script_sig: ScriptBuf::from_bytes(script_sig.clone()), + sequence: Sequence::MAX, + witness: Witness::new(), + }], + output: vec![TxOut { + value: Amount::from_sat(1), + script_pubkey: ScriptBuf::from_bytes(vec![0x51, 0xbb]), + }], + }; + let txids = vec![ + parent.compute_txid().to_byte_array(), + child.compute_txid().to_byte_array(), + ]; + let block = Block { + header: Header { + version: Version::ONE, + prev_blockhash: BlockHash::from_byte_array([0; 32]), + merkle_root: TxMerkleNode::from_byte_array([0; 32]), + time: 1, + bits: CompactTarget::from_consensus(0x207f_ffff), + nonce: 0, + }, + txdata: vec![parent, child], + }; + (block, txids, script_sig) + } + + /// D1: wire planner never builds TxApply; packed ins empty; CreatePin matches wire outs. + #[test] + fn plan_batch_from_wire_skips_tx_apply() { + let (dir, q) = temp_query("plan-from-wire"); + let (block, txids, _script_sig) = wire_parent_child_big_script_sig(); + let parent_spk = block.txdata[0].output[0].script_pubkey.to_bytes(); + let child_spk = block.txdata[1].output[0].script_pubkey.to_bytes(); + let parent_txid = txids[0]; + let plan = q + .archive_plan_batch_from_wire( + &[(Fk(1), &block, txids.as_slice())], + 1, + &crate::InFlightView::empty(), + None, + ) + .expect("wire plan"); + assert_eq!(plan.planned_fks, vec![Fk(1), Fk(2)]); + assert!( + plan.packed.iter().all(|(_, ins)| ins.is_empty()), + "wire planner must not clone script_sig into packed ins" + ); + assert_eq!(plan.batch_pin.len(), 2); + assert_eq!(plan.batch_pin[0].1[0].script, parent_spk); + assert_eq!(plan.batch_pin[1].1[0].script, child_spk); + let cb = plan.edges.get(&1).expect("coinbase edges"); + assert_eq!(cb.len(), 1); + assert!(cb[0].create_fk.is_null()); + let edges = plan.edges.get(&2).expect("child edges"); + assert_eq!(edges.len(), 1); + assert_eq!(edges[0].prev_txid, parent_txid); + assert_eq!(edges[0].vout, 0); + assert_eq!(edges[0].spend_fk, Fk(2)); + assert_eq!(edges[0].create_fk, Fk(1)); + assert!( + plan.body_est >= 10_000, + "body_est must count wire ins, got {}", + plan.body_est + ); + let _ = std::fs::remove_dir_all(&dir); + } + /// packed pin half and batch_pin share one CreatePin Arc (no outs double-store). #[test] fn plan_packed_and_batch_pin_share_create_pin_arc() { @@ -2250,4 +2590,32 @@ mod tests { plan.append(super::ArchiveWritePlan::empty()); assert!(plan.is_empty()); } + + #[test] + fn retain_headers_missing_first_fk_is_corrupt() { + let dummy_pin = std::sync::Arc::new(( + TxRecord { + txid: [5u8; 32], + version: 1, + locktime: 0, + input_start_fk: Fk::NULL, + input_count: 0, + output_start_fk: Fk::NULL, + output_count: 0, + }, + Vec::new(), + )); + let mut plan = super::ArchiveWritePlan::empty(); + plan.planned_fks = vec![Fk(5)]; + plan.per_header_ranges = vec![(Fk(10), Fk(99), 1)]; + plan.packed = vec![(dummy_pin, Vec::new())]; + let err = plan + .retain_headers_needing_body(|_| Ok(false)) + .expect_err("missing first fk must not keep the wrong span"); + let msg = err.to_string(); + assert!( + msg.contains("invariant") && msg.contains("retain first fk"), + "unexpected err: {msg}" + ); + } } diff --git a/docs/algo-review.md b/docs/algo-review.md index f6d04405..d179aca3 100644 --- a/docs/algo-review.md +++ b/docs/algo-review.md @@ -151,9 +151,6 @@ What the node runs. Findings follow. ### Query -- **Q-M1.** `retain_headers_needing_body` (`archive.rs`): missing `first` - fk → `unwrap_or(0)` keeps the wrong span. Should be - `Corrupt("invariant: …")`. - **Q-M3.** `merge_outs` clones script bytes on every RCU retry; empty `checked` always publishes a new Arc (breaks `ptr_eq` / sticky). @@ -345,7 +342,7 @@ Intentional COMPAT Electrum status extra field is **not** counted as High. |-------|------|--------|------------------| | store | 0 | seqlock, flush lost-update, fuse8 OOB, spender cycle, sidecar fsync, runs_io | BDZ fill, bulk_fill, SH N² | | consensus + primitives | 0 | *(none remaining)* | historical MTP walks, rehash txids | -| query | 0 | retain fallback, merge_outs clone | BQ scan, SipHash in-flight | +| query | 0 | merge_outs clone | BQ scan, SipHash in-flight | | net | 0 | compact indexes, random eviction, AddrMan/cmpct unbounded, v2 copies | INV flush, BlockCache | | mempool | 0 | orphan vout, eviction tie, persist order, package feerate | free slot, persist_all | | rpc | 0 | submitblock gate, gettxout, hashps, unbounded batch, blockmintxfee, maxfeerate | GBT depends, longpoll | diff --git a/docs/concurrency.md b/docs/concurrency.md index 23000216..333ad29e 100644 --- a/docs/concurrency.md +++ b/docs/concurrency.md @@ -113,7 +113,7 @@ API tokens: [`COMPAT.md`](../COMPAT.md) (Esplora headers, Electrum JSON-RPC extr 3. Scripts for batch N may run while load does N+1 and write does N−1. Scripts never touch disk. 4. **Load ahead of store tip:** lookup may stamp batch N+1 while write has not advanced tip. Lookup holds a **reserved create-fk HWM** and **in-flight create/out maps** from - uncommitted plans (`WireLoadPipeline` / `archive_plan_batch_from_store`). First height + uncommitted plans (`WireLoadPipeline` / `archive_plan_batch_from_wire`). First height of a batch is the **pipeline path_lo** (tip+1 or last-loaded+1), not only store tip. Write still applies batches in height order; on permanent reject, lookup clears reserved state and re-syncs from `txs.count()`. diff --git a/docs/invariants.md b/docs/invariants.md index ee06f70b..a6ee14fa 100644 --- a/docs/invariants.md +++ b/docs/invariants.md @@ -60,9 +60,10 @@ IO** (head / idx), not load pin. After a lookup wave publishes parent P into live_union, load stamp of a child spending P has **zero** leftover TipOnly for P (`head_need_n=0`). Pack stays on load; do not move `plan_batch` onto lookup. -IBD stamp drops packed `InputRecord`s after plan (SpendEdges + CreatePin -survive freeze). Write encodes ins from `Arc` + those edges. CreatePin -outs stay stamp-time for in-flight. Load still does not head/idx. +IBD stamp does not build `TxApply` / packed ins (`archive_plan_batch_from_wire`). +SpendEdges + CreatePin survive freeze. Write encodes ins from `Arc` + +those edges. CreatePin outs stay stamp-time for in-flight. Load still does +not head/idx. | Stage | Allowed IO | Forbidden | |-------|------------|-----------|