From 28e550229d3480b9529a9e5d72dec62c74db50de Mon Sep 17 00:00:00 2001 From: dicethedev Date: Fri, 31 Jul 2026 13:38:11 +0100 Subject: [PATCH 1/3] fix(storage): store fork-choice votes independently --- crates/blockchain/src/store.rs | 1 + crates/storage/src/store.rs | 149 +++++++++++++++++++++++++++++++-- 2 files changed, 144 insertions(+), 6 deletions(-) diff --git a/crates/blockchain/src/store.rs b/crates/blockchain/src/store.rs index 1a6d3884..e5d75f45 100644 --- a/crates/blockchain/src/store.rs +++ b/crates/blockchain/src/store.rs @@ -699,6 +699,7 @@ fn on_block_core( store .insert_state(block_root, post_state) .expect("DB insert should succeed"); + store.insert_known_attestation_votes(&block.body.attestations); // Block-included attestations are intentionally not counted here. // `lean_attestations_valid_total` tracks the gossip validation pipeline diff --git a/crates/storage/src/store.rs b/crates/storage/src/store.rs index d5cac0ed..c3143c94 100644 --- a/crates/storage/src/store.rs +++ b/crates/storage/src/store.rs @@ -9,7 +9,10 @@ use crate::error::Error; use ethlambda_crypto::signature::ValidatorSignature; use ethlambda_types::{ - attestation::{AggregationBits, AttestationData, HashedAttestationData, bits_is_subset}, + attestation::{ + AggregatedAttestation, AggregationBits, AttestationData, HashedAttestationData, + bits_is_subset, validator_indices, + }, block::{ Block, BlockBody, BlockHeader, MultiMessageAggregate, SignedBlock, SingleMessageAggregate, }, @@ -339,6 +342,7 @@ pub type GossipSignatureSnapshot = Vec<(HashedAttestationData, Vec<(u64, Validat type StorageKey = Vec; type StorageEntry = (StorageKey, Vec); type BlockRootIndexChanges = (Vec, Vec); +type VoteStore = Vec>; /// Bounded buffer for gossip signatures with FIFO eviction. /// @@ -553,6 +557,8 @@ pub struct Store { backend: Arc, new_payloads: Arc>, known_payloads: Arc>, + /// Latest fork-choice vote per validator, independent from bounded proof buffers. + votes: Arc>, /// In-memory gossip signatures, consumed at interval 2 aggregation. gossip_signatures: Arc>, /// LRU memoization of states by block root, shared across `Store` clones. @@ -566,6 +572,10 @@ fn new_state_cache() -> Arc>> { Arc::new(Mutex::new(LruCache::new(capacity))) } +fn new_vote_store() -> Arc> { + Arc::new(Mutex::new(Vec::new())) +} + impl Store { /// Initialize a Store from an anchor state only. /// @@ -638,6 +648,7 @@ impl Store { backend, new_payloads: Arc::new(Mutex::new(PayloadBuffer::new(NEW_PAYLOAD_CAP))), known_payloads: Arc::new(Mutex::new(PayloadBuffer::new(AGGREGATED_PAYLOAD_CAP))), + votes: new_vote_store(), gossip_signatures: Arc::new(Mutex::new(GossipSignatureBuffer::new( GOSSIP_SIGNATURE_CAP, ))), @@ -750,6 +761,7 @@ impl Store { backend, new_payloads: Arc::new(Mutex::new(PayloadBuffer::new(NEW_PAYLOAD_CAP))), known_payloads: Arc::new(Mutex::new(PayloadBuffer::new(AGGREGATED_PAYLOAD_CAP))), + votes: new_vote_store(), gossip_signatures: Arc::new(Mutex::new(GossipSignatureBuffer::new( GOSSIP_SIGNATURE_CAP, ))), @@ -880,11 +892,16 @@ impl Store { let pruned_sigs = self.prune_gossip_signatures(finalized.slot); let pruned_payloads = self.prune_stale_aggregated_payloads(finalized.slot); + let pruned_votes = self.prune_known_votes(finalized.slot); - if pruned_chain > 0 || pruned_sigs > 0 || pruned_payloads > 0 { + if pruned_chain > 0 || pruned_sigs > 0 || pruned_payloads > 0 || pruned_votes > 0 { info!( finalized_slot = finalized.slot, - pruned_chain, pruned_sigs, pruned_payloads, "Pruned finalized data" + pruned_chain, + pruned_sigs, + pruned_payloads, + pruned_votes, + "Pruned finalized data" ); } } @@ -1063,6 +1080,22 @@ impl Store { pruned_new + pruned_known } + /// Prune latest fork-choice votes whose target is at or below `finalized_slot`. + pub fn prune_known_votes(&mut self, finalized_slot: u64) -> usize { + let mut votes = self.votes.lock().unwrap(); + let mut pruned = 0; + for vote in votes.iter_mut() { + let should_prune = vote + .as_ref() + .is_some_and(|data| data.target.slot <= finalized_slot); + if should_prune { + *vote = None; + pruned += 1; + } + } + pruned + } + /// Prune signatures of old finalized blocks, keeping a recent window. /// /// Signatures within [`SIGNATURE_PRUNING_RANGE`] slots of `tip_slot` are @@ -1418,12 +1451,45 @@ impl Store { // ============ Attestation Extraction ============ - /// Extract per-validator latest attestations from known (fork-choice-active) payloads. + fn should_replace_vote(existing: &AttestationData, candidate: &AttestationData) -> bool { + candidate.slot > existing.slot + || (candidate.slot == existing.slot + && candidate.hash_tree_root() > existing.hash_tree_root()) + } + + fn record_vote(votes: &mut VoteStore, validator_id: u64, data: &AttestationData) { + let index = validator_id as usize; + if votes.len() <= index { + votes.resize_with(index + 1, || None); + } + + let should_replace = votes[index] + .as_ref() + .is_none_or(|existing| Self::should_replace_vote(existing, data)); + if should_replace { + votes[index] = Some(data.clone()); + } + } + + fn record_known_votes(&self, data: &AttestationData, validator_ids: I) + where + I: IntoIterator, + { + let mut votes = self.votes.lock().unwrap(); + for validator_id in validator_ids { + Self::record_vote(&mut votes, validator_id, data); + } + } + + /// Extract per-validator latest attestations from known fork-choice votes. pub fn extract_latest_known_attestations(&self) -> HashMap { - self.known_payloads + self.votes .lock() .unwrap() - .extract_latest_attestations() + .iter() + .enumerate() + .filter_map(|(validator_id, vote)| vote.clone().map(|data| (validator_id as u64, data))) + .collect() } /// Extract per-validator latest attestations from new (pending) payloads. @@ -1447,6 +1513,16 @@ impl Store { .extract_latest_attestations() } + /// Insert proof-less fork-choice votes from block-included attestations. + pub fn insert_known_attestation_votes(&mut self, attestations: &[AggregatedAttestation]) { + for attestation in attestations { + self.record_known_votes( + &attestation.data, + validator_indices(&attestation.aggregation_bits), + ); + } + } + // ============ Known Aggregated Payloads ============ // // "Known" aggregated payloads are active in fork choice weight calculations. @@ -1514,6 +1590,7 @@ impl Store { hashed: HashedAttestationData, proof: SingleMessageAggregate, ) { + self.record_known_votes(hashed.data(), proof.participant_indices()); self.known_payloads.lock().unwrap().push(hashed, proof); } @@ -1522,6 +1599,9 @@ impl Store { &mut self, entries: Vec<(HashedAttestationData, SingleMessageAggregate)>, ) { + for (hashed, proof) in &entries { + self.record_known_votes(hashed.data(), proof.participant_indices()); + } self.known_payloads.lock().unwrap().push_batch(entries); } @@ -1554,6 +1634,9 @@ impl Store { /// Drains the new buffer and pushes all entries into the known buffer. pub fn promote_new_aggregated_payloads(&mut self) { let drained = self.new_payloads.lock().unwrap().drain(); + for (hashed, proof) in &drained { + self.record_known_votes(hashed.data(), proof.participant_indices()); + } self.known_payloads.lock().unwrap().push_batch(drained); } @@ -1807,6 +1890,7 @@ mod tests { backend, new_payloads: Arc::new(Mutex::new(PayloadBuffer::new(NEW_PAYLOAD_CAP))), known_payloads: Arc::new(Mutex::new(PayloadBuffer::new(AGGREGATED_PAYLOAD_CAP))), + votes: new_vote_store(), gossip_signatures: Arc::new(Mutex::new(GossipSignatureBuffer::new( GOSSIP_SIGNATURE_CAP, ))), @@ -1821,6 +1905,7 @@ mod tests { backend, new_payloads: Arc::new(Mutex::new(PayloadBuffer::new(NEW_PAYLOAD_CAP))), known_payloads: Arc::new(Mutex::new(PayloadBuffer::new(AGGREGATED_PAYLOAD_CAP))), + votes: new_vote_store(), gossip_signatures: Arc::new(Mutex::new(GossipSignatureBuffer::new( GOSSIP_SIGNATURE_CAP, ))), @@ -2597,6 +2682,58 @@ mod tests { assert_eq!(store.known_aggregated_payloads_count(), 1); } + #[test] + fn known_votes_survive_payload_fifo_eviction() { + let mut store = Store::test_store(); + let vote = make_att_data_for_target(100, root(100)); + let vote_root = vote.hash_tree_root(); + + store.insert_known_aggregated_payload( + HashedAttestationData::new(vote.clone()), + make_proof_for_validator(0), + ); + + for i in 0..=AGGREGATED_PAYLOAD_CAP { + let slot = i as u64 + 1; + let data = make_att_data_for_target(slot, root(1_000 + slot)); + store.insert_known_aggregated_payload( + HashedAttestationData::new(data), + make_proof_for_validator(1), + ); + } + + assert!( + !store + .known_payloads + .lock() + .unwrap() + .data + .contains_key(&vote_root) + ); + assert_eq!(store.extract_latest_known_attestations()[&0], vote); + } + + #[test] + fn prune_known_votes_drops_finalized_targets() { + let mut store = Store::test_store(); + let stale = make_att_data_for_target(2, root(2)); + let fresh = make_att_data_for_target(10, root(10)); + + store.insert_known_aggregated_payload( + HashedAttestationData::new(stale), + make_proof_for_validator(0), + ); + store.insert_known_aggregated_payload( + HashedAttestationData::new(fresh.clone()), + make_proof_for_validator(1), + ); + + assert_eq!(store.prune_known_votes(5), 1); + let votes = store.extract_latest_known_attestations(); + assert!(!votes.contains_key(&0)); + assert_eq!(votes[&1], fresh); + } + /// Build an attestation message at `slot` whose target points at `target_root`, /// distinct from the default zero target so two such datas have different roots. fn make_att_data_for_target(slot: u64, target_root: H256) -> AttestationData { From 87d8b0d1d2361b428108a1752e11ef9e38574c26 Mon Sep 17 00:00:00 2001 From: dicethedev Date: Wed, 5 Aug 2026 11:56:52 +0100 Subject: [PATCH 2/3] fix(storage): keep fork-choice votes independent --- crates/blockchain/src/store.rs | 1 - crates/storage/src/store.rs | 221 +++++++++++++++++++-------------- 2 files changed, 126 insertions(+), 96 deletions(-) diff --git a/crates/blockchain/src/store.rs b/crates/blockchain/src/store.rs index e5d75f45..1a6d3884 100644 --- a/crates/blockchain/src/store.rs +++ b/crates/blockchain/src/store.rs @@ -699,7 +699,6 @@ fn on_block_core( store .insert_state(block_root, post_state) .expect("DB insert should succeed"); - store.insert_known_attestation_votes(&block.body.attestations); // Block-included attestations are intentionally not counted here. // `lean_attestations_valid_total` tracks the gossip validation pipeline diff --git a/crates/storage/src/store.rs b/crates/storage/src/store.rs index c3143c94..77798f62 100644 --- a/crates/storage/src/store.rs +++ b/crates/storage/src/store.rs @@ -307,6 +307,7 @@ impl PayloadBuffer { /// broken toward the larger canonical attestation-data root — the same rule the /// block-level fork-choice tiebreak applies to block roots (leanSpec #1181). The /// pool key is already `hash_tree_root(data)`, so the tie needs no extra hashing. + #[cfg(test)] fn extract_latest_attestations(&self) -> HashMap { let mut ordered: Vec<(&H256, &PayloadEntry)> = self.data.iter().collect(); // Descending by (slot, data_root): the larger tuple is the canonical winner. @@ -342,7 +343,13 @@ pub type GossipSignatureSnapshot = Vec<(HashedAttestationData, Vec<(u64, Validat type StorageKey = Vec; type StorageEntry = (StorageKey, Vec); type BlockRootIndexChanges = (Vec, Vec); -type VoteStore = Vec>; +type VoteStore = HashMap; + +#[derive(Clone, Default)] +struct ForkChoiceState { + known_votes: VoteStore, + new_votes: VoteStore, +} /// Bounded buffer for gossip signatures with FIFO eviction. /// @@ -557,8 +564,8 @@ pub struct Store { backend: Arc, new_payloads: Arc>, known_payloads: Arc>, - /// Latest fork-choice vote per validator, independent from bounded proof buffers. - votes: Arc>, + /// Fork-choice votes, independent from bounded proof/signature buffers. + fork_choice: Arc>, /// In-memory gossip signatures, consumed at interval 2 aggregation. gossip_signatures: Arc>, /// LRU memoization of states by block root, shared across `Store` clones. @@ -572,10 +579,6 @@ fn new_state_cache() -> Arc>> { Arc::new(Mutex::new(LruCache::new(capacity))) } -fn new_vote_store() -> Arc> { - Arc::new(Mutex::new(Vec::new())) -} - impl Store { /// Initialize a Store from an anchor state only. /// @@ -648,7 +651,7 @@ impl Store { backend, new_payloads: Arc::new(Mutex::new(PayloadBuffer::new(NEW_PAYLOAD_CAP))), known_payloads: Arc::new(Mutex::new(PayloadBuffer::new(AGGREGATED_PAYLOAD_CAP))), - votes: new_vote_store(), + fork_choice: Default::default(), gossip_signatures: Arc::new(Mutex::new(GossipSignatureBuffer::new( GOSSIP_SIGNATURE_CAP, ))), @@ -761,7 +764,7 @@ impl Store { backend, new_payloads: Arc::new(Mutex::new(PayloadBuffer::new(NEW_PAYLOAD_CAP))), known_payloads: Arc::new(Mutex::new(PayloadBuffer::new(AGGREGATED_PAYLOAD_CAP))), - votes: new_vote_store(), + fork_choice: Default::default(), gossip_signatures: Arc::new(Mutex::new(GossipSignatureBuffer::new( GOSSIP_SIGNATURE_CAP, ))), @@ -892,16 +895,11 @@ impl Store { let pruned_sigs = self.prune_gossip_signatures(finalized.slot); let pruned_payloads = self.prune_stale_aggregated_payloads(finalized.slot); - let pruned_votes = self.prune_known_votes(finalized.slot); - if pruned_chain > 0 || pruned_sigs > 0 || pruned_payloads > 0 || pruned_votes > 0 { + if pruned_chain > 0 || pruned_sigs > 0 || pruned_payloads > 0 { info!( finalized_slot = finalized.slot, - pruned_chain, - pruned_sigs, - pruned_payloads, - pruned_votes, - "Pruned finalized data" + pruned_chain, pruned_sigs, pruned_payloads, "Pruned finalized data" ); } } @@ -1080,22 +1078,6 @@ impl Store { pruned_new + pruned_known } - /// Prune latest fork-choice votes whose target is at or below `finalized_slot`. - pub fn prune_known_votes(&mut self, finalized_slot: u64) -> usize { - let mut votes = self.votes.lock().unwrap(); - let mut pruned = 0; - for vote in votes.iter_mut() { - let should_prune = vote - .as_ref() - .is_some_and(|data| data.target.slot <= finalized_slot); - if should_prune { - *vote = None; - pruned += 1; - } - } - pruned - } - /// Prune signatures of old finalized blocks, keeping a recent window. /// /// Signatures within [`SIGNATURE_PRUNING_RANGE`] slots of `tip_slot` are @@ -1203,6 +1185,7 @@ impl Store { .expect("put non-finalized chain index"); batch.commit().expect("commit"); + self.record_known_attestation_votes(&block.body.attestations); Ok(()) } @@ -1458,46 +1441,58 @@ impl Store { } fn record_vote(votes: &mut VoteStore, validator_id: u64, data: &AttestationData) { - let index = validator_id as usize; - if votes.len() <= index { - votes.resize_with(index + 1, || None); - } - - let should_replace = votes[index] - .as_ref() + let should_replace = votes + .get(&validator_id) .is_none_or(|existing| Self::should_replace_vote(existing, data)); if should_replace { - votes[index] = Some(data.clone()); + votes.insert(validator_id, data.clone()); } } - fn record_known_votes(&self, data: &AttestationData, validator_ids: I) + fn record_votes(votes: &mut VoteStore, data: &AttestationData, validator_ids: I) where I: IntoIterator, { - let mut votes = self.votes.lock().unwrap(); for validator_id in validator_ids { - Self::record_vote(&mut votes, validator_id, data); + Self::record_vote(votes, validator_id, data); + } + } + + fn record_known_votes(&self, data: &AttestationData, validator_ids: I) + where + I: IntoIterator, + { + let mut fork_choice = self.fork_choice.lock().unwrap(); + Self::record_votes(&mut fork_choice.known_votes, data, validator_ids); + } + + fn record_new_votes(&self, data: &AttestationData, validator_ids: I) + where + I: IntoIterator, + { + let mut fork_choice = self.fork_choice.lock().unwrap(); + Self::record_votes(&mut fork_choice.new_votes, data, validator_ids); + } + + fn record_known_attestation_votes(&self, attestations: &[AggregatedAttestation]) { + let mut fork_choice = self.fork_choice.lock().unwrap(); + for attestation in attestations { + Self::record_votes( + &mut fork_choice.known_votes, + &attestation.data, + validator_indices(&attestation.aggregation_bits), + ); } } /// Extract per-validator latest attestations from known fork-choice votes. pub fn extract_latest_known_attestations(&self) -> HashMap { - self.votes - .lock() - .unwrap() - .iter() - .enumerate() - .filter_map(|(validator_id, vote)| vote.clone().map(|data| (validator_id as u64, data))) - .collect() + self.fork_choice.lock().unwrap().known_votes.clone() } /// Extract per-validator latest attestations from new (pending) payloads. pub fn extract_latest_new_attestations(&self) -> HashMap { - self.new_payloads - .lock() - .unwrap() - .extract_latest_attestations() + self.fork_choice.lock().unwrap().new_votes.clone() } /// Extract per-validator latest attestations from the raw gossip signature @@ -1513,16 +1508,6 @@ impl Store { .extract_latest_attestations() } - /// Insert proof-less fork-choice votes from block-included attestations. - pub fn insert_known_attestation_votes(&mut self, attestations: &[AggregatedAttestation]) { - for attestation in attestations { - self.record_known_votes( - &attestation.data, - validator_indices(&attestation.aggregation_bits), - ); - } - } - // ============ Known Aggregated Payloads ============ // // "Known" aggregated payloads are active in fork choice weight calculations. @@ -1584,16 +1569,6 @@ impl Store { self.new_payloads.lock().unwrap().attestation_data_keys() } - /// Insert a single proof into the known (fork-choice-active) buffer. - pub fn insert_known_aggregated_payload( - &mut self, - hashed: HashedAttestationData, - proof: SingleMessageAggregate, - ) { - self.record_known_votes(hashed.data(), proof.participant_indices()); - self.known_payloads.lock().unwrap().push(hashed, proof); - } - /// Batch-insert proofs into the known buffer. pub fn insert_known_aggregated_payloads_batch( &mut self, @@ -1616,6 +1591,7 @@ impl Store { hashed: HashedAttestationData, proof: SingleMessageAggregate, ) { + self.record_new_votes(hashed.data(), proof.participant_indices()); self.new_payloads.lock().unwrap().push(hashed, proof); } @@ -1624,6 +1600,9 @@ impl Store { &mut self, entries: Vec<(HashedAttestationData, SingleMessageAggregate)>, ) { + for (hashed, proof) in &entries { + self.record_new_votes(hashed.data(), proof.participant_indices()); + } self.new_payloads.lock().unwrap().push_batch(entries); } @@ -1634,8 +1613,12 @@ impl Store { /// Drains the new buffer and pushes all entries into the known buffer. pub fn promote_new_aggregated_payloads(&mut self) { let drained = self.new_payloads.lock().unwrap().drain(); - for (hashed, proof) in &drained { - self.record_known_votes(hashed.data(), proof.participant_indices()); + { + let mut fork_choice = self.fork_choice.lock().unwrap(); + let new_votes = std::mem::take(&mut fork_choice.new_votes); + for (validator_id, data) in new_votes { + Self::record_vote(&mut fork_choice.known_votes, validator_id, &data); + } } self.known_payloads.lock().unwrap().push_batch(drained); } @@ -1882,6 +1865,25 @@ mod tests { } } + fn signed_block_with_attestations( + slot: u64, + parent_root: H256, + attestations: Vec, + ) -> SignedBlock { + SignedBlock { + message: Block { + slot, + proposer_index: 0, + parent_root, + state_root: H256::ZERO, + body: BlockBody { + attestations: attestations.try_into().unwrap(), + }, + }, + proof: MultiMessageAggregate::default(), + } + } + impl Store { /// Create a Store with an in-memory backend for tests. fn test_store() -> Self { @@ -1890,7 +1892,7 @@ mod tests { backend, new_payloads: Arc::new(Mutex::new(PayloadBuffer::new(NEW_PAYLOAD_CAP))), known_payloads: Arc::new(Mutex::new(PayloadBuffer::new(AGGREGATED_PAYLOAD_CAP))), - votes: new_vote_store(), + fork_choice: Default::default(), gossip_signatures: Arc::new(Mutex::new(GossipSignatureBuffer::new( GOSSIP_SIGNATURE_CAP, ))), @@ -1905,7 +1907,7 @@ mod tests { backend, new_payloads: Arc::new(Mutex::new(PayloadBuffer::new(NEW_PAYLOAD_CAP))), known_payloads: Arc::new(Mutex::new(PayloadBuffer::new(AGGREGATED_PAYLOAD_CAP))), - votes: new_vote_store(), + fork_choice: Default::default(), gossip_signatures: Arc::new(Mutex::new(GossipSignatureBuffer::new( GOSSIP_SIGNATURE_CAP, ))), @@ -2001,6 +2003,29 @@ mod tests { assert_eq!(blocks[0].message.hash_tree_root(), block_root); } + #[test] + fn insert_signed_block_records_block_attestation_votes() { + let mut store = Store::test_store(); + let data = make_att_data_for_target(8, root(8)); + let block = signed_block_with_attestations( + 1, + H256::ZERO, + vec![AggregatedAttestation { + aggregation_bits: make_proof_for_validators(&[1, 3]).participants, + data: data.clone(), + }], + ); + let block_root = block.message.hash_tree_root(); + + store + .insert_signed_block(block_root, block) + .expect("insert signed block"); + + let votes = store.extract_latest_known_attestations(); + assert_eq!(votes[&1], data); + assert_eq!(votes[&3], data); + } + #[test] fn prune_old_blocks_within_retention() { let backend = Arc::new(InMemoryBackend::new()); @@ -2320,17 +2345,23 @@ mod tests { make_proof_for_validator(0), ); store.insert_new_aggregated_payload( - HashedAttestationData::new(data), + HashedAttestationData::new(data.clone()), make_proof_for_validator(1), ); assert_eq!(store.new_payloads.lock().unwrap().len(), 1); assert_eq!(store.known_payloads.lock().unwrap().len(), 0); + assert_eq!(store.extract_latest_new_attestations()[&0], data); + assert_eq!(store.extract_latest_new_attestations()[&1], data); + assert!(store.extract_latest_known_attestations().is_empty()); store.promote_new_aggregated_payloads(); assert_eq!(store.new_payloads.lock().unwrap().len(), 0); assert_eq!(store.known_payloads.lock().unwrap().len(), 1); + assert!(store.extract_latest_new_attestations().is_empty()); + assert_eq!(store.extract_latest_known_attestations()[&0], data); + assert_eq!(store.extract_latest_known_attestations()[&1], data); // The known buffer should have 2 proofs for this data assert_eq!( store.known_payloads.lock().unwrap().data[&data_root] @@ -2659,18 +2690,18 @@ mod tests { HashedAttestationData::new(stale.clone()), make_proof_for_validators(&[0]), ); - store.insert_known_aggregated_payload( + store.insert_known_aggregated_payloads_batch(vec![( HashedAttestationData::new(stale), make_proof_for_validators(&[1]), - ); + )]); store.insert_new_aggregated_payload( HashedAttestationData::new(fresh.clone()), make_proof_for_validators(&[2]), ); - store.insert_known_aggregated_payload( + store.insert_known_aggregated_payloads_batch(vec![( HashedAttestationData::new(fresh), make_proof_for_validators(&[3]), - ); + )]); assert_eq!(store.new_aggregated_payloads_count(), 2); assert_eq!(store.known_aggregated_payloads_count(), 2); @@ -2688,18 +2719,18 @@ mod tests { let vote = make_att_data_for_target(100, root(100)); let vote_root = vote.hash_tree_root(); - store.insert_known_aggregated_payload( + store.insert_known_aggregated_payloads_batch(vec![( HashedAttestationData::new(vote.clone()), make_proof_for_validator(0), - ); + )]); for i in 0..=AGGREGATED_PAYLOAD_CAP { let slot = i as u64 + 1; let data = make_att_data_for_target(slot, root(1_000 + slot)); - store.insert_known_aggregated_payload( + store.insert_known_aggregated_payloads_batch(vec![( HashedAttestationData::new(data), make_proof_for_validator(1), - ); + )]); } assert!( @@ -2714,23 +2745,23 @@ mod tests { } #[test] - fn prune_known_votes_drops_finalized_targets() { + fn known_votes_survive_finalized_payload_pruning() { let mut store = Store::test_store(); let stale = make_att_data_for_target(2, root(2)); let fresh = make_att_data_for_target(10, root(10)); - store.insert_known_aggregated_payload( - HashedAttestationData::new(stale), + store.insert_known_aggregated_payloads_batch(vec![( + HashedAttestationData::new(stale.clone()), make_proof_for_validator(0), - ); - store.insert_known_aggregated_payload( + )]); + store.insert_known_aggregated_payloads_batch(vec![( HashedAttestationData::new(fresh.clone()), make_proof_for_validator(1), - ); + )]); - assert_eq!(store.prune_known_votes(5), 1); + assert_eq!(store.prune_stale_aggregated_payloads(5), 1); let votes = store.extract_latest_known_attestations(); - assert!(!votes.contains_key(&0)); + assert_eq!(votes[&0], stale); assert_eq!(votes[&1], fresh); } From 8e666ef5d2142ff724ef963be9a3b4fd33b9b150 Mon Sep 17 00:00:00 2001 From: dicethedev Date: Sun, 9 Aug 2026 13:32:07 +0100 Subject: [PATCH 3/3] fix(storage): inline fork-choice vote recording --- crates/storage/src/store.rs | 165 ++++++++++-------------------------- 1 file changed, 47 insertions(+), 118 deletions(-) diff --git a/crates/storage/src/store.rs b/crates/storage/src/store.rs index 133c3690..52744bc5 100644 --- a/crates/storage/src/store.rs +++ b/crates/storage/src/store.rs @@ -300,32 +300,6 @@ impl PayloadBuffer { } pruned } - - /// Extract per-validator latest attestations from proofs' participation bits. - /// - /// An equivocator can cast two distinct votes at the same `slot`. To keep the - /// extracted head a pure function of pool contents (independent of arrival or - /// insertion order), votes are processed newest-first with an equal-slot tie - /// broken toward the larger canonical attestation-data root — the same rule the - /// block-level fork-choice tiebreak applies to block roots (leanSpec #1181). The - /// pool key is already `hash_tree_root(data)`, so the tie needs no extra hashing. - #[cfg(test)] - fn extract_latest_attestations(&self) -> HashMap { - let mut ordered: Vec<(&H256, &PayloadEntry)> = self.data.iter().collect(); - // Descending by (slot, data_root): the larger tuple is the canonical winner. - ordered.sort_unstable_by(|a, b| (b.1.data.slot, b.0).cmp(&(a.1.data.slot, a.0))); - - let mut result: HashMap = HashMap::new(); - for (_data_root, entry) in ordered { - for proof in &entry.proofs { - for vid in proof.participant_indices() { - // Descending order means the first vote seen for a validator wins. - result.entry(vid).or_insert_with(|| entry.data.clone()); - } - } - } - result - } } /// Gossip signatures grouped by attestation data. @@ -492,10 +466,9 @@ impl GossipSignatureBuffer { /// Extract per-validator latest attestations from the raw signature pool. /// - /// Mirrors `PayloadBuffer::extract_latest_attestations`: votes are processed - /// newest-first with an equal-slot tie broken toward the larger canonical - /// attestation-data root, so the extracted winner is independent of arrival or - /// insertion order (leanSpec #1181). This matches the leanSpec + /// Votes are processed newest-first with an equal-slot tie broken toward the + /// larger canonical attestation-data root, so the extracted winner is + /// independent of arrival or insertion order (leanSpec #1181). This matches the leanSpec /// `location == "signatures"` checker, which folds `attestation_signatures` /// keeping each validator's canonical-precedence winner. fn extract_latest_attestations(&self) -> HashMap { @@ -1489,39 +1462,16 @@ impl Store { } } - fn record_votes(votes: &mut VoteStore, data: &AttestationData, validator_ids: I) - where - I: IntoIterator, - { - for validator_id in validator_ids { - Self::record_vote(votes, validator_id, data); - } - } - - fn record_known_votes(&self, data: &AttestationData, validator_ids: I) - where - I: IntoIterator, - { - let mut fork_choice = self.fork_choice.lock().unwrap(); - Self::record_votes(&mut fork_choice.known_votes, data, validator_ids); - } - - fn record_new_votes(&self, data: &AttestationData, validator_ids: I) - where - I: IntoIterator, - { - let mut fork_choice = self.fork_choice.lock().unwrap(); - Self::record_votes(&mut fork_choice.new_votes, data, validator_ids); - } - fn record_known_attestation_votes(&self, attestations: &[AggregatedAttestation]) { let mut fork_choice = self.fork_choice.lock().unwrap(); for attestation in attestations { - Self::record_votes( - &mut fork_choice.known_votes, - &attestation.data, - validator_indices(&attestation.aggregation_bits), - ); + for validator_id in validator_indices(&attestation.aggregation_bits) { + Self::record_vote( + &mut fork_choice.known_votes, + validator_id, + &attestation.data, + ); + } } } @@ -1614,8 +1564,11 @@ impl Store { &mut self, entries: Vec<(HashedAttestationData, SingleMessageAggregate)>, ) { + let mut fork_choice = self.fork_choice.lock().unwrap(); for (hashed, proof) in &entries { - self.record_known_votes(hashed.data(), proof.participant_indices()); + for validator_id in proof.participant_indices() { + Self::record_vote(&mut fork_choice.known_votes, validator_id, hashed.data()); + } } self.known_payloads.lock().unwrap().push_batch(entries); } @@ -1631,7 +1584,12 @@ impl Store { hashed: HashedAttestationData, proof: SingleMessageAggregate, ) { - self.record_new_votes(hashed.data(), proof.participant_indices()); + { + let mut fork_choice = self.fork_choice.lock().unwrap(); + for validator_id in proof.participant_indices() { + Self::record_vote(&mut fork_choice.new_votes, validator_id, hashed.data()); + } + } self.new_payloads.lock().unwrap().push(hashed, proof); } @@ -1640,8 +1598,11 @@ impl Store { &mut self, entries: Vec<(HashedAttestationData, SingleMessageAggregate)>, ) { + let mut fork_choice = self.fork_choice.lock().unwrap(); for (hashed, proof) in &entries { - self.record_new_votes(hashed.data(), proof.participant_indices()); + for validator_id in proof.participant_indices() { + Self::record_vote(&mut fork_choice.new_votes, validator_id, hashed.data()); + } } self.new_payloads.lock().unwrap().push_batch(entries); } @@ -2070,6 +2031,29 @@ mod tests { assert_eq!(blocks[0].message.hash_tree_root(), block_root); } + #[test] + fn insert_signed_block_records_block_attestation_votes() { + let mut store = Store::test_store(); + let data = make_att_data_for_target(8, root(8)); + let block = signed_block_with_attestations( + 1, + H256::ZERO, + vec![AggregatedAttestation { + aggregation_bits: make_proof_for_validators(&[1, 3]).participants, + data: data.clone(), + }], + ); + let block_root = block.message.hash_tree_root(); + + store + .insert_signed_block(block_root, block) + .expect("insert signed block"); + + let votes = store.extract_latest_known_attestations(); + assert_eq!(votes[&1], data); + assert_eq!(votes[&3], data); + } + #[test] fn prune_old_block_proofs_within_retention() { let backend = Arc::new(InMemoryBackend::new()); @@ -2823,61 +2807,6 @@ mod tests { } } - /// When two aggregations share `slot` but disagree on the target - /// (same-slot equivocation), the vote with the larger canonical - /// attestation-data root must win for the validators present in both, - /// regardless of arrival/insertion order. This is the deterministic - /// fork-choice tiebreak (leanSpec #1181): it makes the extracted head a pure - /// function of pool contents, so two nodes that see the same votes in - /// different orders agree on the same head. - #[test] - fn extract_latest_attestations_canonical_root_wins_on_slot_tie() { - let target_a = H256([0xaa; 32]); - let target_b = H256([0xbb; 32]); - let data_a = make_att_data_for_target(3, target_a); - let data_b = make_att_data_for_target(3, target_b); - // The pool keys the tie on `hash_tree_root(data)`, not on the target root. - let root_a = data_a.hash_tree_root(); - let root_b = data_b.hash_tree_root(); - assert_ne!(root_a, root_b); - - // The larger canonical root is the deterministic winner for shared voters. - let winner_target = if root_a > root_b { target_a } else { target_b }; - - // Deliver the same equivocating votes in both arrival orders; the winner - // for validators present in both aggregations (0, 1) must be identical. - for insert_b_first in [false, true] { - let mut buf = PayloadBuffer::new(10); - if insert_b_first { - buf.push( - HashedAttestationData::new(data_b.clone()), - make_proof_for_validators(&[0, 1, 3, 4]), - ); - buf.push( - HashedAttestationData::new(data_a.clone()), - make_proof_for_validators(&[0, 1, 2]), - ); - } else { - buf.push( - HashedAttestationData::new(data_a.clone()), - make_proof_for_validators(&[0, 1, 2]), - ); - buf.push( - HashedAttestationData::new(data_b.clone()), - make_proof_for_validators(&[0, 1, 3, 4]), - ); - } - let extracted = buf.extract_latest_attestations(); - // Shared validators: deterministic canonical winner, independent of order. - assert_eq!(extracted[&0].target.root, winner_target); - assert_eq!(extracted[&1].target.root, winner_target); - // Exclusive validators keep their only vote. - assert_eq!(extracted[&2].target.root, target_a); - assert_eq!(extracted[&3].target.root, target_b); - assert_eq!(extracted[&4].target.root, target_b); - } - } - /// `drain` must hand back entries in insertion order so that /// `promote_new_aggregated_payloads` lands them in known_payloads in the /// same order, preserving same-slot equivocation semantics through the