Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions binaries/cuprated/src/blockchain/manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,7 @@ pub(crate) async fn init_blockchain_manager(
broadcast_svc: clearnet_interface.broadcast_svc(),
reorg_lock: Arc::clone(&launch_ctx.reorg_lock),
fast_sync_hashes,
node_events: launch_ctx.node_events.clone(),
};

launch_ctx
Expand Down Expand Up @@ -116,6 +117,8 @@ pub struct BlockchainManager {
reorg_lock: Arc<RwLock<()>>,
/// Fast-sync hashes for this node's network.
fast_sync_hashes: &'static [[u8; 32]],
/// Sender for the node event stream.
node_events: crate::events::NodeEventSender,
}

impl BlockchainManager {
Expand Down
25 changes: 17 additions & 8 deletions binaries/cuprated/src/blockchain/manager/handler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ use cuprate_types::{
use crate::{
blockchain::manager::commands::{BlockchainManagerCommand, IncomingBlockOk},
constants::PANIC_CRITICAL_SERVICE_ERROR,
events::NodeEvent,
};

impl super::BlockchainManager {
Expand Down Expand Up @@ -469,14 +470,15 @@ impl super::BlockchainManager {

match reorg_res {
Ok(()) => {
info!(
top_hash = hex::encode(
self.blockchain_context_service
.blockchain_context()
.top_hash
),
"Successfully reorged"
);
let ctx = self.blockchain_context_service.blockchain_context();
let new_top_hash = ctx.top_hash;
let new_chain_height = ctx.chain_height;
info!(top_hash = hex::encode(new_top_hash), "Successfully reorged");
self.node_events.send(NodeEvent::Reorg {
split_height,
new_top_hash,
new_chain_height,
});
Ok(())
}
Err(e) => {
Expand Down Expand Up @@ -635,6 +637,9 @@ impl super::BlockchainManager {
verified_block: VerifiedBlockInformation,
source: BlockSource,
) {
let height = verified_block.height;
let hash = verified_block.block_hash;

// FIXME: this is pretty inefficient, we should probably return the KI map created in the consensus crate.
let spent_key_images = verified_block
.txs
Expand All @@ -656,6 +661,10 @@ impl super::BlockchainManager {
self.add_valid_block_to_blockchain_database(verified_block)
.await;

if matches!(source, BlockSource::Incoming) {
self.node_events.send(NodeEvent::NewBlock { height, hash });
}

if let Some(block_blob) = block_blob {
let chain_height = self
.blockchain_context_service
Expand Down
239 changes: 222 additions & 17 deletions binaries/cuprated/src/blockchain/manager/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ use monero_oxide::{
ed25519::CompressedPoint,
transaction::{Input, Output, Timelock, Transaction, TransactionPrefix},
};
use tokio::sync::{oneshot, watch};
use tokio::sync::{broadcast, oneshot, watch};
use tower::BoxError;

use cuprate_blockchain::config::Config;
Expand All @@ -21,10 +21,11 @@ use crate::{
check_add_genesis, manager::BlockchainManager, manager::BlockchainManagerCommand,
ConsensusBlockchainReadHandle,
},
events::{NodeEvent, NodeEventListener, NodeEventSender, NODE_EVENT_CHANNEL_CAPACITY},
txpool::TxpoolManagerHandle,
};

async fn mock_manager(data_dir: PathBuf) -> BlockchainManager {
async fn mock_manager(data_dir: PathBuf) -> (BlockchainManager, NodeEventListener) {
let config = Config {
blob_dir: data_dir.clone(),
index_dir: data_dir.clone(),
Expand Down Expand Up @@ -66,16 +67,23 @@ async fn mock_manager(data_dir: PathBuf) -> BlockchainManager {
.await
.unwrap();

BlockchainManager {
blockchain_write_handle,
blockchain_read_handle,
txpool_manager_handle: TxpoolManagerHandle::mock(),
blockchain_context_service,
stop_current_block_downloader: Arc::new(Default::default()),
broadcast_svc: BroadcastSvc::mock(),
reorg_lock: Arc::new(Default::default()),
fast_sync_hashes: &[],
}
let node_events = NodeEventSender::new();
let listener = node_events.subscribe();

(
BlockchainManager {
blockchain_write_handle,
blockchain_read_handle,
txpool_manager_handle: TxpoolManagerHandle::mock(),
blockchain_context_service,
stop_current_block_downloader: Arc::new(Default::default()),
broadcast_svc: BroadcastSvc::mock(),
node_events,
reorg_lock: Arc::new(Default::default()),
fast_sync_hashes: &[],
},
listener,
)
}

fn generate_block(context: &BlockchainContext) -> Block {
Expand Down Expand Up @@ -115,10 +123,12 @@ fn generate_block(context: &BlockchainContext) -> Block {
async fn simple_reorg() {
// create 2 managers
let data_dir_1 = tempfile::tempdir().unwrap();
let mut manager_1 = mock_manager(data_dir_1.path().to_path_buf()).await;
let manager_1_with_events = mock_manager(data_dir_1.path().to_path_buf()).await;
let mut manager_1 = manager_1_with_events.0;

let data_dir_2 = tempfile::tempdir().unwrap();
let mut manager_2 = mock_manager(data_dir_2.path().to_path_buf()).await;
let manager_2_with_events = mock_manager(data_dir_2.path().to_path_buf()).await;
let mut manager_2 = manager_2_with_events.0;

// give both managers the same first non-genesis block
let block_1 = generate_block(manager_1.blockchain_context_service.blockchain_context());
Expand Down Expand Up @@ -228,10 +238,12 @@ async fn simple_reorg_block_batch() {

// create 2 managers
let data_dir_1 = tempfile::tempdir().unwrap();
let mut manager_1 = mock_manager(data_dir_1.path().to_path_buf()).await;
let manager_1_with_events = mock_manager(data_dir_1.path().to_path_buf()).await;
let mut manager_1 = manager_1_with_events.0;

let data_dir_2 = tempfile::tempdir().unwrap();
let mut manager_2 = mock_manager(data_dir_2.path().to_path_buf()).await;
let manager_2_with_events = mock_manager(data_dir_2.path().to_path_buf()).await;
let mut manager_2 = manager_2_with_events.0;

// give both managers the same first non-genesis block
let block_1 = generate_block(manager_1.blockchain_context_service.blockchain_context());
Expand Down Expand Up @@ -337,7 +349,8 @@ async fn simple_reorg_block_batch() {
#[tokio::test]
async fn recover_bad_reorg() {
let data_dir_1 = tempfile::tempdir().unwrap();
let mut manager_1 = mock_manager(data_dir_1.path().to_path_buf()).await;
let manager_1_with_events = mock_manager(data_dir_1.path().to_path_buf()).await;
let mut manager_1 = manager_1_with_events.0;

let context_1 = manager_1
.blockchain_context_service
Expand Down Expand Up @@ -442,3 +455,195 @@ async fn recover_bad_reorg() {
manager_1.blockchain_context_service.blockchain_context()
);
}

#[tokio::test]
async fn node_event_delivered_to_prior_subscriber() {
let s = NodeEventSender::new();
let mut l = s.subscribe();

s.send(NodeEvent::NewBlock {
height: 7,
hash: [1_u8; 32],
});

assert_eq!(
l.recv().await.unwrap(),
NodeEvent::NewBlock {
height: 7,
hash: [1_u8; 32],
}
);
}

#[tokio::test]
async fn node_event_two_subscribers_both_receive() {
let s = NodeEventSender::new();
let mut l1 = s.subscribe();
let mut l2 = s.subscribe();
let event = NodeEvent::NewBlock {
height: 3,
hash: [2_u8; 32],
};

s.send(event.clone());

assert_eq!(l1.recv().await.unwrap(), event);
assert_eq!(l2.recv().await.unwrap(), event);
}

#[test]
fn node_event_no_subscriber_is_noop() {
let sender = NodeEventSender::new();

sender.send(NodeEvent::NewBlock {
height: 5,
hash: [9_u8; 32],
});
}

#[test]
fn node_event_late_subscriber_misses_prior() {
let s = NodeEventSender::new();
s.send(NodeEvent::NewBlock {
height: 8,
hash: [4_u8; 32],
});

let mut listener = s.subscribe();

assert_eq!(
listener.try_recv(),
Err(broadcast::error::TryRecvError::Empty)
);
}

#[test]
fn node_event_channel_capacity_is_256() {
assert_eq!(NODE_EVENT_CHANNEL_CAPACITY, 256);
}

#[tokio::test]
async fn new_block_emits_new_block_event() {
let data_dir = tempfile::tempdir().unwrap();
let manager_with_events = mock_manager(data_dir.path().to_path_buf()).await;
let mut manager = manager_with_events.0;
let mut listener = manager_with_events.1;

let block = generate_block(manager.blockchain_context_service.blockchain_context());
let hash = block.hash();

manager
.handle_command(BlockchainManagerCommand::AddBlock {
block,
prepped_txs: HashMap::new(),
response_tx: oneshot::channel().0,
})
.await;

assert_eq!(
listener.try_recv(),
Ok(NodeEvent::NewBlock { height: 1, hash })
);
assert_eq!(
listener.try_recv(),
Err(broadcast::error::TryRecvError::Empty)
);
}

#[tokio::test]
async fn reorg_emits_reorg_event() {
let data_dir_1 = tempfile::tempdir().unwrap();
let manager_1_with_events = mock_manager(data_dir_1.path().to_path_buf()).await;
let mut manager_1 = manager_1_with_events.0;
let mut listener = manager_1_with_events.1;

let data_dir_2 = tempfile::tempdir().unwrap();
let manager_2_with_events = mock_manager(data_dir_2.path().to_path_buf()).await;
let mut manager_2 = manager_2_with_events.0;

let block_1 = generate_block(manager_1.blockchain_context_service.blockchain_context());

manager_1
.handle_command(BlockchainManagerCommand::AddBlock {
block: block_1.clone(),
prepped_txs: HashMap::new(),
response_tx: oneshot::channel().0,
})
.await;

manager_2
.handle_command(BlockchainManagerCommand::AddBlock {
block: block_1,
prepped_txs: HashMap::new(),
response_tx: oneshot::channel().0,
})
.await;

let block_2a = generate_block(manager_1.blockchain_context_service.blockchain_context());
let block_2b = generate_block(manager_2.blockchain_context_service.blockchain_context());

manager_1
.handle_command(BlockchainManagerCommand::AddBlock {
block: block_2a,
prepped_txs: HashMap::new(),
response_tx: oneshot::channel().0,
})
.await;

manager_2
.handle_command(BlockchainManagerCommand::AddBlock {
block: block_2b.clone(),
prepped_txs: HashMap::new(),
response_tx: oneshot::channel().0,
})
.await;

manager_1
.handle_command(BlockchainManagerCommand::AddBlock {
block: block_2b,
prepped_txs: HashMap::new(),
response_tx: oneshot::channel().0,
})
.await;

let block_3 = generate_block(manager_2.blockchain_context_service.blockchain_context());

// Discard the `NewBlock` events emitted while building the initial main chain (block_1 and
// block_2a were added as live `Incoming` blocks). We only want to observe what the reorg itself
// emits.
while listener.try_recv().is_ok() {}

manager_1
.handle_command(BlockchainManagerCommand::AddBlock {
block: block_3,
prepped_txs: HashMap::new(),
response_tx: oneshot::channel().0,
})
.await;

let mut events = Vec::new();
loop {
match listener.try_recv() {
Ok(event) => events.push(event),
Err(broadcast::error::TryRecvError::Empty) => break,
Err(err) => panic!("unexpected broadcast receive error: {err:?}"),
}
}

let chain_context = manager_1
.blockchain_context_service
.blockchain_context()
.clone();

assert!(events.iter().any(|event| matches!(
event,
NodeEvent::Reorg {
split_height: 2,
new_top_hash: _,
new_chain_height
} if *new_chain_height == chain_context.chain_height
)));
assert!(!events
.iter()
.any(|event| matches!(event, NodeEvent::NewBlock { .. })));
}
Loading
Loading