Skip to content
Merged
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
2 changes: 1 addition & 1 deletion crates/hotshot/new-protocol/bench/src/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,7 @@ impl MetricsCollector {
self.view_mut(v).vid_disperse_ns = Some(ts);
},
// Replica: proposal processing
ConsensusInput::ProposalWithVidShare(_sender, p, _) => {
ConsensusInput::Proposal(_sender, p) => {
let v = *p.view_number();
self.view_mut(v).proposal_recv_ns = Some(ts);
},
Expand Down
101 changes: 92 additions & 9 deletions crates/hotshot/new-protocol/src/consensus.rs
Original file line number Diff line number Diff line change
Expand Up @@ -102,11 +102,12 @@ pub enum ConsensusInput<T: NodeType> {
},
EpochChange(EpochChangeMessage<T>),
HeaderCreated(ViewNumber, Commitment<Leaf2<T>>, T::BlockHeader),
ProposalWithVidShare(
T::SignatureKey,
ProposalMessage<T, Validated>,
VidDisperseShare2<T>,
),
/// A validated proposal. Consensus parks it until this node's VID share
/// for the same payload arrives ([`ConsensusInput::VidShare`]) and only
/// processes the two together.
Proposal(T::SignatureKey, ProposalMessage<T, Validated>),
/// This node's validated VID share.
VidShare(VidDisperseShare2<T>),
FetchedProposal(ProposalMessage<T, Validated>),
StateValidated(StateResponse<T>),
StateValidationFailed(StateResponse<T>),
Expand Down Expand Up @@ -156,6 +157,11 @@ pub enum ConsensusOutput<T: NodeType> {
ViewChanged(ViewNumber, EpochNumber),
/// A view timed out with a timeout certificate.
ViewTimedOut(ViewNumber),
/// A validated proposal met this node's VID share.
ProposalPaired {
proposal: SignedProposal<T, Proposal<T>>,
vid_share: VidDisperseShare2<T>,
},
ProposalValidated {
proposal: SignedProposal<T, Proposal<T>>,
sender: T::SignatureKey,
Expand All @@ -178,6 +184,13 @@ pub enum ConsensusOutput<T: NodeType> {
BroadcastVidShare(VidDisperseShare2<T>),
}

type UnpairedProposals<T> = BTreeMap<
(ViewNumber, VidCommitment2),
(<T as NodeType>::SignatureKey, ProposalMessage<T, Validated>),
>;

type UnpairedVidShares<T> = BTreeMap<(ViewNumber, VidCommitment2), VidDisperseShare2<T>>;

/// Views to retain decide inputs (`proposals`, `certs`, `certs2`) behind the
/// decided view, letting a late-broadcast Cert2 decide an older gap view.
pub(crate) const DECIDE_BUFFER: u64 = 20;
Expand All @@ -190,6 +203,8 @@ pub struct Consensus<T: NodeType> {
signed_proposals: BTreeMap<ViewNumber, SignedProposal<T, Proposal<T>>>,
proposed_views: BTreeSet<ViewNumber>,
vid_shares: BTreeMap<ViewNumber, VidDisperseShare2<T>>,
unpaired_proposals: UnpairedProposals<T>,
unpaired_vid_shares: UnpairedVidShares<T>,
states_verified: BTreeMap<ViewNumber, Commitment<Leaf2<T>>>,
blocks_reconstructed: BTreeSet<(ViewNumber, VidCommitment2)>,
blocks: BTreeMap<(ViewNumber, VidCommitment2), T::BlockPayload>,
Expand Down Expand Up @@ -369,6 +384,8 @@ impl<T: NodeType> Consensus<T> {
state_certs: BTreeMap::new(),
upgrade_lock,
vid_shares: BTreeMap::new(),
unpaired_proposals: BTreeMap::new(),
unpaired_vid_shares: BTreeMap::new(),
epoch_height: epoch_height.into(),
}
}
Expand Down Expand Up @@ -659,14 +676,18 @@ impl<T: NodeType> Consensus<T> {
input.view_number()
};
let proto = match input {
ConsensusInput::ProposalWithVidShare(sender, proposal, vid_share) => {
ConsensusInput::Proposal(sender, proposal) => {
debug!(
sender = %KeyPrefix::from(&sender),
block = %proposal.proposal.data.block_header.block_number(),
epoch = %proposal.proposal.data.epoch,
"apply: proposal+vid share"
"apply: proposal"
);
self.handle_proposal_with_vid_share(sender, proposal, vid_share, outbox)
self.pair_proposal(sender, proposal, outbox)
},
ConsensusInput::VidShare(vid_share) => {
debug!("apply: vid share");
self.pair_vid_share(vid_share, outbox)
},
ConsensusInput::FetchedProposal(message) => {
debug!(
Expand Down Expand Up @@ -946,7 +967,10 @@ impl<T: NodeType> Consensus<T> {
match scope {
GcScope::Local(view) => {
let c = Commitment::default_commitment_no_preimage();
let vc = VidCommitment2::default();
self.headers = self.headers.split_off(&(view, c));
self.unpaired_proposals = self.unpaired_proposals.split_off(&(view, vc));
self.unpaired_vid_shares = self.unpaired_vid_shares.split_off(&(view, vc));
self.proposed_views = self.proposed_views.split_off(&view);
self.states_verified = self.states_verified.split_off(&view);
self.timeout_certs = self.timeout_certs.split_off(&view);
Expand Down Expand Up @@ -1009,6 +1033,64 @@ impl<T: NodeType> Consensus<T> {
self.proposals.insert(view, proposal);
}

/// Pair a validated proposal with this node's VID share for the same payload.
///
/// The half arriving first is parked, keyed by (view, payload commitment).
fn pair_proposal(
&mut self,
sender: T::SignatureKey,
proposal: ProposalMessage<T, Validated>,
outbox: &mut Outbox<ConsensusOutput<T>>,
) -> Protocol {
let view = proposal.view_number();
let VidCommitment::V2(commit) = proposal.proposal.data.block_header.payload_commitment()
else {
warn!(%view, "proposal payload commitment is not V2, discarding");
return Protocol::Abort;
};
let Some(vid_share) = self.unpaired_vid_shares.remove(&(view, commit)) else {
self.unpaired_proposals
.insert((view, commit), (sender, proposal));
return Protocol::Abort;
};
self.on_proposal_paired(sender, proposal, vid_share, outbox)
}

/// Pair this node's VID share with a validated proposal for the same payload.
///
/// The half arriving first is parked, keyed by (view, payload commitment).
fn pair_vid_share(
&mut self,
vid_share: VidDisperseShare2<T>,
outbox: &mut Outbox<ConsensusOutput<T>>,
) -> Protocol {
let key = (vid_share.view_number(), vid_share.payload_commitment);
let Some((sender, proposal)) = self.unpaired_proposals.remove(&key) else {
self.unpaired_vid_shares.insert(key, vid_share);
return Protocol::Abort;
};
self.on_proposal_paired(sender, proposal, vid_share, outbox)
}

fn on_proposal_paired(
&mut self,
sender: T::SignatureKey,
proposal: ProposalMessage<T, Validated>,
vid_share: VidDisperseShare2<T>,
outbox: &mut Outbox<ConsensusOutput<T>>,
) -> Protocol {
// Parked halves for this and older views can no longer pair.
let view = proposal.view_number();
let vc = VidCommitment2::default();
self.unpaired_proposals = self.unpaired_proposals.split_off(&(view + 1, vc));
self.unpaired_vid_shares = self.unpaired_vid_shares.split_off(&(view + 1, vc));
outbox.push_back(ConsensusOutput::ProposalPaired {
proposal: proposal.proposal.clone(),
vid_share: vid_share.clone(),
});
self.handle_proposal_with_vid_share(sender, proposal, vid_share, outbox)
}

#[instrument(level = "debug", skip_all)]
fn handle_proposal_with_vid_share(
&mut self,
Expand Down Expand Up @@ -2625,7 +2707,8 @@ impl<T: NodeType> ConsensusInput<T> {
ConsensusInput::AdvanceView(cert) => cert.view_number() + 1,
ConsensusInput::EpochRootCertificates { cert1, .. } => cert1.view_number(),
ConsensusInput::HeaderCreated(view, ..) => *view,
ConsensusInput::ProposalWithVidShare(_, prop, _) => prop.view_number(),
ConsensusInput::Proposal(_, prop) => prop.view_number(),
ConsensusInput::VidShare(share) => share.view_number(),
ConsensusInput::FetchedProposal(prop) => prop.view_number(),
ConsensusInput::StateValidated(response) => response.view,
ConsensusInput::StateValidationFailed(request) => request.view,
Expand Down
108 changes: 29 additions & 79 deletions crates/hotshot/new-protocol/src/coordinator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ use hotshot::{HotShotInitializer, traits::BlockPayload, types::SignatureKey};
use hotshot_types::{
consensus::{ConsensusMetricsValue, ParticipationTracker},
data::{
EpochNumber, Leaf2, VidCommitment, VidCommitment2, VidDisperseShare2, ViewNumber,
EpochNumber, Leaf2, VidCommitment, VidCommitment2, ViewNumber,
vid_disperse::vid_total_weight,
},
epoch_membership::EpochMembershipCoordinator,
Expand Down Expand Up @@ -51,7 +51,7 @@ use crate::{
},
network::Cliquenet,
outbox::Outbox,
proposal::{ProposalValidator, ValidatedProposal, VidShareValidator},
proposal::{ProposalValidator, VidShareValidator},
state::{HeaderRequest, StateEntry, StateManager, StateManagerOutput},
storage::{NewProtocolStorage, Storage},
vid::{VidDisperseRequest, VidDisperser, VidFragmentAccumulator, VidReconstructor},
Expand Down Expand Up @@ -112,10 +112,6 @@ pub struct Coordinator<T: NodeType, S> {
pending_proposal_fetches: PendingProposalFetches<T>,
#[builder(skip)]
requested_missing_proposals: HashSet<ProposalFetchKey<T>>,
#[builder(default)]
cached_validated_proposals: BTreeMap<(ViewNumber, VidCommitment2), ValidatedProposal<T>>,
#[builder(default)]
cached_vid_shares: BTreeMap<(ViewNumber, VidCommitment2), VidDisperseShare2<T>>,
#[builder(skip)]
da_payloads: BTreeMap<(ViewNumber, VidCommitment2), PendingDa<T>>,
metrics: Option<metrics::Metrics>,
Expand Down Expand Up @@ -475,14 +471,7 @@ where
}
Some(item) = self.share_validator.next() => match item {
Ok(vid_share) => {
let view = vid_share.view_number();
let key = (view, vid_share.payload_commitment);
let Some(validated) = self.cached_validated_proposals.remove(&key) else {
// Wait for the proposal
self.cached_vid_shares.insert(key, vid_share);
continue;
};
return self.on_proposal_and_vid_share(validated, vid_share)
return Ok(ConsensusInput::VidShare(vid_share))
},
Err(e) => {
return Err(CoordinatorError::regular(e).context("vid share validation"))
Expand All @@ -496,21 +485,7 @@ where
// Refresh the network's peer set when a proposal is validated.
let epoch = validated.message.proposal.data.epoch;
self.bump_network_epoch(epoch);

let view = validated.message.proposal.data.view_number();
let VidCommitment::V2(commit) =
validated.message.proposal.data.block_header.payload_commitment()
else {
warn!(%view, "proposal payload commitment is not V2, discarding");
continue;
};
let key = (view, commit);
let Some(vid_share) = self.cached_vid_shares.remove(&key) else {
// Wait for the vid share describing this payload.
self.cached_validated_proposals.insert(key, validated);
continue;
};
return self.on_proposal_and_vid_share(validated, vid_share)
return Ok(ConsensusInput::Proposal(validated.sender, validated.message))
}
Err(e) => {
return Err(CoordinatorError::regular(e).context("proposal validation"))
Expand Down Expand Up @@ -746,6 +721,31 @@ where
}
}
},
ConsensusOutput::ProposalPaired {
proposal,
vid_share,
} => {
let view = proposal.data.view_number;
debug!(%node, %view, "proposal paired with vid share");
self.storage.append_vid(vid_share.clone());
self.storage.append_proposal(proposal.data.clone());
if let Some(state_cert) = &proposal.data.state_cert {
self.storage.append_state_cert(
ViewNumber::new(state_cert.light_client_state.view_number),
state_cert.clone(),
);
}
let expected_param = self.expected_vid_param(vid_share.target_epoch);
self.vid_reconstructor.handle_proposal(
view,
vid_share.payload_commitment,
proposal.data.block_header.metadata().clone(),
proposal.data.epoch,
expected_param,
);
self.vid_reconstructor
.handle_vid_share(self.public_key.clone(), vid_share);
},
ConsensusOutput::SendProposal(proposal) => {
let view = proposal.data.view_number;
let epoch = proposal.data.epoch;
Expand Down Expand Up @@ -1318,52 +1318,6 @@ where
init_avidm_gf2_param(total_weight).ok()
}

fn on_proposal_and_vid_share(
&mut self,
validated: ValidatedProposal<T>,
vid_share: VidDisperseShare2<T>,
) -> Result<ConsensusInput<T>, CoordinatorError> {
self.storage.append_vid(vid_share.clone());
self.storage
.append_proposal(validated.message.proposal.data.clone());

if let Some(state_cert) = &validated.message.proposal.data.state_cert {
self.storage.append_state_cert(
ViewNumber::new(state_cert.light_client_state.view_number),
state_cert.clone(),
);
}

let expected_param = self.expected_vid_param(vid_share.target_epoch);
let proposal = &validated.message.proposal.data;
self.vid_reconstructor.handle_proposal(
proposal.view_number(),
vid_share.payload_commitment,
proposal.block_header.metadata().clone(),
proposal.epoch,
expected_param,
);
// This is our own share, addressed to us by the leader and already
// verified by the share validator.
self.vid_reconstructor
.handle_vid_share(self.public_key.clone(), vid_share.clone());

// GC for the cache
let view = validated.message.proposal.data.view_number();
self.cached_vid_shares = self
.cached_vid_shares
.split_off(&(view + 1, VidCommitment2::default()));
self.cached_validated_proposals = self
.cached_validated_proposals
.split_off(&(view + 1, VidCommitment2::default()));

Ok(ConsensusInput::ProposalWithVidShare(
validated.sender,
validated.message,
vid_share,
))
}

fn broadcast(
&self,
message_type: ConsensusMessage<T, Validated>,
Expand Down Expand Up @@ -1715,11 +1669,7 @@ where
self.consensus.gc(scope);
match scope {
GcScope::Local(view) => {
let vc = VidCommitment2::default();
self.block_builder.gc(view);
self.cached_validated_proposals =
self.cached_validated_proposals.split_off(&(view, vc));
self.cached_vid_shares = self.cached_vid_shares.split_off(&(view, vc));
self.vid_disperser.gc(view);
self.vid_fragment_accumulator.gc(view);
// When we enter a new view, we do not want to GC certain data
Expand Down
Loading
Loading