sui_core/
mock_consensus.rs1use crate::authority::authority_per_epoch_store::AuthorityPerEpochStore;
5use crate::authority::{AuthorityState, ExecutionEnv};
6use crate::consensus_adapter::{BlockStatusReceiver, ConsensusClient, SubmitToConsensus};
7
8use consensus_types::block::BlockRef;
9use std::sync::{Arc, Weak};
10use std::time::Duration;
11use sui_types::committee::EpochId;
12use sui_types::error::{SuiError, SuiResult};
13use sui_types::executable_transaction::VerifiedExecutableTransaction;
14use sui_types::messages_consensus::{
15 ConsensusPosition, ConsensusTransaction, ConsensusTransactionKind,
16};
17use sui_types::transaction::VerifiedTransaction;
18use tokio::sync::{mpsc, oneshot};
19use tokio::task::JoinHandle;
20use tracing::debug;
21
22pub struct MockConsensusClient {
23 tx_sender: mpsc::Sender<ConsensusTransaction>,
24 _consensus_handle: JoinHandle<()>,
25}
26
27pub enum ConsensusMode {
28 Noop,
30 DirectSequencing,
32}
33
34impl MockConsensusClient {
35 pub fn new(validator: Weak<AuthorityState>, consensus_mode: ConsensusMode) -> Self {
36 let (tx_sender, tx_receiver) = mpsc::channel(1000000);
37 let _consensus_handle = Self::run(validator, tx_receiver, consensus_mode);
38 Self {
39 tx_sender,
40 _consensus_handle,
41 }
42 }
43
44 pub fn run(
45 validator: Weak<AuthorityState>,
46 tx_receiver: mpsc::Receiver<ConsensusTransaction>,
47 consensus_mode: ConsensusMode,
48 ) -> JoinHandle<()> {
49 tokio::spawn(async move { Self::run_impl(validator, tx_receiver, consensus_mode).await })
50 }
51
52 async fn run_impl(
53 validator: Weak<AuthorityState>,
54 mut tx_receiver: mpsc::Receiver<ConsensusTransaction>,
55 consensus_mode: ConsensusMode,
56 ) {
57 while let Some(tx) = tx_receiver.recv().await {
58 let Some(validator) = validator.upgrade() else {
59 debug!("validator shut down; exiting MockConsensusClient");
60 return;
61 };
62 let epoch_store = validator.epoch_store_for_testing();
63 let env = match consensus_mode {
64 ConsensusMode::Noop => ExecutionEnv::new(),
65 ConsensusMode::DirectSequencing => {
66 let executable_tx = match &tx.kind {
70 ConsensusTransactionKind::UserTransactionV2(tx) => {
71 Some(VerifiedExecutableTransaction::new_from_consensus(
72 VerifiedTransaction::new_unchecked(tx.tx().clone()),
73 0,
74 ))
75 }
76 _ => None,
77 };
78
79 if let Some(exec_tx) = executable_tx {
80 let assigned_versions = epoch_store
82 .assign_shared_object_versions_for_tests(
83 validator.get_object_cache_reader().as_ref(),
84 std::slice::from_ref(&exec_tx),
85 )
86 .unwrap();
87
88 let assigned_version = assigned_versions
89 .into_map()
90 .into_iter()
91 .next()
92 .map(|(_, v)| v)
93 .unwrap_or_default();
94 ExecutionEnv::new().with_assigned_versions(assigned_version)
95 } else {
96 ExecutionEnv::new()
97 }
98 }
99 };
100 match &tx.kind {
101 ConsensusTransactionKind::CertifiedTransaction(_)
103 | ConsensusTransactionKind::UserTransaction(_) => {
104 debug!(
105 "Ignoring deprecated {:?} in MockConsensusClient",
106 std::mem::discriminant(&tx.kind)
107 );
108 }
109 ConsensusTransactionKind::UserTransactionV2(tx) if tx.tx().is_consensus_tx() => {
110 validator.execution_scheduler().enqueue(
111 vec![(
112 VerifiedExecutableTransaction::new_from_consensus(
113 VerifiedTransaction::new_unchecked(tx.tx().clone()),
114 0,
115 )
116 .into(),
117 env,
118 )],
119 &epoch_store,
120 );
121 }
122 _ => {}
123 }
124 }
125 }
126
127 fn submit_impl(
128 &self,
129 transactions: &[ConsensusTransaction],
130 ) -> SuiResult<(Vec<ConsensusPosition>, BlockStatusReceiver)> {
131 assert!(transactions.len() == 1);
133 let transaction = &transactions[0];
134 self.tx_sender
135 .try_send(transaction.clone())
136 .map_err(|_| SuiError::from("MockConsensusClient channel overflowed"))?;
137 Ok((
139 vec![ConsensusPosition {
140 epoch: EpochId::MIN,
141 block: BlockRef::MIN,
142 index: 0,
143 }],
144 with_block_status(consensus_core::BlockStatus::Sequenced(BlockRef::MIN)),
145 ))
146 }
147}
148
149impl SubmitToConsensus for MockConsensusClient {
150 fn submit_to_consensus(
151 &self,
152 transactions: &[ConsensusTransaction],
153 _epoch_store: &Arc<AuthorityPerEpochStore>,
154 ) -> SuiResult {
155 self.submit_impl(transactions).map(|_response| ())
156 }
157
158 fn submit_best_effort(
159 &self,
160 transaction: &ConsensusTransaction,
161 _epoch_store: &Arc<AuthorityPerEpochStore>,
162 _timeout: Duration,
163 ) -> SuiResult {
164 self.submit_impl(std::slice::from_ref(transaction))
165 .map(|_response| ())
166 }
167}
168
169#[async_trait::async_trait]
170impl ConsensusClient for MockConsensusClient {
171 async fn submit(
172 &self,
173 transactions: &[ConsensusTransaction],
174 _epoch_store: &Arc<AuthorityPerEpochStore>,
175 ) -> SuiResult<(Vec<ConsensusPosition>, BlockStatusReceiver)> {
176 self.submit_impl(transactions)
177 }
178}
179
180pub(crate) fn with_block_status(status: consensus_core::BlockStatus) -> BlockStatusReceiver {
181 let (tx, rx) = oneshot::channel();
182 tx.send(status).ok();
183 rx
184}