Files
khadhroony-solana-project/crates/ksp-onchain-transport-lib/unit_tests/ws_helius_transactions.rs
2026-08-23 17:19:44 +02:00

733 lines
41 KiB
Rust

// file: crates/ksp-onchain-transport-lib/unit_tests/ws_helius_transactions.rs
// version: 5
use futures_util::SinkExt; // rust-rules: trait-import
use futures_util::StreamExt; // rust-rules: trait-import
fn pubkey(value: &str) -> ksp_core_lib::Pubkey {
return value.parse::<ksp_core_lib::Pubkey>().expect("fixture public key must parse");
}
fn base_filter() -> crate::HeliusTransactionSubscribeFilter {
return crate::HeliusTransactionSubscribeFilter::new(
std::option::Option::Some(false),
std::option::Option::Some(false),
std::option::Option::Some("fixture-signature-secret-canary".to_owned()),
std::option::Option::Some(std::vec![pubkey("11111111111111111111111111111111")]),
std::option::Option::Some(std::vec![pubkey("SysvarC1ock11111111111111111111111111111111")]),
std::option::Option::Some(std::vec![pubkey("Vote111111111111111111111111111111111111111")]),
std::option::Option::Some(crate::HeliusTokenAccountsFilter::BalanceChanged),
);
}
fn assert_oversized_filter_rejected(filter: crate::HeliusTransactionSubscribeFilter) {
let request = crate::HeliusTransactionSubscribeRequest::new(filter, std::option::Option::None);
let error = request.validate().expect_err("50,001 Helius account filters must fail before I/O");
assert_eq!(error.code(), crate::ERROR_CODE_INVALID_RPC_PARAMETERS);
assert!(error.to_string().contains("invalid_rpc_parameters"));
assert!(!error.to_string().contains("11111111111111111111111111111111"));
}
fn helius_endpoint(url: &str) -> crate::WsEndpointSettings {
return helius_endpoint_with_session(url, crate::WsSessionSettings::default());
}
fn helius_endpoint_with_session(url: &str, session: crate::WsSessionSettings) -> crate::WsEndpointSettings {
return crate::WsEndpointSettings::new(
"local_helius_transaction_fixture",
true,
crate::WsProviderName::new("helius"),
crate::WsClusterName::new("local"),
crate::WsProtocolKind::HeliusLaserStream,
crate::WsEndpointUrl::parse(url).expect("local Helius WebSocket URL must parse"),
session,
);
}
fn reconnect_session_settings(backoff: std::time::Duration) -> crate::WsSessionSettings {
let defaults = crate::WsSessionSettings::default();
return crate::WsSessionSettings::new(
std::time::Duration::from_millis(250),
std::time::Duration::from_millis(200),
crate::WsReconnectSettings::new(2, backoff, backoff),
crate::WsResubscribePolicy::ActiveSubscriptions,
defaults.command_queue_capacity(),
defaults.notification_queue_capacity(),
defaults.max_active_subscriptions(),
defaults.max_pending_requests(),
defaults.max_message_size_bytes(),
defaults.max_frame_size_bytes(),
defaults.max_write_buffer_size_bytes(),
);
}
fn backpressure_session_settings() -> crate::WsSessionSettings {
let defaults = crate::WsSessionSettings::default();
return crate::WsSessionSettings::new(
std::time::Duration::from_millis(250),
std::time::Duration::from_millis(200),
crate::WsReconnectSettings::new(0, std::time::Duration::from_millis(10), std::time::Duration::from_millis(10)),
crate::WsResubscribePolicy::ActiveSubscriptions,
defaults.command_queue_capacity(),
1,
2,
defaults.max_pending_requests(),
defaults.max_message_size_bytes(),
defaults.max_frame_size_bytes(),
defaults.max_write_buffer_size_bytes(),
);
}
fn adversarial_payload_session_settings() -> crate::WsSessionSettings {
let defaults = crate::WsSessionSettings::default();
return crate::WsSessionSettings::new(
std::time::Duration::from_millis(250),
std::time::Duration::from_millis(200),
crate::WsReconnectSettings::new(2, std::time::Duration::from_millis(20), std::time::Duration::from_millis(20)),
crate::WsResubscribePolicy::ActiveSubscriptions,
defaults.command_queue_capacity(),
defaults.notification_queue_capacity(),
defaults.max_active_subscriptions(),
defaults.max_pending_requests(),
256,
128,
defaults.max_write_buffer_size_bytes(),
);
}
async fn bind_local_listener() -> (tokio::net::TcpListener, std::string::String) {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.expect("local listener must bind");
let address = listener.local_addr().expect("local listener must expose address");
return (listener, format!("ws://{address}"));
}
async fn read_request(websocket: &mut tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>) -> serde_json::Value {
let message = websocket.next().await.expect("request message must exist").expect("request message must decode");
let text = message.to_text().expect("request must be text");
return serde_json::from_str(text).expect("request must contain JSON");
}
async fn send_result(websocket: &mut tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>, request: &serde_json::Value, result: serde_json::Value) {
let id = request.get("id").and_then(serde_json::Value::as_u64).expect("request id must be numeric");
let response = serde_json::json!({"jsonrpc":"2.0","id":id,"result":result});
websocket.send(tokio_tungstenite::tungstenite::Message::Text(response.to_string().into())).await.expect("local response must send");
return;
}
async fn send_error(
websocket: &mut tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>,
request: &serde_json::Value,
code: i64,
message: &str,
data: serde_json::Value,
) {
let id = request.get("id").and_then(serde_json::Value::as_u64).expect("request id must be numeric");
let response = serde_json::json!({"jsonrpc":"2.0","id":id,"error":{"code":code,"message":message,"data":data}});
websocket.send(tokio_tungstenite::tungstenite::Message::Text(response.to_string().into())).await.expect("local error response must send");
return;
}
async fn send_notification(websocket: &mut tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>, subscription: u64, result: serde_json::Value) {
let notification = serde_json::json!({"jsonrpc":"2.0","method":"transactionNotification","params":{"subscription":subscription,"result":result}});
websocket.send(tokio_tungstenite::tungstenite::Message::Text(notification.to_string().into())).await.expect("local notification must send");
return;
}
async fn send_root_notification(websocket: &mut tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>, subscription: u64, root: u64) {
let notification = serde_json::json!({"jsonrpc":"2.0","method":"rootNotification","params":{"subscription":subscription,"result":root}});
websocket
.send(tokio_tungstenite::tungstenite::Message::Text(notification.to_string().into()))
.await
.expect("local root notification must send");
return;
}
async fn wait_for_gap_count(session: &crate::HeliusLaserStreamWsSession, expected: u64) {
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(2);
loop {
if session.snapshot().continuity_gap_count() >= expected && session.state() == crate::WsSessionState::Active {
return;
}
assert!(tokio::time::Instant::now() < deadline, "Helius session continuity gap count must advance before timeout");
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
}
}
async fn wait_for_overflow_count(session: &crate::HeliusLaserStreamWsSession, expected: u64) {
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(2);
loop {
if session.snapshot().overflow_count() >= expected {
return;
}
assert!(tokio::time::Instant::now() < deadline, "Helius session overflow count must advance before timeout");
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
}
}
async fn wait_for_session_subscription_count(session: &crate::HeliusLaserStreamWsSession, expected: usize) {
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(2);
loop {
if session.snapshot().subscription_count() == expected {
return;
}
assert!(tokio::time::Instant::now() < deadline, "Helius session subscription count must settle before timeout");
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
}
}
async fn wait_for_subscription_state<T>(subscription: &crate::WsSubscription<T>, expected: crate::WsSubscriptionState) {
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(2);
loop {
if subscription.state() == expected {
return;
}
assert!(tokio::time::Instant::now() < deadline, "Helius logical subscription state must advance before timeout");
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
}
}
async fn wait_for_close_frame(websocket: &mut tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>) {
loop {
let message = websocket.next().await;
match message {
std::option::Option::Some(std::result::Result::Ok(tokio_tungstenite::tungstenite::Message::Close(_))) => return,
std::option::Option::Some(std::result::Result::Ok(_)) => {},
std::option::Option::Some(std::result::Result::Err(_)) | std::option::Option::None => return,
}
}
}
#[test]
fn helius_transaction_filter_and_option_enums_match_documented_wire_labels() {
assert_eq!(crate::HeliusTokenAccountsFilter::None.as_str(), "none");
assert_eq!(crate::HeliusTokenAccountsFilter::BalanceChanged.as_str(), "balanceChanged");
assert_eq!(crate::HeliusTokenAccountsFilter::All.as_str(), "all");
assert_eq!(crate::HeliusTransactionSubscribeEncoding::Base58.as_str(), "base58");
assert_eq!(crate::HeliusTransactionSubscribeEncoding::Base64.as_str(), "base64");
assert_eq!(crate::HeliusTransactionSubscribeEncoding::JsonParsed.as_str(), "jsonParsed");
}
#[test]
fn helius_transaction_subscribe_request_serializes_complete_documented_filter_and_options() {
let filter = base_filter();
let options = crate::HeliusTransactionSubscribeOptions::new(
std::option::Option::Some(crate::SolanaCommitment::Confirmed),
std::option::Option::Some(crate::HeliusTransactionSubscribeEncoding::JsonParsed),
std::option::Option::Some(crate::SolanaTransactionDetails::Accounts),
std::option::Option::Some(true),
std::option::Option::Some(0),
);
let request = crate::HeliusTransactionSubscribeRequest::new(filter, std::option::Option::Some(options));
let params = super::helius_transaction_subscribe_params(&request).expect("complete documented Helius request must validate");
assert_eq!(
params,
std::vec![
serde_json::json!({
"vote": false,
"failed": false,
"signature": "fixture-signature-secret-canary",
"accountInclude": ["11111111111111111111111111111111"],
"accountExclude": ["SysvarC1ock11111111111111111111111111111111"],
"accountRequired": ["Vote111111111111111111111111111111111111111"],
"tokenAccounts": "balanceChanged"
}),
serde_json::json!({
"commitment": "confirmed",
"encoding": "jsonParsed",
"transactionDetails": "accounts",
"showRewards": true,
"maxSupportedTransactionVersion": 0
})
]
);
assert_eq!(request.filter().vote(), std::option::Option::Some(false));
assert_eq!(request.filter().failed(), std::option::Option::Some(false));
assert_eq!(request.filter().signature(), std::option::Option::Some("fixture-signature-secret-canary"));
assert_eq!(request.filter().account_include().map(<[ksp_core_lib::Pubkey]>::len), std::option::Option::Some(1));
assert_eq!(request.filter().account_exclude().map(<[ksp_core_lib::Pubkey]>::len), std::option::Option::Some(1));
assert_eq!(request.filter().account_required().map(<[ksp_core_lib::Pubkey]>::len), std::option::Option::Some(1));
assert_eq!(request.filter().token_accounts(), std::option::Option::Some(crate::HeliusTokenAccountsFilter::BalanceChanged));
let options = request.options().expect("options must remain available");
assert_eq!(options.commitment(), std::option::Option::Some(crate::SolanaCommitment::Confirmed));
assert_eq!(options.encoding(), std::option::Option::Some(crate::HeliusTransactionSubscribeEncoding::JsonParsed));
assert_eq!(options.transaction_details(), std::option::Option::Some(crate::SolanaTransactionDetails::Accounts));
assert_eq!(options.show_rewards(), std::option::Option::Some(true));
assert_eq!(options.max_supported_transaction_version(), std::option::Option::Some(0));
}
#[test]
fn helius_transaction_request_preserves_omitted_explicit_empty_and_explicit_none_states() {
let omitted = crate::HeliusTransactionSubscribeRequest::new(crate::HeliusTransactionSubscribeFilter::default(), std::option::Option::None);
assert_eq!(
super::helius_transaction_subscribe_params(&omitted).expect("fully omitted optional request must validate"),
std::vec![serde_json::json!({})]
);
let explicit = crate::HeliusTransactionSubscribeRequest::new(
crate::HeliusTransactionSubscribeFilter::new(
std::option::Option::None,
std::option::Option::None,
std::option::Option::None,
std::option::Option::Some(std::vec::Vec::new()),
std::option::Option::Some(std::vec::Vec::new()),
std::option::Option::Some(std::vec::Vec::new()),
std::option::Option::Some(crate::HeliusTokenAccountsFilter::None),
),
std::option::Option::Some(crate::HeliusTransactionSubscribeOptions::default()),
);
assert_eq!(
super::helius_transaction_subscribe_params(&explicit).expect("explicit empty Helius request states must validate"),
std::vec![serde_json::json!({"accountInclude":[],"accountExclude":[],"accountRequired":[],"tokenAccounts":"none"}), serde_json::json!({})]
);
}
#[test]
fn helius_transaction_filter_enforces_each_documented_fifty_thousand_account_bound() {
let key = pubkey("11111111111111111111111111111111");
let maximum = std::vec![key; 50_000];
let accepted = crate::HeliusTransactionSubscribeRequest::new(
crate::HeliusTransactionSubscribeFilter::new(
std::option::Option::None,
std::option::Option::None,
std::option::Option::None,
std::option::Option::Some(maximum),
std::option::Option::None,
std::option::Option::None,
std::option::Option::None,
),
std::option::Option::None,
);
assert!(accepted.validate().is_ok());
assert_oversized_filter_rejected(crate::HeliusTransactionSubscribeFilter::new(
std::option::Option::None,
std::option::Option::None,
std::option::Option::None,
std::option::Option::Some(std::vec![key; 50_001]),
std::option::Option::None,
std::option::Option::None,
std::option::Option::None,
));
assert_oversized_filter_rejected(crate::HeliusTransactionSubscribeFilter::new(
std::option::Option::None,
std::option::Option::None,
std::option::Option::None,
std::option::Option::None,
std::option::Option::Some(std::vec![key; 50_001]),
std::option::Option::None,
std::option::Option::None,
));
assert_oversized_filter_rejected(crate::HeliusTransactionSubscribeFilter::new(
std::option::Option::None,
std::option::Option::None,
std::option::Option::None,
std::option::Option::None,
std::option::Option::None,
std::option::Option::Some(std::vec![key; 50_001]),
std::option::Option::None,
));
}
#[test]
fn helius_transaction_details_require_max_supported_version_only_for_accounts_and_full() {
for details in [crate::SolanaTransactionDetails::Full, crate::SolanaTransactionDetails::Accounts] {
let options = crate::HeliusTransactionSubscribeOptions::new(
std::option::Option::None,
std::option::Option::None,
std::option::Option::Some(details),
std::option::Option::None,
std::option::Option::None,
);
let request = crate::HeliusTransactionSubscribeRequest::new(crate::HeliusTransactionSubscribeFilter::default(), std::option::Option::Some(options));
let error = request.validate().expect_err("full/accounts details must require maxSupportedTransactionVersion");
assert_eq!(error.code(), crate::ERROR_CODE_INVALID_RPC_PARAMETERS);
}
for details in [crate::SolanaTransactionDetails::Signatures, crate::SolanaTransactionDetails::None] {
let options = crate::HeliusTransactionSubscribeOptions::new(
std::option::Option::None,
std::option::Option::None,
std::option::Option::Some(details),
std::option::Option::None,
std::option::Option::None,
);
let request = crate::HeliusTransactionSubscribeRequest::new(crate::HeliusTransactionSubscribeFilter::default(), std::option::Option::Some(options));
assert!(request.validate().is_ok());
}
let full_with_version = crate::HeliusTransactionSubscribeOptions::new(
std::option::Option::None,
std::option::Option::None,
std::option::Option::Some(crate::SolanaTransactionDetails::Full),
std::option::Option::None,
std::option::Option::Some(0),
);
let request =
crate::HeliusTransactionSubscribeRequest::new(crate::HeliusTransactionSubscribeFilter::default(), std::option::Option::Some(full_with_version));
assert!(request.validate().is_ok());
}
#[test]
fn helius_transaction_notification_decoder_preserves_full_signature_and_unknown_shapes() {
let full_value = serde_json::json!({
"transaction":{"transaction":["AAAA","base64"],"meta":{"err":null}},
"signature":"full-signature",
"slot":224341380,
"transactionIndex":42
});
let full = super::decode_helius_transaction_notification(full_value.clone()).expect("full Helius notification must decode");
match full {
crate::HeliusTransactionNotification::Full(notification) => {
assert_eq!(notification.transaction(), &full_value["transaction"]);
assert_eq!(notification.signature(), "full-signature");
assert_eq!(notification.slot(), 224341380);
assert_eq!(notification.transaction_index(), 42);
},
_ => panic!("transaction member must select the full Helius notification variant"),
}
let signature_value = serde_json::json!({
"signature":"signature-only",
"slot":224341381,
"transactionIndex":43,
"err":null,
"memo":"memo-canary",
"blockTime":1720000000,
"confirmationStatus":"confirmed"
});
let signature = super::decode_helius_transaction_notification(signature_value).expect("signature Helius notification must decode");
match signature {
crate::HeliusTransactionNotification::Signature(notification) => {
assert_eq!(notification.signature(), "signature-only");
assert_eq!(notification.slot(), 224341381);
assert_eq!(notification.transaction_index(), 43);
assert!(matches!(notification.err(), crate::SolanaWireField::Null));
assert!(matches!(notification.memo(), crate::SolanaWireField::Value(value) if value == "memo-canary"));
assert!(matches!(notification.block_time(), crate::SolanaWireField::Value(1720000000)));
assert!(matches!(notification.confirmation_status(), crate::SolanaWireField::Value(value) if value == "confirmed"));
},
_ => panic!("signature envelope must select the lightweight Helius notification variant"),
}
let unknown_value = serde_json::json!({"futureProviderShape":{"value":7}});
let unknown = super::decode_helius_transaction_notification(unknown_value.clone()).expect("unknown Helius notification must remain forward-compatible");
assert!(matches!(unknown, crate::HeliusTransactionNotification::Unknown(value) if value == unknown_value));
}
#[test]
fn helius_transaction_filter_debug_omits_signature_and_account_values() {
let filter = base_filter();
let request = crate::HeliusTransactionSubscribeRequest::new(filter, std::option::Option::None);
let debug = format!("{request:?}");
assert!(debug.contains("signature_present"));
assert!(debug.contains("account_include_count"));
assert!(!debug.contains("fixture-signature-secret-canary"));
assert!(!debug.contains("11111111111111111111111111111111"));
assert!(!debug.contains("SysvarC1ock11111111111111111111111111111111"));
assert!(!debug.contains("Vote111111111111111111111111111111111111111"));
}
#[test]
fn helius_transaction_notification_debug_omits_raw_provider_payloads() {
let full_value = serde_json::json!({
"transaction":{"raw":"MASSIVE-RAW-PAYLOAD-CANARY"},
"signature":"FULL-SIGNATURE-CANARY",
"slot":77,
"transactionIndex":3
});
let full = super::decode_helius_transaction_notification(full_value).expect("full Helius notification must decode");
let signature_value = serde_json::json!({
"signature":"SIGNATURE-MODE-CANARY",
"slot":78,
"transactionIndex":4,
"err":{"secret":"ERROR-DATA-CANARY"},
"memo":"MEMO-CANARY",
"blockTime":123,
"confirmationStatus":"CONFIRMATION-CANARY"
});
let signature = super::decode_helius_transaction_notification(signature_value).expect("signature Helius notification must decode");
let unknown = super::decode_helius_transaction_notification(serde_json::json!({"provider":"UNKNOWN-PAYLOAD-CANARY"}))
.expect("unknown Helius notification must remain forward-compatible");
let rendered = format!("{full:?} {signature:?} {unknown:?}");
assert!(rendered.contains("transaction: \"<omitted>\""));
assert!(rendered.contains("signature: \"<omitted>\""));
assert!(rendered.contains("err: \"value\""));
assert!(rendered.contains("memo: \"value\""));
assert!(rendered.contains("Unknown(\"<omitted>\")"));
for forbidden in [
"MASSIVE-RAW-PAYLOAD-CANARY",
"FULL-SIGNATURE-CANARY",
"SIGNATURE-MODE-CANARY",
"ERROR-DATA-CANARY",
"MEMO-CANARY",
"CONFIRMATION-CANARY",
"UNKNOWN-PAYLOAD-CANARY",
] {
assert!(!rendered.contains(forbidden));
}
}
#[tokio::test(flavor = "current_thread")]
async fn helius_provider_rpc_application_error_is_safe_and_does_not_fail_session() {
let (listener, url) = bind_local_listener().await;
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.expect("local server must accept Helius client");
let mut websocket = tokio_tungstenite::accept_async(stream).await.expect("local Helius handshake must succeed");
let transaction_subscribe = read_request(&mut websocket).await;
assert_eq!(transaction_subscribe["method"], serde_json::json!("transactionSubscribe"));
send_error(
&mut websocket,
&transaction_subscribe,
-32602,
"PROVIDER-MESSAGE-SECRET-CANARY",
serde_json::json!({"apiKey":"PROVIDER-ERROR-SECRET-CANARY","payload":"X".repeat(4096)}),
)
.await;
let root_subscribe = read_request(&mut websocket).await;
assert_eq!(root_subscribe["method"], serde_json::json!("rootSubscribe"));
send_result(&mut websocket, &root_subscribe, serde_json::json!(72)).await;
send_root_notification(&mut websocket, 72, 88).await;
let root_unsubscribe = read_request(&mut websocket).await;
assert_eq!(root_unsubscribe["method"], serde_json::json!("rootUnsubscribe"));
assert_eq!(root_unsubscribe["params"], serde_json::json!([72]));
send_result(&mut websocket, &root_unsubscribe, serde_json::json!(true)).await;
wait_for_close_frame(&mut websocket).await;
});
let endpoint_url = format!("{url}/?api-key=HELIUS-ENDPOINT-SECRET-CANARY");
let session = crate::HeliusLaserStreamWsSession::connect(helius_endpoint(endpoint_url.as_str())).await.expect("Helius facade must connect");
let request = crate::HeliusTransactionSubscribeRequest::new(base_filter(), std::option::Option::None);
let error = session.transaction_subscribe(&request).await.expect_err("provider application error must reject only the logical subscribe request");
assert_eq!(error.code(), crate::ERROR_CODE_RPC_APPLICATION_ERROR);
assert!(error.context().iter().any(|entry| return entry.key() == "rpc_code" && entry.value() == "-32602"));
assert!(error.context().iter().any(|entry| return entry.key() == "method" && entry.value() == "transactionSubscribe"));
let rendered = format!("{error:?} {error} {session:?} {:?}", session.snapshot());
for forbidden in ["PROVIDER-MESSAGE-SECRET-CANARY", "PROVIDER-ERROR-SECRET-CANARY", "HELIUS-ENDPOINT-SECRET-CANARY", "fixture-signature-secret-canary"] {
assert!(!rendered.contains(forbidden));
}
assert_eq!(session.state(), crate::WsSessionState::Active);
wait_for_session_subscription_count(&session, 0).await;
let mut root = session.root_subscribe().await.expect("session must accept a healthy subscription after provider application error");
assert_eq!(root.recv().await.expect("healthy root notification must arrive").expect("healthy root notification must decode"), 88);
assert!(root.unsubscribe().await.expect("healthy root unsubscribe must complete"));
session.close().await.expect("Helius fixture session must close");
server.await.expect("provider error fixture server must finish");
}
#[tokio::test(flavor = "current_thread")]
async fn helius_notification_method_mismatch_fails_only_transaction_subscription() {
let (listener, url) = bind_local_listener().await;
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.expect("local server must accept Helius client");
let mut websocket = tokio_tungstenite::accept_async(stream).await.expect("local Helius handshake must succeed");
let transaction_subscribe = read_request(&mut websocket).await;
send_result(&mut websocket, &transaction_subscribe, serde_json::json!(41)).await;
let root_subscribe = read_request(&mut websocket).await;
send_result(&mut websocket, &root_subscribe, serde_json::json!(42)).await;
send_root_notification(&mut websocket, 41, 5).await;
let cleanup = read_request(&mut websocket).await;
assert_eq!(cleanup["method"], serde_json::json!("transactionUnsubscribe"));
assert_eq!(cleanup["params"], serde_json::json!([41]));
send_result(&mut websocket, &cleanup, serde_json::json!(true)).await;
send_root_notification(&mut websocket, 42, 99).await;
wait_for_close_frame(&mut websocket).await;
});
let session = crate::HeliusLaserStreamWsSession::connect(helius_endpoint(url.as_str())).await.expect("Helius facade must connect");
let request = crate::HeliusTransactionSubscribeRequest::new(crate::HeliusTransactionSubscribeFilter::default(), std::option::Option::None);
let mut transaction = session.transaction_subscribe(&request).await.expect("transaction subscription must register");
let mut root = session.root_subscribe().await.expect("root subscription must register");
wait_for_subscription_state(&transaction, crate::WsSubscriptionState::Failed).await;
assert_eq!(transaction.terminal_error_code(), std::option::Option::Some(crate::ERROR_CODE_WS_PROTOCOL_ERROR));
assert!(transaction.recv().await.is_none());
assert_eq!(session.state(), crate::WsSessionState::Active);
wait_for_session_subscription_count(&session, 1).await;
assert_eq!(root.recv().await.expect("healthy root notification must arrive").expect("healthy root notification must decode"), 99);
assert_eq!(root.state(), crate::WsSubscriptionState::Active);
session.close().await.expect("Helius fixture session must close");
server.await.expect("notification mismatch fixture server must finish");
}
#[tokio::test(flavor = "current_thread")]
async fn helius_oversized_inbound_payload_reconnects_before_provider_json_decode() {
let (listener, url) = bind_local_listener().await;
let server = tokio::spawn(async move {
let (first_stream, _) = listener.accept().await.expect("initial Helius client must connect");
let mut first = tokio_tungstenite::accept_async(first_stream).await.expect("initial Helius handshake must succeed");
first
.send(tokio_tungstenite::tungstenite::Message::Text("PROVIDER-PAYLOAD-CANARY".repeat(32).into()))
.await
.expect("oversized provider fixture payload must send");
let (replacement_stream, _) = listener.accept().await.expect("replacement Helius client must connect");
let mut replacement = tokio_tungstenite::accept_async(replacement_stream).await.expect("replacement Helius handshake must succeed");
let root_subscribe = read_request(&mut replacement).await;
assert_eq!(root_subscribe["method"], serde_json::json!("rootSubscribe"));
send_result(&mut replacement, &root_subscribe, serde_json::json!(91)).await;
let root_unsubscribe = read_request(&mut replacement).await;
assert_eq!(root_unsubscribe["method"], serde_json::json!("rootUnsubscribe"));
send_result(&mut replacement, &root_unsubscribe, serde_json::json!(true)).await;
wait_for_close_frame(&mut replacement).await;
});
let session = crate::HeliusLaserStreamWsSession::connect(helius_endpoint_with_session(url.as_str(), adversarial_payload_session_settings()))
.await
.expect("Helius facade must connect before adversarial payload");
wait_for_gap_count(&session, 1).await;
assert_eq!(session.state(), crate::WsSessionState::Active);
assert_eq!(session.snapshot().continuity_gap_count(), 1);
let mut root = session.root_subscribe().await.expect("recovered Helius session must remain usable");
assert!(root.unsubscribe().await.expect("recovered root subscription must unsubscribe"));
session.close().await.expect("recovered Helius session must close");
server.await.expect("oversized provider payload fixture server must finish");
}
#[tokio::test(flavor = "current_thread")]
async fn helius_transaction_live_handle_decodes_notification_and_unsubscribes_through_shared_actor() {
let (listener, url) = bind_local_listener().await;
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.expect("local server must accept client");
let mut websocket = tokio_tungstenite::accept_async(stream).await.expect("local WebSocket handshake must succeed");
let subscribe = read_request(&mut websocket).await;
assert_eq!(subscribe["method"], serde_json::json!("transactionSubscribe"));
assert_eq!(
subscribe["params"],
serde_json::json!([
{"failed":false,"accountInclude":["11111111111111111111111111111111"],"tokenAccounts":"balanceChanged"},
{"commitment":"confirmed","encoding":"jsonParsed","transactionDetails":"full","showRewards":false,"maxSupportedTransactionVersion":0}
])
);
send_result(&mut websocket, &subscribe, serde_json::json!(4242)).await;
send_notification(
&mut websocket,
4242,
serde_json::json!({
"transaction":{"transaction":["AAAA","base64"],"meta":{"err":null}},
"signature":"live-signature",
"slot":99,
"transactionIndex":7
}),
)
.await;
let unsubscribe = read_request(&mut websocket).await;
assert_eq!(unsubscribe["method"], serde_json::json!("transactionUnsubscribe"));
assert_eq!(unsubscribe["params"], serde_json::json!([4242]));
send_result(&mut websocket, &unsubscribe, serde_json::json!(true)).await;
wait_for_close_frame(&mut websocket).await;
});
let session = crate::HeliusLaserStreamWsSession::connect(helius_endpoint(url.as_str())).await.expect("Helius facade must connect");
let filter = crate::HeliusTransactionSubscribeFilter::new(
std::option::Option::None,
std::option::Option::Some(false),
std::option::Option::None,
std::option::Option::Some(std::vec![pubkey("11111111111111111111111111111111")]),
std::option::Option::None,
std::option::Option::None,
std::option::Option::Some(crate::HeliusTokenAccountsFilter::BalanceChanged),
);
let options = crate::HeliusTransactionSubscribeOptions::new(
std::option::Option::Some(crate::SolanaCommitment::Confirmed),
std::option::Option::Some(crate::HeliusTransactionSubscribeEncoding::JsonParsed),
std::option::Option::Some(crate::SolanaTransactionDetails::Full),
std::option::Option::Some(false),
std::option::Option::Some(0),
);
let request = crate::HeliusTransactionSubscribeRequest::new(filter, std::option::Option::Some(options));
let mut subscription = session.transaction_subscribe(&request).await.expect("public Helius transaction subscription must register");
assert_eq!(subscription.kind(), crate::WsSubscriptionKind::HeliusTransaction);
let notification = subscription.recv().await.expect("Helius transaction notification must arrive").expect("Helius notification must decode");
match notification {
crate::HeliusTransactionNotification::Full(notification) => {
assert_eq!(notification.signature(), "live-signature");
assert_eq!(notification.slot(), 99);
assert_eq!(notification.transaction_index(), 7);
},
_ => panic!("full live payload must decode as HeliusTransactionNotification::Full"),
}
assert!(subscription.unsubscribe().await.expect("transactionUnsubscribe must complete"));
assert_eq!(subscription.state(), crate::WsSubscriptionState::Closed);
session.close().await.expect("Helius fixture session must close");
server.await.expect("local Helius transaction server must finish");
}
#[tokio::test(flavor = "current_thread")]
async fn helius_transaction_reconnect_remaps_remote_id_and_ignores_late_notification_after_unsubscribe() {
let (listener, url) = bind_local_listener().await;
let server = tokio::spawn(async move {
let (first_stream, _) = listener.accept().await.expect("initial Helius client must connect");
let mut first = tokio_tungstenite::accept_async(first_stream).await.expect("initial Helius handshake must succeed");
let first_subscribe = read_request(&mut first).await;
assert_eq!(first_subscribe["method"], serde_json::json!("transactionSubscribe"));
send_result(&mut first, &first_subscribe, serde_json::json!(41)).await;
send_notification(&mut first, 41, serde_json::json!({"signature":"generation-one","slot":1,"transactionIndex":0})).await;
drop(first);
let (second_stream, _) = listener.accept().await.expect("replacement Helius client must connect");
let mut second = tokio_tungstenite::accept_async(second_stream).await.expect("replacement Helius handshake must succeed");
let second_subscribe = read_request(&mut second).await;
assert_eq!(second_subscribe["method"], serde_json::json!("transactionSubscribe"));
assert_eq!(second_subscribe["params"], first_subscribe["params"]);
send_result(&mut second, &second_subscribe, serde_json::json!(99)).await;
send_notification(&mut second, 99, serde_json::json!({"signature":"generation-two","slot":2,"transactionIndex":1})).await;
let unsubscribe = read_request(&mut second).await;
assert_eq!(unsubscribe["method"], serde_json::json!("transactionUnsubscribe"));
assert_eq!(unsubscribe["params"], serde_json::json!([99]));
send_notification(&mut second, 99, serde_json::json!({"signature":"late-after-cancel","slot":3,"transactionIndex":2})).await;
send_result(&mut second, &unsubscribe, serde_json::json!(true)).await;
wait_for_close_frame(&mut second).await;
});
let settings = reconnect_session_settings(std::time::Duration::from_millis(20));
let session = crate::HeliusLaserStreamWsSession::connect(helius_endpoint_with_session(url.as_str(), settings)).await.expect("Helius facade must connect");
let request = crate::HeliusTransactionSubscribeRequest::new(crate::HeliusTransactionSubscribeFilter::default(), std::option::Option::None);
let mut subscription = session.transaction_subscribe(&request).await.expect("initial Helius transaction subscription must register");
let stable_id = subscription.id();
let first = subscription.recv().await.expect("first generation notification must arrive").expect("first generation notification must decode");
assert!(matches!(first, crate::HeliusTransactionNotification::Signature(ref value) if value.signature() == "generation-one"));
let second = tokio::time::timeout(std::time::Duration::from_secs(2), subscription.recv())
.await
.expect("resubscribed Helius notification must remain bounded")
.expect("resubscribed Helius channel must remain open")
.expect("resubscribed Helius notification must decode");
assert!(matches!(second, crate::HeliusTransactionNotification::Signature(ref value) if value.signature() == "generation-two"));
assert_eq!(subscription.id(), stable_id);
assert_eq!(subscription.state(), crate::WsSubscriptionState::Active);
wait_for_gap_count(&session, 1).await;
assert_eq!(session.snapshot().continuity_gap_count(), 1);
assert!(subscription.unsubscribe().await.expect("Helius transaction cancellation must complete"));
assert_eq!(subscription.state(), crate::WsSubscriptionState::Closed);
assert!(tokio::time::timeout(std::time::Duration::from_millis(100), subscription.recv()).await.expect("closed Helius channel must settle").is_none());
session.close().await.expect("Helius fixture session must close");
server.await.expect("local reconnect Helius server must finish");
}
#[tokio::test(flavor = "current_thread")]
async fn helius_transaction_backpressure_fails_only_slow_subscription_and_uses_transaction_unsubscribe_cleanup() {
let (listener, url) = bind_local_listener().await;
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.expect("local server must accept Helius client");
let mut websocket = tokio_tungstenite::accept_async(stream).await.expect("local Helius handshake must succeed");
let transaction_subscribe = read_request(&mut websocket).await;
assert_eq!(transaction_subscribe["method"], serde_json::json!("transactionSubscribe"));
send_result(&mut websocket, &transaction_subscribe, serde_json::json!(41)).await;
let root_subscribe = read_request(&mut websocket).await;
assert_eq!(root_subscribe["method"], serde_json::json!("rootSubscribe"));
send_result(&mut websocket, &root_subscribe, serde_json::json!(42)).await;
send_notification(&mut websocket, 41, serde_json::json!({"signature":"queued","slot":1,"transactionIndex":0})).await;
send_notification(&mut websocket, 41, serde_json::json!({"signature":"overflow","slot":2,"transactionIndex":1})).await;
let cleanup = read_request(&mut websocket).await;
assert_eq!(cleanup["method"], serde_json::json!("transactionUnsubscribe"));
assert_eq!(cleanup["params"], serde_json::json!([41]));
send_result(&mut websocket, &cleanup, serde_json::json!(true)).await;
send_root_notification(&mut websocket, 42, 99).await;
wait_for_close_frame(&mut websocket).await;
});
let session = crate::HeliusLaserStreamWsSession::connect(helius_endpoint_with_session(url.as_str(), backpressure_session_settings()))
.await
.expect("Helius facade must connect");
let request = crate::HeliusTransactionSubscribeRequest::new(crate::HeliusTransactionSubscribeFilter::default(), std::option::Option::None);
let mut slow = session.transaction_subscribe(&request).await.expect("slow Helius transaction subscription must register");
let mut healthy = session.root_subscribe().await.expect("healthy Helius root subscription must register");
wait_for_subscription_state(&slow, crate::WsSubscriptionState::Failed).await;
wait_for_overflow_count(&session, 1).await;
assert_eq!(slow.terminal_error_code(), std::option::Option::Some(crate::ERROR_CODE_WS_BACKPRESSURE_OVERFLOW));
assert_eq!(session.state(), crate::WsSessionState::Active);
assert_eq!(session.snapshot().subscription_count(), 1);
let queued = slow.recv().await.expect("first Helius notification must remain queued").expect("queued Helius notification must decode");
assert!(matches!(queued, crate::HeliusTransactionNotification::Signature(ref value) if value.signature() == "queued"));
assert!(slow.recv().await.is_none());
assert_eq!(healthy.recv().await.expect("healthy root notification must arrive").expect("healthy root notification must decode"), 99);
assert_eq!(healthy.state(), crate::WsSubscriptionState::Active);
assert_eq!(healthy.terminal_error_code(), std::option::Option::None);
session.close().await.expect("Helius fixture session must close");
server.await.expect("local Helius backpressure server must finish");
}