415 lines
18 KiB
Rust
415 lines
18 KiB
Rust
// file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/admission.rs
|
|
// version: 6
|
|
|
|
fn material(
|
|
network: &str,
|
|
signature_byte: u8,
|
|
transaction_data: &str,
|
|
) -> std::option::Option<(ksp_store_lib::RawNetworkId, ksp_raw_transaction_lib::RawTransactionMaterial)> {
|
|
let network = match ksp_store_lib::RawNetworkId::new(network) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
let signature = ksp_store_lib::RawTransactionSignature::new([signature_byte; 64]);
|
|
let material = ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64(
|
|
network.clone(),
|
|
signature,
|
|
42,
|
|
std::option::Option::Some(1_700_000_000),
|
|
transaction_data,
|
|
ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
|
ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
|
ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
|
);
|
|
return std::option::Option::Some((network, material));
|
|
}
|
|
|
|
fn provenance() -> std::option::Option<ksp_store_lib::RawAcquisitionProvenance> {
|
|
let provider = match ksp_store_lib::RawProvenanceCode::new("deterministic-harness") {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
let protocol = match ksp_store_lib::RawProvenanceCode::new("internal") {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
let method = match ksp_store_lib::RawProvenanceCode::new("admission") {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
let received_at = match ksp_store_lib::RawTimestamp::from_unix_millis(1_700_000_000_000) {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
return std::option::Option::Some(ksp_store_lib::RawAcquisitionProvenance::new(
|
|
provider,
|
|
protocol,
|
|
method,
|
|
ksp_store_lib::RawAcquisitionOrigin::Live,
|
|
received_at,
|
|
));
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_006_bounded_admission_applies_async_backpressure_until_one_slot_is_consumed() {
|
|
let (network, first_material) = match material("mainnet", 1, "AQID") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let (_, second_material) = match material("mainnet", 2, "BAUG") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let first_provenance = match provenance() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let second_provenance = match provenance() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let (mut admission, sender) = crate::RawTransactionAdmission::new(1);
|
|
let first_ingress = crate::RawTransactionIngress {
|
|
material: first_material,
|
|
network: network.clone(),
|
|
provenance: first_provenance,
|
|
source_key: [1; 32],
|
|
};
|
|
assert!(sender.send(first_ingress).await.is_ok());
|
|
let second_sender = sender.clone();
|
|
let second_network = network.clone();
|
|
let completed = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false));
|
|
let completed_task = std::sync::Arc::clone(&completed);
|
|
let second_task = tokio::spawn(async move {
|
|
let second_ingress = crate::RawTransactionIngress {
|
|
material: second_material,
|
|
network: second_network,
|
|
provenance: second_provenance,
|
|
source_key: [2; 32],
|
|
};
|
|
let admitted = second_sender.send(second_ingress).await.is_ok();
|
|
completed_task.store(admitted, std::sync::atomic::Ordering::Release);
|
|
return admitted;
|
|
});
|
|
tokio::task::yield_now().await;
|
|
assert!(!completed.load(std::sync::atomic::Ordering::Acquire));
|
|
let first = admission.receive(&network).await;
|
|
assert!(matches!(first, std::result::Result::Ok(std::option::Option::Some(_))));
|
|
for _ in 0..64 {
|
|
if completed.load(std::sync::atomic::Ordering::Acquire) {
|
|
break;
|
|
}
|
|
tokio::task::yield_now().await;
|
|
}
|
|
assert!(completed.load(std::sync::atomic::Ordering::Acquire));
|
|
let joined = second_task.await;
|
|
assert!(matches!(joined, std::result::Result::Ok(true)));
|
|
let second = admission.receive(&network).await;
|
|
assert!(matches!(second, std::result::Result::Ok(std::option::Option::Some(_))));
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_006_stop_preempts_a_blocked_sender_without_silent_post_stop_admission() {
|
|
let (network, first_material) = match material("mainnet", 3, "BwgJ") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let (_, second_material) = match material("mainnet", 4, "CgsM") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let first_provenance = match provenance() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let second_provenance = match provenance() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let (stop_sender, stop_receiver) = tokio::sync::watch::channel(false);
|
|
let (mut admission, sender) = crate::RawTransactionAdmission::new(1);
|
|
let first_ingress = crate::RawTransactionIngress {
|
|
material: first_material,
|
|
network: network.clone(),
|
|
provenance: first_provenance,
|
|
source_key: [3; 32],
|
|
};
|
|
assert!(sender.send(first_ingress).await.is_ok());
|
|
let blocked_sender = sender.clone();
|
|
let blocked_network = network.clone();
|
|
let blocked_task = tokio::spawn(async move {
|
|
let mut blocked_stop = stop_receiver;
|
|
let blocked_ingress = crate::RawTransactionIngress {
|
|
material: second_material,
|
|
network: blocked_network,
|
|
provenance: second_provenance,
|
|
source_key: [4; 32],
|
|
};
|
|
return tokio::select! {
|
|
biased;
|
|
_changed = blocked_stop.changed() => false,
|
|
result = blocked_sender.send(blocked_ingress) => result.is_ok(),
|
|
};
|
|
});
|
|
tokio::task::yield_now().await;
|
|
stop_sender.send_replace(true);
|
|
let blocked = blocked_task.await;
|
|
assert!(matches!(blocked, std::result::Result::Ok(false)));
|
|
let first = admission.receive(&network).await;
|
|
assert!(matches!(first, std::result::Result::Ok(std::option::Option::Some(_))));
|
|
admission.close();
|
|
let end = admission.receive(&network).await;
|
|
assert!(matches!(end, std::result::Result::Ok(std::option::Option::None)));
|
|
return;
|
|
}
|
|
|
|
#[test]
|
|
fn pre_006_observation_key_domain_and_common_raw_assembly_golden_are_exact() {
|
|
let (network, material) = match material("mainnet", 7, "AQID") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let provenance = match provenance() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let ingress = crate::RawTransactionIngress { material, network: network.clone(), provenance: provenance.clone(), source_key: [9; 32] };
|
|
let acquisition = super::canonicalize_ingress(&network, ingress);
|
|
assert!(acquisition.is_ok());
|
|
let acquisition = match acquisition {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return,
|
|
};
|
|
assert_eq!(acquisition.transaction().reference().network(), &network);
|
|
assert_eq!(acquisition.transaction().reference().signature().as_bytes(), &[7; 64]);
|
|
assert_eq!(acquisition.transaction().slot(), 42);
|
|
assert_eq!(acquisition.transaction().payload().format_id().as_str(), ksp_raw_transaction_lib::RAW_TRANSACTION_FORMAT_ID);
|
|
assert_eq!(acquisition.transaction().payload().format_version(), ksp_raw_transaction_lib::RAW_TRANSACTION_FORMAT_VERSION);
|
|
assert_eq!(acquisition.transaction().payload().bytes(), br#"{"transaction":["AQID","base64"]}"#);
|
|
assert_eq!(acquisition.observation().provenance(), &provenance);
|
|
assert_eq!(acquisition.observation().transaction(), acquisition.transaction().reference());
|
|
assert_eq!(
|
|
acquisition.observation().observation_key().as_bytes(),
|
|
&[
|
|
248, 209, 97, 220, 168, 168, 167, 42, 13, 161, 16, 252, 135, 136, 46, 231, 19, 104, 72, 70, 28, 106, 163, 43, 9, 227, 10, 97, 136, 185, 93, 204,
|
|
],
|
|
);
|
|
return;
|
|
}
|
|
|
|
#[test]
|
|
fn pre_006_network_guards_reject_ingress_or_material_mismatch_without_echoing_values() {
|
|
let (mainnet, mainnet_material) = match material("mainnet", 8, "DQ4P") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let (devnet, devnet_material) = match material("devnet", 9, "EBES") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let first_provenance = match provenance() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let second_provenance = match provenance() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let wrong_ingress = crate::RawTransactionIngress { material: mainnet_material, network: devnet, provenance: first_provenance, source_key: [10; 32] };
|
|
let wrong_ingress = super::canonicalize_ingress(&mainnet, wrong_ingress);
|
|
assert!(wrong_ingress.is_err());
|
|
let wrong_material = crate::RawTransactionIngress {
|
|
material: devnet_material,
|
|
network: mainnet.clone(),
|
|
provenance: second_provenance,
|
|
source_key: [11; 32],
|
|
};
|
|
let wrong_material = super::canonicalize_ingress(&mainnet, wrong_material);
|
|
assert!(wrong_material.is_err());
|
|
for result in [wrong_ingress, wrong_material] {
|
|
let error = match result {
|
|
std::result::Result::Ok(_) => continue,
|
|
std::result::Result::Err(value) => value,
|
|
};
|
|
assert_eq!(error.code(), crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID);
|
|
let rendered = std::format!("{error:?}");
|
|
assert!(!rendered.contains("mainnet"));
|
|
assert!(!rendered.contains("devnet"));
|
|
}
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_009_full_queue_dequeue_marks_source_neutral_backpressure_observation() {
|
|
let (network, material) = match material("mainnet", 13, "AQID") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let provenance = match provenance() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let (mut admission, sender) = crate::RawTransactionAdmission::new(1);
|
|
let ingress = crate::RawTransactionIngress { material, network: network.clone(), provenance, source_key: [13; 32] };
|
|
assert!(sender.send(ingress).await.is_ok());
|
|
let first = admission.receive(&network).await;
|
|
assert!(matches!(first, std::result::Result::Ok(std::option::Option::Some(_))));
|
|
assert!(admission.take_backpressure_wait_observed());
|
|
admission.close();
|
|
let second = admission.receive(&network).await;
|
|
assert!(matches!(second, std::result::Result::Ok(std::option::Option::None)));
|
|
assert!(!admission.take_backpressure_wait_observed());
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn v0_3_13_pre_008_distinct_source_keys_produce_distinct_idempotent_observation_keys() {
|
|
let (network, first_material) = match material("mainnet", 51, "AQID") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let (_, second_material) = match material("mainnet", 51, "AQID") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let first_provenance = match provenance() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let second_provenance = match provenance() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let (mut admission, sender) = crate::RawTransactionAdmission::new(2);
|
|
let first_ingress = crate::RawTransactionIngress {
|
|
material: first_material,
|
|
network: network.clone(),
|
|
provenance: first_provenance,
|
|
source_key: [11_u8; 32],
|
|
};
|
|
let second_ingress = crate::RawTransactionIngress {
|
|
material: second_material,
|
|
network: network.clone(),
|
|
provenance: second_provenance,
|
|
source_key: [12_u8; 32],
|
|
};
|
|
if sender.send(first_ingress).await.is_err() || sender.send(second_ingress).await.is_err() {
|
|
return;
|
|
}
|
|
std::mem::drop(sender);
|
|
let first = match admission.receive(&network).await {
|
|
std::result::Result::Ok(std::option::Option::Some(value)) => value,
|
|
_ => return,
|
|
};
|
|
let second = match admission.receive(&network).await {
|
|
std::result::Result::Ok(std::option::Option::Some(value)) => value,
|
|
_ => return,
|
|
};
|
|
assert_eq!(first.transaction().reference(), second.transaction().reference());
|
|
assert_eq!(first.transaction().payload().content_hash(), second.transaction().payload().content_hash());
|
|
assert_ne!(first.observation().observation_key(), second.observation().observation_key());
|
|
return;
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn v0_3_13_pre_009_duplicate_storm_cannot_starve_second_ready_source_on_bounded_admission() {
|
|
let (network, initial_material) = match material("mainnet", 60, "AQID") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let (_, peer_material) = match material("mainnet", 99, "BwgJ") {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let initial_provenance = match provenance() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let storm_provenance = match provenance() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let peer_provenance = match provenance() {
|
|
std::option::Option::Some(value) => value,
|
|
std::option::Option::None => return,
|
|
};
|
|
let (mut admission, sender) = crate::RawTransactionAdmission::new(1);
|
|
assert!(
|
|
sender
|
|
.send(crate::RawTransactionIngress {
|
|
material: initial_material,
|
|
network: network.clone(),
|
|
provenance: initial_provenance,
|
|
source_key: [60_u8; 32],
|
|
})
|
|
.await
|
|
.is_ok()
|
|
);
|
|
let storm_sender = sender.clone();
|
|
let storm_network = network.clone();
|
|
let storm_task = tokio::spawn(async move {
|
|
for _ in 0..8 {
|
|
let material = ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64(
|
|
storm_network.clone(),
|
|
ksp_store_lib::RawTransactionSignature::new([61_u8; 64]),
|
|
42,
|
|
std::option::Option::Some(1_700_000_000),
|
|
"BAUG",
|
|
ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
|
ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
|
ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
|
);
|
|
if storm_sender
|
|
.send(crate::RawTransactionIngress {
|
|
material,
|
|
network: storm_network.clone(),
|
|
provenance: storm_provenance.clone(),
|
|
source_key: [61_u8; 32],
|
|
})
|
|
.await
|
|
.is_err()
|
|
{
|
|
return false;
|
|
}
|
|
}
|
|
return true;
|
|
});
|
|
tokio::task::yield_now().await;
|
|
let peer_sender = sender.clone();
|
|
let peer_network = network.clone();
|
|
let peer_task = tokio::spawn(async move {
|
|
return peer_sender
|
|
.send(crate::RawTransactionIngress { material: peer_material, network: peer_network, provenance: peer_provenance, source_key: [99_u8; 32] })
|
|
.await
|
|
.is_ok();
|
|
});
|
|
tokio::task::yield_now().await;
|
|
assert_eq!(admission.queue_depth(), 1);
|
|
let first = match admission.receive(&network).await {
|
|
std::result::Result::Ok(std::option::Option::Some(value)) => value,
|
|
_ => return,
|
|
};
|
|
assert_eq!(first.transaction().reference().signature().as_bytes(), &[60_u8; 64]);
|
|
let second = match admission.receive(&network).await {
|
|
std::result::Result::Ok(std::option::Option::Some(value)) => value,
|
|
_ => return,
|
|
};
|
|
assert_eq!(second.transaction().reference().signature().as_bytes(), &[61_u8; 64]);
|
|
let third = match admission.receive(&network).await {
|
|
std::result::Result::Ok(std::option::Option::Some(value)) => value,
|
|
_ => return,
|
|
};
|
|
assert_eq!(third.transaction().reference().signature().as_bytes(), &[99_u8; 64]);
|
|
for _ in 0..7 {
|
|
let next = admission.receive(&network).await;
|
|
assert!(matches!(next, std::result::Result::Ok(std::option::Option::Some(_))));
|
|
assert!(admission.queue_depth() <= 1);
|
|
}
|
|
let peer_joined = peer_task.await;
|
|
assert!(matches!(peer_joined, std::result::Result::Ok(true)));
|
|
let storm_joined = storm_task.await;
|
|
assert!(matches!(storm_joined, std::result::Result::Ok(true)));
|
|
return;
|
|
}
|