1358 lines
81 KiB
Rust
1358 lines
81 KiB
Rust
// 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<tokio::net::TcpStream>) -> serde_json::Value {
|
|
let message = websocket.next().await.expect("request message must exist").expect("request message must decode");
|
|
let text = message.to_text().expect("request must be text");
|
|
return serde_json::from_str(text).expect("request must contain JSON");
|
|
}
|
|
|
|
async fn send_result(websocket: &mut tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>, request: &serde_json::Value, result: serde_json::Value) {
|
|
let id = request.get("id").and_then(serde_json::Value::as_u64).expect("request id must be numeric");
|
|
let response = serde_json::json!({"jsonrpc":"2.0","id":id,"result":result});
|
|
websocket.send(tokio_tungstenite::tungstenite::Message::Text(response.to_string().into())).await.expect("local response must send");
|
|
}
|
|
|
|
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<tokio::net::TcpStream>) {
|
|
loop {
|
|
let message = websocket.next().await;
|
|
match message {
|
|
std::option::Option::Some(std::result::Result::Ok(tokio_tungstenite::tungstenite::Message::Close(_))) => return,
|
|
std::option::Option::Some(std::result::Result::Ok(_)) => {},
|
|
std::option::Option::Some(std::result::Result::Err(_)) | std::option::Option::None => return,
|
|
}
|
|
}
|
|
}
|
|
|
|
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::<tokio::time::Instant>);
|
|
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<tokio::net::TcpStream>,
|
|
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<T>(subscription: &crate::WsSubscription<T>, 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");
|
|
}
|