// file: crates/ksp-onchain-transport-lib/unit_tests/ws_transactions.rs // version: 2 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_transactions", 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) -> serde_json::Value { let message = websocket.next().await.expect("request message must exist").expect("request message must decode"); let text = message.to_text().expect("request must be text"); return serde_json::from_str(text).expect("request must contain JSON"); } async fn send_result(websocket: &mut tokio_tungstenite::WebSocketStream, request: &serde_json::Value, result: serde_json::Value) { let id = request.get("id").and_then(serde_json::Value::as_u64).expect("request id must be numeric"); let response = serde_json::json!({"jsonrpc":"2.0","id":id,"result":result}); websocket.send(tokio_tungstenite::tungstenite::Message::Text(response.to_string().into())).await.expect("local response must send"); } async fn wait_for_close_frame(websocket: &mut tokio_tungstenite::WebSocketStream) { loop { let message = websocket.next().await; match message { std::option::Option::Some(std::result::Result::Ok(tokio_tungstenite::tungstenite::Message::Close(_))) => return, std::option::Option::Some(std::result::Result::Ok(_)) => {}, std::option::Option::Some(std::result::Result::Err(_)) | std::option::Option::None => return, } } } #[test] fn logs_subscribe_filters_preserve_all_all_with_votes_and_exactly_one_mention() { let pubkey = "11111111111111111111111111111111".parse::().expect("mention fixture must parse"); assert_eq!(crate::SolanaLogsSubscribeFilter::All.to_json_value(), serde_json::json!("all")); assert_eq!(crate::SolanaLogsSubscribeFilter::AllWithVotes.to_json_value(), serde_json::json!("allWithVotes")); assert_eq!(crate::SolanaLogsSubscribeFilter::Mentions(pubkey).to_json_value(), serde_json::json!({"mentions":["11111111111111111111111111111111"]})); } #[test] fn logs_notification_decoder_preserves_context_signature_nullable_error_and_ordered_logs() { let success = super::decode_logs_notification( "logsSubscribe", serde_json::json!({ "context":{"slot":81,"apiVersion":"4.2.1"}, "value":{"signature":"fixture-signature","err":null,"logs":["first","second"]} }), ) .expect("successful logs notification must decode"); assert_eq!(success.context().slot(), 81); assert_eq!(success.value().signature(), "fixture-signature"); assert!(success.value().err().is_none()); assert_eq!(success.value().logs(), &["first".to_owned(), "second".to_owned()]); let failed = super::decode_logs_notification( "logsSubscribe", serde_json::json!({"context":{"slot":82},"value":{"signature":"fixture-signature-2","err":{"InstructionError":[0,"Custom"]},"logs":[]}}), ) .expect("failed logs notification must preserve transaction error wire value"); assert_eq!(failed.context().slot(), 82); assert!(failed.value().err().is_some()); let missing_err = super::decode_logs_notification("logsSubscribe", serde_json::json!({"context":{"slot":83},"value":{"signature":"fixture-signature-3","logs":[]}})); assert!(missing_err.is_err()); } #[tokio::test(flavor = "current_thread")] async fn stable_logs_wrapper_uses_exact_filter_config_notification_and_handle_unsubscribe() { let (listener, url) = bind_local_listener().await; let server = tokio::spawn(async move { let (stream, _) = listener.accept().await.expect("local server must accept client"); let mut websocket = tokio_tungstenite::accept_async(stream).await.expect("local WebSocket handshake must succeed"); let subscribe = read_request(&mut websocket).await; assert_eq!(subscribe["method"], serde_json::json!("logsSubscribe")); assert_eq!(subscribe["params"], serde_json::json!([{"mentions":["11111111111111111111111111111111"]},{"commitment":"finalized"}])); send_result(&mut websocket, &subscribe, serde_json::json!(88)).await; let id = subscribe.get("id").and_then(serde_json::Value::as_u64).expect("subscribe request id must exist"); assert!(id > 0); let notification = serde_json::json!({ "jsonrpc":"2.0", "method":"logsNotification", "params":{ "result":{"context":{"slot":900},"value":{"signature":"fixture-signature","err":null,"logs":["Program fixture success"]}}, "subscription":88 } }); websocket.send(tokio_tungstenite::tungstenite::Message::Text(notification.to_string().into())).await.expect("logs notification must send"); let unsubscribe = read_request(&mut websocket).await; assert_eq!(unsubscribe["method"], serde_json::json!("logsUnsubscribe")); assert_eq!(unsubscribe["params"], serde_json::json!([88])); send_result(&mut websocket, &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 mention = "11111111111111111111111111111111".parse::().expect("mention fixture must parse"); let filter = crate::SolanaLogsSubscribeFilter::Mentions(mention); let config = crate::SolanaCommitmentConfig::new(std::option::Option::Some(crate::SolanaCommitment::Finalized)); let mut subscription = session.logs_subscribe(&filter, std::option::Option::Some(&config)).await.expect("logsSubscribe must register"); assert_eq!(subscription.kind(), crate::WsSubscriptionKind::Logs); let notification = subscription.recv().await.expect("logs notification must arrive").expect("logs notification must decode"); assert_eq!(notification.context().slot(), 900); assert_eq!(notification.value().signature(), "fixture-signature"); assert_eq!(notification.value().logs(), &["Program fixture success".to_owned()]); assert!(subscription.unsubscribe().await.expect("logs unsubscribe must complete")); session.close().await.expect("session close must complete"); server.await.expect("local server task must complete"); } #[test] fn signature_subscribe_config_and_decoder_preserve_all_documented_wire_variants() { let config = crate::SolanaSignatureSubscribeConfig::new(std::option::Option::Some(crate::SolanaCommitment::Confirmed), std::option::Option::Some(false)); assert_eq!(config.commitment(), std::option::Option::Some(crate::SolanaCommitment::Confirmed)); assert_eq!(config.enable_received_notification(), std::option::Option::Some(false)); assert_eq!(config.to_json_value(), serde_json::json!({"commitment":"confirmed","enableReceivedNotification":false})); let received = super::decode_signature_notification("signatureSubscribe", serde_json::json!({"context":{"slot":90},"value":"receivedSignature"})) .expect("receivedSignature notification must decode"); assert_eq!(received.context().slot(), 90); assert_eq!(*received.value(), crate::SolanaSignatureNotification::ReceivedSignature); assert!(!received.value().is_terminal()); let success = super::decode_signature_notification("signatureSubscribe", serde_json::json!({"context":{"slot":91},"value":{"err":null}})) .expect("terminal successful signature notification must decode"); assert!(success.value().is_terminal()); assert!(success.value().err().is_none()); let failure = super::decode_signature_notification( "signatureSubscribe", serde_json::json!({"context":{"slot":92},"value":{"err":{"InstructionError":[0,"Custom"]}}}), ) .expect("terminal failed signature notification must decode"); assert!(failure.value().is_terminal()); assert!(failure.value().err().is_some()); assert!(super::decode_signature_notification("signatureSubscribe", serde_json::json!({"context":{"slot":93},"value":"futureVariant"})).is_err()); assert!(super::decode_signature_notification("signatureSubscribe", serde_json::json!({"context":{"slot":94},"value":{}})).is_err()); } #[tokio::test(flavor = "current_thread")] async fn signature_unsubscribe_before_terminal_notification_uses_current_remote_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 subscribe = read_request(&mut websocket).await; assert_eq!(subscribe["method"], serde_json::json!("signatureSubscribe")); assert_eq!(subscribe["params"], serde_json::json!(["fixture-signature"])); send_result(&mut websocket, &subscribe, serde_json::json!(301)).await; let unsubscribe = read_request(&mut websocket).await; assert_eq!(unsubscribe["method"], serde_json::json!("signatureUnsubscribe")); assert_eq!(unsubscribe["params"], serde_json::json!([301])); send_result(&mut websocket, &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 empty_config = crate::SolanaSignatureSubscribeConfig::default(); let mut subscription = session.signature_subscribe("fixture-signature", std::option::Option::Some(&empty_config)).await.expect("signatureSubscribe must register"); assert!(subscription.unsubscribe().await.expect("signature unsubscribe must complete")); assert_eq!(subscription.state(), crate::WsSubscriptionState::Closed); session.close().await.expect("session close must complete"); server.await.expect("local server task must complete"); } #[tokio::test(flavor = "current_thread")] async fn signature_terminal_notification_closes_handle_and_is_not_resubscribed_after_reconnect() { let (listener, url) = bind_local_listener().await; let (send_terminal_tx, send_terminal_rx) = tokio::sync::oneshot::channel(); let (replacement_ready_tx, replacement_ready_rx) = tokio::sync::oneshot::channel(); 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 WebSocket handshake must succeed"); let subscribe = read_request(&mut websocket).await; assert_eq!(subscribe["method"], serde_json::json!("signatureSubscribe")); assert_eq!(subscribe["params"], serde_json::json!(["fixture-signature",{"commitment":"finalized","enableReceivedNotification":true}])); send_result(&mut websocket, &subscribe, serde_json::json!(401)).await; let received = serde_json::json!({ "jsonrpc":"2.0", "method":"signatureNotification", "params":{"result":{"context":{"slot":100},"value":"receivedSignature"},"subscription":401} }); websocket .send(tokio_tungstenite::tungstenite::Message::Text(received.to_string().into())) .await .expect("receivedSignature notification must send"); send_terminal_rx.await.expect("client must observe early signature notification before terminal send"); let terminal = serde_json::json!({ "jsonrpc":"2.0", "method":"signatureNotification", "params":{"result":{"context":{"slot":101},"value":{"err":null}},"subscription":401} }); websocket .send(tokio_tungstenite::tungstenite::Message::Text(terminal.to_string().into())) .await .expect("terminal signature notification must send"); let unexpected_cleanup = tokio::time::timeout(std::time::Duration::from_millis(100), websocket.next()).await; assert!(unexpected_cleanup.is_err(), "server-terminal signature notification must not trigger signatureUnsubscribe"); drop(websocket); 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 WebSocket handshake must succeed"); let unexpected = tokio::time::timeout(std::time::Duration::from_millis(100), replacement.next()).await; assert!(unexpected.is_err(), "terminal signature subscription must not be replayed after reconnect"); replacement_ready_tx.send(()).expect("replacement-ready signal must send"); wait_for_close_frame(&mut replacement).await; }); let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed"); let config = crate::SolanaSignatureSubscribeConfig::new(std::option::Option::Some(crate::SolanaCommitment::Finalized), std::option::Option::Some(true)); let mut subscription = session.signature_subscribe("fixture-signature", std::option::Option::Some(&config)).await.expect("signatureSubscribe must register"); let received = subscription.recv().await.expect("receivedSignature must arrive").expect("receivedSignature must decode"); assert_eq!(*received.value(), crate::SolanaSignatureNotification::ReceivedSignature); assert_eq!(subscription.state(), crate::WsSubscriptionState::Active); send_terminal_tx.send(()).expect("terminal-send signal must reach fixture"); let terminal = subscription.recv().await.expect("terminal signature notification must arrive").expect("terminal signature notification must decode"); assert!(terminal.value().is_terminal()); assert!(terminal.value().err().is_none()); let closed = tokio::time::timeout(std::time::Duration::from_secs(1), async { loop { if subscription.state() == crate::WsSubscriptionState::Closed { return; } tokio::task::yield_now().await; } }) .await; assert!(closed.is_ok()); assert!(subscription.recv().await.is_none()); assert!(subscription.terminal_error_code().is_none()); assert!(!subscription.unsubscribe().await.expect("already terminal signature unsubscribe must be local-only")); replacement_ready_rx.await.expect("replacement connection must be observed without signature replay"); assert_eq!(session.snapshot().subscription_count(), 0); assert!(session.snapshot().continuity_gap_count() >= 1); session.close().await.expect("session close must complete"); server.await.expect("local server task must complete"); }