iota_core/
mock_consensus.rs1use std::sync::{Arc, Weak};
6
7use iota_types::{
8 error::{IotaError, IotaResult},
9 executable_transaction::VerifiedExecutableTransaction,
10 messages_consensus::{ConsensusTransaction, ConsensusTransactionKind},
11 transaction::{SenderSignedTransactionAPI, VerifiedCertificate},
12};
13use prometheus_filtered::Registry;
14use starfish_core::BlockRef;
15use tokio::{
16 sync::{mpsc, oneshot},
17 task::JoinHandle,
18};
19use tracing::debug;
20
21use crate::{
22 authority::{
23 AuthorityMetrics, AuthorityState, ExecutionEnv,
24 authority_per_epoch_store::AuthorityPerEpochStore,
25 },
26 checkpoints::CheckpointServiceNoop,
27 consensus_adapter::{BlockStatusReceiver, ConsensusClient, SubmitToConsensus},
28 consensus_handler::SequencedConsensusTransaction,
29 execution_scheduler::ExecutionSchedulerAPI,
30};
31pub struct MockConsensusClient {
32 tx_sender: mpsc::Sender<ConsensusTransaction>,
33 _consensus_handle: JoinHandle<()>,
34}
35
36pub enum ConsensusMode {
37 Noop,
39 DirectSequencing,
41}
42
43impl MockConsensusClient {
44 pub fn new(validator: Weak<AuthorityState>, consensus_mode: ConsensusMode) -> Self {
45 let (tx_sender, tx_receiver) = mpsc::channel(1000000);
46 let _consensus_handle = Self::run(validator, tx_receiver, consensus_mode);
47 Self {
48 tx_sender,
49 _consensus_handle,
50 }
51 }
52
53 pub fn run(
54 validator: Weak<AuthorityState>,
55 tx_receiver: mpsc::Receiver<ConsensusTransaction>,
56 consensus_mode: ConsensusMode,
57 ) -> JoinHandle<()> {
58 tokio::spawn(async move { Self::run_impl(validator, tx_receiver, consensus_mode).await })
59 }
60
61 async fn run_impl(
62 validator: Weak<AuthorityState>,
63 mut tx_receiver: mpsc::Receiver<ConsensusTransaction>,
64 consensus_mode: ConsensusMode,
65 ) {
66 let checkpoint_service = Arc::new(CheckpointServiceNoop {});
67 let authority_metrics = Arc::new(AuthorityMetrics::new(&Registry::new()));
68 while let Some(tx) = tx_receiver.recv().await {
69 let Some(validator) = validator.upgrade() else {
70 debug!("validator shut down; exiting MockConsensusClient");
71 return;
72 };
73 let epoch_store = validator.epoch_store_for_testing();
74 let assigned_versions = match consensus_mode {
75 ConsensusMode::Noop => None,
76 ConsensusMode::DirectSequencing => {
77 let (_, assigned_versions) = epoch_store
78 .process_consensus_transactions_for_tests(
79 vec![SequencedConsensusTransaction::new_test(tx.clone())],
80 &checkpoint_service,
81 validator.get_object_cache_reader().as_ref(),
82 &authority_metrics,
83 true,
84 validator.as_ref(),
85 )
86 .await
87 .unwrap();
88 Some(assigned_versions.into_map())
89 }
90 };
91 if let ConsensusTransactionKind::CertifiedTransaction(tx) = tx.kind {
92 if tx.contains_shared_object() {
93 let transaction = VerifiedExecutableTransaction::new_from_certificate(
94 VerifiedCertificate::new_unchecked(*tx),
95 );
96 let env = ExecutionEnv::new().with_assigned_versions(
97 assigned_versions
98 .and_then(|mut map| map.remove(&transaction.key()))
99 .unwrap_or_default(),
100 );
101 validator
102 .execution_scheduler()
103 .enqueue(vec![(transaction.into(), env)], &epoch_store);
104 }
105 }
106 }
107 }
108
109 fn submit_impl(
110 &self,
111 transactions: &[ConsensusTransaction],
112 ) -> IotaResult<BlockStatusReceiver> {
113 assert!(transactions.len() == 1);
115 let transaction = &transactions[0];
116 self.tx_sender
117 .try_send(transaction.clone())
118 .map_err(|_| IotaError::from("MockConsensusClient channel overflowed"))?;
119 Ok(with_block_status(starfish_core::BlockStatus::Sequenced(
120 starfish_core::GenericTransactionRef::BlockRef(BlockRef::MIN),
121 )))
122 }
123}
124
125impl SubmitToConsensus for MockConsensusClient {
126 fn submit_to_consensus(
127 &self,
128 transactions: &[ConsensusTransaction],
129 _epoch_store: &Arc<AuthorityPerEpochStore>,
130 ) -> IotaResult {
131 self.submit_impl(transactions).map(|_response| ())
132 }
133}
134
135#[async_trait::async_trait]
136impl ConsensusClient for MockConsensusClient {
137 async fn submit(
138 &self,
139 transactions: &[ConsensusTransaction],
140 _epoch_store: &Arc<AuthorityPerEpochStore>,
141 ) -> IotaResult<BlockStatusReceiver> {
142 self.submit_impl(transactions)
143 }
144}
145
146pub(crate) fn with_block_status(status: starfish_core::BlockStatus) -> BlockStatusReceiver {
147 let (tx, rx) = oneshot::channel();
148 tx.send(status.into()).ok();
149 rx
150}