v0.3.10-pre.005-fix.002
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-job-backfill-lib/tests/hardening.rs
|
||||
// version: 5
|
||||
// version: 6
|
||||
|
||||
//! Adversarial, security, visibility and external-boundary hardening canaries for `pre.010`.
|
||||
|
||||
@@ -275,7 +275,7 @@ fn pre_010_manifest_dependency_surface_is_exact_and_backend_neutral() {
|
||||
"tokio",
|
||||
])
|
||||
);
|
||||
assert_eq!(dev, std::collections::BTreeSet::from(["tokio"]));
|
||||
assert_eq!(dev, std::collections::BTreeSet::from(["tokio", "tokio-tungstenite"]));
|
||||
assert!(manifest.contains("ksp-store-lib = { path = \"../ksp-store-lib\", default-features = false }"));
|
||||
assert!(manifest.contains("tokio = { workspace = true, features = [\"macros\", \"sync\"] }"));
|
||||
assert!(!manifest.contains("ksp-store-postgres-lib"));
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-job-backfill-lib/tests/ws_raw_parity.rs
|
||||
// version: 2
|
||||
// version: 3
|
||||
|
||||
//! Cross-layer parity canaries qualifying standard `blockSubscribe` and Helius `transactionSubscribe` for RAW v1 ingestion.
|
||||
|
||||
@@ -179,7 +179,7 @@ fn raw_from_block_transaction(
|
||||
},
|
||||
_ => return std::option::Option::None,
|
||||
};
|
||||
let material = ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64_with_embedded_signature(
|
||||
let material = match ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64_with_embedded_signature(
|
||||
network,
|
||||
slot,
|
||||
block_time,
|
||||
@@ -187,8 +187,10 @@ fn raw_from_block_transaction(
|
||||
map_meta(transaction.meta()),
|
||||
map_version(transaction.version()),
|
||||
ksp_raw_transaction_lib::RawTransactionWireField::Value(transaction_index),
|
||||
)
|
||||
.ok()?;
|
||||
) {
|
||||
Ok(material) => material,
|
||||
Err(_) => return std::option::Option::None,
|
||||
};
|
||||
return ksp_raw_transaction_lib::canonicalize_raw_transaction(material).ok();
|
||||
}
|
||||
|
||||
@@ -202,7 +204,7 @@ fn raw_from_confirmed_transaction(
|
||||
},
|
||||
_ => return std::option::Option::None,
|
||||
};
|
||||
let material = ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64_with_embedded_signature(
|
||||
let material = match ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64_with_embedded_signature(
|
||||
network,
|
||||
transaction.slot(),
|
||||
transaction.block_time(),
|
||||
@@ -214,8 +216,10 @@ fn raw_from_confirmed_transaction(
|
||||
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()?;
|
||||
) {
|
||||
Ok(material) => material,
|
||||
Err(_) => return std::option::Option::None,
|
||||
};
|
||||
return ksp_raw_transaction_lib::canonicalize_raw_transaction(material).ok();
|
||||
}
|
||||
|
||||
@@ -223,19 +227,50 @@ 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" {
|
||||
let object = match notification.transaction().as_object() {
|
||||
std::option::Option::Some(object) => object,
|
||||
std::option::Option::None => return std::option::Option::None,
|
||||
};
|
||||
let encoded = match object.get("transaction") {
|
||||
std::option::Option::Some(value) => match value.as_array() {
|
||||
std::option::Option::Some(encoded) => encoded,
|
||||
std::option::Option::None => return std::option::Option::None,
|
||||
},
|
||||
std::option::Option::None => return std::option::Option::None,
|
||||
};
|
||||
if encoded.len() != 2 {
|
||||
return std::option::Option::None;
|
||||
}
|
||||
let data = encoded.first()?.as_str()?;
|
||||
let encoding = match encoded.get(1) {
|
||||
std::option::Option::Some(value) => match value.as_str() {
|
||||
std::option::Option::Some(encoding) => encoding,
|
||||
std::option::Option::None => return std::option::Option::None,
|
||||
},
|
||||
std::option::Option::None => return std::option::Option::None,
|
||||
};
|
||||
if encoding != "base64" {
|
||||
return std::option::Option::None;
|
||||
}
|
||||
let data = match encoded.first() {
|
||||
std::option::Option::Some(value) => match value.as_str() {
|
||||
std::option::Option::Some(data) => data,
|
||||
std::option::Option::None => return std::option::Option::None,
|
||||
},
|
||||
std::option::Option::None => return std::option::Option::None,
|
||||
};
|
||||
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 signature = match ksp_raw_transaction_lib::parse_raw_transaction_signature(notification.signature()) {
|
||||
Ok(signature) => signature,
|
||||
Err(_) => return std::option::Option::None,
|
||||
};
|
||||
let transaction_index = match u32::try_from(notification.transaction_index()) {
|
||||
Ok(transaction_index) => transaction_index,
|
||||
Err(_) => return std::option::Option::None,
|
||||
};
|
||||
let material = ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64(
|
||||
network,
|
||||
signature,
|
||||
@@ -276,14 +311,25 @@ async fn pre_005_standard_block_subscribe_full_base64_matches_http_get_block_raw
|
||||
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?;
|
||||
let (stream, _) = match listener.accept().await {
|
||||
Ok(accepted) => accepted,
|
||||
Err(_) => return Err("WS parity server accept failed".to_owned()),
|
||||
};
|
||||
let mut websocket = match tokio_tungstenite::accept_async(stream).await {
|
||||
Ok(websocket) => websocket,
|
||||
Err(_) => return Err("WS parity handshake failed".to_owned()),
|
||||
};
|
||||
let subscribe = match read_ws_request(&mut websocket).await {
|
||||
Ok(subscribe) => subscribe,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
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(
|
||||
if let Err(error) = send_ws_result(&mut websocket, &subscribe, serde_json::json!(501)).await {
|
||||
return Err(error);
|
||||
}
|
||||
if let Err(error) = send_ws_json(
|
||||
&mut websocket,
|
||||
serde_json::json!({
|
||||
"jsonrpc":"2.0",
|
||||
@@ -291,12 +337,20 @@ async fn pre_005_standard_block_subscribe_full_base64_matches_http_get_block_raw
|
||||
"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?;
|
||||
.await
|
||||
{
|
||||
return Err(error);
|
||||
}
|
||||
let unsubscribe = match read_ws_request(&mut websocket).await {
|
||||
Ok(unsubscribe) => unsubscribe,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
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?;
|
||||
if let Err(error) = send_ws_result(&mut websocket, &unsubscribe, serde_json::json!(true)).await {
|
||||
return Err(error);
|
||||
}
|
||||
wait_for_ws_close(&mut websocket).await;
|
||||
return Ok::<(), std::string::String>(());
|
||||
});
|
||||
@@ -364,14 +418,25 @@ async fn pre_005_helius_full_base64_without_block_time_and_version_requires_http
|
||||
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?;
|
||||
let (stream, _) = match listener.accept().await {
|
||||
Ok(accepted) => accepted,
|
||||
Err(_) => return Err("Helius parity server accept failed".to_owned()),
|
||||
};
|
||||
let mut websocket = match tokio_tungstenite::accept_async(stream).await {
|
||||
Ok(websocket) => websocket,
|
||||
Err(_) => return Err("Helius parity handshake failed".to_owned()),
|
||||
};
|
||||
let subscribe = match read_ws_request(&mut websocket).await {
|
||||
Ok(subscribe) => subscribe,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
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(
|
||||
if let Err(error) = send_ws_result(&mut websocket, &subscribe, serde_json::json!(601)).await {
|
||||
return Err(error);
|
||||
}
|
||||
if let Err(error) = send_ws_json(
|
||||
&mut websocket,
|
||||
serde_json::json!({
|
||||
"jsonrpc":"2.0",
|
||||
@@ -387,12 +452,20 @@ async fn pre_005_helius_full_base64_without_block_time_and_version_requires_http
|
||||
}
|
||||
}),
|
||||
)
|
||||
.await?;
|
||||
let unsubscribe = read_ws_request(&mut websocket).await?;
|
||||
.await
|
||||
{
|
||||
return Err(error);
|
||||
}
|
||||
let unsubscribe = match read_ws_request(&mut websocket).await {
|
||||
Ok(unsubscribe) => unsubscribe,
|
||||
Err(error) => return Err(error),
|
||||
};
|
||||
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?;
|
||||
if let Err(error) = send_ws_result(&mut websocket, &unsubscribe, serde_json::json!(true)).await {
|
||||
return Err(error);
|
||||
}
|
||||
wait_for_ws_close(&mut websocket).await;
|
||||
return Ok::<(), std::string::String>(());
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user