v0.2.9-pre.007
This commit is contained in:
@@ -1,62 +0,0 @@
|
||||
// file: crates/ksp-onchain-transport-lib/unit_tests/client.rs
|
||||
// version: 3
|
||||
|
||||
fn endpoint(enabled: bool, url_text: &str) -> crate::HttpEndpointSettings {
|
||||
let url = crate::HttpEndpointUrl::parse(url_text).expect("test endpoint URL must parse");
|
||||
let role = crate::HttpEndpointRoleSettings::new(
|
||||
crate::HttpRoleName::new("default"),
|
||||
true,
|
||||
std::vec![crate::HttpRequestKind::wildcard()],
|
||||
10,
|
||||
crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None),
|
||||
);
|
||||
return crate::HttpEndpointSettings::new(
|
||||
"endpoint",
|
||||
enabled,
|
||||
crate::HttpProviderName::new("provider"),
|
||||
crate::HttpClusterName::new("devnet"),
|
||||
url,
|
||||
std::time::Duration::from_secs(1),
|
||||
std::time::Duration::from_secs(2),
|
||||
std::option::Option::Some(4),
|
||||
std::vec![role],
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn endpoint_client_snapshot_never_contains_url_or_secret_material() {
|
||||
let client = crate::HttpEndpointClient::new(endpoint(true, "https://provider.invalid/rpc?api-key=SECRET-CANARY")).expect("client must build");
|
||||
let snapshot = client.snapshot();
|
||||
let rendered = format!("{snapshot:?} {client:?}");
|
||||
assert_eq!(snapshot.availability(), crate::HttpEndpointAvailability::Available);
|
||||
assert!(!rendered.contains("SECRET-CANARY"));
|
||||
assert!(!rendered.contains("provider.invalid"));
|
||||
assert!(!rendered.contains("https://"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn disabled_endpoint_client_is_visible_but_not_selectable() {
|
||||
let client = crate::HttpEndpointClient::new(endpoint(false, "https://api.devnet.solana.com")).expect("disabled client must still build");
|
||||
assert_eq!(client.snapshot().availability(), crate::HttpEndpointAvailability::Disabled);
|
||||
assert!(!client.supports(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance")));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn endpoint_client_matches_exact_and_wildcard_capabilities() {
|
||||
let client = crate::HttpEndpointClient::new(endpoint(true, "https://api.devnet.solana.com")).expect("client must build");
|
||||
assert!(client.supports(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance")));
|
||||
assert!(!client.supports(&crate::HttpRoleName::new("write"), &crate::HttpRequestKind::new("get_balance")));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn endpoint_role_snapshot_exposes_safe_resilience_state() {
|
||||
let client = crate::HttpEndpointClient::new(endpoint(true, "https://api.devnet.solana.com")).expect("client must build");
|
||||
let snapshot = client.snapshot();
|
||||
let role = &snapshot.roles()[0];
|
||||
assert_eq!(role.availability(), crate::HttpEndpointAvailability::Available);
|
||||
assert_eq!(role.in_flight_requests(), std::option::Option::None);
|
||||
assert_eq!(role.cooldown_remaining(), std::option::Option::None);
|
||||
assert_eq!(role.success_count(), 0);
|
||||
assert_eq!(role.failure_count(), 0);
|
||||
assert_eq!(role.rate_limit_count(), 0);
|
||||
}
|
||||
@@ -1,141 +0,0 @@
|
||||
// file: crates/ksp-onchain-transport-lib/unit_tests/executor.rs
|
||||
// version: 2
|
||||
|
||||
fn pool_for_url(url: &str, request_timeout: std::time::Duration, max_retries: u32) -> crate::HttpTransportPool {
|
||||
let role = crate::HttpEndpointRoleSettings::new(
|
||||
crate::HttpRoleName::new("default"),
|
||||
true,
|
||||
std::vec![crate::HttpRequestKind::wildcard()],
|
||||
10,
|
||||
crate::HttpRoleLimits::new(
|
||||
std::option::Option::None,
|
||||
std::option::Option::None,
|
||||
std::option::Option::None,
|
||||
std::option::Option::Some(std::time::Duration::from_millis(1)),
|
||||
),
|
||||
);
|
||||
let endpoint = crate::HttpEndpointSettings::new(
|
||||
"fixture",
|
||||
true,
|
||||
crate::HttpProviderName::new("fixture"),
|
||||
crate::HttpClusterName::new("local"),
|
||||
crate::HttpEndpointUrl::parse(url).expect("fixture URL must parse"),
|
||||
std::time::Duration::from_millis(100),
|
||||
request_timeout,
|
||||
std::option::Option::Some(1),
|
||||
std::vec![role],
|
||||
);
|
||||
let settings = crate::HttpTransportSettings::new(
|
||||
std::vec![endpoint],
|
||||
crate::HttpRetrySettings::new(max_retries, std::time::Duration::from_millis(1), std::time::Duration::from_millis(2)),
|
||||
);
|
||||
return crate::HttpTransportPool::new(settings).expect("fixture pool must build");
|
||||
}
|
||||
|
||||
fn health_method() -> &'static crate::HttpRpcMethodDescriptor {
|
||||
return crate::find_http_rpc_method("getHealth").expect("getHealth descriptor must exist");
|
||||
}
|
||||
|
||||
fn serve_rate_limit_then_success() -> (std::string::String, std::thread::JoinHandle<usize>) {
|
||||
let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("fixture listener must bind");
|
||||
let address = listener.local_addr().expect("fixture listener address must resolve");
|
||||
let handle = std::thread::spawn(move || {
|
||||
let mut count = 0_usize;
|
||||
while count < 2 {
|
||||
let (mut stream, _) = listener.accept().expect("fixture server must accept request");
|
||||
let _ = read_request(&mut stream);
|
||||
let response = if count == 0 {
|
||||
"HTTP/1.1 429 Too Many Requests\r\nRetry-After: 0\r\nContent-Length: 0\r\nConnection: close\r\n\r\n".to_owned()
|
||||
} else {
|
||||
let body = include_str!("../fixtures/http/get_health.success.json");
|
||||
format!("HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", body.len(), body)
|
||||
};
|
||||
std::io::Write::write_all(&mut stream, response.as_bytes()).expect("fixture response must write");
|
||||
count = count.saturating_add(1);
|
||||
}
|
||||
return count;
|
||||
});
|
||||
return (format!("http://{address}"), handle);
|
||||
}
|
||||
|
||||
fn serve_timeout() -> (std::string::String, std::thread::JoinHandle<()>) {
|
||||
let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("fixture listener must bind");
|
||||
let address = listener.local_addr().expect("fixture listener address must resolve");
|
||||
let handle = std::thread::spawn(move || {
|
||||
let (mut stream, _) = listener.accept().expect("fixture server must accept request");
|
||||
let _ = read_request(&mut stream);
|
||||
std::thread::sleep(std::time::Duration::from_millis(100));
|
||||
return;
|
||||
});
|
||||
return (format!("http://{address}"), handle);
|
||||
}
|
||||
|
||||
fn read_request(stream: &mut std::net::TcpStream) -> std::string::String {
|
||||
let mut bytes = std::vec::Vec::new();
|
||||
let mut buffer = [0_u8; 1024];
|
||||
loop {
|
||||
let count = std::io::Read::read(stream, &mut buffer).expect("fixture request must read");
|
||||
if count == 0 {
|
||||
break;
|
||||
}
|
||||
bytes.extend_from_slice(&buffer[..count]);
|
||||
if request_complete(bytes.as_slice()) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
return std::string::String::from_utf8(bytes).expect("fixture request must be UTF-8");
|
||||
}
|
||||
|
||||
fn request_complete(bytes: &[u8]) -> bool {
|
||||
let text = match std::str::from_utf8(bytes) {
|
||||
std::result::Result::Ok(text) => text,
|
||||
std::result::Result::Err(_) => return false,
|
||||
};
|
||||
let header_end = match text.find("\r\n\r\n") {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return false,
|
||||
};
|
||||
let mut content_length = 0_usize;
|
||||
for line in text[..header_end].lines() {
|
||||
let (name, value) = match line.split_once(':') {
|
||||
std::option::Option::Some(parts) => parts,
|
||||
std::option::Option::None => continue,
|
||||
};
|
||||
if name.eq_ignore_ascii_case("content-length") {
|
||||
content_length = value.trim().parse::<usize>().expect("content length must parse");
|
||||
}
|
||||
}
|
||||
return bytes.len() >= header_end.saturating_add(4).saturating_add(content_length);
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn executor_applies_retry_after_and_retries_http_429_for_retry_safe_method() {
|
||||
let (url, handle) = serve_rate_limit_then_success();
|
||||
let pool = pool_for_url(url.as_str(), std::time::Duration::from_millis(500), 1);
|
||||
let result = pool
|
||||
.execute_standard_rpc(&crate::HttpRoleName::new("default"), health_method(), std::vec::Vec::new())
|
||||
.await
|
||||
.expect("retry-safe request must recover from one 429");
|
||||
assert_eq!(result, serde_json::json!("ok"));
|
||||
assert_eq!(handle.join().expect("fixture server must join"), 2);
|
||||
let snapshot = pool.snapshot();
|
||||
assert_eq!(snapshot.endpoints()[0].roles()[0].rate_limit_count(), 1);
|
||||
assert_eq!(snapshot.endpoints()[0].roles()[0].success_count(), 1);
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn executor_maps_reqwest_timeout_to_ksp_timeout_error_without_endpoint_secret_leak() {
|
||||
const SECRET_CANARY: &str = "SECRET-REQWEST-URL-CANARY";
|
||||
let (url, handle) = serve_timeout();
|
||||
let sensitive_url = format!("{url}/rpc?api-key={SECRET_CANARY}");
|
||||
let pool = pool_for_url(sensitive_url.as_str(), std::time::Duration::from_millis(20), 0);
|
||||
let error = pool
|
||||
.execute_standard_rpc(&crate::HttpRoleName::new("default"), health_method(), std::vec::Vec::new())
|
||||
.await
|
||||
.expect_err("timed out request must fail");
|
||||
assert_eq!(error.code(), crate::ERROR_CODE_TIMEOUT);
|
||||
assert!(!format!("{error:?}").contains(SECRET_CANARY));
|
||||
let source = std::error::Error::source(&error).expect("transport timeout should preserve a sanitized reqwest source");
|
||||
assert!(!format!("{source:?}").contains(SECRET_CANARY));
|
||||
handle.join().expect("fixture server must join");
|
||||
}
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/unit_tests/grpc_subscribe.rs
|
||||
// version: 2
|
||||
// version: 3
|
||||
|
||||
fn filter_name(value: &str) -> crate::YellowstoneSubscribeFilterName {
|
||||
return crate::YellowstoneSubscribeFilterName::new(value).expect("fixture filter name must validate");
|
||||
@@ -55,7 +55,7 @@ fn yellowstone_subscribe_empty_and_named_empty_maps_encode_exactly() {
|
||||
assert_eq!(wire.slots.get("slots"), std::option::Option::Some(&yellowstone_grpc_proto::geyser::SubscribeRequestFilterSlots::default()));
|
||||
assert_eq!(
|
||||
wire.transactions.get("transactions"),
|
||||
std::option::Option::Some(&yellowstone_grpc_proto::geyser::SubscribeRequestFilterTransactions::default())
|
||||
std::option::Option::Some(&yellowstone_grpc_proto::geyser::SubscribeRequestFilterTransactions::default()),
|
||||
);
|
||||
assert_eq!(
|
||||
wire.transactions_status.get("transaction-status"),
|
||||
@@ -365,3 +365,191 @@ fn yellowstone_slot_update_preserves_all_current_statuses_and_bounds_dead_error(
|
||||
};
|
||||
assert!(super::decode_slot_update(oversized).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn yellowstone_transaction_filters_encode_complete_current_wire_and_redact_selectors() {
|
||||
let mut filter = crate::YellowstoneSubscribeTransactionFilter::new();
|
||||
filter.set_vote(std::option::Option::Some(false));
|
||||
filter.set_failed(std::option::Option::Some(true));
|
||||
let signature = crate::YellowstoneTransactionSignatureSelector::new("1".repeat(64)).expect("signature selector must validate");
|
||||
filter.set_signature(std::option::Option::Some(signature));
|
||||
assert!(filter.push_account_include(ksp_core_lib::Pubkey::new_from_array([1_u8; 32])).is_ok());
|
||||
assert!(filter.push_account_exclude(ksp_core_lib::Pubkey::new_from_array([2_u8; 32])).is_ok());
|
||||
assert!(filter.push_account_required(ksp_core_lib::Pubkey::new_from_array([3_u8; 32])).is_ok());
|
||||
let cuckoo =
|
||||
crate::YellowstoneCuckooFilter::new(vec![0_u8; 16], 4, 4, 8, 7, crate::YellowstoneCuckooHashAlgorithm::SipHash).expect("cuckoo filter must validate");
|
||||
filter.set_cuckoo_account_include(std::option::Option::Some(cuckoo));
|
||||
filter.set_token_accounts(std::option::Option::Some(crate::YellowstoneTokenAccountExpansion::BalanceChanged));
|
||||
let wire = filter.to_wire();
|
||||
assert_eq!(wire.vote, std::option::Option::Some(false));
|
||||
assert_eq!(wire.failed, std::option::Option::Some(true));
|
||||
assert!(wire.signature.is_some());
|
||||
assert_eq!(wire.account_include.len(), 1);
|
||||
assert_eq!(wire.account_exclude.len(), 1);
|
||||
assert_eq!(wire.account_required.len(), 1);
|
||||
assert!(wire.cuckoo_account_include.is_some());
|
||||
assert_eq!(wire.token_accounts, std::option::Option::Some(yellowstone_grpc_proto::geyser::TokenAccountExpansionControlFlag::BalanceChanged as i32));
|
||||
let debug = format!("{filter:?}");
|
||||
assert!(debug.contains("account_include_count"));
|
||||
assert!(!debug.contains(&"1".repeat(32)));
|
||||
assert!(!debug.contains(&ksp_core_lib::Pubkey::new_from_array([1_u8; 32]).to_string()));
|
||||
assert!(crate::YellowstoneTransactionSignatureSelector::new("contains-0-O-I-l").is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn yellowstone_transaction_update_decodes_current_storage_wire_including_v1_config_and_meta() {
|
||||
let confirmed = yellowstone_grpc_proto::geyser::SubscribeUpdateTransactionInfo {
|
||||
signature: vec![9_u8; 64],
|
||||
is_vote: false,
|
||||
transaction: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::Transaction {
|
||||
signatures: vec![vec![9_u8; 64], vec![8_u8; 64]],
|
||||
message: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::Message {
|
||||
header: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::MessageHeader {
|
||||
num_required_signatures: 2,
|
||||
num_readonly_signed_accounts: 1,
|
||||
num_readonly_unsigned_accounts: 1,
|
||||
}),
|
||||
account_keys: vec![vec![1_u8; 32], vec![2_u8; 32]],
|
||||
recent_blockhash: vec![3_u8; 32],
|
||||
instructions: vec![yellowstone_grpc_proto::solana::storage::confirmed_block::CompiledInstruction {
|
||||
program_id_index: 1,
|
||||
accounts: vec![0_u8, 1],
|
||||
data: vec![4_u8, 5, 6],
|
||||
}],
|
||||
versioned: true,
|
||||
address_table_lookups: vec![yellowstone_grpc_proto::solana::storage::confirmed_block::MessageAddressTableLookup {
|
||||
account_key: vec![4_u8; 32],
|
||||
writable_indexes: vec![1_u8, 2],
|
||||
readonly_indexes: vec![3_u8],
|
||||
}],
|
||||
config: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::TransactionConfig {
|
||||
priority_fee: std::option::Option::Some(7),
|
||||
compute_unit_limit: std::option::Option::Some(8),
|
||||
loaded_accounts_data_size_limit: std::option::Option::Some(9),
|
||||
heap_size: std::option::Option::Some(10),
|
||||
}),
|
||||
}),
|
||||
}),
|
||||
meta: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::TransactionStatusMeta {
|
||||
err: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::TransactionError { err: vec![11_u8, 12] }),
|
||||
fee: 5_000,
|
||||
pre_balances: vec![100, 200],
|
||||
post_balances: vec![90, 210],
|
||||
inner_instructions: vec![yellowstone_grpc_proto::solana::storage::confirmed_block::InnerInstructions {
|
||||
index: 0,
|
||||
instructions: vec![yellowstone_grpc_proto::solana::storage::confirmed_block::InnerInstruction {
|
||||
program_id_index: 1,
|
||||
accounts: vec![0_u8],
|
||||
data: vec![13_u8, 14],
|
||||
stack_height: std::option::Option::Some(2),
|
||||
}],
|
||||
}],
|
||||
inner_instructions_none: false,
|
||||
log_messages: vec!["Program log: fixture".to_owned()],
|
||||
log_messages_none: false,
|
||||
pre_token_balances: vec![yellowstone_grpc_proto::solana::storage::confirmed_block::TokenBalance {
|
||||
account_index: 0,
|
||||
mint: "mint-fixture".to_owned(),
|
||||
ui_token_amount: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::UiTokenAmount {
|
||||
ui_amount: 1.5,
|
||||
decimals: 6,
|
||||
amount: "1500000".to_owned(),
|
||||
ui_amount_string: "1.5".to_owned(),
|
||||
}),
|
||||
owner: "owner-fixture".to_owned(),
|
||||
program_id: "program-fixture".to_owned(),
|
||||
}],
|
||||
post_token_balances: vec![],
|
||||
rewards: vec![yellowstone_grpc_proto::solana::storage::confirmed_block::Reward {
|
||||
pubkey: ksp_core_lib::Pubkey::new_from_array([5_u8; 32]).to_string(),
|
||||
lamports: 17,
|
||||
post_balance: 18,
|
||||
reward_type: yellowstone_grpc_proto::solana::storage::confirmed_block::RewardType::Staking as i32,
|
||||
commission: "5".to_owned(),
|
||||
commission_bps: "500".to_owned(),
|
||||
}],
|
||||
loaded_writable_addresses: vec![vec![6_u8; 32]],
|
||||
loaded_readonly_addresses: vec![vec![7_u8; 32]],
|
||||
return_data: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::ReturnData {
|
||||
program_id: vec![8_u8; 32],
|
||||
data: vec![15_u8, 16],
|
||||
}),
|
||||
return_data_none: false,
|
||||
compute_units_consumed: std::option::Option::Some(123),
|
||||
cost_units: std::option::Option::Some(456),
|
||||
}),
|
||||
index: 3,
|
||||
};
|
||||
let wire = yellowstone_grpc_proto::geyser::SubscribeUpdate {
|
||||
filters: vec!["transactions-main".to_owned()],
|
||||
update_oneof: std::option::Option::Some(yellowstone_grpc_proto::geyser::subscribe_update::UpdateOneof::Transaction(
|
||||
yellowstone_grpc_proto::geyser::SubscribeUpdateTransaction { transaction: std::option::Option::Some(confirmed), slot: 42 },
|
||||
)),
|
||||
created_at: std::option::Option::Some(yellowstone_grpc_proto::prost_types::Timestamp { seconds: 100, nanos: 200 }),
|
||||
};
|
||||
let update = super::decode_transaction_update(wire).expect("transaction update fixture must decode");
|
||||
assert_eq!(update.slot(), 42);
|
||||
assert_eq!(update.filters()[0].as_str(), "transactions-main");
|
||||
assert_eq!(update.transaction().signature().as_bytes(), &[9_u8; 64]);
|
||||
assert_eq!(format!("{:?}", update.transaction().signature()), "YellowstoneTransactionSignature(<redacted>)");
|
||||
assert_eq!(update.transaction().index(), 3);
|
||||
assert_eq!(update.transaction().transaction().signatures().len(), 2);
|
||||
assert!(update.transaction().transaction().message().versioned());
|
||||
let config = update.transaction().transaction().message().config().expect("v1 config must be preserved");
|
||||
assert_eq!(config.priority_fee(), std::option::Option::Some(7));
|
||||
assert_eq!(config.heap_size(), std::option::Option::Some(10));
|
||||
assert_eq!(update.transaction().meta().fee(), 5_000);
|
||||
assert_eq!(update.transaction().meta().error().expect("error must be present").as_bytes(), &[11_u8, 12]);
|
||||
assert_eq!(update.transaction().meta().inner_instructions()[0].instructions()[0].stack_height(), std::option::Option::Some(2));
|
||||
assert_eq!(update.transaction().meta().pre_token_balances()[0].ui_token_amount().expect("token amount must be present").amount(), "1500000");
|
||||
let token_debug = format!("{:?}", update.transaction().meta().pre_token_balances()[0]);
|
||||
assert!(!token_debug.contains("mint-fixture"));
|
||||
assert!(!token_debug.contains("owner-fixture"));
|
||||
assert!(!token_debug.contains("program-fixture"));
|
||||
assert_eq!(update.transaction().meta().rewards()[0].reward_type(), crate::YellowstoneRewardType::Staking);
|
||||
assert_eq!(update.transaction().meta().loaded_writable_addresses()[0], ksp_core_lib::Pubkey::new_from_array([6_u8; 32]));
|
||||
assert_eq!(update.transaction().meta().return_data().expect("return data must be present").data(), &[15_u8, 16]);
|
||||
assert_eq!(update.transaction().meta().compute_units_consumed(), std::option::Option::Some(123));
|
||||
assert_eq!(update.transaction().meta().cost_units(), std::option::Option::Some(456));
|
||||
let debug = format!("{update:?}");
|
||||
assert!(!debug.contains("Program log: fixture"));
|
||||
assert!(!debug.contains("11, 12"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn yellowstone_transaction_status_update_preserves_error_and_rejects_malformed_signature() {
|
||||
let wire = yellowstone_grpc_proto::geyser::SubscribeUpdate {
|
||||
filters: vec!["status".to_owned()],
|
||||
update_oneof: std::option::Option::Some(yellowstone_grpc_proto::geyser::subscribe_update::UpdateOneof::TransactionStatus(
|
||||
yellowstone_grpc_proto::geyser::SubscribeUpdateTransactionStatus {
|
||||
slot: 55,
|
||||
signature: vec![2_u8; 64],
|
||||
is_vote: true,
|
||||
index: 4,
|
||||
err: std::option::Option::Some(yellowstone_grpc_proto::solana::storage::confirmed_block::TransactionError { err: vec![99_u8] }),
|
||||
},
|
||||
)),
|
||||
created_at: std::option::Option::None,
|
||||
};
|
||||
let update = super::decode_transaction_status_update(wire).expect("transaction-status update must decode");
|
||||
assert_eq!(update.slot(), 55);
|
||||
assert!(update.is_vote());
|
||||
assert_eq!(update.index(), 4);
|
||||
assert_eq!(update.signature().as_bytes(), &[2_u8; 64]);
|
||||
assert_eq!(update.error().expect("status error must be present").as_bytes(), &[99_u8]);
|
||||
assert!(!format!("{update:?}").contains("99"));
|
||||
let malformed = yellowstone_grpc_proto::geyser::SubscribeUpdate {
|
||||
filters: vec!["status".to_owned()],
|
||||
update_oneof: std::option::Option::Some(yellowstone_grpc_proto::geyser::subscribe_update::UpdateOneof::TransactionStatus(
|
||||
yellowstone_grpc_proto::geyser::SubscribeUpdateTransactionStatus {
|
||||
slot: 1,
|
||||
signature: vec![0_u8; 63],
|
||||
is_vote: false,
|
||||
index: 0,
|
||||
err: std::option::Option::None,
|
||||
},
|
||||
)),
|
||||
created_at: std::option::Option::None,
|
||||
};
|
||||
assert!(super::decode_transaction_status_update(malformed).is_err());
|
||||
}
|
||||
|
||||
@@ -1,336 +0,0 @@
|
||||
// file: crates/ksp-onchain-transport-lib/unit_tests/pool.rs
|
||||
// version: 4
|
||||
|
||||
fn role(name: &str, priority: u32, request_kinds: std::vec::Vec<crate::HttpRequestKind>) -> crate::HttpEndpointRoleSettings {
|
||||
return crate::HttpEndpointRoleSettings::new(
|
||||
crate::HttpRoleName::new(name),
|
||||
true,
|
||||
request_kinds,
|
||||
priority,
|
||||
crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None),
|
||||
);
|
||||
}
|
||||
|
||||
fn endpoint(name: &str, enabled: bool, priority: u32, request_kinds: std::vec::Vec<crate::HttpRequestKind>) -> crate::HttpEndpointSettings {
|
||||
return crate::HttpEndpointSettings::new(
|
||||
name,
|
||||
enabled,
|
||||
crate::HttpProviderName::new("provider"),
|
||||
crate::HttpClusterName::new("devnet"),
|
||||
crate::HttpEndpointUrl::parse(format!("https://{name}.invalid/rpc?token=SECRET-CANARY")).expect("test URL must parse"),
|
||||
std::time::Duration::from_secs(1),
|
||||
std::time::Duration::from_secs(2),
|
||||
std::option::Option::Some(4),
|
||||
std::vec![role("default", priority, request_kinds)],
|
||||
);
|
||||
}
|
||||
|
||||
fn non_zero(value: u32) -> std::num::NonZeroU32 {
|
||||
return std::num::NonZeroU32::new(value).expect("test limit must be non-zero");
|
||||
}
|
||||
|
||||
fn limited_endpoint(
|
||||
name: &str,
|
||||
priority: u32,
|
||||
requests_per_second: std::option::Option<u32>,
|
||||
burst_capacity: std::option::Option<u32>,
|
||||
max_concurrent_requests: std::option::Option<u32>,
|
||||
cooldown: std::option::Option<std::time::Duration>,
|
||||
) -> crate::HttpEndpointSettings {
|
||||
let limits = crate::HttpRoleLimits::new(
|
||||
requests_per_second.map(|value| return non_zero(value)),
|
||||
burst_capacity.map(|value| return non_zero(value)),
|
||||
max_concurrent_requests.map(|value| return non_zero(value)),
|
||||
cooldown,
|
||||
);
|
||||
let role = crate::HttpEndpointRoleSettings::new(crate::HttpRoleName::new("default"), true, std::vec![crate::HttpRequestKind::wildcard()], priority, limits);
|
||||
return crate::HttpEndpointSettings::new(
|
||||
name,
|
||||
true,
|
||||
crate::HttpProviderName::new("provider"),
|
||||
crate::HttpClusterName::new("devnet"),
|
||||
crate::HttpEndpointUrl::parse(format!("https://{name}.invalid/rpc?token=SECRET-CANARY")).expect("test URL must parse"),
|
||||
std::time::Duration::from_secs(1),
|
||||
std::time::Duration::from_secs(2),
|
||||
std::option::Option::Some(4),
|
||||
std::vec![role],
|
||||
);
|
||||
}
|
||||
|
||||
fn settings(endpoints: std::vec::Vec<crate::HttpEndpointSettings>) -> crate::HttpTransportSettings {
|
||||
return crate::HttpTransportSettings::new(
|
||||
endpoints,
|
||||
crate::HttpRetrySettings::new(2, std::time::Duration::from_millis(10), std::time::Duration::from_millis(50)),
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pool_prefers_lowest_priority_tier() {
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![
|
||||
endpoint("secondary", true, 20, std::vec![crate::HttpRequestKind::wildcard()]),
|
||||
endpoint("primary", true, 10, std::vec![crate::HttpRequestKind::wildcard()]),
|
||||
]))
|
||||
.expect("pool must build");
|
||||
let selection = pool
|
||||
.select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance"))
|
||||
.expect("selection must succeed");
|
||||
assert_eq!(selection.endpoint_name(), "primary");
|
||||
assert_eq!(selection.priority(), 10);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pool_round_robins_fairly_inside_best_priority_tier() {
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![
|
||||
endpoint("one", true, 10, std::vec![crate::HttpRequestKind::wildcard()]),
|
||||
endpoint("two", true, 10, std::vec![crate::HttpRequestKind::wildcard()]),
|
||||
endpoint("fallback", true, 20, std::vec![crate::HttpRequestKind::wildcard()]),
|
||||
]))
|
||||
.expect("pool must build");
|
||||
let role = crate::HttpRoleName::new("default");
|
||||
let kind = crate::HttpRequestKind::new("get_balance");
|
||||
let first = pool.select_for_request_kind(&role, &kind).expect("first selection must succeed");
|
||||
let second = pool.select_for_request_kind(&role, &kind).expect("second selection must succeed");
|
||||
let third = pool.select_for_request_kind(&role, &kind).expect("third selection must succeed");
|
||||
assert_eq!(first.endpoint_name(), "one");
|
||||
assert_eq!(second.endpoint_name(), "two");
|
||||
assert_eq!(third.endpoint_name(), "one");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn disabled_best_priority_endpoint_falls_back_to_next_tier() {
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![
|
||||
endpoint("disabled-primary", false, 1, std::vec![crate::HttpRequestKind::wildcard()]),
|
||||
endpoint("fallback", true, 20, std::vec![crate::HttpRequestKind::wildcard()]),
|
||||
]))
|
||||
.expect("pool must build");
|
||||
let selection = pool
|
||||
.select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance"))
|
||||
.expect("fallback must be selected");
|
||||
assert_eq!(selection.endpoint_name(), "fallback");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pool_filters_role_and_capability_before_priority() {
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![
|
||||
endpoint("wrong-capability", true, 1, std::vec![crate::HttpRequestKind::new("send_transaction")]),
|
||||
endpoint("matching", true, 50, std::vec![crate::HttpRequestKind::new("get_balance")]),
|
||||
]))
|
||||
.expect("pool must build");
|
||||
let selection = pool
|
||||
.select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance"))
|
||||
.expect("matching capability must be selected");
|
||||
assert_eq!(selection.endpoint_name(), "matching");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pool_returns_structured_error_when_no_endpoint_matches() {
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![endpoint("read-only", true, 10, std::vec![crate::HttpRequestKind::new("get_balance")],)]))
|
||||
.expect("pool must build");
|
||||
let error = pool
|
||||
.select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("send_transaction"))
|
||||
.expect_err("unsupported request kind must fail selection");
|
||||
assert_eq!(error.code(), crate::ERROR_CODE_ENDPOINT_SELECTION_FAILED);
|
||||
assert!(!format!("{error:?}").contains("SECRET-CANARY"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn standard_method_selection_uses_registry_request_kind() {
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![endpoint("balance", true, 10, std::vec![crate::HttpRequestKind::new("get_balance")],)]))
|
||||
.expect("pool must build");
|
||||
let method = crate::find_http_rpc_method("getBalance").expect("audited method must exist");
|
||||
let selection = pool.select_for_method(&crate::HttpRoleName::new("default"), method).expect("standard method must route");
|
||||
assert_eq!(selection.request_kind().as_str(), "get_balance");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pool_snapshot_is_safe_and_preserves_disabled_endpoints() {
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![
|
||||
endpoint("enabled", true, 10, std::vec![crate::HttpRequestKind::wildcard()]),
|
||||
endpoint("disabled", false, 10, std::vec![crate::HttpRequestKind::wildcard()]),
|
||||
]))
|
||||
.expect("pool must build");
|
||||
let snapshot = pool.snapshot();
|
||||
let rendered = format!("{snapshot:?} {pool:?}");
|
||||
assert_eq!(snapshot.endpoint_count(), 2);
|
||||
assert_eq!(snapshot.available_endpoint_count(), 1);
|
||||
assert!(!rendered.contains("SECRET-CANARY"));
|
||||
assert!(!rendered.contains(".invalid/rpc"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn disabled_role_is_excluded_before_priority_selection() {
|
||||
let base = endpoint("disabled-role", true, 1, std::vec![crate::HttpRequestKind::wildcard()]);
|
||||
let disabled_role = crate::HttpEndpointRoleSettings::new(
|
||||
crate::HttpRoleName::new("default"),
|
||||
false,
|
||||
std::vec![crate::HttpRequestKind::wildcard()],
|
||||
1,
|
||||
crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None),
|
||||
);
|
||||
let enabled_non_matching_role = crate::HttpEndpointRoleSettings::new(
|
||||
crate::HttpRoleName::new("maintenance"),
|
||||
true,
|
||||
std::vec![crate::HttpRequestKind::wildcard()],
|
||||
1,
|
||||
crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None),
|
||||
);
|
||||
let disabled_role_endpoint = crate::HttpEndpointSettings::new(
|
||||
base.name(),
|
||||
true,
|
||||
base.provider().clone(),
|
||||
base.cluster().clone(),
|
||||
base.url().clone(),
|
||||
base.connect_timeout(),
|
||||
base.request_timeout(),
|
||||
base.max_idle_connections_per_host(),
|
||||
std::vec![disabled_role, enabled_non_matching_role],
|
||||
);
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![
|
||||
disabled_role_endpoint,
|
||||
endpoint("fallback", true, 20, std::vec![crate::HttpRequestKind::wildcard()]),
|
||||
]))
|
||||
.expect("pool must build");
|
||||
let selection = pool
|
||||
.select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance"))
|
||||
.expect("enabled fallback role must be selected");
|
||||
assert_eq!(selection.endpoint_name(), "fallback");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn removed_standard_method_is_rejected_before_endpoint_routing() {
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![endpoint("wildcard", true, 10, std::vec![crate::HttpRequestKind::wildcard()],)]))
|
||||
.expect("pool must build");
|
||||
let method = crate::find_http_rpc_method("confirmTransaction").expect("historical method must exist");
|
||||
let error = pool.select_for_method(&crate::HttpRoleName::new("default"), method).expect_err("removed standard method must be rejected before routing");
|
||||
assert_eq!(error.code(), crate::ERROR_CODE_METHOD_REMOVED);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn runtime_concurrency_saturation_falls_back_to_lower_priority_tier() {
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![
|
||||
limited_endpoint("primary", 1, std::option::Option::None, std::option::Option::None, std::option::Option::Some(1), std::option::Option::None),
|
||||
limited_endpoint("fallback", 20, std::option::Option::None, std::option::Option::None, std::option::Option::Some(1), std::option::Option::None),
|
||||
]))
|
||||
.expect("pool must build");
|
||||
let role = crate::HttpRoleName::new("default");
|
||||
let kind = crate::HttpRequestKind::new("get_balance");
|
||||
let first = pool.acquire_for_request_kind(&role, &kind).await.expect("first request must acquire primary");
|
||||
assert_eq!(first.selection().endpoint_name(), "primary");
|
||||
let second = pool.acquire_for_request_kind(&role, &kind).await.expect("second request must fall back while primary is saturated");
|
||||
assert_eq!(second.selection().endpoint_name(), "fallback");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn runtime_token_bucket_exhaustion_falls_back_without_busy_waiting() {
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![
|
||||
limited_endpoint("primary", 1, std::option::Option::Some(1), std::option::Option::Some(1), std::option::Option::None, std::option::Option::None),
|
||||
limited_endpoint("fallback", 20, std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None),
|
||||
]))
|
||||
.expect("pool must build");
|
||||
let role = crate::HttpRoleName::new("default");
|
||||
let kind = crate::HttpRequestKind::new("get_balance");
|
||||
let first = pool.acquire_for_request_kind(&role, &kind).await.expect("first request must consume primary token");
|
||||
assert_eq!(first.selection().endpoint_name(), "primary");
|
||||
drop(first);
|
||||
let second = pool.acquire_for_request_kind(&role, &kind).await.expect("fallback must be used while primary token bucket refills");
|
||||
assert_eq!(second.selection().endpoint_name(), "fallback");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn provider_cooldown_excludes_rate_limited_role_and_uses_fallback() {
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![
|
||||
limited_endpoint(
|
||||
"primary",
|
||||
1,
|
||||
std::option::Option::None,
|
||||
std::option::Option::None,
|
||||
std::option::Option::None,
|
||||
std::option::Option::Some(std::time::Duration::from_millis(50)),
|
||||
),
|
||||
limited_endpoint("fallback", 20, std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None),
|
||||
]))
|
||||
.expect("pool must build");
|
||||
let role = crate::HttpRoleName::new("default");
|
||||
let kind = crate::HttpRequestKind::new("get_balance");
|
||||
let primary = pool.acquire_for_request_kind(&role, &kind).await.expect("primary must be acquired");
|
||||
assert_eq!(primary.selection().endpoint_name(), "primary");
|
||||
let pause = primary.record_rate_limited(std::option::Option::None);
|
||||
assert_eq!(pause, std::time::Duration::from_millis(50));
|
||||
drop(primary);
|
||||
let fallback = pool.acquire_for_request_kind(&role, &kind).await.expect("fallback must be selected during primary cooldown");
|
||||
assert_eq!(fallback.selection().endpoint_name(), "fallback");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn admission_waits_for_released_concurrency_without_holding_a_sync_mutex_across_await() {
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![limited_endpoint(
|
||||
"primary",
|
||||
1,
|
||||
std::option::Option::None,
|
||||
std::option::Option::None,
|
||||
std::option::Option::Some(1),
|
||||
std::option::Option::None,
|
||||
)]))
|
||||
.expect("pool must build");
|
||||
let role = crate::HttpRoleName::new("default");
|
||||
let kind = crate::HttpRequestKind::new("get_balance");
|
||||
let first = pool.acquire_for_request_kind(&role, &kind).await.expect("first permit must be acquired");
|
||||
let release_task = tokio::spawn(async move {
|
||||
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
|
||||
drop(first);
|
||||
});
|
||||
let second = pool
|
||||
.acquire_for_request_kind_with_timeout(&role, &kind, std::time::Duration::from_millis(100))
|
||||
.await
|
||||
.expect("second permit must wake after concurrency release");
|
||||
assert_eq!(second.selection().endpoint_name(), "primary");
|
||||
release_task.await.expect("release task must complete");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn admission_timeout_is_bounded_when_concurrency_never_becomes_available() {
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![limited_endpoint(
|
||||
"primary",
|
||||
1,
|
||||
std::option::Option::None,
|
||||
std::option::Option::None,
|
||||
std::option::Option::Some(1),
|
||||
std::option::Option::None,
|
||||
)]))
|
||||
.expect("pool must build");
|
||||
let role = crate::HttpRoleName::new("default");
|
||||
let kind = crate::HttpRequestKind::new("get_balance");
|
||||
let _held = pool.acquire_for_request_kind(&role, &kind).await.expect("first permit must be acquired");
|
||||
let error = pool
|
||||
.acquire_for_request_kind_with_timeout(&role, &kind, std::time::Duration::from_millis(20))
|
||||
.await
|
||||
.expect_err("second permit must time out while concurrency remains saturated");
|
||||
assert_eq!(error.code(), crate::ERROR_CODE_TIMEOUT);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn passive_health_snapshot_moves_from_degraded_back_to_available_after_success() {
|
||||
let pool = crate::HttpTransportPool::new(settings(std::vec![limited_endpoint(
|
||||
"primary",
|
||||
1,
|
||||
std::option::Option::None,
|
||||
std::option::Option::None,
|
||||
std::option::Option::None,
|
||||
std::option::Option::None,
|
||||
)]))
|
||||
.expect("pool must build");
|
||||
let role = crate::HttpRoleName::new("default");
|
||||
let kind = crate::HttpRequestKind::new("get_balance");
|
||||
let first = pool.acquire_for_request_kind(&role, &kind).await.expect("request permit must be acquired");
|
||||
first.record_failure();
|
||||
drop(first);
|
||||
let degraded = pool.snapshot();
|
||||
assert_eq!(degraded.endpoints()[0].availability(), crate::HttpEndpointAvailability::Degraded);
|
||||
assert_eq!(degraded.endpoints()[0].roles()[0].failure_count(), 1);
|
||||
let second = pool.acquire_for_request_kind(&role, &kind).await.expect("degraded endpoint remains eligible for passive recovery");
|
||||
second.record_success();
|
||||
drop(second);
|
||||
let recovered = pool.snapshot();
|
||||
assert_eq!(recovered.endpoints()[0].availability(), crate::HttpEndpointAvailability::Available);
|
||||
assert_eq!(recovered.endpoints()[0].roles()[0].success_count(), 1);
|
||||
}
|
||||
@@ -1,186 +0,0 @@
|
||||
// file: crates/ksp-onchain-transport-lib/unit_tests/resilience.rs
|
||||
// version: 2
|
||||
|
||||
fn non_zero(value: u32) -> std::num::NonZeroU32 {
|
||||
return std::num::NonZeroU32::new(value).expect("test limit must be non-zero");
|
||||
}
|
||||
|
||||
fn retry_settings() -> crate::HttpRetrySettings {
|
||||
return crate::HttpRetrySettings::new(4, std::time::Duration::from_millis(100), std::time::Duration::from_millis(500));
|
||||
}
|
||||
|
||||
fn method(name: &str) -> &'static crate::HttpRpcMethodDescriptor {
|
||||
return crate::find_http_rpc_method(name).expect("audited test method must exist");
|
||||
}
|
||||
|
||||
fn role_limits(
|
||||
requests_per_second: std::option::Option<u32>,
|
||||
burst_capacity: std::option::Option<u32>,
|
||||
max_concurrent_requests: std::option::Option<u32>,
|
||||
cooldown: std::option::Option<std::time::Duration>,
|
||||
) -> crate::HttpEndpointRoleSettings {
|
||||
return crate::HttpEndpointRoleSettings::new(
|
||||
crate::HttpRoleName::new("default"),
|
||||
true,
|
||||
std::vec![crate::HttpRequestKind::wildcard()],
|
||||
10,
|
||||
crate::HttpRoleLimits::new(
|
||||
requests_per_second.map(|value| return non_zero(value)),
|
||||
burst_capacity.map(|value| return non_zero(value)),
|
||||
max_concurrent_requests.map(|value| return non_zero(value)),
|
||||
cooldown,
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn retry_backoff_is_exponential_and_bounded() {
|
||||
let settings = retry_settings();
|
||||
assert_eq!(super::retry_backoff(&settings, 1), std::time::Duration::from_millis(100));
|
||||
assert_eq!(super::retry_backoff(&settings, 2), std::time::Duration::from_millis(200));
|
||||
assert_eq!(super::retry_backoff(&settings, 3), std::time::Duration::from_millis(400));
|
||||
assert_eq!(super::retry_backoff(&settings, 4), std::time::Duration::from_millis(500));
|
||||
assert_eq!(super::retry_backoff(&settings, 32), std::time::Duration::from_millis(500));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn retry_safe_timeout_is_retried_until_budget_is_exhausted() {
|
||||
let settings = retry_settings();
|
||||
let first = crate::evaluate_transport_retry(
|
||||
method("getBalance"),
|
||||
&settings,
|
||||
crate::HttpRetryCause::Timeout,
|
||||
crate::HttpDispatchState::DispatchedAmbiguous,
|
||||
0,
|
||||
std::option::Option::None,
|
||||
);
|
||||
assert_eq!(first, crate::HttpRetryDecision::RetryAfter(std::time::Duration::from_millis(100)));
|
||||
let exhausted = crate::evaluate_transport_retry(
|
||||
method("getBalance"),
|
||||
&settings,
|
||||
crate::HttpRetryCause::Timeout,
|
||||
crate::HttpDispatchState::DispatchedAmbiguous,
|
||||
settings.max_retries(),
|
||||
std::option::Option::None,
|
||||
);
|
||||
assert_eq!(exhausted, crate::HttpRetryDecision::Stop);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_submission_never_retries_after_ambiguous_dispatch() {
|
||||
let decision = crate::evaluate_transport_retry(
|
||||
method("sendTransaction"),
|
||||
&retry_settings(),
|
||||
crate::HttpRetryCause::Connection,
|
||||
crate::HttpDispatchState::DispatchedAmbiguous,
|
||||
0,
|
||||
std::option::Option::None,
|
||||
);
|
||||
assert_eq!(decision, crate::HttpRetryDecision::Stop);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_submission_can_retry_when_transport_proves_no_dispatch() {
|
||||
let decision = crate::evaluate_transport_retry(
|
||||
method("sendTransaction"),
|
||||
&retry_settings(),
|
||||
crate::HttpRetryCause::Connection,
|
||||
crate::HttpDispatchState::NotDispatched,
|
||||
0,
|
||||
std::option::Option::None,
|
||||
);
|
||||
assert!(decision.should_retry());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rpc_application_and_invalid_response_are_not_transport_retries() {
|
||||
for cause in [crate::HttpRetryCause::RpcApplication, crate::HttpRetryCause::InvalidResponse, crate::HttpRetryCause::Request] {
|
||||
let decision = crate::evaluate_transport_retry(
|
||||
method("getBalance"),
|
||||
&retry_settings(),
|
||||
cause,
|
||||
crate::HttpDispatchState::NotDispatched,
|
||||
0,
|
||||
std::option::Option::None,
|
||||
);
|
||||
assert_eq!(decision, crate::HttpRetryDecision::Stop);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn provider_retry_after_can_extend_backoff_but_is_defensively_bounded() {
|
||||
let settings = retry_settings();
|
||||
let extended = crate::evaluate_transport_retry(
|
||||
method("getBalance"),
|
||||
&settings,
|
||||
crate::HttpRetryCause::RateLimited,
|
||||
crate::HttpDispatchState::DispatchedAmbiguous,
|
||||
0,
|
||||
std::option::Option::Some(std::time::Duration::from_secs(3)),
|
||||
);
|
||||
assert_eq!(extended.delay(), std::option::Option::Some(std::time::Duration::from_secs(3)));
|
||||
let bounded = crate::evaluate_transport_retry(
|
||||
method("getBalance"),
|
||||
&settings,
|
||||
crate::HttpRetryCause::RateLimited,
|
||||
crate::HttpDispatchState::DispatchedAmbiguous,
|
||||
0,
|
||||
std::option::Option::Some(std::time::Duration::from_secs(600)),
|
||||
);
|
||||
assert_eq!(bounded.delay(), std::option::Option::Some(std::time::Duration::from_secs(60)));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn token_bucket_consumes_burst_then_refills_from_elapsed_time() {
|
||||
let start = std::time::Instant::now();
|
||||
let mut bucket = super::HttpTokenBucketState::new(2, 2, start);
|
||||
assert!(bucket.try_consume_at(start).is_none());
|
||||
assert!(bucket.try_consume_at(start).is_none());
|
||||
assert!(bucket.try_consume_at(start).is_some());
|
||||
let later = start.checked_add(std::time::Duration::from_millis(500)).expect("test instant must advance");
|
||||
assert!(bucket.try_consume_at(later).is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn absent_burst_capacity_defaults_to_one_second_of_rps_capacity() {
|
||||
let role = role_limits(std::option::Option::Some(2), std::option::Option::None, std::option::Option::None, std::option::Option::None);
|
||||
let runtime = std::sync::Arc::new(crate::HttpRoleRuntime::new(&role, std::sync::Arc::new(tokio::sync::Notify::new())));
|
||||
let now = std::time::Instant::now();
|
||||
let first = runtime.try_acquire(now);
|
||||
let second = runtime.try_acquire(now);
|
||||
let third = runtime.try_acquire(now);
|
||||
assert!(matches!(first, crate::RoleAdmissionAttempt::Ready(_)));
|
||||
assert!(matches!(second, crate::RoleAdmissionAttempt::Ready(_)));
|
||||
assert!(matches!(third, crate::RoleAdmissionAttempt::BlockedUntil(_)));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn concurrency_semaphore_releases_capacity_when_permit_is_dropped() {
|
||||
let role = role_limits(std::option::Option::None, std::option::Option::None, std::option::Option::Some(1), std::option::Option::None);
|
||||
let runtime = std::sync::Arc::new(crate::HttpRoleRuntime::new(&role, std::sync::Arc::new(tokio::sync::Notify::new())));
|
||||
let now = std::time::Instant::now();
|
||||
let first = runtime.try_acquire(now);
|
||||
let held = match first {
|
||||
crate::RoleAdmissionAttempt::Ready(permit) => permit,
|
||||
_ => panic!("first concurrency permit must be available"),
|
||||
};
|
||||
assert!(matches!(runtime.try_acquire(now), crate::RoleAdmissionAttempt::ConcurrencySaturated));
|
||||
drop(held);
|
||||
assert!(matches!(runtime.try_acquire(now), crate::RoleAdmissionAttempt::Ready(_)));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rate_limit_cooldown_marks_role_and_caps_provider_delay() {
|
||||
let role = role_limits(
|
||||
std::option::Option::None,
|
||||
std::option::Option::None,
|
||||
std::option::Option::None,
|
||||
std::option::Option::Some(std::time::Duration::from_millis(10)),
|
||||
);
|
||||
let runtime = crate::HttpRoleRuntime::new(&role, std::sync::Arc::new(tokio::sync::Notify::new()));
|
||||
let pause = runtime.record_rate_limited(std::option::Option::Some(std::time::Duration::from_secs(600)));
|
||||
assert_eq!(pause, std::time::Duration::from_secs(60));
|
||||
assert_eq!(runtime.rate_limit_count(), 1);
|
||||
assert_eq!(runtime.failure_count(), 1);
|
||||
assert_eq!(runtime.availability(std::time::Instant::now()), crate::HttpEndpointAvailability::RateLimited);
|
||||
}
|
||||
@@ -1,207 +0,0 @@
|
||||
// file: crates/ksp-onchain-transport-lib/unit_tests/settings.rs
|
||||
// version: 2
|
||||
|
||||
fn non_zero(value: u32) -> std::num::NonZeroU32 {
|
||||
return std::num::NonZeroU32::new(value).expect("test non-zero value must remain non-zero");
|
||||
}
|
||||
|
||||
fn valid_settings(url_text: &str) -> crate::HttpTransportSettings {
|
||||
let url = crate::HttpEndpointUrl::parse(url_text).expect("test URL must be valid");
|
||||
let limits = crate::HttpRoleLimits::new(
|
||||
std::option::Option::Some(non_zero(10)),
|
||||
std::option::Option::Some(non_zero(20)),
|
||||
std::option::Option::Some(non_zero(4)),
|
||||
std::option::Option::Some(std::time::Duration::from_millis(500)),
|
||||
);
|
||||
let role = crate::HttpEndpointRoleSettings::new(crate::HttpRoleName::new("default"), true, std::vec![crate::HttpRequestKind::wildcard()], 100, limits);
|
||||
let endpoint = crate::HttpEndpointSettings::new(
|
||||
"devnet_public",
|
||||
true,
|
||||
crate::HttpProviderName::new("solana-public"),
|
||||
crate::HttpClusterName::new("devnet"),
|
||||
url,
|
||||
std::time::Duration::from_secs(5),
|
||||
std::time::Duration::from_secs(15),
|
||||
std::option::Option::Some(8),
|
||||
std::vec![role],
|
||||
);
|
||||
return crate::HttpTransportSettings::new(
|
||||
std::vec![endpoint],
|
||||
crate::HttpRetrySettings::new(2, std::time::Duration::from_millis(100), std::time::Duration::from_secs(2)),
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn endpoint_url_accepts_http_and_https() {
|
||||
assert!(crate::HttpEndpointUrl::parse("https://api.devnet.solana.com").is_ok());
|
||||
assert!(crate::HttpEndpointUrl::parse("http://127.0.0.1:8899").is_ok());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn endpoint_url_rejects_non_http_schemes() {
|
||||
let result = crate::HttpEndpointUrl::parse("ws://api.devnet.solana.com");
|
||||
let error = result.expect_err("WebSocket URL must not be accepted by HTTP settings");
|
||||
assert_eq!(error.code(), crate::ERROR_CODE_INVALID_SETTINGS);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn endpoint_url_debug_redacts_secret_material() {
|
||||
let url = crate::HttpEndpointUrl::parse("https://provider.invalid/rpc?api-key=SECRET-CANARY").expect("test URL must parse");
|
||||
let rendered = format!("{url:?}");
|
||||
assert!(rendered.contains("<redacted>"));
|
||||
assert!(!rendered.contains("SECRET-CANARY"));
|
||||
assert!(!rendered.contains("provider.invalid"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn valid_transport_settings_pass_validation() {
|
||||
let settings = valid_settings("https://api.devnet.solana.com");
|
||||
assert!(settings.validate().is_ok());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transport_settings_debug_does_not_leak_endpoint_url() {
|
||||
let settings = valid_settings("https://provider.invalid/rpc?api-key=SECRET-CANARY");
|
||||
let rendered = format!("{settings:?}");
|
||||
assert!(!rendered.contains("SECRET-CANARY"));
|
||||
assert!(!rendered.contains("provider.invalid"));
|
||||
assert!(rendered.contains("HttpEndpointUrl(<redacted>)"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transport_settings_require_one_enabled_endpoint() {
|
||||
let url = crate::HttpEndpointUrl::parse("https://api.devnet.solana.com").expect("test URL must parse");
|
||||
let role = crate::HttpEndpointRoleSettings::new(
|
||||
crate::HttpRoleName::new("default"),
|
||||
true,
|
||||
std::vec![crate::HttpRequestKind::wildcard()],
|
||||
100,
|
||||
crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None),
|
||||
);
|
||||
let endpoint = crate::HttpEndpointSettings::new(
|
||||
"disabled",
|
||||
false,
|
||||
crate::HttpProviderName::new("provider"),
|
||||
crate::HttpClusterName::new("devnet"),
|
||||
url,
|
||||
std::time::Duration::from_secs(1),
|
||||
std::time::Duration::from_secs(1),
|
||||
std::option::Option::None,
|
||||
std::vec![role],
|
||||
);
|
||||
let settings = crate::HttpTransportSettings::new(
|
||||
std::vec![endpoint],
|
||||
crate::HttpRetrySettings::new(1, std::time::Duration::from_millis(1), std::time::Duration::from_millis(2)),
|
||||
);
|
||||
let error = settings.validate().expect_err("all-disabled settings must fail");
|
||||
assert_eq!(error.code(), crate::ERROR_CODE_INVALID_SETTINGS);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transport_settings_reject_duplicate_endpoint_names() {
|
||||
let first = valid_settings("https://one.invalid");
|
||||
let second = valid_settings("https://two.invalid");
|
||||
let settings = crate::HttpTransportSettings::new(std::vec![first.endpoints()[0].clone(), second.endpoints()[0].clone()], first.retry().clone());
|
||||
let error = settings.validate().expect_err("duplicate endpoint names must fail");
|
||||
assert_eq!(error.code(), crate::ERROR_CODE_INVALID_SETTINGS);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transport_settings_reject_duplicate_roles() {
|
||||
let base = valid_settings("https://api.devnet.solana.com");
|
||||
let endpoint = &base.endpoints()[0];
|
||||
let duplicated_endpoint = crate::HttpEndpointSettings::new(
|
||||
endpoint.name(),
|
||||
true,
|
||||
endpoint.provider().clone(),
|
||||
endpoint.cluster().clone(),
|
||||
endpoint.url().clone(),
|
||||
endpoint.connect_timeout(),
|
||||
endpoint.request_timeout(),
|
||||
endpoint.max_idle_connections_per_host(),
|
||||
std::vec![endpoint.roles()[0].clone(), endpoint.roles()[0].clone()],
|
||||
);
|
||||
let settings = crate::HttpTransportSettings::new(std::vec![duplicated_endpoint], base.retry().clone());
|
||||
assert!(settings.validate().is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transport_settings_reject_wildcard_mixed_with_specific_kind() {
|
||||
let base = valid_settings("https://api.devnet.solana.com");
|
||||
let endpoint = &base.endpoints()[0];
|
||||
let role = crate::HttpEndpointRoleSettings::new(
|
||||
crate::HttpRoleName::new("default"),
|
||||
true,
|
||||
std::vec![crate::HttpRequestKind::wildcard(), crate::HttpRequestKind::new("get_balance")],
|
||||
100,
|
||||
endpoint.roles()[0].limits().clone(),
|
||||
);
|
||||
let modified_endpoint = crate::HttpEndpointSettings::new(
|
||||
endpoint.name(),
|
||||
true,
|
||||
endpoint.provider().clone(),
|
||||
endpoint.cluster().clone(),
|
||||
endpoint.url().clone(),
|
||||
endpoint.connect_timeout(),
|
||||
endpoint.request_timeout(),
|
||||
endpoint.max_idle_connections_per_host(),
|
||||
std::vec![role],
|
||||
);
|
||||
let settings = crate::HttpTransportSettings::new(std::vec![modified_endpoint], base.retry().clone());
|
||||
assert!(settings.validate().is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transport_settings_reject_burst_without_rps() {
|
||||
let base = valid_settings("https://api.devnet.solana.com");
|
||||
let endpoint = &base.endpoints()[0];
|
||||
let role = crate::HttpEndpointRoleSettings::new(
|
||||
crate::HttpRoleName::new("default"),
|
||||
true,
|
||||
std::vec![crate::HttpRequestKind::wildcard()],
|
||||
100,
|
||||
crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::Some(non_zero(2)), std::option::Option::None, std::option::Option::None),
|
||||
);
|
||||
let modified_endpoint = crate::HttpEndpointSettings::new(
|
||||
endpoint.name(),
|
||||
true,
|
||||
endpoint.provider().clone(),
|
||||
endpoint.cluster().clone(),
|
||||
endpoint.url().clone(),
|
||||
endpoint.connect_timeout(),
|
||||
endpoint.request_timeout(),
|
||||
endpoint.max_idle_connections_per_host(),
|
||||
std::vec![role],
|
||||
);
|
||||
let settings = crate::HttpTransportSettings::new(std::vec![modified_endpoint], base.retry().clone());
|
||||
assert!(settings.validate().is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transport_settings_reject_reversed_retry_backoff() {
|
||||
let base = valid_settings("https://api.devnet.solana.com");
|
||||
let settings = crate::HttpTransportSettings::new(
|
||||
base.endpoints().to_vec(),
|
||||
crate::HttpRetrySettings::new(2, std::time::Duration::from_secs(2), std::time::Duration::from_secs(1)),
|
||||
);
|
||||
assert!(settings.validate().is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn transport_settings_reject_zero_request_timeout() {
|
||||
let base = valid_settings("https://api.devnet.solana.com");
|
||||
let endpoint = &base.endpoints()[0];
|
||||
let modified_endpoint = crate::HttpEndpointSettings::new(
|
||||
endpoint.name(),
|
||||
true,
|
||||
endpoint.provider().clone(),
|
||||
endpoint.cluster().clone(),
|
||||
endpoint.url().clone(),
|
||||
endpoint.connect_timeout(),
|
||||
std::time::Duration::ZERO,
|
||||
endpoint.max_idle_connections_per_host(),
|
||||
endpoint.roles().to_vec(),
|
||||
);
|
||||
let settings = crate::HttpTransportSettings::new(std::vec![modified_endpoint], base.retry().clone());
|
||||
assert!(settings.validate().is_err());
|
||||
}
|
||||
Reference in New Issue
Block a user