// file: crates/ksp-onchain-transport-lib/unit_tests/ws_helius_transactions.rs // version: 4 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::().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(), ); } 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) -> 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, 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_notification(websocket: &mut tokio_tungstenite::WebSocketStream, 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, 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 { 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_subscription_state(subscription: &crate::WsSubscription, 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) { 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")); } #[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"); }