From 946477f65c1b309648b8a7b0a5e020a64b8e6e33 Mon Sep 17 00:00:00 2001 From: dicethedev Date: Fri, 3 Jul 2026 08:55:51 +0100 Subject: [PATCH 1/8] feat: add block root slot index --- crates/net/p2p/src/req_resp/handlers.rs | 25 +-- crates/storage/src/api/tables.rs | 6 +- crates/storage/src/backend/rocksdb.rs | 1 + crates/storage/src/store.rs | 236 +++++++++++++++++++++++- 4 files changed, 239 insertions(+), 29 deletions(-) diff --git a/crates/net/p2p/src/req_resp/handlers.rs b/crates/net/p2p/src/req_resp/handlers.rs index ad041766..0fe1d615 100644 --- a/crates/net/p2p/src/req_resp/handlers.rs +++ b/crates/net/p2p/src/req_resp/handlers.rs @@ -1,4 +1,4 @@ -use std::collections::{HashMap, HashSet}; +use std::collections::HashSet; use ethlambda_storage::Store; use libp2p::{PeerId, request_response}; @@ -269,29 +269,10 @@ fn canonical_blocks_by_range(store: &Store, start_slot: u64, count: u64) -> Vec< return Vec::new(); }; - let mut roots_by_slot = HashMap::new(); - let mut current_root = store.head(); - - while !current_root.is_zero() { - let Some(header) = store.get_block_header(¤t_root) else { - break; - }; - - if header.slot < start_slot { - break; - } - - if header.slot <= end_slot { - roots_by_slot.insert(header.slot, current_root); - } - - current_root = header.parent_root; - } - (start_slot..=end_slot) .filter_map(|slot| { - let root = roots_by_slot.get(&slot)?; - store.get_signed_block(root) + let root = store.get_block_root_by_slot(slot)?; + store.get_signed_block(&root) }) .collect() } diff --git a/crates/storage/src/api/tables.rs b/crates/storage/src/api/tables.rs index 5884f1f9..40653218 100644 --- a/crates/storage/src/api/tables.rs +++ b/crates/storage/src/api/tables.rs @@ -10,6 +10,8 @@ pub enum Table { /// Stored separately from blocks because the genesis block has no signatures. /// All other blocks must have an entry in this table. BlockSignatures, + /// Canonical block index: slot -> block root + BlockRoots, /// State storage: H256 -> State States, /// Metadata: string keys -> various scalar values @@ -23,10 +25,11 @@ pub enum Table { } /// All table variants. -pub const ALL_TABLES: [Table; 6] = [ +pub const ALL_TABLES: [Table; 7] = [ Table::BlockHeaders, Table::BlockBodies, Table::BlockSignatures, + Table::BlockRoots, Table::States, Table::Metadata, Table::LiveChain, @@ -39,6 +42,7 @@ impl Table { Table::BlockHeaders => "block_headers", Table::BlockBodies => "block_bodies", Table::BlockSignatures => "block_signatures", + Table::BlockRoots => "block_roots", Table::States => "states", Table::Metadata => "metadata", Table::LiveChain => "live_chain", diff --git a/crates/storage/src/backend/rocksdb.rs b/crates/storage/src/backend/rocksdb.rs index e278c8fe..4f3be940 100644 --- a/crates/storage/src/backend/rocksdb.rs +++ b/crates/storage/src/backend/rocksdb.rs @@ -16,6 +16,7 @@ fn cf_name(table: Table) -> &'static str { Table::BlockHeaders => "block_headers", Table::BlockBodies => "block_bodies", Table::BlockSignatures => "block_signatures", + Table::BlockRoots => "block_roots", Table::States => "states", Table::Metadata => "metadata", Table::LiveChain => "live_chain", diff --git a/crates/storage/src/store.rs b/crates/storage/src/store.rs index fb06021f..fd0610a4 100644 --- a/crates/storage/src/store.rs +++ b/crates/storage/src/store.rs @@ -475,12 +475,17 @@ fn decode_live_chain_key(bytes: &[u8]) -> (u64, H256) { (slot, root) } +fn encode_block_root_key(slot: u64) -> Vec { + slot.to_be_bytes().to_vec() +} + /// Fork choice store backed by a pluggable storage backend. /// /// The Store maintains all state required for fork choice and block processing: /// /// - **Metadata**: time, config, head, safe_target, justified/finalized checkpoints /// - **Blocks**: headers and bodies stored separately for efficient header-only queries +/// - **BlockRoots**: canonical block roots indexed by slot /// - **States**: beacon states indexed by block root /// - **Attestations**: latest known and pending ("new") attestations per validator /// - **Signatures**: gossip signatures and aggregated proofs for signature verification @@ -561,15 +566,19 @@ impl Store { ); return None; } - info!("Loaded store from persisted DB state"); - Some(Self { + let store = Self { backend, new_payloads: Arc::new(Mutex::new(PayloadBuffer::new(NEW_PAYLOAD_CAP))), known_payloads: Arc::new(Mutex::new(PayloadBuffer::new(AGGREGATED_PAYLOAD_CAP))), gossip_signatures: Arc::new(Mutex::new(GossipSignatureBuffer::new( GOSSIP_SIGNATURE_CAP, ))), - }) + }; + if store.get_block_root_by_slot(store.head_slot()) != Some(store.head()) { + store.rebuild_block_root_index(); + } + info!("Loaded store from persisted DB state"); + Some(store) } /// Internal helper to initialize the store with anchor data. @@ -631,6 +640,16 @@ impl Store { .put_batch(Table::BlockHeaders, header_entries) .expect("put block header"); + batch + .put_batch( + Table::BlockRoots, + vec![( + encode_block_root_key(anchor_state.latest_block_header.slot), + anchor_block_root.to_ssz(), + )], + ) + .expect("put block root index"); + // Block body (if provided) if let Some(body) = anchor_body { let body_entries = vec![(anchor_block_root.to_ssz(), body.to_ssz())]; @@ -754,6 +773,9 @@ impl Store { pub fn update_checkpoints(&mut self, checkpoints: ForkCheckpoints) { // Read old finalized slot before updating metadata let old_finalized_slot = self.latest_finalized().slot; + let old_head = self.head(); + let (block_root_deletes, block_root_entries) = + self.block_root_index_changes(old_head, checkpoints.head); let mut entries = vec![(KEY_HEAD.to_vec(), checkpoints.head.to_ssz())]; @@ -767,6 +789,12 @@ impl Store { let mut batch = self.backend.begin_write().expect("write batch"); batch.put_batch(Table::Metadata, entries).expect("put"); + batch + .delete_batch(Table::BlockRoots, block_root_deletes) + .expect("delete old canonical block roots"); + batch + .put_batch(Table::BlockRoots, block_root_entries) + .expect("put canonical block roots"); batch.commit().expect("commit"); // Lightweight pruning that should happen immediately on finalization advance: @@ -810,6 +838,98 @@ impl Store { // ============ Blocks ============ + fn block_root_index_changes( + &self, + mut old_root: H256, + mut new_root: H256, + ) -> (Vec>, Vec<(Vec, Vec)>) { + let mut deletes = Vec::new(); + let mut entries = Vec::new(); + + while old_root != new_root { + if old_root.is_zero() { + let header = self + .get_block_header(&new_root) + .expect("new canonical block header exists"); + entries.push((encode_block_root_key(header.slot), new_root.to_ssz())); + new_root = header.parent_root; + continue; + } + if new_root.is_zero() { + let header = self + .get_block_header(&old_root) + .expect("old canonical block header exists"); + deletes.push(encode_block_root_key(header.slot)); + old_root = header.parent_root; + continue; + } + + let old_header = self + .get_block_header(&old_root) + .expect("old canonical block header exists"); + let new_header = self + .get_block_header(&new_root) + .expect("new canonical block header exists"); + + match old_header.slot.cmp(&new_header.slot) { + std::cmp::Ordering::Greater => { + deletes.push(encode_block_root_key(old_header.slot)); + old_root = old_header.parent_root; + } + std::cmp::Ordering::Less => { + entries.push((encode_block_root_key(new_header.slot), new_root.to_ssz())); + new_root = new_header.parent_root; + } + std::cmp::Ordering::Equal => { + deletes.push(encode_block_root_key(old_header.slot)); + entries.push((encode_block_root_key(new_header.slot), new_root.to_ssz())); + old_root = old_header.parent_root; + new_root = new_header.parent_root; + } + } + } + + (deletes, entries) + } + + fn rebuild_block_root_index(&self) { + let view = self.backend.begin_read().expect("read view"); + let old_keys = view + .prefix_iterator(Table::BlockRoots, &[]) + .expect("iterator") + .filter_map(Result::ok) + .map(|(key, _)| key.to_vec()) + .collect(); + drop(view); + + let mut entries = Vec::new(); + let mut root = self.head(); + while !root.is_zero() { + let Some(header) = self.get_block_header(&root) else { + break; + }; + entries.push((encode_block_root_key(header.slot), root.to_ssz())); + root = header.parent_root; + } + + let mut batch = self.backend.begin_write().expect("write batch"); + batch + .delete_batch(Table::BlockRoots, old_keys) + .expect("clear block root index"); + batch + .put_batch(Table::BlockRoots, entries) + .expect("rebuild block root index"); + batch.commit().expect("commit"); + } + + /// Return the canonical block root at `slot`. + pub fn get_block_root_by_slot(&self, slot: u64) -> Option { + let view = self.backend.begin_read().expect("read view"); + view.get(Table::BlockRoots, &encode_block_root_key(slot)) + .expect("get block root") + .map(|bytes| H256::from_ssz_bytes(&bytes).expect("valid block root")) + } + /// Get block data for fork choice: root -> (slot, parent_root). /// /// Iterates only the LiveChain table, avoiding Block deserialization. @@ -987,15 +1107,27 @@ impl Store { let protected: HashSet> = protected_roots.iter().map(|r| r.to_ssz()).collect(); - let keys_to_delete: Vec> = entries + let blocks_to_delete: Vec<(Vec, u64)> = entries .into_iter() .skip(BLOCKS_TO_KEEP) .filter(|(key, _)| !protected.contains(key)) - .map(|(key, _)| key) .collect(); - let count = keys_to_delete.len(); + let count = blocks_to_delete.len(); if count > 0 { + let view = self.backend.begin_read().expect("read view"); + let block_root_keys: Vec> = blocks_to_delete + .iter() + .filter_map(|(root, slot)| { + let slot_key = encode_block_root_key(*slot); + (view.get(Table::BlockRoots, &slot_key).expect("get") == Some(root.clone())) + .then_some(slot_key) + }) + .collect(); + drop(view); + + let keys_to_delete: Vec> = + blocks_to_delete.into_iter().map(|(root, _)| root).collect(); let mut batch = self.backend.begin_write().expect("write batch"); batch .delete_batch(Table::BlockHeaders, keys_to_delete.clone()) @@ -1006,6 +1138,9 @@ impl Store { batch .delete_batch(Table::BlockSignatures, keys_to_delete) .expect("delete old block signatures"); + batch + .delete_batch(Table::BlockRoots, block_root_keys) + .expect("delete pruned block roots"); batch.commit().expect("commit"); } count @@ -1440,6 +1575,12 @@ mod tests { batch .put_batch(Table::BlockSignatures, vec![(key, vec![0u8; 4])]) .expect("put sigs"); + batch + .put_batch( + Table::BlockRoots, + vec![(encode_block_root_key(slot), root.to_ssz())], + ) + .expect("put block root"); batch.commit().expect("commit"); } @@ -1475,6 +1616,19 @@ mod tests { H256::from(bytes) } + fn signed_block(slot: u64, parent_root: H256) -> SignedBlock { + SignedBlock { + message: Block { + slot, + proposer_index: 0, + parent_root, + state_root: H256::ZERO, + body: BlockBody::default(), + }, + proof: MultiMessageAggregate::default(), + } + } + impl Store { /// Create a Store with an in-memory backend for tests. fn test_store() -> Self { @@ -1505,6 +1659,71 @@ mod tests { // ============ Block Pruning Tests ============ + #[test] + fn block_root_index_tracks_canonical_chain_across_reorgs() { + let backend = Arc::new(InMemoryBackend::new()); + let mut store = Store::from_anchor_state(backend, State::from_genesis(0, vec![])); + let anchor_root = store.head(); + + let block_1 = signed_block(1, anchor_root); + let root_1 = block_1.message.hash_tree_root(); + store.insert_signed_block(root_1, block_1); + + let block_3 = signed_block(3, root_1); + let root_3 = block_3.message.hash_tree_root(); + store.insert_signed_block(root_3, block_3); + store.update_checkpoints(ForkCheckpoints::head_only(root_3)); + + assert_eq!(store.get_block_root_by_slot(0), Some(anchor_root)); + assert_eq!(store.get_block_root_by_slot(1), Some(root_1)); + assert_eq!(store.get_block_root_by_slot(2), None); + assert_eq!(store.get_block_root_by_slot(3), Some(root_3)); + + let side_block_2 = signed_block(2, anchor_root); + let side_root_2 = side_block_2.message.hash_tree_root(); + store.insert_signed_block(side_root_2, side_block_2); + + let side_block_4 = signed_block(4, side_root_2); + let side_root_4 = side_block_4.message.hash_tree_root(); + store.insert_signed_block(side_root_4, side_block_4); + store.update_checkpoints(ForkCheckpoints::head_only(side_root_4)); + + assert_eq!(store.get_block_root_by_slot(0), Some(anchor_root)); + assert_eq!(store.get_block_root_by_slot(1), None); + assert_eq!(store.get_block_root_by_slot(2), Some(side_root_2)); + assert_eq!(store.get_block_root_by_slot(3), None); + assert_eq!(store.get_block_root_by_slot(4), Some(side_root_4)); + } + + #[test] + fn from_db_state_rebuilds_block_root_index() { + let backend = Arc::new(InMemoryBackend::new()); + let mut store = + Store::from_anchor_state(backend.clone(), State::from_genesis(12345, vec![])); + + let block = signed_block(1, store.head()); + let block_root = block.message.hash_tree_root(); + store.insert_signed_block(block_root, block); + store.update_checkpoints(ForkCheckpoints::head_only(block_root)); + + let view = backend.begin_read().expect("read view"); + let keys = view + .prefix_iterator(Table::BlockRoots, &[]) + .expect("iterator") + .filter_map(Result::ok) + .map(|(key, _)| key.to_vec()) + .collect(); + drop(view); + let mut batch = backend.begin_write().expect("write batch"); + batch + .delete_batch(Table::BlockRoots, keys) + .expect("clear block roots"); + batch.commit().expect("commit"); + + let restored = Store::from_db_state(backend, 12345).expect("restore store"); + assert_eq!(restored.get_block_root_by_slot(1), Some(block_root)); + } + #[test] fn prune_old_blocks_within_retention() { let backend = Arc::new(InMemoryBackend::new()); @@ -1552,6 +1771,10 @@ mod tests { count_entries(backend.as_ref(), Table::BlockSignatures), BLOCKS_TO_KEEP ); + assert_eq!( + count_entries(backend.as_ref(), Table::BlockRoots), + BLOCKS_TO_KEEP + ); // Oldest blocks (slots 0..10) should be gone for i in 0..10u64 { @@ -1682,6 +1905,7 @@ mod tests { .put_batch( Table::Metadata, vec![ + (KEY_HEAD.to_vec(), finalized.root.to_ssz()), (KEY_LATEST_FINALIZED.to_vec(), finalized.to_ssz()), (KEY_LATEST_JUSTIFIED.to_vec(), justified.to_ssz()), ], From 18ebd3d4d738c1e02dadcb228988cb17d5f3cf9b Mon Sep 17 00:00:00 2001 From: dicethedev Date: Fri, 3 Jul 2026 09:20:13 +0100 Subject: [PATCH 2/8] fix(storage): restore block root slot index after main merge --- crates/storage/src/api/tables.rs | 2 +- crates/storage/src/store.rs | 92 +++++++++++++------------------- 2 files changed, 39 insertions(+), 55 deletions(-) diff --git a/crates/storage/src/api/tables.rs b/crates/storage/src/api/tables.rs index 97fb64fa..300b5293 100644 --- a/crates/storage/src/api/tables.rs +++ b/crates/storage/src/api/tables.rs @@ -34,7 +34,7 @@ pub enum Table { } /// All table variants. -pub const ALL_TABLES: [Table; 7] = [ +pub const ALL_TABLES: [Table; 8] = [ Table::BlockHeaders, Table::BlockBodies, Table::BlockSignatures, diff --git a/crates/storage/src/store.rs b/crates/storage/src/store.rs index 3276e2bb..524b9417 100644 --- a/crates/storage/src/store.rs +++ b/crates/storage/src/store.rs @@ -337,6 +337,10 @@ struct GossipDataEntry { /// Gossip signatures snapshot: (hashed_attestation_data, Vec<(validator_id, signature)>). pub type GossipSignatureSnapshot = Vec<(HashedAttestationData, Vec<(u64, ValidatorSignature)>)>; +type StorageKey = Vec; +type StorageEntry = (StorageKey, Vec); +type BlockRootIndexChanges = (Vec, Vec); + /// Bounded buffer for gossip signatures with FIFO eviction. /// /// Groups signatures by attestation data (via data_root). Each distinct @@ -626,7 +630,12 @@ impl Store { GOSSIP_SIGNATURE_CAP, ))), state_cache: new_state_cache(), - }) + }; + if store.get_block_root_by_slot(store.head_slot()) != Some(store.head()) { + store.rebuild_block_root_index(); + } + info!("Loaded store from persisted DB state"); + Some(store) } /// Internal helper to initialize the store with anchor data. @@ -900,7 +909,7 @@ impl Store { &self, mut old_root: H256, mut new_root: H256, - ) -> (Vec>, Vec<(Vec, Vec)>) { + ) -> BlockRootIndexChanges { let mut deletes = Vec::new(); let mut entries = Vec::new(); @@ -1129,19 +1138,6 @@ impl Store { let count = keys_to_delete.len(); if count > 0 { - let view = self.backend.begin_read().expect("read view"); - let block_root_keys: Vec> = blocks_to_delete - .iter() - .filter_map(|(root, slot)| { - let slot_key = encode_block_root_key(*slot); - (view.get(Table::BlockRoots, &slot_key).expect("get") == Some(root.clone())) - .then_some(slot_key) - }) - .collect(); - drop(view); - - let keys_to_delete: Vec> = - blocks_to_delete.into_iter().map(|(root, _)| root).collect(); let mut batch = self.backend.begin_write().expect("write batch"); batch .delete_batch(Table::BlockSignatures, keys_to_delete) @@ -1805,12 +1801,18 @@ mod tests { let block_1 = signed_block(1, anchor_root); let root_1 = block_1.message.hash_tree_root(); - store.insert_signed_block(root_1, block_1); + store + .insert_signed_block(root_1, block_1) + .expect("insert block 1"); let block_3 = signed_block(3, root_1); let root_3 = block_3.message.hash_tree_root(); - store.insert_signed_block(root_3, block_3); - store.update_checkpoints(ForkCheckpoints::head_only(root_3)); + store + .insert_signed_block(root_3, block_3) + .expect("insert block 3"); + store + .update_checkpoints(ForkCheckpoints::head_only(root_3)) + .expect("update head to block 3"); assert_eq!(store.get_block_root_by_slot(0), Some(anchor_root)); assert_eq!(store.get_block_root_by_slot(1), Some(root_1)); @@ -1819,12 +1821,18 @@ mod tests { let side_block_2 = signed_block(2, anchor_root); let side_root_2 = side_block_2.message.hash_tree_root(); - store.insert_signed_block(side_root_2, side_block_2); + store + .insert_signed_block(side_root_2, side_block_2) + .expect("insert side block 2"); let side_block_4 = signed_block(4, side_root_2); let side_root_4 = side_block_4.message.hash_tree_root(); - store.insert_signed_block(side_root_4, side_block_4); - store.update_checkpoints(ForkCheckpoints::head_only(side_root_4)); + store + .insert_signed_block(side_root_4, side_block_4) + .expect("insert side block 4"); + store + .update_checkpoints(ForkCheckpoints::head_only(side_root_4)) + .expect("update head to side block 4"); assert_eq!(store.get_block_root_by_slot(0), Some(anchor_root)); assert_eq!(store.get_block_root_by_slot(1), None); @@ -1841,8 +1849,12 @@ mod tests { let block = signed_block(1, store.head()); let block_root = block.message.hash_tree_root(); - store.insert_signed_block(block_root, block); - store.update_checkpoints(ForkCheckpoints::head_only(block_root)); + store + .insert_signed_block(block_root, block) + .expect("insert block"); + store + .update_checkpoints(ForkCheckpoints::head_only(block_root)) + .expect("update head"); let view = backend.begin_read().expect("read view"); let keys = view @@ -1882,24 +1894,9 @@ mod tests { // cutoff = 10: slots 0..9 pruned, slots 10..12 kept (within the window). assert_eq!(pruned, 10); - assert_eq!( - count_entries(backend.as_ref(), Table::BlockHeaders), - BLOCKS_TO_KEEP - ); - assert_eq!( - count_entries(backend.as_ref(), Table::BlockBodies), - BLOCKS_TO_KEEP - ); - assert_eq!( - count_entries(backend.as_ref(), Table::BlockSignatures), - BLOCKS_TO_KEEP - ); - assert_eq!( - count_entries(backend.as_ref(), Table::BlockRoots), - BLOCKS_TO_KEEP - ); + assert_eq!(count_entries(backend.as_ref(), Table::BlockSignatures), 3); - // Oldest blocks (slots 0..10) should be gone + // Oldest signatures are gone, but headers, bodies, and roots stay queryable. for i in 0..10u64 { assert!(!has_signature(backend.as_ref(), i, &root(i))); } @@ -1910,6 +1907,7 @@ mod tests { // Headers and bodies are always retained for the whole history. assert_eq!(count_entries(backend.as_ref(), Table::BlockHeaders), 13); assert_eq!(count_entries(backend.as_ref(), Table::BlockBodies), 13); + assert_eq!(count_entries(backend.as_ref(), Table::BlockRoots), 13); } #[test] @@ -2007,20 +2005,6 @@ mod tests { // Hot path: the just-imported state is memoized in the cache. assert_eq!(store.get_state(&r1).unwrap().to_ssz(), s1.to_ssz()); - /// Set up finalized and justified checkpoints in metadata. - fn set_checkpoints(backend: &dyn StorageBackend, finalized: Checkpoint, justified: Checkpoint) { - let mut batch = backend.begin_write().expect("write batch"); - batch - .put_batch( - Table::Metadata, - vec![ - (KEY_HEAD.to_vec(), finalized.root.to_ssz()), - (KEY_LATEST_FINALIZED.to_vec(), finalized.to_ssz()), - (KEY_LATEST_JUSTIFIED.to_vec(), justified.to_ssz()), - ], - ) - .expect("put checkpoints"); - batch.commit().expect("commit"); // A cold store (empty cache, shared backend) reconstructs from the diff, // byte-identically. let cold = Store::test_store_with_backend(backend.clone()); From 70235d0a6de99707d9f1a4a80e62c771ee217412 Mon Sep 17 00:00:00 2001 From: dicethedev Date: Fri, 17 Jul 2026 12:53:02 +0100 Subject: [PATCH 3/8] fix(storage): resolve block root index merge conflicts --- crates/net/p2p/src/req_resp/handlers.rs | 25 +------------ crates/storage/src/store.rs | 47 ++++++++++++------------- 2 files changed, 23 insertions(+), 49 deletions(-) diff --git a/crates/net/p2p/src/req_resp/handlers.rs b/crates/net/p2p/src/req_resp/handlers.rs index 49c295d8..c1d4a9bc 100644 --- a/crates/net/p2p/src/req_resp/handlers.rs +++ b/crates/net/p2p/src/req_resp/handlers.rs @@ -272,30 +272,7 @@ fn canonical_blocks_by_range(store: &Store, start_slot: u64, count: u64) -> Vec< (start_slot..=end_slot) .filter_map(|slot| { let root = store.get_block_root_by_slot(slot)?; - store.get_signed_block(&root) - let mut roots_by_slot = HashMap::new(); - let mut current_root = store.head().expect("head block exists"); - - while !current_root.is_zero() { - let Ok(Some(header)) = store.get_block_header(¤t_root) else { - break; - }; - - if header.slot < start_slot { - break; - } - - if header.slot <= end_slot { - roots_by_slot.insert(header.slot, current_root); - } - - current_root = header.parent_root; - } - - (start_slot..=end_slot) - .filter_map(|slot| { - let root = roots_by_slot.get(&slot)?; - store.get_signed_block(root).ok().flatten() + store.get_signed_block(&root).ok().flatten() }) .collect() } diff --git a/crates/storage/src/store.rs b/crates/storage/src/store.rs index 4f8e1b37..f1bb9b3e 100644 --- a/crates/storage/src/store.rs +++ b/crates/storage/src/store.rs @@ -635,8 +635,6 @@ impl Store { return Ok(None); } let store = Self { - info!("Loaded store from persisted DB state"); - Ok(Some(Self { backend, new_payloads: Arc::new(Mutex::new(PayloadBuffer::new(NEW_PAYLOAD_CAP))), known_payloads: Arc::new(Mutex::new(PayloadBuffer::new(AGGREGATED_PAYLOAD_CAP))), @@ -645,12 +643,12 @@ impl Store { ))), state_cache: new_state_cache(), }; - if store.get_block_root_by_slot(store.head_slot()) != Some(store.head()) { - store.rebuild_block_root_index(); + let head = store.head()?; + if store.get_block_root_by_slot(store.head_slot()) != Some(head) { + store.rebuild_block_root_index()?; } info!("Loaded store from persisted DB state"); - Some(store) - })) + Ok(Some(store)) } /// Internal helper to initialize the store with anchor data. @@ -848,14 +846,10 @@ impl Store { /// When finalization advances, prunes the LiveChain index. pub fn update_checkpoints(&mut self, checkpoints: ForkCheckpoints) -> Result<(), Error> { // Read old finalized slot before updating metadata - let old_finalized_slot = self.latest_finalized().slot; - let old_head = self.head(); + let old_finalized_slot = self.latest_finalized()?.slot; + let old_head = self.head()?; let (block_root_deletes, block_root_entries) = - self.block_root_index_changes(old_head, checkpoints.head); - let old_finalized_slot = self - .latest_finalized() - .expect("Failed to get latest finalized checkpoint") - .slot; + self.block_root_index_changes(old_head, checkpoints.head)?; let mut entries = vec![(KEY_HEAD.to_vec(), checkpoints.head.to_ssz())]; @@ -934,14 +928,14 @@ impl Store { &self, mut old_root: H256, mut new_root: H256, - ) -> BlockRootIndexChanges { + ) -> Result { let mut deletes = Vec::new(); let mut entries = Vec::new(); while old_root != new_root { if old_root.is_zero() { let header = self - .get_block_header(&new_root) + .get_block_header(&new_root)? .expect("new canonical block header exists"); entries.push((encode_block_root_key(header.slot), new_root.to_ssz())); new_root = header.parent_root; @@ -949,7 +943,7 @@ impl Store { } if new_root.is_zero() { let header = self - .get_block_header(&old_root) + .get_block_header(&old_root)? .expect("old canonical block header exists"); deletes.push(encode_block_root_key(header.slot)); old_root = header.parent_root; @@ -957,10 +951,10 @@ impl Store { } let old_header = self - .get_block_header(&old_root) + .get_block_header(&old_root)? .expect("old canonical block header exists"); let new_header = self - .get_block_header(&new_root) + .get_block_header(&new_root)? .expect("new canonical block header exists"); match old_header.slot.cmp(&new_header.slot) { @@ -981,10 +975,10 @@ impl Store { } } - (deletes, entries) + Ok((deletes, entries)) } - fn rebuild_block_root_index(&self) { + fn rebuild_block_root_index(&self) -> Result<(), Error> { let view = self.backend.begin_read().expect("read view"); let old_keys = view .prefix_iterator(Table::BlockRoots, &[]) @@ -995,9 +989,9 @@ impl Store { drop(view); let mut entries = Vec::new(); - let mut root = self.head(); + let mut root = self.head()?; while !root.is_zero() { - let Some(header) = self.get_block_header(&root) else { + let Some(header) = self.get_block_header(&root)? else { break; }; entries.push((encode_block_root_key(header.slot), root.to_ssz())); @@ -1012,6 +1006,7 @@ impl Store { .put_batch(Table::BlockRoots, entries) .expect("rebuild block root index"); batch.commit().expect("commit"); + Ok(()) } /// Return the canonical block root at `slot`. @@ -1860,7 +1855,7 @@ mod tests { fn block_root_index_tracks_canonical_chain_across_reorgs() { let backend = Arc::new(InMemoryBackend::new()); let mut store = Store::from_anchor_state(backend, State::from_genesis(0, vec![])); - let anchor_root = store.head(); + let anchor_root = store.head().expect("head root"); let block_1 = signed_block(1, anchor_root); let root_1 = block_1.message.hash_tree_root(); @@ -1910,7 +1905,7 @@ mod tests { let mut store = Store::from_anchor_state(backend.clone(), State::from_genesis(12345, vec![])); - let block = signed_block(1, store.head()); + let block = signed_block(1, store.head().expect("head root")); let block_root = block.message.hash_tree_root(); store .insert_signed_block(block_root, block) @@ -1933,7 +1928,9 @@ mod tests { .expect("clear block roots"); batch.commit().expect("commit"); - let restored = Store::from_db_state(backend, 12345).expect("restore store"); + let restored = Store::from_db_state(backend, 12345) + .expect("restore store") + .expect("store exists"); assert_eq!(restored.get_block_root_by_slot(1), Some(block_root)); } From 2b9654ad29afc4e7424b3419fd0184da83273ea5 Mon Sep 17 00:00:00 2001 From: dicethedev Date: Tue, 21 Jul 2026 06:02:07 +0100 Subject: [PATCH 4/8] fix(storage): address block root index review feedback --- crates/net/p2p/src/req_resp/handlers.rs | 9 +-- crates/storage/src/store.rs | 73 ++++++++----------------- 2 files changed, 26 insertions(+), 56 deletions(-) diff --git a/crates/net/p2p/src/req_resp/handlers.rs b/crates/net/p2p/src/req_resp/handlers.rs index c1d4a9bc..7ee31a25 100644 --- a/crates/net/p2p/src/req_resp/handlers.rs +++ b/crates/net/p2p/src/req_resp/handlers.rs @@ -269,12 +269,9 @@ fn canonical_blocks_by_range(store: &Store, start_slot: u64, count: u64) -> Vec< return Vec::new(); }; - (start_slot..=end_slot) - .filter_map(|slot| { - let root = store.get_block_root_by_slot(slot)?; - store.get_signed_block(&root).ok().flatten() - }) - .collect() + store + .get_signed_blocks_by_slot_range(start_slot, end_slot) + .unwrap_or_default() } async fn handle_blocks_by_root_response( diff --git a/crates/storage/src/store.rs b/crates/storage/src/store.rs index f1bb9b3e..881d8918 100644 --- a/crates/storage/src/store.rs +++ b/crates/storage/src/store.rs @@ -643,10 +643,6 @@ impl Store { ))), state_cache: new_state_cache(), }; - let head = store.head()?; - if store.get_block_root_by_slot(store.head_slot()) != Some(head) { - store.rebuild_block_root_index()?; - } info!("Loaded store from persisted DB state"); Ok(Some(store)) } @@ -978,37 +974,6 @@ impl Store { Ok((deletes, entries)) } - fn rebuild_block_root_index(&self) -> Result<(), Error> { - let view = self.backend.begin_read().expect("read view"); - let old_keys = view - .prefix_iterator(Table::BlockRoots, &[]) - .expect("iterator") - .filter_map(Result::ok) - .map(|(key, _)| key.to_vec()) - .collect(); - drop(view); - - let mut entries = Vec::new(); - let mut root = self.head()?; - while !root.is_zero() { - let Some(header) = self.get_block_header(&root)? else { - break; - }; - entries.push((encode_block_root_key(header.slot), root.to_ssz())); - root = header.parent_root; - } - - let mut batch = self.backend.begin_write().expect("write batch"); - batch - .delete_batch(Table::BlockRoots, old_keys) - .expect("clear block root index"); - batch - .put_batch(Table::BlockRoots, entries) - .expect("rebuild block root index"); - batch.commit().expect("commit"); - Ok(()) - } - /// Return the canonical block root at `slot`. pub fn get_block_root_by_slot(&self, slot: u64) -> Option { let view = self.backend.begin_read().expect("read view"); @@ -1305,6 +1270,28 @@ impl Store { })) } + /// Return canonical signed blocks for the slot range `[start_slot, end_slot]`. + /// + /// Missing slots or blocks are skipped. This keeps the current request + /// behavior while centralizing the slot-index lookup so the storage backend + /// can optimize range reads later. + pub fn get_signed_blocks_by_slot_range( + &self, + start_slot: u64, + end_slot: u64, + ) -> Result, Error> { + let mut blocks = Vec::new(); + for slot in start_slot..=end_slot { + let Some(root) = self.get_block_root_by_slot(slot) else { + continue; + }; + if let Some(block) = self.get_signed_block(&root)? { + blocks.push(block); + } + } + Ok(blocks) + } + // ============ States ============ /// Returns the state for the given block root. @@ -1900,7 +1887,7 @@ mod tests { } #[test] - fn from_db_state_rebuilds_block_root_index() { + fn from_db_state_preserves_block_root_index() { let backend = Arc::new(InMemoryBackend::new()); let mut store = Store::from_anchor_state(backend.clone(), State::from_genesis(12345, vec![])); @@ -1914,20 +1901,6 @@ mod tests { .update_checkpoints(ForkCheckpoints::head_only(block_root)) .expect("update head"); - let view = backend.begin_read().expect("read view"); - let keys = view - .prefix_iterator(Table::BlockRoots, &[]) - .expect("iterator") - .filter_map(Result::ok) - .map(|(key, _)| key.to_vec()) - .collect(); - drop(view); - let mut batch = backend.begin_write().expect("write batch"); - batch - .delete_batch(Table::BlockRoots, keys) - .expect("clear block roots"); - batch.commit().expect("commit"); - let restored = Store::from_db_state(backend, 12345) .expect("restore store") .expect("store exists"); From 558de7301e8b370c960b0a3ebc35464a583126e2 Mon Sep 17 00:00:00 2001 From: dicethedev Date: Tue, 21 Jul 2026 23:36:35 +0100 Subject: [PATCH 5/8] reviews --- crates/net/p2p/src/req_resp/handlers.rs | 15 ++- crates/storage/src/store.rs | 121 ++++++++++++++++-------- 2 files changed, 91 insertions(+), 45 deletions(-) diff --git a/crates/net/p2p/src/req_resp/handlers.rs b/crates/net/p2p/src/req_resp/handlers.rs index 7ee31a25..489ebfba 100644 --- a/crates/net/p2p/src/req_resp/handlers.rs +++ b/crates/net/p2p/src/req_resp/handlers.rs @@ -269,9 +269,18 @@ fn canonical_blocks_by_range(store: &Store, start_slot: u64, count: u64) -> Vec< return Vec::new(); }; - store - .get_signed_blocks_by_slot_range(start_slot, end_slot) - .unwrap_or_default() + match store.get_signed_blocks_by_slot_range(start_slot, end_slot) { + Ok(blocks) => blocks, + Err(err) => { + warn!( + start_slot, + end_slot, + ?err, + "Failed to get signed blocks by slot range" + ); + Vec::new() + } + } } async fn handle_blocks_by_root_response( diff --git a/crates/storage/src/store.rs b/crates/storage/src/store.rs index 881d8918..d2fb5629 100644 --- a/crates/storage/src/store.rs +++ b/crates/storage/src/store.rs @@ -4,7 +4,7 @@ use std::sync::{Arc, LazyLock, Mutex}; use lru::LruCache; -use crate::api::{StorageBackend, StorageWriteBatch, Table}; +use crate::api::{StorageBackend, StorageReadView, StorageWriteBatch, Table}; use crate::error::Error; use ethlambda_types::{ @@ -930,28 +930,44 @@ impl Store { while old_root != new_root { if old_root.is_zero() { - let header = self - .get_block_header(&new_root)? - .expect("new canonical block header exists"); + let Some(header) = self.get_block_header(&new_root)? else { + warn!( + ?new_root, + "Skipping block root index update for missing new head header" + ); + break; + }; entries.push((encode_block_root_key(header.slot), new_root.to_ssz())); new_root = header.parent_root; continue; } if new_root.is_zero() { - let header = self - .get_block_header(&old_root)? - .expect("old canonical block header exists"); + let Some(header) = self.get_block_header(&old_root)? else { + warn!( + ?old_root, + "Skipping block root index update for missing old head header" + ); + break; + }; deletes.push(encode_block_root_key(header.slot)); old_root = header.parent_root; continue; } - let old_header = self - .get_block_header(&old_root)? - .expect("old canonical block header exists"); - let new_header = self - .get_block_header(&new_root)? - .expect("new canonical block header exists"); + let Some(old_header) = self.get_block_header(&old_root)? else { + warn!( + ?old_root, + "Skipping block root index update for missing old head header" + ); + break; + }; + let Some(new_header) = self.get_block_header(&new_root)? else { + warn!( + ?new_root, + "Skipping block root index update for missing new head header" + ); + break; + }; match old_header.slot.cmp(&new_header.slot) { std::cmp::Ordering::Greater => { @@ -974,14 +990,6 @@ impl Store { Ok((deletes, entries)) } - /// Return the canonical block root at `slot`. - pub fn get_block_root_by_slot(&self, slot: u64) -> Option { - let view = self.backend.begin_read().expect("read view"); - view.get(Table::BlockRoots, &encode_block_root_key(slot)) - .expect("get block root") - .map(|bytes| H256::from_ssz_bytes(&bytes).expect("valid block root")) - } - /// Get block data for fork choice: root -> (slot, parent_root). /// /// Iterates only the LiveChain table, avoiding Block deserialization. @@ -1233,20 +1241,20 @@ impl Store { /// longer be served with its proof) rather than as a fabricated block. pub fn get_signed_block(&self, root: &H256) -> Result, Error> { let view = self.backend.begin_read().expect("read view"); + Ok(Self::signed_block_from_view(view.as_ref(), root)) + } + + fn signed_block_from_view(view: &dyn StorageReadView, root: &H256) -> Option { let key = root.to_ssz(); - let Some(header_bytes) = view.get(Table::BlockHeaders, &key).expect("get") else { - return Ok(None); - }; + let header_bytes = view.get(Table::BlockHeaders, &key).expect("get")?; let header = BlockHeader::from_ssz_bytes(&header_bytes).expect("valid header"); // Use empty body if header indicates empty, otherwise fetch from DB let body = if header.body_root == *EMPTY_BODY_ROOT { BlockBody::default() } else { - let Some(body_bytes) = view.get(Table::BlockBodies, &key).expect("get") else { - return Ok(None); - }; + let body_bytes = view.get(Table::BlockBodies, &key).expect("get")?; BlockBody::from_ssz_bytes(&body_bytes).expect("valid body") }; @@ -1259,15 +1267,15 @@ impl Store { // other slot a missing proof (pruned finalized block, or genuine // corruption) surfaces as `None` rather than a fabricated block. None if header.slot == 0 => MultiMessageAggregate::default(), - None => return Ok(None), + None => return None, }; let block = Block::from_header_and_body(header, body); - Ok(Some(SignedBlock { + Some(SignedBlock { message: block, proof, - })) + }) } /// Return canonical signed blocks for the slot range `[start_slot, end_slot]`. @@ -1280,12 +1288,17 @@ impl Store { start_slot: u64, end_slot: u64, ) -> Result, Error> { + let view = self.backend.begin_read().expect("read view"); let mut blocks = Vec::new(); for slot in start_slot..=end_slot { - let Some(root) = self.get_block_root_by_slot(slot) else { + let Some(root_bytes) = view + .get(Table::BlockRoots, &encode_block_root_key(slot)) + .expect("get block root") + else { continue; }; - if let Some(block) = self.get_signed_block(&root)? { + let root = H256::from_ssz_bytes(&root_bytes).expect("valid block root"); + if let Some(block) = Self::signed_block_from_view(view.as_ref(), &root) { blocks.push(block); } } @@ -1786,6 +1799,14 @@ mod tests { .is_some() } + /// Return the canonical block root at `slot` for storage-index assertions. + fn block_root_by_slot(backend: &dyn StorageBackend, slot: u64) -> Option { + let view = backend.begin_read().expect("read view"); + view.get(Table::BlockRoots, &encode_block_root_key(slot)) + .expect("get block root") + .map(|bytes| H256::from_ssz_bytes(&bytes).expect("valid block root")) + } + /// Generate a deterministic H256 root from an index. fn root(index: u64) -> H256 { let mut bytes = [0u8; 32]; @@ -1859,10 +1880,13 @@ mod tests { .update_checkpoints(ForkCheckpoints::head_only(root_3)) .expect("update head to block 3"); - assert_eq!(store.get_block_root_by_slot(0), Some(anchor_root)); - assert_eq!(store.get_block_root_by_slot(1), Some(root_1)); - assert_eq!(store.get_block_root_by_slot(2), None); - assert_eq!(store.get_block_root_by_slot(3), Some(root_3)); + assert_eq!( + block_root_by_slot(store.backend.as_ref(), 0), + Some(anchor_root) + ); + assert_eq!(block_root_by_slot(store.backend.as_ref(), 1), Some(root_1)); + assert_eq!(block_root_by_slot(store.backend.as_ref(), 2), None); + assert_eq!(block_root_by_slot(store.backend.as_ref(), 3), Some(root_3)); let side_block_2 = signed_block(2, anchor_root); let side_root_2 = side_block_2.message.hash_tree_root(); @@ -1879,11 +1903,20 @@ mod tests { .update_checkpoints(ForkCheckpoints::head_only(side_root_4)) .expect("update head to side block 4"); - assert_eq!(store.get_block_root_by_slot(0), Some(anchor_root)); - assert_eq!(store.get_block_root_by_slot(1), None); - assert_eq!(store.get_block_root_by_slot(2), Some(side_root_2)); - assert_eq!(store.get_block_root_by_slot(3), None); - assert_eq!(store.get_block_root_by_slot(4), Some(side_root_4)); + assert_eq!( + block_root_by_slot(store.backend.as_ref(), 0), + Some(anchor_root) + ); + assert_eq!(block_root_by_slot(store.backend.as_ref(), 1), None); + assert_eq!( + block_root_by_slot(store.backend.as_ref(), 2), + Some(side_root_2) + ); + assert_eq!(block_root_by_slot(store.backend.as_ref(), 3), None); + assert_eq!( + block_root_by_slot(store.backend.as_ref(), 4), + Some(side_root_4) + ); } #[test] @@ -1904,7 +1937,11 @@ mod tests { let restored = Store::from_db_state(backend, 12345) .expect("restore store") .expect("store exists"); - assert_eq!(restored.get_block_root_by_slot(1), Some(block_root)); + let blocks = restored + .get_signed_blocks_by_slot_range(1, 1) + .expect("get blocks by slot range"); + assert_eq!(blocks.len(), 1); + assert_eq!(blocks[0].message.hash_tree_root(), block_root); } #[test] From 4993516f553e00d600943cef057227580151a69b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Tom=C3=A1s=20Gr=C3=BCner?= <47506558+MegaRedHand@users.noreply.github.com> Date: Fri, 24 Jul 2026 16:17:45 -0300 Subject: [PATCH 6/8] refactor(p2p): use inspect_err for BlocksByRange lookup failure The match existed only to log the error and fall back to an empty vec, which is exactly what the codebase's error-handling idiom expresses with inspect_err + unwrap_or_default. --- crates/net/p2p/src/req_resp/handlers.rs | 13 ++++++------- 1 file changed, 6 insertions(+), 7 deletions(-) diff --git a/crates/net/p2p/src/req_resp/handlers.rs b/crates/net/p2p/src/req_resp/handlers.rs index 489ebfba..810e658b 100644 --- a/crates/net/p2p/src/req_resp/handlers.rs +++ b/crates/net/p2p/src/req_resp/handlers.rs @@ -269,18 +269,17 @@ fn canonical_blocks_by_range(store: &Store, start_slot: u64, count: u64) -> Vec< return Vec::new(); }; - match store.get_signed_blocks_by_slot_range(start_slot, end_slot) { - Ok(blocks) => blocks, - Err(err) => { + store + .get_signed_blocks_by_slot_range(start_slot, end_slot) + .inspect_err(|err| { warn!( start_slot, end_slot, ?err, "Failed to get signed blocks by slot range" - ); - Vec::new() - } - } + ) + }) + .unwrap_or_default() } async fn handle_blocks_by_root_response( From 43f0ee29b715834632d92e7886503130837cc426 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Tom=C3=A1s=20Gr=C3=BCner?= <47506558+MegaRedHand@users.noreply.github.com> Date: Fri, 24 Jul 2026 17:17:00 -0300 Subject: [PATCH 7/8] refactor(storage): simplify the block root index diff walk The two-branch walk special-cased zero roots and refetched both headers every iteration, so the same "fetch a header or give up" block appeared four times. Carrying each branch's header across iterations and refetching only the side that stepped drops both special cases: genesis' parent has no header, so it already flows through the missing-header path. Equal slots no longer need an arm of their own either. The old branch steps first, dropping it below the new one, which then steps on the following iteration, yielding the same deletes and entries in the same order. A missing header now returns Error::UnexpectedMissingBlockHeader rather than warning, so update_checkpoints aborts before writing anything instead of committing a partial diff that describes neither branch. --- crates/storage/src/error.rs | 4 ++ crates/storage/src/store.rs | 82 ++++++++++++------------------------- 2 files changed, 31 insertions(+), 55 deletions(-) diff --git a/crates/storage/src/error.rs b/crates/storage/src/error.rs index 8333a2cf..3f802aeb 100644 --- a/crates/storage/src/error.rs +++ b/crates/storage/src/error.rs @@ -1,5 +1,9 @@ +use ethlambda_types::primitives::H256; + #[derive(Debug, thiserror::Error)] pub enum Error { #[error("storage error: {0}")] Storage(#[from] crate::api::Error), + #[error("unexpected missing block header for root {0}")] + UnexpectedMissingBlockHeader(H256), } diff --git a/crates/storage/src/store.rs b/crates/storage/src/store.rs index d2fb5629..83f1e8fc 100644 --- a/crates/storage/src/store.rs +++ b/crates/storage/src/store.rs @@ -920,6 +920,13 @@ impl Store { // ============ Blocks ============ + /// `BlockRoots` index diff between the branch ending at `old_root` and the one + /// ending at `new_root`: slot keys to delete (canonical only on the old branch) + /// and slot -> root entries to write (canonical on the new branch). + /// + /// Both branches must be walkable down to their common ancestor. A root with no + /// header, genesis' zero parent included, means they have none in common, and + /// yields [`Error::UnexpectedMissingBlockHeader`] instead of a partial diff. fn block_root_index_changes( &self, mut old_root: H256, @@ -928,62 +935,27 @@ impl Store { let mut deletes = Vec::new(); let mut entries = Vec::new(); - while old_root != new_root { - if old_root.is_zero() { - let Some(header) = self.get_block_header(&new_root)? else { - warn!( - ?new_root, - "Skipping block root index update for missing new head header" - ); - break; - }; - entries.push((encode_block_root_key(header.slot), new_root.to_ssz())); - new_root = header.parent_root; - continue; - } - if new_root.is_zero() { - let Some(header) = self.get_block_header(&old_root)? else { - warn!( - ?old_root, - "Skipping block root index update for missing old head header" - ); - break; - }; - deletes.push(encode_block_root_key(header.slot)); - old_root = header.parent_root; - continue; - } - - let Some(old_header) = self.get_block_header(&old_root)? else { - warn!( - ?old_root, - "Skipping block root index update for missing old head header" - ); - break; - }; - let Some(new_header) = self.get_block_header(&new_root)? else { - warn!( - ?new_root, - "Skipping block root index update for missing new head header" - ); - break; - }; + let mut old_header = self + .get_block_header(&old_root)? + .ok_or(Error::UnexpectedMissingBlockHeader(old_root))?; + let mut new_header = self + .get_block_header(&new_root)? + .ok_or(Error::UnexpectedMissingBlockHeader(new_root))?; - match old_header.slot.cmp(&new_header.slot) { - std::cmp::Ordering::Greater => { - deletes.push(encode_block_root_key(old_header.slot)); - old_root = old_header.parent_root; - } - std::cmp::Ordering::Less => { - entries.push((encode_block_root_key(new_header.slot), new_root.to_ssz())); - new_root = new_header.parent_root; - } - std::cmp::Ordering::Equal => { - deletes.push(encode_block_root_key(old_header.slot)); - entries.push((encode_block_root_key(new_header.slot), new_root.to_ssz())); - old_root = old_header.parent_root; - new_root = new_header.parent_root; - } + // Walk both branches back toward their common ancestor, until we find the common ancestor. + while old_root != new_root { + if old_header.slot < new_header.slot { + entries.push((encode_block_root_key(new_header.slot), new_root.to_ssz())); + new_root = new_header.parent_root; + new_header = self + .get_block_header(&new_root)? + .ok_or(Error::UnexpectedMissingBlockHeader(new_root))?; + } else { + deletes.push(encode_block_root_key(old_header.slot)); + old_root = old_header.parent_root; + old_header = self + .get_block_header(&old_root)? + .ok_or(Error::UnexpectedMissingBlockHeader(old_root))?; } } From 0cad37a2b706e5de54895376c5c8f6fff1d3257e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Tom=C3=A1s=20Gr=C3=BCner?= <47506558+MegaRedHand@users.noreply.github.com> Date: Fri, 24 Jul 2026 18:10:55 -0300 Subject: [PATCH 8/8] test(events): drop the orphan-head diff test The canonical slot -> root index makes `update_checkpoints` reject a head whose block header is not stored, so this test can no longer set up its fixture: it pointed the head at an unstored root to reach the warn-and-skip path in `diff_and_emit`. That defensive branch stays; only its test goes. --- crates/blockchain/src/events.rs | 36 --------------------------------- 1 file changed, 36 deletions(-) diff --git a/crates/blockchain/src/events.rs b/crates/blockchain/src/events.rs index 0f3367ec..94862418 100644 --- a/crates/blockchain/src/events.rs +++ b/crates/blockchain/src/events.rs @@ -417,42 +417,6 @@ mod tests { assert!(matches!(rx.try_recv(), Err(TryRecvError::Empty))); } - /// A head whose header cannot be read is skipped (warn), while checkpoint - /// moves still emit. - #[test] - fn chain_event_diff_skips_head_with_missing_header() { - let mut store = test_store(); - let bus = EventBus::new(8); - let mut rx = bus.subscribe(); - - let snapshot = ChainEventSnapshot::capture(&store); - - // Point the head at a root with no stored header; advance finalized to - // a real block so its event still fires. - let orphan_head = H256([7u8; 32]); - let genesis = store.head().expect("store head exists"); - let finalized_root = H256([8u8; 32]); - let finalized_state = H256([88u8; 32]); - insert_test_block(&mut store, finalized_root, 1, genesis, finalized_state); - let finalized = Checkpoint { - root: finalized_root, - slot: 1, - }; - store - .update_checkpoints(ForkCheckpoints::new(orphan_head, None, Some(finalized))) - .expect("update_checkpoints should succeed"); - - snapshot.diff_and_emit(&store, &bus, 1); - - match rx.try_recv().unwrap() { - ChainEvent::FinalizedCheckpoint { slot, block, state } => { - assert_eq!((slot, block, state), (1, finalized_root, finalized_state)); - } - other => panic!("expected finalized_checkpoint only, got: {other:?}"), - } - assert!(matches!(rx.try_recv(), Err(TryRecvError::Empty))); - } - /// A head far below the wall-clock slot (catch-up/backfill) emits no /// `head` event, but `justified_checkpoint`/`finalized_checkpoint` still /// fire since only `head` is gated.