Skip to main content

cuprated/blockchain/
interface.rs

1//! The blockchain manager interface.
2//!
3//! This module contains all the functions to mutate the blockchain's state in any way, through the
4//! blockchain manager.
5use std::{
6    collections::{HashMap, HashSet},
7    sync::{Arc, Mutex},
8};
9
10use monero_oxide::{block::Block, transaction::Transaction};
11use tokio::sync::{mpsc, oneshot};
12use tower::{Service, ServiceExt};
13
14use cuprate_blockchain::{service::BlockchainReadHandle, BlockchainError};
15use cuprate_consensus::{block::BlockVerificationError, transactions::new_tx_verification_data};
16use cuprate_txpool::service::{
17    interface::{TxpoolReadRequest, TxpoolReadResponse},
18    TxpoolReadHandle,
19};
20use cuprate_types::blockchain::{BlockchainReadRequest, BlockchainResponse};
21
22use crate::blockchain::{
23    manager::{BlockchainManagerCommand, IncomingBlockOk},
24    IncomingBlockError,
25};
26
27/// Handle for the blockchain manager.
28///
29/// Created by `init_blockchain_manager`.
30#[derive(Clone)]
31pub struct BlockchainManagerHandle {
32    /// The channel used to send [`BlockchainManagerCommand`]s to the blockchain manager.
33    command_tx: mpsc::Sender<BlockchainManagerCommand>,
34    /// A [`HashSet`] of block hashes that the blockchain manager is currently handling.
35    ///
36    /// This prevents sending the same block to the blockchain manager from multiple connections
37    /// before one of them actually gets added to the chain, allowing peers to do other things.
38    ///
39    /// This is used over something like a dashmap as we expect a lot of collisions in a short amount of
40    /// time for new blocks, so we would lose the benefit of sharded locks. A dashmap is made up of `RwLocks`
41    /// which are also more expensive than `Mutex`s.
42    blocks_being_handled: Arc<Mutex<HashSet<[u8; 32]>>>,
43}
44
45impl BlockchainManagerHandle {
46    /// Create a new handle and command receiver pair.
47    pub(crate) fn new() -> (Self, mpsc::Receiver<BlockchainManagerCommand>) {
48        let (command_tx, command_rx) = mpsc::channel(3);
49        (
50            Self {
51                command_tx,
52                blocks_being_handled: Arc::new(Mutex::new(HashSet::new())),
53            },
54            command_rx,
55        )
56    }
57
58    /// Returns `true` if the given block hash is currently being handled.
59    pub fn is_block_being_handled(&self, hash: &[u8; 32]) -> bool {
60        self.blocks_being_handled.lock().unwrap().contains(hash)
61    }
62
63    /// Try to add a new block to the blockchain.
64    ///
65    /// On success returns `IncomingBlockOk`.
66    ///
67    /// # Errors
68    ///
69    /// This function will return an error if:
70    ///  - the block was invalid
71    ///  - we are missing transactions
72    ///  - the block's parent is unknown
73    ///  - the blockchain manager command channel is closed
74    pub async fn handle_incoming_block(
75        &self,
76        block: Block,
77        mut given_txs: HashMap<[u8; 32], Transaction>,
78        blockchain_read_handle: &mut BlockchainReadHandle,
79        txpool_read_handle: &mut TxpoolReadHandle,
80    ) -> Result<IncomingBlockOk, IncomingBlockError> {
81        if given_txs.len() > block.transactions.len() {
82            return Err(IncomingBlockError::TooManyTxs);
83        }
84
85        if !block_exists(block.header.previous, blockchain_read_handle).await? {
86            return Err(IncomingBlockError::Orphan);
87        }
88
89        let block_hash = block.hash();
90
91        if block_exists(block_hash, blockchain_read_handle).await? {
92            return Ok(IncomingBlockOk::AlreadyHave);
93        }
94
95        let TxpoolReadResponse::TxsForBlock { mut txs, missing } = txpool_read_handle
96            .ready()
97            .await?
98            .call(TxpoolReadRequest::TxsForBlock(block.transactions.clone()))
99            .await?
100        else {
101            unreachable!()
102        };
103
104        if !missing.is_empty() {
105            let needed_hashes = missing.iter().map(|index| block.transactions[*index]);
106
107            for needed_hash in needed_hashes {
108                let Some(tx) = given_txs.remove(&needed_hash) else {
109                    // We return back the indexes of all txs missing from our pool, not taking into account the txs
110                    // that were given with the block, as these txs will be dropped. It is not worth it to try to add
111                    // these txs to the pool as this will only happen with a misbehaving peer or if the txpool reaches
112                    // the size limit.
113                    return Err(IncomingBlockError::UnknownTransactions(block_hash, missing));
114                };
115
116                txs.insert(
117                    needed_hash,
118                    new_tx_verification_data(tx).map_err(BlockVerificationError::invalid_pow)?,
119                );
120            }
121        }
122
123        // Add the blocks hash to the blocks being handled.
124        if !self.blocks_being_handled.lock().unwrap().insert(block_hash) {
125            // If another place is already adding this block then we can stop.
126            return Ok(IncomingBlockOk::AlreadyHave);
127        }
128
129        // We must remove the block hash from `blocks_being_handled`.
130        let blocks = Arc::clone(&self.blocks_being_handled);
131        let _guard = {
132            struct RemoveFromBlocksBeingHandled {
133                block_hash: [u8; 32],
134                blocks: Arc<Mutex<HashSet<[u8; 32]>>>,
135            }
136            impl Drop for RemoveFromBlocksBeingHandled {
137                fn drop(&mut self) {
138                    self.blocks.lock().unwrap().remove(&self.block_hash);
139                }
140            }
141            RemoveFromBlocksBeingHandled { block_hash, blocks }
142        };
143
144        let (response_tx, response_rx) = oneshot::channel();
145
146        self.command_tx
147            .send(BlockchainManagerCommand::AddBlock {
148                block,
149                prepped_txs: txs,
150                response_tx,
151            })
152            .await
153            .map_err(|_| IncomingBlockError::ChannelClosed)?;
154
155        response_rx
156            .await
157            .map_err(|_| IncomingBlockError::ChannelClosed)?
158    }
159
160    /// Pop blocks from the top of the blockchain.
161    ///
162    /// # Errors
163    ///
164    /// Will error if the blockchain manager channel is closed.
165    pub async fn pop_blocks(&self, numb_blocks: usize) -> Result<(), anyhow::Error> {
166        let (response_tx, response_rx) = oneshot::channel();
167
168        self.command_tx
169            .send(BlockchainManagerCommand::PopBlocks {
170                numb_blocks,
171                response_tx,
172            })
173            .await?;
174
175        Ok(response_rx.await?)
176    }
177}
178
179/// Check if we have a block with the given hash.
180async fn block_exists(
181    block_hash: [u8; 32],
182    blockchain_read_handle: &mut BlockchainReadHandle,
183) -> Result<bool, BlockchainError> {
184    let BlockchainResponse::FindBlock(chain) = blockchain_read_handle
185        .ready()
186        .await?
187        .call(BlockchainReadRequest::FindBlock(block_hash))
188        .await?
189    else {
190        unreachable!();
191    };
192
193    Ok(chain.is_some())
194}