Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<Block>` + 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<Block>`
+ 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
Expand Down
63 changes: 42 additions & 21 deletions crates/rbitcoin-consensus/src/confirm_run/lookup.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -310,11 +304,6 @@ pub(super) fn wire_lookup_phase(
}
}

let mut with_fk: Vec<(
rbitcoin_primitives::Fk,
rbitcoin_store::HeaderRecord,
Vec<rbitcoin_query::TxApply>,
)> = Vec::with_capacity(blocks.len());
let mut wire_blocks: Vec<Arc<Block>> = Vec::with_capacity(blocks.len());
let mut metas: Vec<BodyMeta> = Vec::with_capacity(blocks.len());

Expand Down Expand Up @@ -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)?
Expand All @@ -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,
Expand All @@ -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<rbitcoin_primitives::Fk> = 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()
Expand All @@ -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,
Expand Down Expand Up @@ -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);
}
}
49 changes: 33 additions & 16 deletions crates/rbitcoin-consensus/src/confirm_run/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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).
///
Expand Down Expand Up @@ -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<rbitcoin_query::TxApply>,
)> = Vec::with_capacity(blocks.len());
let mut wire_blocks: Vec<Arc<Block>> = Vec::with_capacity(blocks.len());
let mut metas: Vec<BodyMeta> = Vec::with_capacity(blocks.len());

Expand Down Expand Up @@ -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
Expand All @@ -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,
Expand All @@ -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<rbitcoin_primitives::Fk> = 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()
Expand All @@ -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,
Expand Down
60 changes: 14 additions & 46 deletions crates/rbitcoin-consensus/src/confirm_run/pin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Loading
Loading