Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions crates/blockchain/src/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
149 changes: 143 additions & 6 deletions crates/storage/src/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
},
Expand Down Expand Up @@ -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>>;
Comment thread
dicethedev marked this conversation as resolved.
Outdated

/// Bounded buffer for gossip signatures with FIFO eviction.
///
Expand Down Expand Up @@ -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 thread
dicethedev marked this conversation as resolved.
Outdated
/// 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.
Expand All @@ -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 thread
dicethedev marked this conversation as resolved.
Outdated

impl Store {
/// Initialize a Store from an anchor state only.
///
Expand Down Expand Up @@ -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,
))),
Expand Down Expand Up @@ -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,
))),
Expand Down Expand Up @@ -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"
);
}
}
Expand Down Expand Up @@ -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 {
Comment thread
dicethedev marked this conversation as resolved.
Outdated
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 thread
MegaRedHand marked this conversation as resolved.
Outdated
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
Expand Down Expand Up @@ -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.
Expand All @@ -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]) {
Comment thread
dicethedev marked this conversation as resolved.
Outdated
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.
Expand Down Expand Up @@ -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);
}

Expand All @@ -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);
}

Expand Down Expand Up @@ -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);
}

Expand Down Expand Up @@ -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,
))),
Expand All @@ -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,
))),
Expand Down Expand Up @@ -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 {
Expand Down