-
Notifications
You must be signed in to change notification settings - Fork 28
fix(storage): store fork-choice votes independently #552
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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<u8>; | ||
| type StorageEntry = (StorageKey, Vec<u8>); | ||
| type BlockRootIndexChanges = (Vec<StorageKey>, Vec<StorageEntry>); | ||
| type VoteStore = Vec<Option<AttestationData>>; | ||
|
|
||
| /// Bounded buffer for gossip signatures with FIFO eviction. | ||
| /// | ||
|
|
@@ -553,6 +557,8 @@ pub struct Store { | |
| backend: Arc<dyn StorageBackend>, | ||
| new_payloads: Arc<Mutex<PayloadBuffer>>, | ||
| known_payloads: Arc<Mutex<PayloadBuffer>>, | ||
| /// Latest fork-choice vote per validator, independent from bounded proof buffers. | ||
| votes: Arc<Mutex<VoteStore>>, | ||
|
Comment on lines
+560
to
+561
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Let's wrap the Also, another thing I forgot about is that we have to store both known and new votes (each in a separate map). Here we are storing known votes only. |
||
| /// In-memory gossip signatures, consumed at interval 2 aggregation. | ||
| gossip_signatures: Arc<Mutex<GossipSignatureBuffer>>, | ||
| /// LRU memoization of states by block root, shared across `Store` clones. | ||
|
|
@@ -566,6 +572,10 @@ fn new_state_cache() -> Arc<Mutex<LruCache<H256, State>>> { | |
| Arc::new(Mutex::new(LruCache::new(capacity))) | ||
| } | ||
|
|
||
| fn new_vote_store() -> Arc<Mutex<VoteStore>> { | ||
| Arc::new(Mutex::new(Vec::new())) | ||
| } | ||
|
Comment on lines
+575
to
+577
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This is just |
||
|
|
||
| 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 { | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Pruning isn't needed for votes, since we always keep the last vote of each validator only |
||
| 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); | ||
|
Comment on lines
+1087
to
+1090
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
When a validator's latest attestation targets a finalized checkpoint but its head remains on a live post-finalization branch, Knowledge Base Used: Prompt To Fix With AIThis is a comment left during a code review.
Path: crates/storage/src/store.rs
Line: 1087-1090
Comment:
**Target-based pruning drops live votes**
When a validator's latest attestation targets a finalized checkpoint but its head remains on a live post-finalization branch, `prune_known_votes` removes the vote based solely on `target.slot`. LMD-GHOST weights `head.root`, so subsequent head calculations lose that validator's weight and can select the wrong branch until a replacement vote arrives.
**Knowledge Base Used:**
- [Blockchain core: fork choice, state transition, block building, sync](https://app.greptile.com/lambdaclass/-/custom-context/knowledge-base/lambdaclass/ethlambda/-/docs/blockchain-core.md)
- [Storage (`ethlambda_storage`)](https://app.greptile.com/lambdaclass/-/custom-context/knowledge-base/lambdaclass/ethlambda/-/docs/storage.md)
---
For each issue above, determine whether it is valid and should be fixed. If so, fix it directly. |
||
| 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<I>(&self, data: &AttestationData, validator_ids: I) | ||
| where | ||
| I: IntoIterator<Item = u64>, | ||
| { | ||
| 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<u64, AttestationData> { | ||
| 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]) { | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Can we embed this inside |
||
| 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 { | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Let's make this a
HashMap<usize, AttestationData>instead. That'll make this easier to reason about, and avoid performance edge-cases.