1mod proposal_task;
17pub use proposal_task::ProposalTask;
18
19use crate::{
20 Gateway,
21 MAX_BATCH_DELAY,
22 MAX_LEADER_CERTIFICATE_DELAY,
23 MAX_WORKERS,
24 MIN_BATCH_DELAY,
25 PRIMARY_PING_INTERVAL,
26 Sync,
27 Transport,
28 WORKER_PING_INTERVAL,
29 Worker,
30 events::{BatchPropose, BatchSignature, Event},
31 helpers::{
32 PrimaryReceiver,
33 PrimarySender,
34 Proposal,
35 ProposalCache,
36 SignedProposals,
37 Storage,
38 assign_to_worker,
39 assign_to_workers,
40 fmt_id,
41 init_sync_channels,
42 init_worker_channels,
43 now,
44 },
45 spawn_blocking,
46 sync::SyncCallback,
47};
48
49use snarkos_account::Account;
50use snarkos_node_bft_events::PrimaryPing;
51use snarkos_node_bft_ledger_service::LedgerService;
52#[cfg(test)]
53use snarkos_node_network::ConnectionMode;
54use snarkos_node_network::PeerPoolHandling;
55use snarkos_node_sync::{BlockSync, DUMMY_SELF_IP, Ping};
56use snarkos_utilities::{CallbackHandle, NodeDataDir};
57
58use snarkvm::{
59 console::{
60 prelude::*,
61 types::{Address, Field},
62 },
63 ledger::{
64 block::Transaction,
65 narwhal::{BatchCertificate, BatchHeader, Data, Transmission, TransmissionID},
66 puzzle::{Solution, SolutionID},
67 },
68 prelude::{Signature, committee::Committee},
69 utilities::flatten_error,
70};
71
72use anyhow::Context;
73use colored::Colorize;
74use futures::stream::{FuturesUnordered, StreamExt};
75use indexmap::{IndexMap, IndexSet};
76#[cfg(feature = "locktick")]
77use locktick::{
78 parking_lot::{Mutex, RwLock},
79 tokio::RwLock as TRwLock,
80};
81#[cfg(not(feature = "locktick"))]
82use parking_lot::{Mutex, RwLock};
83#[cfg(not(feature = "serial"))]
84use rayon::prelude::*;
85use std::{
86 collections::{HashMap, HashSet},
87 future::Future,
88 net::SocketAddr,
89 pin::Pin,
90 sync::{Arc, OnceLock},
91 time::Instant,
92};
93#[cfg(not(feature = "locktick"))]
94use tokio::sync::RwLock as TRwLock;
95use tokio::{sync::Notify, task::JoinHandle};
96
97#[derive(Debug, PartialEq, Eq)]
99pub enum ProposedBatchState<N: Network> {
100 None,
102 Certifying(Box<Proposal<N>>),
104 Certified(Field<N>),
107}
108
109impl<N: Network> Default for ProposedBatchState<N> {
110 fn default() -> Self {
111 Self::None
112 }
113}
114
115impl<N: Network> ProposedBatchState<N> {
116 pub fn is_none(&self) -> bool {
118 matches!(self, Self::None)
119 }
120
121 pub fn is_proposed(&self) -> bool {
123 matches!(self, Self::Certifying(_))
124 }
125
126 pub fn as_proposal(&self) -> Option<&Proposal<N>> {
128 match self {
129 Self::Certifying(p) => Some(p.as_ref()),
130 _ => None,
131 }
132 }
133}
134
135pub type ProposedBatch<N> = RwLock<ProposedBatchState<N>>;
137
138#[async_trait::async_trait]
141pub trait PrimaryCallback<N: Network>: Send + std::marker::Sync {
142 fn try_advance_to_next_round(&self, current_round: u64) -> bool;
150
151 async fn add_new_certificate(&self, certificate: BatchCertificate<N>) -> Result<()>;
153}
154
155#[derive(Clone)]
158pub struct Primary<N: Network> {
159 sync: Sync<N>,
161 gateway: Gateway<N>,
163 storage: Storage<N>,
165 ledger: Arc<dyn LedgerService<N>>,
167 workers: Arc<OnceLock<Vec<Worker<N>>>>,
169
170 primary_callback: Arc<CallbackHandle<Arc<dyn PrimaryCallback<N>>>>,
172
173 proposed_batch: Arc<ProposedBatch<N>>,
175
176 #[cfg(feature = "metrics")]
179 batch_propose_start: Arc<Mutex<Option<Instant>>>,
180
181 latest_proposal_timestamp: Arc<TRwLock<Option<(u64, i64)>>>,
185
186 signed_proposals: Arc<RwLock<SignedProposals<N>>>,
188
189 handles: Arc<Mutex<Vec<JoinHandle<()>>>>,
191
192 node_data_dir: NodeDataDir,
194
195 proposal_task: ProposalTask<N>,
197
198 round_increment_notify: Arc<Notify>,
201}
202
203impl<N: Network> Primary<N> {
204 pub const MAX_TRANSMISSIONS_TOLERANCE: usize = BatchHeader::<N>::MAX_TRANSMISSIONS_PER_BATCH * 2;
206
207 #[allow(clippy::too_many_arguments)]
209 pub fn new(
210 account: Account<N>,
211 storage: Storage<N>,
212 ledger: Arc<dyn LedgerService<N>>,
213 block_sync: Arc<BlockSync<N>>,
214 ip: Option<SocketAddr>,
215 trusted_validators: &[SocketAddr],
216 trusted_peers_only: bool,
217 node_data_dir: NodeDataDir,
218 dev: Option<u16>,
219 ) -> Result<Self> {
220 let gateway = Gateway::new(
222 account,
223 storage.clone(),
224 ledger.clone(),
225 ip,
226 trusted_validators,
227 trusted_peers_only,
228 node_data_dir.clone(),
229 dev,
230 )?;
231 let sync = Sync::new(gateway.clone(), storage.clone(), ledger.clone(), block_sync);
233
234 Ok(Self {
236 sync,
237 gateway,
238 storage,
239 ledger,
240 node_data_dir,
241 workers: Default::default(),
242 primary_callback: Default::default(),
243 proposed_batch: Default::default(),
244 #[cfg(feature = "metrics")]
245 batch_propose_start: Default::default(),
246 latest_proposal_timestamp: Default::default(),
247 signed_proposals: Default::default(),
248 handles: Default::default(),
249 proposal_task: Default::default(),
250 round_increment_notify: Default::default(),
251 })
252 }
253
254 async fn load_proposal_cache(&self) -> Result<()> {
256 match ProposalCache::<N>::exists(&self.node_data_dir) {
258 true => match ProposalCache::<N>::load(self.gateway.account().address(), &self.node_data_dir) {
260 Ok(proposal_cache) => {
261 let (latest_certificate_round, proposed_batch, signed_proposals, pending_certificates) =
263 proposal_cache.into();
264
265 *self.latest_proposal_timestamp.write().await = Some((latest_certificate_round, now()));
266 *self.proposed_batch.write() = match proposed_batch {
267 Some(p) => ProposedBatchState::Certifying(Box::new(p)),
268 None => ProposedBatchState::None,
269 };
270 *self.signed_proposals.write() = signed_proposals;
271
272 for certificate in pending_certificates {
274 let batch_id = certificate.batch_id();
275 if let Err(err) = self.sync_with_certificate_from_peer::<true>(DUMMY_SELF_IP, certificate).await
279 {
280 let err = err.context(format!(
281 "Failed to load stored certificate {} from proposal cache",
282 fmt_id(batch_id)
283 ));
284 warn!("{}", &flatten_error(err));
285 }
286 }
287 Ok(())
288 }
289 Err(err) => Err(err.context("Failed to read the signed proposals from the file system")),
290 },
291 false => Ok(()),
293 }
294 }
295
296 pub async fn run(
298 &self,
299 ping: Option<Arc<Ping<N>>>,
300 primary_callback: Option<Arc<dyn PrimaryCallback<N>>>,
301 sync_callback: Option<Arc<dyn SyncCallback<N>>>,
302 primary_sender: PrimarySender<N>,
303 primary_receiver: PrimaryReceiver<N>,
304 ) -> Result<()> {
305 info!("Starting the primary instance of the memory pool...");
306
307 if let Some(callback) = primary_callback {
309 self.primary_callback.set(callback)?;
310 }
311
312 let mut worker_senders = IndexMap::new();
314 let mut workers = Vec::new();
316 for id in 0..MAX_WORKERS {
318 let (tx_worker, rx_worker) = init_worker_channels();
320 let worker = Worker::new(
322 id,
323 Arc::new(self.gateway.clone()),
324 self.storage.clone(),
325 self.ledger.clone(),
326 self.proposed_batch.clone(),
327 )?;
328 worker.run(rx_worker);
330 workers.push(worker);
332 worker_senders.insert(id, tx_worker);
334 }
335 if self.workers.set(workers).is_err() {
337 bail!("Workers already set. `Primary::run` cannot be called more than once.");
338 }
339
340 let (sync_sender, sync_receiver) = init_sync_channels();
342 self.sync.initialize(sync_callback)?;
344 self.load_proposal_cache().await?;
346 self.sync.run(ping, sync_receiver).await?;
348 self.gateway.run(primary_sender, worker_senders, Some(sync_sender)).await;
350 self.start_handlers(primary_receiver);
353
354 Ok(())
355 }
356
357 pub fn current_round(&self) -> u64 {
359 self.storage.current_round()
360 }
361
362 pub fn is_synced(&self) -> bool {
364 self.sync.is_synced()
365 }
366
367 pub const fn gateway(&self) -> &Gateway<N> {
369 &self.gateway
370 }
371
372 pub const fn storage(&self) -> &Storage<N> {
374 &self.storage
375 }
376
377 pub const fn ledger(&self) -> &Arc<dyn LedgerService<N>> {
379 &self.ledger
380 }
381
382 pub fn num_workers(&self) -> u8 {
384 u8::try_from(self.workers.get().expect("Primary is not running yet").len()).expect("Too many workers")
385 }
386
387 pub fn workers(&self) -> &[Worker<N>] {
389 self.workers.get().expect("Primary is not running yet")
390 }
391}
392
393impl<N: Network> Primary<N> {
394 pub fn num_unconfirmed_transmissions(&self) -> usize {
396 self.workers().iter().map(|worker| worker.num_transmissions()).sum()
397 }
398
399 pub fn num_unconfirmed_ratifications(&self) -> usize {
401 self.workers().iter().map(|worker| worker.num_ratifications()).sum()
402 }
403
404 pub fn num_unconfirmed_solutions(&self) -> usize {
406 self.workers().iter().map(|worker| worker.num_solutions()).sum()
407 }
408
409 pub fn num_unconfirmed_transactions(&self) -> usize {
411 self.workers().iter().map(|worker| worker.num_transactions()).sum()
412 }
413}
414
415impl<N: Network> Primary<N> {
416 pub fn worker_transmission_ids(&self) -> impl '_ + Iterator<Item = TransmissionID<N>> {
418 self.workers().iter().flat_map(|worker| worker.transmission_ids())
419 }
420
421 pub fn worker_transmissions(&self) -> impl '_ + Iterator<Item = (TransmissionID<N>, Transmission<N>)> {
423 self.workers().iter().flat_map(|worker| worker.transmissions())
424 }
425
426 pub fn worker_solutions(&self) -> impl '_ + Iterator<Item = (SolutionID<N>, Data<Solution<N>>)> {
428 self.workers().iter().flat_map(|worker| worker.solutions())
429 }
430
431 pub fn worker_transactions(&self) -> impl '_ + Iterator<Item = (N::TransactionID, Data<Transaction<N>>)> {
433 self.workers().iter().flat_map(|worker| worker.transactions())
434 }
435}
436
437impl<N: Network> Primary<N> {
438 pub fn clear_worker_solutions(&self) {
440 self.workers().iter().for_each(Worker::clear_solutions);
441 }
442}
443
444#[async_trait::async_trait]
445impl<N: Network> proposal_task::BatchPropose for Primary<N> {
446 fn current_round(&self) -> u64 {
447 Primary::current_round(self)
448 }
449
450 fn wait_for_synced_if_syncing(&self) -> Option<futures::future::BoxFuture<'_, ()>> {
451 self.sync.wait_for_synced_if_syncing()
452 }
453
454 fn is_synced(&self) -> bool {
455 self.sync.is_synced()
456 }
457
458 async fn propose_batch(&self) -> Result<bool> {
471 let mut lock_guard = self.latest_proposal_timestamp.write().await;
476
477 if let Err(err) = self
479 .check_proposed_batch_for_expiration()
480 .with_context(|| "Failed to check the proposed batch for expiration")
481 {
482 warn!("{}", flatten_error(&err));
483 return Ok(false);
484 }
485
486 let round = self.current_round();
488 let previous_round = round.saturating_sub(1);
490
491 ensure!(round > 0, "Round 0 cannot have transaction batches");
495
496 if let Some((latest_round, _)) = &*lock_guard
498 && round < *latest_round
499 {
500 warn!("Cannot propose a batch for round {round} - the latest proposal cache round is {latest_round}");
501 return Ok(false);
502 }
503
504 match &*self.proposed_batch.read() {
506 ProposedBatchState::Certifying(proposal) => {
507 if round < proposal.round()
509 || proposal
510 .batch_header()
511 .previous_certificate_ids()
512 .iter()
513 .any(|id| !self.storage.contains_certificate(*id))
514 {
515 warn!(
516 "Cannot propose a batch for round {} - the current storage (round {round}) is not caught up to the proposed batch.",
517 proposal.round(),
518 );
519 return Ok(false);
520 }
521 let event = Event::BatchPropose(proposal.batch_header().clone().into());
524 for address in proposal.nonsigners(&self.ledger.get_committee_lookback_for_round(proposal.round())?) {
526 match self.gateway.resolver().read().get_peer_ip_for_address(address) {
528 Some(peer_ip) => {
530 let (gateway, event_, round) = (self.gateway.clone(), event.clone(), proposal.round());
531 tokio::spawn(async move {
532 debug!("Resending batch proposal for round {round} to peer '{peer_ip}'");
533 if gateway.send(peer_ip, event_).await.is_none() {
535 warn!("Failed to resend batch proposal for round {round} to peer '{peer_ip}'");
536 }
537 });
538 }
539 None => continue,
540 }
541 }
542 debug!("Proposed batch for round {} is still valid", proposal.round());
543 return Ok(false);
544 }
545 ProposedBatchState::Certified(_) => {
547 debug!("Cannot propose a batch for round {round} - a batch is currently being certified");
548 return Ok(false);
549 }
550 ProposedBatchState::None => {
551 }
553 }
554
555 #[cfg(feature = "metrics")]
556 metrics::gauge(metrics::bft::PROPOSAL_ROUND, round as f64);
557
558 if let Some((_, latest_timestamp)) = &*lock_guard
560 && !self.check_own_proposal_timestamp(previous_round, *latest_timestamp, now())?
561 {
562 return Ok(false);
563 }
564
565 if self.storage.contains_certificate_in_round_from(round, self.gateway.account().address()) {
567 if let Some(cb) = &*self.primary_callback.get_ref() {
569 match cb.try_advance_to_next_round(self.current_round()) {
570 true => (), false => return Ok(false),
572 }
573 }
574 debug!("Primary is safely skipping {}", format!("(round {round} was already certified)").dimmed());
575 return Ok(false);
576 }
577
578 if let Some((latest_round, _)) = &*lock_guard
584 && *latest_round == round
585 {
586 debug!("Primary is safely skipping a batch proposal - round {round} already proposed");
587 return Ok(false);
588 }
589
590 let committee_lookback = self.ledger.get_committee_lookback_for_round(round)?;
592 {
594 let mut connected_validators = self.gateway.connected_addresses();
596 connected_validators.insert(self.gateway.account().address());
598 if !committee_lookback.is_quorum_threshold_reached(&connected_validators) {
600 debug!(
601 "Primary is safely skipping a batch proposal for round {round} {}",
602 "(please connect to more validators)".dimmed()
603 );
604 trace!("Primary is connected to {} validators", connected_validators.len() - 1);
605 return Ok(false);
606 }
607 }
608
609 let previous_certificates = self.storage.get_certificates_for_round(previous_round);
611
612 let mut is_ready = previous_round == 0;
615 if previous_round > 0 {
617 let Ok(previous_committee_lookback) = self.ledger.get_committee_lookback_for_round(previous_round) else {
619 bail!("Cannot propose a batch for round {round}: the committee lookback is not known yet")
620 };
621 let authors = previous_certificates.iter().map(BatchCertificate::author).collect();
623 if previous_committee_lookback.is_quorum_threshold_reached(&authors) {
625 is_ready = true;
626 }
627 #[cfg(feature = "test_network")]
628 {
629 if let Some(dev_committee) = self.ledger.dev_committee_for_round(previous_round)? {
631 if round <= dev_committee.starting_round() {
632 is_ready = true;
633 }
634 }
635 }
636 }
637 if !is_ready {
639 debug!(
640 "Primary is safely skipping a batch proposal for round {round} {}",
641 format!("(previous round {previous_round} has not reached quorum)").dimmed()
642 );
643 return Ok(false);
644 }
645
646 let mut transmissions: IndexMap<_, _> = Default::default();
648 let mut proposal_cost = 0u64;
650 debug_assert_eq!(MAX_WORKERS, 1);
654
655 'outer: for worker in self.workers().iter() {
656 let mut num_worker_transmissions = 0usize;
657
658 while let Some((id, transmission)) = worker.remove_front() {
659 if transmissions.len() >= BatchHeader::<N>::MAX_TRANSMISSIONS_PER_BATCH {
661 worker.insert_front(id, transmission);
663 break 'outer;
664 }
665
666 if num_worker_transmissions >= Worker::<N>::MAX_TRANSMISSIONS_PER_WORKER {
668 worker.insert_front(id, transmission);
670 continue 'outer;
671 }
672
673 if self.ledger.contains_transmission(&id).unwrap_or(true) {
675 trace!("Proposing - Skipping transmission '{}' - Already in ledger", fmt_id(id));
676 continue;
677 }
678
679 if !transmissions.is_empty() && self.storage.contains_transmission(id) {
683 trace!("Proposing - Skipping transmission '{}' - Already in storage", fmt_id(id));
684 continue;
685 }
686
687 match (id, transmission.clone()) {
689 (TransmissionID::Solution(solution_id, checksum), Transmission::Solution(solution)) => {
690 if !matches!(solution.to_checksum::<N>(), Ok(solution_checksum) if solution_checksum == checksum)
692 {
693 trace!("Proposing - Skipping solution '{}' - Checksum mismatch", fmt_id(solution_id));
694 continue;
695 }
696 if let Err(e) = self.ledger.check_solution_basic(solution_id, solution).await {
698 trace!("Proposing - Skipping solution '{}' - {e}", fmt_id(solution_id));
699 continue;
700 }
701 }
702 (TransmissionID::Transaction(transaction_id, checksum), Transmission::Transaction(transaction)) => {
703 if !matches!(transaction.to_checksum::<N>(), Ok(transaction_checksum) if transaction_checksum == checksum )
705 {
706 trace!("Proposing - Skipping transaction '{}' - Checksum mismatch", fmt_id(transaction_id));
707 continue;
708 }
709
710 let transaction = spawn_blocking!({
712 match transaction {
713 Data::Object(transaction) => Ok(transaction),
714 Data::Buffer(bytes) => Ok(Transaction::<N>::read_le(
715 &mut bytes.take(N::LATEST_MAX_TRANSACTION_SIZE() as u64),
716 )?),
717 }
718 })?;
719
720 let current_block_height = self.ledger.latest_block_height();
722 let consensus_version = N::CONSENSUS_VERSION(current_block_height)?;
723
724 let Ok(cost) = self.ledger.transaction_spend_in_microcredits(&transaction, consensus_version)
727 else {
728 debug!(
729 "Proposing - Skipping and discarding transaction '{}' - Unable to compute transaction spent cost",
730 fmt_id(transaction_id)
731 );
732 continue;
733 };
734
735 if let Err(e) = self.ledger.check_transaction_basic(transaction_id, transaction).await {
737 trace!("Proposing - Skipping transaction '{}' - {e}", fmt_id(transaction_id));
738 continue;
739 }
740
741 let Some(next_proposal_cost) = proposal_cost.checked_add(cost) else {
744 debug!(
745 "Proposing - Skipping and discarding transaction '{}' - Proposal cost overflowed",
746 fmt_id(transaction_id)
747 );
748 continue;
749 };
750
751 let batch_spend_limit = BatchHeader::<N>::batch_spend_limit(current_block_height);
753 if next_proposal_cost > batch_spend_limit {
754 debug!(
755 "Proposing - Skipping transaction '{}' - Batch spend limit surpassed ({next_proposal_cost} > {})",
756 fmt_id(transaction_id),
757 batch_spend_limit
758 );
759
760 worker.insert_front(id, transmission);
762 break 'outer;
763 }
764
765 proposal_cost = next_proposal_cost;
767 }
768
769 (TransmissionID::Ratification, Transmission::Ratification) => continue,
772 _ => continue,
774 }
775
776 transmissions.insert(id, transmission);
778 num_worker_transmissions = num_worker_transmissions.saturating_add(1);
779 }
780 }
781
782 let current_timestamp = now();
784
785 info!("Proposing a batch with {} transmissions for round {round}...", transmissions.len());
787
788 *lock_guard = Some((round, current_timestamp));
790 let private_key = *self.gateway.account().private_key();
792 let committee_id = committee_lookback.id();
794 let transmission_ids = transmissions.keys().copied().collect();
796 let previous_certificate_ids = previous_certificates.into_iter().map(|c| c.id()).collect();
798 let (batch_header, proposal) = spawn_blocking!(BatchHeader::new(
800 &private_key,
801 round,
802 current_timestamp,
803 committee_id,
804 transmission_ids,
805 previous_certificate_ids,
806 &mut rand::rng()
807 ))
808 .and_then(|batch_header| {
809 Proposal::new(committee_lookback, batch_header.clone(), transmissions.clone())
810 .map(|proposal| (batch_header, proposal))
811 })
812 .inspect_err(|_| {
813 if let Err(err) = self.reinsert_transmissions_into_workers(transmissions) {
815 error!("{}", flatten_error(err.context("Failed to reinsert transmissions")));
816 }
817 })?;
818
819 self.gateway.broadcast(Event::BatchPropose(batch_header.into()));
821 *self.proposed_batch.write() = ProposedBatchState::Certifying(Box::new(proposal));
823 #[cfg(feature = "metrics")]
825 {
826 *self.batch_propose_start.lock() = Some(Instant::now());
827 }
828
829 Ok(true)
830 }
831}
832
833impl<N: Network> Primary<N> {
834 async fn process_batch_propose_from_peer(&self, peer_ip: SocketAddr, batch_propose: BatchPropose<N>) -> Result<()> {
844 let BatchPropose { round: batch_round, batch_header } = batch_propose;
845
846 let batch_header = spawn_blocking!(batch_header.deserialize_blocking())?;
848 if batch_round != batch_header.round() {
850 self.gateway.disconnect(peer_ip);
852 bail!("Malicious peer - proposed round {batch_round}, but sent batch for round {}", batch_header.round());
853 }
854
855 let batch_author = batch_header.author();
857
858 match self.gateway.resolve_to_aleo_addr(peer_ip) {
860 Some(address) => {
862 if address != batch_author {
863 self.gateway.disconnect(peer_ip);
865 bail!("Malicious peer - proposed batch from a different validator ({batch_author})");
866 }
867 }
868 None => bail!("Batch proposal from a disconnected validator"),
869 }
870 if !self.gateway.is_authorized_validator_address(batch_author) {
872 self.gateway.disconnect(peer_ip);
874 bail!("Malicious peer - proposed batch from a non-committee member ({batch_author})");
875 }
876 if self.gateway.account().address() == batch_author {
878 bail!("Invalid peer - proposed batch from myself ({batch_author})");
879 }
880
881 let expected_committee_id = self.ledger.get_committee_lookback_for_round(batch_round)?.id();
887 if expected_committee_id != batch_header.committee_id() {
888 self.gateway.disconnect(peer_ip);
890 bail!(
891 "Malicious peer - proposed batch has a different committee ID ({expected_committee_id} != {})",
892 batch_header.committee_id()
893 );
894 }
895
896 if let Some((signed_round, signed_batch_id, signature)) =
898 self.signed_proposals.read().get(&batch_author).copied()
899 {
900 if signed_round > batch_header.round() {
903 bail!(
904 "Peer ({batch_author}) proposed a batch for a previous round ({}), latest signed round: {signed_round}",
905 batch_header.round()
906 );
907 }
908
909 if signed_round == batch_header.round() && signed_batch_id != batch_header.batch_id() {
911 bail!("Peer ({batch_author}) proposed another batch for the same round ({signed_round})");
912 }
913 if signed_round == batch_header.round() && signed_batch_id == batch_header.batch_id() {
916 let gateway = self.gateway.clone();
917 tokio::spawn(async move {
918 debug!("Resending a signature for a batch in round {batch_round} from '{peer_ip}'");
919 let event = Event::BatchSignature(BatchSignature::new(batch_header.batch_id(), signature));
920 if gateway.send(peer_ip, event).await.is_none() {
922 warn!("Failed to resend a signature for a batch in round {batch_round} to '{peer_ip}'");
923 }
924 });
925 return Ok(());
927 }
928 }
929
930 if self.storage.contains_batch(batch_header.batch_id()) {
933 debug!(
934 "Primary is safely skipping a batch proposal from '{peer_ip}' - {}",
935 format!("batch for round {batch_round} already exists in storage").dimmed()
936 );
937 return Ok(());
938 }
939
940 let previous_round = batch_round.saturating_sub(1);
942 if let Err(err) = self.check_peer_proposal_timestamp(previous_round, batch_author, batch_header.timestamp()) {
944 self.gateway.disconnect(peer_ip);
946 return Err(err.context(format!("Malicious behavior of peer '{peer_ip}'")));
947 }
948
949 if batch_header.contains(TransmissionID::Ratification) {
951 self.gateway.disconnect(peer_ip);
953 bail!(
954 "Malicious peer - proposed batch contains an unsupported ratification transmissionID from '{peer_ip}'",
955 );
956 }
957
958 let mut missing_transmissions =
960 self.sync_with_batch_header_from_peer::<false, true>(peer_ip, &batch_header).await?;
961
962 if let Err(err) = cfg_iter_mut!(&mut missing_transmissions).try_for_each(|(transmission_id, transmission)| {
964 self.ledger.ensure_transmission_is_well_formed(*transmission_id, transmission)
966 }) {
967 let err = err.context(format!(
968 "Batch propose at round {batch_round} from '{peer_ip}' contains an invalid transmission"
969 ));
970 debug!("{}", flatten_error(err));
971 return Ok(());
972 }
973
974 if let Err(e) = self.ensure_is_signing_round(batch_round) {
978 debug!("{e} from '{peer_ip}'");
980 return Ok(());
981 }
982
983 let (storage, header) = (self.storage.clone(), batch_header.clone());
985
986 let Some(missing_transmissions) =
988 spawn_blocking!(storage.check_batch_header(&header, missing_transmissions, Default::default()))?
989 else {
990 return Ok(());
991 };
992
993 self.insert_missing_transmissions_into_workers(peer_ip, missing_transmissions.into_iter())?;
995
996 let batch_id = batch_header.batch_id();
1000 let account = self.gateway.account().clone();
1002 let signature = spawn_blocking!(account.sign(&[batch_id], &mut rand::rng()))?;
1003
1004 match self.signed_proposals.write().0.entry(batch_author) {
1010 std::collections::hash_map::Entry::Occupied(mut entry) => {
1011 if entry.get().0 == batch_round {
1016 return Ok(());
1017 }
1018 entry.insert((batch_round, batch_id, signature));
1020 }
1021 std::collections::hash_map::Entry::Vacant(entry) => {
1023 entry.insert((batch_round, batch_id, signature));
1025 }
1026 };
1027
1028 let self_ = self.clone();
1030 tokio::spawn(async move {
1031 let event = Event::BatchSignature(BatchSignature::new(batch_id, signature));
1032 if self_.gateway.send(peer_ip, event).await.is_some() {
1034 debug!("Signed a batch for round {batch_round} from '{peer_ip}'");
1035 }
1036 });
1037
1038 Ok(())
1039 }
1040
1041 fn add_signature_to_batch(
1052 &self,
1053 state: ProposedBatchState<N>,
1054 peer_ip: SocketAddr,
1055 batch_id: Field<N>,
1056 signature: Signature<N>,
1057 ) -> (Result<Option<Proposal<N>>>, ProposedBatchState<N>) {
1058 match state {
1059 ProposedBatchState::Certifying(mut proposal) if proposal.batch_id() == batch_id => {
1060 let inner: Result<bool> = (|| {
1063 let committee_lookback = self.ledger.get_committee_lookback_for_round(proposal.round())?;
1064 let Some(signer) = self.gateway.resolve_to_aleo_addr(peer_ip) else {
1065 bail!("Signature is from a disconnected validator");
1066 };
1067 let new_signature = proposal.add_signature(signer, signature, &committee_lookback)?;
1068 if new_signature {
1069 info!("Received a batch signature for round {} from '{peer_ip}'", proposal.round());
1070 Ok(proposal.is_quorum_threshold_reached(&committee_lookback))
1071 } else {
1072 debug!(
1073 "Received duplicated signature from '{peer_ip}' for batch \
1074 {batch_id} in round {round}",
1075 round = proposal.round()
1076 );
1077 Ok(false)
1078 }
1079 })();
1080 match inner {
1081 Ok(true) => {
1082 let certified_id = proposal.batch_id();
1083 (Ok(Some(*proposal)), ProposedBatchState::Certified(certified_id))
1084 }
1085 Ok(false) => (Ok(None), ProposedBatchState::Certifying(proposal)),
1086 Err(e) => (Err(e), ProposedBatchState::Certifying(proposal)),
1087 }
1088 }
1089 ProposedBatchState::Certifying(proposal) => {
1090 if self.storage.contains_batch(batch_id) {
1092 debug!(
1093 "Primary is safely skipping a batch signature from {peer_ip} for \
1094 round {} - batch is already certified",
1095 proposal.round()
1096 );
1097 (Ok(None), ProposedBatchState::Certifying(proposal))
1098 } else {
1099 let expected_id = proposal.batch_id();
1100 let round = proposal.round();
1101 (
1102 Err(anyhow!("Unknown batch ID '{batch_id}', expected '{expected_id}' for round {round}")),
1103 ProposedBatchState::Certifying(proposal),
1104 )
1105 }
1106 }
1107 ProposedBatchState::Certified(id) if id == batch_id => {
1108 debug!(
1110 "Skipping batch signature from {peer_ip} for batch '{batch_id}' - \
1111 already received sufficient signatures"
1112 );
1113 (Ok(None), ProposedBatchState::Certified(id))
1114 }
1115 ProposedBatchState::Certified(id) => {
1116 let result = if self.storage.contains_batch(batch_id) {
1117 warn!("Received signature for an older batch {batch_id}");
1119 Ok(None)
1120 } else {
1121 Err(anyhow!("Unknown batch ID '{batch_id}'"))
1122 };
1123
1124 (result, ProposedBatchState::Certified(id))
1125 }
1126 ProposedBatchState::None => {
1127 let result = if self.storage.contains_batch(batch_id) {
1128 warn!("Received signature for an older batch {batch_id}");
1130 Ok(None)
1131 } else {
1132 Err(anyhow!("Unknown batch ID '{batch_id}'"))
1133 };
1134
1135 (result, ProposedBatchState::None)
1136 }
1137 }
1138 }
1139
1140 async fn process_batch_signature_from_peer(
1149 &self,
1150 peer_ip: SocketAddr,
1151 batch_signature: BatchSignature<N>,
1152 ) -> Result<()> {
1153 self.check_proposed_batch_for_expiration()?;
1155
1156 let BatchSignature { batch_id, signature } = batch_signature;
1158
1159 let signer = signature.to_address();
1161
1162 match self.gateway.resolve_to_aleo_addr(peer_ip) {
1164 Some(address) => {
1166 if address != signer {
1167 self.gateway.disconnect(peer_ip);
1169 bail!("Malicious peer - batch signature is from a different validator ({signer})");
1170 }
1171 }
1172 None => bail!("Batch signature from a disconnected validator"),
1173 }
1174 if self.gateway.account().address() == signer {
1176 bail!("Invalid peer - received a batch signature from myself ({signer})");
1177 }
1178
1179 let self_ = self.clone();
1180 let Some(proposal) = spawn_blocking!({
1181 let mut proposed_batch = self_.proposed_batch.write();
1183
1184 let (result, new_state) =
1185 self_.add_signature_to_batch(std::mem::take(&mut *proposed_batch), peer_ip, batch_id, signature);
1186 *proposed_batch = new_state;
1187 result
1188 })?
1189 else {
1190 return Ok(());
1191 };
1192
1193 info!("Quorum threshold reached - Preparing to certify our batch for round {}...", proposal.round());
1196
1197 let committee_lookback = self.ledger.get_committee_lookback_for_round(proposal.round())?;
1199 if let Err(e) = self.store_and_broadcast_certificate(&proposal, &committee_lookback).await {
1202 self.reinsert_transmissions_into_workers(proposal.into_transmissions())?;
1204 return Err(e);
1205 }
1206
1207 #[cfg(feature = "metrics")]
1208 metrics::increment_gauge(metrics::bft::CERTIFIED_BATCHES, 1.0);
1209 Ok(())
1210 }
1211
1212 async fn process_batch_certificate_from_peer(
1219 &self,
1220 peer_ip: SocketAddr,
1221 certificate: BatchCertificate<N>,
1222 ) -> Result<()> {
1223 if !self.gateway.is_authorized_validator_ip(peer_ip) {
1225 self.gateway.disconnect(peer_ip);
1227 bail!("Malicious peer - Received a batch certificate from an unauthorized validator IP ({peer_ip})");
1228 }
1229 if self.storage.contains_certificate(certificate.id()) {
1231 return Ok(());
1232 } else if !self.storage.contains_unprocessed_certificate(certificate.id()) {
1234 self.storage.insert_unprocessed_certificate(certificate.clone())?;
1235 }
1236
1237 let author = certificate.author();
1239 let certificate_round = certificate.round();
1241 let committee_id = certificate.committee_id();
1243
1244 if self.gateway.account().address() == author {
1246 bail!("Received a batch certificate for myself ({author})");
1247 }
1248
1249 self.storage.check_incoming_certificate(&certificate)?;
1251
1252 self.sync_with_certificate_from_peer::<false>(peer_ip, certificate).await?;
1264
1265 let committee_lookback = self.ledger.get_committee_lookback_for_round(certificate_round)?;
1270
1271 let authors = self.storage.get_certificate_authors_for_round(certificate_round);
1273 let is_quorum = committee_lookback.is_quorum_threshold_reached(&authors);
1275
1276 let expected_committee_id = committee_lookback.id();
1278 if expected_committee_id != committee_id {
1279 self.gateway.disconnect(peer_ip);
1281 bail!("Batch certificate has a different committee ID ({expected_committee_id} != {committee_id})");
1282 }
1283
1284 let should_advance = match &*self.latest_proposal_timestamp.read().await {
1288 Some((latest_round, _)) => *latest_round < certificate_round,
1290 None => true,
1292 };
1293
1294 let current_round = self.current_round();
1296
1297 if is_quorum && should_advance && certificate_round >= current_round {
1299 self.round_increment_notify.notify_one();
1301 }
1302 Ok(())
1303 }
1304}
1305
1306impl<N: Network> Primary<N> {
1307 fn start_handlers(&self, primary_receiver: PrimaryReceiver<N>) {
1316 let PrimaryReceiver {
1317 mut rx_batch_propose,
1318 mut rx_batch_signature,
1319 mut rx_batch_certified,
1320 mut rx_primary_ping,
1321 mut rx_unconfirmed_solution,
1322 mut rx_unconfirmed_transaction,
1323 } = primary_receiver;
1324
1325 let self_ = self.clone();
1327 self.spawn(async move {
1328 loop {
1329 tokio::time::sleep(PRIMARY_PING_INTERVAL).await;
1331
1332 let self__ = self_.clone();
1334 let block_locators = match spawn_blocking!(self__.sync.get_block_locators()) {
1335 Ok(block_locators) => block_locators,
1336 Err(e) => {
1337 warn!("Failed to retrieve block locators - {e}");
1338 continue;
1339 }
1340 };
1341
1342 let primary_certificate = {
1344 let primary_address = self_.gateway.account().address();
1346
1347 let mut certificate = None;
1349 let mut current_round = self_.current_round();
1350 while certificate.is_none() {
1351 if current_round == 0 {
1353 break;
1354 }
1355 if let Some(primary_certificate) =
1357 self_.storage.get_certificate_for_round_with_author(current_round, primary_address)
1358 {
1359 certificate = Some(primary_certificate);
1360 } else {
1362 current_round = current_round.saturating_sub(1);
1363 }
1364 }
1365
1366 match certificate {
1368 Some(certificate) => certificate,
1369 None => continue,
1371 }
1372 };
1373
1374 let primary_ping = PrimaryPing::from((<Event<N>>::VERSION, block_locators, primary_certificate));
1376 self_.gateway.broadcast(Event::PrimaryPing(primary_ping));
1378 }
1379 });
1380
1381 let self_ = self.clone();
1383 self.spawn(async move {
1384 while let Some((peer_ip, primary_certificate)) = rx_primary_ping.recv().await {
1385 if self_.sync.is_synced() {
1387 trace!("Processing new primary ping from '{peer_ip}'");
1388 } else {
1389 trace!("Skipping a primary ping from '{peer_ip}' {}", "(node is syncing)".dimmed());
1390 continue;
1391 }
1392
1393 {
1395 let self_ = self_.clone();
1396 tokio::spawn(async move {
1397 let Ok(primary_certificate) = spawn_blocking!(primary_certificate.deserialize_blocking())
1399 else {
1400 warn!("Failed to deserialize primary certificate in 'PrimaryPing' from '{peer_ip}'");
1401 return;
1402 };
1403 let id = fmt_id(primary_certificate.id());
1405 let round = primary_certificate.round();
1406 if let Err(e) = self_.process_batch_certificate_from_peer(peer_ip, primary_certificate).await {
1407 warn!("Cannot process a primary certificate '{id}' at round {round} in a 'PrimaryPing' from '{peer_ip}' - {e}");
1408 }
1409 });
1410 }
1411 }
1412 });
1413
1414 let self_ = self.clone();
1416 self.spawn(async move {
1417 loop {
1418 tokio::time::sleep(WORKER_PING_INTERVAL).await;
1419 if !self_.sync.is_synced() {
1421 trace!("Skipping worker ping(s) {}", "(node is syncing)".dimmed());
1422 continue;
1423 }
1424 for worker in self_.workers() {
1426 worker.broadcast_ping();
1427 }
1428 }
1429 });
1430
1431 let proposal_task = self.proposal_task.clone();
1433 let self_ = self.clone();
1434 self.spawn(async move { proposal_task.run(self_).await });
1435
1436 let self_ = self.clone();
1438 self.spawn(async move {
1439 while let Some((peer_ip, batch_propose)) = rx_batch_propose.recv().await {
1440 if !self_.sync.is_synced() {
1442 trace!("Skipping a batch proposal from '{peer_ip}' {}", "(node is syncing)".dimmed());
1443 continue;
1444 }
1445
1446 let self_ = self_.clone();
1448 tokio::spawn(async move {
1449 let round = batch_propose.round;
1451 if let Err(err) = self_.process_batch_propose_from_peer(peer_ip, batch_propose).await {
1452 let err = err.context(format!("Cannot sign a batch at round {round} from '{peer_ip}'"));
1453 warn!("{}", flatten_error(err));
1454 }
1455 });
1456 }
1457 });
1458
1459 let self_ = self.clone();
1461 self.spawn(async move {
1462 while let Some((peer_ip, batch_signature)) = rx_batch_signature.recv().await {
1463 if !self_.sync.is_synced() {
1465 trace!("Skipping a batch signature from '{peer_ip}' {}", "(node is syncing)".dimmed());
1466 continue;
1467 }
1468 let id = fmt_id(batch_signature.batch_id);
1474 if let Err(err) = self_.process_batch_signature_from_peer(peer_ip, batch_signature).await {
1475 let err = err.context(format!("Cannot store a signature for batch '{id}' from '{peer_ip}'"));
1476 warn!("{}", flatten_error(err));
1477 }
1478 }
1479 });
1480
1481 let self_ = self.clone();
1483 self.spawn(async move {
1484 while let Some((peer_ip, batch_certificate)) = rx_batch_certified.recv().await {
1485 if !self_.sync.is_synced() {
1487 trace!("Skipping a certified batch from '{peer_ip}' {}", "(node is syncing)".dimmed());
1488 continue;
1489 }
1490 let self_ = self_.clone();
1492 tokio::spawn(async move {
1493 let Ok(batch_certificate) = spawn_blocking!(batch_certificate.deserialize_blocking()) else {
1495 warn!("Failed to deserialize the batch certificate from '{peer_ip}'");
1496 return;
1497 };
1498 let id = fmt_id(batch_certificate.id());
1500 let round = batch_certificate.round();
1501 if let Err(err) = self_.process_batch_certificate_from_peer(peer_ip, batch_certificate).await {
1502 warn!(
1503 "{}",
1504 flatten_error(err.context(format!(
1505 "Cannot store a certificate '{id}' for round {round} from '{peer_ip}'"
1506 )))
1507 );
1508 }
1509 });
1510 }
1511 });
1512
1513 let self_ = self.clone();
1516 self.spawn(async move {
1517 loop {
1518 let round_start = Instant::now();
1519 let current_round = self_.current_round();
1520
1521 while self_.current_round() == current_round {
1523 let mut futures: Vec<Pin<Box<dyn Future<Output = ()> + Send>>> =
1524 vec![Box::pin(self_.round_increment_notify.notified())];
1525
1526 if let Some(remaining_delay) = MAX_BATCH_DELAY.checked_sub(round_start.elapsed())
1527 && !remaining_delay.is_zero()
1528 {
1529 futures.push(Box::pin(tokio::time::sleep(remaining_delay)));
1530 }
1531 futures.push(Box::pin(tokio::time::sleep(MAX_LEADER_CERTIFICATE_DELAY)));
1536 if !self_.sync.is_synced() {
1537 futures.push(Box::pin(self_.sync.wait_for_synced()));
1538 }
1539 let _ = futures::future::select_all(futures).await;
1540
1541 if !self_.sync.is_synced() {
1542 trace!("Skipping round increment {}", "(node is syncing)".dimmed());
1543 continue;
1544 }
1545
1546 let next_round = current_round.saturating_add(1);
1547 let is_quorum_threshold_reached = {
1548 let authors = self_.storage.get_certificate_authors_for_round(current_round);
1549 if authors.is_empty() {
1550 continue;
1551 }
1552 let Ok(committee_lookback) = self_.ledger.get_committee_lookback_for_round(current_round)
1553 else {
1554 warn!("Failed to retrieve the committee lookback for round {current_round}");
1555 continue;
1556 };
1557 committee_lookback.is_quorum_threshold_reached(&authors)
1558 };
1559
1560 if is_quorum_threshold_reached {
1561 debug!("Quorum threshold reached for round {current_round}");
1562 if let Err(err) = self_.try_increment_to_the_next_round(next_round).await {
1563 warn!("{}", flatten_error(err.context("Failed to increment to the next round")));
1564 }
1565 }
1566 }
1567 }
1568 });
1569
1570 let self_ = self.clone();
1572 self.spawn(async move {
1573 while let Some((solution_id, solution, callback)) = rx_unconfirmed_solution.recv().await {
1574 let Ok(checksum) = solution.to_checksum::<N>() else {
1576 error!("Failed to compute the checksum for the unconfirmed solution");
1577 continue;
1578 };
1579 let Ok(worker_id) = assign_to_worker((solution_id, checksum), self_.num_workers()) else {
1581 error!("Unable to determine the worker ID for the unconfirmed solution");
1582 continue;
1583 };
1584 let self_ = self_.clone();
1585 tokio::spawn(async move {
1586 let worker = &self_.workers()[worker_id as usize];
1588 let result = worker.process_unconfirmed_solution(solution_id, solution).await;
1590 callback.send(result).ok();
1592 });
1593 }
1594 });
1595
1596 let self_ = self.clone();
1598 self.spawn(async move {
1599 while let Some((transaction_id, transaction, callback)) = rx_unconfirmed_transaction.recv().await {
1600 trace!("Primary - Received an unconfirmed transaction '{}'", fmt_id(transaction_id));
1601 let Ok(checksum) = transaction.to_checksum::<N>() else {
1603 error!("Failed to compute the checksum for the unconfirmed transaction");
1604 continue;
1605 };
1606 let Ok(worker_id) = assign_to_worker::<N>((&transaction_id, &checksum), self_.num_workers()) else {
1608 error!("Unable to determine the worker ID for the unconfirmed transaction");
1609 continue;
1610 };
1611 let self_ = self_.clone();
1612 tokio::spawn(async move {
1613 let worker = &self_.workers().get(worker_id as usize).expect("Invalid worker ID");
1615 let result = worker.process_unconfirmed_transaction(transaction_id, transaction).await;
1617 callback.send(result).ok();
1619 });
1620 }
1621 });
1622 }
1623
1624 fn check_proposed_batch_for_expiration(&self) -> Result<()> {
1626 let is_expired = match &*self.proposed_batch.read() {
1629 ProposedBatchState::Certifying(proposal) => proposal.round() < self.current_round(),
1630 _ => false,
1631 };
1632 if is_expired {
1634 let old = std::mem::replace(&mut *self.proposed_batch.write(), ProposedBatchState::None);
1636 if let ProposedBatchState::Certifying(proposal) = old {
1637 debug!("Cleared expired proposal for round {}", proposal.round());
1638 self.reinsert_transmissions_into_workers(proposal.into_transmissions())?;
1639 }
1640 }
1641 Ok(())
1642 }
1643
1644 async fn try_increment_to_the_next_round(&self, next_round: u64) -> Result<()> {
1646 if self.current_round() + self.storage.max_gc_rounds() >= next_round {
1648 let mut fast_forward_round = self.current_round();
1649 while fast_forward_round < next_round.saturating_sub(1) {
1651 fast_forward_round = self.storage.increment_to_next_round(fast_forward_round)?;
1653 *self.proposed_batch.write() = ProposedBatchState::None;
1655 }
1656 }
1657
1658 let current_round = self.current_round();
1660 if current_round < next_round {
1662 let is_ready = if let Some(cb) = self.primary_callback.get() {
1664 cb.try_advance_to_next_round(current_round)
1665 }
1666 else {
1668 self.storage.increment_to_next_round(current_round)?;
1670 true
1672 };
1673
1674 if is_ready && self.is_synced() {
1676 debug!("Primary is ready to propose the next round");
1677 self.proposal_task.signal();
1678 } else {
1679 debug!("Primary is not ready to propose the next round");
1680 }
1681 }
1682 Ok(())
1683 }
1684
1685 fn ensure_is_signing_round(&self, batch_round: u64) -> Result<()> {
1689 let current_round = self.current_round();
1691 if current_round + self.storage.max_gc_rounds() <= batch_round {
1693 bail!("Round {batch_round} is too far in the future")
1694 }
1695 if current_round > batch_round + 1 {
1699 bail!("Primary is on round {current_round}, and no longer signing for round {batch_round}")
1700 }
1701 if let ProposedBatchState::Certifying(proposal) = &*self.proposed_batch.read()
1703 && proposal.round() > batch_round
1704 {
1705 bail!("Our primary at round {} is no longer signing for round {batch_round}", proposal.round())
1706 }
1707 Ok(())
1708 }
1709
1710 fn check_peer_proposal_timestamp(&self, previous_round: u64, author: Address<N>, timestamp: i64) -> Result<()> {
1713 ensure!(author != self.gateway.account().address(), "Peer cannot propose a batch that is authored by myself");
1714
1715 let previous_timestamp = match self.storage.get_certificate_for_round_with_author(previous_round, author) {
1717 Some(certificate) => certificate.timestamp(),
1719 None => return Ok(()),
1721 };
1722
1723 let elapsed = timestamp
1725 .checked_sub(previous_timestamp)
1726 .ok_or_else(|| anyhow!("Timestamp cannot be before the previous certificate at round {previous_round}"))?;
1727 match elapsed < MIN_BATCH_DELAY.as_secs() as i64 {
1729 true => bail!("Timestamp is too soon after the previous certificate at round {previous_round}"),
1730 false => Ok(()),
1731 }
1732 }
1733
1734 fn check_own_proposal_timestamp(
1742 &self,
1743 previous_round: u64,
1744 previous_timestamp: i64,
1745 timestamp: i64,
1746 ) -> Result<bool> {
1747 let elapsed = timestamp
1749 .checked_sub(previous_timestamp)
1750 .ok_or_else(|| anyhow!("Timestamp cannot be before the previous certificate at round {previous_round}"))?;
1751
1752 Ok(elapsed >= MIN_BATCH_DELAY.as_secs() as i64)
1753 }
1754
1755 async fn store_and_broadcast_certificate(&self, proposal: &Proposal<N>, committee: &Committee<N>) -> Result<()> {
1757 let (certificate, transmissions) = tokio::task::block_in_place(|| proposal.to_certificate(committee))?;
1759
1760 let transmissions = transmissions.into_iter().collect::<HashMap<_, _>>();
1763
1764 let round = certificate.round();
1766 let num_transmissions = certificate.transmission_ids().len();
1767
1768 let (storage, certificate_) = (self.storage.clone(), certificate.clone());
1770 spawn_blocking!(storage.insert_certificate(certificate_, transmissions, Default::default()))?;
1771 debug!("Stored a batch certificate for round {}", certificate.round());
1772 *self.proposed_batch.write() = ProposedBatchState::None;
1775
1776 if let Some(cb) = self.primary_callback.get() {
1778 cb.add_new_certificate(certificate.clone()).await.with_context(|| {
1780 format!("Failed to insert our newly certified batch for round {round} into the DAG")
1781 })?;
1782 }
1783 self.gateway.broadcast(Event::BatchCertified(certificate.into()));
1785
1786 info!("Our batch with {num_transmissions} transmissions for round {round} was certified!");
1788
1789 #[cfg(feature = "metrics")]
1791 if let Some(start) = self.batch_propose_start.lock().take() {
1792 metrics::histogram(metrics::bft::BATCH_CERTIFICATION_LATENCY, start.elapsed().as_secs_f64());
1793 }
1794
1795 self.round_increment_notify.notify_one();
1797
1798 Ok(())
1799 }
1800
1801 fn insert_missing_transmissions_into_workers(
1803 &self,
1804 peer_ip: SocketAddr,
1805 transmissions: impl Iterator<Item = (TransmissionID<N>, Transmission<N>)>,
1806 ) -> Result<()> {
1807 assign_to_workers(self.workers(), transmissions, |worker, transmission_id, transmission| {
1809 worker.process_transmission_from_peer(peer_ip, transmission_id, transmission);
1810 })
1811 }
1812
1813 fn reinsert_transmissions_into_workers(
1815 &self,
1816 transmissions: IndexMap<TransmissionID<N>, Transmission<N>>,
1817 ) -> Result<()> {
1818 assign_to_workers(self.workers(), transmissions.into_iter(), |worker, transmission_id, transmission| {
1820 worker.reinsert(transmission_id, transmission);
1821 })
1822 }
1823
1824 #[async_recursion::async_recursion]
1834 async fn sync_with_certificate_from_peer<const IS_SYNCING: bool>(
1835 &self,
1836 peer_ip: SocketAddr,
1837 certificate: BatchCertificate<N>,
1838 ) -> Result<()> {
1839 let batch_header = certificate.batch_header();
1841 let batch_round = batch_header.round();
1843
1844 if batch_round <= self.storage.gc_round() {
1846 return Ok(());
1847 }
1848 if self.storage.contains_certificate(certificate.id()) {
1850 return Ok(());
1851 }
1852
1853 if !IS_SYNCING && !self.is_synced() {
1855 bail!(
1856 "Failed to process certificate `{}` at round {batch_round} from '{peer_ip}' (node is syncing)",
1857 fmt_id(certificate.id())
1858 );
1859 }
1860
1861 let missing_transmissions =
1863 self.sync_with_batch_header_from_peer::<IS_SYNCING, false>(peer_ip, batch_header).await?;
1864
1865 if !self.storage.contains_certificate(certificate.id()) {
1867 let (storage, certificate_) = (self.storage.clone(), certificate.clone());
1869 spawn_blocking!(storage.insert_certificate(certificate_, missing_transmissions, Default::default()))?;
1870 debug!("Stored a batch certificate for round {batch_round} from '{peer_ip}'");
1871 if let Some(cb) = self.primary_callback.get() {
1873 cb.add_new_certificate(certificate).await.with_context(|| "Failed to update the DAG from sync")?;
1874 }
1875 self.round_increment_notify.notify_one();
1877 }
1878 Ok(())
1879 }
1880
1881 async fn sync_with_batch_header_from_peer<const IS_SYNCING: bool, const CHECK_PREVIOUS_CERTIFICATES: bool>(
1883 &self,
1884 peer_ip: SocketAddr,
1885 batch_header: &BatchHeader<N>,
1886 ) -> Result<HashMap<TransmissionID<N>, Transmission<N>>> {
1887 let batch_round = batch_header.round();
1889
1890 if batch_round <= self.storage.gc_round() {
1892 bail!("Round {batch_round} is too far in the past")
1893 }
1894
1895 if !IS_SYNCING && !self.is_synced() {
1897 bail!(
1898 "Failed to process batch header `{}` at round {batch_round} from '{peer_ip}' (node is syncing)",
1899 fmt_id(batch_header.batch_id())
1900 );
1901 }
1902
1903 let is_quorum_threshold_reached = {
1905 let authors = self.storage.get_certificate_authors_for_round(batch_round);
1906 let committee_lookback = self.ledger.get_committee_lookback_for_round(batch_round)?;
1907 committee_lookback.is_quorum_threshold_reached(&authors)
1908 };
1909
1910 let is_behind_schedule = is_quorum_threshold_reached && batch_round > self.current_round();
1915 let is_peer_far_in_future = batch_round > self.current_round() + self.storage.max_gc_rounds();
1917 if is_behind_schedule || is_peer_far_in_future {
1919 self.try_increment_to_the_next_round(batch_round)
1921 .await
1922 .with_context(|| "Failed to fast forward current round")?;
1923 }
1924
1925 let missing_transmissions_handle = self.fetch_missing_transmissions(peer_ip, batch_header);
1927
1928 let missing_previous_certificates_handle = self.fetch_missing_previous_certificates(peer_ip, batch_header);
1930
1931 let (missing_transmissions, missing_previous_certificates) = tokio::try_join!(
1933 missing_transmissions_handle,
1934 missing_previous_certificates_handle,
1935 ).with_context(|| format!("Failed to fetch missing transmissions and previous certificates for round {batch_round} from '{peer_ip}"))?;
1936
1937 for batch_certificate in missing_previous_certificates {
1941 if CHECK_PREVIOUS_CERTIFICATES {
1946 self.storage.check_incoming_certificate(&batch_certificate)?;
1947 }
1948 self.sync_with_certificate_from_peer::<IS_SYNCING>(peer_ip, batch_certificate).await?;
1950 }
1951 Ok(missing_transmissions)
1952 }
1953
1954 async fn fetch_missing_transmissions(
1957 &self,
1958 peer_ip: SocketAddr,
1959 batch_header: &BatchHeader<N>,
1960 ) -> Result<HashMap<TransmissionID<N>, Transmission<N>>> {
1961 if batch_header.round() <= self.storage.gc_round() {
1963 return Ok(Default::default());
1964 }
1965
1966 if self.storage.contains_batch(batch_header.batch_id()) {
1968 trace!("Batch for round {} from peer has already been processed", batch_header.round());
1969 return Ok(Default::default());
1970 }
1971
1972 let workers = self.workers.clone();
1974
1975 let mut fetch_transmissions = FuturesUnordered::new();
1977
1978 let num_workers = self.num_workers();
1980 for transmission_id in batch_header.transmission_ids() {
1982 if !self.storage.contains_transmission(*transmission_id) {
1984 let Ok(worker_id) = assign_to_worker(*transmission_id, num_workers) else {
1986 bail!("Unable to assign transmission ID '{transmission_id}' to a worker")
1987 };
1988 let Some(worker) = workers.get().expect("No workers set").get(worker_id as usize) else {
1990 bail!("Unable to find worker {worker_id}")
1991 };
1992 fetch_transmissions.push(worker.get_or_fetch_transmission(peer_ip, *transmission_id));
1994 }
1995 }
1996
1997 let mut transmissions = HashMap::with_capacity(fetch_transmissions.len());
1999 while let Some(result) = fetch_transmissions.next().await {
2001 let (transmission_id, transmission) = result?;
2003 transmissions.insert(transmission_id, transmission);
2005 }
2006 Ok(transmissions)
2008 }
2009
2010 async fn fetch_missing_previous_certificates(
2012 &self,
2013 peer_ip: SocketAddr,
2014 batch_header: &BatchHeader<N>,
2015 ) -> Result<HashSet<BatchCertificate<N>>> {
2016 let round = batch_header.round();
2018 if round == 1 || round <= self.storage.gc_round() + 1 {
2020 return Ok(Default::default());
2021 }
2022
2023 let missing_previous_certificates =
2025 self.fetch_missing_certificates(peer_ip, round, batch_header.previous_certificate_ids()).await?;
2026 if !missing_previous_certificates.is_empty() {
2027 debug!(
2028 "Fetched {} missing previous certificates for round {round} from '{peer_ip}'",
2029 missing_previous_certificates.len(),
2030 );
2031 }
2032 Ok(missing_previous_certificates)
2034 }
2035
2036 async fn fetch_missing_certificates(
2038 &self,
2039 peer_ip: SocketAddr,
2040 round: u64,
2041 certificate_ids: &IndexSet<Field<N>>,
2042 ) -> Result<HashSet<BatchCertificate<N>>> {
2043 let mut fetch_certificates = FuturesUnordered::new();
2045 let mut missing_certificates = HashSet::default();
2047 for certificate_id in certificate_ids {
2049 if self.ledger.contains_certificate(certificate_id)? {
2051 continue;
2052 }
2053 if self.storage.contains_certificate(*certificate_id) {
2055 continue;
2056 }
2057 if let Some(certificate) = self.storage.get_unprocessed_certificate(*certificate_id) {
2059 missing_certificates.insert(certificate);
2060 } else {
2061 trace!("Primary - Found a new certificate ID for round {round} from '{peer_ip}'");
2063 fetch_certificates.push(self.sync.send_certificate_request(peer_ip, *certificate_id));
2066 }
2067 }
2068
2069 match fetch_certificates.is_empty() {
2071 true => return Ok(missing_certificates),
2072 false => trace!(
2073 "Fetching {} missing certificates for round {round} from '{peer_ip}'...",
2074 fetch_certificates.len(),
2075 ),
2076 }
2077
2078 while let Some(result) = fetch_certificates.next().await {
2080 missing_certificates.insert(result?);
2082 }
2083 Ok(missing_certificates)
2085 }
2086}
2087
2088impl<N: Network> Primary<N> {
2089 fn spawn<T: Future<Output = ()> + Send + 'static>(&self, future: T) {
2091 self.handles.lock().push(tokio::spawn(future));
2092 }
2093
2094 pub async fn shut_down(&self) {
2096 info!("Shutting down the primary...");
2097 self.primary_callback.clear();
2099 self.sync.shut_down().await;
2101 self.workers().iter().for_each(|worker| worker.shut_down());
2103 self.handles.lock().drain(..).for_each(|handle| handle.abort());
2105 let proposal_cache = {
2107 let proposal = match std::mem::replace(&mut *self.proposed_batch.write(), ProposedBatchState::None) {
2111 ProposedBatchState::Certifying(p) => Some(*p),
2112 _ => None,
2113 };
2114 let signed_proposals = self.signed_proposals.read().clone();
2115 let latest_round = proposal
2116 .as_ref()
2117 .map(Proposal::round)
2118 .unwrap_or(self.latest_proposal_timestamp.read().await.map(|(round, _)| round).unwrap_or(0));
2119 let pending_certificates = self.storage.get_pending_certificates();
2120 ProposalCache::new(latest_round, proposal, signed_proposals, pending_certificates)
2121 };
2122 if let Err(err) = proposal_cache.store(&self.node_data_dir) {
2123 error!("{}", flatten_error(err.context("Failed to store the current proposal cache")));
2124 }
2125 self.gateway.shut_down().await;
2127 }
2128}
2129
2130#[cfg(test)]
2131mod tests {
2132 use super::{proposal_task::BatchPropose as _, *};
2133
2134 use snarkos_node_bft_ledger_service::MockLedgerService;
2135 use snarkos_node_bft_storage_service::BFTMemoryService;
2136 use snarkos_node_sync::{BlockSync, locators::test_helpers::sample_block_locators};
2137 use snarkvm::{
2138 ledger::{
2139 committee::{Committee, MIN_VALIDATOR_STAKE},
2140 test_helpers::sample_execution_transaction_with_fee,
2141 },
2142 prelude::{Address, Signature},
2143 };
2144
2145 use bytes::Bytes;
2146 use indexmap::IndexSet;
2147 use rand::RngExt;
2148
2149 type CurrentNetwork = snarkvm::prelude::MainnetV0;
2150
2151 fn sample_committee(rng: &mut TestRng) -> (Vec<(SocketAddr, Account<CurrentNetwork>)>, Committee<CurrentNetwork>) {
2152 const COMMITTEE_SIZE: usize = 4;
2154 let mut accounts = Vec::with_capacity(COMMITTEE_SIZE);
2155 let mut members = IndexMap::new();
2156
2157 for i in 0..COMMITTEE_SIZE {
2158 let socket_addr = format!("127.0.0.1:{}", 5000 + i).parse().unwrap();
2159 let account = Account::new(rng).unwrap();
2160
2161 members.insert(account.address(), (MIN_VALIDATOR_STAKE, true, rng.random_range(0..100)));
2162 accounts.push((socket_addr, account));
2163 }
2164
2165 (accounts, Committee::<CurrentNetwork>::new(1, members).unwrap())
2166 }
2167
2168 fn primary_with_committee(
2170 account_index: usize,
2171 accounts: &[(SocketAddr, Account<CurrentNetwork>)],
2172 committee: Committee<CurrentNetwork>,
2173 height: u32,
2174 ) -> Primary<CurrentNetwork> {
2175 let ledger = Arc::new(MockLedgerService::new_at_height(committee, height));
2176 let storage = Storage::new(ledger.clone(), Arc::new(BFTMemoryService::new()), 10).unwrap();
2177
2178 let account = accounts[account_index].1.clone();
2180 let block_sync = Arc::new(BlockSync::new(ledger.clone(), ConnectionMode::Gateway));
2181 let primary =
2182 Primary::new(account, storage, ledger, block_sync, None, &[], false, NodeDataDir::new_test(None), None)
2183 .unwrap();
2184
2185 let worker = Worker::new(
2187 0, Arc::new(primary.gateway.clone()),
2189 primary.storage.clone(),
2190 primary.ledger.clone(),
2191 primary.proposed_batch.clone(),
2192 )
2193 .unwrap();
2194 let _ = primary.workers.set(vec![worker]);
2195 for a in accounts.iter().skip(account_index) {
2196 primary.gateway.insert_connected_peer(a.0, a.0, a.1.address());
2197 }
2198
2199 primary
2200 }
2201
2202 fn primary_without_handlers(
2203 rng: &mut TestRng,
2204 ) -> (Primary<CurrentNetwork>, Vec<(SocketAddr, Account<CurrentNetwork>)>) {
2205 let (accounts, committee) = sample_committee(rng);
2206 let primary = primary_with_committee(
2207 0, &accounts,
2209 committee,
2210 CurrentNetwork::CONSENSUS_HEIGHT(ConsensusVersion::V1).unwrap(),
2211 );
2212
2213 (primary, accounts)
2214 }
2215
2216 fn sample_unconfirmed_solution(rng: &mut TestRng) -> (SolutionID<CurrentNetwork>, Data<Solution<CurrentNetwork>>) {
2218 let solution_id = rng.random::<u64>().into();
2220 let size = rng.random_range(1024..10 * 1024);
2222 let vec: Vec<u8> = (0..size).map(|_| rng.random::<u8>()).collect();
2224 let solution = Data::Buffer(Bytes::from(vec));
2225 (solution_id, solution)
2227 }
2228
2229 fn sample_unconfirmed_transaction(
2231 rng: &mut TestRng,
2232 ) -> (<CurrentNetwork as Network>::TransactionID, Data<Transaction<CurrentNetwork>>) {
2233 let transaction = sample_execution_transaction_with_fee(false, rng, 0);
2234 let id = transaction.id();
2235
2236 (id, Data::Object(transaction))
2237 }
2238
2239 fn create_test_proposal(
2241 author: &Account<CurrentNetwork>,
2242 committee: Committee<CurrentNetwork>,
2243 round: u64,
2244 previous_certificate_ids: IndexSet<Field<CurrentNetwork>>,
2245 timestamp: i64,
2246 num_transactions: u64,
2247 rng: &mut TestRng,
2248 ) -> Proposal<CurrentNetwork> {
2249 let mut transmission_ids = IndexSet::new();
2250 let mut transmissions = IndexMap::new();
2251
2252 let (solution_id, solution) = sample_unconfirmed_solution(rng);
2254 let solution_checksum = solution.to_checksum::<CurrentNetwork>().unwrap();
2255 let solution_transmission_id = (solution_id, solution_checksum).into();
2256 transmission_ids.insert(solution_transmission_id);
2257 transmissions.insert(solution_transmission_id, Transmission::Solution(solution));
2258
2259 for _ in 0..num_transactions {
2261 let (transaction_id, transaction) = sample_unconfirmed_transaction(rng);
2262 let transaction_checksum = transaction.to_checksum::<CurrentNetwork>().unwrap();
2263 let transaction_transmission_id = (&transaction_id, &transaction_checksum).into();
2264 transmission_ids.insert(transaction_transmission_id);
2265 transmissions.insert(transaction_transmission_id, Transmission::Transaction(transaction));
2266 }
2267
2268 let private_key = author.private_key();
2270 let batch_header = BatchHeader::new(
2272 private_key,
2273 round,
2274 timestamp,
2275 committee.id(),
2276 transmission_ids,
2277 previous_certificate_ids,
2278 rng,
2279 )
2280 .unwrap();
2281 Proposal::new(committee, batch_header, transmissions).unwrap()
2283 }
2284
2285 fn peer_signatures_for_proposal(
2288 primary: &Primary<CurrentNetwork>,
2289 accounts: &[(SocketAddr, Account<CurrentNetwork>)],
2290 rng: &mut TestRng,
2291 ) -> Vec<(SocketAddr, BatchSignature<CurrentNetwork>)> {
2292 let mut signatures = Vec::with_capacity(accounts.len() - 1);
2294 for (socket_addr, account) in accounts {
2295 if account.address() == primary.gateway.account().address() {
2296 continue;
2297 }
2298 let batch_id = primary.proposed_batch.read().as_proposal().unwrap().batch_id();
2299 let signature = account.sign(&[batch_id], rng).unwrap();
2300 signatures.push((*socket_addr, BatchSignature::new(batch_id, signature)));
2301 }
2302
2303 signatures
2304 }
2305
2306 fn peer_signatures_for_batch(
2308 primary_address: Address<CurrentNetwork>,
2309 accounts: &[(SocketAddr, Account<CurrentNetwork>)],
2310 batch_id: Field<CurrentNetwork>,
2311 rng: &mut TestRng,
2312 ) -> IndexSet<Signature<CurrentNetwork>> {
2313 let mut signatures = IndexSet::new();
2314 for (_, account) in accounts {
2315 if account.address() == primary_address {
2316 continue;
2317 }
2318 let signature = account.sign(&[batch_id], rng).unwrap();
2319 signatures.insert(signature);
2320 }
2321 signatures
2322 }
2323
2324 fn create_batch_certificate(
2326 primary_address: Address<CurrentNetwork>,
2327 accounts: &[(SocketAddr, Account<CurrentNetwork>)],
2328 round: u64,
2329 previous_certificate_ids: IndexSet<Field<CurrentNetwork>>,
2330 rng: &mut TestRng,
2331 ) -> (BatchCertificate<CurrentNetwork>, HashMap<TransmissionID<CurrentNetwork>, Transmission<CurrentNetwork>>) {
2332 let timestamp = now();
2333
2334 let author =
2335 accounts.iter().find(|&(_, acct)| acct.address() == primary_address).map(|(_, acct)| acct.clone()).unwrap();
2336 let private_key = author.private_key();
2337
2338 let committee_id = Field::rand(rng);
2339 let (solution_id, solution) = sample_unconfirmed_solution(rng);
2340 let (transaction_id, transaction) = sample_unconfirmed_transaction(rng);
2341 let solution_checksum = solution.to_checksum::<CurrentNetwork>().unwrap();
2342 let transaction_checksum = transaction.to_checksum::<CurrentNetwork>().unwrap();
2343
2344 let solution_transmission_id = (solution_id, solution_checksum).into();
2345 let transaction_transmission_id = (&transaction_id, &transaction_checksum).into();
2346
2347 let transmission_ids = [solution_transmission_id, transaction_transmission_id].into();
2348 let transmissions = [
2349 (solution_transmission_id, Transmission::Solution(solution)),
2350 (transaction_transmission_id, Transmission::Transaction(transaction)),
2351 ]
2352 .into();
2353
2354 let batch_header = BatchHeader::new(
2355 private_key,
2356 round,
2357 timestamp,
2358 committee_id,
2359 transmission_ids,
2360 previous_certificate_ids,
2361 rng,
2362 )
2363 .unwrap();
2364 let signatures = peer_signatures_for_batch(primary_address, accounts, batch_header.batch_id(), rng);
2365 let certificate = BatchCertificate::<CurrentNetwork>::from(batch_header, signatures).unwrap();
2366 (certificate, transmissions)
2367 }
2368
2369 fn store_certificate_chain(
2371 primary: &Primary<CurrentNetwork>,
2372 accounts: &[(SocketAddr, Account<CurrentNetwork>)],
2373 round: u64,
2374 rng: &mut TestRng,
2375 ) -> IndexSet<Field<CurrentNetwork>> {
2376 let mut previous_certificates = IndexSet::<Field<CurrentNetwork>>::new();
2377 let mut next_certificates = IndexSet::<Field<CurrentNetwork>>::new();
2378 for cur_round in 1..round {
2379 for (_, account) in accounts.iter() {
2380 let (certificate, transmissions) = create_batch_certificate(
2381 account.address(),
2382 accounts,
2383 cur_round,
2384 previous_certificates.clone(),
2385 rng,
2386 );
2387 next_certificates.insert(certificate.id());
2388 assert!(primary.storage.insert_certificate(certificate, transmissions, Default::default()).is_ok());
2389 }
2390
2391 assert!(primary.storage.increment_to_next_round(cur_round).is_ok());
2392 previous_certificates = next_certificates;
2393 next_certificates = IndexSet::<Field<CurrentNetwork>>::new();
2394 }
2395
2396 previous_certificates
2397 }
2398
2399 fn map_account_addresses(primary: &Primary<CurrentNetwork>, accounts: &[(SocketAddr, Account<CurrentNetwork>)]) {
2402 for (addr, acct) in accounts.iter().skip(1) {
2404 primary.gateway.resolver().write().insert_peer(*addr, *addr, Some(acct.address()));
2405 }
2406 }
2407
2408 #[test_log::test(tokio::test)]
2409 async fn test_propose_batch() {
2410 let mut rng = TestRng::default();
2411 let (primary, _) = primary_without_handlers(&mut rng);
2412
2413 assert!(primary.proposed_batch.read().is_none());
2415
2416 let (solution_id, solution) = sample_unconfirmed_solution(&mut rng);
2418 let (transaction_id, transaction) = sample_unconfirmed_transaction(&mut rng);
2419
2420 primary.workers()[0].process_unconfirmed_solution(solution_id, solution).await.unwrap();
2422 primary.workers()[0].process_unconfirmed_transaction(transaction_id, transaction).await.unwrap();
2423
2424 assert!(primary.propose_batch().await.is_ok());
2426 assert!(primary.proposed_batch.read().is_proposed());
2427 }
2428
2429 #[test_log::test(tokio::test)]
2430 async fn test_propose_batch_with_no_transmissions() {
2431 let mut rng = TestRng::default();
2432 let (primary, _) = primary_without_handlers(&mut rng);
2433
2434 assert!(primary.proposed_batch.read().is_none());
2436
2437 assert!(primary.propose_batch().await.is_ok());
2439 assert!(primary.proposed_batch.read().is_proposed());
2440 }
2441
2442 #[test_log::test(tokio::test)]
2443 async fn test_propose_batch_in_round() {
2444 let round = 3;
2445 let mut rng = TestRng::default();
2446 let (primary, accounts) = primary_without_handlers(&mut rng);
2447
2448 store_certificate_chain(&primary, &accounts, round, &mut rng);
2450
2451 tokio::time::sleep(MIN_BATCH_DELAY).await;
2453
2454 let (solution_id, solution) = sample_unconfirmed_solution(&mut rng);
2456 let (transaction_id, transaction) = sample_unconfirmed_transaction(&mut rng);
2457
2458 primary.workers()[0].process_unconfirmed_solution(solution_id, solution).await.unwrap();
2460 primary.workers()[0].process_unconfirmed_transaction(transaction_id, transaction).await.unwrap();
2461
2462 assert!(primary.propose_batch().await.is_ok());
2464 assert!(primary.proposed_batch.read().is_proposed());
2465 }
2466
2467 #[test_log::test(tokio::test)]
2468 async fn test_propose_batch_skip_transmissions_from_previous_certificates() {
2469 let round = 3;
2470 let prev_round = round - 1;
2471 let mut rng = TestRng::default();
2472 let (primary, accounts) = primary_without_handlers(&mut rng);
2473 let peer_account = &accounts[1];
2474 let peer_ip = peer_account.0;
2475
2476 store_certificate_chain(&primary, &accounts, round, &mut rng);
2478
2479 let previous_certificate_ids: IndexSet<_> = primary.storage.get_certificate_ids_for_round(prev_round);
2481
2482 let mut num_transmissions_in_previous_round = 0;
2484
2485 let (solution_commitment, solution) = sample_unconfirmed_solution(&mut rng);
2487 let (transaction_id, transaction) = sample_unconfirmed_transaction(&mut rng);
2488 let solution_checksum = solution.to_checksum::<CurrentNetwork>().unwrap();
2489 let transaction_checksum = transaction.to_checksum::<CurrentNetwork>().unwrap();
2490
2491 primary.workers()[0].process_unconfirmed_solution(solution_commitment, solution).await.unwrap();
2493 primary.workers()[0].process_unconfirmed_transaction(transaction_id, transaction).await.unwrap();
2494
2495 assert_eq!(primary.workers()[0].num_transmissions(), 2);
2497
2498 for (_, account) in accounts.iter() {
2500 let (certificate, transmissions) = create_batch_certificate(
2501 account.address(),
2502 &accounts,
2503 round,
2504 previous_certificate_ids.clone(),
2505 &mut rng,
2506 );
2507
2508 for (transmission_id, transmission) in transmissions.iter() {
2510 primary.workers()[0].process_transmission_from_peer(peer_ip, *transmission_id, transmission.clone());
2511 }
2512
2513 num_transmissions_in_previous_round += transmissions.len();
2515 primary.storage.insert_certificate(certificate, transmissions, Default::default()).unwrap();
2516 }
2517
2518 tokio::time::sleep(MIN_BATCH_DELAY).await;
2520
2521 assert!(primary.storage.increment_to_next_round(round).is_ok());
2523
2524 assert_eq!(primary.workers()[0].num_transmissions(), num_transmissions_in_previous_round + 2);
2526
2527 assert!(primary.propose_batch().await.is_ok());
2529
2530 let proposed_transmissions = primary.proposed_batch.read().as_proposal().unwrap().transmissions().clone();
2532 assert_eq!(proposed_transmissions.len(), 2);
2533 assert!(proposed_transmissions.contains_key(&TransmissionID::Solution(solution_commitment, solution_checksum)));
2534 assert!(
2535 proposed_transmissions.contains_key(&TransmissionID::Transaction(transaction_id, transaction_checksum))
2536 );
2537 }
2538
2539 #[test_log::test(tokio::test)]
2540 async fn test_propose_batch_over_spend_limit() {
2541 let mut rng = TestRng::default();
2542
2543 let (accounts, committee) = sample_committee(&mut rng);
2545 let primary = primary_with_committee(
2546 0,
2547 &accounts,
2548 committee.clone(),
2549 CurrentNetwork::CONSENSUS_HEIGHT(ConsensusVersion::V4).unwrap(),
2550 );
2551
2552 assert!(primary.proposed_batch.read().is_none());
2554 primary.workers().iter().for_each(|worker| assert!(worker.transmissions().is_empty()));
2556
2557 let (solution_id, solution) = sample_unconfirmed_solution(&mut rng);
2559 primary.workers()[0].process_unconfirmed_solution(solution_id, solution).await.unwrap();
2560
2561 for _i in 0..5 {
2562 let (transaction_id, transaction) = sample_unconfirmed_transaction(&mut rng);
2563 primary.workers()[0].process_unconfirmed_transaction(transaction_id, transaction).await.unwrap();
2565 }
2566
2567 assert!(primary.propose_batch().await.is_ok());
2569 assert_eq!(primary.proposed_batch.read().as_proposal().unwrap().transmissions().len(), 3);
2571 assert_eq!(primary.workers().iter().map(|worker| worker.transmissions().len()).sum::<usize>(), 3);
2573 }
2574
2575 #[test_log::test(tokio::test)]
2576 async fn test_batch_propose_from_peer() {
2577 let mut rng = TestRng::default();
2578 let (primary, accounts) = primary_without_handlers(&mut rng);
2579
2580 let round = 1;
2582 let peer_account = &accounts[1];
2583 let peer_ip = peer_account.0;
2584 let timestamp = now() + MIN_BATCH_DELAY.as_secs() as i64;
2585 let proposal = create_test_proposal(
2586 &peer_account.1,
2587 primary.ledger.current_committee().unwrap(),
2588 round,
2589 Default::default(),
2590 timestamp,
2591 1,
2592 &mut rng,
2593 );
2594
2595 for (transmission_id, transmission) in proposal.transmissions() {
2597 primary.workers()[0].process_transmission_from_peer(peer_ip, *transmission_id, transmission.clone())
2598 }
2599
2600 primary.gateway.resolver().write().insert_peer(peer_ip, peer_ip, Some(peer_account.1.address()));
2602
2603 primary.sync.testing_only_update_peer_locators_testing_only(peer_ip, sample_block_locators(20)).unwrap();
2606 primary.sync.testing_only_set_sync_height_testing_only(20);
2607
2608 assert!(
2610 primary.process_batch_propose_from_peer(peer_ip, (*proposal.batch_header()).clone().into()).await.is_ok()
2611 );
2612 }
2613
2614 #[test_log::test(tokio::test)]
2615 async fn test_batch_propose_from_peer_when_not_synced() {
2616 let mut rng = TestRng::default();
2617 let (primary, accounts) = primary_without_handlers(&mut rng);
2618
2619 let round = 1;
2621 let peer_account = &accounts[1];
2622 let peer_ip = peer_account.0;
2623 let timestamp = now() + MIN_BATCH_DELAY.as_secs() as i64;
2624 let proposal = create_test_proposal(
2625 &peer_account.1,
2626 primary.ledger.current_committee().unwrap(),
2627 round,
2628 Default::default(),
2629 timestamp,
2630 1,
2631 &mut rng,
2632 );
2633
2634 for (transmission_id, transmission) in proposal.transmissions() {
2636 primary.workers()[0].process_transmission_from_peer(peer_ip, *transmission_id, transmission.clone())
2637 }
2638
2639 primary.gateway.resolver().write().insert_peer(peer_ip, peer_ip, Some(peer_account.1.address()));
2641
2642 primary.sync.testing_only_update_peer_locators_testing_only(peer_ip, sample_block_locators(20)).unwrap();
2644
2645 assert!(
2647 primary.process_batch_propose_from_peer(peer_ip, (*proposal.batch_header()).clone().into()).await.is_err()
2648 );
2649 }
2650
2651 #[test_log::test(tokio::test)]
2652 async fn test_batch_propose_from_peer_in_round() {
2653 let round = 2;
2654 let mut rng = TestRng::default();
2655 let (primary, accounts) = primary_without_handlers(&mut rng);
2656
2657 let previous_certificates = store_certificate_chain(&primary, &accounts, round, &mut rng);
2659
2660 let peer_account = &accounts[1];
2662 let peer_ip = peer_account.0;
2663 let timestamp = now() + MIN_BATCH_DELAY.as_secs() as i64;
2664 let proposal = create_test_proposal(
2665 &peer_account.1,
2666 primary.ledger.current_committee().unwrap(),
2667 round,
2668 previous_certificates,
2669 timestamp,
2670 1,
2671 &mut rng,
2672 );
2673
2674 for (transmission_id, transmission) in proposal.transmissions() {
2676 primary.workers()[0].process_transmission_from_peer(peer_ip, *transmission_id, transmission.clone())
2677 }
2678
2679 primary.gateway.resolver().write().insert_peer(peer_ip, peer_ip, Some(peer_account.1.address()));
2681
2682 primary.sync.testing_only_update_peer_locators_testing_only(peer_ip, sample_block_locators(20)).unwrap();
2685 primary.sync.testing_only_set_sync_height_testing_only(20);
2686
2687 primary.process_batch_propose_from_peer(peer_ip, (*proposal.batch_header()).clone().into()).await.unwrap();
2689 }
2690
2691 #[test_log::test(tokio::test)]
2692 async fn test_batch_propose_from_peer_wrong_round() {
2693 let mut rng = TestRng::default();
2694 let (primary, accounts) = primary_without_handlers(&mut rng);
2695
2696 let round = 1;
2698 let peer_account = &accounts[1];
2699 let peer_ip = peer_account.0;
2700 let timestamp = now() + MIN_BATCH_DELAY.as_secs() as i64;
2701 let proposal = create_test_proposal(
2702 &peer_account.1,
2703 primary.ledger.current_committee().unwrap(),
2704 round,
2705 Default::default(),
2706 timestamp,
2707 1,
2708 &mut rng,
2709 );
2710
2711 for (transmission_id, transmission) in proposal.transmissions() {
2713 primary.workers()[0].process_transmission_from_peer(peer_ip, *transmission_id, transmission.clone())
2714 }
2715
2716 primary.gateway.resolver().write().insert_peer(peer_ip, peer_ip, Some(peer_account.1.address()));
2718 primary.sync.testing_only_update_peer_locators_testing_only(peer_ip, sample_block_locators(20)).unwrap();
2720 primary.sync.testing_only_set_sync_height_testing_only(20);
2721
2722 assert!(
2724 primary
2725 .process_batch_propose_from_peer(peer_ip, BatchPropose {
2726 round: round + 1,
2727 batch_header: Data::Object(proposal.batch_header().clone())
2728 })
2729 .await
2730 .is_err()
2731 );
2732 }
2733
2734 #[test_log::test(tokio::test)]
2735 async fn test_batch_propose_from_peer_in_round_wrong_round() {
2736 let round = 4;
2737 let mut rng = TestRng::default();
2738 let (primary, accounts) = primary_without_handlers(&mut rng);
2739
2740 let previous_certificates = store_certificate_chain(&primary, &accounts, round, &mut rng);
2742
2743 let peer_account = &accounts[1];
2745 let peer_ip = peer_account.0;
2746 let timestamp = now() + MIN_BATCH_DELAY.as_secs() as i64;
2747 let proposal = create_test_proposal(
2748 &peer_account.1,
2749 primary.ledger.current_committee().unwrap(),
2750 round,
2751 previous_certificates,
2752 timestamp,
2753 1,
2754 &mut rng,
2755 );
2756
2757 for (transmission_id, transmission) in proposal.transmissions() {
2759 primary.workers()[0].process_transmission_from_peer(peer_ip, *transmission_id, transmission.clone())
2760 }
2761
2762 primary.gateway.resolver().write().insert_peer(peer_ip, peer_ip, Some(peer_account.1.address()));
2764 primary.sync.testing_only_update_peer_locators_testing_only(peer_ip, sample_block_locators(0)).unwrap();
2766 primary.sync.testing_only_set_sync_height_testing_only(0);
2767
2768 assert!(
2770 primary
2771 .process_batch_propose_from_peer(peer_ip, BatchPropose {
2772 round: round + 1,
2773 batch_header: Data::Object(proposal.batch_header().clone())
2774 })
2775 .await
2776 .is_err()
2777 );
2778 }
2779
2780 #[test_log::test(tokio::test)]
2782 async fn test_batch_propose_from_peer_with_past_timestamp() {
2783 let round = 2;
2784 let mut rng = TestRng::default();
2785 let (primary, accounts) = primary_without_handlers(&mut rng);
2786
2787 let previous_certificates = store_certificate_chain(&primary, &accounts, round, &mut rng);
2789
2790 let peer_account = &accounts[1];
2792 let peer_ip = peer_account.0;
2793
2794 let last_timestamp = primary
2798 .storage
2799 .get_certificate_for_round_with_author(round - 1, peer_account.1.address())
2800 .expect("No previous proposal exists")
2801 .timestamp();
2802 let invalid_timestamp = last_timestamp + (MIN_BATCH_DELAY.as_secs() as i64) - 1;
2803
2804 let proposal = create_test_proposal(
2805 &peer_account.1,
2806 primary.ledger.current_committee().unwrap(),
2807 round,
2808 previous_certificates,
2809 invalid_timestamp,
2810 1,
2811 &mut rng,
2812 );
2813
2814 for (transmission_id, transmission) in proposal.transmissions() {
2816 primary.workers()[0].process_transmission_from_peer(peer_ip, *transmission_id, transmission.clone())
2817 }
2818
2819 primary.gateway.resolver().write().insert_peer(peer_ip, peer_ip, Some(peer_account.1.address()));
2821 primary.sync.testing_only_update_peer_locators_testing_only(peer_ip, sample_block_locators(0)).unwrap();
2823 primary.sync.testing_only_set_sync_height_testing_only(0);
2824
2825 assert!(
2827 primary.process_batch_propose_from_peer(peer_ip, (*proposal.batch_header()).clone().into()).await.is_err()
2828 );
2829 }
2830
2831 #[test_log::test(tokio::test)]
2832 async fn test_propose_batch_with_storage_round_behind_proposal_lock() {
2833 let round = 3;
2834 let mut rng = TestRng::default();
2835 let (primary, _) = primary_without_handlers(&mut rng);
2836
2837 assert!(primary.proposed_batch.read().is_none());
2839
2840 let (solution_id, solution) = sample_unconfirmed_solution(&mut rng);
2842 let (transaction_id, transaction) = sample_unconfirmed_transaction(&mut rng);
2843
2844 primary.workers()[0].process_unconfirmed_solution(solution_id, solution).await.unwrap();
2846 primary.workers()[0].process_unconfirmed_transaction(transaction_id, transaction).await.unwrap();
2847
2848 let (old_proposal_round, old_proposal_timestamp) = primary
2850 .latest_proposal_timestamp
2851 .read()
2852 .await
2853 .map(|(round, timestamp)| (round, timestamp))
2854 .unwrap_or((0, 0));
2855 *primary.latest_proposal_timestamp.write().await =
2856 Some((round + 1, old_proposal_timestamp + MIN_BATCH_DELAY.as_secs() as i64));
2857
2858 assert!(primary.propose_batch().await.is_ok());
2860 assert!(primary.proposed_batch.read().is_none());
2861
2862 *primary.latest_proposal_timestamp.write().await = Some((old_proposal_round, old_proposal_timestamp));
2864
2865 assert!(primary.propose_batch().await.is_ok());
2867 assert!(primary.proposed_batch.read().is_proposed());
2868 }
2869
2870 #[test_log::test(tokio::test)]
2871 async fn test_propose_batch_with_storage_round_behind_proposal() {
2872 let round = 5;
2873 let mut rng = TestRng::default();
2874 let (primary, accounts) = primary_without_handlers(&mut rng);
2875
2876 let previous_certificates = store_certificate_chain(&primary, &accounts, round, &mut rng);
2878
2879 let timestamp = now();
2881 let proposal = create_test_proposal(
2882 primary.gateway.account(),
2883 primary.ledger.current_committee().unwrap(),
2884 round + 1,
2885 previous_certificates,
2886 timestamp,
2887 1,
2888 &mut rng,
2889 );
2890
2891 *primary.proposed_batch.write() = ProposedBatchState::Certifying(Box::new(proposal));
2893
2894 assert!(primary.propose_batch().await.is_ok());
2896 assert!(primary.proposed_batch.read().is_proposed());
2897 assert!(primary.proposed_batch.read().as_proposal().unwrap().round() > primary.current_round());
2898 }
2899
2900 #[test_log::test(tokio::test(flavor = "multi_thread"))]
2901 async fn test_batch_signature_from_peer() {
2902 let mut rng = TestRng::default();
2903 let (primary, accounts) = primary_without_handlers(&mut rng);
2904 map_account_addresses(&primary, &accounts);
2905
2906 let round = 1;
2908 let timestamp = now() + MIN_BATCH_DELAY.as_secs() as i64;
2909 let proposal = create_test_proposal(
2910 primary.gateway.account(),
2911 primary.ledger.current_committee().unwrap(),
2912 round,
2913 Default::default(),
2914 timestamp,
2915 1,
2916 &mut rng,
2917 );
2918
2919 *primary.proposed_batch.write() = ProposedBatchState::Certifying(Box::new(proposal));
2921
2922 let signatures = peer_signatures_for_proposal(&primary, &accounts, &mut rng);
2924
2925 for (socket_addr, signature) in signatures {
2927 primary.process_batch_signature_from_peer(socket_addr, signature).await.unwrap();
2928 }
2929
2930 assert!(primary.storage.contains_certificate_in_round_from(round, primary.gateway.account().address()));
2932 primary.try_increment_to_the_next_round(round + 1).await.unwrap();
2934 assert_eq!(primary.current_round(), round + 1);
2936 }
2937
2938 #[test_log::test(tokio::test(flavor = "multi_thread"))]
2939 async fn test_batch_signature_from_peer_in_round() {
2940 let round = 5;
2941 let mut rng = TestRng::default();
2942 let (primary, accounts) = primary_without_handlers(&mut rng);
2943 map_account_addresses(&primary, &accounts);
2944
2945 let previous_certificates = store_certificate_chain(&primary, &accounts, round, &mut rng);
2947
2948 let timestamp = now();
2950 let proposal = create_test_proposal(
2951 primary.gateway.account(),
2952 primary.ledger.current_committee().unwrap(),
2953 round,
2954 previous_certificates,
2955 timestamp,
2956 1,
2957 &mut rng,
2958 );
2959
2960 *primary.proposed_batch.write() = ProposedBatchState::Certifying(Box::new(proposal));
2962
2963 let signatures = peer_signatures_for_proposal(&primary, &accounts, &mut rng);
2965
2966 for (socket_addr, signature) in signatures {
2968 primary.process_batch_signature_from_peer(socket_addr, signature).await.unwrap();
2969 }
2970
2971 assert!(primary.storage.contains_certificate_in_round_from(round, primary.gateway.account().address()));
2973 primary.try_increment_to_the_next_round(round + 1).await.unwrap();
2975 assert_eq!(primary.current_round(), round + 1);
2977 }
2978
2979 #[test_log::test(tokio::test)]
2980 async fn test_batch_signature_from_peer_no_quorum() {
2981 let mut rng = TestRng::default();
2982 let (primary, accounts) = primary_without_handlers(&mut rng);
2983 map_account_addresses(&primary, &accounts);
2984
2985 let round = 1;
2987 let timestamp = now() + MIN_BATCH_DELAY.as_secs() as i64;
2988 let proposal = create_test_proposal(
2989 primary.gateway.account(),
2990 primary.ledger.current_committee().unwrap(),
2991 round,
2992 Default::default(),
2993 timestamp,
2994 1,
2995 &mut rng,
2996 );
2997
2998 *primary.proposed_batch.write() = ProposedBatchState::Certifying(Box::new(proposal));
3000
3001 let signatures = peer_signatures_for_proposal(&primary, &accounts, &mut rng);
3003
3004 let (socket_addr, signature) = signatures.first().unwrap();
3006 primary.process_batch_signature_from_peer(*socket_addr, *signature).await.unwrap();
3007
3008 assert!(!primary.storage.contains_certificate_in_round_from(round, primary.gateway.account().address()));
3010 assert_eq!(primary.current_round(), round);
3012 }
3013
3014 #[test_log::test(tokio::test)]
3015 async fn test_batch_signature_from_peer_in_round_no_quorum() {
3016 let round = 7;
3017 let mut rng = TestRng::default();
3018 let (primary, accounts) = primary_without_handlers(&mut rng);
3019 map_account_addresses(&primary, &accounts);
3020
3021 let previous_certificates = store_certificate_chain(&primary, &accounts, round, &mut rng);
3023
3024 let timestamp = now() + MIN_BATCH_DELAY.as_secs() as i64;
3026 let proposal = create_test_proposal(
3027 primary.gateway.account(),
3028 primary.ledger.current_committee().unwrap(),
3029 round,
3030 previous_certificates,
3031 timestamp,
3032 1,
3033 &mut rng,
3034 );
3035
3036 *primary.proposed_batch.write() = ProposedBatchState::Certifying(Box::new(proposal));
3038
3039 let signatures = peer_signatures_for_proposal(&primary, &accounts, &mut rng);
3041
3042 let (socket_addr, signature) = signatures.first().unwrap();
3044 primary.process_batch_signature_from_peer(*socket_addr, *signature).await.unwrap();
3045
3046 assert!(!primary.storage.contains_certificate_in_round_from(round, primary.gateway.account().address()));
3048 assert_eq!(primary.current_round(), round);
3050 }
3051
3052 #[test_log::test(tokio::test)]
3057 async fn test_batch_signature_from_peer_batch_being_certified() {
3058 let mut rng = TestRng::default();
3059 let (primary, accounts) = primary_without_handlers(&mut rng);
3060 map_account_addresses(&primary, &accounts);
3061
3062 let round = 1;
3064 let timestamp = now() + MIN_BATCH_DELAY.as_secs() as i64;
3065 let proposal = create_test_proposal(
3066 primary.gateway.account(),
3067 primary.ledger.current_committee().unwrap(),
3068 round,
3069 Default::default(),
3070 timestamp,
3071 1,
3072 &mut rng,
3073 );
3074 let batch_id = proposal.batch_id();
3075
3076 *primary.proposed_batch.write() = ProposedBatchState::Certified(batch_id);
3078
3079 let (socket_addr, account) =
3081 accounts.iter().find(|(_, a)| a.address() != primary.gateway.account().address()).unwrap();
3082 let signature = account.sign(&[batch_id], &mut rng).unwrap();
3083 let batch_signature = BatchSignature::new(batch_id, signature);
3084
3085 assert!(primary.process_batch_signature_from_peer(*socket_addr, batch_signature).await.is_ok());
3087 assert!(matches!(&*primary.proposed_batch.read(), ProposedBatchState::Certified(id) if *id == batch_id));
3089 }
3090
3091 #[test_log::test(tokio::test)]
3094 async fn test_batch_signature_from_peer_unknown_id_while_certifying() {
3095 let mut rng = TestRng::default();
3096 let (primary, accounts) = primary_without_handlers(&mut rng);
3097 map_account_addresses(&primary, &accounts);
3098
3099 let round = 1;
3101 let timestamp = now() + MIN_BATCH_DELAY.as_secs() as i64;
3102 let proposal_a = create_test_proposal(
3103 primary.gateway.account(),
3104 primary.ledger.current_committee().unwrap(),
3105 round,
3106 Default::default(),
3107 timestamp,
3108 1,
3109 &mut rng,
3110 );
3111 let proposal_b = create_test_proposal(
3112 primary.gateway.account(),
3113 primary.ledger.current_committee().unwrap(),
3114 round,
3115 Default::default(),
3116 timestamp,
3117 1,
3118 &mut rng,
3119 );
3120 let batch_id_a = proposal_a.batch_id();
3121 let batch_id_b = proposal_b.batch_id();
3122 assert_ne!(batch_id_a, batch_id_b);
3123
3124 *primary.proposed_batch.write() = ProposedBatchState::Certified(batch_id_a);
3126
3127 let (socket_addr, account) =
3129 accounts.iter().find(|(_, a)| a.address() != primary.gateway.account().address()).unwrap();
3130 let signature = account.sign(&[batch_id_b], &mut rng).unwrap();
3131 let batch_signature = BatchSignature::new(batch_id_b, signature);
3132
3133 assert!(primary.process_batch_signature_from_peer(*socket_addr, batch_signature).await.is_err());
3135 }
3136
3137 #[test_log::test(tokio::test(flavor = "multi_thread"))]
3140 async fn test_batch_signature_from_peer_already_certified() {
3141 let mut rng = TestRng::default();
3142 let (primary, accounts) = primary_without_handlers(&mut rng);
3143 map_account_addresses(&primary, &accounts);
3144
3145 let round = 1;
3147 let timestamp = now() + MIN_BATCH_DELAY.as_secs() as i64;
3148 let old_proposal = create_test_proposal(
3149 primary.gateway.account(),
3150 primary.ledger.current_committee().unwrap(),
3151 round,
3152 Default::default(),
3153 timestamp,
3154 1,
3155 &mut rng,
3156 );
3157 let old_batch_id = old_proposal.batch_id();
3158 *primary.proposed_batch.write() = ProposedBatchState::Certifying(Box::new(old_proposal));
3159 let signatures = peer_signatures_for_proposal(&primary, &accounts, &mut rng);
3160 for (socket_addr, signature) in signatures {
3161 primary.process_batch_signature_from_peer(socket_addr, signature).await.unwrap();
3162 }
3163 assert!(primary.storage.contains_certificate_in_round_from(round, primary.gateway.account().address()));
3165
3166 let new_proposal = create_test_proposal(
3168 primary.gateway.account(),
3169 primary.ledger.current_committee().unwrap(),
3170 round,
3171 Default::default(),
3172 timestamp,
3173 1,
3174 &mut rng,
3175 );
3176 assert_ne!(new_proposal.batch_id(), old_batch_id);
3177 *primary.proposed_batch.write() = ProposedBatchState::Certifying(Box::new(new_proposal));
3178
3179 let (socket_addr, account) =
3181 accounts.iter().find(|(_, a)| a.address() != primary.gateway.account().address()).unwrap();
3182 let signature = account.sign(&[old_batch_id], &mut rng).unwrap();
3183 let batch_signature = BatchSignature::new(old_batch_id, signature);
3184
3185 assert!(primary.process_batch_signature_from_peer(*socket_addr, batch_signature).await.is_ok());
3187 }
3188
3189 #[test_log::test(tokio::test)]
3190 async fn test_insert_certificate_with_aborted_transmissions() {
3191 let round = 3;
3192 let prev_round = round - 1;
3193 let mut rng = TestRng::default();
3194 let (primary, accounts) = primary_without_handlers(&mut rng);
3195 let peer_account = &accounts[1];
3196 let peer_ip = peer_account.0;
3197
3198 store_certificate_chain(&primary, &accounts, round, &mut rng);
3200
3201 let previous_certificate_ids: IndexSet<_> = primary.storage.get_certificate_ids_for_round(prev_round);
3203
3204 let (solution_commitment, solution) = sample_unconfirmed_solution(&mut rng);
3206 let (transaction_id, transaction) = sample_unconfirmed_transaction(&mut rng);
3207
3208 primary.workers()[0].process_unconfirmed_solution(solution_commitment, solution).await.unwrap();
3210 primary.workers()[0].process_unconfirmed_transaction(transaction_id, transaction).await.unwrap();
3211
3212 assert_eq!(primary.workers()[0].num_transmissions(), 2);
3214
3215 let account = accounts[0].1.clone();
3217 let (certificate, transmissions) =
3218 create_batch_certificate(account.address(), &accounts, round, previous_certificate_ids.clone(), &mut rng);
3219 let certificate_id = certificate.id();
3220
3221 let mut aborted_transmissions = HashSet::new();
3223 let mut transmissions_without_aborted = HashMap::new();
3224 for (transmission_id, transmission) in transmissions.clone() {
3225 match rng.random::<bool>() || aborted_transmissions.is_empty() {
3226 true => {
3227 aborted_transmissions.insert(transmission_id);
3229 }
3230 false => {
3231 transmissions_without_aborted.insert(transmission_id, transmission);
3233 }
3234 };
3235 }
3236
3237 for (transmission_id, transmission) in transmissions_without_aborted.iter() {
3239 primary.workers()[0].process_transmission_from_peer(peer_ip, *transmission_id, transmission.clone());
3240 }
3241
3242 assert!(
3244 primary
3245 .storage
3246 .check_certificate(&certificate, transmissions_without_aborted.clone(), Default::default())
3247 .is_err()
3248 );
3249 assert!(
3250 primary
3251 .storage
3252 .insert_certificate(certificate.clone(), transmissions_without_aborted.clone(), Default::default())
3253 .is_err()
3254 );
3255
3256 primary
3258 .storage
3259 .insert_certificate(certificate, transmissions_without_aborted, aborted_transmissions.clone())
3260 .unwrap();
3261
3262 assert!(primary.storage.contains_certificate(certificate_id));
3264 for aborted_transmission_id in aborted_transmissions {
3266 assert!(primary.storage.contains_transmission(aborted_transmission_id));
3267 assert!(primary.storage.get_transmission(aborted_transmission_id).is_none());
3268 }
3269 }
3270
3271 #[test]
3277 fn test_add_signature_to_batch_none_state() {
3278 let mut rng = TestRng::default();
3279 let (primary, accounts) = primary_without_handlers(&mut rng);
3280
3281 let peer_ip = accounts[1].0;
3282 let batch_id = Field::rand(&mut rng);
3283 let signature = accounts[1].1.sign(&[batch_id], &mut rng).unwrap();
3284
3285 let (result, new_state) =
3286 primary.add_signature_to_batch(ProposedBatchState::None, peer_ip, batch_id, signature);
3287
3288 assert!(result.is_err());
3289 assert_eq!(new_state, ProposedBatchState::None);
3290 }
3291
3292 #[test]
3294 fn test_add_signature_to_batch_certified_matching_id() {
3295 let mut rng = TestRng::default();
3296 let (primary, accounts) = primary_without_handlers(&mut rng);
3297
3298 let peer_ip = accounts[1].0;
3299 let batch_id = Field::rand(&mut rng);
3300 let signature = accounts[1].1.sign(&[batch_id], &mut rng).unwrap();
3301
3302 let (result, new_state) =
3303 primary.add_signature_to_batch(ProposedBatchState::Certified(batch_id), peer_ip, batch_id, signature);
3304
3305 assert!(result.unwrap().is_none());
3306 assert_eq!(new_state, ProposedBatchState::Certified(batch_id));
3307 }
3308
3309 #[test]
3311 fn test_add_signature_to_batch_certified_different_id() {
3312 let mut rng = TestRng::default();
3313 let (primary, accounts) = primary_without_handlers(&mut rng);
3314
3315 let peer_ip = accounts[1].0;
3316 let certified_id = Field::rand(&mut rng);
3317 let other_id = Field::rand(&mut rng);
3318 let signature = accounts[1].1.sign(&[other_id], &mut rng).unwrap();
3319
3320 let (result, new_state) =
3321 primary.add_signature_to_batch(ProposedBatchState::Certified(certified_id), peer_ip, other_id, signature);
3322
3323 assert!(result.is_err());
3324 assert_eq!(new_state, ProposedBatchState::Certified(certified_id));
3325 }
3326
3327 #[tokio::test(flavor = "multi_thread")]
3330 async fn test_add_signature_to_batch_certifying_different_id_in_storage() {
3331 let round = 1;
3332 let mut rng = TestRng::default();
3333 let (primary, accounts) = primary_without_handlers(&mut rng);
3334 map_account_addresses(&primary, &accounts);
3335
3336 let proposal = create_test_proposal(
3338 primary.gateway.account(),
3339 primary.ledger.current_committee().unwrap(),
3340 round,
3341 Default::default(),
3342 now(),
3343 0,
3344 &mut rng,
3345 );
3346 let proposal_batch_id = proposal.batch_id();
3347
3348 let (certificate, transmissions) =
3350 create_batch_certificate(accounts[1].1.address(), &accounts, round, Default::default(), &mut rng);
3351 let stored_batch_id = certificate.batch_id();
3352 primary.storage.insert_certificate(certificate, transmissions, Default::default()).unwrap();
3353
3354 let peer_ip = accounts[1].0;
3355 let signature = accounts[1].1.sign(&[stored_batch_id], &mut rng).unwrap();
3356
3357 let (result, new_state) = primary.add_signature_to_batch(
3358 ProposedBatchState::Certifying(Box::new(proposal)),
3359 peer_ip,
3360 stored_batch_id,
3361 signature,
3362 );
3363
3364 assert!(result.unwrap().is_none());
3365 assert_eq!(new_state.as_proposal().unwrap().batch_id(), proposal_batch_id);
3367 }
3368
3369 #[test]
3372 fn test_add_signature_to_batch_certifying_different_id_unknown() {
3373 let mut rng = TestRng::default();
3374 let (primary, accounts) = primary_without_handlers(&mut rng);
3375
3376 let proposal = create_test_proposal(
3377 primary.gateway.account(),
3378 primary.ledger.current_committee().unwrap(),
3379 1,
3380 Default::default(),
3381 now(),
3382 0,
3383 &mut rng,
3384 );
3385 let proposal_batch_id = proposal.batch_id();
3386
3387 let peer_ip = accounts[1].0;
3388 let unknown_id = Field::rand(&mut rng);
3389 let signature = accounts[1].1.sign(&[unknown_id], &mut rng).unwrap();
3390
3391 let (result, new_state) = primary.add_signature_to_batch(
3392 ProposedBatchState::Certifying(Box::new(proposal)),
3393 peer_ip,
3394 unknown_id,
3395 signature,
3396 );
3397
3398 assert!(result.is_err());
3399 assert_eq!(new_state.as_proposal().unwrap().batch_id(), proposal_batch_id);
3400 }
3401
3402 #[test]
3404 fn test_add_signature_to_batch_certifying_matching_no_quorum() {
3405 let mut rng = TestRng::default();
3406 let (primary, accounts) = primary_without_handlers(&mut rng);
3407 map_account_addresses(&primary, &accounts);
3408
3409 let proposal = create_test_proposal(
3410 primary.gateway.account(),
3411 primary.ledger.current_committee().unwrap(),
3412 1,
3413 Default::default(),
3414 now(),
3415 0,
3416 &mut rng,
3417 );
3418 let batch_id = proposal.batch_id();
3419
3420 let peer_ip = accounts[1].0;
3422 let signature = accounts[1].1.sign(&[batch_id], &mut rng).unwrap();
3423
3424 let (result, new_state) = primary.add_signature_to_batch(
3425 ProposedBatchState::Certifying(Box::new(proposal)),
3426 peer_ip,
3427 batch_id,
3428 signature,
3429 );
3430
3431 assert!(result.unwrap().is_none());
3432 assert_eq!(new_state.as_proposal().unwrap().batch_id(), batch_id);
3433 }
3434
3435 #[test]
3438 fn test_add_signature_to_batch_certifying_matching_quorum_reached() {
3439 let mut rng = TestRng::default();
3440 let (primary, accounts) = primary_without_handlers(&mut rng);
3441 map_account_addresses(&primary, &accounts);
3442
3443 let proposal = create_test_proposal(
3444 primary.gateway.account(),
3445 primary.ledger.current_committee().unwrap(),
3446 1,
3447 Default::default(),
3448 now(),
3449 0,
3450 &mut rng,
3451 );
3452 let batch_id = proposal.batch_id();
3453
3454 let peers: Vec<_> =
3456 accounts.iter().filter(|(_, a)| a.address() != primary.gateway.account().address()).collect();
3457 let mut state = ProposedBatchState::Certifying(Box::new(proposal));
3458 let mut final_result = None;
3459
3460 for (peer_ip, peer_account) in &peers {
3461 let signature = peer_account.sign(&[batch_id], &mut rng).unwrap();
3462 let (result, new_state) = primary.add_signature_to_batch(state, *peer_ip, batch_id, signature);
3463 state = new_state;
3464 if result.as_ref().unwrap().is_some() {
3465 final_result = Some(result);
3466 break;
3467 }
3468 }
3469
3470 let proposal = final_result.expect("quorum should be reached").unwrap().unwrap();
3472 assert_eq!(proposal.batch_id(), batch_id);
3473 assert_eq!(state, ProposedBatchState::Certified(batch_id));
3474 }
3475}