#![forbid(unsafe_code)]
use crate::{
Error, EventNotificationListener, EventNotificationSender, EventSubscriptionService,
ReconfigNotificationListener,
};
use aptos_infallible::RwLock;
use aptos_types::{
account_address::AccountAddress,
contract_event::ContractEvent,
event::EventKey,
on_chain_config,
on_chain_config::{OnChainConfig, ON_CHAIN_CONFIG_REGISTRY},
transaction::{Transaction, Version, WriteSetPayload},
};
use aptos_vm::AptosVM;
use aptosdb::AptosDB;
use claim::{assert_lt, assert_matches, assert_ok};
use executor_test_helpers::bootstrap_genesis;
use futures::{FutureExt, StreamExt};
use move_deps::move_core_types::language_storage::TypeTag;
use serde::{Deserialize, Serialize};
use std::{convert::TryInto, sync::Arc};
use storage_interface::DbReaderWriter;
#[test]
fn test_all_configs_returned() {
let mut event_service = create_event_subscription_service();
let mut listener_1 = event_service.subscribe_to_reconfigurations().unwrap();
let mut listener_2 = event_service.subscribe_to_reconfigurations().unwrap();
let mut listener_3 = event_service.subscribe_to_reconfigurations().unwrap();
let version = 0;
let epoch = 1;
let num_reconfigs = 10;
for _ in 0..num_reconfigs {
notify_initial_configs(&mut event_service, version);
verify_reconfig_notifications_received(
vec![&mut listener_1, &mut listener_2, &mut listener_3],
version,
epoch,
);
}
let reconfig_event = create_test_event(on_chain_config::new_epoch_event_key());
for _ in 0..num_reconfigs {
notify_events(&mut event_service, version, vec![reconfig_event.clone()]);
verify_reconfig_notifications_received(
vec![&mut listener_1, &mut listener_2, &mut listener_3],
version,
epoch,
);
}
}
#[test]
fn test_reconfig_notification_no_queuing() {
let mut event_service = create_event_subscription_service();
let mut listener_1 = event_service.subscribe_to_reconfigurations().unwrap();
let mut listener_2 = event_service.subscribe_to_reconfigurations().unwrap();
let reconfig_event = create_test_event(on_chain_config::new_epoch_event_key());
let num_reconfigs = 10;
for _ in 0..num_reconfigs {
notify_events(&mut event_service, 0, vec![reconfig_event.clone()]);
}
let notification_count = count_reconfig_notifications(&mut listener_1);
assert_eq!(notification_count, 1);
let num_reconfigs = 5;
for _ in 0..num_reconfigs {
notify_initial_configs(&mut event_service, 0);
}
let notification_count = count_reconfig_notifications(&mut listener_1);
assert_eq!(notification_count, 1);
let notification_count = count_reconfig_notifications(&mut listener_2);
assert_eq!(notification_count, 1);
}
#[test]
fn test_dynamic_subscribers() {
let mut event_service = create_event_subscription_service();
let event_key_1 = create_random_event_key();
let reconfig_event_key = on_chain_config::new_epoch_event_key();
let mut event_listener_1 = event_service
.subscribe_to_events(vec![event_key_1])
.unwrap();
let mut reconfig_listener_1 = event_service.subscribe_to_reconfigurations().unwrap();
let event_1 = create_test_event(event_key_1);
notify_events(&mut event_service, 0, vec![event_1.clone()]);
verify_event_notification_received(vec![&mut event_listener_1], 0, vec![event_1.clone()]);
let mut event_listener_2 = event_service
.subscribe_to_events(vec![event_key_1, reconfig_event_key])
.unwrap();
let reconfig_event = create_test_event(reconfig_event_key);
notify_events(
&mut event_service,
0,
vec![event_1.clone(), reconfig_event.clone()],
);
verify_event_notification_received(vec![&mut event_listener_1], 0, vec![event_1.clone()]);
verify_event_notification_received(
vec![&mut event_listener_2],
0,
vec![event_1, reconfig_event.clone()],
);
verify_reconfig_notifications_received(vec![&mut reconfig_listener_1], 0, 1);
let mut reconfig_listener_2 = event_service.subscribe_to_reconfigurations().unwrap();
notify_events(&mut event_service, 0, vec![reconfig_event.clone()]);
verify_event_notification_received(vec![&mut event_listener_2], 0, vec![reconfig_event]);
verify_reconfig_notifications_received(
vec![&mut reconfig_listener_1, &mut reconfig_listener_2],
0,
1,
);
verify_no_event_notifications(vec![&mut event_listener_1]);
}
#[test]
fn test_event_and_reconfig_subscribers() {
let mut event_service = create_event_subscription_service();
let event_key_1 = create_random_event_key();
let event_key_2 = create_random_event_key();
let reconfig_event_key = on_chain_config::new_epoch_event_key();
let mut event_listener_1 = event_service
.subscribe_to_events(vec![event_key_1])
.unwrap();
let mut event_listener_2 = event_service
.subscribe_to_events(vec![event_key_1, event_key_2])
.unwrap();
let mut event_listener_3 = event_service
.subscribe_to_events(vec![reconfig_event_key])
.unwrap();
let mut reconfig_listener_1 = event_service.subscribe_to_reconfigurations().unwrap();
let mut reconfig_listener_2 = event_service.subscribe_to_reconfigurations().unwrap();
notify_initial_configs(&mut event_service, 0);
verify_reconfig_notifications_received(
vec![&mut reconfig_listener_1, &mut reconfig_listener_2],
0,
1,
);
verify_no_event_notifications(vec![
&mut event_listener_1,
&mut event_listener_2,
&mut event_listener_3,
]);
let reconfig_event = create_test_event(reconfig_event_key);
notify_events(&mut event_service, 0, vec![reconfig_event.clone()]);
verify_reconfig_notifications_received(
vec![&mut reconfig_listener_1, &mut reconfig_listener_2],
0,
1,
);
verify_event_notification_received(
vec![&mut event_listener_3],
0,
vec![reconfig_event.clone()],
);
verify_no_event_notifications(vec![&mut event_listener_1, &mut event_listener_2]);
let event_1 = create_test_event(event_key_1);
let event_2 = create_test_event(event_key_2);
notify_events(
&mut event_service,
0,
vec![event_1.clone(), event_2.clone(), reconfig_event.clone()],
);
verify_reconfig_notifications_received(
vec![&mut reconfig_listener_1, &mut reconfig_listener_2],
0,
1,
);
verify_event_notification_received(vec![&mut event_listener_1], 0, vec![event_1.clone()]);
verify_event_notification_received(
vec![&mut event_listener_2],
0,
vec![event_1.clone(), event_2],
);
verify_event_notification_received(vec![&mut event_listener_3], 0, vec![reconfig_event]);
notify_events(&mut event_service, 0, vec![event_1.clone()]);
verify_event_notification_received(
vec![&mut event_listener_1, &mut event_listener_2],
0,
vec![event_1],
);
verify_no_reconfig_notifications(vec![&mut reconfig_listener_1, &mut reconfig_listener_2]);
verify_no_event_notifications(vec![&mut event_listener_3]);
}
#[test]
fn test_event_notification_queuing() {
let mut event_service = create_event_subscription_service();
let event_key_1 = create_random_event_key();
let event_key_2 = on_chain_config::new_epoch_event_key();
let event_key_3 = create_random_event_key();
let mut listener_1 = event_service
.subscribe_to_events(vec![event_key_1])
.unwrap();
let mut listener_2 = event_service
.subscribe_to_events(vec![event_key_2])
.unwrap();
let num_events = 1000;
let event_1 = create_test_event(event_key_1);
for version in 0..num_events {
notify_events(&mut event_service, version, vec![event_1.clone()]);
}
let notification_count = count_event_notifications_and_ensure_ordering(&mut listener_1);
assert!(notification_count < num_events);
let num_events = 50;
let event_2 = create_test_event(event_key_1);
for version in 0..num_events {
notify_events(
&mut event_service,
version,
vec![event_2.clone(), event_1.clone()],
);
}
let notification_count = count_event_notifications_and_ensure_ordering(&mut listener_1);
assert_eq!(notification_count, num_events);
let num_events = 1000;
let event_3 = create_test_event(event_key_3);
for version in 0..num_events {
notify_events(&mut event_service, version, vec![event_3.clone()]);
}
verify_no_event_notifications(vec![&mut listener_1, &mut listener_2]);
}
#[test]
fn test_event_subscribers() {
let mut event_service = create_event_subscription_service();
let event_key_1 = create_random_event_key();
let event_key_2 = create_random_event_key();
let event_key_3 = on_chain_config::new_epoch_event_key();
let event_key_4 = create_random_event_key();
let event_key_5 = create_random_event_key();
let mut listener_1 = event_service
.subscribe_to_events(vec![event_key_1])
.unwrap();
let mut listener_2 = event_service
.subscribe_to_events(vec![event_key_1, event_key_2])
.unwrap();
let mut listener_3 = event_service
.subscribe_to_events(vec![event_key_2, event_key_3, event_key_4])
.unwrap();
let version = 99;
let event_1 = create_test_event(event_key_1);
notify_events(&mut event_service, version, vec![event_1.clone()]);
verify_event_notification_received(
vec![&mut listener_1, &mut listener_2],
version,
vec![event_1],
);
verify_no_event_notifications(vec![&mut listener_1, &mut listener_2, &mut listener_3]);
let version = 200;
let event_2 = create_test_event(event_key_2);
let event_3 = create_test_event(event_key_3);
notify_events(
&mut event_service,
version,
vec![event_2.clone(), event_3.clone()],
);
verify_no_event_notifications(vec![&mut listener_1]);
verify_event_notification_received(vec![&mut listener_2], version, vec![event_2.clone()]);
verify_event_notification_received(vec![&mut listener_3], version, vec![event_2, event_3]);
let event_5 = create_test_event(event_key_5);
notify_events(&mut event_service, version, vec![event_5]);
verify_no_event_notifications(vec![&mut listener_1]);
}
#[test]
fn test_no_events_no_subscribers() {
let mut event_service = create_event_subscription_service();
notify_events(&mut event_service, 1, vec![]);
notify_events(
&mut event_service,
1,
vec![create_test_event(create_random_event_key())],
);
assert_matches!(
event_service.subscribe_to_events(vec![]),
Err(Error::CannotSubscribeToZeroEventKeys)
);
let _event_listener = event_service.subscribe_to_events(vec![create_random_event_key()]);
let _reconfig_listener = event_service.subscribe_to_reconfigurations();
notify_events(&mut event_service, 1, vec![]);
}
#[test]
fn test_missing_configs() {
let mut config_registry = ON_CHAIN_CONFIG_REGISTRY.to_owned();
config_registry.push(TestOnChainConfig::CONFIG_ID);
let mut event_service = EventSubscriptionService::new(&config_registry, create_database());
let mut reconfig_listener = event_service.subscribe_to_reconfigurations().unwrap();
assert_ok!(event_service.notify_reconfiguration_subscribers(0));
if let Some(reconfig_notification) = reconfig_listener.select_next_some().now_or_never() {
let returned_configs = reconfig_notification.on_chain_configs.configs();
assert_eq!(
returned_configs.keys().len(),
ON_CHAIN_CONFIG_REGISTRY.len()
);
for config in ON_CHAIN_CONFIG_REGISTRY {
assert!(returned_configs.contains_key(config));
}
assert!(!returned_configs.contains_key(&TestOnChainConfig::CONFIG_ID));
} else {
panic!("Expected a reconfiguration notification but got None!");
}
}
#[derive(Clone, Debug, Deserialize, PartialEq, Eq, PartialOrd, Ord, Serialize)]
pub struct TestOnChainConfig {
pub some_value: u64,
}
impl OnChainConfig for TestOnChainConfig {
const MODULE_IDENTIFIER: &'static str = "test_on_chain_config";
const TYPE_IDENTIFIER: &'static str = "TestOnChainConfig";
}
fn count_event_notifications_and_ensure_ordering(listener: &mut EventNotificationListener) -> u64 {
let mut notification_received = true;
let mut notification_count = 0;
let mut last_version_received: i64 = -1;
while notification_received {
if let Some(event_notification) = listener.select_next_some().now_or_never() {
notification_count += 1;
assert_lt!(
last_version_received,
event_notification.version.try_into().unwrap()
);
last_version_received = event_notification.version.try_into().unwrap();
} else {
notification_received = false;
}
}
notification_count
}
fn count_reconfig_notifications(listener: &mut ReconfigNotificationListener) -> u64 {
let mut notification_received = true;
let mut notification_count = 0;
while notification_received {
if listener.select_next_some().now_or_never().is_some() {
notification_count += 1;
} else {
notification_received = false;
}
}
notification_count
}
fn verify_no_event_notifications(listeners: Vec<&mut EventNotificationListener>) {
for listener in listeners {
assert!(listener.select_next_some().now_or_never().is_none());
}
}
fn verify_no_reconfig_notifications(listeners: Vec<&mut ReconfigNotificationListener>) {
for listener in listeners {
assert!(listener.select_next_some().now_or_never().is_none());
}
}
fn verify_event_notification_received(
listeners: Vec<&mut EventNotificationListener>,
expected_version: Version,
expected_events: Vec<ContractEvent>,
) {
for listener in listeners {
if let Some(event_notification) = listener.select_next_some().now_or_never() {
assert_eq!(event_notification.version, expected_version);
assert_eq!(event_notification.subscribed_events, expected_events);
} else {
panic!("Expected an event notification but got None!");
}
}
}
fn verify_reconfig_notifications_received(
listeners: Vec<&mut ReconfigNotificationListener>,
expected_version: Version,
expected_epoch: u64,
) {
for listener in listeners {
if let Some(reconfig_notification) = listener.select_next_some().now_or_never() {
assert_eq!(reconfig_notification.version, expected_version);
assert_eq!(
reconfig_notification.on_chain_configs.epoch(),
expected_epoch
);
let returned_configs = reconfig_notification.on_chain_configs.configs();
for config in ON_CHAIN_CONFIG_REGISTRY {
assert!(returned_configs.contains_key(config));
}
} else {
panic!("Expected a reconfiguration notification but got None!");
}
}
}
fn notify_initial_configs(event_service: &mut EventSubscriptionService, version: Version) {
assert_ok!(event_service.notify_initial_configs(version));
}
fn notify_events(
event_service: &mut EventSubscriptionService,
version: Version,
events: Vec<ContractEvent>,
) {
assert_ok!(event_service.notify_events(version, events));
}
fn create_test_event(event_key: EventKey) -> ContractEvent {
ContractEvent::new(event_key, 0, TypeTag::Bool, bcs::to_bytes(&0).unwrap())
}
fn create_random_event_key() -> EventKey {
EventKey::new(0, AccountAddress::random())
}
fn create_event_subscription_service() -> EventSubscriptionService {
EventSubscriptionService::new(ON_CHAIN_CONFIG_REGISTRY, create_database())
}
fn create_database() -> Arc<RwLock<DbReaderWriter>> {
let (genesis, _) = vm_genesis::test_genesis_change_set_and_validators(Some(1));
let db_path = aptos_temppath::TempPath::new();
assert_ok!(db_path.create_as_dir());
let (_, db_rw) = DbReaderWriter::wrap(AptosDB::new_for_test(db_path.path()));
let genesis_txn = Transaction::GenesisTransaction(WriteSetPayload::Direct(genesis));
assert_ok!(bootstrap_genesis::<AptosVM>(&db_rw, &genesis_txn));
Arc::new(RwLock::new(db_rw))
}