use {
crate::rpc_subscriptions::RpcSubscriptions,
crossbeam_channel::{Receiver, RecvTimeoutError, Sender},
solana_clock::{BankId, Slot},
solana_hash::Hash,
solana_rpc_client_api::response::{SlotTransactionStats, SlotUpdate},
solana_runtime::{
bank::Bank, bank_forks::BankForks, dependency_tracker::DependencyTracker,
prioritization_fee_cache::PrioritizationFeeCache,
},
solana_time_utils::timestamp,
std::{
collections::HashSet,
sync::{
Arc, RwLock,
atomic::{AtomicBool, Ordering},
},
thread::{self, Builder, JoinHandle},
time::Duration,
},
};
pub struct OptimisticallyConfirmedBank {
pub bank: Arc<Bank>,
}
impl OptimisticallyConfirmedBank {
pub fn locked_from_bank_forks_root(bank_forks: &RwLock<BankForks>) -> Arc<RwLock<Self>> {
Arc::new(RwLock::new(Self {
bank: bank_forks.read().unwrap().root_bank(),
}))
}
}
#[derive(Clone)]
pub enum BankNotification {
OptimisticallyConfirmed(Slot, Hash),
Frozen(Arc<Bank>),
NewRootBank(Arc<Bank>),
NewRootedChain(Vec<(Slot, BankId)>, Slot),
}
#[derive(Clone, Debug)]
pub enum SlotNotification {
OptimisticallyConfirmed(Slot, BankId),
Frozen((Slot, Slot, BankId)),
Root((Slot, Slot, BankId)),
}
impl std::fmt::Debug for BankNotification {
fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
match self {
BankNotification::OptimisticallyConfirmed(slot, hash) => {
write!(f, "OptimisticallyConfirmed({slot:?}, {hash:?})")
}
BankNotification::Frozen(bank) => write!(f, "Frozen({})", bank.slot()),
BankNotification::NewRootBank(bank) => write!(f, "Root({})", bank.slot()),
BankNotification::NewRootedChain(chain, parent) => {
write!(f, "RootedChain({chain:?}, parent: {parent})")
}
}
}
}
pub type BankNotificationWithDependencyWork = (
BankNotification,
Option<u64>, );
pub type BankNotificationReceiver = Receiver<BankNotificationWithDependencyWork>;
pub type BankNotificationSender = Sender<BankNotificationWithDependencyWork>;
#[derive(Clone)]
pub struct BankNotificationSenderConfig {
pub sender: BankNotificationSender,
pub should_send_parents: bool,
pub dependency_tracker: Option<Arc<DependencyTracker>>,
}
pub type SlotNotificationReceiver = Receiver<SlotNotification>;
pub type SlotNotificationSender = Sender<SlotNotification>;
type PendingOptimisticallyConfirmedBanks = HashSet<(Slot, Hash)>;
pub struct OptimisticallyConfirmedBankTracker {
thread_hdl: JoinHandle<()>,
}
impl OptimisticallyConfirmedBankTracker {
pub fn new(
receiver: BankNotificationReceiver,
exit: Arc<AtomicBool>,
bank_forks: Arc<RwLock<BankForks>>,
optimistically_confirmed_bank: Arc<RwLock<OptimisticallyConfirmedBank>>,
subscriptions: Arc<RpcSubscriptions>,
slot_notification_subscribers: Option<Arc<RwLock<Vec<SlotNotificationSender>>>>,
prioritization_fee_cache: Option<Arc<PrioritizationFeeCache>>,
dependency_tracker: Option<Arc<DependencyTracker>>,
) -> Self {
let mut pending_optimistically_confirmed_banks = HashSet::new();
let mut last_notified_confirmed_slot: Slot = 0;
let mut highest_confirmed_slot: Slot = 0;
let mut newest_root_slot: Slot = 0;
let thread_hdl = Builder::new()
.name("solOpConfBnkTrk".to_string())
.spawn(move || {
loop {
if exit.load(Ordering::Relaxed) {
break;
}
if let Err(RecvTimeoutError::Disconnected) = Self::recv_notification(
&receiver,
&bank_forks,
&optimistically_confirmed_bank,
&subscriptions,
&mut pending_optimistically_confirmed_banks,
&mut last_notified_confirmed_slot,
&mut highest_confirmed_slot,
&mut newest_root_slot,
&slot_notification_subscribers,
prioritization_fee_cache.as_deref(),
&dependency_tracker,
) {
break;
}
}
})
.unwrap();
Self { thread_hdl }
}
#[allow(clippy::too_many_arguments)]
fn recv_notification(
receiver: &Receiver<BankNotificationWithDependencyWork>,
bank_forks: &RwLock<BankForks>,
optimistically_confirmed_bank: &RwLock<OptimisticallyConfirmedBank>,
subscriptions: &RpcSubscriptions,
pending_optimistically_confirmed_banks: &mut PendingOptimisticallyConfirmedBanks,
last_notified_confirmed_slot: &mut Slot,
highest_confirmed_slot: &mut Slot,
newest_root_slot: &mut Slot,
slot_notification_subscribers: &Option<Arc<RwLock<Vec<SlotNotificationSender>>>>,
prioritization_fee_cache: Option<&PrioritizationFeeCache>,
dependency_tracker: &Option<Arc<DependencyTracker>>,
) -> Result<(), RecvTimeoutError> {
let notification = receiver.recv_timeout(Duration::from_secs(1))?;
Self::process_notification(
notification,
bank_forks,
optimistically_confirmed_bank,
subscriptions,
pending_optimistically_confirmed_banks,
last_notified_confirmed_slot,
highest_confirmed_slot,
newest_root_slot,
slot_notification_subscribers,
prioritization_fee_cache,
dependency_tracker,
);
Ok(())
}
fn notify_slot_status(
slot_notification_subscribers: &Option<Arc<RwLock<Vec<SlotNotificationSender>>>>,
notification: SlotNotification,
) {
if let Some(slot_notification_subscribers) = slot_notification_subscribers {
for sender in slot_notification_subscribers.read().unwrap().iter() {
match sender.send(notification.clone()) {
Ok(_) => {}
Err(err) => {
info!("Failed to send notification {notification:?}, error: {err:?}");
}
}
}
}
}
fn notify_or_defer(
subscriptions: &RpcSubscriptions,
bank_forks: &RwLock<BankForks>,
bank: &Bank,
pending_confirmation: Option<(Slot, Hash)>,
last_notified_confirmed_slot: &mut Slot,
pending_optimistically_confirmed_banks: &mut PendingOptimisticallyConfirmedBanks,
slot_notification_subscribers: &Option<Arc<RwLock<Vec<SlotNotificationSender>>>>,
prioritization_fee_cache: Option<&PrioritizationFeeCache>,
) {
if bank.is_frozen() {
if bank.slot() > *last_notified_confirmed_slot {
debug!(
"notify_or_defer notifying via notify_gossip_subscribers for slot {:?}",
bank.slot()
);
subscriptions.notify_gossip_subscribers(bank.slot());
*last_notified_confirmed_slot = bank.slot();
Self::notify_slot_status(
slot_notification_subscribers,
SlotNotification::OptimisticallyConfirmed(bank.slot(), bank.bank_id()),
);
if let Some(prioritization_fee_cache) = prioritization_fee_cache {
prioritization_fee_cache.finalize_priority_fee(bank.slot(), bank.bank_id());
}
}
} else if let Some((slot, hash)) = pending_confirmation
&& bank.slot() == slot
&& bank.slot() > bank_forks.read().unwrap().root()
{
pending_optimistically_confirmed_banks.insert((slot, hash));
debug!("notify_or_defer defer notifying for slot {slot:?}");
}
}
fn notify_or_defer_confirmed_banks(
subscriptions: &RpcSubscriptions,
bank_forks: &RwLock<BankForks>,
bank: Arc<Bank>,
slot_threshold: Slot,
pending_confirmation: Option<(Slot, Hash)>,
last_notified_confirmed_slot: &mut Slot,
pending_optimistically_confirmed_banks: &mut PendingOptimisticallyConfirmedBanks,
slot_notification_subscribers: &Option<Arc<RwLock<Vec<SlotNotificationSender>>>>,
prioritization_fee_cache: Option<&PrioritizationFeeCache>,
) {
for confirmed_bank in bank.parents_inclusive().iter().rev() {
if confirmed_bank.slot() > slot_threshold {
debug!(
"Calling notify_or_defer for confirmed_bank {:?}",
confirmed_bank.slot()
);
Self::notify_or_defer(
subscriptions,
bank_forks,
confirmed_bank,
pending_confirmation,
last_notified_confirmed_slot,
pending_optimistically_confirmed_banks,
slot_notification_subscribers,
prioritization_fee_cache,
);
}
}
}
fn notify_new_root_slots(
roots: &mut [(Slot, BankId)],
oldest_parent: Slot,
newest_root_slot: &mut Slot,
slot_notification_subscribers: &Option<Arc<RwLock<Vec<SlotNotificationSender>>>>,
) {
if slot_notification_subscribers.is_none() {
return;
}
roots.sort_unstable_by_key(|(root, _bank_id)| *root);
assert!(!roots.is_empty());
let mut parent = oldest_parent;
for (root, bank_id) in roots.iter() {
if *root > *newest_root_slot {
debug!("Doing SlotNotification::Root for root {root}, parent: {parent}");
Self::notify_slot_status(
slot_notification_subscribers,
SlotNotification::Root((*root, parent, *bank_id)),
);
*newest_root_slot = *root;
}
parent = *root;
}
}
#[allow(clippy::too_many_arguments)]
pub fn process_notification(
(notification, dependency_work): BankNotificationWithDependencyWork,
bank_forks: &RwLock<BankForks>,
optimistically_confirmed_bank: &RwLock<OptimisticallyConfirmedBank>,
subscriptions: &RpcSubscriptions,
pending_optimistically_confirmed_banks: &mut PendingOptimisticallyConfirmedBanks,
last_notified_confirmed_slot: &mut Slot,
highest_confirmed_slot: &mut Slot,
newest_root_slot: &mut Slot,
slot_notification_subscribers: &Option<Arc<RwLock<Vec<SlotNotificationSender>>>>,
prioritization_fee_cache: Option<&PrioritizationFeeCache>,
dependency_tracker: &Option<Arc<DependencyTracker>>,
) {
debug!("received bank notification: {notification:?} event: {dependency_work:?}");
if let Some(tracker) = dependency_tracker.as_ref()
&& let Some(dependency_work) = dependency_work
{
tracker.wait_for_dependency(dependency_work);
}
match notification {
BankNotification::OptimisticallyConfirmed(slot, hash) => {
let bank = bank_forks.read().unwrap().get(slot);
if let Some(bank) = bank {
if bank.is_frozen() {
if bank.hash() != hash {
if slot > bank_forks.read().unwrap().root() {
pending_optimistically_confirmed_banks.insert((slot, hash));
debug!(
"defer notifying optimistic confirmation for slot {slot}: \
local bank hash {} does not match optimistic confirmation \
hash {hash}",
bank.hash()
);
}
} else {
let mut w_optimistically_confirmed_bank =
optimistically_confirmed_bank.write().unwrap();
if bank.slot() > w_optimistically_confirmed_bank.bank.slot() {
w_optimistically_confirmed_bank.bank = bank.clone();
}
if slot > *highest_confirmed_slot {
Self::notify_or_defer_confirmed_banks(
subscriptions,
bank_forks,
bank,
*highest_confirmed_slot,
None,
last_notified_confirmed_slot,
pending_optimistically_confirmed_banks,
slot_notification_subscribers,
prioritization_fee_cache,
);
*highest_confirmed_slot = slot;
}
drop(w_optimistically_confirmed_bank);
}
} else if slot > bank_forks.read().unwrap().root() {
pending_optimistically_confirmed_banks.insert((slot, hash));
debug!("defer notifying optimistic confirmation for slot {slot}");
} else {
inc_new_counter_info!(
"dropped-already-rooted-optimistic-bank-notification",
1
);
}
} else if slot > bank_forks.read().unwrap().root() {
pending_optimistically_confirmed_banks.insert((slot, hash));
} else {
inc_new_counter_info!("dropped-already-rooted-optimistic-bank-notification", 1);
}
subscriptions.notify_slot_update(SlotUpdate::OptimisticConfirmation {
slot,
timestamp: timestamp(),
});
}
BankNotification::Frozen(bank) => {
let frozen_slot = bank.slot();
if let Some(parent) = bank.parent() {
let num_successful_transactions = bank
.transaction_count()
.saturating_sub(parent.transaction_count());
subscriptions.notify_slot_update(SlotUpdate::Frozen {
slot: frozen_slot,
timestamp: timestamp(),
stats: SlotTransactionStats {
num_transaction_entries: bank.transaction_entries_count(),
num_successful_transactions,
num_failed_transactions: bank.transaction_error_count(),
max_transactions_per_entry: bank.transactions_per_entry_max(),
},
});
Self::notify_slot_status(
slot_notification_subscribers,
SlotNotification::Frozen((bank.slot(), bank.parent_slot(), bank.bank_id())),
);
}
if pending_optimistically_confirmed_banks.remove(&(bank.slot(), bank.hash())) {
debug!(
"Calling notify_gossip_subscribers to send deferred notification \
{frozen_slot:?}"
);
Self::notify_or_defer_confirmed_banks(
subscriptions,
bank_forks,
bank.clone(),
*last_notified_confirmed_slot,
None,
last_notified_confirmed_slot,
pending_optimistically_confirmed_banks,
slot_notification_subscribers,
prioritization_fee_cache,
);
if frozen_slot > *highest_confirmed_slot {
*highest_confirmed_slot = frozen_slot;
}
let mut w_optimistically_confirmed_bank =
optimistically_confirmed_bank.write().unwrap();
if frozen_slot > w_optimistically_confirmed_bank.bank.slot() {
w_optimistically_confirmed_bank.bank = bank;
}
drop(w_optimistically_confirmed_bank);
}
}
BankNotification::NewRootBank(bank) => {
let root_slot = bank.slot();
let mut w_optimistically_confirmed_bank =
optimistically_confirmed_bank.write().unwrap();
if root_slot > w_optimistically_confirmed_bank.bank.slot() {
w_optimistically_confirmed_bank.bank = bank;
}
drop(w_optimistically_confirmed_bank);
pending_optimistically_confirmed_banks.retain(|&(slot, _hash)| slot > root_slot);
}
BankNotification::NewRootedChain(mut roots, oldest_parent) => {
Self::notify_new_root_slots(
&mut roots,
oldest_parent,
newest_root_slot,
slot_notification_subscribers,
);
}
}
}
pub fn close(self) -> thread::Result<()> {
self.join()
}
pub fn join(self) -> thread::Result<()> {
self.thread_hdl.join()
}
}
#[cfg(test)]
mod tests {
use {
super::*,
crossbeam_channel::bounded,
solana_ledger::genesis_utils::{GenesisConfigInfo, create_genesis_config},
solana_runtime::{bank::SlotLeader, commitment::BlockCommitmentCache, dependency_tracker},
std::sync::atomic::AtomicU64,
};
fn get_root_notifications(receiver: &Receiver<SlotNotification>) -> Vec<SlotNotification> {
let mut notifications = Vec::new();
while let Ok(notification) = receiver.recv_timeout(Duration::from_millis(100)) {
notifications.push(notification);
}
notifications
}
fn root_slot_notifications(root_bank: &Arc<Bank>) -> (Vec<(Slot, BankId)>, Slot) {
let mut rooted_banks = root_bank.parents();
let oldest_parent = rooted_banks
.last()
.map(|last| last.parent_slot())
.unwrap_or_else(|| root_bank.parent_slot());
rooted_banks.push(root_bank.clone());
(
rooted_banks
.iter()
.map(|bank| (bank.slot(), bank.bank_id()))
.collect(),
oldest_parent,
)
}
#[test]
fn test_process_notification() {
let exit = Arc::new(AtomicBool::new(false));
let GenesisConfigInfo { genesis_config, .. } = create_genesis_config(100);
let bank = Bank::new_for_tests(&genesis_config);
let bank_forks = BankForks::new_rw_arc(bank);
let bank0 = bank_forks.read().unwrap().get(0).unwrap();
let bank1 = Bank::new_from_parent(bank0, SlotLeader::default(), 1);
bank_forks.write().unwrap().insert(bank1);
let bank1 = bank_forks.read().unwrap().get(1).unwrap();
let bank2 = Bank::new_from_parent(bank1, SlotLeader::default(), 2);
bank_forks.write().unwrap().insert(bank2);
let bank2 = bank_forks.read().unwrap().get(2).unwrap();
let bank3 = Bank::new_from_parent(bank2, SlotLeader::default(), 3);
bank_forks.write().unwrap().insert(bank3);
let bank1_hash = bank_forks.read().unwrap().get(1).unwrap().hash();
let bank2 = bank_forks.read().unwrap().get(2).unwrap();
let bank2_hash = bank2.hash();
let bank3_pending_hash = Hash::new_unique();
let optimistically_confirmed_bank: Arc<RwLock<OptimisticallyConfirmedBank>> =
OptimisticallyConfirmedBank::locked_from_bank_forks_root(&bank_forks);
let block_commitment_cache = Arc::new(RwLock::new(BlockCommitmentCache::default()));
let max_complete_transaction_status_slot = Arc::new(AtomicU64::default());
let subscriptions = Arc::new(RpcSubscriptions::new_for_tests(
exit,
max_complete_transaction_status_slot,
bank_forks.clone(),
block_commitment_cache,
optimistically_confirmed_bank.clone(),
));
let mut pending_optimistically_confirmed_banks = PendingOptimisticallyConfirmedBanks::new();
assert_eq!(optimistically_confirmed_bank.read().unwrap().bank.slot(), 0);
let mut highest_confirmed_slot: Slot = 0;
let mut newest_root_slot: Slot = 0;
let mut last_notified_confirmed_slot: Slot = 0;
OptimisticallyConfirmedBankTracker::process_notification(
(
BankNotification::OptimisticallyConfirmed(2, bank2_hash),
None,
),
&bank_forks,
&optimistically_confirmed_bank,
&subscriptions,
&mut pending_optimistically_confirmed_banks,
&mut last_notified_confirmed_slot,
&mut highest_confirmed_slot,
&mut newest_root_slot,
&None,
None,
&None, );
assert_eq!(optimistically_confirmed_bank.read().unwrap().bank.slot(), 2);
assert_eq!(highest_confirmed_slot, 2);
OptimisticallyConfirmedBankTracker::process_notification(
(
BankNotification::OptimisticallyConfirmed(1, bank1_hash),
None,
),
&bank_forks,
&optimistically_confirmed_bank,
&subscriptions,
&mut pending_optimistically_confirmed_banks,
&mut last_notified_confirmed_slot,
&mut highest_confirmed_slot,
&mut newest_root_slot,
&None,
None,
&None, );
assert_eq!(optimistically_confirmed_bank.read().unwrap().bank.slot(), 2);
assert_eq!(highest_confirmed_slot, 2);
OptimisticallyConfirmedBankTracker::process_notification(
(
BankNotification::OptimisticallyConfirmed(3, bank3_pending_hash),
None,
),
&bank_forks,
&optimistically_confirmed_bank,
&subscriptions,
&mut pending_optimistically_confirmed_banks,
&mut last_notified_confirmed_slot,
&mut highest_confirmed_slot,
&mut newest_root_slot,
&None,
None,
&None, );
assert_eq!(optimistically_confirmed_bank.read().unwrap().bank.slot(), 2);
assert_eq!(pending_optimistically_confirmed_banks.len(), 1);
assert!(pending_optimistically_confirmed_banks.contains(&(3, bank3_pending_hash)));
assert_eq!(highest_confirmed_slot, 2);
let bank3 = bank_forks.read().unwrap().get(3).unwrap();
bank3.freeze();
assert!(pending_optimistically_confirmed_banks.remove(&(3, bank3_pending_hash)));
pending_optimistically_confirmed_banks.insert((3, bank3.hash()));
OptimisticallyConfirmedBankTracker::process_notification(
(
BankNotification::Frozen(bank3),
None,
),
&bank_forks,
&optimistically_confirmed_bank,
&subscriptions,
&mut pending_optimistically_confirmed_banks,
&mut last_notified_confirmed_slot,
&mut highest_confirmed_slot,
&mut newest_root_slot,
&None,
None,
&None, );
assert_eq!(optimistically_confirmed_bank.read().unwrap().bank.slot(), 3);
assert_eq!(highest_confirmed_slot, 3);
assert_eq!(pending_optimistically_confirmed_banks.len(), 0);
let bank3 = bank_forks.read().unwrap().get(3).unwrap();
let bank4 = Bank::new_from_parent(bank3, SlotLeader::default(), 4);
bank_forks.write().unwrap().insert(bank4);
let bank4_hash = Hash::new_unique();
OptimisticallyConfirmedBankTracker::process_notification(
(
BankNotification::OptimisticallyConfirmed(4, bank4_hash),
None,
),
&bank_forks,
&optimistically_confirmed_bank,
&subscriptions,
&mut pending_optimistically_confirmed_banks,
&mut last_notified_confirmed_slot,
&mut highest_confirmed_slot,
&mut newest_root_slot,
&None,
None,
&None, );
assert_eq!(optimistically_confirmed_bank.read().unwrap().bank.slot(), 3);
assert_eq!(pending_optimistically_confirmed_banks.len(), 1);
assert!(pending_optimistically_confirmed_banks.contains(&(4, bank4_hash)));
assert_eq!(highest_confirmed_slot, 3);
let bank4 = bank_forks.read().unwrap().get(4).unwrap();
let bank5 = Bank::new_from_parent(bank4, SlotLeader::default(), 5);
bank_forks.write().unwrap().insert(bank5);
let bank5 = bank_forks.read().unwrap().get(5).unwrap();
let mut bank_notification_senders = Vec::new();
let (sender, receiver) = bounded(1024);
bank_notification_senders.push(sender);
let subscribers = Some(Arc::new(RwLock::new(bank_notification_senders)));
let (parent_roots, oldest_parent) = root_slot_notifications(&bank5);
OptimisticallyConfirmedBankTracker::process_notification(
(
BankNotification::NewRootBank(bank5),
None,
),
&bank_forks,
&optimistically_confirmed_bank,
&subscriptions,
&mut pending_optimistically_confirmed_banks,
&mut last_notified_confirmed_slot,
&mut highest_confirmed_slot,
&mut newest_root_slot,
&subscribers,
None,
&None, );
assert_eq!(optimistically_confirmed_bank.read().unwrap().bank.slot(), 5);
assert_eq!(pending_optimistically_confirmed_banks.len(), 0);
assert!(!pending_optimistically_confirmed_banks.contains(&(4, bank4_hash)));
assert_eq!(highest_confirmed_slot, 3);
assert_eq!(newest_root_slot, 0);
OptimisticallyConfirmedBankTracker::process_notification(
(
BankNotification::NewRootedChain(parent_roots, oldest_parent),
None,
),
&bank_forks,
&optimistically_confirmed_bank,
&subscriptions,
&mut pending_optimistically_confirmed_banks,
&mut last_notified_confirmed_slot,
&mut highest_confirmed_slot,
&mut newest_root_slot,
&subscribers,
None,
&None, );
assert_eq!(newest_root_slot, 5);
let notifications = get_root_notifications(&receiver);
assert_eq!(notifications.len(), 5);
let bank5 = bank_forks.read().unwrap().get(5).unwrap();
let bank6 = Bank::new_from_parent(bank5, SlotLeader::default(), 6);
bank_forks.write().unwrap().insert(bank6);
let bank5 = bank_forks.read().unwrap().get(5).unwrap();
let bank7 = Bank::new_from_parent(bank5, SlotLeader::default(), 7);
bank_forks.write().unwrap().insert(bank7);
bank_forks.write().unwrap().set_root(7, None, None);
let bank6_hash = Hash::new_unique();
OptimisticallyConfirmedBankTracker::process_notification(
(
BankNotification::OptimisticallyConfirmed(6, bank6_hash),
None,
),
&bank_forks,
&optimistically_confirmed_bank,
&subscriptions,
&mut pending_optimistically_confirmed_banks,
&mut last_notified_confirmed_slot,
&mut highest_confirmed_slot,
&mut newest_root_slot,
&None,
None,
&None, );
assert_eq!(optimistically_confirmed_bank.read().unwrap().bank.slot(), 5);
assert_eq!(pending_optimistically_confirmed_banks.len(), 0);
assert!(!pending_optimistically_confirmed_banks.contains(&(6, bank6_hash)));
assert_eq!(highest_confirmed_slot, 3);
assert_eq!(newest_root_slot, 5);
let bank7 = bank_forks.read().unwrap().get(7).unwrap();
let (parent_roots, oldest_parent) = root_slot_notifications(&bank7);
OptimisticallyConfirmedBankTracker::process_notification(
(
BankNotification::NewRootBank(bank7),
None,
),
&bank_forks,
&optimistically_confirmed_bank,
&subscriptions,
&mut pending_optimistically_confirmed_banks,
&mut last_notified_confirmed_slot,
&mut highest_confirmed_slot,
&mut newest_root_slot,
&subscribers,
None,
&None, );
assert_eq!(optimistically_confirmed_bank.read().unwrap().bank.slot(), 7);
assert_eq!(pending_optimistically_confirmed_banks.len(), 0);
assert!(!pending_optimistically_confirmed_banks.contains(&(6, bank6_hash)));
assert_eq!(highest_confirmed_slot, 3);
assert_eq!(newest_root_slot, 5);
OptimisticallyConfirmedBankTracker::process_notification(
(
BankNotification::NewRootedChain(parent_roots, oldest_parent),
None,
),
&bank_forks,
&optimistically_confirmed_bank,
&subscriptions,
&mut pending_optimistically_confirmed_banks,
&mut last_notified_confirmed_slot,
&mut highest_confirmed_slot,
&mut newest_root_slot,
&subscribers,
None,
&None, );
assert_eq!(newest_root_slot, 7);
let notifications = get_root_notifications(&receiver);
assert_eq!(notifications.len(), 1);
}
#[test]
fn test_event_synchronization() {
let exit = Arc::new(AtomicBool::new(false));
let dependency_tracker: Arc<DependencyTracker> =
Arc::new(dependency_tracker::DependencyTracker::default());
let work_id_1 = 345;
let work_id_2 = 678;
let tracker_clone = dependency_tracker.clone();
let handle = thread::spawn(move || {
let GenesisConfigInfo { genesis_config, .. } = create_genesis_config(100);
let bank = Bank::new_for_tests(&genesis_config);
let bank_forks = BankForks::new_rw_arc(bank);
let bank0 = bank_forks.read().unwrap().get(0).unwrap();
let bank1 = Bank::new_from_parent(bank0, SlotLeader::default(), 1);
bank_forks.write().unwrap().insert(bank1);
let bank1_pending_hash = Hash::new_unique();
let mut pending_optimistically_confirmed_banks =
PendingOptimisticallyConfirmedBanks::new();
let max_complete_transaction_status_slot = Arc::new(AtomicU64::default());
let block_commitment_cache = Arc::new(RwLock::new(BlockCommitmentCache::default()));
let mut highest_confirmed_slot: Slot = 0;
let mut newest_root_slot: Slot = 0;
let mut last_notified_confirmed_slot: Slot = 0;
let optimistically_confirmed_bank: Arc<RwLock<OptimisticallyConfirmedBank>> =
OptimisticallyConfirmedBank::locked_from_bank_forks_root(&bank_forks);
let subscriptions = Arc::new(RpcSubscriptions::new_for_tests(
exit,
max_complete_transaction_status_slot,
bank_forks.clone(),
block_commitment_cache,
optimistically_confirmed_bank.clone(),
));
OptimisticallyConfirmedBankTracker::process_notification(
(
BankNotification::OptimisticallyConfirmed(1, bank1_pending_hash),
Some(work_id_1),
),
&bank_forks,
&optimistically_confirmed_bank,
&subscriptions,
&mut pending_optimistically_confirmed_banks,
&mut last_notified_confirmed_slot,
&mut highest_confirmed_slot,
&mut newest_root_slot,
&None,
None,
&Some(tracker_clone.clone()),
);
assert_eq!(optimistically_confirmed_bank.read().unwrap().bank.slot(), 0);
assert_eq!(highest_confirmed_slot, 0);
assert_eq!(pending_optimistically_confirmed_banks.len(), 1);
assert!(pending_optimistically_confirmed_banks.contains(&(1, bank1_pending_hash)));
let bank1 = bank_forks.read().unwrap().get(1).unwrap();
bank1.freeze();
assert!(pending_optimistically_confirmed_banks.remove(&(1, bank1_pending_hash)));
pending_optimistically_confirmed_banks.insert((1, bank1.hash()));
OptimisticallyConfirmedBankTracker::process_notification(
(
BankNotification::Frozen(bank1),
Some(work_id_2),
),
&bank_forks,
&optimistically_confirmed_bank,
&subscriptions,
&mut pending_optimistically_confirmed_banks,
&mut last_notified_confirmed_slot,
&mut highest_confirmed_slot,
&mut newest_root_slot,
&None,
None,
&Some(tracker_clone),
);
assert_eq!(optimistically_confirmed_bank.read().unwrap().bank.slot(), 1);
assert_eq!(highest_confirmed_slot, 1);
assert_eq!(pending_optimistically_confirmed_banks.len(), 0);
});
dependency_tracker.mark_this_and_all_previous_work_processed(work_id_1);
dependency_tracker.mark_this_and_all_previous_work_processed(work_id_2);
handle.join().unwrap();
}
}