Files
khadhroony-solana-project/crates/ksp-onchain-transport-lib/unit_tests/ws_cluster.rs
2026-08-23 00:12:39 +02:00

101 lines
6.1 KiB
Rust

// file: crates/ksp-onchain-transport-lib/unit_tests/ws_cluster.rs
// version: 1
use futures_util::SinkExt; // rust-rules: trait-import
use futures_util::StreamExt; // rust-rules: trait-import
fn local_endpoint(url: &str) -> crate::WsEndpointSettings {
return crate::WsEndpointSettings::new(
"local_ws_cluster",
true,
crate::WsProviderName::new("local-fixture"),
crate::WsClusterName::new("local"),
crate::WsProtocolKind::SolanaStandard,
crate::WsEndpointUrl::parse(url).expect("local test WebSocket URL must parse"),
crate::WsSessionSettings::default(),
);
}
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 send_notification(websocket: &mut tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>, method: &str, remote_id: u64, result: serde_json::Value) {
let notification = serde_json::json!({"jsonrpc":"2.0","method":method,"params":{"result":result,"subscription":remote_id}});
websocket.send(tokio_tungstenite::tungstenite::Message::Text(notification.to_string().into())).await.expect("local notification must send");
}
async fn wait_for_close_frame(websocket: &mut tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>) {
loop {
let message = websocket.next().await;
match message {
std::option::Option::Some(std::result::Result::Ok(tokio_tungstenite::tungstenite::Message::Close(_))) => return,
std::option::Option::Some(std::result::Result::Ok(_)) => {},
std::option::Option::Some(std::result::Result::Err(_)) | std::option::Option::None => return,
}
}
}
#[test]
fn slot_notification_decoder_preserves_slot_parent_and_root() {
let notification =
super::decode_slot_notification("slotSubscribe", serde_json::json!({"slot":76,"parent":75,"root":44})).expect("slot notification must decode");
assert_eq!(notification.slot(), 76);
assert_eq!(notification.parent(), 75);
assert_eq!(notification.root(), 44);
assert!(super::decode_slot_notification("slotSubscribe", serde_json::json!({"slot":76,"parent":75})).is_err());
}
#[tokio::test(flavor = "current_thread")]
async fn stable_slot_and_root_wrappers_use_no_params_decode_exact_notifications_and_unsubscribe_by_handle() {
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 slot_subscribe = read_request(&mut websocket).await;
assert_eq!(slot_subscribe["method"], serde_json::json!("slotSubscribe"));
assert_eq!(slot_subscribe["params"], serde_json::json!([]));
send_result(&mut websocket, &slot_subscribe, serde_json::json!(201)).await;
send_notification(&mut websocket, "slotNotification", 201, serde_json::json!({"slot":76,"parent":75,"root":44})).await;
let root_subscribe = read_request(&mut websocket).await;
assert_eq!(root_subscribe["method"], serde_json::json!("rootSubscribe"));
assert_eq!(root_subscribe["params"], serde_json::json!([]));
send_result(&mut websocket, &root_subscribe, serde_json::json!(202)).await;
send_notification(&mut websocket, "rootNotification", 202, serde_json::json!(42)).await;
let slot_unsubscribe = read_request(&mut websocket).await;
assert_eq!(slot_unsubscribe["method"], serde_json::json!("slotUnsubscribe"));
assert_eq!(slot_unsubscribe["params"], serde_json::json!([201]));
send_result(&mut websocket, &slot_unsubscribe, serde_json::json!(true)).await;
let root_unsubscribe = read_request(&mut websocket).await;
assert_eq!(root_unsubscribe["method"], serde_json::json!("rootUnsubscribe"));
assert_eq!(root_unsubscribe["params"], serde_json::json!([202]));
send_result(&mut websocket, &root_unsubscribe, serde_json::json!(true)).await;
wait_for_close_frame(&mut websocket).await;
});
let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed");
let mut slot_subscription = session.slot_subscribe().await.expect("slotSubscribe must register");
let slot = slot_subscription.recv().await.expect("slot notification must arrive").expect("slot notification must decode");
assert_eq!((slot.slot(), slot.parent(), slot.root()), (76, 75, 44));
let mut root_subscription = session.root_subscribe().await.expect("rootSubscribe must register");
let root = root_subscription.recv().await.expect("root notification must arrive").expect("root notification must decode");
assert_eq!(root, 42);
assert!(slot_subscription.unsubscribe().await.expect("slot unsubscribe must complete"));
assert!(root_subscription.unsubscribe().await.expect("root unsubscribe must complete"));
session.close().await.expect("session close must complete");
server.await.expect("local server task must complete");
}