// file: crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs // version: 11 use futures_util::SinkExt; // rust-rules: trait-import use futures_util::StreamExt; // rust-rules: trait-import fn local_endpoint(url: &str) -> crate::WsEndpointSettings { return local_endpoint_with_session(url, crate::WsSessionSettings::default()); } fn local_endpoint_with_session(url: &str, session: crate::WsSessionSettings) -> crate::WsEndpointSettings { return crate::WsEndpointSettings::new( "local_ws", true, crate::WsProviderName::new("local-fixture"), crate::WsClusterName::new("local"), crate::WsProtocolKind::SolanaStandard, crate::WsEndpointUrl::parse(url).expect("local test WebSocket URL must parse"), session, ); } fn helius_local_endpoint(url: &str) -> crate::WsEndpointSettings { return helius_local_endpoint_with_session(url, crate::WsSessionSettings::default()); } fn helius_local_endpoint_with_session(url: &str, session: crate::WsSessionSettings) -> crate::WsEndpointSettings { return crate::WsEndpointSettings::new( "local_helius_ws", true, crate::WsProviderName::new("helius-fixture"), crate::WsClusterName::new("local"), crate::WsProtocolKind::HeliusLaserStream, crate::WsEndpointUrl::parse(url).expect("local Helius test WebSocket URL must parse"), session, ); } fn session_settings( command_timeout: std::time::Duration, close_timeout: std::time::Duration, max_pending_requests: usize, max_message_size_bytes: usize, max_frame_size_bytes: usize, max_write_buffer_size_bytes: usize, ) -> crate::WsSessionSettings { let defaults = crate::WsSessionSettings::default(); return crate::WsSessionSettings::new( command_timeout, close_timeout, defaults.reconnect().clone(), defaults.resubscribe(), defaults.command_queue_capacity(), defaults.notification_queue_capacity(), defaults.max_active_subscriptions(), max_pending_requests, max_message_size_bytes, max_frame_size_bytes, max_write_buffer_size_bytes, ); } fn reconnect_session_settings(max_retries: u32, backoff: std::time::Duration, resubscribe: crate::WsResubscribePolicy) -> 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(max_retries, backoff, backoff), resubscribe, 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 subscription_session_settings( notification_queue_capacity: usize, max_active_subscriptions: usize, reconnect: crate::WsReconnectSettings, ) -> crate::WsSessionSettings { let defaults = crate::WsSessionSettings::default(); return crate::WsSessionSettings::new( std::time::Duration::from_millis(250), std::time::Duration::from_millis(200), reconnect, crate::WsResubscribePolicy::ActiveSubscriptions, defaults.command_queue_capacity(), notification_queue_capacity, max_active_subscriptions, 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"); } async fn wait_for_state(session: &crate::WsSession, expected: crate::WsSessionState) { let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(1); loop { if session.state() == expected { return; } assert!(tokio::time::Instant::now() < deadline, "session did not reach expected state: {expected:?}"); tokio::time::sleep(std::time::Duration::from_millis(5)).await; } } async fn wait_for_pending_count(session: &crate::WsSession, expected: usize) { let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(1); loop { if session.snapshot().pending_request_count() == expected { return; } assert!(tokio::time::Instant::now() < deadline, "session did not reach pending count {expected}"); tokio::time::sleep(std::time::Duration::from_millis(5)).await; } } async fn wait_for_gap_count(session: &crate::WsSession, 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, "session did not recover with continuity gap count {expected}"); tokio::time::sleep(std::time::Duration::from_millis(5)).await; } } async fn wait_for_overflow_count(session: &crate::WsSession, expected: u64) { let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(1); loop { if session.snapshot().overflow_count() == expected { return; } assert!(tokio::time::Instant::now() < deadline, "session did not reach overflow count {expected}"); tokio::time::sleep(std::time::Duration::from_millis(5)).await; } } async fn wait_for_subscription_count(session: &crate::WsSession, expected: usize) { let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(1); loop { if session.snapshot().subscription_count() == expected { return; } assert!(tokio::time::Instant::now() < deadline, "session did not reach subscription count {expected}"); 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, } } } async fn yield_runtime_steps() { for _ in 0..16 { tokio::task::yield_now().await; } return; } async fn wait_for_observed_ping_with_io_progress(ping_rx: &mut tokio::sync::mpsc::UnboundedReceiver<()>) -> std::time::Duration { let started = tokio::time::Instant::now(); tokio::time::resume(); let observed = tokio::time::timeout(std::time::Duration::from_secs(1), ping_rx.recv()).await; tokio::time::pause(); let elapsed = tokio::time::Instant::now().duration_since(started); assert!( matches!(observed, std::result::Result::Ok(std::option::Option::Some(()))), "heartbeat Ping was not observed during bounded real I/O progress" ); return elapsed; } #[tokio::test(flavor = "current_thread")] async fn websocket_session_connects_and_round_trips_internal_json_rpc() { 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 request = read_request(&mut websocket).await; assert_eq!(request.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("getVersion")); send_result(&mut websocket, &request, serde_json::json!({"solana-core":"fixture"})).await; }); let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed"); assert_eq!(session.state(), crate::WsSessionState::Active); assert_eq!(session.snapshot().pending_request_count(), 0); let result = session.execute_json_rpc("getVersion", std::vec::Vec::new()).await.expect("fixture JSON-RPC call must succeed"); assert_eq!(result.get("solana-core").and_then(serde_json::Value::as_str), std::option::Option::Some("fixture")); tokio::task::yield_now().await; assert_eq!(session.snapshot().pending_request_count(), 0); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_sessions_are_physical_and_distinct_even_for_the_same_url() { let (listener, url) = bind_local_listener().await; let server = tokio::spawn(async move { let first = listener.accept().await.expect("first local client must connect"); let first_ws = tokio_tungstenite::accept_async(first.0).await.expect("first handshake must succeed"); let second = listener.accept().await.expect("second local client must connect"); let second_ws = tokio_tungstenite::accept_async(second.0).await.expect("second handshake must succeed"); return (first_ws, second_ws); }); let first = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("first client must connect"); let second = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("second client must connect"); assert_ne!(first.id(), second.id()); assert_eq!(first.snapshot().endpoint_name(), second.snapshot().endpoint_name()); let _ = server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_pending_map_dispatches_out_of_order_responses_by_request_id() { 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 first = read_request(&mut websocket).await; let second = read_request(&mut websocket).await; send_result(&mut websocket, &second, serde_json::json!("second")).await; send_result(&mut websocket, &first, serde_json::json!("first")).await; }); let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed"); let first_session = session.clone(); let second_session = session.clone(); let (first, second) = tokio::join!( first_session.execute_json_rpc("firstMethod", std::vec::Vec::new()), second_session.execute_json_rpc("secondMethod", std::vec::Vec::new()) ); assert_eq!(first.expect("first response must dispatch"), serde_json::json!("first")); assert_eq!(second.expect("second response must dispatch"), serde_json::json!("second")); tokio::task::yield_now().await; assert_eq!(session.snapshot().pending_request_count(), 0); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_rpc_application_error_does_not_fail_the_physical_session() { 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 request = read_request(&mut websocket).await; 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":-32602,"message":"fixture invalid params"}}); websocket.send(tokio_tungstenite::tungstenite::Message::Text(response.to_string().into())).await.expect("local error response must send"); let request = read_request(&mut websocket).await; send_result(&mut websocket, &request, serde_json::json!(true)).await; }); let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed"); let error = session.execute_json_rpc("badMethod", std::vec::Vec::new()).await.expect_err("RPC application error must surface"); assert_eq!(error.code(), crate::ERROR_CODE_RPC_APPLICATION_ERROR); assert_eq!(session.state(), crate::WsSessionState::Active); let result = session.execute_json_rpc("goodMethod", std::vec::Vec::new()).await.expect("session must remain usable"); assert_eq!(result, serde_json::json!(true)); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_connection_errors_do_not_echo_url_credentials() { let defaults = crate::WsSessionSettings::default(); let short_session = crate::WsSessionSettings::new( std::time::Duration::from_millis(150), defaults.close_timeout(), defaults.reconnect().clone(), defaults.resubscribe(), 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(), ); let endpoint = crate::WsEndpointSettings::new( "redaction_fixture", true, crate::WsProviderName::new("fixture-provider"), crate::WsClusterName::new("local"), crate::WsProtocolKind::SolanaStandard, crate::WsEndpointUrl::parse("ws://user:password@127.0.0.1:1/private?api-key=SECRET-CANARY").expect("test URL must parse"), short_session, ); let error = crate::WsSession::connect(endpoint).await.expect_err("unreachable local endpoint must fail"); let rendered = format!("{error:?}"); assert!(!rendered.contains("SECRET-CANARY")); assert!(!rendered.contains("password")); assert!(!rendered.contains("private?")); } #[tokio::test(flavor = "current_thread")] async fn websocket_pending_request_capacity_rejects_only_the_excess_request() { let (listener, url) = bind_local_listener().await; let (first_seen_tx, first_seen_rx) = tokio::sync::oneshot::channel(); let (release_tx, release_rx) = tokio::sync::oneshot::channel(); 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 first = read_request(&mut websocket).await; let _ = first_seen_tx.send(()); let _ = release_rx.await; send_result(&mut websocket, &first, serde_json::json!("first")).await; }); let settings = session_settings(std::time::Duration::from_secs(1), std::time::Duration::from_millis(200), 1, 64 * 1024, 64 * 1024, 64 * 1024); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); let first_session = session.clone(); let first = tokio::spawn(async move { return first_session.execute_json_rpc("firstMethod", std::vec::Vec::new()).await; }); first_seen_rx.await.expect("server must observe first request"); wait_for_pending_count(&session, 1).await; let second_error = session.execute_json_rpc("secondMethod", std::vec::Vec::new()).await.expect_err("second pending request must be rejected"); assert_eq!(second_error.code(), crate::ERROR_CODE_WS_BACKPRESSURE_OVERFLOW); assert_eq!(session.state(), crate::WsSessionState::Active); release_tx.send(()).expect("first request release signal must send"); assert_eq!(first.await.expect("first request task must join").expect("first request must complete"), serde_json::json!("first")); wait_for_pending_count(&session, 0).await; server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_oversized_outbound_request_is_rejected_before_socket_write() { 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 request = read_request(&mut websocket).await; assert_eq!(request.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("smallMethod")); send_result(&mut websocket, &request, serde_json::json!(true)).await; }); let settings = session_settings(std::time::Duration::from_secs(1), std::time::Duration::from_millis(200), 8, 256, 128, 1024); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); let oversized = serde_json::json!({"value":"X".repeat(512)}); let error = session.execute_json_rpc("oversizedMethod", std::vec![oversized]).await.expect_err("oversized outbound request must be rejected"); assert_eq!(error.code(), crate::ERROR_CODE_INVALID_RPC_PARAMETERS); assert_eq!(session.state(), crate::WsSessionState::Active); let result = session.execute_json_rpc("smallMethod", std::vec::Vec::new()).await.expect("small request must still succeed"); assert_eq!(result, serde_json::json!(true)); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_oversized_inbound_frame_triggers_reconnect_before_json_decode() { let (listener, url) = bind_local_listener().await; let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.expect("local server must accept initial client"); let mut websocket = tokio_tungstenite::accept_async(stream).await.expect("initial local WebSocket handshake must succeed"); websocket.send(tokio_tungstenite::tungstenite::Message::Text("X".repeat(512).into())).await.expect("oversized fixture message must send"); let (replacement_stream, _) = listener.accept().await.expect("local server must accept replacement client"); let mut replacement = tokio_tungstenite::accept_async(replacement_stream).await.expect("replacement local WebSocket handshake must succeed"); let _ = replacement.next().await; let _ = replacement.flush().await; }); let settings = session_settings(std::time::Duration::from_secs(1), std::time::Duration::from_millis(200), 8, 128, 64, 1024); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); wait_for_gap_count(&session, 1).await; assert_eq!(session.state(), crate::WsSessionState::Active); session.close().await.expect("recovered session must close cleanly"); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_pending_request_timeout_purges_capacity_without_failing_session() { 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 _ = read_request(&mut websocket).await; tokio::time::sleep(std::time::Duration::from_millis(400)).await; }); let settings = session_settings(std::time::Duration::from_millis(100), std::time::Duration::from_millis(100), 1, 64 * 1024, 64 * 1024, 64 * 1024); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); let wait = tokio::time::timeout(std::time::Duration::from_secs(1), session.execute_json_rpc("timeoutMethod", std::vec::Vec::new())).await; let error = wait.expect("request timeout fixture must remain bounded").expect_err("remote silence must time out the request"); assert_eq!(error.code(), crate::ERROR_CODE_TIMEOUT); wait_for_pending_count(&session, 0).await; assert_eq!(session.state(), crate::WsSessionState::Active); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_ping_flushes_automatic_pong_and_keeps_session_active() { let (listener, url) = bind_local_listener().await; let (pong_tx, pong_rx) = tokio::sync::oneshot::channel(); 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"); websocket.send(tokio_tungstenite::tungstenite::Message::Ping(b"ping-canary".to_vec().into())).await.expect("fixture ping must send"); let pong = websocket.next().await.expect("pong frame must arrive").expect("pong frame must decode"); match pong { tokio_tungstenite::tungstenite::Message::Pong(payload) => assert_eq!(payload.as_ref(), b"ping-canary"), other => panic!("expected Pong frame, got {other:?}"), } let _ = pong_tx.send(()); let request = read_request(&mut websocket).await; send_result(&mut websocket, &request, serde_json::json!(true)).await; }); let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed"); tokio::time::timeout(std::time::Duration::from_secs(1), pong_rx) .await .expect("automatic Pong must remain bounded") .expect("fixture server must observe Pong"); assert_eq!(session.state(), crate::WsSessionState::Active); let result = session.execute_json_rpc("afterPing", std::vec::Vec::new()).await.expect("session must remain usable after Ping/Pong"); assert_eq!(result, serde_json::json!(true)); server.await.expect("local server task must complete"); } #[test] fn helius_heartbeat_policy_is_provider_owned_and_fixed_to_sixty_seconds() { assert!(super::helius_heartbeat_enabled(crate::WsProtocolKind::HeliusLaserStream)); assert!(!super::helius_heartbeat_enabled(crate::WsProtocolKind::SolanaStandard)); assert_eq!(super::HELIUS_WS_HEARTBEAT_INTERVAL, std::time::Duration::from_secs(60)); } #[tokio::test(flavor = "current_thread")] async fn helius_heartbeat_sends_ping_at_sixty_seconds_and_rearms() { let (listener, url) = bind_local_listener().await; let (ping_tx, mut ping_rx) = tokio::sync::mpsc::unbounded_channel::<()>(); let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.expect("local Helius server must accept client"); let mut websocket = tokio_tungstenite::accept_async(stream).await.expect("local Helius handshake must succeed"); loop { let message = websocket.next().await; match message { std::option::Option::Some(std::result::Result::Ok(tokio_tungstenite::tungstenite::Message::Ping(_))) => { ping_tx.send(()).expect("heartbeat observation channel must remain open"); websocket.flush().await.expect("automatic heartbeat Pong must flush"); }, 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, } } }); let session = crate::HeliusLaserStreamWsSession::connect(helius_local_endpoint(url.as_str())).await.expect("Helius fixture session must connect"); tokio::time::pause(); yield_runtime_steps().await; tokio::time::advance(std::time::Duration::from_secs(59)).await; yield_runtime_steps().await; assert!(matches!(ping_rx.try_recv(), std::result::Result::Err(tokio::sync::mpsc::error::TryRecvError::Empty))); tokio::time::advance(std::time::Duration::from_secs(1)).await; let first_io_elapsed = wait_for_observed_ping_with_io_progress(&mut ping_rx).await; let pre_second_interval = std::time::Duration::from_secs(59).checked_sub(first_io_elapsed).expect("bounded first Ping I/O progress must remain below one second"); tokio::time::advance(pre_second_interval).await; yield_runtime_steps().await; assert!(matches!(ping_rx.try_recv(), std::result::Result::Err(tokio::sync::mpsc::error::TryRecvError::Empty))); tokio::time::advance(std::time::Duration::from_secs(1) + first_io_elapsed).await; let _second_io_elapsed = wait_for_observed_ping_with_io_progress(&mut ping_rx).await; session.close().await.expect("Helius heartbeat fixture session must close"); tokio::time::resume(); server.await.expect("Helius heartbeat fixture server must complete"); } #[tokio::test(flavor = "current_thread")] async fn standard_session_never_emits_helius_provider_heartbeat() { let (listener, url) = bind_local_listener().await; let (ping_tx, mut ping_rx) = tokio::sync::mpsc::unbounded_channel::<()>(); let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.expect("local standard server must accept client"); let mut websocket = tokio_tungstenite::accept_async(stream).await.expect("local standard handshake must succeed"); loop { let message = websocket.next().await; match message { std::option::Option::Some(std::result::Result::Ok(tokio_tungstenite::tungstenite::Message::Ping(_))) => { ping_tx.send(()).expect("standard heartbeat observation channel must remain open"); websocket.flush().await.expect("automatic Pong must flush"); }, 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, } } }); let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("standard fixture session must connect"); tokio::time::pause(); tokio::time::advance(std::time::Duration::from_secs(180)).await; yield_runtime_steps().await; assert!(matches!(ping_rx.try_recv(), std::result::Result::Err(tokio::sync::mpsc::error::TryRecvError::Empty))); session.close().await.expect("standard fixture session must close"); tokio::time::resume(); server.await.expect("standard heartbeat absence fixture server must complete"); } #[tokio::test(flavor = "current_thread")] async fn helius_explicit_close_cancels_heartbeat_before_deadline() { let (listener, url) = bind_local_listener().await; let (ping_tx, mut ping_rx) = tokio::sync::mpsc::unbounded_channel::<()>(); let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.expect("local Helius close server must accept client"); let mut websocket = tokio_tungstenite::accept_async(stream).await.expect("local Helius close handshake must succeed"); loop { let message = websocket.next().await; match message { std::option::Option::Some(std::result::Result::Ok(tokio_tungstenite::tungstenite::Message::Ping(_))) => { ping_tx.send(()).expect("close heartbeat observation channel must remain open"); websocket.flush().await.expect("automatic heartbeat Pong must flush"); }, 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, } } }); let session = crate::HeliusLaserStreamWsSession::connect(helius_local_endpoint(url.as_str())).await.expect("Helius close fixture session must connect"); tokio::time::pause(); yield_runtime_steps().await; tokio::time::advance(std::time::Duration::from_secs(30)).await; yield_runtime_steps().await; assert!(matches!(ping_rx.try_recv(), std::result::Result::Err(tokio::sync::mpsc::error::TryRecvError::Empty))); session.close().await.expect("Helius close fixture session must close before heartbeat deadline"); yield_runtime_steps().await; assert!(!matches!(ping_rx.try_recv(), std::result::Result::Ok(()))); tokio::time::resume(); server.await.expect("Helius close fixture server must complete"); } #[tokio::test(flavor = "current_thread")] async fn helius_heartbeat_write_failure_maps_to_existing_reconnect_failure_outcome() { let (client_io, peer_io) = tokio::io::duplex(64); let mut websocket = tokio_tungstenite::WebSocketStream::from_raw_socket(client_io, tokio_tungstenite::tungstenite::protocol::Role::Client, std::option::Option::None).await; drop(peer_io); let endpoint = helius_local_endpoint("ws://127.0.0.1:65535"); let (_shutdown_tx, mut shutdown_rx) = tokio::sync::watch::channel(std::option::Option::None::); let session_id = crate::WsSessionId::new(std::num::NonZeroU64::new(1).expect("fixture session id must be non-zero")); let outcome = super::send_helius_heartbeat(session_id, &endpoint, &mut websocket, &mut shutdown_rx).await; match outcome { super::WsActorIoOutcome::Failed { code, pending_message } => { assert_eq!(code, crate::ERROR_CODE_WS_CONNECTION_FAILED); assert_eq!(pending_message, "WebSocket connection failed while writing Helius heartbeat Ping"); }, _ => panic!("heartbeat write failure must use the existing reconnect failure outcome"), } } #[tokio::test(flavor = "current_thread")] async fn helius_heartbeat_is_rearmed_from_successful_reconnect() { let (listener, url) = bind_local_listener().await; let (disconnect_tx, disconnect_rx) = tokio::sync::oneshot::channel::<()>(); let (replacement_tx, replacement_rx) = tokio::sync::oneshot::channel::<()>(); let (ping_tx, mut ping_rx) = tokio::sync::mpsc::unbounded_channel::<()>(); let server = tokio::spawn(async move { let (first_stream, _) = listener.accept().await.expect("first Helius connection must be accepted"); let mut first = tokio_tungstenite::accept_async(first_stream).await.expect("first Helius handshake must succeed"); disconnect_rx.await.expect("disconnect trigger must arrive"); first.send(tokio_tungstenite::tungstenite::Message::Close(std::option::Option::None)).await.expect("first Helius close must send"); let (replacement_stream, _) = listener.accept().await.expect("replacement Helius connection must be accepted"); let mut replacement = tokio_tungstenite::accept_async(replacement_stream).await.expect("replacement Helius handshake must succeed"); replacement_tx.send(()).expect("replacement observation must be delivered"); loop { let message = replacement.next().await; match message { std::option::Option::Some(std::result::Result::Ok(tokio_tungstenite::tungstenite::Message::Ping(_))) => { ping_tx.send(()).expect("replacement heartbeat observation channel must remain open"); replacement.flush().await.expect("replacement automatic Pong must flush"); }, 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, } } }); let settings = reconnect_session_settings(1, std::time::Duration::from_millis(20), crate::WsResubscribePolicy::ActiveSubscriptions); let session = crate::HeliusLaserStreamWsSession::connect(helius_local_endpoint_with_session(url.as_str(), settings)) .await .expect("Helius reconnect fixture session must connect"); disconnect_tx.send(()).expect("disconnect trigger must send"); tokio::time::timeout(std::time::Duration::from_secs(1), replacement_rx) .await .expect("replacement connection must remain bounded") .expect("replacement connection must be observed"); let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(1); loop { if session.state() == crate::WsSessionState::Active && session.snapshot().continuity_gap_count() == 1 { break; } assert!(tokio::time::Instant::now() < deadline, "Helius fixture must recover before heartbeat rearm check"); tokio::time::sleep(std::time::Duration::from_millis(5)).await; } tokio::time::pause(); tokio::time::advance(std::time::Duration::from_secs(59)).await; yield_runtime_steps().await; assert!(matches!(ping_rx.try_recv(), std::result::Result::Err(tokio::sync::mpsc::error::TryRecvError::Empty))); tokio::time::advance(std::time::Duration::from_secs(1)).await; yield_runtime_steps().await; assert!(matches!(ping_rx.try_recv(), std::result::Result::Ok(()))); session.close().await.expect("reconnected Helius fixture session must close"); tokio::time::resume(); server.await.expect("Helius reconnect heartbeat fixture server must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_remote_close_consumes_bounded_reconnect_budget_before_terminal_failure() { 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"); websocket.send(tokio_tungstenite::tungstenite::Message::Close(std::option::Option::None)).await.expect("fixture Close must send"); }); let settings = reconnect_session_settings(1, std::time::Duration::from_millis(20), crate::WsResubscribePolicy::ActiveSubscriptions); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); wait_for_state(&session, crate::WsSessionState::Failed).await; assert_eq!(session.snapshot().continuity_gap_count(), 1); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_explicit_close_cancels_pending_and_finishes_with_hostile_peer() { let (listener, url) = bind_local_listener().await; let (request_seen_tx, request_seen_rx) = tokio::sync::oneshot::channel(); 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 _ = read_request(&mut websocket).await; let _ = request_seen_tx.send(()); tokio::time::sleep(std::time::Duration::from_millis(500)).await; }); let settings = session_settings(std::time::Duration::from_secs(1), std::time::Duration::from_millis(150), 8, 64 * 1024, 64 * 1024, 64 * 1024); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); let request_session = session.clone(); let pending = tokio::spawn(async move { return request_session.execute_json_rpc("neverRespond", std::vec::Vec::new()).await; }); request_seen_rx.await.expect("hostile server must observe request"); wait_for_pending_count(&session, 1).await; let close_result = tokio::time::timeout(std::time::Duration::from_millis(300), session.close()).await; close_result.expect("explicit close must remain bounded").expect("explicit close must finish cleanly"); assert_eq!(session.state(), crate::WsSessionState::Closed); assert_eq!(session.snapshot().pending_request_count(), 0); let pending_error = pending.await.expect("pending request task must join").expect_err("shutdown must cancel pending request"); assert_eq!(pending_error.code(), crate::ERROR_CODE_WS_SESSION_CLOSED); server.await.expect("hostile server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_repeated_connect_close_cycles_are_bounded() { const CYCLES: usize = 8; let (listener, url) = bind_local_listener().await; let server = tokio::spawn(async move { for _ in 0..CYCLES { 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 message = websocket.next().await.expect("client Close must arrive").expect("client Close must decode"); assert!(matches!(message, tokio_tungstenite::tungstenite::Message::Close(_))); let _ = websocket.flush().await; } }); let settings = session_settings(std::time::Duration::from_secs(1), std::time::Duration::from_millis(200), 8, 64 * 1024, 64 * 1024, 64 * 1024); for _ in 0..CYCLES { let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings.clone())).await.expect("client must connect"); session.close().await.expect("client close must succeed"); assert_eq!(session.state(), crate::WsSessionState::Closed); } tokio::time::timeout(std::time::Duration::from_secs(2), server) .await .expect("repeated close server must remain bounded") .expect("repeated close server task must join"); } #[tokio::test(flavor = "current_thread")] async fn websocket_reconnect_resubscribes_in_local_id_order_and_remaps_remote_ids() { let (listener, url) = bind_local_listener().await; let server = tokio::spawn(async move { let (first_stream, _) = listener.accept().await.expect("initial client must connect"); let mut first = tokio_tungstenite::accept_async(first_stream).await.expect("initial handshake must succeed"); let account = read_request(&mut first).await; assert_eq!(account.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("accountSubscribe")); send_result(&mut first, &account, serde_json::json!(101)).await; let logs = read_request(&mut first).await; assert_eq!(logs.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("logsSubscribe")); send_result(&mut first, &logs, serde_json::json!(202)).await; drop(first); let (second_stream, _) = listener.accept().await.expect("replacement client must connect"); let mut second = tokio_tungstenite::accept_async(second_stream).await.expect("replacement handshake must succeed"); let restored_account = read_request(&mut second).await; assert_eq!(restored_account.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("accountSubscribe")); assert_eq!(restored_account.get("params"), account.get("params")); send_result(&mut second, &restored_account, serde_json::json!(301)).await; let restored_logs = read_request(&mut second).await; assert_eq!(restored_logs.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("logsSubscribe")); assert_eq!(restored_logs.get("params"), logs.get("params")); send_result(&mut second, &restored_logs, serde_json::json!(302)).await; send_notification(&mut second, "logsNotification", 302, serde_json::json!({"generation": 2, "family": "logs"})).await; send_notification(&mut second, "accountNotification", 301, serde_json::json!({"generation": 2, "family": "account"})).await; let request = read_request(&mut second).await; send_result(&mut second, &request, serde_json::json!(true)).await; }); let settings = reconnect_session_settings(2, std::time::Duration::from_millis(20), crate::WsResubscribePolicy::ActiveSubscriptions); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); let mut account = session .subscribe_typed(crate::WsSubscriptionKind::Account, std::vec![serde_json::json!("account-canary")], |value| { return std::result::Result::Ok(value); }) .await .expect("account subscribe must succeed"); let mut logs = session .subscribe_typed(crate::WsSubscriptionKind::Logs, std::vec![serde_json::json!({"mentions": ["log-canary"]})], |value| { return std::result::Result::Ok(value); }) .await .expect("logs subscribe must succeed"); let account_id = account.id(); let logs_id = logs.id(); let logs_value = tokio::time::timeout(std::time::Duration::from_secs(2), logs.recv()) .await .expect("logs resubscribe notification must remain bounded") .expect("logs channel must remain open") .expect("logs notification must decode"); let account_value = tokio::time::timeout(std::time::Duration::from_secs(2), account.recv()) .await .expect("account resubscribe notification must remain bounded") .expect("account channel must remain open") .expect("account notification must decode"); assert_eq!(logs_value, serde_json::json!({"generation": 2, "family": "logs"})); assert_eq!(account_value, serde_json::json!({"generation": 2, "family": "account"})); assert_eq!(account.id(), account_id); assert_eq!(logs.id(), logs_id); assert_eq!(account.state(), crate::WsSubscriptionState::Active); assert_eq!(logs.state(), crate::WsSubscriptionState::Active); assert_eq!(session.snapshot().continuity_gap_count(), 1); assert!(session.snapshot().subscriptions().iter().all(|snapshot| return snapshot.remote_bound())); let result = session.execute_json_rpc("afterResubscribe", std::vec::Vec::new()).await.expect("reconnected session must remain usable"); assert_eq!(result, serde_json::json!(true)); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_unsubscribe_during_reconnect_prevents_resubscribe() { let (listener, url) = bind_local_listener().await; let (dropped_tx, dropped_rx) = tokio::sync::oneshot::channel(); let server = tokio::spawn(async move { let (first_stream, _) = listener.accept().await.expect("initial client must connect"); let mut first = tokio_tungstenite::accept_async(first_stream).await.expect("initial handshake must succeed"); let subscribe = read_request(&mut first).await; send_result(&mut first, &subscribe, serde_json::json!(41)).await; drop(first); let _ = dropped_tx.send(()); let (second_stream, _) = listener.accept().await.expect("replacement client must connect"); let mut second = tokio_tungstenite::accept_async(second_stream).await.expect("replacement handshake must succeed"); let first_request = read_request(&mut second).await; assert_eq!(first_request.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("afterReconnectCancel")); send_result(&mut second, &first_request, serde_json::json!(true)).await; }); let settings = reconnect_session_settings(2, std::time::Duration::from_millis(120), crate::WsResubscribePolicy::ActiveSubscriptions); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); let mut subscription = session .subscribe_typed(crate::WsSubscriptionKind::Root, std::vec::Vec::new(), |value| { return std::result::Result::Ok(value); }) .await .expect("root subscribe must succeed"); dropped_rx.await.expect("fixture must drop initial socket"); wait_for_state(&session, crate::WsSessionState::Reconnecting { attempt: 1 }).await; assert!(!subscription.unsubscribe().await.expect("local cancellation during reconnect must succeed")); assert_eq!(subscription.state(), crate::WsSubscriptionState::Closed); wait_for_gap_count(&session, 1).await; assert_eq!(session.snapshot().subscription_count(), 0); let result = session .execute_json_rpc("afterReconnectCancel", std::vec::Vec::new()) .await .expect("session must remain usable after reconnect cancellation"); assert_eq!(result, serde_json::json!(true)); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_late_resubscribe_ack_after_unsubscribe_is_cleaned_up_without_reactivation() { let (listener, url) = bind_local_listener().await; let (resubscribe_seen_tx, resubscribe_seen_rx) = tokio::sync::oneshot::channel(); let (release_ack_tx, release_ack_rx) = tokio::sync::oneshot::channel(); let server = tokio::spawn(async move { let (first_stream, _) = listener.accept().await.expect("initial client must connect"); let mut first = tokio_tungstenite::accept_async(first_stream).await.expect("initial handshake must succeed"); let subscribe = read_request(&mut first).await; send_result(&mut first, &subscribe, serde_json::json!(41)).await; drop(first); let (second_stream, _) = listener.accept().await.expect("replacement client must connect"); let mut second = tokio_tungstenite::accept_async(second_stream).await.expect("replacement handshake must succeed"); let resubscribe = read_request(&mut second).await; assert_eq!(resubscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("rootSubscribe")); let _ = resubscribe_seen_tx.send(()); release_ack_rx.await.expect("test must release stale acknowledgement"); send_result(&mut second, &resubscribe, serde_json::json!(99)).await; let cleanup = read_request(&mut second).await; assert_eq!(cleanup.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("rootUnsubscribe")); assert_eq!(cleanup.get("params"), std::option::Option::Some(&serde_json::json!([99]))); send_result(&mut second, &cleanup, serde_json::json!(true)).await; let request = read_request(&mut second).await; assert_eq!(request.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("afterStaleAck")); send_result(&mut second, &request, serde_json::json!(true)).await; }); let settings = reconnect_session_settings(2, std::time::Duration::from_millis(20), crate::WsResubscribePolicy::ActiveSubscriptions); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); let mut subscription = session .subscribe_typed(crate::WsSubscriptionKind::Root, std::vec::Vec::new(), |value| { return std::result::Result::Ok(value); }) .await .expect("root subscribe must succeed"); resubscribe_seen_rx.await.expect("fixture must observe replacement subscribe"); assert!(!subscription.unsubscribe().await.expect("local cancellation must win while resubscribe acknowledgement is pending")); assert_eq!(subscription.state(), crate::WsSubscriptionState::Closed); release_ack_tx.send(()).expect("fixture stale acknowledgement must be released"); wait_for_gap_count(&session, 1).await; assert_eq!(subscription.state(), crate::WsSubscriptionState::Closed); assert_eq!(session.snapshot().subscription_count(), 0); let result = session.execute_json_rpc("afterStaleAck", std::vec::Vec::new()).await.expect("session must remain usable after stale ack cleanup"); assert_eq!(result, serde_json::json!(true)); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_reconnect_budget_resets_after_full_active_recovery() { let (listener, url) = bind_local_listener().await; let server = tokio::spawn(async move { let (first_stream, _) = listener.accept().await.expect("initial client must connect"); let mut first = tokio_tungstenite::accept_async(first_stream).await.expect("initial handshake must succeed"); let first_request = read_request(&mut first).await; send_result(&mut first, &first_request, serde_json::json!(0)).await; drop(first); let (second_stream, _) = listener.accept().await.expect("first replacement must connect"); let mut second = tokio_tungstenite::accept_async(second_stream).await.expect("first replacement handshake must succeed"); let second_request = read_request(&mut second).await; send_result(&mut second, &second_request, serde_json::json!(1)).await; drop(second); let (third_stream, _) = listener.accept().await.expect("second replacement must connect"); let mut third = tokio_tungstenite::accept_async(third_stream).await.expect("second replacement handshake must succeed"); let third_request = read_request(&mut third).await; send_result(&mut third, &third_request, serde_json::json!(2)).await; }); let settings = reconnect_session_settings(1, std::time::Duration::from_millis(20), crate::WsResubscribePolicy::ActiveSubscriptions); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); let first = session.execute_json_rpc("cycleZero", std::vec::Vec::new()).await.expect("initial request must succeed"); assert_eq!(first, serde_json::json!(0)); wait_for_gap_count(&session, 1).await; let second = session.execute_json_rpc("cycleOne", std::vec::Vec::new()).await.expect("first recovered request must succeed"); assert_eq!(second, serde_json::json!(1)); wait_for_gap_count(&session, 2).await; let third = session.execute_json_rpc("cycleTwo", std::vec::Vec::new()).await.expect("second recovered request must succeed"); assert_eq!(third, serde_json::json!(2)); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_resubscribe_never_keeps_session_but_fails_logical_subscription() { let (listener, url) = bind_local_listener().await; let server = tokio::spawn(async move { let (first_stream, _) = listener.accept().await.expect("initial client must connect"); let mut first = tokio_tungstenite::accept_async(first_stream).await.expect("initial handshake must succeed"); let subscribe = read_request(&mut first).await; send_result(&mut first, &subscribe, serde_json::json!(71)).await; drop(first); let (second_stream, _) = listener.accept().await.expect("replacement client must connect"); let mut second = tokio_tungstenite::accept_async(second_stream).await.expect("replacement handshake must succeed"); let request = read_request(&mut second).await; assert_eq!(request.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("afterNeverPolicy")); send_result(&mut second, &request, serde_json::json!(true)).await; }); let settings = reconnect_session_settings(2, std::time::Duration::from_millis(20), crate::WsResubscribePolicy::Never); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); let subscription = session .subscribe_typed(crate::WsSubscriptionKind::Root, std::vec::Vec::new(), |value| { return std::result::Result::Ok(value); }) .await .expect("root subscribe must succeed"); wait_for_gap_count(&session, 1).await; assert_eq!(subscription.state(), crate::WsSubscriptionState::Failed); assert_eq!(subscription.terminal_error_code(), std::option::Option::Some(crate::ERROR_CODE_WS_CONNECTION_FAILED)); assert_eq!(session.snapshot().subscription_count(), 0); let result = session.execute_json_rpc("afterNeverPolicy", std::vec::Vec::new()).await.expect("physical session must reconnect without resubscribe"); assert_eq!(result, serde_json::json!(true)); server.await.expect("local server task must complete"); } #[test] fn websocket_reconnect_backoff_is_exponential_and_bounded() { let settings = crate::WsReconnectSettings::new(5, std::time::Duration::from_millis(10), std::time::Duration::from_millis(40)); assert_eq!(super::reconnect_backoff(&settings, 1), std::time::Duration::from_millis(10)); assert_eq!(super::reconnect_backoff(&settings, 2), std::time::Duration::from_millis(20)); assert_eq!(super::reconnect_backoff(&settings, 3), std::time::Duration::from_millis(40)); assert_eq!(super::reconnect_backoff(&settings, 4), std::time::Duration::from_millis(40)); } #[tokio::test(flavor = "current_thread")] async fn websocket_shutdown_interrupts_reconnect_backoff_without_new_connection() { let (listener, url) = bind_local_listener().await; let (dropped_tx, dropped_rx) = tokio::sync::oneshot::channel(); let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.expect("initial client must connect"); let websocket = tokio_tungstenite::accept_async(stream).await.expect("initial handshake must succeed"); drop(websocket); let _ = dropped_tx.send(()); tokio::time::sleep(std::time::Duration::from_millis(500)).await; }); let defaults = crate::WsSessionSettings::default(); let settings = crate::WsSessionSettings::new( std::time::Duration::from_millis(250), std::time::Duration::from_millis(150), crate::WsReconnectSettings::new(5, std::time::Duration::from_secs(2), std::time::Duration::from_secs(2)), 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(), ); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); dropped_rx.await.expect("fixture must drop physical connection"); wait_for_state(&session, crate::WsSessionState::Reconnecting { attempt: 1 }).await; tokio::time::timeout(std::time::Duration::from_millis(300), session.close()) .await .expect("shutdown must interrupt reconnect backoff") .expect("shutdown during reconnect must finish cleanly"); assert_eq!(session.state(), crate::WsSessionState::Closed); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_resubscribe_application_error_fails_only_one_subscription() { let (listener, url) = bind_local_listener().await; let server = tokio::spawn(async move { let (first_stream, _) = listener.accept().await.expect("initial client must connect"); let mut first = tokio_tungstenite::accept_async(first_stream).await.expect("initial handshake must succeed"); let account = read_request(&mut first).await; send_result(&mut first, &account, serde_json::json!(11)).await; let logs = read_request(&mut first).await; send_result(&mut first, &logs, serde_json::json!(12)).await; drop(first); let (second_stream, _) = listener.accept().await.expect("replacement client must connect"); let mut second = tokio_tungstenite::accept_async(second_stream).await.expect("replacement handshake must succeed"); let failed = read_request(&mut second).await; let failed_id = failed.get("id").and_then(serde_json::Value::as_u64).expect("resubscribe id must be numeric"); let error = serde_json::json!({"jsonrpc":"2.0","id":failed_id,"error":{"code":-32602,"message":"fixture resubscribe rejection"}}); second.send(tokio_tungstenite::tungstenite::Message::Text(error.to_string().into())).await.expect("resubscribe application error must send"); let restored = read_request(&mut second).await; assert_eq!(restored.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("logsSubscribe")); send_result(&mut second, &restored, serde_json::json!(88)).await; send_notification(&mut second, "logsNotification", 88, serde_json::json!({"restored": true})).await; let request = read_request(&mut second).await; send_result(&mut second, &request, serde_json::json!(true)).await; }); let settings = reconnect_session_settings(2, std::time::Duration::from_millis(20), crate::WsResubscribePolicy::ActiveSubscriptions); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); let account = session .subscribe_typed(crate::WsSubscriptionKind::Account, std::vec![serde_json::json!("account")], |value| { return std::result::Result::Ok(value); }) .await .expect("account subscribe must succeed"); let mut logs = session .subscribe_typed(crate::WsSubscriptionKind::Logs, std::vec![serde_json::json!("all")], |value| { return std::result::Result::Ok(value); }) .await .expect("logs subscribe must succeed"); let notification = tokio::time::timeout(std::time::Duration::from_secs(2), logs.recv()) .await .expect("restored logs notification must remain bounded") .expect("logs channel must remain open") .expect("logs notification must decode"); assert_eq!(notification, serde_json::json!({"restored": true})); assert_eq!(account.state(), crate::WsSubscriptionState::Failed); assert_eq!(account.terminal_error_code(), std::option::Option::Some(crate::ERROR_CODE_RPC_APPLICATION_ERROR)); assert_eq!(logs.state(), crate::WsSubscriptionState::Active); assert_eq!(session.state(), crate::WsSessionState::Active); assert_eq!(session.snapshot().continuity_gap_count(), 1); let result = session.execute_json_rpc("afterPartialResubscribe", std::vec::Vec::new()).await.expect("session must remain usable"); assert_eq!(result, serde_json::json!(true)); server.await.expect("local server task must complete"); } async fn send_notification( websocket: &mut tokio_tungstenite::WebSocketStream, method: &str, remote_subscription_id: u64, result: serde_json::Value, ) { let notification = serde_json::json!({ "jsonrpc": "2.0", "method": method, "params": { "result": result, "subscription": remote_subscription_id } }); websocket.send(tokio_tungstenite::tungstenite::Message::Text(notification.to_string().into())).await.expect("local notification must send"); } async fn wait_for_subscription_state(subscription: &crate::WsSubscription, expected: crate::WsSubscriptionState) { let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(1); loop { if subscription.state() == expected { return; } assert!(tokio::time::Instant::now() < deadline, "subscription did not reach expected state: {expected:?}"); tokio::time::sleep(std::time::Duration::from_millis(5)).await; } } #[tokio::test(flavor = "current_thread")] async fn websocket_generic_subscription_registers_remote_id_dispatches_typed_notification_and_unsubscribes() { 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.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("slotSubscribe")); send_result(&mut websocket, &subscribe, serde_json::json!(41)).await; send_notification(&mut websocket, "slotNotification", 41, serde_json::json!({"slot": 9001})).await; let unsubscribe = read_request(&mut websocket).await; assert_eq!(unsubscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("slotUnsubscribe")); assert_eq!(unsubscribe.get("params"), std::option::Option::Some(&serde_json::json!([41]))); send_result(&mut websocket, &unsubscribe, serde_json::json!(true)).await; }); let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed"); let mut subscription = session .subscribe_typed(crate::WsSubscriptionKind::Slot, std::vec::Vec::new(), |value| { return std::result::Result::Ok(value); }) .await .expect("generic slot subscribe must succeed"); assert_eq!(subscription.id().get(), 1); assert_eq!(subscription.kind(), crate::WsSubscriptionKind::Slot); assert_eq!(subscription.state(), crate::WsSubscriptionState::Active); let snapshot = session.snapshot(); assert_eq!(snapshot.subscription_count(), 1); assert_eq!(snapshot.subscriptions()[0].id(), subscription.id()); assert!(snapshot.subscriptions()[0].remote_bound()); let notification = subscription.recv().await.expect("typed notification channel must remain open").expect("notification must decode"); assert_eq!(notification, serde_json::json!({"slot": 9001})); assert!(subscription.unsubscribe().await.expect("unsubscribe must succeed")); assert_eq!(subscription.state(), crate::WsSubscriptionState::Closed); tokio::task::yield_now().await; assert_eq!(session.snapshot().subscription_count(), 0); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_remote_subscription_mapping_routes_multiple_families_to_stable_local_ids() { 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 first = read_request(&mut websocket).await; send_result(&mut websocket, &first, serde_json::json!(77)).await; let second = read_request(&mut websocket).await; send_result(&mut websocket, &second, serde_json::json!(12)).await; send_notification(&mut websocket, "logsNotification", 12, serde_json::json!({"family": "logs"})).await; send_notification(&mut websocket, "accountNotification", 77, serde_json::json!({"family": "account"})).await; }); let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed"); let mut account = session .subscribe_typed(crate::WsSubscriptionKind::Account, std::vec![serde_json::json!("account")], |value| { return std::result::Result::Ok(value); }) .await .expect("account subscribe must succeed"); let mut logs = session .subscribe_typed(crate::WsSubscriptionKind::Logs, std::vec![serde_json::json!("all")], |value| { return std::result::Result::Ok(value); }) .await .expect("logs subscribe must succeed"); assert_eq!(account.id().get(), 1); assert_eq!(logs.id().get(), 2); assert_eq!(logs.recv().await.expect("logs notification must exist").expect("logs notification must decode"), serde_json::json!({"family": "logs"})); assert_eq!( account.recv().await.expect("account notification must exist").expect("account notification must decode"), serde_json::json!({"family": "account"}) ); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_unknown_remote_subscription_notification_is_ignored_without_affecting_registered_subscription() { 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; send_result(&mut websocket, &subscribe, serde_json::json!(3)).await; send_notification(&mut websocket, "rootNotification", 999, serde_json::json!(100)).await; send_notification(&mut websocket, "rootNotification", 3, serde_json::json!(101)).await; let request = read_request(&mut websocket).await; send_result(&mut websocket, &request, serde_json::json!(true)).await; }); let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed"); let mut subscription = session .subscribe_typed(crate::WsSubscriptionKind::Root, std::vec::Vec::new(), |value| { return std::result::Result::Ok(value); }) .await .expect("root subscribe must succeed"); let notification = subscription.recv().await.expect("valid notification must exist").expect("valid notification must decode"); assert_eq!(notification, serde_json::json!(101)); assert_eq!(subscription.state(), crate::WsSubscriptionState::Active); assert_eq!(session.state(), crate::WsSessionState::Active); let result = session.execute_json_rpc("afterUnknownRemote", std::vec::Vec::new()).await.expect("physical session must remain usable"); assert_eq!(result, serde_json::json!(true)); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_notification_method_mismatch_fails_only_the_logical_subscription() { 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; send_result(&mut websocket, &subscribe, serde_json::json!(5)).await; send_notification(&mut websocket, "rootNotification", 5, serde_json::json!(1)).await; let cleanup = read_request(&mut websocket).await; assert_eq!(cleanup.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("slotUnsubscribe")); assert_eq!(cleanup.get("params"), std::option::Option::Some(&serde_json::json!([5]))); let request = read_request(&mut websocket).await; send_result(&mut websocket, &request, serde_json::json!(true)).await; }); let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed"); let mut subscription = session .subscribe_typed(crate::WsSubscriptionKind::Slot, std::vec::Vec::new(), |value| { return std::result::Result::Ok(value); }) .await .expect("slot subscribe must succeed"); wait_for_subscription_state(&subscription, crate::WsSubscriptionState::Failed).await; assert_eq!(subscription.terminal_error_code(), std::option::Option::Some(crate::ERROR_CODE_WS_PROTOCOL_ERROR)); assert!(subscription.recv().await.is_none()); assert_eq!(session.state(), crate::WsSessionState::Active); let result = session.execute_json_rpc("afterMismatch", std::vec::Vec::new()).await.expect("physical session must remain usable"); assert_eq!(result, serde_json::json!(true)); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_typed_notification_decode_failure_fails_only_one_subscription_and_surfaces_error() { 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; send_result(&mut websocket, &subscribe, serde_json::json!(8)).await; send_notification(&mut websocket, "rootNotification", 8, serde_json::json!("not-a-slot")).await; let cleanup = read_request(&mut websocket).await; assert_eq!(cleanup.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("rootUnsubscribe")); assert_eq!(cleanup.get("params"), std::option::Option::Some(&serde_json::json!([8]))); let request = read_request(&mut websocket).await; send_result(&mut websocket, &request, serde_json::json!(true)).await; }); let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed"); let mut subscription = session .subscribe_typed(crate::WsSubscriptionKind::Root, std::vec::Vec::new(), |value| { return match value.as_u64() { std::option::Option::Some(slot) => std::result::Result::Ok(slot), std::option::Option::None => { std::result::Result::Err(ksp_core_lib::Error::new(crate::ERROR_CODE_INVALID_RESPONSE, "fixture root notification must be numeric")) }, }; }) .await .expect("root subscribe must succeed"); let error = subscription.recv().await.expect("decode error must be delivered").expect_err("fixture payload must fail typed decoder"); assert_eq!(error.code(), crate::ERROR_CODE_INVALID_RESPONSE); wait_for_subscription_state(&subscription, crate::WsSubscriptionState::Failed).await; assert_eq!(subscription.terminal_error_code(), std::option::Option::Some(crate::ERROR_CODE_INVALID_RESPONSE)); assert_eq!(session.state(), crate::WsSessionState::Active); let result = session.execute_json_rpc("afterDecodeFailure", std::vec::Vec::new()).await.expect("physical session must remain usable"); assert_eq!(result, serde_json::json!(true)); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_notification_queue_overflow_fails_only_slow_subscription_and_cleans_remote_binding() { 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 slow_subscribe = read_request(&mut websocket).await; assert_eq!(slow_subscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("slotSubscribe")); send_result(&mut websocket, &slow_subscribe, serde_json::json!(41)).await; let healthy_subscribe = read_request(&mut websocket).await; assert_eq!(healthy_subscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("rootSubscribe")); send_result(&mut websocket, &healthy_subscribe, serde_json::json!(42)).await; send_notification(&mut websocket, "slotNotification", 41, serde_json::json!(1)).await; send_notification(&mut websocket, "slotNotification", 41, serde_json::json!(2)).await; let cleanup = read_request(&mut websocket).await; assert_eq!(cleanup.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("slotUnsubscribe")); assert_eq!(cleanup.get("params"), std::option::Option::Some(&serde_json::json!([41]))); send_notification(&mut websocket, "rootNotification", 42, serde_json::json!(99)).await; wait_for_close_frame(&mut websocket).await; }); let reconnect = crate::WsReconnectSettings::new(0, std::time::Duration::from_millis(10), std::time::Duration::from_millis(10)); let settings = subscription_session_settings(1, 2, reconnect); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); let mut slow = session .subscribe_typed(crate::WsSubscriptionKind::Slot, std::vec::Vec::new(), |value| { return std::result::Result::Ok(value); }) .await .expect("slow subscription must register"); let mut healthy = session .subscribe_typed(crate::WsSubscriptionKind::Root, std::vec::Vec::new(), |value| { return std::result::Result::Ok(value); }) .await .expect("healthy 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 first_slow = slow.recv().await.expect("first queued notification must remain observable").expect("first queued notification must decode"); assert_eq!(first_slow, serde_json::json!(1)); assert!(slow.recv().await.is_none()); let healthy_value = healthy.recv().await.expect("healthy notification must remain available").expect("healthy notification must decode"); assert_eq!(healthy_value, serde_json::json!(99)); assert_eq!(healthy.state(), crate::WsSubscriptionState::Active); assert_eq!(healthy.terminal_error_code(), std::option::Option::None); session.close().await.expect("session close must remain bounded"); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_active_subscription_limit_rejects_excess_without_leaking_capacity() { 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 first_subscribe = read_request(&mut websocket).await; assert_eq!(first_subscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("slotSubscribe")); send_result(&mut websocket, &first_subscribe, serde_json::json!(11)).await; let first_unsubscribe = read_request(&mut websocket).await; assert_eq!(first_unsubscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("slotUnsubscribe")); assert_eq!(first_unsubscribe.get("params"), std::option::Option::Some(&serde_json::json!([11]))); send_result(&mut websocket, &first_unsubscribe, serde_json::json!(true)).await; let replacement_subscribe = read_request(&mut websocket).await; assert_eq!(replacement_subscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("rootSubscribe")); send_result(&mut websocket, &replacement_subscribe, serde_json::json!(12)).await; let replacement_unsubscribe = read_request(&mut websocket).await; assert_eq!(replacement_unsubscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("rootUnsubscribe")); assert_eq!(replacement_unsubscribe.get("params"), std::option::Option::Some(&serde_json::json!([12]))); send_result(&mut websocket, &replacement_unsubscribe, serde_json::json!(true)).await; wait_for_close_frame(&mut websocket).await; }); let reconnect = crate::WsReconnectSettings::new(0, std::time::Duration::from_millis(10), std::time::Duration::from_millis(10)); let settings = subscription_session_settings(1, 1, reconnect); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); let mut first = session .subscribe_typed(crate::WsSubscriptionKind::Slot, std::vec::Vec::new(), |value| { return std::result::Result::Ok(value); }) .await .expect("first subscription must register"); let excess = session .subscribe_typed(crate::WsSubscriptionKind::Root, std::vec::Vec::new(), |value| { return std::result::Result::Ok(value); }) .await .expect_err("subscription above configured active limit must be rejected"); assert_eq!(excess.code(), crate::ERROR_CODE_WS_BACKPRESSURE_OVERFLOW); assert_eq!(session.snapshot().overflow_count(), 0); assert_eq!(session.snapshot().subscription_count(), 1); assert!(first.unsubscribe().await.expect("first unsubscribe must succeed")); assert_eq!(first.state(), crate::WsSubscriptionState::Closed); assert_eq!(first.terminal_error_code(), std::option::Option::None); wait_for_subscription_count(&session, 0).await; let mut replacement = session .subscribe_typed(crate::WsSubscriptionKind::Root, std::vec::Vec::new(), |value| { return std::result::Result::Ok(value); }) .await .expect("capacity must become reusable after terminal cleanup"); assert_eq!(replacement.id().get(), 2); assert!(replacement.unsubscribe().await.expect("replacement unsubscribe must succeed")); wait_for_subscription_count(&session, 0).await; session.close().await.expect("session close must remain bounded"); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn websocket_dropped_notification_receiver_triggers_remote_cleanup_and_releases_capacity() { let (listener, url) = bind_local_listener().await; let (trigger_tx, trigger_rx) = tokio::sync::oneshot::channel::<()>(); 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 first_subscribe = read_request(&mut websocket).await; send_result(&mut websocket, &first_subscribe, serde_json::json!(71)).await; trigger_rx.await.expect("client must signal receiver drop"); send_notification(&mut websocket, "slotNotification", 71, serde_json::json!(1)).await; let cleanup = read_request(&mut websocket).await; assert_eq!(cleanup.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("slotUnsubscribe")); assert_eq!(cleanup.get("params"), std::option::Option::Some(&serde_json::json!([71]))); let replacement_subscribe = read_request(&mut websocket).await; assert_eq!(replacement_subscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("rootSubscribe")); send_result(&mut websocket, &replacement_subscribe, serde_json::json!(72)).await; let replacement_unsubscribe = read_request(&mut websocket).await; send_result(&mut websocket, &replacement_unsubscribe, serde_json::json!(true)).await; wait_for_close_frame(&mut websocket).await; }); let reconnect = crate::WsReconnectSettings::new(0, std::time::Duration::from_millis(10), std::time::Duration::from_millis(10)); let settings = subscription_session_settings(1, 1, reconnect); let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); let subscription = session .subscribe_typed(crate::WsSubscriptionKind::Slot, std::vec::Vec::new(), |value| { return std::result::Result::Ok(value); }) .await .expect("first subscription must register"); drop(subscription); trigger_tx.send(()).expect("receiver-drop trigger must send"); wait_for_subscription_count(&session, 0).await; assert_eq!(session.snapshot().overflow_count(), 0); let mut replacement = session .subscribe_typed(crate::WsSubscriptionKind::Root, std::vec::Vec::new(), |value| { return std::result::Result::Ok(value); }) .await .expect("capacity must be reusable after dropped receiver cleanup"); assert!(replacement.unsubscribe().await.expect("replacement unsubscribe must succeed")); session.close().await.expect("session close must remain bounded"); server.await.expect("local server task must complete"); }