diff --git a/Cargo.toml b/Cargo.toml index f9784c9..f46df22 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -6,7 +6,7 @@ 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", "crates/ksp-worker-raw-transaction-ingest-lib"] [workspace.package] -version = "0.3.11-pre.5" +version = "0.3.11-pre.6" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/src/admission.rs b/crates/ksp-worker-raw-transaction-ingest-lib/src/admission.rs new file mode 100644 index 0000000..e97537e --- /dev/null +++ b/crates/ksp-worker-raw-transaction-ingest-lib/src/admission.rs @@ -0,0 +1,93 @@ +// file: crates/ksp-worker-raw-transaction-ingest-lib/src/admission.rs +// version: 1 + +use sha2::Digest; // rust-rules: trait-import + +const RAW_TRANSACTION_INGEST_OBSERVATION_DOMAIN: &[u8] = b"ksp.raw_transaction_ingest.observation.v1\0"; + +/// Crate-private source-neutral ingress admitted by the central bounded Worker queue. +pub(crate) struct RawTransactionIngress { + /// Complete Common RAW material supplied by one private source task. + pub(crate) material: ksp_raw_transaction_lib::RawTransactionMaterial, + /// Logical network expected to match both Worker settings and canonical material. + pub(crate) network: ksp_store_lib::RawNetworkId, + /// Safe source-independent provenance attached to the acquisition observation. + pub(crate) provenance: ksp_store_lib::RawAcquisitionProvenance, + /// Opaque deterministic source-owned key material used only for Worker observation-key derivation. + pub(crate) source_key: [u8; 32], +} + +/// Receiver side of the bounded central RAW transaction admission queue owned by the Worker supervisor. +pub(crate) struct RawTransactionAdmission { + receiver: tokio::sync::mpsc::Receiver, +} + +impl crate::RawTransactionAdmission { + /// Creates one bounded admission queue and returns its crate-private source sender. + #[must_use] + pub(crate) fn new(capacity: usize) -> (crate::RawTransactionAdmission, tokio::sync::mpsc::Sender) { + let (sender, receiver) = tokio::sync::mpsc::channel(capacity); + return (Self { receiver }, sender); + } + + /// Closes new admissions while preserving already queued ingress for deterministic drain. + pub(crate) fn close(&mut self) { + self.receiver.close(); + return; + } + + /// Receives and converts one queued ingress into a common RAW transaction acquisition. + pub(crate) async fn receive( + &mut self, + expected_network: &ksp_store_lib::RawNetworkId, + ) -> ksp_core_lib::Result> { + let ingress = match self.receiver.recv().await { + std::option::Option::Some(value) => value, + std::option::Option::None => return std::result::Result::Ok(std::option::Option::None), + }; + let acquisition = canonicalize_ingress(expected_network, ingress); + return match acquisition { + std::result::Result::Ok(value) => std::result::Result::Ok(std::option::Option::Some(value)), + std::result::Result::Err(error) => std::result::Result::Err(error), + }; + } +} + +fn canonicalize_ingress( + expected_network: &ksp_store_lib::RawNetworkId, + ingress: crate::RawTransactionIngress, +) -> ksp_core_lib::Result { + if &ingress.network != expected_network { + return std::result::Result::Err(crate::runtime_error("admission.network_mismatch")); + } + let transaction = ksp_raw_transaction_lib::canonicalize_raw_transaction(ingress.material); + let transaction = match transaction { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(crate::runtime_error("admission.canonicalization_invalid")), + }; + if transaction.reference().network() != expected_network { + return std::result::Result::Err(crate::runtime_error("admission.material_network_mismatch")); + } + let observation_key = observation_key(transaction.reference(), &ingress.source_key); + return std::result::Result::Ok(ksp_raw_transaction_lib::assemble_raw_transaction_acquisition(transaction, observation_key, ingress.provenance)); +} + +fn hash_bytes(hasher: &mut sha2::Sha256, value: &[u8]) { + hasher.update((value.len() as u64).to_be_bytes()); + hasher.update(value); + return; +} + +fn observation_key(reference: &ksp_store_lib::RawTransactionReference, source_key: &[u8; 32]) -> ksp_store_lib::RawObservationKey { + let mut hasher = sha2::Sha256::new(); + hasher.update(RAW_TRANSACTION_INGEST_OBSERVATION_DOMAIN); + hash_bytes(&mut hasher, reference.network().as_str().as_bytes()); + hash_bytes(&mut hasher, reference.signature().as_bytes()); + hash_bytes(&mut hasher, source_key); + let bytes: [u8; 32] = hasher.finalize().into(); + return ksp_store_lib::RawObservationKey::new(bytes); +} + +#[cfg(test)] +#[path = "../unit_tests/admission.rs"] +mod tests; diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs b/crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs index 3220126..b8e3e84 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs -// version: 4 +// version: 5 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -8,10 +8,11 @@ //! Source-neutral runtime foundation for continuous KSP RAW transaction ingestion. //! //! This tranche owns the concrete Worker family identity, validated technical settings -//! and the caller-runtime-owned lifecycle with private child-task supervision. Bounded -//! admission, persistence and latest-value snapshots remain in their dedicated prereleases; no -//! live source or Transport dependency exists here. +//! and the caller-runtime-owned lifecycle with private child-task supervision. This tranche also +//! owns bounded source-neutral admission plus common RAW canonicalization/assembly; persistence and +//! latest-value snapshots remain in later prereleases, and no live source or Transport dependency exists. +mod admission; mod error; mod identity; mod runtime; @@ -50,6 +51,10 @@ pub use self::settings::MIN_RAW_TRANSACTION_INGEST_SHUTDOWN_DRAIN_TIMEOUT; /// Validated source-neutral runtime settings for one continuous RAW transaction ingest Worker. pub use self::settings::RawTransactionIngestSettings; +/// Receiver-side owner of the private bounded RAW transaction admission queue. +pub(crate) use self::admission::RawTransactionAdmission; +/// Crate-private source-neutral ingress sent through the bounded central admission queue. +pub(crate) use self::admission::RawTransactionIngress; /// Creates one runtime-domain error without copying runtime/provider/Store values into diagnostics. pub(crate) use self::error::runtime_error; /// Creates one settings-domain error without copying caller-supplied values into diagnostics. diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs b/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs index c51de2b..3939604 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs -// version: 2 +// version: 3 /// Runtime-neutral boxed future resolving after one RAW transaction ingest Worker has fully reached a terminal lifecycle state. pub type RawTransactionIngestTerminalFuture<'a> = @@ -85,6 +85,19 @@ fn current_runtime_handle() -> ksp_core_lib::Result { }; } +async fn drain_admission(settings: &crate::RawTransactionIngestSettings, admission: &mut crate::RawTransactionAdmission) -> bool { + let mut clean = true; + admission.close(); + loop { + let received = admission.receive(settings.network()).await; + match received { + std::result::Result::Ok(std::option::Option::Some(_acquisition)) => {}, + std::result::Result::Ok(std::option::Option::None) => return clean, + std::result::Result::Err(_) => clean = false, + } + } +} + async fn drain_children(children: &mut tokio::task::JoinSet<()>) -> bool { let mut clean = true; while let std::option::Option::Some(joined) = children.join_next().await { @@ -119,13 +132,16 @@ fn finish_stopped(lifecycle: &mut ksp_worker_api::WorkerLifecycle, sender: &toki } async fn run_supervisor( + settings: crate::RawTransactionIngestSettings, mut lifecycle: ksp_worker_api::WorkerLifecycle, _store_guard: std::option::Option>, mut stop_receiver: tokio::sync::watch::Receiver, terminal_sender: tokio::sync::watch::Sender, source_spawner: Spawner, ) where - Spawner: FnOnce(&mut tokio::task::JoinSet<()>, tokio::sync::watch::Receiver) + std::marker::Send + 'static, + Spawner: FnOnce(&mut tokio::task::JoinSet<()>, tokio::sync::watch::Receiver, tokio::sync::mpsc::Sender) + + std::marker::Send + + 'static, { if *stop_receiver.borrow() { finish_stopped(&mut lifecycle, &terminal_sender); @@ -137,9 +153,12 @@ async fn run_supervisor( } terminal_sender.send_replace(lifecycle.state()); let mut children = tokio::task::JoinSet::new(); - source_spawner(&mut children, stop_receiver.clone()); - let supervised_clean = supervise_until_stop(&mut stop_receiver, &mut children).await; - if !supervised_clean { + let (mut admission, admission_sender) = crate::RawTransactionAdmission::new(settings.admission_queue_capacity()); + source_spawner(&mut children, stop_receiver.clone(), admission_sender); + let supervised_clean = supervise_until_stop(&settings, &mut stop_receiver, &mut children, &mut admission).await; + let admission_clean = drain_admission(&settings, &mut admission).await; + let children_clean = drain_children(&mut children).await; + if !supervised_clean || !admission_clean || !children_clean { finish_faulted(&mut lifecycle, &terminal_sender); return; } @@ -152,7 +171,7 @@ fn start_foundation( runtime: tokio::runtime::Handle, store_guard: std::option::Option>, ) -> ksp_core_lib::Result { - return start_foundation_with_source_spawner(settings, runtime, store_guard, |_children, _stop_receiver| {}); + return start_foundation_with_source_spawner(settings, runtime, store_guard, |_children, _stop_receiver, _admission_sender| {}); } fn start_foundation_with_source_spawner( @@ -162,7 +181,9 @@ fn start_foundation_with_source_spawner( source_spawner: Spawner, ) -> ksp_core_lib::Result where - Spawner: FnOnce(&mut tokio::task::JoinSet<()>, tokio::sync::watch::Receiver) + std::marker::Send + 'static, + Spawner: FnOnce(&mut tokio::task::JoinSet<()>, tokio::sync::watch::Receiver, tokio::sync::mpsc::Sender) + + std::marker::Send + + 'static, { let kind = match ksp_worker_api::WorkerKindCode::new(crate::RAW_TRANSACTION_INGEST_WORKER_KIND_CODE) { std::result::Result::Ok(value) => value, @@ -176,17 +197,23 @@ where let (stop_sender, stop_receiver) = tokio::sync::watch::channel(false); let (terminal_sender, terminal_receiver) = tokio::sync::watch::channel(lifecycle.state()); let handle = crate::RawTransactionIngestHandle { stop_sender, stop_token, terminal_receiver }; - std::mem::drop(runtime.spawn(run_supervisor(lifecycle, store_guard, stop_receiver, terminal_sender, source_spawner))); + std::mem::drop(runtime.spawn(run_supervisor(settings, lifecycle, store_guard, stop_receiver, terminal_sender, source_spawner))); return std::result::Result::Ok(handle); } -async fn supervise_until_stop(stop_receiver: &mut tokio::sync::watch::Receiver, children: &mut tokio::task::JoinSet<()>) -> bool { +async fn supervise_until_stop( + settings: &crate::RawTransactionIngestSettings, + stop_receiver: &mut tokio::sync::watch::Receiver, + children: &mut tokio::task::JoinSet<()>, + admission: &mut crate::RawTransactionAdmission, +) -> bool { + let mut admission_open = true; let mut clean = true; loop { if *stop_receiver.borrow() { break; } - if children.is_empty() { + if children.is_empty() && !admission_open { let changed = stop_receiver.changed().await; if changed.is_err() || *stop_receiver.borrow() { break; @@ -200,15 +227,21 @@ async fn supervise_until_stop(stop_receiver: &mut tokio::sync::watch::Receiver { + joined = children.join_next(), if !children.is_empty() => { if let std::option::Option::Some(std::result::Result::Err(_)) = joined { clean = false; } } + received = admission.receive(settings.network()), if admission_open => { + match received { + std::result::Result::Ok(std::option::Option::Some(_acquisition)) => {}, + std::result::Result::Ok(std::option::Option::None) => admission_open = false, + std::result::Result::Err(_) => clean = false, + } + } } } - let drained_clean = drain_children(children).await; - return clean && drained_clean; + return clean; } fn validate_store_network(settings: &crate::RawTransactionIngestSettings, store_network: &ksp_store_lib::RawNetworkId) -> ksp_core_lib::Result<()> { diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/tests/dependency_boundary.rs b/crates/ksp-worker-raw-transaction-ingest-lib/tests/dependency_boundary.rs index 39321e3..332fcc7 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/tests/dependency_boundary.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/tests/dependency_boundary.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/tests/dependency_boundary.rs -// version: 4 +// version: 5 //! Dependency firewall canaries for the RAW transaction ingest Worker foundation. @@ -52,22 +52,37 @@ fn pre_002_manifest_keeps_forbidden_layers_and_live_sources_out() { } #[test] -fn pre_005_source_surface_owns_private_supervisor_without_admission_persistence_or_live_sources() { +fn pre_006_source_surface_opens_bounded_admission_and_common_raw_without_persistence_or_live_sources() { let root = include_str!("../src/lib.rs"); let runtime = include_str!("../src/runtime.rs"); - for required in ["tokio::task::JoinSet", "run_supervisor", "supervise_until_stop", "drain_children", "start_foundation_with_source_spawner"] { - assert!(runtime.contains(required), "required pre.005 private supervisor contract missing: {required}"); + let admission = include_str!("../src/admission.rs"); + for required in [ + "tokio::task::JoinSet", + "tokio::sync::mpsc::channel", + "canonicalize_raw_transaction", + "assemble_raw_transaction_acquisition", + "ksp.raw_transaction_ingest.observation.v1", + "RawTransactionIngress", + ] { + assert!( + root.contains(required) || runtime.contains(required) || admission.contains(required), + "required pre.006 bounded admission/common RAW contract missing: {required}" + ); } for forbidden in [ - "pub use self::runtime::JoinSet", - "mpsc::", - "canonicalize_raw_transaction", + "pub use self::admission::RawTransactionAdmission", + "pub use self::admission::RawTransactionIngress", + "unbounded_channel", "persist_raw_transaction", + "RawTransactionAcquisitionMode", "ksp_onchain_transport_lib::", "RawTransactionIngestSnapshot", "WorkerSnapshotSource", ] { - assert!(!root.contains(forbidden) && !runtime.contains(forbidden), "pre.005 opened later runtime scope too early: {forbidden}"); + assert!( + !root.contains(forbidden) && !runtime.contains(forbidden) && !admission.contains(forbidden), + "pre.006 opened later runtime scope too early: {forbidden}" + ); } return; } diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/admission.rs b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/admission.rs new file mode 100644 index 0000000..3cf16c1 --- /dev/null +++ b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/admission.rs @@ -0,0 +1,241 @@ +// file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/admission.rs +// version: 1 + +fn material( + network: &str, + signature_byte: u8, + transaction_data: &str, +) -> std::option::Option<(ksp_store_lib::RawNetworkId, ksp_raw_transaction_lib::RawTransactionMaterial)> { + let network = match ksp_store_lib::RawNetworkId::new(network) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::option::Option::None, + }; + let signature = ksp_store_lib::RawTransactionSignature::new([signature_byte; 64]); + let material = ksp_raw_transaction_lib::RawTransactionMaterial::binary_base64( + network.clone(), + signature, + 42, + std::option::Option::Some(1_700_000_000), + transaction_data, + ksp_raw_transaction_lib::RawTransactionWireField::Omitted, + ksp_raw_transaction_lib::RawTransactionWireField::Omitted, + ksp_raw_transaction_lib::RawTransactionWireField::Omitted, + ); + return std::option::Option::Some((network, material)); +} + +fn provenance() -> std::option::Option { + let provider = match ksp_store_lib::RawProvenanceCode::new("deterministic-harness") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::option::Option::None, + }; + let protocol = match ksp_store_lib::RawProvenanceCode::new("internal") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::option::Option::None, + }; + let method = match ksp_store_lib::RawProvenanceCode::new("admission") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::option::Option::None, + }; + let received_at = match ksp_store_lib::RawTimestamp::from_unix_millis(1_700_000_000_000) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::option::Option::None, + }; + return std::option::Option::Some(ksp_store_lib::RawAcquisitionProvenance::new( + provider, + protocol, + method, + ksp_store_lib::RawAcquisitionOrigin::Live, + received_at, + )); +} + +#[tokio::test(flavor = "current_thread")] +async fn pre_006_bounded_admission_applies_async_backpressure_until_one_slot_is_consumed() { + let (network, first_material) = match material("mainnet", 1, "AQID") { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let (_, second_material) = match material("mainnet", 2, "BAUG") { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let first_provenance = match provenance() { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let second_provenance = match provenance() { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let (mut admission, sender) = crate::RawTransactionAdmission::new(1); + let first_ingress = crate::RawTransactionIngress { + material: first_material, + network: network.clone(), + provenance: first_provenance, + source_key: [1; 32], + }; + assert!(sender.send(first_ingress).await.is_ok()); + let second_sender = sender.clone(); + let second_network = network.clone(); + let completed = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)); + let completed_task = std::sync::Arc::clone(&completed); + let second_task = tokio::spawn(async move { + let second_ingress = crate::RawTransactionIngress { + material: second_material, + network: second_network, + provenance: second_provenance, + source_key: [2; 32], + }; + let admitted = second_sender.send(second_ingress).await.is_ok(); + completed_task.store(admitted, std::sync::atomic::Ordering::Release); + return admitted; + }); + tokio::task::yield_now().await; + assert!(!completed.load(std::sync::atomic::Ordering::Acquire)); + let first = admission.receive(&network).await; + assert!(matches!(first, std::result::Result::Ok(std::option::Option::Some(_)))); + for _ in 0..64 { + if completed.load(std::sync::atomic::Ordering::Acquire) { + break; + } + tokio::task::yield_now().await; + } + assert!(completed.load(std::sync::atomic::Ordering::Acquire)); + let joined = second_task.await; + assert!(matches!(joined, std::result::Result::Ok(true))); + let second = admission.receive(&network).await; + assert!(matches!(second, std::result::Result::Ok(std::option::Option::Some(_)))); + return; +} + +#[tokio::test(flavor = "current_thread")] +async fn pre_006_stop_preempts_a_blocked_sender_without_silent_post_stop_admission() { + let (network, first_material) = match material("mainnet", 3, "BwgJ") { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let (_, second_material) = match material("mainnet", 4, "CgsM") { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let first_provenance = match provenance() { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let second_provenance = match provenance() { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let (stop_sender, stop_receiver) = tokio::sync::watch::channel(false); + let (mut admission, sender) = crate::RawTransactionAdmission::new(1); + let first_ingress = crate::RawTransactionIngress { + material: first_material, + network: network.clone(), + provenance: first_provenance, + source_key: [3; 32], + }; + assert!(sender.send(first_ingress).await.is_ok()); + let blocked_sender = sender.clone(); + let blocked_network = network.clone(); + let blocked_task = tokio::spawn(async move { + let mut blocked_stop = stop_receiver; + let blocked_ingress = crate::RawTransactionIngress { + material: second_material, + network: blocked_network, + provenance: second_provenance, + source_key: [4; 32], + }; + return tokio::select! { + biased; + _changed = blocked_stop.changed() => false + result = blocked_sender.send(blocked_ingress) => result.is_ok() + }; + }); + tokio::task::yield_now().await; + stop_sender.send_replace(true); + let blocked = blocked_task.await; + assert!(matches!(blocked, std::result::Result::Ok(false))); + let first = admission.receive(&network).await; + assert!(matches!(first, std::result::Result::Ok(std::option::Option::Some(_)))); + admission.close(); + let end = admission.receive(&network).await; + assert!(matches!(end, std::result::Result::Ok(std::option::Option::None))); + return; +} + +#[test] +fn pre_006_observation_key_domain_and_common_raw_assembly_golden_are_exact() { + let (network, material) = match material("mainnet", 7, "AQID") { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let provenance = match provenance() { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let ingress = crate::RawTransactionIngress { material, network: network.clone(), provenance: provenance.clone(), source_key: [9; 32] }; + let acquisition = super::canonicalize_ingress(&network, ingress); + assert!(acquisition.is_ok()); + let acquisition = match acquisition { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + assert_eq!(acquisition.transaction().reference().network(), &network); + assert_eq!(acquisition.transaction().reference().signature().as_bytes(), &[7; 64]); + assert_eq!(acquisition.transaction().slot(), 42); + assert_eq!(acquisition.transaction().payload().format_id().as_str(), ksp_raw_transaction_lib::RAW_TRANSACTION_FORMAT_ID); + assert_eq!(acquisition.transaction().payload().format_version(), ksp_raw_transaction_lib::RAW_TRANSACTION_FORMAT_VERSION); + assert_eq!(acquisition.transaction().payload().bytes(), br#"{"transaction":["AQID","base64"]}"#); + assert_eq!(acquisition.observation().provenance(), &provenance); + assert_eq!(acquisition.observation().transaction(), acquisition.transaction().reference()); + assert_eq!( + acquisition.observation().observation_key().as_bytes(), + &[ + 248, 209, 97, 220, 168, 168, 167, 42, 13, 161, 16, 252, 135, 136, 46, 231, 19, 104, 72, 70, 28, 106, 163, 43, 9, 227, 10, 97, 136, 185, 93, 204, + ], + ); + return; +} + +#[test] +fn pre_006_network_guards_reject_ingress_or_material_mismatch_without_echoing_values() { + let (mainnet, mainnet_material) = match material("mainnet", 8, "DQ4P") { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let (devnet, devnet_material) = match material("devnet", 9, "EBES") { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let first_provenance = match provenance() { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let second_provenance = match provenance() { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let wrong_ingress = crate::RawTransactionIngress { material: mainnet_material, network: devnet, provenance: first_provenance, source_key: [10; 32] }; + let wrong_ingress = super::canonicalize_ingress(&mainnet, wrong_ingress); + assert!(wrong_ingress.is_err()); + let wrong_material = crate::RawTransactionIngress { + material: devnet_material, + network: mainnet.clone(), + provenance: second_provenance, + source_key: [11; 32], + }; + let wrong_material = super::canonicalize_ingress(&mainnet, wrong_material); + assert!(wrong_material.is_err()); + for result in [wrong_ingress, wrong_material] { + let error = match result { + std::result::Result::Ok(_) => continue, + std::result::Result::Err(value) => value, + }; + assert_eq!(error.code(), crate::ERROR_CODE_RAW_TRANSACTION_INGEST_RUNTIME_INVALID); + let rendered = std::format!("{error:?}"); + assert!(!rendered.contains("mainnet")); + assert!(!rendered.contains("devnet")); + } + return; +} diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs index ef61f89..7596de4 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs -// version: 2 +// version: 3 struct ActiveTaskGuard { active: std::sync::Arc, @@ -166,7 +166,7 @@ async fn pre_005_supervisor_joins_all_cooperative_source_tasks_before_terminal() settings, tokio::runtime::Handle::current(), std::option::Option::None, - move |children, stop_receiver| { + move |children, stop_receiver, _admission_sender| { for _ in 0..3 { let active = std::sync::Arc::clone(&source_active); let mut child_stop = stop_receiver.clone(); @@ -211,7 +211,7 @@ async fn pre_005_supervisor_reaps_completed_source_task_and_still_joins_live_chi settings, tokio::runtime::Handle::current(), std::option::Option::None, - move |children, stop_receiver| { + move |children, stop_receiver, _admission_sender| { let completed = std::sync::Arc::clone(&source_completed); let _completed_abort_handle = children.spawn(async move { completed.fetch_add(1, std::sync::atomic::Ordering::AcqRel); diff --git a/deltas/0.3.11/pre.006.md b/deltas/0.3.11/pre.006.md new file mode 100644 index 0000000..def5922 --- /dev/null +++ b/deltas/0.3.11/pre.006.md @@ -0,0 +1,180 @@ + + + +# Delta `0.3.11-pre.006` — admission bornée et canonicalisation Common RAW + +## Base requise + +```text +0.3.11-pre.005 +workspace.package.version = 0.3.11-pre.5 +``` + +Le gate opérateur du 8 septembre 2026 est vert sur `fmt`, audits Rust/Markdown, `cargo check --workspace`, Clippy strict, les 18 tests de la crate Worker, doc-tests et les arbres Cargo normal/features. + +## Objectif + +Matérialiser uniquement la responsabilité `pre.006` du plan `032` : admission `mpsc` centrale bornée, backpressure async, ingress privé, canonicalisation via `ksp-raw-transaction-lib`, observation key Worker déterministe et assembly du `RawTransactionAcquisition`. Aucune persistence Store et aucune source réseau. + +## Version + +```text +identifiant de livraison : 0.3.11-pre.006 +workspace.package.version : 0.3.11-pre.6 +``` + +## Admission privée bornée + +Le supervisor crée un seul channel central dimensionné par `RawTransactionIngestSettings::admission_queue_capacity()` : + +```text +tokio::sync::mpsc::channel(...) +``` + +Le seam source reçoit le sender `mpsc` borné et le signal de stop. Le contrat prouvé est : + +```text +send().await pour le backpressure +select! biased : stop avant send dans le harness source +aucun unbounded_channel +receiver close au shutdown +``` + +Au shutdown, le receiver ferme les nouvelles admissions et consume les entrées déjà en queue avant le join final des tâches enfants. + +## Ingress et Common RAW + +L'ingress crate-private, absent de l'API publique, transporte : + +```text +RawTransactionMaterial +RawNetworkId attendu +RawAcquisitionProvenance +source key opaque [u8; 32] +``` + +Le pipeline est : + +```text +network guard +ksp_raw_transaction_lib::canonicalize_raw_transaction(...) +material network guard +observation key Worker +ksp_raw_transaction_lib::assemble_raw_transaction_acquisition(...) +``` + +Aucune canonicalisation RAW v1 n'est réimplémentée dans le Worker. + +## Observation key V1 + +Domaine exact : + +```text +ksp.raw_transaction_ingest.observation.v1\0 +``` + +Le SHA-256 encode avec longueur `u64` big-endian : + +```text +network +signature +source key +``` + +Golden harness : + +```text +network : mainnet +signature : 64 x 0x07 +source key : 32 x 0x09 +SHA-256 : f8d161dca8a8a72a0da110fc87882ee7136848461c6aa32b09e30a6188b95dcc +``` + +Aucun timestamp local, run id ou ordre d'admission ne participe à la clé. + +## Tests ajoutés ou étendus + +```text +capacity=1 : second send bloqué jusqu'à consommation d'un slot +stop : sender bloqué rejeté sans admission post-stop +payload common RAW canonique exact +format id/version exacts +observation-key golden exact +assembly transaction/provenance/reference exact +mismatch network ingress/material rejetés sans fuite de valeurs +dependency boundary : mpsc/common RAW requis, persistence/Transport/snapshot toujours interdits +``` + +## Fichiers ajoutés + +```text +crates/ksp-worker-raw-transaction-ingest-lib/src/admission.rs +crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/admission.rs +deltas/0.3.11/pre.006.md +``` + +## Fichiers modifiés + +```text +Cargo.toml +crates/ksp-worker-raw-transaction-ingest-lib/src/lib.rs +crates/ksp-worker-raw-transaction-ingest-lib/src/runtime.rs +crates/ksp-worker-raw-transaction-ingest-lib/tests/dependency_boundary.rs +crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime.rs +docs/plans/032-V0_3_11_RAW_TRANSACTION_INGEST_WORKER_FOUNDATION_PLAN.md +docs/validation/028-V0_3_11_RAW_TRANSACTION_INGEST_WORKER_FOUNDATION.md +``` + +## Fichiers supprimés + +Aucun. + +## Hors scope conservé + +```text +RawTransactionWrite / persistence Store +RawTransactionAcquisitionMode::Normal +persistence concurrency +Store outcome mapping / idempotence / conflict +snapshots concrets +source live / Transport +Config +retry / reconnect / gap repair +``` + +## Validations exécutées dans l'environnement d'assemblage + +```text +General Rust rule audit: clean +Rust export completeness audit: 0 candidate(s) +KSP workspace Rust rule audit: clean +Markdown table audit: clean (340 table(s), 781 file(s)) + +scan production : aucun ?, unwrap, expect, panic, unbounded_channel +scan scope : aucune persistence Store, aucun Transport, aucun snapshot concret +pipeline : exactement 1 mpsc::channel, 1 canonicalize_raw_transaction, 1 assemble_raw_transaction_acquisition +observation golden indépendant : f8d161dca8a8a72a0da110fc87882ee7136848461c6aa32b09e30a6188b95dcc +``` + +L'environnement d'assemblage ne fournit ni `cargo`, ni `rustc`, ni `rustfmt`. Aucun gate Cargo de `pre.006` 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-worker-raw-transaction-ingest-lib +cargo tree -p ksp-worker-raw-transaction-ingest-lib --edges normal +cargo tree -p ksp-worker-raw-transaction-ingest-lib -e features +``` + +## Décision + +`pre.007` reste bloquée jusqu'à validation opérateur verte de l'admission bornée, du backpressure, du golden observation et de l'assembly Common RAW. + +## Questions ouvertes + +Aucune nouvelle question architecturale. La persistence réelle et les outcomes Store restent dans `pre.007`, les snapshots dans `pre.008`, et le hardening fault/drain dans `pre.009`. diff --git a/docs/plans/032-V0_3_11_RAW_TRANSACTION_INGEST_WORKER_FOUNDATION_PLAN.md b/docs/plans/032-V0_3_11_RAW_TRANSACTION_INGEST_WORKER_FOUNDATION_PLAN.md index 02a7434..fd91813 100644 --- a/docs/plans/032-V0_3_11_RAW_TRANSACTION_INGEST_WORKER_FOUNDATION_PLAN.md +++ b/docs/plans/032-V0_3_11_RAW_TRANSACTION_INGEST_WORKER_FOUNDATION_PLAN.md @@ -1,5 +1,5 @@ - + # Plan v0.3.11 — fondation runtime du Worker RawTransaction ingest @@ -344,7 +344,7 @@ Les `JoinHandle`/`JoinSet` restent strictement privés. Dropper le handle public Un seul channel central : ```text -tokio::sync::mpsc::channel(admission_queue_capacity) +tokio::sync::mpsc::channel(admission_queue_capacity) ``` Le channel est borné. `send().await` applique le backpressure ; aucun `unbounded_channel`, aucun drop silencieux et aucune event queue publique ne sont admis. @@ -361,13 +361,14 @@ Un `watch` privé conserve uniquement la dernière ## 9. Seam/harness déterministe sans réseau -Le runtime de fondation est prouvé par un harness privé qui produit des `PrivateRawTransactionIngress` contrôlés. +Le runtime de fondation est prouvé par un harness privé qui produit des `RawTransactionIngress` crate-private contrôlés. Une entrée de harness contient conceptuellement : ```text RawTransactionMaterial -RawObservationKey déterministe Worker-owned +RawNetworkId attendu +source key privée déterministe [u8; 32] RawAcquisitionProvenance sûre ``` @@ -413,9 +414,9 @@ La Worker observation key est un SHA-256 KSP-owned versionné et domain-separate ksp.raw_transaction_ingest.observation.v1 ``` -Les bytes exacts/golden seront figés avec le premier code de persistence, à partir d'éléments stables fournis par la source technique et de l'identité `(network, signature)`. Aucun timestamp local aléatoire ni run id ne peut rendre une rediffusion identique non idempotente. +Les bytes exacts/golden sont figés en `pre.006` à partir de l'identité `(network, signature)` et d'une source key privée opaque de 32 bytes. Le SHA-256 est domain-separated puis encode chaque composant avec une longueur `u64` big-endian. Aucun timestamp local aléatoire, run id ni ordre d'admission ne peut rendre une rediffusion identique non idempotente. -Les éléments source exacts ne sont pas figés publiquement avant `0.3.12`; le harness utilise une source key privée déterministe afin de prouver l'idempotence du domaine Worker. +Les éléments provider/source utilisés pour produire cette source key ne sont pas figés publiquement avant `0.3.12`; le harness utilise une source key privée déterministe afin de prouver l'idempotence du domaine Worker sans ouvrir une API source prématurée. ### 10.3 Outcomes @@ -719,12 +720,14 @@ Les méthodes `snapshot_source()` et `worker_snapshot_source()` prévues par la Budget cible : **15–20 min**. Introduire le supervisor, les JoinSet/joins privés, intégrer le wake-up de stop déjà matérialisé à l'ownership des tâches enfants, ajouter le test source seam et prouver qu'aucune tâche enfant ne survit au terminal. Pas encore de persistence réelle. -État après matérialisation : **implémenté, gate opérateur requis**. Le task racine est désormais le supervisor privé et possède un `tokio::task::JoinSet<()>`. Un seam privé de spawn de sources reçoit le `JoinSet` et un clone du `watch` de stop ; la production utilise un seam vide tandis que les unit tests injectent des tâches coopératives déterministes. Le supervisor récole les enfants terminés pendant l'exécution puis, au stop ou à la fermeture du channel de contrôle, attend la fin de tous les enfants avant de publier `Stopped`. Un `JoinError` observé marque la supervision comme non clean ; si la fermeture de supervision est ensuite engagée, le terminal est projeté vers `Faulted(worker_raw_transaction_ingest.runtime_invalid)`. Le traitement immédiat/prioritaire des faults reste différé à `pre.009`. Aucun timeout/abort forcé n'est encore introduit : ce hardening reste réservé à `pre.009`. +État après matérialisation : **implémenté et gate opérateur validé**. Le task racine est désormais le supervisor privé et possède un `tokio::task::JoinSet<()>`. Un seam privé de spawn de sources reçoit le `JoinSet` et un clone du `watch` de stop ; la production utilise un seam vide tandis que les unit tests injectent des tâches coopératives déterministes. Le supervisor récole les enfants terminés pendant l'exécution puis, au stop ou à la fermeture du channel de contrôle, attend la fin de tous les enfants avant de publier `Stopped`. Un `JoinError` observé marque la supervision comme non clean ; si la fermeture de supervision est ensuite engagée, le terminal est projeté vers `Faulted(worker_raw_transaction_ingest.runtime_invalid)`. Le traitement immédiat/prioritaire des faults reste différé à `pre.009`. Aucun timeout/abort forcé n'est encore introduit : ce hardening reste réservé à `pre.009`. Le gate communiqué le 8 septembre 2026 est vert sur `fmt`, audits, `check`, Clippy strict, 18 tests de crate, doc-tests et les deux arbres Cargo. ### `pre.006` — admission bornée + canonicalisation common Budget cible : **15–20 min**. `mpsc` borné, backpressure, ingress privé, common `RawTransactionMaterial -> RawTransaction`, observation-key domain/golden et assembly. Pas de source réseau. +État après matérialisation : **implémenté, gate opérateur requis**. Le supervisor possède un channel central `tokio::sync::mpsc` borné par `admission_queue_capacity`. Le seam source reçoit directement un clone du sender borné ainsi que le `watch` de stop ; le harness prouve `send().await` sous saturation et un `select! biased` stop-before-send, sans helper de production mort avant les vraies sources. L'ingress reste crate-private et transporte material source-neutral, network attendu, provenance sûre et une source key opaque de 32 bytes. Le Worker canonicalise exclusivement via `ksp-raw-transaction-lib`, vérifie le réseau avant/après canonicalisation, dérive l'observation key sous `ksp.raw_transaction_ingest.observation.v1` à partir de `(network, signature, source_key)`, puis assemble un `RawTransactionAcquisition` via le common RAW. Au stop, le receiver ferme les nouvelles admissions puis canonicalise les entrées déjà admises avant le join final. Aucune persistence Store, source live, Transport, snapshot concret ou policy provider n'est introduite. + ### `pre.007` — persistence Store + idempotence/conflict Budget cible : **15–20 min**. Port privé Store, mode Normal, concurrency bornée, outcomes new/idempotent/purged, observation distincte, Store error et content conflict terminal. diff --git a/docs/validation/028-V0_3_11_RAW_TRANSACTION_INGEST_WORKER_FOUNDATION.md b/docs/validation/028-V0_3_11_RAW_TRANSACTION_INGEST_WORKER_FOUNDATION.md index ca95f8c..3c723f6 100644 --- a/docs/validation/028-V0_3_11_RAW_TRANSACTION_INGEST_WORKER_FOUNDATION.md +++ b/docs/validation/028-V0_3_11_RAW_TRANSACTION_INGEST_WORKER_FOUNDATION.md @@ -1,5 +1,5 @@ - + # Validation v0.3.11 — fondation runtime du Worker RawTransaction ingest @@ -852,3 +852,165 @@ cargo tree -p ksp-worker-raw-transaction-ingest-lib -e features Critère de passage : le supervisor/JoinSet privé, le seam source déterministe et le join complet avant terminal sont verts ; aucun `mpsc`, common RAW assembly, Store write, snapshot concret, source live ou Transport n'apparaît. +## 17. Fermeture opérateur `pre.005` et matérialisation `pre.006` + +### 17.1 Fermeture opérateur de `pre.005` + +Le journal opérateur communiqué le 8 septembre 2026 ferme `0.3.11-pre.005` (`workspace.package.version = 0.3.11-pre.5`) : + +```text +cargo fmt --all : terminé sans erreur +python3 scripts/audit_rust_workspace_rules.py : clean, export completeness 0 +python3 scripts/audit_markdown_tables.py ... : clean (340 tables, 780 files) +cargo check --workspace : terminé sans erreur +cargo clippy --workspace --all-targets --all-features -- -D warnings : terminé sans erreur +cargo test -p ksp-worker-raw-transaction-ingest-lib : 11 unit + 3 dependency-boundary + 4 public API PASS, 0 échec +doc-tests : PASS +cargo tree normal : conforme au firewall attendu +cargo tree features : Tokio limité à macros/rt/sync/time côté Worker et Store sans backend Worker direct +``` + +`pre.006` peut donc être ouverte. + +### 17.2 Admission centrale bornée + +Le supervisor crée exactement un channel central : + +```text +tokio::sync::mpsc::channel(admission_queue_capacity) +``` + +`RawTransactionIngress` reste crate-private et n'appartient pas à l'API publique. Le receiver owner `RawTransactionAdmission` est également crate-private ; le seam source reçoit directement un clone du `tokio::sync::mpsc::Sender` avec son `watch` de stop. + +Le contrat d'admission est : + +```text +mpsc::Sender::send(...).await pour le backpressure +source harness : select! biased avec stop avant send +receiver close au shutdown +aucun unbounded_channel +aucun drop silencieux sous saturation normale +``` + +Le supervisor continue à observer le stop en priorité. Lors de la fermeture, le receiver appelle `close()`, refuse toute nouvelle admission, puis canonicalise les entrées déjà présentes dans le buffer avant le join final des tâches enfants. + +### 17.3 Ingress privé et common RAW + +L'ingress de fondation porte uniquement : + +```text +RawTransactionMaterial +RawNetworkId attendu +RawAcquisitionProvenance sûre +source key opaque [u8; 32] +``` + +Le pipeline `pre.006` est : + +```text +vérifier ingress.network == settings.network +canonicalize_raw_transaction(material) +vérifier transaction.reference.network == settings.network +dériver RawObservationKey Worker-owned +assemble_raw_transaction_acquisition(transaction, observation_key, provenance) +``` + +Aucune logique de canonicalisation RAW v1 n'est recodée dans le Worker. Une erreur common RAW est réduite à un contexte runtime statique et n'échoe aucun payload, réseau ou signature. + +### 17.4 Observation key privée V1 + +Le domaine exact est : + +```text +ksp.raw_transaction_ingest.observation.v1\0 +``` + +L'entrée SHA-256 est domain-separated puis encode avec longueur `u64` big-endian : + +```text +network UTF-8 +signature 64 bytes +source key 32 bytes +``` + +Aucun timestamp local, run id ou ordre d'admission n'entre dans la clé. Une rediffusion du même triplet produit donc la même `RawObservationKey`; le contrat provider/source public reste reporté à `0.3.12`. + +Golden harness figé : + +```text +network : mainnet +signature : 64 x 0x07 +source key : 32 x 0x09 +SHA-256 : f8d161dca8a8a72a0da110fc87882ee7136848461c6aa32b09e30a6188b95dcc +``` + +### 17.5 Preuves ajoutées + +Les tests unitaires `pre.006` couvrent : + +```text +capacity = 1 : second send reste bloqué tant que le premier slot n'est pas consommé +stop : un sender bloqué retourne false sans admission post-stop +common RAW : payload canonique exact + format/version exacts +observation key : golden exact +assembly : transaction/provenance/reference alignés +network guard : mismatch ingress et mismatch material rejetés sans fuite de valeurs +``` + +Le test de dependency boundary exige désormais `mpsc`, common canonicalization, common assembly et le domaine observation, tout en interdisant toujours persistence, mode Store, Transport, snapshots et `unbounded_channel`. + +### 17.6 Frontière volontaire de `pre.006` + +La tranche n'introduit encore aucun : + +```text +RawTransactionWrite::persist_raw_transaction_acquisition +RawTransactionAcquisitionMode::Normal +persistence concurrency +Store outcome mapping +content conflict terminal +store failure terminal +snapshot concret / compteurs publics +source live / Transport +Config +retry/reconnect/gap repair +``` + +`pre.007` reste propriétaire de la persistence et des outcomes. `pre.008` reste propriétaire des snapshots. `pre.009` reste propriétaire du fault ordering immédiat et du drain timeout/abort. + +### 17.7 Preuves locales d'assemblage `pre.006` + +Exécuté 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 (340 table(s), 781 file(s)) + +scan statique production : aucun ?, unwrap, expect, panic, unbounded_channel +scan statique scope : aucune persistence Store, aucun Transport, aucun snapshot concret +scan statique pipeline : exactement un mpsc::channel, un appel common canonicalize et un appel common assembly +golden SHA-256 indépendant : f8d161dca8a8a72a0da110fc87882ee7136848461c6aa32b09e30a6188b95dcc +``` + +L'environnement d'assemblage ne fournit ni `cargo`, ni `rustc`, ni `rustfmt`. Aucun gate Cargo `pre.006` n'est déclaré PASS localement. + +### 17.8 Gate opérateur demandé pour `pre.006` + +```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-worker-raw-transaction-ingest-lib +cargo tree -p ksp-worker-raw-transaction-ingest-lib --edges normal +cargo tree -p ksp-worker-raw-transaction-ingest-lib -e features +``` + +Critère de passage : admission bornée/backpressure, stop-preemption, canonicalisation common, golden observation et assembly sont verts ; aucune persistence Store, source live, Transport ou snapshot concret n'apparaît. +