Skip to main content

hotshot_new_protocol/
block.rs

1use std::{
2    collections::{BTreeMap, HashMap, HashSet},
3    sync::Arc,
4    time::Duration,
5};
6
7use committable::{Commitment, Committable};
8use hotshot::traits::{BlockPayload, ValidatedState as _};
9use hotshot_types::{
10    consensus::PayloadWithMetadata,
11    data::{
12        EpochNumber, Leaf2, VidCommitment, ViewNumber, vid_commitment,
13        vid_disperse::vid_total_weight,
14    },
15    epoch_membership::EpochMembershipCoordinator,
16    message::UpgradeLock,
17    traits::{
18        EncodeBytes,
19        block_contents::{BuilderFee, Transaction},
20        node_implementation::NodeType,
21        signature_key::BuilderSignatureKey,
22    },
23    utils::BuilderCommitment,
24};
25use tokio::{
26    task::{AbortHandle, JoinSet},
27    time::sleep,
28};
29use tracing::{error, warn};
30
31use crate::{
32    consensus::ConsensusInput,
33    helpers::proposal_commitment,
34    message::{DedupManifest, Proposal, TransactionMessage},
35    state::HeaderRequest,
36};
37
38#[derive(Debug, thiserror::Error)]
39pub enum BlockError {
40    #[error("payload construction failed: {0}")]
41    PayloadConstruction(String),
42
43    #[error("stake table unavailable")]
44    StakeTableUnavailable,
45
46    #[error("builder signature failed")]
47    BuilderSignature,
48}
49
50#[derive(Clone, Eq, PartialEq, Debug)]
51pub struct BlockAndHeaderRequest<T: NodeType> {
52    pub view: ViewNumber,
53    pub epoch: EpochNumber,
54    pub parent_proposal: Proposal<T>,
55}
56
57pub struct BlockBuilderOutput<T: NodeType> {
58    pub view: ViewNumber,
59    pub epoch: EpochNumber,
60    pub payload: PayloadWithMetadata<T>,
61    pub parent_proposal: Proposal<T>,
62    pub builder_commitment: BuilderCommitment,
63    pub builder_fee: BuilderFee<T>,
64    pub payload_commitment: VidCommitment,
65    pub manifest: DedupManifest<T>,
66}
67
68pub struct BlockBuilderConfig {
69    pub max_retry_bytes: u64,
70    pub max_leader_bytes: u64,
71    pub ttl: u64,
72    pub dedup_window_size: u64,
73}
74
75impl Default for BlockBuilderConfig {
76    fn default() -> Self {
77        Self {
78            max_retry_bytes: 100 * 1024 * 1024,
79            max_leader_bytes: 2 * 1024 * 1024,
80            ttl: 50,
81            dedup_window_size: 10,
82        }
83    }
84}
85
86struct RetryEntry<T: NodeType> {
87    tx: T::Transaction,
88    valid_until: ViewNumber,
89    size: u64,
90}
91
92pub struct BlockBuilder<T: NodeType> {
93    instance: Arc<T::InstanceState>,
94    membership: EpochMembershipCoordinator<T>,
95    retry_pending: HashMap<Commitment<T::Transaction>, RetryEntry<T>>,
96    retry_total_bytes: u64,
97    leader_buffer: HashMap<Commitment<T::Transaction>, T::Transaction>,
98    leader_total_bytes: u64,
99    dedups: BTreeMap<ViewNumber, HashSet<Commitment<T::Transaction>>>,
100    config: BlockBuilderConfig,
101    upgrade_lock: UpgradeLock<T>,
102    current_view: ViewNumber,
103    // Keyed by (view, parent_proposal commitment) so that two requests for
104    // the same view but different parents (e.g. one from
105    // `handle_proposal_with_vid_share` and one from
106    // `handle_timeout_certificate`) don't dedup against each other.
107    calculations: BTreeMap<(ViewNumber, Commitment<Leaf2<T>>), AbortHandle>,
108    tasks: JoinSet<Result<BlockBuilderOutput<T>, BlockError>>,
109}
110
111impl<T: NodeType> BlockBuilder<T> {
112    pub fn new(
113        instance: Arc<T::InstanceState>,
114        membership: EpochMembershipCoordinator<T>,
115        config: BlockBuilderConfig,
116        upgrade_lock: UpgradeLock<T>,
117    ) -> Self {
118        Self {
119            instance,
120            membership,
121            config,
122            upgrade_lock,
123            retry_pending: HashMap::new(),
124            retry_total_bytes: 0,
125            leader_buffer: HashMap::new(),
126            leader_total_bytes: 0,
127            dedups: BTreeMap::new(),
128            current_view: ViewNumber::genesis(),
129            calculations: BTreeMap::new(),
130            tasks: JoinSet::new(),
131        }
132    }
133
134    pub fn request_block(&mut self, request: BlockAndHeaderRequest<T>) {
135        let view = request.view;
136        let parent_commitment = proposal_commitment(&request.parent_proposal);
137        if self.calculations.contains_key(&(view, parent_commitment)) {
138            return;
139        }
140        let Ok(version) = self.upgrade_lock.version(view) else {
141            warn!(%view, "unsupported version");
142            return;
143        };
144        let epoch = request.epoch;
145        let buffer = std::mem::take(&mut self.leader_buffer);
146        self.leader_total_bytes = 0;
147        let instance = self.instance.clone();
148        let membership = self.membership.clone();
149
150        let handle = self.tasks.spawn(async move {
151            // Throttle empty block production: when no transactions are pending,
152            // sleep so the coordinator's event queue doesnot overflow
153            // because if there are no transactions then the block production is way too fast
154            if buffer.is_empty() {
155                sleep(Duration::from_secs(1)).await;
156            }
157            let (hashes, txs): (Vec<_>, Vec<_>) = buffer.into_iter().unzip();
158            let manifest = DedupManifest {
159                view,
160                epoch,
161                hashes,
162            };
163
164            let validated_state =
165                T::ValidatedState::from_header(&request.parent_proposal.block_header);
166            let (payload, metadata) =
167                T::BlockPayload::from_transactions(txs, &validated_state, &instance)
168                    .await
169                    .map_err(|e| BlockError::PayloadConstruction(e.to_string()))?;
170            let payload: PayloadWithMetadata<T> = PayloadWithMetadata { payload, metadata };
171
172            let payload_bytes = payload.payload.encode();
173            let metadata_bytes = payload.metadata.encode();
174
175            let total_weight = {
176                let target_mem = membership
177                    .stake_table_for_epoch(Some(epoch))
178                    .map_err(|_| BlockError::StakeTableUnavailable)?;
179                vid_total_weight(target_mem.stake_table(), Some(epoch))
180            };
181            let payload_commitment = {
182                vid_commitment(
183                    payload_bytes.as_ref(),
184                    metadata_bytes.as_ref(),
185                    total_weight,
186                    version,
187                )
188            };
189
190            let builder_commitment = payload.payload.builder_commitment(&payload.metadata);
191            let (builder_key, builder_private_key) =
192                T::BuilderSignatureKey::generated_from_seed_indexed([0u8; 32], 0);
193            let block_size = payload_bytes.len() as u64;
194            let offered_fee = block_size;
195            let builder_fee = BuilderFee {
196                fee_amount: offered_fee,
197                fee_account: builder_key,
198                fee_signature: T::BuilderSignatureKey::sign_fee(
199                    &builder_private_key,
200                    offered_fee,
201                    &payload.metadata,
202                )
203                .map_err(|_| BlockError::BuilderSignature)?,
204            };
205            Ok(BlockBuilderOutput {
206                view,
207                epoch,
208                payload,
209                parent_proposal: request.parent_proposal,
210                builder_commitment,
211                builder_fee,
212                payload_commitment,
213                manifest,
214            })
215        });
216        self.calculations.insert((view, parent_commitment), handle);
217    }
218
219    pub async fn next(&mut self) -> Option<Result<BlockBuilderOutput<T>, BlockError>> {
220        loop {
221            match self.tasks.join_next().await {
222                Some(Ok(result)) => return Some(result),
223                Some(Err(err)) => {
224                    if err.is_panic() {
225                        error!(%err, "block builder task panicked");
226                    }
227                },
228                None => return None,
229            }
230        }
231    }
232
233    pub fn gc(&mut self, view_number: ViewNumber) {
234        self.calculations.retain(|(view, _), handle| {
235            if *view < view_number {
236                handle.abort();
237                false
238            } else {
239                true
240            }
241        });
242    }
243
244    pub fn outstanding_transactions(&self) -> (usize, usize) {
245        (self.retry_pending.len(), self.retry_total_bytes as usize)
246    }
247
248    pub fn on_submit_transaction(&mut self, tx: T::Transaction) {
249        let hash = tx.commit();
250
251        if self.retry_pending.contains_key(&hash) {
252            return;
253        }
254
255        let size = tx.minimum_block_size();
256        if self.retry_total_bytes + size > self.config.max_retry_bytes {
257            warn!("retry buffer full, rejecting {hash}");
258            return;
259        }
260
261        let valid_until = self.current_view + self.config.ttl;
262
263        self.retry_total_bytes += size;
264        self.retry_pending.insert(
265            hash,
266            RetryEntry {
267                tx,
268                valid_until,
269                size,
270            },
271        );
272    }
273
274    pub fn on_transactions(&mut self, msg: TransactionMessage<T>) {
275        for tx in msg.transactions {
276            let hash = tx.commit();
277
278            if self.dedups.values().any(|hs| hs.contains(&hash)) {
279                continue;
280            }
281
282            if self.leader_buffer.contains_key(&hash) {
283                continue;
284            }
285
286            let size = tx.minimum_block_size();
287            if self.leader_total_bytes + size > self.config.max_leader_bytes {
288                continue;
289            }
290
291            self.leader_total_bytes += size;
292            self.leader_buffer.insert(hash, tx);
293        }
294    }
295
296    pub fn on_dedup_manifest(&mut self, manifest: DedupManifest<T>) {
297        let DedupManifest { view, hashes, .. } = manifest;
298
299        for hash in &hashes {
300            if let Some(tx) = self.leader_buffer.remove(hash) {
301                self.leader_total_bytes -= tx.minimum_block_size();
302            }
303        }
304
305        let lower_bound: ViewNumber = self
306            .current_view
307            .saturating_sub(self.config.dedup_window_size)
308            .into();
309
310        if view >= lower_bound {
311            self.dedups.entry(view).or_default().extend(hashes);
312        }
313
314        self.dedups = self.dedups.split_off(&lower_bound);
315    }
316
317    pub fn on_view_changed(&mut self, view: ViewNumber) -> Vec<T::Transaction> {
318        self.current_view = view;
319
320        let mut expired_bytes = 0u64;
321        self.retry_pending.retain(|_, entry| {
322            if view > entry.valid_until {
323                expired_bytes += entry.size;
324                false
325            } else {
326                true
327            }
328        });
329        self.retry_total_bytes -= expired_bytes;
330
331        self.retry_pending
332            .values()
333            .map(|entry| entry.tx.clone())
334            .collect()
335    }
336
337    pub fn on_block_reconstructed(&mut self, tx_commitments: Vec<Commitment<T::Transaction>>) {
338        for hash in tx_commitments {
339            if let Some(entry) = self.retry_pending.remove(&hash) {
340                self.retry_total_bytes = self.retry_total_bytes.saturating_sub(entry.size);
341            }
342        }
343    }
344
345    #[cfg(test)]
346    pub(crate) fn drain(
347        &mut self,
348        view: ViewNumber,
349        epoch: EpochNumber,
350    ) -> (Vec<T::Transaction>, DedupManifest<T>) {
351        let (hashes, txs) = self.leader_buffer.drain().unzip();
352        self.leader_total_bytes = 0;
353
354        let manifest = DedupManifest {
355            view,
356            epoch,
357            hashes,
358        };
359
360        (txs, manifest)
361    }
362}
363
364impl<T: NodeType> From<&BlockBuilderOutput<T>> for HeaderRequest<T> {
365    fn from(output: &BlockBuilderOutput<T>) -> Self {
366        HeaderRequest {
367            view: output.view,
368            epoch: output.epoch,
369            parent_proposal: output.parent_proposal.clone(),
370            payload_commitment: output.payload_commitment,
371            builder_commitment: output.builder_commitment.clone(),
372            metadata: output.payload.metadata.clone(),
373            builder_fee: output.builder_fee.clone(),
374        }
375    }
376}
377
378impl<T: NodeType> From<BlockBuilderOutput<T>> for ConsensusInput<T> {
379    fn from(output: BlockBuilderOutput<T>) -> Self {
380        ConsensusInput::BlockBuilt {
381            view: output.view,
382            epoch: output.epoch,
383            payload: output.payload.payload,
384            metadata: output.payload.metadata,
385            payload_commitment: output.payload_commitment,
386        }
387    }
388}