Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
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
383 changes: 108 additions & 275 deletions lean_client/Cargo.lock

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion lean_client/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -253,7 +253,7 @@ indexmap = "2"
http-body-util = "0.1"
http_api_utils = { git = "https://github.com/grandinetech/grandine", rev = "c4b676e3daa0ddcb86d0bcb321b7166d59f920f9" }
k256 = "0.13"
lean-multisig = { git = "https://github.com/leanEthereum/leanVM.git", rev = "a5909d18647de6aed38640c098d9177fab2bf36a" }
leanvm = { git = "https://github.com/leanEthereum/leanVM.git", rev = "48a904208d682848dac0e18ef8b01ebfc40df9ad" }
postcard = { version = "1.1.3", features = ["alloc"] }
libp2p = { git = "https://github.com/libp2p/rust-libp2p.git", rev = "91e8931e275bcd1c72791d18b09fea8b77209baf", default-features = false, features = [
'dns',
Expand Down
12 changes: 12 additions & 0 deletions lean_client/fork_choice/src/block_cache.rs
Original file line number Diff line number Diff line change
Expand Up @@ -157,6 +157,18 @@ impl BlockCache {
results.into_iter().collect()
}

pub fn prune_finalized(&mut self, finalized_slot: Slot) {
let stale: Vec<H256> = self
.blocks
.iter()
.filter(|(_, pending)| pending.slot <= finalized_slot)
.map(|(root, _)| *root)
.collect();
for root in stale {
self.remove(&root);
}
}

pub fn clear(&mut self) {
self.blocks.clear();
self.insertion_order.clear();
Expand Down
21 changes: 21 additions & 0 deletions lean_client/fork_choice/src/reaggregate.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,9 @@ pub const MAX_REAGGREGATIONS_PER_BLOCK: usize = 4;

pub struct ReaggregateContext {
pubkeys_per_component: Vec<Vec<PublicKey>>,
/// The `(message, slot)` binding of each component, index-aligned with
/// `pubkeys_per_component`: one per attestation, then the proposer's.
components: Vec<(H256, u32)>,
candidates: Vec<(H256, AggregationBits)>,
}

Expand Down Expand Up @@ -41,6 +44,7 @@ pub fn select_candidates(

let mut pubkeys_per_component: Vec<Vec<PublicKey>> =
Vec::with_capacity(attestations.len_u64() as usize + 1);
let mut components: Vec<(H256, u32)> = Vec::with_capacity(attestations.len_u64() as usize + 1);

for att in attestations.into_iter() {
let validator_ids = att.aggregation_bits.to_validator_indices();
Expand All @@ -62,6 +66,7 @@ pub fn select_candidates(
pks.push(v.attestation_pubkey.clone());
}
pubkeys_per_component.push(pks);
components.push((att.data.hash_tree_root(), att.data.slot.0 as u32));
}

let proposer = match validators.get(proposer_index) {
Expand All @@ -76,6 +81,7 @@ pub fn select_candidates(
}
};
pubkeys_per_component.push(vec![proposer.proposal_pubkey.clone()]);
components.push((block.hash_tree_root(), block.slot.0 as u32));

let latest_justified_slot = store.latest_justified.slot;
let mut candidates: Vec<(H256, AggregationBits)> = Vec::new();
Expand Down Expand Up @@ -147,6 +153,7 @@ pub fn select_candidates(

Some(ReaggregateContext {
pubkeys_per_component,
components,
candidates,
})
}
Expand All @@ -166,6 +173,7 @@ pub fn compute_recoveries(
for (data_root, participants) in context.candidates {
match signed_block.proof.split_by_message(
&pubkeys_per_component_view,
&context.components,
data_root,
log_inv_rate,
) {
Expand Down Expand Up @@ -233,6 +241,19 @@ pub fn run_sync(store: &mut Store, signed_block: &SignedBlock, log_inv_rate: usi
}

pub fn run_in_executor(store: Arc<RwLock<Store>>, signed_block: SignedBlock, log_inv_rate: usize) {
if xmss::prover_busy() {
debug!(
block_root = %signed_block.block.hash_tree_root(),
"reaggregate skipped: prover busy"
);
METRICS.get().map(|m| {
m.grandine_reaggregate_split_outcomes_total
.with_label_values(&["skipped_prover_busy"])
.inc()
});
return;
}

let context = {
let s = store.read();
let parent_validators = match s.states.get(&signed_block.block.parent_root) {
Expand Down
45 changes: 15 additions & 30 deletions lean_client/networking/src/network/service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1789,36 +1789,21 @@ where
tokio::spawn(async move {
for block in blocks {
let slot = block.block.slot.0;
match chain_sink.try_send(ChainMessage::ProcessBlock {
signed_block: block,
is_trusted: false,
should_gossip: false,
cached_post_state: None,
}) {
Ok(()) => {}
Err(tokio::sync::mpsc::error::TrySendError::Full(
_,
)) => {
warn!(
slot,
protocol = "blocks_by_range",
"Dropping RPC chunk: chain channel full"
);
METRICS.get().map(|m| {
m.lean_chain_message_drop_total
.with_label_values(&["blocks_by_range"])
.inc()
});
}
Err(
tokio::sync::mpsc::error::TrySendError::Closed(_),
) => {
warn!(
slot,
"Failed to forward range block to chain: channel closed"
);
break;
}
if let Err(err) = chain_sink
.send(ChainMessage::ProcessBlock {
signed_block: block,
is_trusted: false,
should_gossip: false,
cached_post_state: None,
})
.await
{
warn!(
slot,
?err,
"Failed to forward range block to chain"
);
break;
}
}
});
Expand Down
Loading
Loading