Skip to main content

iota_core/
mock_consensus.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5use 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    // ConsensusClient does absolutely nothing when receiving a transaction
38    Noop,
39    // ConsensusClient directly sequences the transaction into the store.
40    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        // TODO: maybe support multi-transactions and remove this check
114        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}