Skip to main content

sui_core/
mock_consensus.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4use 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    // ConsensusClient does absolutely nothing when receiving a transaction
29    Noop,
30    // ConsensusClient directly sequences the transaction into the store.
31    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                    // Extract the executable transaction from the consensus transaction.
67                    // DEPRECATED: CertifiedTransaction and UserTransaction are no longer
68                    // produced (MFP uses UserTransactionV2).
69                    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                        // Use the simpler assign_shared_object_versions_for_tests API
81                        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                // DEPRECATED: CertifiedTransaction and UserTransaction are no longer produced.
102                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        // TODO: maybe support multi-transactions and remove this check
132        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        // TODO(fastpath): Add some way to simulate consensus positions across blocks
138        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}