// file: crates/ksp-onchain-transport-lib/unit_tests/ws_helius_standard.rs // version: 1 use futures_util::SinkExt; // rust-rules: trait-import use futures_util::StreamExt; // rust-rules: trait-import fn helius_endpoint(url: &str) -> crate::WsEndpointSettings { return crate::WsEndpointSettings::new( "local_helius_standard_fixture", true, crate::WsProviderName::new("helius"), crate::WsClusterName::new("local"), crate::WsProtocolKind::HeliusLaserStream, crate::WsEndpointUrl::parse(url).expect("local Helius 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) -> 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 expect_pair( websocket: &mut tokio_tungstenite::WebSocketStream, subscribe_method: &str, expected_params: serde_json::Value, unsubscribe_method: &str, remote_id: u64, ) { let subscribe = read_request(websocket).await; assert_eq!(subscribe["method"], serde_json::Value::String(subscribe_method.to_owned())); assert_eq!(subscribe["params"], expected_params); send_result(websocket, &subscribe, serde_json::json!(remote_id)).await; let unsubscribe = read_request(websocket).await; assert_eq!(unsubscribe["method"], serde_json::Value::String(unsubscribe_method.to_owned())); assert_eq!(unsubscribe["params"], serde_json::json!([remote_id])); send_result(websocket, &unsubscribe, serde_json::json!(true)).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, } } } #[tokio::test(flavor = "current_thread")] async fn helius_facade_reuses_exact_standard_wire_for_all_six_supported_families() { 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"); expect_pair( &mut websocket, "accountSubscribe", serde_json::json!(["11111111111111111111111111111111", {"encoding":"base64","commitment":"confirmed"}]), "accountUnsubscribe", 101, ) .await; expect_pair( &mut websocket, "programSubscribe", serde_json::json!(["11111111111111111111111111111111", {"encoding":"jsonParsed","filters":[{"dataSize":80}],"withContext":true}]), "programUnsubscribe", 102, ) .await; expect_pair(&mut websocket, "logsSubscribe", serde_json::json!(["all", {"commitment":"finalized"}]), "logsUnsubscribe", 103).await; expect_pair( &mut websocket, "signatureSubscribe", serde_json::json!(["fixture-signature", {"commitment":"confirmed","enableReceivedNotification":true}]), "signatureUnsubscribe", 104, ) .await; expect_pair(&mut websocket, "slotSubscribe", serde_json::json!([]), "slotUnsubscribe", 105).await; expect_pair(&mut websocket, "rootSubscribe", serde_json::json!([]), "rootUnsubscribe", 106).await; wait_for_close_frame(&mut websocket).await; }); let session = crate::HeliusLaserStreamWsSession::connect(helius_endpoint(url.as_str())).await.expect("Helius facade must connect"); let pubkey = "11111111111111111111111111111111".parse::().expect("fixture pubkey must parse"); let account_config = crate::SolanaAccountSubscribeConfig::new( std::option::Option::Some(crate::SolanaAccountEncoding::Base64), std::option::Option::None, std::option::Option::Some(crate::SolanaCommitment::Confirmed), ); let mut account = session.account_subscribe(&pubkey, std::option::Option::Some(&account_config)).await.expect("Helius accountSubscribe must register"); assert!(account.unsubscribe().await.expect("Helius accountUnsubscribe must complete")); let program_config = crate::SolanaProgramSubscribeConfig::new( crate::SolanaAccountSubscribeConfig::new( std::option::Option::Some(crate::SolanaAccountEncoding::JsonParsed), std::option::Option::None, std::option::Option::None, ), std::vec![crate::SolanaProgramAccountFilter::DataSize(80)], std::option::Option::Some(true), ); let mut program = session.program_subscribe(&pubkey, std::option::Option::Some(&program_config)).await.expect("Helius programSubscribe must register"); assert!(program.unsubscribe().await.expect("Helius programUnsubscribe must complete")); let logs_config = crate::SolanaCommitmentConfig::new(std::option::Option::Some(crate::SolanaCommitment::Finalized)); let mut logs = session .logs_subscribe(&crate::SolanaLogsSubscribeFilter::All, std::option::Option::Some(&logs_config)) .await .expect("Helius logsSubscribe must register"); assert!(logs.unsubscribe().await.expect("Helius logsUnsubscribe must complete")); let signature_config = crate::SolanaSignatureSubscribeConfig::new(std::option::Option::Some(crate::SolanaCommitment::Confirmed), std::option::Option::Some(true)); let mut signature = session .signature_subscribe("fixture-signature", std::option::Option::Some(&signature_config)) .await .expect("Helius signatureSubscribe must register"); assert!(signature.unsubscribe().await.expect("Helius signatureUnsubscribe must complete")); let mut slot = session.slot_subscribe().await.expect("Helius slotSubscribe must register"); assert!(slot.unsubscribe().await.expect("Helius slotUnsubscribe must complete")); let mut root = session.root_subscribe().await.expect("Helius rootSubscribe must register"); assert!(root.unsubscribe().await.expect("Helius rootUnsubscribe must complete")); session.close().await.expect("Helius facade close must complete"); server.await.expect("local Helius peer task must complete"); }