cuprated/blockchain/
interface.rs1use 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#[derive(Clone)]
31pub struct BlockchainManagerHandle {
32 command_tx: mpsc::Sender<BlockchainManagerCommand>,
34 blocks_being_handled: Arc<Mutex<HashSet<[u8; 32]>>>,
43}
44
45impl BlockchainManagerHandle {
46 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 pub fn is_block_being_handled(&self, hash: &[u8; 32]) -> bool {
60 self.blocks_being_handled.lock().unwrap().contains(hash)
61 }
62
63 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 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 if !self.blocks_being_handled.lock().unwrap().insert(block_hash) {
125 return Ok(IncomingBlockOk::AlreadyHave);
127 }
128
129 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 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
179async 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}