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 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 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}