From 28e550229d3480b9529a9e5d72dec62c74db50de Mon Sep 17 00:00:00 2001 From: dicethedev Date: Fri, 31 Jul 2026 13:38:11 +0100 Subject: [PATCH 1/2] 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/2] 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); }