1#![forbid(unsafe_code)]
17
18mod transactions_queue;
19use transactions_queue::TransactionsQueue;
20
21#[macro_use]
22extern crate tracing;
23
24#[cfg(feature = "metrics")]
25extern crate snarkos_node_metrics as metrics;
26
27use snarkos_account::Account;
28use snarkos_node_bft::{
29 BFT,
30 MAX_BATCH_DELAY,
31 Primary,
32 helpers::{
33 ConsensusReceiver,
34 ConsensusSender,
35 PrimaryReceiver,
36 PrimarySender,
37 Storage as NarwhalStorage,
38 fmt_id,
39 init_consensus_channels,
40 init_primary_channels,
41 },
42 spawn_blocking,
43};
44use snarkos_node_bft_ledger_service::LedgerService;
45use snarkos_node_bft_storage_service::BFTPersistentStorage;
46use snarkos_node_sync::{BlockSync, Ping};
47use snarkos_utilities::NodeDataDir;
48
49use snarkvm::{
50 ledger::{
51 CheckBlockError,
52 block::Transaction,
53 narwhal::{BatchHeader, Data, Subdag, Transmission, TransmissionID},
54 puzzle::{Solution, SolutionID},
55 },
56 prelude::*,
57 utilities::flatten_error,
58};
59
60use aleo_std::StorageMode;
61use anyhow::{Context, Result, bail};
62use cfg_if::cfg_if;
63use colored::Colorize;
64use indexmap::IndexMap;
65#[cfg(feature = "locktick")]
66use locktick::parking_lot::{Mutex, RwLock};
67use lru::LruCache;
68#[cfg(not(feature = "locktick"))]
69use parking_lot::{Mutex, RwLock};
70use std::{future::Future, net::SocketAddr, num::NonZeroUsize, sync::Arc};
71use tokio::{
72 sync::{Notify, oneshot},
73 task::JoinHandle,
74};
75
76#[cfg(feature = "metrics")]
77use std::collections::HashMap;
78
79const CAPACITY_FOR_DEPLOYMENTS: usize = 1 << 10;
82const CAPACITY_FOR_EXECUTIONS: usize = 1 << 10;
85const CAPACITY_FOR_SOLUTIONS: usize = 1 << 10;
88const MAX_DEPLOYMENTS_PER_INTERVAL: usize = 1;
91
92type BftRunArgs<N> = Arc<Mutex<Option<(ConsensusSender<N>, PrimaryReceiver<N>)>>>;
93
94#[derive(Clone)]
101pub struct Consensus<N: Network> {
102 ledger: Arc<dyn LedgerService<N>>,
104 bft: BFT<N>,
106 primary_sender: PrimarySender<N>,
108 solutions_queue: Arc<Mutex<LruCache<SolutionID<N>, Solution<N>>>>,
110 transactions_queue: Arc<RwLock<TransactionsQueue<N>>>,
112 seen_solutions: Arc<Mutex<LruCache<SolutionID<N>, ()>>>,
114 seen_transactions: Arc<Mutex<LruCache<N::TransactionID, ()>>>,
116 #[cfg(feature = "metrics")]
117 transmissions_tracker: Arc<Mutex<HashMap<TransmissionID<N>, i64>>>,
118 handles: Arc<Mutex<Vec<JoinHandle<()>>>>,
120 ping: Arc<Ping<N>>,
122 block_sync: Arc<BlockSync<N>>,
124 block_commit_notify: Arc<Notify>,
126 consensus_receiver: Arc<Mutex<Option<ConsensusReceiver<N>>>>,
128 bft_run_args: BftRunArgs<N>,
130}
131
132impl<N: Network> Consensus<N> {
133 #[allow(clippy::too_many_arguments)]
135 pub async fn new(
136 account: Account<N>,
137 ledger: Arc<dyn LedgerService<N>>,
138 block_sync: Arc<BlockSync<N>>,
139 ip: Option<SocketAddr>,
140 trusted_validators: &[SocketAddr],
141 trusted_peers_only: bool,
142 storage_mode: StorageMode,
143 node_data_dir: NodeDataDir,
144 ping: Arc<Ping<N>>,
145 dev: Option<u16>,
146 ) -> Result<Self> {
147 let (primary_sender, primary_receiver) = init_primary_channels::<N>();
149 let transmissions = Arc::new(BFTPersistentStorage::open(storage_mode.clone())?);
151 let storage = NarwhalStorage::new(ledger.clone(), transmissions, BatchHeader::<N>::MAX_GC_ROUNDS as u64)
153 .with_context(|| "Failed to initialize the BFT storage")?;
154 let bft = BFT::new(
156 account,
157 storage,
158 ledger.clone(),
159 block_sync.clone(),
160 ip,
161 trusted_validators,
162 trusted_peers_only,
163 node_data_dir,
164 dev,
165 )?;
166 let (consensus_sender, consensus_receiver) = init_consensus_channels();
168
169 let _self = Self {
171 ledger,
172 bft,
173 block_sync,
174 primary_sender,
175 solutions_queue: Arc::new(Mutex::new(LruCache::new(NonZeroUsize::new(CAPACITY_FOR_SOLUTIONS).unwrap()))),
176 transactions_queue: Default::default(),
177 seen_solutions: Arc::new(Mutex::new(LruCache::new(NonZeroUsize::new(1 << 16).unwrap()))),
178 seen_transactions: Arc::new(Mutex::new(LruCache::new(NonZeroUsize::new(1 << 16).unwrap()))),
179 #[cfg(feature = "metrics")]
180 transmissions_tracker: Default::default(),
181 handles: Default::default(),
182 ping: ping.clone(),
183 block_commit_notify: Arc::new(Notify::new()),
184 consensus_receiver: Arc::new(Mutex::new(Some(consensus_receiver))),
185 bft_run_args: Arc::new(Mutex::new(Some((consensus_sender, primary_receiver)))),
186 };
187
188 Ok(_self)
189 }
190
191 pub const fn bft(&self) -> &BFT<N> {
193 &self.bft
194 }
195
196 pub fn contains_transaction(&self, transaction_id: &N::TransactionID) -> bool {
197 self.transactions_queue.read().contains(transaction_id)
198 }
199}
200
201impl<N: Network> Consensus<N> {
202 pub fn num_unconfirmed_transmissions(&self) -> usize {
204 self.bft.num_unconfirmed_transmissions()
205 }
206
207 pub fn num_unconfirmed_ratifications(&self) -> usize {
209 self.bft.num_unconfirmed_ratifications()
210 }
211
212 pub fn num_unconfirmed_solutions(&self) -> usize {
214 self.bft.num_unconfirmed_solutions()
215 }
216
217 pub fn num_unconfirmed_transactions(&self) -> usize {
219 self.bft.num_unconfirmed_transactions()
220 }
221}
222
223impl<N: Network> Consensus<N> {
224 pub fn unconfirmed_transmission_ids(&self) -> impl '_ + Iterator<Item = TransmissionID<N>> {
226 self.worker_transmission_ids().chain(self.inbound_transmission_ids())
227 }
228
229 pub fn unconfirmed_transmissions(&self) -> impl '_ + Iterator<Item = (TransmissionID<N>, Transmission<N>)> {
231 self.worker_transmissions().chain(self.inbound_transmissions())
232 }
233
234 pub fn unconfirmed_solutions(&self) -> impl '_ + Iterator<Item = (SolutionID<N>, Data<Solution<N>>)> {
236 self.worker_solutions().chain(self.inbound_solutions())
237 }
238
239 pub fn unconfirmed_transactions(&self) -> impl '_ + Iterator<Item = (N::TransactionID, Data<Transaction<N>>)> {
241 self.worker_transactions().chain(self.inbound_transactions())
242 }
243}
244
245impl<N: Network> Consensus<N> {
246 pub fn worker_transmission_ids(&self) -> impl '_ + Iterator<Item = TransmissionID<N>> {
248 self.bft.worker_transmission_ids()
249 }
250
251 pub fn worker_transmissions(&self) -> impl '_ + Iterator<Item = (TransmissionID<N>, Transmission<N>)> {
253 self.bft.worker_transmissions()
254 }
255
256 pub fn worker_solutions(&self) -> impl '_ + Iterator<Item = (SolutionID<N>, Data<Solution<N>>)> {
258 self.bft.worker_solutions()
259 }
260
261 pub fn worker_transactions(&self) -> impl '_ + Iterator<Item = (N::TransactionID, Data<Transaction<N>>)> {
263 self.bft.worker_transactions()
264 }
265}
266
267impl<N: Network> Consensus<N> {
268 pub fn inbound_transmission_ids(&self) -> impl '_ + Iterator<Item = TransmissionID<N>> {
270 self.inbound_transmissions().map(|(id, _)| id)
271 }
272
273 pub fn inbound_transmissions(&self) -> impl '_ + Iterator<Item = (TransmissionID<N>, Transmission<N>)> {
275 self.inbound_transactions()
276 .map(|(id, tx)| {
277 (
278 TransmissionID::Transaction(id, tx.to_checksum::<N>().unwrap_or_default()),
279 Transmission::Transaction(tx),
280 )
281 })
282 .chain(self.inbound_solutions().map(|(id, solution)| {
283 (
284 TransmissionID::Solution(id, solution.to_checksum::<N>().unwrap_or_default()),
285 Transmission::Solution(solution),
286 )
287 }))
288 }
289
290 pub fn inbound_solutions(&self) -> impl '_ + Iterator<Item = (SolutionID<N>, Data<Solution<N>>)> {
292 self.solutions_queue.lock().clone().into_iter().map(|(id, solution)| (id, Data::Object(solution)))
294 }
295
296 pub fn inbound_transactions(&self) -> impl '_ + Iterator<Item = (N::TransactionID, Data<Transaction<N>>)> {
298 self.transactions_queue.read().transactions().map(|(id, tx)| (id, Data::Object(tx)))
300 }
301}
302
303impl<N: Network> Consensus<N> {
304 pub async fn add_unconfirmed_solution(&self, solution: Solution<N>) -> Result<()> {
307 let checksum = Data::<Solution<N>>::Buffer(solution.to_bytes_le()?.into()).to_checksum::<N>()?;
309 {
311 let solution_id = solution.id();
312
313 if self.seen_solutions.lock().put(solution_id, ()).is_some() {
315 return Ok(());
317 }
318 if self.ledger.contains_transmission(&TransmissionID::Solution(solution_id, checksum))? {
320 bail!("Solution '{}' exists in the ledger {}", fmt_id(solution_id), "(skipping)".dimmed());
321 }
322 #[cfg(feature = "metrics")]
323 {
324 metrics::increment_gauge(metrics::consensus::UNCONFIRMED_SOLUTIONS, 1f64);
325 let timestamp = snarkos_node_bft::helpers::now();
326 self.transmissions_tracker.lock().insert(TransmissionID::Solution(solution.id(), checksum), timestamp);
327 }
328 trace!("Received unconfirmed solution '{}' in the queue", fmt_id(solution_id));
330 if self.solutions_queue.lock().put(solution_id, solution).is_some() {
331 bail!("Solution '{}' exists in the memory pool", fmt_id(solution_id));
332 }
333 }
334
335 self.process_unconfirmed_solutions().await
337 }
338
339 async fn process_unconfirmed_solutions(&self) -> Result<()> {
342 let num_unconfirmed_solutions = self.num_unconfirmed_solutions();
344 let num_unconfirmed_transmissions = self.num_unconfirmed_transmissions();
345 if num_unconfirmed_solutions >= N::MAX_SOLUTIONS
346 || num_unconfirmed_transmissions >= Primary::<N>::MAX_TRANSMISSIONS_TOLERANCE
347 {
348 return Ok(());
349 }
350 let solutions = {
352 let capacity = N::MAX_SOLUTIONS.saturating_sub(num_unconfirmed_solutions);
354 let mut queue = self.solutions_queue.lock();
356 let num_solutions = queue.len().min(capacity);
358 (0..num_solutions).filter_map(|_| queue.pop_lru().map(|(_, solution)| solution)).collect::<Vec<_>>()
360 };
361 for solution in solutions.into_iter() {
363 let solution_id = solution.id();
364 trace!("Adding unconfirmed solution '{}' to the memory pool...", fmt_id(solution_id));
365 match self.primary_sender.send_unconfirmed_solution(solution_id, Data::Object(solution)).await {
367 Ok(true) => {}
368 Ok(false) => debug!(
369 "Unable to add unconfirmed solution '{}' to the memory pool. Already exists.",
370 fmt_id(solution_id)
371 ),
372 Err(err) => {
373 let err = err.context(format!(
374 "Unable to add unconfirmed solution '{}' to the memory pool",
375 fmt_id(solution_id)
376 ));
377
378 if self.bft.is_synced() && self.ledger.latest_block_height() % N::NUM_BLOCKS_PER_EPOCH > 10 {
380 warn!("{}", flatten_error(err));
381 } else {
382 trace!("{}", flatten_error(err));
383 }
384 }
385 }
386 }
387 Ok(())
388 }
389
390 pub async fn add_unconfirmed_transaction(&self, transaction: Transaction<N>) -> Result<()> {
393 let checksum = Data::<Transaction<N>>::Buffer(transaction.to_bytes_le()?.into()).to_checksum::<N>()?;
395 {
397 let transaction_id = transaction.id();
398
399 if transaction.is_fee() {
401 bail!("Transaction '{}' is a fee transaction {}", fmt_id(transaction_id), "(skipping)".dimmed());
402 }
403 if self.seen_transactions.lock().put(transaction_id, ()).is_some() {
405 return Ok(());
407 }
408 if self.ledger.contains_transmission(&TransmissionID::Transaction(transaction_id, checksum))? {
410 bail!("Transaction '{}' exists in the ledger {}", fmt_id(transaction_id), "(skipping)".dimmed());
411 }
412 if self.contains_transaction(&transaction_id) {
414 bail!("Transaction '{}' exists in the memory pool", fmt_id(transaction_id));
415 }
416 #[cfg(feature = "metrics")]
417 {
418 metrics::increment_gauge(metrics::consensus::UNCONFIRMED_TRANSACTIONS, 1f64);
419 let timestamp = snarkos_node_bft::helpers::now();
420 self.transmissions_tracker
421 .lock()
422 .insert(TransmissionID::Transaction(transaction.id(), checksum), timestamp);
423 }
424 trace!("Received unconfirmed transaction '{}' in the queue", fmt_id(transaction_id));
426 let priority_fee = transaction.priority_fee_amount()?;
427 self.transactions_queue.write().insert(transaction_id, transaction, priority_fee)?;
428 }
429
430 self.process_unconfirmed_transactions().await
432 }
433
434 async fn process_unconfirmed_transactions(&self) -> Result<()> {
437 let num_unconfirmed_transmissions = self.num_unconfirmed_transmissions();
439 if num_unconfirmed_transmissions >= Primary::<N>::MAX_TRANSMISSIONS_TOLERANCE {
440 return Ok(());
441 }
442 let transactions = {
444 let capacity = Primary::<N>::MAX_TRANSMISSIONS_TOLERANCE.saturating_sub(num_unconfirmed_transmissions);
446 let mut tx_queue = self.transactions_queue.write();
448 let num_deployments = tx_queue.deployments.len().min(capacity).min(MAX_DEPLOYMENTS_PER_INTERVAL);
450 let num_executions = tx_queue.executions.len().min(capacity.saturating_sub(num_deployments));
452 let selector_iter = (0..num_deployments).map(|_| true).interleave((0..num_executions).map(|_| false));
455 selector_iter
457 .filter_map(
458 |select_deployment| {
459 if select_deployment { tx_queue.deployments.pop() } else { tx_queue.executions.pop() }
460 },
461 )
462 .map(|(_, tx)| tx)
463 .collect_vec()
464 };
465 for transaction in transactions.into_iter() {
467 let transaction_id = transaction.id();
468 let tx_type_str = match transaction {
470 Transaction::Deploy(..) => "deployment",
471 Transaction::Execute(..) => "execution",
472 Transaction::Fee(..) => "fee",
473 };
474 trace!("Adding unconfirmed {tx_type_str} transaction '{}' to the memory pool...", fmt_id(transaction_id));
475 match self.primary_sender.send_unconfirmed_transaction(transaction_id, Data::Object(transaction)).await {
477 Ok(true) => {}
478 Ok(false) => debug!(
479 "Unable to add unconfirmed {tx_type_str} transaction '{}' to the memory pool. Already exists.",
480 fmt_id(transaction_id)
481 ),
482 Err(err) => {
483 if self.bft.is_synced() {
485 let err = err.context(format!(
486 "Unable to add unconfirmed {tx_type_str} transaction '{}' to the memory pool",
487 fmt_id(transaction_id)
488 ));
489 warn!("{}", flatten_error(err));
490 }
491 }
492 }
493 }
494 Ok(())
495 }
496}
497
498impl<N: Network> Consensus<N> {
499 pub async fn start_consensus_handlers(&self) -> Result<()> {
505 info!("Starting the consensus instance...");
506
507 let Some((consensus_sender, primary_receiver)) = self.bft_run_args.lock().take() else {
509 bail!("The consensus handlers have already been started");
510 };
511 self.bft
512 .clone()
513 .run(Some(self.ping.clone()), Some(consensus_sender), self.primary_sender.clone(), primary_receiver)
514 .await?;
515
516 let Some(consensus_receiver) = self.consensus_receiver.lock().take() else {
518 bail!("The consensus handlers have already been started");
519 };
520 let ConsensusReceiver { mut rx_consensus_subdag } = consensus_receiver;
521
522 let self_ = self.clone();
524 self.spawn(async move {
525 while let Some((committed_subdag, transmissions, callback)) = rx_consensus_subdag.recv().await {
526 self_.process_bft_subdag(committed_subdag, transmissions, callback).await;
527 }
528 });
529
530 let self_ = self.clone();
538 self.spawn(async move {
539 loop {
540 tokio::select! {
542 _ = self_.block_commit_notify.notified() => {}
543 _ = tokio::time::sleep(MAX_BATCH_DELAY) => {}
544 }
545 if let Err(err) = self_.process_unconfirmed_transactions().await {
547 warn!("{}", flatten_error(err.context("Cannot process unconfirmed transactions")));
548 }
549 if let Err(err) = self_.process_unconfirmed_solutions().await {
551 warn!("{}", flatten_error(err.context("Cannot process unconfirmed solutions")));
552 }
553 }
554 });
555
556 Ok(())
557 }
558
559 async fn process_bft_subdag(
566 &self,
567 subdag: Subdag<N>,
568 transmissions: IndexMap<TransmissionID<N>, Transmission<N>>,
569 callback: oneshot::Sender<Result<bool>>,
570 ) {
571 let self_ = self.clone();
573 let transmissions_ = transmissions.clone();
574 let result = spawn_blocking! { self_.try_advance_to_next_block(subdag, transmissions_).with_context(|| "Unable to advance to the next block") };
575
576 match result {
578 Ok(true) => {
579 self.block_commit_notify.notify_one();
581 }
582 Ok(false) | Err(_) => self.reinsert_transmissions(transmissions).await,
583 }
584
585 callback.send(result).ok();
586 }
587
588 fn try_advance_to_next_block(
595 &self,
596 subdag: Subdag<N>,
597 transmissions: IndexMap<TransmissionID<N>, Transmission<N>>,
598 ) -> Result<bool> {
599 #[cfg(feature = "metrics")]
600 let start = subdag.leader_certificate().batch_header().timestamp();
601 #[cfg(feature = "metrics")]
602 let num_committed_certificates = subdag.values().map(|c| c.len()).sum::<usize>();
603 #[cfg(feature = "metrics")]
604 let current_block_timestamp = self.ledger.latest_block().header().metadata().timestamp();
605
606 let ledger_update = self.ledger.begin_ledger_update()?;
608
609 let prepare_instant = std::time::Instant::now();
610 let block = match ledger_update.prepare_advance_to_next_quorum_block(subdag, transmissions) {
611 Ok(block) => block,
612 Err(err) => return Err(err.into_anyhow()),
613 };
614 let prepare_elapsed = prepare_instant.elapsed();
615 trace!("prepare_advance_to_next_quorum_block took {:.3}s", prepare_elapsed.as_secs_f64());
616 #[cfg(feature = "metrics")]
617 metrics::histogram(metrics::consensus::PREPARE_ADVANCE_SECS, prepare_elapsed.as_secs_f64());
618
619 let check_instant = std::time::Instant::now();
620 cfg_if! {
621 if #[cfg(feature = "test_network")] {
622 let result = if self.ledger.dev_committee_for_round(block.round())?.is_some() {
624 Ok(block)
625 } else {
626 ledger_update.check_next_block(block)
627 };
628 } else {
629 let result = ledger_update.check_next_block(block);
630 }
631 }
632
633 let block = match result {
634 Ok(block) => block,
635 Err(CheckBlockError::BlockAlreadyExists { .. }) => {
636 debug!("The given block hash already exists in the ledger");
637 return Ok(false);
638 }
639 Err(CheckBlockError::InvalidHeight { .. }) => {
640 debug!("The ledger advanced while we were constructing the next block");
641 return Ok(false);
642 }
643 Err(CheckBlockError::InvalidRound { new, previous }) => {
644 debug!("The subDAG round is too low. Expected >{previous}, got {new}");
645 return Ok(false);
646 }
647 Err(err) => return Err(err.into_anyhow()),
648 };
649
650 let check_elapsed = check_instant.elapsed();
651 trace!("check_next_block took {:.3}s", check_elapsed.as_secs_f64());
652 #[cfg(feature = "metrics")]
653 metrics::histogram(metrics::consensus::CHECK_NEXT_BLOCK_SECS, check_elapsed.as_secs_f64());
654
655 let block_height = block.height();
656
657 let advance_instant = std::time::Instant::now();
659 ledger_update.advance_to_next_block(&block)?;
660 let advance_elapsed = advance_instant.elapsed();
661 trace!("advance_to_next_block took {:.3}s", advance_elapsed.as_secs_f64());
662 #[cfg(feature = "metrics")]
663 metrics::histogram(metrics::consensus::ADVANCE_TO_NEXT_BLOCK_SECS, advance_elapsed.as_secs_f64());
664
665 #[cfg(feature = "telemetry")]
666 let latest_committee = self.ledger.get_committee_lookback_for_round(self.ledger.latest_round());
670
671 if block_height.is_multiple_of(N::NUM_BLOCKS_PER_EPOCH) {
673 self.solutions_queue.lock().clear();
675 self.bft.primary().clear_worker_solutions();
677 }
678
679 match self.block_sync.get_block_locators() {
681 Ok(locators) => self.ping.update_block_locators(locators),
682 Err(err) => error!(
683 "{}",
684 flatten_error(err.context("Failed to generate new block locators after block advancement"))
685 ),
686 }
687
688 self.block_sync.set_sync_height(block_height);
690
691 #[cfg(feature = "metrics")]
695 {
696 let now_utc = snarkos_node_bft::helpers::now_utc();
697 let elapsed = std::time::Duration::from_secs((now_utc.unix_timestamp() - start) as u64);
698 let next_block_timestamp = block.header().metadata().timestamp();
699 let next_block_utc = snarkos_node_bft::helpers::to_utc_datetime(next_block_timestamp);
700 let block_latency = next_block_timestamp - current_block_timestamp;
701 let block_lag = (now_utc - next_block_utc).whole_milliseconds();
702
703 let proof_target = block.header().proof_target();
704 let coinbase_target = block.header().coinbase_target();
705 let cumulative_proof_target = block.header().cumulative_proof_target();
706
707 metrics::add_transmission_latency_metric(&self.transmissions_tracker, &block);
709
710 metrics::gauge(metrics::consensus::COMMITTED_CERTIFICATES, num_committed_certificates as f64);
711 metrics::histogram(metrics::consensus::CERTIFICATE_COMMIT_LATENCY, elapsed.as_secs_f64());
712 metrics::histogram(metrics::consensus::BLOCK_LATENCY, block_latency as f64);
713 metrics::histogram(metrics::consensus::BLOCK_LAG, block_lag as f64);
714 metrics::gauge(metrics::blocks::PROOF_TARGET, proof_target as f64);
715 metrics::gauge(metrics::blocks::COINBASE_TARGET, coinbase_target as f64);
716 metrics::gauge(metrics::blocks::CUMULATIVE_PROOF_TARGET, cumulative_proof_target as f64);
717
718 #[cfg(feature = "telemetry")]
720 {
721 match latest_committee {
722 Ok(latest_committee) => {
723 let participation_scores = self
725 .bft()
726 .primary()
727 .gateway()
728 .validator_telemetry()
729 .get_participation_scores(&latest_committee);
730
731 for (address, (certificate_score, signature_score)) in participation_scores {
733 let address_str = address.to_string();
734 metrics::gauge_label(
735 metrics::consensus::VALIDATOR_CERTIFICATE_PARTICIPATION,
736 "validator_address",
737 address_str.clone(),
738 certificate_score,
739 );
740 metrics::gauge_label(
741 metrics::consensus::VALIDATOR_SIGNATURE_PARTICIPATION,
742 "validator_address",
743 address_str,
744 signature_score,
745 );
746 }
747 }
748 Err(err) => warn!("{}", flatten_error(err.context("Failed to get latest committee for telemetry"))),
749 }
750 }
751 }
752
753 Ok(true)
754 }
755
756 async fn reinsert_transmissions(&self, transmissions: IndexMap<TransmissionID<N>, Transmission<N>>) {
758 for (transmission_id, transmission) in transmissions.into_iter() {
760 match self.reinsert_transmission(transmission_id, transmission).await {
762 Ok(true) => {}
763 Ok(false) => debug!(
764 "Unable to reinsert transmission {}:{} into the memory pool. Already exists.",
765 fmt_id(transmission_id),
766 fmt_id(transmission_id.checksum().unwrap_or_default()).dimmed()
767 ),
768 Err(err) => {
769 let err = err.context(format!(
770 "Unable to reinsert transmission {}.{} into the memory pool",
771 fmt_id(transmission_id),
772 fmt_id(transmission_id.checksum().unwrap_or_default()).dimmed()
773 ));
774 warn!("{}", flatten_error(err));
775 }
776 }
777 }
778 }
779
780 async fn reinsert_transmission(
787 &self,
788 transmission_id: TransmissionID<N>,
789 transmission: Transmission<N>,
790 ) -> Result<bool> {
791 let (callback, callback_receiver) = oneshot::channel();
793 match (transmission_id, transmission) {
795 (TransmissionID::Ratification, Transmission::Ratification) => return Ok(true),
796 (TransmissionID::Solution(solution_id, _), Transmission::Solution(solution)) => {
797 self.primary_sender.tx_unconfirmed_solution.send((solution_id, solution, callback)).await?;
799 }
800 (TransmissionID::Transaction(transaction_id, _), Transmission::Transaction(transaction)) => {
801 self.primary_sender.tx_unconfirmed_transaction.send((transaction_id, transaction, callback)).await?;
803 }
804 _ => bail!("Mismatching `(transmission_id, transmission)` pair in consensus"),
805 }
806 callback_receiver.await?
808 }
809
810 fn spawn<T: Future<Output = ()> + Send + 'static>(&self, future: T) {
812 self.handles.lock().push(tokio::spawn(future));
813 }
814
815 pub async fn shut_down(&self) {
817 info!("Shutting down consensus...");
818 self.bft.shut_down().await;
820 self.handles.lock().iter().for_each(|handle| handle.abort());
822 }
823}