431 lines
24 KiB
Rust
431 lines
24 KiB
Rust
// file: crates/ksp-job-backfill-lib/tests/ws_raw_parity.rs
|
|
// version: 1
|
|
|
|
//! Cross-layer parity canaries qualifying standard `blockSubscribe` and Helius `transactionSubscribe` for RAW v1 ingestion.
|
|
|
|
use futures_util::SinkExt; // rust-rules: trait-import
|
|
use futures_util::StreamExt; // rust-rules: trait-import
|
|
|
|
const FIXTURE_BLOCK_TIME: i64 = 1_787_072_400;
|
|
const FIXTURE_SLOT: u64 = 430_000_123;
|
|
const ZERO_SIGNATURE_TEXT: &str = "1111111111111111111111111111111111111111111111111111111111111111";
|
|
const ZERO_TRANSACTION_BASE64: &str = "AQAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA";
|
|
|
|
type FixtureWebSocket = tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>;
|
|
|
|
fn http_pool_for_url(url: &str) -> ksp_core_lib::Result<ksp_onchain_transport_lib::HttpTransportPool> {
|
|
let endpoint_url = match ksp_onchain_transport_lib::HttpEndpointUrl::parse(url) {
|
|
Ok(endpoint_url) => endpoint_url,
|
|
Err(error) => return Err(error),
|
|
};
|
|
let role = ksp_onchain_transport_lib::HttpEndpointRoleSettings::new(
|
|
ksp_onchain_transport_lib::HttpRoleName::new("default"),
|
|
true,
|
|
std::vec![ksp_onchain_transport_lib::HttpRequestKind::wildcard()],
|
|
10,
|
|
ksp_onchain_transport_lib::HttpRoleLimits::new(
|
|
std::option::Option::None,
|
|
std::option::Option::None,
|
|
std::option::Option::None,
|
|
std::option::Option::None,
|
|
),
|
|
);
|
|
let endpoint = ksp_onchain_transport_lib::HttpEndpointSettings::new(
|
|
"ws-parity-http-fixture",
|
|
true,
|
|
ksp_onchain_transport_lib::HttpProviderName::new("fixture-http-provider"),
|
|
ksp_onchain_transport_lib::HttpClusterName::new("devnet"),
|
|
endpoint_url,
|
|
std::time::Duration::from_secs(1),
|
|
std::time::Duration::from_secs(1),
|
|
std::option::Option::Some(1),
|
|
std::vec![role],
|
|
);
|
|
return ksp_onchain_transport_lib::HttpTransportPool::new(ksp_onchain_transport_lib::HttpTransportSettings::new(
|
|
std::vec![endpoint],
|
|
ksp_onchain_transport_lib::HttpRetrySettings::new(0, std::time::Duration::from_millis(1), std::time::Duration::from_millis(1)),
|
|
));
|
|
}
|
|
|
|
fn serve_http_once(body: &'static str) -> std::io::Result<(std::string::String, std::thread::JoinHandle<std::io::Result<()>>)> {
|
|
let listener = match std::net::TcpListener::bind("127.0.0.1:0") {
|
|
Ok(listener) => listener,
|
|
Err(error) => return Err(error),
|
|
};
|
|
let address = match listener.local_addr() {
|
|
Ok(address) => address,
|
|
Err(error) => return Err(error),
|
|
};
|
|
let handle = std::thread::spawn(move || {
|
|
let (mut stream, _) = match listener.accept() {
|
|
Ok(accepted) => accepted,
|
|
Err(error) => return Err(error),
|
|
};
|
|
let mut bytes = std::vec::Vec::new();
|
|
let mut buffer = [0_u8; 1024];
|
|
loop {
|
|
let count = match std::io::Read::read(&mut stream, &mut buffer) {
|
|
Ok(count) => count,
|
|
Err(error) => return Err(error),
|
|
};
|
|
if count == 0 {
|
|
break;
|
|
}
|
|
bytes.extend_from_slice(&buffer[..count]);
|
|
if bytes.windows(4).any(|window| return window == b"\r\n\r\n") {
|
|
break;
|
|
}
|
|
}
|
|
let response = format!("HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", body.len(), body);
|
|
if let Err(error) = std::io::Write::write_all(&mut stream, response.as_bytes()) {
|
|
return Err(error);
|
|
}
|
|
return Ok(());
|
|
});
|
|
return Ok((format!("http://{address}"), handle));
|
|
}
|
|
|
|
fn ws_endpoint(url: &str, protocol: ksp_onchain_transport_lib::WsProtocolKind) -> ksp_core_lib::Result<ksp_onchain_transport_lib::WsEndpointSettings> {
|
|
let parsed = match ksp_onchain_transport_lib::WsEndpointUrl::parse(url) {
|
|
Ok(parsed) => parsed,
|
|
Err(error) => return Err(error),
|
|
};
|
|
return Ok(ksp_onchain_transport_lib::WsEndpointSettings::new(
|
|
"ws-raw-parity-fixture",
|
|
true,
|
|
ksp_onchain_transport_lib::WsProviderName::new("fixture-ws-provider"),
|
|
ksp_onchain_transport_lib::WsClusterName::new("devnet"),
|
|
protocol,
|
|
parsed,
|
|
ksp_onchain_transport_lib::WsSessionSettings::default(),
|
|
));
|
|
}
|
|
|
|
async fn read_ws_request(websocket: &mut FixtureWebSocket) -> std::result::Result<serde_json::Value, std::string::String> {
|
|
let message = match websocket.next().await {
|
|
std::option::Option::Some(std::result::Result::Ok(message)) => message,
|
|
std::option::Option::Some(std::result::Result::Err(_)) => return Err("fixture WebSocket request decode failed".to_owned()),
|
|
std::option::Option::None => return Err("fixture WebSocket request was absent".to_owned()),
|
|
};
|
|
let text = match message.to_text() {
|
|
Ok(text) => text,
|
|
Err(_) => return Err("fixture WebSocket request was not text".to_owned()),
|
|
};
|
|
return serde_json::from_str(text).map_err(|_| return "fixture WebSocket request JSON was invalid".to_owned());
|
|
}
|
|
|
|
async fn send_ws_json(websocket: &mut FixtureWebSocket, value: serde_json::Value) -> std::result::Result<(), std::string::String> {
|
|
return websocket
|
|
.send(tokio_tungstenite::tungstenite::Message::Text(value.to_string().into()))
|
|
.await
|
|
.map_err(|_| return "fixture WebSocket response send failed".to_owned());
|
|
}
|
|
|
|
async fn send_ws_result(
|
|
websocket: &mut FixtureWebSocket,
|
|
request: &serde_json::Value,
|
|
result: serde_json::Value,
|
|
) -> std::result::Result<(), std::string::String> {
|
|
let id = match request.get("id").and_then(serde_json::Value::as_u64) {
|
|
std::option::Option::Some(id) => id,
|
|
std::option::Option::None => return Err("fixture request id was absent".to_owned()),
|
|
};
|
|
return send_ws_json(websocket, serde_json::json!({"jsonrpc":"2.0","id":id,"result":result})).await;
|
|
}
|
|
|
|
async fn wait_for_ws_close(websocket: &mut FixtureWebSocket) {
|
|
loop {
|
|
match websocket.next().await {
|
|
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,
|
|
}
|
|
}
|
|
}
|
|
|
|
fn map_meta(field: &ksp_onchain_transport_lib::SolanaWireField<serde_json::Value>) -> ksp_raw_transaction_lib::RawTransactionWireField<serde_json::Value> {
|
|
return match field {
|
|
ksp_onchain_transport_lib::SolanaWireField::Omitted => ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
|
ksp_onchain_transport_lib::SolanaWireField::Null => ksp_raw_transaction_lib::RawTransactionWireField::Null,
|
|
ksp_onchain_transport_lib::SolanaWireField::Value(value) => ksp_raw_transaction_lib::RawTransactionWireField::Value(value.clone()),
|
|
};
|
|
}
|
|
|
|
fn map_version(
|
|
field: &ksp_onchain_transport_lib::SolanaWireField<ksp_onchain_transport_lib::SolanaTransactionVersion>,
|
|
) -> ksp_raw_transaction_lib::RawTransactionWireField<ksp_raw_transaction_lib::RawTransactionVersion> {
|
|
return match field {
|
|
ksp_onchain_transport_lib::SolanaWireField::Omitted => ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
|
ksp_onchain_transport_lib::SolanaWireField::Null => ksp_raw_transaction_lib::RawTransactionWireField::Null,
|
|
ksp_onchain_transport_lib::SolanaWireField::Value(ksp_onchain_transport_lib::SolanaTransactionVersion::Legacy) => {
|
|
ksp_raw_transaction_lib::RawTransactionWireField::Value(ksp_raw_transaction_lib::RawTransactionVersion::Legacy)
|
|
},
|
|
ksp_onchain_transport_lib::SolanaWireField::Value(ksp_onchain_transport_lib::SolanaTransactionVersion::Number(number)) => {
|
|
ksp_raw_transaction_lib::RawTransactionWireField::Value(ksp_raw_transaction_lib::RawTransactionVersion::Number(*number))
|
|
},
|
|
};
|
|
}
|
|
|
|
fn raw_from_block_transaction(
|
|
network: ksp_store_lib::RawNetworkId,
|
|
slot: u64,
|
|
block_time: std::option::Option<i64>,
|
|
transaction: &ksp_onchain_transport_lib::SolanaBlockTransaction,
|
|
transaction_index: u32,
|
|
) -> std::option::Option<ksp_store_lib::RawTransaction> {
|
|
let data = match transaction.transaction() {
|
|
ksp_onchain_transport_lib::SolanaEncodedTransaction::Binary { data, encoding: ksp_onchain_transport_lib::SolanaTransactionBinaryEncoding::Base64 } => {
|
|
data.as_str()
|
|
},
|
|
_ => return std::option::Option::None,
|
|
};
|
|
let material = ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64_with_embedded_signature(
|
|
network,
|
|
slot,
|
|
block_time,
|
|
data,
|
|
map_meta(transaction.meta()),
|
|
map_version(transaction.version()),
|
|
ksp_raw_transaction_lib::RawTransactionWireField::Value(transaction_index),
|
|
)
|
|
.ok()?;
|
|
return ksp_raw_transaction_lib::canonicalize_raw_transaction(material).ok();
|
|
}
|
|
|
|
fn raw_from_confirmed_transaction(
|
|
network: ksp_store_lib::RawNetworkId,
|
|
transaction: &ksp_onchain_transport_lib::SolanaConfirmedTransaction,
|
|
) -> std::option::Option<ksp_store_lib::RawTransaction> {
|
|
let data = match transaction.transaction() {
|
|
ksp_onchain_transport_lib::SolanaEncodedTransaction::Binary { data, encoding: ksp_onchain_transport_lib::SolanaTransactionBinaryEncoding::Base64 } => {
|
|
data.as_str()
|
|
},
|
|
_ => return std::option::Option::None,
|
|
};
|
|
let material = ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64_with_embedded_signature(
|
|
network,
|
|
transaction.slot(),
|
|
transaction.block_time(),
|
|
data,
|
|
map_meta(transaction.meta()),
|
|
map_version(transaction.version()),
|
|
match transaction.transaction_index() {
|
|
ksp_onchain_transport_lib::SolanaWireField::Omitted => ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
|
ksp_onchain_transport_lib::SolanaWireField::Null => ksp_raw_transaction_lib::RawTransactionWireField::Null,
|
|
ksp_onchain_transport_lib::SolanaWireField::Value(value) => ksp_raw_transaction_lib::RawTransactionWireField::Value(*value),
|
|
},
|
|
)
|
|
.ok()?;
|
|
return ksp_raw_transaction_lib::canonicalize_raw_transaction(material).ok();
|
|
}
|
|
|
|
fn raw_from_helius_full(
|
|
network: ksp_store_lib::RawNetworkId,
|
|
notification: &ksp_onchain_transport_lib::HeliusFullTransactionNotification,
|
|
) -> std::option::Option<ksp_store_lib::RawTransaction> {
|
|
let object = notification.transaction().as_object()?;
|
|
let encoded = object.get("transaction")?.as_array()?;
|
|
if encoded.len() != 2 || encoded.get(1)?.as_str()? != "base64" {
|
|
return std::option::Option::None;
|
|
}
|
|
let data = encoded.first()?.as_str()?;
|
|
let meta = match object.get("meta") {
|
|
std::option::Option::None => ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
|
std::option::Option::Some(serde_json::Value::Null) => ksp_raw_transaction_lib::RawTransactionWireField::Null,
|
|
std::option::Option::Some(value) => ksp_raw_transaction_lib::RawTransactionWireField::Value(value.clone()),
|
|
};
|
|
let signature = ksp_raw_transaction_lib::parse_raw_transaction_signature(notification.signature()).ok()?;
|
|
let transaction_index = u32::try_from(notification.transaction_index()).ok()?;
|
|
let material = ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64(
|
|
network,
|
|
signature,
|
|
notification.slot(),
|
|
std::option::Option::None,
|
|
data,
|
|
meta,
|
|
ksp_raw_transaction_lib::RawTransactionWireField::Omitted,
|
|
ksp_raw_transaction_lib::RawTransactionWireField::Value(transaction_index),
|
|
);
|
|
return ksp_raw_transaction_lib::canonicalize_raw_transaction(material).ok();
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_005_standard_block_subscribe_full_base64_matches_http_get_block_raw_v1_for_legacy_transaction() {
|
|
const HTTP_BODY: &str = concat!(
|
|
"{\"jsonrpc\":\"2.0\",\"result\":{\"previousBlockhash\":\"previous\",\"blockhash\":\"block\",\"parentSlot\":430000122,",
|
|
"\"rewards\":[],\"numRewardPartitions\":0,\"blockTime\":1787072400,\"blockHeight\":410000000,",
|
|
"\"transactions\":[{\"transaction\":[\"AQAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA",
|
|
"AAAAAA\",\"base64\"],\"meta\":{\"err\":null,\"fee\":5000},\"version\":\"legacy\"}]},\"id\":1}",
|
|
);
|
|
let (http_url, http_server) = serve_http_once(HTTP_BODY).expect("HTTP parity fixture must start");
|
|
let http_pool = http_pool_for_url(http_url.as_str()).expect("HTTP parity pool must build");
|
|
let http_config = ksp_onchain_transport_lib::SolanaGetBlockConfig::new(
|
|
std::option::Option::Some(ksp_onchain_transport_lib::SolanaCommitment::Confirmed),
|
|
std::option::Option::Some(ksp_onchain_transport_lib::SolanaTransactionEncoding::Base64),
|
|
std::option::Option::Some(ksp_onchain_transport_lib::SolanaTransactionDetails::Full),
|
|
std::option::Option::Some(0),
|
|
std::option::Option::Some(false),
|
|
);
|
|
let http = http_pool
|
|
.get_block_observed(&ksp_onchain_transport_lib::HttpRoleName::new("default"), FIXTURE_SLOT, std::option::Option::Some(&http_config))
|
|
.await
|
|
.expect("HTTP block parity fixture must decode");
|
|
let http_block = http.value().as_ref().expect("HTTP parity block must be present").clone();
|
|
http_server.join().expect("HTTP parity fixture must join").expect("HTTP parity fixture must complete");
|
|
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.expect("WS parity listener must bind");
|
|
let address = listener.local_addr().expect("WS parity listener address must resolve");
|
|
let ws_block = serde_json::from_str::<serde_json::Value>(HTTP_BODY).expect("HTTP parity JSON must parse")["result"].clone();
|
|
let server = tokio::spawn(async move {
|
|
let (stream, _) = listener.accept().await.map_err(|_| return "WS parity server accept failed".to_owned())?;
|
|
let mut websocket = tokio_tungstenite::accept_async(stream).await.map_err(|_| return "WS parity handshake failed".to_owned())?;
|
|
let subscribe = read_ws_request(&mut websocket).await?;
|
|
if subscribe.get("method") != std::option::Option::Some(&serde_json::json!("blockSubscribe")) {
|
|
return Err("WS parity server received unexpected subscribe method".to_owned());
|
|
}
|
|
send_ws_result(&mut websocket, &subscribe, serde_json::json!(501)).await?;
|
|
send_ws_json(
|
|
&mut websocket,
|
|
serde_json::json!({
|
|
"jsonrpc":"2.0",
|
|
"method":"blockNotification",
|
|
"params":{"subscription":501,"result":{"context":{"slot":FIXTURE_SLOT},"value":{"slot":FIXTURE_SLOT,"block":ws_block,"err":null}}}
|
|
}),
|
|
)
|
|
.await?;
|
|
let unsubscribe = read_ws_request(&mut websocket).await?;
|
|
if unsubscribe.get("method") != std::option::Option::Some(&serde_json::json!("blockUnsubscribe")) {
|
|
return Err("WS parity server received unexpected unsubscribe method".to_owned());
|
|
}
|
|
send_ws_result(&mut websocket, &unsubscribe, serde_json::json!(true)).await?;
|
|
wait_for_ws_close(&mut websocket).await;
|
|
return Ok::<(), std::string::String>(());
|
|
});
|
|
let endpoint = ws_endpoint(format!("ws://{address}").as_str(), ksp_onchain_transport_lib::WsProtocolKind::SolanaStandard)
|
|
.expect("standard WS parity endpoint must build");
|
|
let session = ksp_onchain_transport_lib::SolanaStandardWsSession::connect(endpoint).await.expect("standard WS parity session must connect");
|
|
let ws_config = ksp_onchain_transport_lib::SolanaBlockSubscribeConfig::new(
|
|
std::option::Option::Some(ksp_onchain_transport_lib::SolanaCommitment::Confirmed),
|
|
std::option::Option::Some(ksp_onchain_transport_lib::SolanaTransactionEncoding::Base64),
|
|
std::option::Option::Some(ksp_onchain_transport_lib::SolanaTransactionDetails::Full),
|
|
std::option::Option::Some(0),
|
|
std::option::Option::Some(false),
|
|
);
|
|
let mut subscription = session
|
|
.block_subscribe(&ksp_onchain_transport_lib::SolanaBlockSubscribeFilter::All, std::option::Option::Some(&ws_config))
|
|
.await
|
|
.expect("standard blockSubscribe parity subscription must register");
|
|
let notification = subscription
|
|
.recv()
|
|
.await
|
|
.expect("standard blockSubscribe parity notification must arrive")
|
|
.expect("standard blockSubscribe parity notification must decode");
|
|
let ws_block = notification.value().block().expect("standard blockSubscribe parity block must be present");
|
|
assert_eq!(ws_block, &http_block);
|
|
let network = ksp_store_lib::RawNetworkId::new("devnet").expect("parity fixture network must be valid");
|
|
let http_transaction =
|
|
http_block.transactions().value().expect("HTTP parity transactions must be present").first().expect("HTTP parity transaction must exist");
|
|
let ws_transaction = ws_block.transactions().value().expect("WS parity transactions must be present").first().expect("WS parity transaction must exist");
|
|
let http_raw = raw_from_block_transaction(network.clone(), FIXTURE_SLOT, http_block.block_time(), http_transaction, 0)
|
|
.expect("HTTP block transaction must canonicalize");
|
|
let ws_raw = raw_from_block_transaction(network, notification.value().slot(), ws_block.block_time(), ws_transaction, 0)
|
|
.expect("WS block transaction must canonicalize");
|
|
assert_eq!(ws_raw.reference(), http_raw.reference());
|
|
assert_eq!(ws_raw.slot(), http_raw.slot());
|
|
assert_eq!(ws_raw.block_time(), http_raw.block_time());
|
|
assert_eq!(ws_raw.payload().bytes(), http_raw.payload().bytes());
|
|
assert_eq!(ws_raw.payload().content_hash(), http_raw.payload().content_hash());
|
|
assert!(subscription.unsubscribe().await.expect("standard blockSubscribe parity unsubscribe must complete"));
|
|
session.close().await.expect("standard WS parity session must close");
|
|
server.await.expect("standard WS parity server task must join").expect("standard WS parity server must complete");
|
|
}
|
|
|
|
#[tokio::test(flavor = "current_thread")]
|
|
async fn pre_005_helius_full_base64_without_block_time_and_version_requires_http_hydration_for_raw_v1() {
|
|
const HTTP_BODY: &str = concat!(
|
|
"{\"jsonrpc\":\"2.0\",\"result\":{\"slot\":430000123,\"blockTime\":1787072400,",
|
|
"\"transaction\":[\"AQAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA\",\"base64\"],",
|
|
"\"meta\":{\"err\":null,\"fee\":5000},\"version\":\"legacy\",\"transactionIndex\":0},\"id\":1}",
|
|
);
|
|
let (http_url, http_server) = serve_http_once(HTTP_BODY).expect("Helius hydration HTTP fixture must start");
|
|
let http_pool = http_pool_for_url(http_url.as_str()).expect("Helius hydration HTTP pool must build");
|
|
let http_config = ksp_onchain_transport_lib::SolanaGetTransactionConfig::new(
|
|
std::option::Option::Some(ksp_onchain_transport_lib::SolanaCommitment::Confirmed),
|
|
std::option::Option::Some(ksp_onchain_transport_lib::SolanaTransactionEncoding::Base64),
|
|
std::option::Option::Some(0),
|
|
);
|
|
let http = http_pool
|
|
.get_transaction_observed(&ksp_onchain_transport_lib::HttpRoleName::new("default"), ZERO_SIGNATURE_TEXT, std::option::Option::Some(&http_config))
|
|
.await
|
|
.expect("Helius hydration HTTP fixture must decode");
|
|
let http_transaction = http.value().as_ref().expect("Helius hydration HTTP transaction must be present");
|
|
let network = ksp_store_lib::RawNetworkId::new("devnet").expect("Helius hydration fixture network must be valid");
|
|
let http_raw = raw_from_confirmed_transaction(network.clone(), http_transaction).expect("hydrated HTTP transaction must canonicalize");
|
|
http_server.join().expect("Helius hydration HTTP fixture must join").expect("Helius hydration HTTP fixture must complete");
|
|
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.expect("Helius parity listener must bind");
|
|
let address = listener.local_addr().expect("Helius parity listener address must resolve");
|
|
let server = tokio::spawn(async move {
|
|
let (stream, _) = listener.accept().await.map_err(|_| return "Helius parity server accept failed".to_owned())?;
|
|
let mut websocket = tokio_tungstenite::accept_async(stream).await.map_err(|_| return "Helius parity handshake failed".to_owned())?;
|
|
let subscribe = read_ws_request(&mut websocket).await?;
|
|
if subscribe.get("method") != std::option::Option::Some(&serde_json::json!("transactionSubscribe")) {
|
|
return Err("Helius parity server received unexpected subscribe method".to_owned());
|
|
}
|
|
send_ws_result(&mut websocket, &subscribe, serde_json::json!(601)).await?;
|
|
send_ws_json(
|
|
&mut websocket,
|
|
serde_json::json!({
|
|
"jsonrpc":"2.0",
|
|
"method":"transactionNotification",
|
|
"params":{
|
|
"subscription":601,
|
|
"result":{
|
|
"transaction":{"transaction":[ZERO_TRANSACTION_BASE64,"base64"],"meta":{"err":null,"fee":5000}},
|
|
"signature":ZERO_SIGNATURE_TEXT,
|
|
"slot":FIXTURE_SLOT,
|
|
"transactionIndex":0
|
|
}
|
|
}
|
|
}),
|
|
)
|
|
.await?;
|
|
let unsubscribe = read_ws_request(&mut websocket).await?;
|
|
if unsubscribe.get("method") != std::option::Option::Some(&serde_json::json!("transactionUnsubscribe")) {
|
|
return Err("Helius parity server received unexpected unsubscribe method".to_owned());
|
|
}
|
|
send_ws_result(&mut websocket, &unsubscribe, serde_json::json!(true)).await?;
|
|
wait_for_ws_close(&mut websocket).await;
|
|
return Ok::<(), std::string::String>(());
|
|
});
|
|
let endpoint = ws_endpoint(format!("ws://{address}").as_str(), ksp_onchain_transport_lib::WsProtocolKind::HeliusLaserStream)
|
|
.expect("Helius parity endpoint must build");
|
|
let session = ksp_onchain_transport_lib::HeliusLaserStreamWsSession::connect(endpoint).await.expect("Helius parity session must connect");
|
|
let request = ksp_onchain_transport_lib::HeliusTransactionSubscribeRequest::new(
|
|
ksp_onchain_transport_lib::HeliusTransactionSubscribeFilter::default(),
|
|
std::option::Option::Some(ksp_onchain_transport_lib::HeliusTransactionSubscribeOptions::new(
|
|
std::option::Option::Some(ksp_onchain_transport_lib::SolanaCommitment::Confirmed),
|
|
std::option::Option::Some(ksp_onchain_transport_lib::HeliusTransactionSubscribeEncoding::Base64),
|
|
std::option::Option::Some(ksp_onchain_transport_lib::SolanaTransactionDetails::Full),
|
|
std::option::Option::Some(false),
|
|
std::option::Option::Some(0),
|
|
)),
|
|
);
|
|
let mut subscription = session.transaction_subscribe(&request).await.expect("Helius parity transaction subscription must register");
|
|
let notification = subscription.recv().await.expect("Helius parity notification must arrive").expect("Helius parity notification must decode");
|
|
assert!(matches!(¬ification, ksp_onchain_transport_lib::HeliusTransactionNotification::Full(_)));
|
|
let full = match notification {
|
|
ksp_onchain_transport_lib::HeliusTransactionNotification::Full(full) => full,
|
|
_ => return,
|
|
};
|
|
let helius_raw = raw_from_helius_full(network, &full).expect("best-effort Helius full projection must canonicalize for negative parity proof");
|
|
assert_eq!(helius_raw.reference(), http_raw.reference());
|
|
assert_eq!(helius_raw.slot(), http_raw.slot());
|
|
assert!(helius_raw.block_time().is_none());
|
|
assert_eq!(http_raw.block_time().map(|value| return value.unix_millis()), std::option::Option::Some(FIXTURE_BLOCK_TIME * 1_000));
|
|
assert_ne!(helius_raw.block_time(), http_raw.block_time());
|
|
assert_ne!(helius_raw.payload().bytes(), http_raw.payload().bytes());
|
|
assert_ne!(helius_raw.payload().content_hash(), http_raw.payload().content_hash());
|
|
assert!(subscription.unsubscribe().await.expect("Helius parity transaction unsubscribe must complete"));
|
|
session.close().await.expect("Helius parity session must close");
|
|
server.await.expect("Helius parity server task must join").expect("Helius parity server must complete");
|
|
}
|