1#[cfg(feature = "docs")]
12pub mod documentation;
13
14use committable::Committable;
15use futures::future::{Either, select};
16use hotshot_types::{
17 drb::{DrbResult, INITIAL_DRB_RESULT, drb_difficulty_selector},
18 epoch_membership::EpochMembershipCoordinator,
19 message::UpgradeLock,
20 simple_certificate::{CertificatePair, LightClientStateUpdateCertificateV2},
21 traits::{
22 block_contents::BlockHeader, election::Membership, network::BroadcastDelay,
23 signature_key::StateSignatureKey, storage::Storage,
24 },
25 utils::{epoch_from_block_number, is_ge_epoch_root},
26};
27use rand::Rng;
28
29pub mod traits;
31pub mod types;
33
34pub mod tasks;
35use hotshot_types::data::QuorumProposalWrapper;
36use versions::{EPOCH_VERSION, Upgrade};
37
38pub mod helpers;
40
41use std::{
42 collections::{BTreeMap, HashMap},
43 num::NonZeroUsize,
44 sync::Arc,
45 time::Duration,
46};
47
48use alloy::primitives::U256;
49use async_broadcast::{InactiveReceiver, Receiver, Sender, broadcast};
50use async_lock::RwLock;
51use async_trait::async_trait;
52use futures::join;
53use hotshot_task::task::{ConsensusTaskRegistry, NetworkTaskRegistry};
54use hotshot_task_impls::{events::HotShotEvent, helpers::broadcast_event};
55pub use hotshot_types::error::HotShotError;
58use hotshot_types::{
59 HotShotConfig,
60 consensus::{
61 Consensus, ConsensusMetricsValue, OuterConsensus, PayloadWithMetadata, VidShares, View,
62 ViewInner,
63 },
64 constants::{EVENT_CHANNEL_SIZE, EXTERNAL_EVENT_CHANNEL_SIZE},
65 data::{EpochNumber, Leaf2, ViewNumber},
66 event::{EventType, LeafInfo},
67 message::{DataMessage, Message, MessageKind, Proposal},
68 simple_certificate::{NextEpochQuorumCertificate2, QuorumCertificate2, UpgradeCertificate},
69 stake_table::HSStakeTable,
70 storage_metrics::StorageMetricsValue,
71 traits::{
72 consensus_api::ConsensusApi, network::ConnectedNetwork, node_implementation::NodeType,
73 signature_key::SignatureKey, states::ValidatedState,
74 },
75 utils::{genesis_epoch_from_version, option_epoch_from_block_number},
76};
77use hotshot_utils::warn;
78pub use rand;
80use tokio::{spawn, time::sleep};
81use tracing::{Instrument, debug, error_span, info, instrument, trace};
82
83use crate::{
86 tasks::{add_consensus_tasks, add_network_tasks},
87 traits::NodeImplementation,
88 types::{Event, SystemContextHandle},
89};
90
91pub const H_512: usize = 64;
93pub const H_256: usize = 32;
95
96pub struct SystemContext<TYPES: NodeType, I: NodeImplementation<TYPES>> {
98 public_key: TYPES::SignatureKey,
100
101 private_key: <TYPES::SignatureKey as SignatureKey>::PrivateKey,
103
104 state_private_key: <TYPES::StateSignatureKey as StateSignatureKey>::StatePrivateKey,
106
107 pub config: HotShotConfig<TYPES>,
109
110 pub network: Arc<I::Network>,
112
113 pub membership_coordinator: EpochMembershipCoordinator<TYPES>,
115
116 metrics: Arc<ConsensusMetricsValue>,
118
119 consensus: OuterConsensus<TYPES>,
121
122 instance_state: Arc<TYPES::InstanceState>,
124
125 start_view: ViewNumber,
127
128 start_epoch: Option<EpochNumber>,
130
131 output_event_stream: (Sender<Event<TYPES>>, InactiveReceiver<Event<TYPES>>),
133
134 pub(crate) external_event_stream: (Sender<Event<TYPES>>, InactiveReceiver<Event<TYPES>>),
136
137 anchored_leaf: Leaf2<TYPES>,
139
140 #[allow(clippy::type_complexity)]
142 internal_event_stream: (
143 Sender<Arc<HotShotEvent<TYPES>>>,
144 InactiveReceiver<Arc<HotShotEvent<TYPES>>>,
145 ),
146
147 pub id: u64,
149
150 pub storage: I::Storage,
152
153 pub storage_metrics: Arc<StorageMetricsValue>,
155
156 pub upgrade_lock: UpgradeLock<TYPES>,
158}
159impl<TYPES: NodeType, I: NodeImplementation<TYPES>> Clone for SystemContext<TYPES, I> {
160 #![allow(deprecated)]
161 fn clone(&self) -> Self {
162 Self {
163 public_key: self.public_key.clone(),
164 private_key: self.private_key.clone(),
165 state_private_key: self.state_private_key.clone(),
166 config: self.config.clone(),
167 network: Arc::clone(&self.network),
168 membership_coordinator: self.membership_coordinator.clone(),
169 metrics: Arc::clone(&self.metrics),
170 consensus: self.consensus.clone(),
171 instance_state: Arc::clone(&self.instance_state),
172 start_view: self.start_view,
173 start_epoch: self.start_epoch,
174 output_event_stream: self.output_event_stream.clone(),
175 external_event_stream: self.external_event_stream.clone(),
176 anchored_leaf: self.anchored_leaf.clone(),
177 internal_event_stream: self.internal_event_stream.clone(),
178 id: self.id,
179 storage: self.storage.clone(),
180 storage_metrics: Arc::clone(&self.storage_metrics),
181 upgrade_lock: self.upgrade_lock.clone(),
182 }
183 }
184}
185
186impl<TYPES: NodeType, I: NodeImplementation<TYPES>> SystemContext<TYPES, I> {
187 #![allow(deprecated)]
188 #[allow(clippy::too_many_arguments)]
195 pub async fn new(
196 public_key: TYPES::SignatureKey,
197 private_key: <TYPES::SignatureKey as SignatureKey>::PrivateKey,
198 state_private_key: <TYPES::StateSignatureKey as StateSignatureKey>::StatePrivateKey,
199 nonce: u64,
200 config: HotShotConfig<TYPES>,
201 upgrade: versions::Upgrade,
202 memberships: EpochMembershipCoordinator<TYPES>,
203 network: Arc<I::Network>,
204 initializer: HotShotInitializer<TYPES>,
205 consensus_metrics: ConsensusMetricsValue,
206 storage: I::Storage,
207 storage_metrics: StorageMetricsValue,
208 ) -> Arc<Self> {
209 let internal_chan = broadcast(EVENT_CHANNEL_SIZE);
210 let external_chan = broadcast(EXTERNAL_EVENT_CHANNEL_SIZE);
211
212 Self::new_from_channels(
213 public_key,
214 private_key,
215 state_private_key,
216 nonce,
217 config,
218 upgrade,
219 memberships,
220 network,
221 initializer,
222 consensus_metrics,
223 storage,
224 storage_metrics,
225 internal_chan,
226 external_chan,
227 )
228 .await
229 }
230
231 #[allow(clippy::too_many_arguments, clippy::type_complexity)]
239 pub async fn new_from_channels(
240 public_key: TYPES::SignatureKey,
241 private_key: <TYPES::SignatureKey as SignatureKey>::PrivateKey,
242 state_private_key: <TYPES::StateSignatureKey as StateSignatureKey>::StatePrivateKey,
243 nonce: u64,
244 config: HotShotConfig<TYPES>,
245 upgrade: versions::Upgrade,
246 membership_coordinator: EpochMembershipCoordinator<TYPES>,
247 network: Arc<I::Network>,
248 initializer: HotShotInitializer<TYPES>,
249 consensus_metrics: ConsensusMetricsValue,
250 storage: I::Storage,
251 storage_metrics: StorageMetricsValue,
252 internal_channel: (
253 Sender<Arc<HotShotEvent<TYPES>>>,
254 Receiver<Arc<HotShotEvent<TYPES>>>,
255 ),
256 external_channel: (Sender<Event<TYPES>>, Receiver<Event<TYPES>>),
257 ) -> Arc<Self> {
258 debug!("Creating a new hotshot");
259
260 tracing::warn!("Starting consensus with HotShotConfig:\n\n {config:?}");
261
262 let consensus_metrics = Arc::new(consensus_metrics);
263 let storage_metrics = Arc::new(storage_metrics);
264 let anchored_leaf = initializer.anchor_leaf;
265 let instance_state = initializer.instance_state;
266
267 let (internal_tx, internal_rx) = internal_channel;
268 let (mut external_tx, external_rx) = external_channel;
269
270 let mut internal_rx = internal_rx.new_receiver();
271
272 let mut external_rx = external_rx.new_receiver();
273
274 internal_rx.set_overflow(true);
277 external_rx.set_overflow(true);
279
280 tracing::warn!(
281 "Starting consensus with versions:\n\n Base: {:?}\nUpgrade: {:?}.",
282 upgrade.base,
283 upgrade.target
284 );
285 tracing::warn!(
286 "Loading previously decided upgrade certificate from storage: {:?}",
287 initializer.decided_upgrade_certificate
288 );
289
290 let upgrade_lock = UpgradeLock::<TYPES>::from_certificate(
291 upgrade,
292 &initializer.decided_upgrade_certificate,
293 );
294
295 let current_version = if let Some(cert) = initializer.decided_upgrade_certificate {
296 cert.data.new_version
297 } else {
298 upgrade.base
299 };
300
301 debug!("Setting DRB difficulty selector in membership");
302 let drb_difficulty_selector = drb_difficulty_selector(&config);
303
304 membership_coordinator.set_drb_difficulty_selector(drb_difficulty_selector);
305
306 for da_committee in &config.da_committees {
307 if current_version >= da_committee.start_version {
308 membership_coordinator.membership().add_da_committee(
309 da_committee.start_epoch.into(),
310 da_committee.committee.clone(),
311 );
312 }
313 }
314
315 let validated_state = initializer.anchor_state;
318
319 load_start_epoch_info(
320 &membership_coordinator,
321 &initializer.start_epoch_info,
322 config.epoch_height,
323 config.epoch_start_block,
324 )
325 .await;
326
327 let epoch = initializer.high_qc.data.block_number.map(|block_number| {
329 EpochNumber::new(epoch_from_block_number(
330 block_number + 1,
331 config.epoch_height,
332 ))
333 });
334
335 let mut validated_state_map = BTreeMap::default();
337 validated_state_map.insert(
338 anchored_leaf.view_number(),
339 View {
340 view_inner: ViewInner::Leaf {
341 leaf: anchored_leaf.commit(),
342 state: Arc::clone(&validated_state),
343 delta: initializer.anchor_state_delta,
344 epoch,
345 },
346 },
347 );
348 for (view_num, inner) in initializer.undecided_state {
349 validated_state_map.insert(view_num, inner);
350 }
351
352 let mut saved_leaves = HashMap::new();
353 let mut saved_payloads = BTreeMap::new();
354 saved_leaves.insert(anchored_leaf.commit(), anchored_leaf.clone());
355
356 for (_, leaf) in initializer.undecided_leaves {
357 saved_leaves.insert(leaf.commit(), leaf.clone());
358 }
359 if let Some(payload) = anchored_leaf.block_payload() {
360 let metadata = anchored_leaf.block_header().metadata().clone();
361 saved_payloads.insert(
362 anchored_leaf.view_number(),
363 Arc::new(PayloadWithMetadata { payload, metadata }),
364 );
365 }
366 let high_qc_block_number = initializer.high_qc.data.block_number;
367 let (stake_table, success_threshold) =
368 if let Ok(epoch_membership) = membership_coordinator.stake_table_for_epoch(epoch) {
369 (
370 HSStakeTable::from_iter(epoch_membership.stake_table()),
371 epoch_membership.success_threshold(),
372 )
373 } else {
374 tracing::warn!(
375 "Failed to get stake table for epoch {:?} while creating vote participation",
376 epoch
377 );
378 (HSStakeTable::default(), U256::MAX)
379 };
380
381 let recovered_locked_view = initializer
382 .high_qc
383 .view_number
384 .max(anchored_leaf.view_number());
385 let consensus = Consensus::new(
386 validated_state_map,
387 Some(initializer.saved_vid_shares),
388 anchored_leaf.view_number(),
389 epoch,
390 recovered_locked_view,
391 anchored_leaf.view_number(),
392 initializer.last_actioned_view,
393 initializer.saved_proposals,
394 saved_leaves,
395 saved_payloads,
396 initializer.high_qc,
397 initializer.next_epoch_high_qc,
398 Arc::clone(&consensus_metrics),
399 config.epoch_height,
400 initializer.state_cert,
401 config.drb_difficulty,
402 config.drb_upgrade_difficulty,
403 stake_table,
404 success_threshold,
405 );
406
407 let consensus = Arc::new(RwLock::new(consensus));
408
409 if let Some(epoch) = epoch {
410 tracing::info!(
411 "Triggering catchup for epoch {} and next epoch {}",
412 epoch,
413 epoch + 1
414 );
415 let _ = membership_coordinator.membership_for_epoch(Some(epoch));
417 let _ = membership_coordinator.membership_for_epoch(Some(epoch + 1));
418 if let Some(high_qc_block_number) = high_qc_block_number
421 && is_ge_epoch_root(high_qc_block_number, config.epoch_height)
422 {
423 let _ = membership_coordinator.stake_table_for_epoch(Some(epoch + 2));
424 }
425
426 if let Ok(drb_result) = storage.load_drb_result(epoch + 1).await {
427 info!(target: "announce::drb", epoch = %(epoch + 1), "writing drb result for epoch");
428 if let Ok(mem) = membership_coordinator.stake_table_for_epoch(Some(epoch + 1)) {
429 mem.add_drb_result(drb_result);
430 }
431 }
432 }
433
434 external_tx.set_await_active(false);
437
438 let inner: Arc<SystemContext<TYPES, I>> = Arc::new(SystemContext {
439 id: nonce,
440 consensus: OuterConsensus::new(consensus),
441 instance_state: Arc::new(instance_state),
442 public_key,
443 private_key,
444 state_private_key,
445 config,
446 start_view: initializer.start_view,
447 start_epoch: initializer.start_epoch,
448 network,
449 membership_coordinator,
450 metrics: Arc::clone(&consensus_metrics),
451 internal_event_stream: (internal_tx, internal_rx.deactivate()),
452 output_event_stream: (external_tx.clone(), external_rx.clone().deactivate()),
453 external_event_stream: (external_tx, external_rx.deactivate()),
454 anchored_leaf: anchored_leaf.clone(),
455 storage,
456 storage_metrics,
457 upgrade_lock,
458 });
459
460 inner
461 }
462
463 #[instrument(skip_all, target = "SystemContext", fields(id = self.id))]
468 pub async fn start_consensus(&self) {
469 #[cfg(all(feature = "rewind", not(debug_assertions)))]
470 compile_error!("Cannot run rewind in production builds!");
471
472 debug!("Starting Consensus");
473 let consensus = self.consensus.read().await;
474
475 let first_epoch = option_epoch_from_block_number(
476 self.upgrade_lock.upgrade().base >= EPOCH_VERSION,
477 self.config.epoch_start_block,
478 self.config.epoch_height,
479 );
480 let initial_view_change_epoch = self.start_epoch.max(first_epoch);
484 #[allow(clippy::panic)]
485 self.internal_event_stream
486 .0
487 .broadcast_direct(Arc::new(HotShotEvent::ViewChange(
488 self.start_view,
489 initial_view_change_epoch,
490 )))
491 .await
492 .unwrap_or_else(|_| {
493 panic!(
494 "Genesis Broadcast failed; event = ViewChange({:?}, {:?})",
495 self.start_view, initial_view_change_epoch,
496 )
497 });
498
499 let event_stream = self.internal_event_stream.0.clone();
501 let next_view_timeout = self.config.next_view_timeout;
502 let start_view = self.start_view;
503 let start_epoch = self.start_epoch;
504
505 spawn(
508 async move {
509 sleep(Duration::from_millis(next_view_timeout)).await;
510 broadcast_event(
511 Arc::new(HotShotEvent::Timeout(start_view, start_epoch)),
512 &event_stream,
513 )
514 .await;
515 }
516 .instrument(error_span!("initial timeout", view = *start_view)),
517 );
518 #[allow(clippy::panic)]
519 self.internal_event_stream
520 .0
521 .broadcast_direct(Arc::new(HotShotEvent::Qc2Formed(either::Left(
522 consensus.high_qc().clone(),
523 ))))
524 .await
525 .unwrap_or_else(|_| {
526 panic!(
527 "Genesis Broadcast failed; event = Qc2Formed(either::Left({:?}))",
528 consensus.high_qc()
529 )
530 });
531
532 {
533 if self.anchored_leaf.view_number() == ViewNumber::genesis() {
536 let (validated_state, state_delta) =
537 TYPES::ValidatedState::genesis(&self.instance_state);
538
539 let qc = QuorumCertificate2::genesis(
540 &validated_state,
541 self.instance_state.as_ref(),
542 self.upgrade_lock.upgrade(),
543 )
544 .await;
545
546 broadcast_event(
547 Event {
548 view_number: self.anchored_leaf.view_number(),
549 event: EventType::Decide {
550 leaf_chain: Arc::new(vec![LeafInfo::new(
551 self.anchored_leaf.clone(),
552 Arc::new(validated_state),
553 Some(Arc::new(state_delta)),
554 None,
555 None,
556 )]),
557 committing_qc: Arc::new(CertificatePair::non_epoch_change(qc)),
558 deciding_qc: None,
559 block_size: None,
560 },
561 },
562 &self.external_event_stream.0,
563 )
564 .await;
565 }
566 }
567 }
568
569 async fn send_external_event(&self, event: Event<TYPES>) {
571 debug!(?event, "send_external_event");
572 broadcast_event(event, &self.external_event_stream.0).await;
573 }
574
575 #[instrument(skip(self), err, target = "SystemContext", fields(id = self.id))]
581 pub async fn publish_transaction_async(
582 &self,
583 transaction: TYPES::Transaction,
584 ) -> Result<(), HotShotError<TYPES>> {
585 trace!("Adding transaction to our own queue");
586
587 let api = self.clone();
588
589 let consensus_reader = api.consensus.read().await;
590 let view_number = consensus_reader.cur_view();
591 let epoch = consensus_reader.cur_epoch();
592 drop(consensus_reader);
593
594 let message_kind: DataMessage<TYPES> =
596 DataMessage::SubmitTransaction(transaction.clone(), view_number);
597 let message = Message {
598 sender: api.public_key.clone(),
599 kind: MessageKind::from(message_kind),
600 };
601
602 let serialized_message = self.upgrade_lock.serialize(&message).map_err(|err| {
603 HotShotError::FailedToSerialize(format!("failed to serialize transaction: {err}"))
604 })?;
605
606 let membership = match api.membership_coordinator.membership_for_epoch(epoch) {
607 Ok(m) => m,
608 Err(e) => return Err(HotShotError::InvalidState(e.message)),
609 };
610
611 spawn(async move {
612 let memberships_da_committee_members = membership
613 .da_committee_members(view_number)
614 .cloned()
615 .collect();
616
617 join! {
618 api
626 .network.da_broadcast_message(
627 view_number.u64().into(),
628 serialized_message,
629 memberships_da_committee_members,
630 BroadcastDelay::None,
631 ),
632 api
633 .send_external_event(Event {
634 view_number,
635 event: EventType::Transactions {
636 transactions: vec![transaction],
637 },
638 }),
639 }
640 });
641 Ok(())
642 }
643
644 #[must_use]
646 pub fn consensus(&self) -> Arc<RwLock<Consensus<TYPES>>> {
647 Arc::clone(&self.consensus.inner_consensus)
648 }
649
650 pub fn instance_state(&self) -> Arc<TYPES::InstanceState> {
652 Arc::clone(&self.instance_state)
653 }
654
655 #[instrument(skip_all, target = "SystemContext", fields(id = self.id))]
659 pub async fn decided_leaf(&self) -> Leaf2<TYPES> {
660 self.consensus.read().await.decided_leaf()
661 }
662
663 #[must_use]
669 #[instrument(skip_all, target = "SystemContext", fields(id = self.id))]
670 pub fn try_decided_leaf(&self) -> Option<Leaf2<TYPES>> {
671 self.consensus.try_read().map(|guard| guard.decided_leaf())
672 }
673
674 #[instrument(skip_all, target = "SystemContext", fields(id = self.id))]
679 pub async fn decided_state(&self) -> Arc<TYPES::ValidatedState> {
680 Arc::clone(&self.consensus.read().await.decided_state())
681 }
682
683 #[instrument(skip_all, target = "SystemContext", fields(id = self.id))]
691 pub async fn state(&self, view: ViewNumber) -> Option<Arc<TYPES::ValidatedState>> {
692 self.consensus.read().await.state(view).cloned()
693 }
694
695 #[allow(clippy::too_many_arguments)]
709 pub async fn init(
710 public_key: TYPES::SignatureKey,
711 private_key: <TYPES::SignatureKey as SignatureKey>::PrivateKey,
712 state_private_key: <TYPES::StateSignatureKey as StateSignatureKey>::StatePrivateKey,
713 node_id: u64,
714 config: HotShotConfig<TYPES>,
715 upgrade: versions::Upgrade,
716 memberships: EpochMembershipCoordinator<TYPES>,
717 network: Arc<I::Network>,
718 initializer: HotShotInitializer<TYPES>,
719 consensus_metrics: ConsensusMetricsValue,
720 storage: I::Storage,
721 storage_metrics: StorageMetricsValue,
722 ) -> Result<
723 (
724 SystemContextHandle<TYPES, I>,
725 Sender<Arc<HotShotEvent<TYPES>>>,
726 Receiver<Arc<HotShotEvent<TYPES>>>,
727 ),
728 HotShotError<TYPES>,
729 > {
730 let hotshot = Self::new(
731 public_key,
732 private_key,
733 state_private_key,
734 node_id,
735 config,
736 upgrade,
737 memberships,
738 network,
739 initializer,
740 consensus_metrics,
741 storage,
742 storage_metrics,
743 )
744 .await;
745 let handle = Arc::clone(&hotshot).run_tasks().await;
746 let (tx, rx) = hotshot.internal_event_stream.clone();
747
748 Ok((handle, tx, rx.activate()))
749 }
750 #[must_use]
752 pub fn next_view_timeout(&self) -> u64 {
753 self.config.next_view_timeout
754 }
755}
756
757impl<TYPES: NodeType, I: NodeImplementation<TYPES>> SystemContext<TYPES, I> {
758 pub async fn run_tasks(&self) -> SystemContextHandle<TYPES, I> {
762 let consensus_registry = ConsensusTaskRegistry::new();
763 let network_registry = NetworkTaskRegistry::new();
764
765 let output_event_stream = self.external_event_stream.clone();
766 let internal_event_stream = self.internal_event_stream.clone();
767
768 let mut handle = SystemContextHandle {
769 consensus_registry,
770 network_registry,
771 output_event_stream: output_event_stream.clone(),
772 internal_event_stream: internal_event_stream.clone(),
773 hotshot: self.clone().into(),
774 storage: self.storage.clone(),
775 network: Arc::clone(&self.network),
776 membership_coordinator: self.membership_coordinator.clone(),
777 epoch_height: self.config.epoch_height,
778 };
779
780 add_network_tasks::<TYPES, I>(&mut handle).await;
781 add_consensus_tasks::<TYPES, I>(&mut handle).await;
782
783 handle
784 }
785}
786
787type Channel<S> = (Sender<Arc<S>>, Receiver<Arc<S>>);
789
790#[async_trait]
792pub trait TwinsHandlerState<TYPES, I>
793where
794 Self: std::fmt::Debug + Send + Sync,
795 TYPES: NodeType,
796 I: NodeImplementation<TYPES>,
797{
798 async fn send_handler(
800 &mut self,
801 event: &HotShotEvent<TYPES>,
802 ) -> Vec<Either<HotShotEvent<TYPES>, HotShotEvent<TYPES>>>;
803
804 async fn recv_handler(
806 &mut self,
807 event: &Either<HotShotEvent<TYPES>, HotShotEvent<TYPES>>,
808 ) -> Vec<HotShotEvent<TYPES>>;
809
810 fn fuse_channels(
814 &'static mut self,
815 left: Channel<HotShotEvent<TYPES>>,
816 right: Channel<HotShotEvent<TYPES>>,
817 ) -> Channel<HotShotEvent<TYPES>> {
818 let send_state = Arc::new(RwLock::new(self));
819 let recv_state = Arc::clone(&send_state);
820
821 let (left_sender, mut left_receiver) = (left.0, left.1);
822 let (right_sender, mut right_receiver) = (right.0, right.1);
823
824 let (sender_to_network, network_task_receiver) = broadcast(EVENT_CHANNEL_SIZE);
826 let (network_task_sender, mut receiver_from_network): Channel<HotShotEvent<TYPES>> =
828 broadcast(EVENT_CHANNEL_SIZE);
829
830 let _recv_loop_handle = spawn(async move {
831 loop {
832 let msg = match select(left_receiver.recv(), right_receiver.recv()).await {
833 Either::Left(msg) => Either::Left(msg.0.unwrap().as_ref().clone()),
834 Either::Right(msg) => Either::Right(msg.0.unwrap().as_ref().clone()),
835 };
836
837 let mut state = recv_state.write().await;
838 let mut result = state.recv_handler(&msg).await;
839
840 while let Some(event) = result.pop() {
841 let _ = sender_to_network.broadcast(event.into()).await;
842 }
843 }
844 });
845
846 let _send_loop_handle = spawn(async move {
847 loop {
848 if let Ok(msg) = receiver_from_network.recv().await {
849 let mut state = send_state.write().await;
850
851 let mut result = state.send_handler(&msg).await;
852
853 while let Some(event) = result.pop() {
854 match event {
855 Either::Left(msg) => {
856 let _ = left_sender.broadcast(msg.into()).await;
857 },
858 Either::Right(msg) => {
859 let _ = right_sender.broadcast(msg.into()).await;
860 },
861 }
862 }
863 }
864 }
865 });
866
867 (network_task_sender, network_task_receiver)
868 }
869
870 #[allow(clippy::too_many_arguments)]
871 async fn spawn_twin_handles(
875 &'static mut self,
876 public_key: TYPES::SignatureKey,
877 private_key: <TYPES::SignatureKey as SignatureKey>::PrivateKey,
878 state_private_key: <TYPES::StateSignatureKey as StateSignatureKey>::StatePrivateKey,
879 nonce: u64,
880 config: HotShotConfig<TYPES>,
881 upgrade: versions::Upgrade,
882 memberships: EpochMembershipCoordinator<TYPES>,
883 network: Arc<I::Network>,
884 initializer: HotShotInitializer<TYPES>,
885 consensus_metrics: ConsensusMetricsValue,
886 storage: I::Storage,
887 storage_metrics: StorageMetricsValue,
888 ) -> (SystemContextHandle<TYPES, I>, SystemContextHandle<TYPES, I>) {
889 let epoch_height = config.epoch_height;
890 let left_system_context = SystemContext::new(
891 public_key.clone(),
892 private_key.clone(),
893 state_private_key.clone(),
894 nonce,
895 config.clone(),
896 upgrade,
897 memberships.clone(),
898 Arc::clone(&network),
899 initializer.clone(),
900 consensus_metrics.clone(),
901 storage.clone(),
902 storage_metrics.clone(),
903 )
904 .await;
905 let right_system_context = SystemContext::new(
906 public_key,
907 private_key,
908 state_private_key,
909 nonce,
910 config,
911 upgrade,
912 memberships,
913 network,
914 initializer,
915 consensus_metrics,
916 storage,
917 storage_metrics,
918 )
919 .await;
920
921 let left_consensus_registry = ConsensusTaskRegistry::new();
923 let left_network_registry = NetworkTaskRegistry::new();
924
925 let right_consensus_registry = ConsensusTaskRegistry::new();
926 let right_network_registry = NetworkTaskRegistry::new();
927
928 let (left_external_sender, left_external_receiver) = broadcast(EXTERNAL_EVENT_CHANNEL_SIZE);
930 let left_external_event_stream =
931 (left_external_sender, left_external_receiver.deactivate());
932
933 let (right_external_sender, right_external_receiver) =
934 broadcast(EXTERNAL_EVENT_CHANNEL_SIZE);
935 let right_external_event_stream =
936 (right_external_sender, right_external_receiver.deactivate());
937
938 let (left_internal_sender, left_internal_receiver) = broadcast(EVENT_CHANNEL_SIZE);
940 let left_internal_event_stream = (
941 left_internal_sender.clone(),
942 left_internal_receiver.clone().deactivate(),
943 );
944
945 let (right_internal_sender, right_internal_receiver) = broadcast(EVENT_CHANNEL_SIZE);
946 let right_internal_event_stream = (
947 right_internal_sender.clone(),
948 right_internal_receiver.clone().deactivate(),
949 );
950
951 let mut left_handle = SystemContextHandle::<_, I> {
953 consensus_registry: left_consensus_registry,
954 network_registry: left_network_registry,
955 output_event_stream: left_external_event_stream.clone(),
956 internal_event_stream: left_internal_event_stream.clone(),
957 hotshot: Arc::clone(&left_system_context),
958 storage: left_system_context.storage.clone(),
959 network: Arc::clone(&left_system_context.network),
960 membership_coordinator: left_system_context.membership_coordinator.clone(),
961 epoch_height,
962 };
963
964 let mut right_handle = SystemContextHandle::<_, I> {
965 consensus_registry: right_consensus_registry,
966 network_registry: right_network_registry,
967 output_event_stream: right_external_event_stream.clone(),
968 internal_event_stream: right_internal_event_stream.clone(),
969 hotshot: Arc::clone(&right_system_context),
970 storage: right_system_context.storage.clone(),
971 network: Arc::clone(&right_system_context.network),
972 membership_coordinator: right_system_context.membership_coordinator.clone(),
973 epoch_height,
974 };
975
976 add_consensus_tasks::<TYPES, I>(&mut left_handle).await;
978 add_consensus_tasks::<TYPES, I>(&mut right_handle).await;
979
980 let fused_internal_event_stream = self.fuse_channels(
982 (left_internal_sender, left_internal_receiver),
983 (right_internal_sender, right_internal_receiver),
984 );
985
986 left_handle.internal_event_stream = (
988 fused_internal_event_stream.0,
989 fused_internal_event_stream.1.deactivate(),
990 );
991
992 add_network_tasks::<TYPES, I>(&mut left_handle).await;
994
995 left_handle.internal_event_stream = left_internal_event_stream.clone();
997
998 (left_handle, right_handle)
999 }
1000}
1001
1002#[derive(Debug)]
1003pub struct RandomTwinsHandler;
1006
1007#[async_trait]
1008impl<TYPES: NodeType, I: NodeImplementation<TYPES>> TwinsHandlerState<TYPES, I>
1009 for RandomTwinsHandler
1010{
1011 async fn send_handler(
1012 &mut self,
1013 event: &HotShotEvent<TYPES>,
1014 ) -> Vec<Either<HotShotEvent<TYPES>, HotShotEvent<TYPES>>> {
1015 let random: bool = rand::thread_rng().r#gen();
1016
1017 #[allow(clippy::match_bool)]
1018 match random {
1019 true => vec![Either::Left(event.clone())],
1020 false => vec![Either::Right(event.clone())],
1021 }
1022 }
1023
1024 async fn recv_handler(
1025 &mut self,
1026 event: &Either<HotShotEvent<TYPES>, HotShotEvent<TYPES>>,
1027 ) -> Vec<HotShotEvent<TYPES>> {
1028 match event {
1029 Either::Left(msg) | Either::Right(msg) => vec![msg.clone()],
1030 }
1031 }
1032}
1033
1034#[derive(Debug)]
1037pub struct DoubleTwinsHandler;
1038
1039#[async_trait]
1040impl<TYPES: NodeType, I: NodeImplementation<TYPES>> TwinsHandlerState<TYPES, I>
1041 for DoubleTwinsHandler
1042{
1043 async fn send_handler(
1044 &mut self,
1045 event: &HotShotEvent<TYPES>,
1046 ) -> Vec<Either<HotShotEvent<TYPES>, HotShotEvent<TYPES>>> {
1047 vec![Either::Left(event.clone()), Either::Right(event.clone())]
1048 }
1049
1050 async fn recv_handler(
1051 &mut self,
1052 event: &Either<HotShotEvent<TYPES>, HotShotEvent<TYPES>>,
1053 ) -> Vec<HotShotEvent<TYPES>> {
1054 match event {
1055 Either::Left(msg) | Either::Right(msg) => vec![msg.clone()],
1056 }
1057 }
1058}
1059
1060#[async_trait]
1061impl<TYPES: NodeType, I: NodeImplementation<TYPES>> ConsensusApi<TYPES, I>
1062 for SystemContextHandle<TYPES, I>
1063{
1064 fn total_nodes(&self) -> NonZeroUsize {
1065 self.hotshot.config.num_nodes_with_stake
1066 }
1067
1068 fn builder_timeout(&self) -> Duration {
1069 self.hotshot.config.builder_timeout
1070 }
1071
1072 async fn send_event(&self, event: Event<TYPES>) {
1073 debug!(?event, "send_event");
1074 broadcast_event(event, &self.hotshot.external_event_stream.0).await;
1075 }
1076
1077 fn public_key(&self) -> &TYPES::SignatureKey {
1078 &self.hotshot.public_key
1079 }
1080
1081 fn private_key(&self) -> &<TYPES::SignatureKey as SignatureKey>::PrivateKey {
1082 &self.hotshot.private_key
1083 }
1084
1085 fn state_private_key(
1086 &self,
1087 ) -> &<TYPES::StateSignatureKey as StateSignatureKey>::StatePrivateKey {
1088 &self.hotshot.state_private_key
1089 }
1090}
1091
1092#[derive(Clone, Debug, PartialEq)]
1093pub struct InitializerEpochInfo<TYPES: NodeType> {
1094 pub epoch: EpochNumber,
1095 pub drb_result: DrbResult,
1096 pub block_header: Option<TYPES::BlockHeader>,
1098}
1099
1100pub type InitializerAnchor<TYPES> = (
1104 Leaf2<TYPES>,
1105 Option<Arc<<TYPES as NodeType>::ValidatedState>>,
1106 Option<Arc<<<TYPES as NodeType>::ValidatedState as ValidatedState<TYPES>>::Delta>>,
1107);
1108
1109#[derive(Clone, Debug)]
1111pub struct HotShotInitializer<TYPES: NodeType> {
1112 instance_state: TYPES::InstanceState,
1114
1115 epoch_height: u64,
1117
1118 epoch_start_block: u64,
1120
1121 anchor_leaf: Leaf2<TYPES>,
1123
1124 anchor_state: Arc<TYPES::ValidatedState>,
1126
1127 anchor_state_delta: Option<Arc<<TYPES::ValidatedState as ValidatedState<TYPES>>::Delta>>,
1129
1130 start_view: ViewNumber,
1132
1133 last_actioned_view: ViewNumber,
1136
1137 start_epoch: Option<EpochNumber>,
1139
1140 high_qc: QuorumCertificate2<TYPES>,
1144
1145 next_epoch_high_qc: Option<NextEpochQuorumCertificate2<TYPES>>,
1147
1148 saved_proposals: BTreeMap<ViewNumber, Proposal<TYPES, QuorumProposalWrapper<TYPES>>>,
1150
1151 decided_upgrade_certificate: Option<UpgradeCertificate<TYPES>>,
1153
1154 undecided_leaves: BTreeMap<ViewNumber, Leaf2<TYPES>>,
1157
1158 undecided_state: BTreeMap<ViewNumber, View<TYPES>>,
1160
1161 saved_vid_shares: VidShares<TYPES>,
1163
1164 state_cert: Option<LightClientStateUpdateCertificateV2<TYPES>>,
1166
1167 start_epoch_info: Vec<InitializerEpochInfo<TYPES>>,
1169}
1170
1171impl<TYPES: NodeType> HotShotInitializer<TYPES> {
1172 pub async fn from_genesis(
1176 instance_state: TYPES::InstanceState,
1177 epoch_height: u64,
1178 epoch_start_block: u64,
1179 start_epoch_info: Vec<InitializerEpochInfo<TYPES>>,
1180 upgrade: Upgrade,
1181 ) -> Result<Self, HotShotError<TYPES>> {
1182 let (validated_state, state_delta) = TYPES::ValidatedState::genesis(&instance_state);
1183 let high_qc = QuorumCertificate2::genesis(&validated_state, &instance_state, upgrade).await;
1184
1185 Ok(Self {
1186 anchor_leaf: Leaf2::genesis(&validated_state, &instance_state, upgrade.base).await,
1187 anchor_state: Arc::new(validated_state),
1188 anchor_state_delta: Some(Arc::new(state_delta)),
1189 start_view: ViewNumber::new(0),
1190 start_epoch: genesis_epoch_from_version(upgrade.base),
1191 last_actioned_view: ViewNumber::new(0),
1192 saved_proposals: BTreeMap::new(),
1193 high_qc,
1194 next_epoch_high_qc: None,
1195 decided_upgrade_certificate: None,
1196 undecided_leaves: BTreeMap::new(),
1197 undecided_state: BTreeMap::new(),
1198 instance_state,
1199 saved_vid_shares: BTreeMap::new(),
1200 epoch_height,
1201 state_cert: None,
1202 epoch_start_block,
1203 start_epoch_info,
1204 })
1205 }
1206
1207 #[must_use]
1209 fn update_undecided(self) -> Self {
1210 let mut undecided_leaves = self.undecided_leaves.clone();
1211 let mut undecided_state = self.undecided_state.clone();
1212
1213 for proposal in self.saved_proposals.values() {
1214 if proposal.data.view_number() <= self.anchor_leaf.view_number() {
1216 continue;
1217 }
1218
1219 undecided_leaves.insert(
1220 proposal.data.view_number(),
1221 Leaf2::from_quorum_proposal(&proposal.data),
1222 );
1223 }
1224
1225 for leaf in undecided_leaves.values() {
1226 let view_inner = ViewInner::Leaf {
1227 leaf: leaf.commit(),
1228 state: Arc::new(TYPES::ValidatedState::from_header(leaf.block_header())),
1229 delta: None,
1230 epoch: leaf.epoch(self.epoch_height),
1231 };
1232 let view = View { view_inner };
1233
1234 undecided_state.insert(leaf.view_number(), view);
1235 }
1236
1237 Self {
1238 undecided_leaves,
1239 undecided_state,
1240 ..self
1241 }
1242 }
1243
1244 #[allow(clippy::too_many_arguments)]
1251 pub fn load(
1252 instance_state: TYPES::InstanceState,
1253 epoch_height: u64,
1254 epoch_start_block: u64,
1255 start_epoch_info: Vec<InitializerEpochInfo<TYPES>>,
1256 (anchor_leaf, anchor_state, anchor_state_delta): InitializerAnchor<TYPES>,
1257 (start_view, start_epoch): (ViewNumber, Option<EpochNumber>),
1258 (high_qc, next_epoch_high_qc): (
1259 QuorumCertificate2<TYPES>,
1260 Option<NextEpochQuorumCertificate2<TYPES>>,
1261 ),
1262 last_actioned_view: ViewNumber,
1263 saved_proposals: BTreeMap<ViewNumber, Proposal<TYPES, QuorumProposalWrapper<TYPES>>>,
1264 saved_vid_shares: VidShares<TYPES>,
1265 decided_upgrade_certificate: Option<UpgradeCertificate<TYPES>>,
1266 state_cert: Option<LightClientStateUpdateCertificateV2<TYPES>>,
1267 ) -> Self {
1268 let anchor_state = anchor_state.unwrap_or_else(|| {
1269 Arc::new(TYPES::ValidatedState::from_header(
1270 anchor_leaf.block_header(),
1271 ))
1272 });
1273
1274 let initializer = Self {
1275 instance_state,
1276 epoch_height,
1277 epoch_start_block,
1278 anchor_leaf,
1279 anchor_state,
1280 anchor_state_delta,
1281 high_qc,
1282 start_view,
1283 start_epoch,
1284 last_actioned_view,
1285 saved_proposals,
1286 saved_vid_shares,
1287 next_epoch_high_qc,
1288 decided_upgrade_certificate,
1289 undecided_leaves: BTreeMap::new(),
1290 undecided_state: BTreeMap::new(),
1291 state_cert,
1292 start_epoch_info,
1293 };
1294
1295 initializer.update_undecided()
1296 }
1297
1298 pub fn instance_state(&self) -> &TYPES::InstanceState {
1300 &self.instance_state
1301 }
1302
1303 pub fn epoch_height(&self) -> u64 {
1305 self.epoch_height
1306 }
1307
1308 pub fn epoch_start_block(&self) -> u64 {
1310 self.epoch_start_block
1311 }
1312
1313 pub fn anchor_leaf(&self) -> &Leaf2<TYPES> {
1315 &self.anchor_leaf
1316 }
1317
1318 pub fn anchor_state(&self) -> &Arc<TYPES::ValidatedState> {
1320 &self.anchor_state
1321 }
1322
1323 pub fn start_view(&self) -> ViewNumber {
1325 self.start_view
1326 }
1327
1328 pub fn last_actioned_view(&self) -> ViewNumber {
1330 self.last_actioned_view
1331 }
1332
1333 pub fn high_qc(&self) -> &QuorumCertificate2<TYPES> {
1335 &self.high_qc
1336 }
1337
1338 pub fn next_epoch_high_qc(&self) -> Option<&NextEpochQuorumCertificate2<TYPES>> {
1340 self.next_epoch_high_qc.as_ref()
1341 }
1342
1343 pub fn saved_proposals(
1345 &self,
1346 ) -> &BTreeMap<ViewNumber, Proposal<TYPES, QuorumProposalWrapper<TYPES>>> {
1347 &self.saved_proposals
1348 }
1349
1350 pub fn decided_upgrade_certificate(&self) -> Option<&UpgradeCertificate<TYPES>> {
1352 self.decided_upgrade_certificate.as_ref()
1353 }
1354
1355 pub fn state_cert(&self) -> Option<&LightClientStateUpdateCertificateV2<TYPES>> {
1357 self.state_cert.as_ref()
1358 }
1359}
1360
1361async fn load_start_epoch_info<TYPES: NodeType>(
1362 coordinator: &EpochMembershipCoordinator<TYPES>,
1363 start_epoch_info: &Vec<InitializerEpochInfo<TYPES>>,
1364 epoch_height: u64,
1365 epoch_start_block: u64,
1366) {
1367 let membership = coordinator.membership();
1368 let first_epoch_number =
1369 EpochNumber::new(epoch_from_block_number(epoch_start_block, epoch_height));
1370
1371 tracing::warn!("Calling set_first_epoch for epoch {first_epoch_number}");
1372 membership.set_first_epoch(first_epoch_number, INITIAL_DRB_RESULT);
1373
1374 let mut sorted_epoch_info = start_epoch_info.clone();
1375 sorted_epoch_info.sort_by_key(|info| info.epoch);
1376 for epoch_info in sorted_epoch_info {
1377 if let Some(block_header) = &epoch_info.block_header {
1378 tracing::warn!("Calling add_epoch_root for epoch {}", epoch_info.epoch);
1379
1380 coordinator
1381 .add_epoch_root(block_header.clone())
1382 .await
1383 .unwrap_or_else(|err| {
1384 tracing::error!(
1386 "Failed to add epoch root for epoch {}: {err}",
1387 epoch_info.epoch
1388 );
1389 });
1390 }
1391 }
1392
1393 for epoch_info in start_epoch_info {
1394 tracing::warn!("Calling add_drb_result for epoch {}", epoch_info.epoch);
1395 membership.add_drb_result(epoch_info.epoch, epoch_info.drb_result);
1396 }
1397}