From 0683b70abe1443c109682bee11fb23827be47741 Mon Sep 17 00:00:00 2001 From: SinuS Von SifriduS Date: Sun, 6 Sep 2026 21:59:15 +0200 Subject: [PATCH] v0.3.10-pre.005 --- Cargo.toml | 4 +- crates/ksp-job-backfill-lib/Cargo.toml | 3 +- .../tests/dependency_boundary.rs | 21 +- .../tests/ws_raw_parity.rs | 430 ++++++++++++++++++ deltas/0.3.10/pre.005.md | 232 ++++++++++ .../011-RAW_TRANSACTION_ACQUISITION.md | 73 +-- ...031-V0_3_10_RAW_TRANSACTION_INGEST_PLAN.md | 60 +-- .../027-V0_3_10_RAW_TRANSACTION_INGEST.md | 62 ++- 8 files changed, 815 insertions(+), 70 deletions(-) create mode 100644 crates/ksp-job-backfill-lib/tests/ws_raw_parity.rs create mode 100644 deltas/0.3.10/pre.005.md diff --git a/Cargo.toml b/Cargo.toml index 58a64a0..5b17edc 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,12 @@ # file: Cargo.toml -# version: 495 +# version: 496 [workspace] resolver = "3" members = ["crates/ksp-app-backfill-desk", "crates/ksp-app-config-desk", "crates/ksp-app-solprices-desk", "crates/ksp-app-store-desk", "crates/ksp-app-wallet-desk", "crates/ksp-config-lib", "crates/ksp-core-lib", "crates/ksp-interface-lib", "crates/ksp-job-api", "crates/ksp-job-backfill-lib", "crates/ksp-logging-lib", "crates/ksp-offchain-transport-lib", "crates/ksp-onchain-transport-lib", "crates/ksp-program-api", "crates/ksp-raw-transaction-lib", "crates/ksp-store-api", "crates/ksp-store-lib", "crates/ksp-store-postgres-lib", "crates/ksp-wallet-lib", "crates/ksp-worker-api"] [workspace.package] -version = "0.3.10-pre.4.fix.1" +version = "0.3.10-pre.5" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-job-backfill-lib/Cargo.toml b/crates/ksp-job-backfill-lib/Cargo.toml index 474dcfd..cbc02c8 100644 --- a/crates/ksp-job-backfill-lib/Cargo.toml +++ b/crates/ksp-job-backfill-lib/Cargo.toml @@ -1,5 +1,5 @@ # file: crates/ksp-job-backfill-lib/Cargo.toml -# version: 5 +# version: 6 [package] name = "ksp-job-backfill-lib" @@ -21,6 +21,7 @@ tokio = { workspace = true, features = ["macros", "sync"] } [dev-dependencies] tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } +tokio-tungstenite = { workspace = true, features = ["handshake"] } [lints] workspace = true diff --git a/crates/ksp-job-backfill-lib/tests/dependency_boundary.rs b/crates/ksp-job-backfill-lib/tests/dependency_boundary.rs index 1380eaa..e099133 100644 --- a/crates/ksp-job-backfill-lib/tests/dependency_boundary.rs +++ b/crates/ksp-job-backfill-lib/tests/dependency_boundary.rs @@ -1,5 +1,5 @@ // file: crates/ksp-job-backfill-lib/tests/dependency_boundary.rs -// version: 7 +// version: 8 //! Dependency firewall canaries through the concrete cancellation and latest-value runtime tranche. @@ -40,6 +40,25 @@ fn pre_003_manifest_adds_common_raw_edge_without_crossing_store_or_runtime_bound return; } +#[test] +fn pre_005_websocket_parity_dependency_stays_test_only() { + let manifest = include_str!("../Cargo.toml"); + let parts = manifest.splitn(2, "[dev-dependencies]").collect::>(); + let (dependencies, dev_dependencies) = match parts.as_slice() { + [dependencies, dev_dependencies] => (*dependencies, *dev_dependencies), + _ => { + assert_eq!(parts.len(), 2, "Backfill manifest must retain exactly one dev-dependencies section"); + return; + }, + }; + assert!(!dependencies.contains("tokio-tungstenite"), "WebSocket parity client must not become a production Backfill dependency"); + assert!( + dev_dependencies.contains("tokio-tungstenite = { workspace = true, features = [\"handshake\"] }"), + "WebSocket parity client must remain an explicit test-only dependency", + ); + return; +} + #[test] fn pre_009_production_sources_keep_transport_store_and_scheduler_ownership_separate() { let neutral_sources = [ diff --git a/crates/ksp-job-backfill-lib/tests/ws_raw_parity.rs b/crates/ksp-job-backfill-lib/tests/ws_raw_parity.rs new file mode 100644 index 0000000..9122c20 --- /dev/null +++ b/crates/ksp-job-backfill-lib/tests/ws_raw_parity.rs @@ -0,0 +1,430 @@ +// 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; + +fn http_pool_for_url(url: &str) -> ksp_core_lib::Result { + 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>)> { + 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 { + 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 { + 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) -> ksp_raw_transaction_lib::RawTransactionWireField { + 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_raw_transaction_lib::RawTransactionWireField { + 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, + transaction: &ksp_onchain_transport_lib::SolanaBlockTransaction, + transaction_index: u32, +) -> std::option::Option { + 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 { + 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 { + 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::(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"); +} diff --git a/deltas/0.3.10/pre.005.md b/deltas/0.3.10/pre.005.md new file mode 100644 index 0000000..0dfc5fe --- /dev/null +++ b/deltas/0.3.10/pre.005.md @@ -0,0 +1,232 @@ + + + +# Delta `0.3.10-pre.005` — qualification RAW WS standard et Helius + +## Base requise + +```text +livraison : 0.3.10-pre.004-fix.001 +workspace.package.version = 0.3.10-pre.4.fix.1 +``` + +Le rejeu opérateur communiqué le 2026-09-06 ferme intégralement le fix : audits statiques propres, `cargo check --workspace` PASS, Clippy strict PASS, tests common PASS, 387 tests Transport plus suites publiques/release/doctests PASS, et suite complète Backfill PASS. Les arbres Cargo n'ont pas été rejoués sur le fix car aucune dépendance ni feature n'y avait changé. + +## Objectif + +Qualifier les deux surfaces WebSocket déjà présentes sans ouvrir le Worker : + +```text +Solana standard blockSubscribe full/base64 + -> prouver ou refuser la production RAW directe + +Helius transactionSubscribe full/base64 + -> prouver ou refuser la production RAW directe +``` + +Aucun adapter productif n'est placé dans Transport ou Backfill. `TR-C2` reste Worker-owned. + +## Version + +```text +workspace.package.version = 0.3.10-pre.5 +``` + +Le header du `Cargo.toml` racine passe de `495` à `496`. Le workspace reste à `20` membres. + +## Qualification positive — `blockSubscribe` + +Le canari cross-layer `crates/ksp-job-backfill-lib/tests/ws_raw_parity.rs` monte un serveur HTTP local et un serveur WebSocket standard local sur une fixture commune. + +Configuration qualifiée : + +```text +commitment = confirmed +encoding = base64 +transactionDetails = full +maxSupportedTransactionVersion = 0 +showRewards = false +``` + +Le test compare : + +```text +HTTP getBlock observed + vs +WS blockSubscribe +``` + +puis exige : + +```text +même SolanaConfirmedBlock +même RawTransactionReference +même slot +même block_time +mêmes bytes RAW v1 +même content_hash +``` + +La qualification directe est volontairement bornée au sous-ensemble legacy/v0 prouvé. Transaction V1 reste au gate `pre.006`. + +## Qualification négative — Helius `transactionSubscribe` + +Le même canari ouvre une session Helius locale et une hydration HTTP de référence pour la même transaction. + +L'enveloppe full/base64 qualifiée conserve : + +```text +transaction wire +meta +signature +slot +transactionIndex +``` + +mais ne porte pas, dans le contrat typé actuel : + +```text +blockTime +version +``` + +Le test construit donc la meilleure projection RAW possible sans inventer ces champs et prouve : + +```text +même identité que HTTP +même slot que HTTP +block_time différent +bytes RAW v1 différents +content_hash différent +``` + +Décision normative : + +```text +Helius transactionSubscribe full/base64 + -> signal live riche + -> HTTP hydration obligatoire + -> persistence RAW depuis le matériau HTTP complet +``` + +Aucun RAW v2 automatique et aucune normalisation ad hoc ne sont introduits. + +## Dependency firewall + +Pour les serveurs WebSocket locaux du canari, le manifest Backfill ajoute uniquement : + +```text +[dev-dependencies] +tokio-tungstenite = { workspace = true, features = ["handshake"] } +``` + +`tests/dependency_boundary.rs` vérifie explicitement que `tokio-tungstenite` n'apparaît pas dans `[dependencies]` et reste dans `[dev-dependencies]`. + +Il n'existe toujours aucun edge productif : + +```text +ksp-onchain-transport-lib -> ksp-raw-transaction-lib +ksp-raw-transaction-lib -> ksp-onchain-transport-lib +``` + +## Documentation + +`docs/architecture/011-RAW_TRANSACTION_ACQUISITION.md` distingue désormais : + +```text +blockSubscribe full/base64 legacy-v0 -> direct après parité prouvée +Helius transactionSubscribe full/base64 -> signal + hydration HTTP +``` + +Le plan `031` et la validation `027` enregistrent la même qualification, la fermeture opérateur de `pre.004-fix.001` et le maintien de Transaction V1 dans `pre.006`. + +## Fraîcheur externe revalidée + +Les surfaces changeantes ont été revalidées le 2026-09-06 : + +```text +Solana blockSubscribe : méthode unstable, confirmed/finalized, base64, full, maxSupportedTransactionVersion disponible +Helius transactionSubscribe : full transaction streaming disponible ; enveloppe KSP qualifiée sans blockTime/version +``` + +Sources documentaires : + +```text +https://solana.com/docs/rpc/websocket/blocksubscribe +https://solana.com/docs/rpc/json-structures +https://www.helius.dev/docs/rpc/websocket/transaction-subscribe +https://www.helius.dev/blog/introducing-next-generation-enhanced-websockets +``` + +Ces faits d'accès/protocole restent documentaires et ne deviennent pas des constantes provider dans le Worker. + +## Fichiers ajoutés + +```text +crates/ksp-job-backfill-lib/tests/ws_raw_parity.rs +deltas/0.3.10/pre.005.md +``` + +## Fichiers modifiés + +```text +Cargo.toml +crates/ksp-job-backfill-lib/Cargo.toml +crates/ksp-job-backfill-lib/tests/dependency_boundary.rs +docs/architecture/011-RAW_TRANSACTION_ACQUISITION.md +docs/plans/031-V0_3_10_RAW_TRANSACTION_INGEST_PLAN.md +docs/validation/027-V0_3_10_RAW_TRANSACTION_INGEST.md +``` + +## Fichiers supprimés + +Aucun. + +## Validations exécutées dans l'environnement d'assemblage + +```text +python3 scripts/audit_rust_workspace_rules.py +General Rust rule audit: clean +Rust export completeness audit: 0 candidate(s) +KSP workspace Rust rule audit: clean + +python3 scripts/audit_markdown_tables.py README.md RULES.md ROADMAP.md CHANGELOG.md docs prompts crates deltas +Markdown table audit: clean (339 table(s), 759 file(s)) +``` + +L'environnement d'assemblage ne fournit ni `cargo`, ni `rustc`, ni `rustfmt` ; aucun gate Cargo de `pre.005` n'est déclaré PASS localement. + +## Gate opérateur demandé + +```bash +cargo fmt --all +python3 scripts/audit_rust_workspace_rules.py +python3 scripts/audit_markdown_tables.py README.md RULES.md ROADMAP.md CHANGELOG.md docs prompts crates deltas +cargo check --workspace +cargo clippy --workspace --all-targets --all-features -- -D warnings +cargo test -p ksp-raw-transaction-lib +cargo test -p ksp-onchain-transport-lib +cargo test -p ksp-job-backfill-lib +cargo tree -p ksp-job-backfill-lib --edges normal +cargo tree -p ksp-job-backfill-lib --edges dev +cargo tree -p ksp-job-backfill-lib -e features +``` + +Les arbres sont requis cette fois-ci parce que `pre.005` ajoute une dev-dependency et doit prouver qu'elle ne contamine pas le graphe productif. + +## Décisions prises + +```text +blockSubscribe full/base64 legacy-v0 est qualifié RAW-direct +Helius transactionSubscribe full/base64 n'est pas qualifié RAW-direct +Helius doit hydrater par HTTP avant persistence RAW +aucun champ manquant n'est inventé +aucun RAW v2 automatique +TR-C2 reste Worker-owned +aucun code de production ajouté dans cette tranche +Transaction V1 reste à pre.006 +``` + +## Questions ouvertes + +Aucune question ne bloque le gate déterministe de `pre.005`. `pre.006` reste dédiée aux adapters Yellowstone transactions/blocks et au gate Transaction V1 avant toute persistance directe de ces matériaux. diff --git a/docs/architecture/011-RAW_TRANSACTION_ACQUISITION.md b/docs/architecture/011-RAW_TRANSACTION_ACQUISITION.md index fcba4af..2d691d6 100644 --- a/docs/architecture/011-RAW_TRANSACTION_ACQUISITION.md +++ b/docs/architecture/011-RAW_TRANSACTION_ACQUISITION.md @@ -1,5 +1,5 @@ - + # Acquisition et alimentation `RawTransaction` @@ -238,29 +238,29 @@ Une source concrète peut fournir plusieurs capabilities. Une capability peut ê ## 5. Matrice des méthodes et de leur applicabilité -| Famille / méthode | Contenu obtenu | RAW complet | Usage Worker live | Usage Job Backfill | Remarque | -|-----------------------------------------------------------------|---------------------------------------|------------------------------|------------------------------------------------------------|---------------------------------------------|------------------------------------------------------| -| HTTP `getTransaction(signature)` | transaction depuis signature connue | oui si disponible | oui, hydration d'un signal live | oui, hydration historique | dépend de la rétention du RPC | -| HTTP `getSignaturesForAddress` + `getTransaction` | discovery adressée puis transaction | oui après hydration | non comme campagne de scan | oui, stratégie historique principale | naturellement paramétré par adresse/programme | -| HTTP `getBlocks` / `getBlocksWithLimit` | slots confirmés | non | oui pour suivre/recoller le frontier courant si nécessaire | oui pour énumérer une plage historique | discovery par slots | -| HTTP `getBlock(slot)` full | bloc et transactions | oui | oui pour live polling ou repair de continuité | oui pour historique par bloc | nécessite provenance observed en pool multi-endpoint | -| HTTP `getSlot` / `getFirstAvailableBlock` / `minimumLedgerSlot` | bornes de ledger | non | oui, continuité | oui, admission d'une campagne | aucune transaction directe | -| WS `logsSubscribe` | signature + logs | non | oui, discovery live + hydration | non pour historique pur | `mentions` standard limité à un pubkey | -| WS `signatureSubscribe` | statut d'une signature connue | non | oui, confirmation ciblée interne | possible pour une requête ciblée en attente | one-shot | -| WS `blockSubscribe` full | bloc + transactions | oui | oui | non sans mécanisme de replay historique | méthode standard instable | -| Helius `transactionSubscribe` full | transaction + meta | oui | oui | non comme source historique | extension provider | -| Yellowstone `transactions` | transaction exécutée + meta | oui | oui | oui si replay borné demandé/disponible | filters server-side | -| Yellowstone `blocks` avec transactions | bloc + transactions | oui | oui | oui si replay borné demandé/disponible | utile au live et au backfill | -| Yellowstone `transactions_status` | signature/status/error | non | oui + hydration | oui si replay/filtre borné | discovery/statut uniquement | -| Yellowstone slots / `blocks_meta` | continuité de slots | non | oui | auxiliaire | aucune transaction full | -| Yellowstone `from_slot` / replay | reprise d'un stream depuis slot | dépend du stream | oui uniquement pour continuité du Worker | oui pour campagne historique bornée | profondeur provider-specific | -| API provider historique adressée | historique signatures ou transactions | provider-dependent | non comme campagne historique | oui | exemple Helius `getTransactionsForAddress` | -| RPC archive standard | méthodes HTTP sur ledger ancien | oui selon méthode | non pour recherche historique | oui | provider ou self-host | -| Old Faithful / `yellowstone-faithful` | historique par RPC/index | oui pour données disponibles | non | oui | archive spécialisée | -| transaction body pré-exécution | signature/corps sans meta finale | non | oui, EARLY + hydration | non principal | faible latence | -| shreds / deshred | fragments ou transaction reconstruite | non sans meta d'exécution | oui, EARLY + hydration | non principal | adapter spécialisé | -| Agave RPC auto-hébergé | méthodes standard | selon méthode | oui | oui | mêmes rôles que RPC standard | -| Agave + Yellowstone auto-hébergé | streams Yellowstone | selon stream | oui | oui avec rétention/replay opérateur | capability opérateur | +| Famille / méthode | Contenu obtenu | RAW complet | Usage Worker live | Usage Job Backfill | Remarque | +|-----------------------------------------------------------------|---------------------------------------|------------------------------|------------------------------------------------------------|---------------------------------------------|------------------------------------------------------------| +| HTTP `getTransaction(signature)` | transaction depuis signature connue | oui si disponible | oui, hydration d'un signal live | oui, hydration historique | dépend de la rétention du RPC | +| HTTP `getSignaturesForAddress` + `getTransaction` | discovery adressée puis transaction | oui après hydration | non comme campagne de scan | oui, stratégie historique principale | naturellement paramétré par adresse/programme | +| HTTP `getBlocks` / `getBlocksWithLimit` | slots confirmés | non | oui pour suivre/recoller le frontier courant si nécessaire | oui pour énumérer une plage historique | discovery par slots | +| HTTP `getBlock(slot)` full | bloc et transactions | oui | oui pour live polling ou repair de continuité | oui pour historique par bloc | nécessite provenance observed en pool multi-endpoint | +| HTTP `getSlot` / `getFirstAvailableBlock` / `minimumLedgerSlot` | bornes de ledger | non | oui, continuité | oui, admission d'une campagne | aucune transaction directe | +| WS `logsSubscribe` | signature + logs | non | oui, discovery live + hydration | non pour historique pur | `mentions` standard limité à un pubkey | +| WS `signatureSubscribe` | statut d'une signature connue | non | oui, confirmation ciblée interne | possible pour une requête ciblée en attente | one-shot | +| WS `blockSubscribe` full/base64, legacy-v0 | bloc + transactions | oui, parité RAW v1 prouvée | oui, direct pour le sous-ensemble qualifié | non sans mécanisme de replay historique | méthode standard instable ; capability validator-dependent | +| Helius `transactionSubscribe` full/base64 | transaction + meta + signature/index | non sans hydration | oui, signal riche puis hydration HTTP | non comme source historique | absence de blockTime/version dans l’enveloppe qualifiée | +| Yellowstone `transactions` | transaction exécutée + meta | oui | oui | oui si replay borné demandé/disponible | filters server-side | +| Yellowstone `blocks` avec transactions | bloc + transactions | oui | oui | oui si replay borné demandé/disponible | utile au live et au backfill | +| Yellowstone `transactions_status` | signature/status/error | non | oui + hydration | oui si replay/filtre borné | discovery/statut uniquement | +| Yellowstone slots / `blocks_meta` | continuité de slots | non | oui | auxiliaire | aucune transaction full | +| Yellowstone `from_slot` / replay | reprise d'un stream depuis slot | dépend du stream | oui uniquement pour continuité du Worker | oui pour campagne historique bornée | profondeur provider-specific | +| API provider historique adressée | historique signatures ou transactions | provider-dependent | non comme campagne historique | oui | exemple Helius `getTransactionsForAddress` | +| RPC archive standard | méthodes HTTP sur ledger ancien | oui selon méthode | non pour recherche historique | oui | provider ou self-host | +| Old Faithful / `yellowstone-faithful` | historique par RPC/index | oui pour données disponibles | non | oui | archive spécialisée | +| transaction body pré-exécution | signature/corps sans meta finale | non | oui, EARLY + hydration | non principal | faible latence | +| shreds / deshred | fragments ou transaction reconstruite | non sans meta d'exécution | oui, EARLY + hydration | non principal | adapter spécialisé | +| Agave RPC auto-hébergé | méthodes standard | selon méthode | oui | oui | mêmes rôles que RPC standard | +| Agave + Yellowstone auto-hébergé | streams Yellowstone | selon stream | oui | oui avec rétention/replay opérateur | capability opérateur | Cette matrice est intentionnellement **usage-first**. HTTP n'est pas « Backfill » et gRPC n'est pas « Worker » : leur rôle dépend de l'opération effectuée. @@ -273,9 +273,8 @@ Sources capables de produire directement un matériau transactionnel complet : ```text Yellowstone transactions Yellowstone blocks avec transactions -WS blockSubscribe full -Helius transactionSubscribe full -provider stream compatible full transaction +WS blockSubscribe full/base64 legacy-v0 qualifié +provider stream compatible full transaction après parité prouvée ``` Chemin logique : @@ -562,6 +561,8 @@ Helius transactionSubscribe Le runtime WS possède des queues bornées et une logique de reconnect/resubscribe. Un reconnect ne constitue toutefois pas un replay adressable. +La qualification `pre.005` ferme deux cas distincts : `blockSubscribe` standard en `full/base64` avec `maxSupportedTransactionVersion = 0` produit, sur fixture identique, les mêmes bytes/hash RAW v1 que `getBlock`; Helius `transactionSubscribe` full/base64 conserve l’identité, le slot, le transaction wire, la meta et l’index mais ne transporte pas `blockTime` ni `version` dans l’enveloppe qualifiée. Helius reste donc un signal live riche à hydrater par HTTP avant persistence RAW directe. + ### 12.3 Yellowstone gRPC Le moteur KSP expose déjà les familles nécessaires : @@ -700,8 +701,8 @@ Yellowstone transactions Yellowstone blocks Yellowstone status + hydration WS logsSubscribe + HTTP getTransaction -WS blockSubscribe full -Helius transactionSubscribe full +WS blockSubscribe full/base64 legacy-v0 direct qualifié +Helius transactionSubscribe full/base64 + HTTP hydration HTTP live block polling HTTP hydration replay Yellowstone pour continuité du run @@ -710,13 +711,13 @@ sources EARLY via adapter extensible ### 14.3 Gaps Transport Worker -| ID | Adaptation | Motif | -|--------|-----------------------------------------------------------------------------------|------------------------------------------------------| -| `TR-B` | fermé : `get_block_observed` symétrique de `get_transaction_observed` | provenance exacte en pool HTTP | -| `TR-C` | projection source-neutral des transactions full WS/Yellowstone | éviter plusieurs canonicalizers filaires | -| `TR-D` | métadonnées sûres d'acquisition live au moment de la conversion | observation uniforme | -| `TR-E` | conserver/exploiter `from_slot`, replay info et snapshots de continuité existants | ne pas créer un second moteur Yellowstone | -| `TR-F` | adapters EARLY uniquement quand leur protocole est réellement implémenté | réserver la capability sans tout coder immédiatement | +| ID | Adaptation | Motif | +|--------|---------------------------------------------------------------------------------------|------------------------------------------------------| +| `TR-B` | fermé : `get_block_observed` symétrique de `get_transaction_observed` | provenance exacte en pool HTTP | +| `TR-C` | projection source-neutral des transactions full WS/Yellowstone ; Helius via hydration | éviter plusieurs canonicalizers filaires | +| `TR-D` | métadonnées sûres d'acquisition live au moment de la conversion | observation uniforme | +| `TR-E` | conserver/exploiter `from_slot`, replay info et snapshots de continuité existants | ne pas créer un second moteur Yellowstone | +| `TR-F` | adapters EARLY uniquement quand leur protocole est réellement implémenté | réserver la capability sans tout coder immédiatement | ### 14.4 Config Worker diff --git a/docs/plans/031-V0_3_10_RAW_TRANSACTION_INGEST_PLAN.md b/docs/plans/031-V0_3_10_RAW_TRANSACTION_INGEST_PLAN.md index 11a71a0..50a2783 100644 --- a/docs/plans/031-V0_3_10_RAW_TRANSACTION_INGEST_PLAN.md +++ b/docs/plans/031-V0_3_10_RAW_TRANSACTION_INGEST_PLAN.md @@ -1,5 +1,5 @@ - + # Plan v0.3.10 — RAW Transaction commune + Worker d’ingestion multi-source @@ -356,19 +356,19 @@ La priorité `pre.001` n’est pas d’affirmer que la parité existe déjà ; e ### 6.1 Matrice capabilities Worker V1 -| Capability | Méthode / protocole | Surface KSP `0.3.9` | Gap exact | Full / hydration | Provenance disponible | Continuité / replay | Preuve accessible | Adaptation Transport | Adaptation Config | Fixture | Live smoke | -|:-------------------------------|:-----------------------------|:-------------------------------------------------|:--------------------------------------------------------------------|:----------------------------------------|:---------------------------------------------------------------------|:--------------------------------------------------------------|:----------------------------------------------------------|:-----------------------------------------|:--------------------------------------------------|:------------------------|:-----------------------------| -| HTTP transaction hydration | `getTransaction` JSON-RPC | `get_transaction_observed` | aucun gap de provenance ; projection vers common à extraire | full après réponse non-null | provider + endpoint gagnant | aucune continuité native | Mainnet/Devnet/Testnet publics | non, hors adapter common | non | oui | oui opt-in | -| HTTP live blocks | `getSlot` + `getBlock` | wrappers typés présents | `TR-B`: pas de `get_block_observed` | full par bloc | actuellement insuffisante en pool | polling depuis run frontier ; repair borné | publics | oui `get_block_observed` | rôle HTTP worker si profil dédié nécessaire | oui | oui opt-in | -| WS logs + hydration | `logsSubscribe` + HTTP | wrapper logs + sessions | composer discovery + hydration sans perdre le statut de source | hydration requise | WS session safe + HTTP observed ; Store garde endpoint full-material | resubscribe seulement, aucun replay WS | Solana public/Helius standard | adapter Worker + TR-D, pas nouveau actor | profils existants suffisants pour première preuve | oui | oui opt-in | -| WS signature + hydration | `signatureSubscribe` + HTTP | wrapper one-shot présent | utile ciblé mais pas source globale principale | hydration requise | idem | one-shot ; aucun replay | publics | aucun P0 spécifique | non | oui | secondaire | -| WS blockSubscribe | `blockSubscribe` standard | wrapper + `SolanaConfirmedBlock` | projection transactions de block vers common + parité | full si parité prouvée | session/provider logique ; pas de winner HTTP | reconnect/resubscribe, pas replay | validator public selon support ; Helius explicitement non | TR-C adapter/parité | aucune exigée si endpoint supporte | oui | opt-in, capability-dependent | -| Helius full transaction WS | `transactionSubscribe` | typed provider envelope, nested transaction JSON | transformer le nested payload en matériau common sans provider leak | full si parité prouvée, sinon hydration | provider/session/filter sûrs | reconnexion, pas historique garanti WSS | clé/tier opérateur requis selon plan | TR-C + TR-D | profil Helius éventuellement nécessaire plus tard | oui | opt-in/tier-dependent | -| Yellowstone transactions | `Subscribe.transactions` | DTO full KSP, stream/reconnect présents | projection body+meta structurés vers RAW v1 byte-identique | full si parité prouvée | provider + endpoint settings + filter + created_at disponibles | `from_slot`, replay info, reconnect snapshot | OrbitFlare Devnet / PublicNode selon accès | TR-C + TR-D + TR-E | profils V3 déjà présents | oui | oui opt-in | -| Yellowstone blocks | `Subscribe.blocks` | DTO block full + transactions | même gate de projection ; préserver slot/index/block_time | full si parité prouvée | idem | `from_slot`; upstream a corrigé replay blocks en juillet 2026 | mêmes providers | TR-C + TR-D + TR-E | non P0 | oui | oui opt-in | -| Yellowstone status + hydration | `transactions_status` + HTTP | DTO status + HTTP observed | composer signal + hydration | hydration requise | status filter/created_at + HTTP winner | replay selon provider + HTTP repair | mêmes providers | TR-D + TR-E | non P0 | oui | oui opt-in | -| Yellowstone block meta / slots | `blocks_meta`, `slots` | DTO + stream présents | aucun RawTransaction direct ; alimenter continuity tracker | non, auxiliaire | filter/slot/status/created_at | continuity evidence ; replay provider-specific | mêmes providers | TR-E côté Worker | non | oui | oui opt-in | -| EARLY source | provider/shred adapter | aucune surface générique unique | `TR-F`, uniquement protocole réellement implémenté | hydration sauf preuve meta complète | provider-specific | provider-specific | généralement payant | différé | différé | seulement si implémenté | seulement si accessible | +| Capability | Méthode / protocole | Surface KSP `0.3.9` | Gap exact | Full / hydration | Provenance disponible | Continuité / replay | Preuve accessible | Adaptation Transport | Adaptation Config | Fixture | Live smoke | +|:-------------------------------|:-----------------------------|:-------------------------------------------------|:----------------------------------------------------------------------|:--------------------------------------|:---------------------------------------------------------------------|:--------------------------------------------------------------|:----------------------------------------------------------|:-----------------------------------------|:--------------------------------------------------|:------------------------|:-----------------------------| +| HTTP transaction hydration | `getTransaction` JSON-RPC | `get_transaction_observed` | aucun gap de provenance ; projection vers common à extraire | full après réponse non-null | provider + endpoint gagnant | aucune continuité native | Mainnet/Devnet/Testnet publics | non, hors adapter common | non | oui | oui opt-in | +| HTTP live blocks | `getSlot` + `getBlock` | wrappers typés présents | `TR-B`: pas de `get_block_observed` | full par bloc | actuellement insuffisante en pool | polling depuis run frontier ; repair borné | publics | oui `get_block_observed` | rôle HTTP worker si profil dédié nécessaire | oui | oui opt-in | +| WS logs + hydration | `logsSubscribe` + HTTP | wrapper logs + sessions | composer discovery + hydration sans perdre le statut de source | hydration requise | WS session safe + HTTP observed ; Store garde endpoint full-material | resubscribe seulement, aucun replay WS | Solana public/Helius standard | adapter Worker + TR-D, pas nouveau actor | profils existants suffisants pour première preuve | oui | oui opt-in | +| WS signature + hydration | `signatureSubscribe` + HTTP | wrapper one-shot présent | utile ciblé mais pas source globale principale | hydration requise | idem | one-shot ; aucun replay | publics | aucun P0 spécifique | non | oui | secondaire | +| WS blockSubscribe | `blockSubscribe` standard | wrapper + `SolanaConfirmedBlock` | parité fermée pour full/base64 legacy-v0 ; adapter productif Worker | full direct pour sous-ensemble prouvé | session/provider logique ; pas de winner HTTP | reconnect/resubscribe, pas replay | validator public selon support ; Helius explicitement non | TR-C adapter Worker | aucune exigée si endpoint supporte | oui | opt-in, capability-dependent | +| Helius full transaction WS | `transactionSubscribe` | typed provider envelope, nested transaction JSON | signal riche prouvé incomplet pour RAW v1 : blockTime/version absents | hydration HTTP obligatoire | provider/session/filter sûrs + winner HTTP du matériau complet | reconnexion, pas historique garanti WSS | clé/tier opérateur requis selon plan | TR-D + hydration Worker | profil Helius éventuellement nécessaire plus tard | oui | opt-in/tier-dependent | +| Yellowstone transactions | `Subscribe.transactions` | DTO full KSP, stream/reconnect présents | projection body+meta structurés vers RAW v1 byte-identique | full si parité prouvée | provider + endpoint settings + filter + created_at disponibles | `from_slot`, replay info, reconnect snapshot | OrbitFlare Devnet / PublicNode selon accès | TR-C + TR-D + TR-E | profils V3 déjà présents | oui | oui opt-in | +| Yellowstone blocks | `Subscribe.blocks` | DTO block full + transactions | même gate de projection ; préserver slot/index/block_time | full si parité prouvée | idem | `from_slot`; upstream a corrigé replay blocks en juillet 2026 | mêmes providers | TR-C + TR-D + TR-E | non P0 | oui | oui opt-in | +| Yellowstone status + hydration | `transactions_status` + HTTP | DTO status + HTTP observed | composer signal + hydration | hydration requise | status filter/created_at + HTTP winner | replay selon provider + HTTP repair | mêmes providers | TR-D + TR-E | non P0 | oui | oui opt-in | +| Yellowstone block meta / slots | `blocks_meta`, `slots` | DTO + stream présents | aucun RawTransaction direct ; alimenter continuity tracker | non, auxiliaire | filter/slot/status/created_at | continuity evidence ; replay provider-specific | mêmes providers | TR-E côté Worker | non | oui | oui opt-in | +| EARLY source | provider/shred adapter | aucune surface générique unique | `TR-F`, uniquement protocole réellement implémenté | hydration sauf preuve meta complète | provider-specific | provider-specific | généralement payant | différé | différé | seulement si implémenté | seulement si accessible | ### 6.2 `TR-B` — `get_block_observed` @@ -479,17 +479,17 @@ Transaction V1 / champ config : proto >= 12.6 nécessaire, donc 12.7 courant est ### 7.3 Provider/access matrix revalidée -| Provider / source | État courant utile | Accès de preuve KSP | Décision `0.3.10` | -|:---------------------|:----------------------------------------------------------------------------------------------------|:-----------------------------------------------------|:-----------------------------------------------------------------------| -| Solana public RPC/WS | HTTP/WS standards accessibles ; blockSubscribe reste capability validator-dependent | sans secret, opt-in live | P0 HTTP/WS | -| PublicNode | Mainnet/Testnet RPC/WS/Yellowstone affichés ; archive access proposé, profondeur replay non prouvée | profils KSP + smokes ignored existants | P0 Yellowstone live ; replay seulement après preuve | -| OrbitFlare | Free/Developer : gRPC Devnet ; Yellowstone full fidelity ; Mainnet gRPC payant/add-on | profil Devnet + x-token operator | P0 Devnet live ; Mainnet non requis | -| Helius standard WSS | tx extension disponible ; blockSubscribe explicitement non supporté ; 10 min inactivity | dépend de clé/tier opérateur | txSubscribe adapter fixture + smoke opt-in si accès | -| Helius LaserStream | reconnect + replay jusqu’à 24 h documentés | Mainnet/tiers selon compte ; ne pas supposer gratuit | architecture compatible, live non bloquant pour release si tier absent | -| QuickNode | Yellowstone Scale+ ; port 443/x-token ; port 10000 legacy sunset 1 octobre 2026 | pas de tier gRPC actuel | bloqué tier ; aucun profil 10000 à introduire | -| Alchemy | Yellowstone Mainnet/Devnet PAYG/Enterprise | pas de tier actuel | bloqué ; replay non constant | -| Jito ShredStream | shutdown annoncé pour le 5 septembre 2026, migration DoubleZero recommandée | aucun besoin | `SUNSET`, ne pas implémenter TR-F dessus | -| EARLY autres | protocoles/prix variables | variable | adapter seulement après protocole+accès+RAW role prouvés | +| Provider / source | État courant utile | Accès de preuve KSP | Décision `0.3.10` | +|:---------------------|:----------------------------------------------------------------------------------------------------------------------|:-----------------------------------------------------|:-----------------------------------------------------------------------| +| Solana public RPC/WS | HTTP/WS standards accessibles ; blockSubscribe reste capability validator-dependent | sans secret, opt-in live | P0 HTTP/WS | +| PublicNode | Mainnet/Testnet RPC/WS/Yellowstone affichés ; archive access proposé, profondeur replay non prouvée | profils KSP + smokes ignored existants | P0 Yellowstone live ; replay seulement après preuve | +| OrbitFlare | Free/Developer : gRPC Devnet ; Yellowstone full fidelity ; Mainnet gRPC payant/add-on | profil Devnet + x-token operator | P0 Devnet live ; Mainnet non requis | +| Helius standard WSS | tx extension disponible ; blockSubscribe explicitement non supporté ; enveloppe full sans blockTime/version qualifiés | dépend de clé/tier opérateur | signal + hydration HTTP ; smoke opt-in si accès | +| Helius LaserStream | reconnect + replay jusqu’à 24 h documentés | Mainnet/tiers selon compte ; ne pas supposer gratuit | architecture compatible, live non bloquant pour release si tier absent | +| QuickNode | Yellowstone Scale+ ; port 443/x-token ; port 10000 legacy sunset 1 octobre 2026 | pas de tier gRPC actuel | bloqué tier ; aucun profil 10000 à introduire | +| Alchemy | Yellowstone Mainnet/Devnet PAYG/Enterprise | pas de tier actuel | bloqué ; replay non constant | +| Jito ShredStream | shutdown annoncé pour le 5 septembre 2026, migration DoubleZero recommandée | aucun besoin | `SUNSET`, ne pas implémenter TR-F dessus | +| EARLY autres | protocoles/prix variables | variable | adapter seulement après protocole+accès+RAW role prouvés | Alchemy reste volontairement `À REVALIDER` sur la profondeur replay : sa page dédiée Historical Replay annonce environ `432000` slots / `48 h`, tandis que son overview/SubscribeRequest mentionne encore `6000` slots. KSP ne doit encoder **aucune** constante Alchemy ; le runtime se fie à une capability/replay info observée ou à un paramètre configuré/prouvé. @@ -1062,9 +1062,17 @@ Conformément à `TR-C2`, aucun adapter productif Transport DTO -> `RawTransacti Le premier rejeu opérateur de `pre.004` a validé les audits, `cargo check`, les tests common, les 387 tests Transport et le canari cross-layer, mais Clippy strict a rejeté sept `expect()` placés dans deux helpers de fixture hors du corps `#[tokio::test]`. `pre.004-fix.001` corrige uniquement ces helpers en retournant explicitement des `Result`; aucun code de production, contrat RAW, surface Transport, golden ou dépendance n’est modifié. +Le rejeu opérateur du fix le 2026-09-06 est intégralement vert : audits statiques, `cargo check --workspace`, Clippy strict, tests common, 387 tests Transport avec suites publiques/release/doctests, et suite complète `ksp-job-backfill-lib`. Les arbres n’ont pas été rejoués car `pre.004-fix.001` n’a modifié aucune dépendance ni feature. + ### `pre.005` — parité WS/Helius full -Fermer projection `blockSubscribe` + Helius `transactionSubscribe` vers matériau common, golden parity ou fallback hydration explicitement qualifié. +**Statut : réalisé.** + +Deux canaris cross-layer test-only ferment la qualification sans ajouter de code de production. `blockSubscribe` standard en `confirmed/full/base64`, `maxSupportedTransactionVersion = 0`, `showRewards = false` reproduit le même `SolanaConfirmedBlock` que `getBlock` puis les mêmes référence, slot, block time, bytes RAW v1 et content hash pour une transaction legacy. Ce sous-ensemble legacy/v0 est donc qualifié direct ; Transaction V1 reste explicitement au gate `pre.006`. + +Helius `transactionSubscribe` full/base64 conserve la signature, le slot, le transaction wire, la meta et `transactionIndex`, mais la notification qualifiée ne fournit ni `blockTime` ni `version`. Le canari construit la meilleure projection possible puis prouve qu’elle diverge de l’hydration HTTP sur block time, bytes RAW v1 et content hash malgré une identité/slot identiques. Helius est donc verrouillé comme signal live riche + hydration HTTP, jamais comme producteur RAW-direct dans ce contrat. + +`tokio-tungstenite` est ajouté uniquement en `dev-dependencies` du Backfill pour les serveurs WebSocket locaux déterministes ; un canari de dependency boundary interdit sa migration vers les dépendances de production. `TR-C2` reste Worker-owned. ### `pre.006` — parité Yellowstone transactions/blocks diff --git a/docs/validation/027-V0_3_10_RAW_TRANSACTION_INGEST.md b/docs/validation/027-V0_3_10_RAW_TRANSACTION_INGEST.md index 6a8b174..f2156d9 100644 --- a/docs/validation/027-V0_3_10_RAW_TRANSACTION_INGEST.md +++ b/docs/validation/027-V0_3_10_RAW_TRANSACTION_INGEST.md @@ -1,5 +1,5 @@ - + # Validation v0.3.10 — RAW Transaction commune + Worker d’ingestion @@ -228,8 +228,8 @@ Familles à prouver séparément : ```text HTTP getTransaction HTTP getBlock transaction -WS blockSubscribe transaction -Helius transactionSubscribe full +WS blockSubscribe transaction legacy/v0 full/base64 +Helius transactionSubscribe full/base64 -> negative parity / hydration Yellowstone transaction Yellowstone block transaction Transaction V1 @@ -796,5 +796,59 @@ Version technique du fix : workspace.package.version = 0.3.10-pre.4.fix.1 ``` -Le gate complet de `pre.004` doit être rejoué après application du fix avant ouverture de `pre.005`. +Le rejeu du fix communiqué le 2026-09-06 est intégralement vert : + +```text +cargo fmt --all : exécuté +General Rust rule audit: clean +Rust export completeness audit: 0 candidate(s) +KSP workspace Rust rule audit: clean +Markdown table audit: clean (339 table(s), 758 file(s)) +cargo check --workspace : PASS +cargo clippy --workspace --all-targets --all-features -- -D warnings : PASS +cargo test -p ksp-raw-transaction-lib : PASS, 12 unit + 9 intégration, 0 échec +cargo test -p ksp-onchain-transport-lib : PASS, 387 unit + 51 public API + 43 release completeness + 4 doctests, 0 échec ; smokes live ignorés comme prévu +cargo test -p ksp-job-backfill-lib : PASS, 51 unit + 21 intégration, 0 échec +``` + +Les arbres Cargo n’ont pas été rejoués sur ce fix car aucune dépendance ni feature n’a changé. `pre.004-fix.001` est donc fermé et `pre.005` peut s’ouvrir. + +### 13.7 Qualification WS/Helius de `pre.005` + +Le canari `ksp-job-backfill-lib/tests/ws_raw_parity.rs` utilise uniquement des serveurs HTTP/WS locaux déterministes. Aucun smoke réseau ni secret provider n’est requis. + +Cas positif standard : + +```text +HTTP getBlock observed +WS blockSubscribe confirmed/full/base64/maxSupportedTransactionVersion=0 + -> même SolanaConfirmedBlock + -> même RawTransactionReference + -> même slot + -> même block_time + -> mêmes bytes RAW v1 + -> même content_hash +``` + +Ce gate qualifie uniquement le sous-ensemble legacy/v0. Transaction V1 reste au gate `pre.006`. + +Cas négatif Helius : + +```text +Helius transactionSubscribe full/base64 + -> même signature/identité que HTTP + -> même slot + -> transaction wire + meta + transactionIndex présents + -> blockTime absent + -> version absente + +meilleure projection Helius + != hydration HTTP complète sur block_time + != bytes RAW v1 + != content_hash +``` + +Le résultat est normatif pour le Worker : Helius `transactionSubscribe` est un signal live riche qui doit être hydraté par HTTP avant persistence RAW. Aucun champ absent n’est inventé et aucun RAW v2 n’est introduit. + +La dépendance `tokio-tungstenite` nécessaire aux fixtures WS est strictement `dev-dependencies`; `tests/dependency_boundary.rs` verrouille cette frontière. Aucun edge de production Transport -> common n’est ajouté.