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/hotshot/new-protocol/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ tracing-subscriber = { workspace = true }
url = { workspace = true }
vec1 = { workspace = true }
versions = { workspace = true }
vid = { workspace = true }

[dev-dependencies]
quickcheck = { workspace = true }
Expand Down
5 changes: 0 additions & 5 deletions crates/hotshot/new-protocol/bench/src/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,6 @@ pub struct ViewMetrics {
pub is_leader: bool,
pub header_created_ns: Option<i128>,
pub block_built_ns: Option<i128>,
pub vid_disperse_ns: Option<i128>,
pub proposal_sent_ns: Option<i128>,
pub proposal_recv_ns: Option<i128>,
pub state_validated_ns: Option<i128>,
Expand Down Expand Up @@ -65,10 +64,6 @@ impl MetricsCollector {
let v = **view;
self.view_mut(v).block_built_ns = Some(ts);
},
ConsensusInput::VidDisperseCreated(view, _) => {
let v = **view;
self.view_mut(v).vid_disperse_ns = Some(ts);
},
// Replica: proposal processing
ConsensusInput::ProposalWithVidShare(_sender, p, _) => {
let v = *p.view_number();
Expand Down
19 changes: 6 additions & 13 deletions crates/hotshot/new-protocol/bench/src/node.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ use hotshot_new_protocol::{
outbox::Outbox,
proposal::{ProposalValidator, VidShareValidator},
state::StateManager,
vid::{VidDisperser, VidReconstructor},
vid::VidReconstructor,
vote::VoteCollector,
};
use hotshot_types::{
Expand Down Expand Up @@ -142,20 +142,15 @@ async fn build_coordinator(

let epoch_manager = EpochManager::new(epoch_height, membership.clone());

let vid_disperser = VidDisperser::new(
membership.clone(),
network.sender().clone(),
public_key,
private_key.clone(),
);

let vid_reconstructor = VidReconstructor::new();

let block_config = BlockBuilderConfig::default();
let block_builder = BlockBuilder::new(
instance.clone(),
membership.clone(),
block_config,
network.sender().clone(),
public_key,
private_key.clone(),
BlockBuilderConfig::default(),
upgrade_lock.clone(),
);

Expand Down Expand Up @@ -203,7 +198,6 @@ async fn build_coordinator(
.timeout_collector(timeout_collector)
.timeout_one_honest_collector(timeout_one_honest_collector)
.epoch_root_collector(epoch_root_collector)
.vid_disperser(vid_disperser)
.vid_reconstructor(vid_reconstructor)
.epoch_manager(epoch_manager)
.block_builder(block_builder)
Expand Down Expand Up @@ -290,8 +284,7 @@ async fn run_instrumented(mut coordinator: BenchCoordinator, cfg: &NodeConfig) -
let block_input = ConsensusInput::BlockBuilt {
view: req.view,
epoch: req.epoch,
payload: block.block,
metadata: block.metadata,
payload: std::sync::Arc::new(block.block),
payload_commitment: block.payload_commitment,
};
metrics.on_input(&block_input);
Expand Down
117 changes: 82 additions & 35 deletions crates/hotshot/new-protocol/src/block.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,23 +7,19 @@ use std::{
use committable::{Commitment, Committable};
use hotshot::traits::{BlockPayload, ValidatedState as _};
use hotshot_types::{
consensus::PayloadWithMetadata,
data::{
EpochNumber, Leaf2, VidCommitment, ViewNumber, vid_commitment,
vid_disperse::vid_total_weight,
},
data::{EpochNumber, Leaf2, VidCommitment, VidDisperse2, ViewNumber},
epoch_membership::EpochMembershipCoordinator,
message::UpgradeLock,
traits::{
EncodeBytes,
block_contents::{BuilderFee, Transaction},
node_implementation::NodeType,
signature_key::BuilderSignatureKey,
signature_key::{BuilderSignatureKey, SignatureKey},
},
utils::BuilderCommitment,
};
use tokio::{
task::{AbortHandle, JoinSet},
task::{AbortHandle, JoinSet, spawn_blocking},
time::sleep,
};
use tracing::{error, warn};
Expand All @@ -32,7 +28,9 @@ use crate::{
consensus::ConsensusInput,
helpers::proposal_commitment,
message::{DedupManifest, Proposal, TransactionMessage},
network::Sender,
state::HeaderRequest,
vid::fanout,
};

#[derive(Debug, thiserror::Error)]
Expand All @@ -43,6 +41,8 @@ pub enum BlockError {
StakeTableUnavailable,
#[error("builder signature failed")]
BuilderSignature,
#[error("vid dispersal failed: {0}")]
VidDisperse(String),
}

#[derive(Clone, Eq, PartialEq, Debug)]
Expand All @@ -55,7 +55,8 @@ pub struct BlockAndHeaderRequest<T: NodeType> {
pub struct BlockBuilderOutput<T: NodeType> {
pub view: ViewNumber,
pub epoch: EpochNumber,
pub payload: PayloadWithMetadata<T>,
pub payload: Arc<T::BlockPayload>,
pub metadata: <T::BlockPayload as BlockPayload<T>>::Metadata,
pub parent_proposal: Proposal<T>,
pub builder_commitment: BuilderCommitment,
pub builder_fee: BuilderFee<T>,
Expand Down Expand Up @@ -90,6 +91,9 @@ struct RetryEntry<T: NodeType> {
pub struct BlockBuilder<T: NodeType> {
instance: Arc<T::InstanceState>,
membership: EpochMembershipCoordinator<T>,
network: Sender<T>,
public_key: T::SignatureKey,
private_key: <T::SignatureKey as SignatureKey>::PrivateKey,
retry_pending: HashMap<Commitment<T::Transaction>, RetryEntry<T>>,
retry_total_bytes: u64,
leader_buffer: HashMap<Commitment<T::Transaction>, T::Transaction>,
Expand All @@ -108,15 +112,22 @@ pub struct BlockBuilder<T: NodeType> {
}

impl<T: NodeType> BlockBuilder<T> {
#[allow(clippy::too_many_arguments)]
pub fn new(
instance: Arc<T::InstanceState>,
membership: EpochMembershipCoordinator<T>,
network: Sender<T>,
public_key: T::SignatureKey,
private_key: <T::SignatureKey as SignatureKey>::PrivateKey,
config: BlockBuilderConfig,
upgrade_lock: UpgradeLock<T>,
) -> Self {
Self {
instance,
membership,
network,
public_key,
private_key,
config,
upgrade_lock,
retry_pending: HashMap::new(),
Expand All @@ -137,15 +148,18 @@ impl<T: NodeType> BlockBuilder<T> {
if self.calculations.contains_key(&(view, parent_commitment)) {
return;
}
let Ok(version) = self.upgrade_lock.version(view) else {
if self.upgrade_lock.version(view).is_err() {
warn!(%view, "unsupported version");
return;
};
}
let epoch = request.epoch;
let buffer = std::mem::take(&mut self.leader_buffer);
self.leader_total_bytes = 0;
let instance = self.instance.clone();
let membership = self.membership.clone();
let network = self.network.clone();
let public_key = self.public_key.clone();
let private_key = self.private_key.clone();

let handle = self.tasks.spawn(async move {
// Throttle empty block production: when no transactions are pending,
Expand All @@ -167,45 +181,79 @@ impl<T: NodeType> BlockBuilder<T> {
T::BlockPayload::from_transactions(txs, &validated_state, &instance)
.await
.map_err(|e| BlockError::PayloadConstruction(e.to_string()))?;
let payload: PayloadWithMetadata<T> = PayloadWithMetadata { payload, metadata };

let payload_bytes = payload.payload.encode();
let metadata_bytes = payload.metadata.encode();

let total_weight = {
let target_mem = membership
.stake_table_for_epoch(Some(epoch))
.map_err(|_| BlockError::StakeTableUnavailable)?;
vid_total_weight(target_mem.stake_table(), Some(epoch))
};
let payload_commitment = {
vid_commitment(
payload_bytes.as_ref(),
metadata_bytes.as_ref(),
total_weight,
version,
)
};
let payload_bytes = payload.encode();
let metadata_bytes = metadata.encode();
let block_size = payload_bytes.len() as u64;

let builder_commitment = payload.payload.builder_commitment(&payload.metadata);
// Erasure-code every namespace exactly once and derive the payload
// commitment from that same computation (this crate always disperses
// V2/AvidmGf2 shares). Runs on a blocking thread; the proposal is not
// gated on the fanout that follows.
let (commitment, per_bucket, param, recipients, num_namespaces) =
spawn_blocking(move || -> Result<_, BlockError> {
let params = VidDisperse2::<T>::disperse_params(
payload_bytes,
metadata_bytes.as_ref(),
&membership,
Some(epoch),
)
.map_err(|e| BlockError::VidDisperse(e.to_string()))?;
let num_namespaces = params.ns_table.len();
let (commitment, per_bucket) = fanout::encode(&params)
.map_err(|e| BlockError::VidDisperse(e.to_string()))?;
Ok((
commitment,
per_bucket,
params.param,
params.recipients,
num_namespaces,
))
})
.await
.map_err(|e| BlockError::VidDisperse(e.to_string()))??;

// Fan the shares out in the background, including the leader's own
// share (delivered via unicast loopback, which is how the leader
// obtains the share it votes with). A critical send failure is
// logged, not fatal.
spawn_blocking(move || {
if let Err(err) = fanout::fan_out::<T>(
per_bucket,
commitment,
param,
recipients,
num_namespaces,
view,
epoch,
network,
public_key,
private_key,
) {
error!(%view, %err, "vid share fanout failed");
}
});
Comment thread
mrain marked this conversation as resolved.
Outdated
let payload_commitment = VidCommitment::V2(commitment);

let builder_commitment = payload.builder_commitment(&metadata);
let (builder_key, builder_private_key) =
T::BuilderSignatureKey::generated_from_seed_indexed([0u8; 32], 0);
let block_size = payload_bytes.len() as u64;
let offered_fee = block_size;
let builder_fee = BuilderFee {
fee_amount: offered_fee,
fee_account: builder_key,
fee_signature: T::BuilderSignatureKey::sign_fee(
&builder_private_key,
offered_fee,
&payload.metadata,
&metadata,
)
.map_err(|_| BlockError::BuilderSignature)?,
};
Ok(BlockBuilderOutput {
view,
epoch,
payload,
payload: Arc::new(payload),
metadata,
parent_proposal: request.parent_proposal,
builder_commitment,
builder_fee,
Expand Down Expand Up @@ -373,7 +421,7 @@ impl<T: NodeType> From<&BlockBuilderOutput<T>> for HeaderRequest<T> {
parent_proposal: output.parent_proposal.clone(),
payload_commitment: output.payload_commitment,
builder_commitment: output.builder_commitment.clone(),
metadata: output.payload.metadata.clone(),
metadata: output.metadata.clone(),
builder_fee: output.builder_fee.clone(),
}
}
Expand All @@ -384,8 +432,7 @@ impl<T: NodeType> From<BlockBuilderOutput<T>> for ConsensusInput<T> {
ConsensusInput::BlockBuilt {
view: output.view,
epoch: output.epoch,
payload: output.payload.payload,
metadata: output.payload.metadata,
payload: output.payload,
payload_commitment: output.payload_commitment,
}
}
Expand Down
Loading
Loading