From 8c4b458724b01a6fec02ca31fa77a9caa0f2927e Mon Sep 17 00:00:00 2001 From: SinuS Von SifriduS Date: Sat, 12 Sep 2026 14:20:02 +0200 Subject: [PATCH] v0.3.14-pre.010 --- Cargo.toml | 4 +- .../src/runtime_resources.rs | 217 +++++++++++++- .../tests/hardening.rs | 33 ++- .../tests/release_completeness.rs | 23 +- .../unit_tests/runtime_resources.rs | 81 +++++- deltas/0.3.14/pre.010.md | 273 ++++++++++++++++++ 6 files changed, 620 insertions(+), 11 deletions(-) create mode 100644 deltas/0.3.14/pre.010.md diff --git a/Cargo.toml b/Cargo.toml index 1203b10..d7f11e3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,12 @@ # file: Cargo.toml -# version: 581 +# version: 582 [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", "crates/ksp-worker-raw-transaction-ingest-lib"] [workspace.package] -version = "0.3.14-pre.9.fix.2" +version = "0.3.14-pre.10" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs b/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs index 7833f29..02a9ca1 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs -// version: 38 +// version: 39 use sha2::Digest; // rust-rules: trait-import @@ -18,6 +18,7 @@ pub const MIN_RAW_TRANSACTION_INGEST_HTTP_POLL_INTERVAL: std::time::Duration = s /// Minimum number of confirmed blocks accepted during one HTTP live polling cycle. pub const MIN_RAW_TRANSACTION_INGEST_HTTP_POLL_MAX_BLOCKS_PER_CYCLE: u16 = 1; +const MAX_RAW_TRANSACTION_INGEST_REPAIR_BURST: usize = 1; const RAW_TRANSACTION_INGEST_HELIUS_TRANSACTION_FILTER_FINGERPRINT_DOMAIN: &[u8] = b"ksp.raw_transaction_ingest.helius_transaction.filter.v1\0"; const RAW_TRANSACTION_INGEST_HELIUS_TRANSACTION_HTTP_PROTOCOL: &str = "helius_ws_http"; const RAW_TRANSACTION_INGEST_HELIUS_TRANSACTION_HTTP_SOURCE_KEY_DOMAIN: &[u8] = b"ksp.raw_transaction_ingest.helius_transaction_http.source_key.v1\0"; @@ -36,6 +37,197 @@ const RAW_TRANSACTION_INGEST_YELLOWSTONE_FILTER_FINGERPRINT_DOMAIN: &[u8] = b"ks const RAW_TRANSACTION_INGEST_YELLOWSTONE_HTTP_PROTOCOL: &str = "yellowstone_http"; const RAW_TRANSACTION_INGEST_YELLOWSTONE_HTTP_SOURCE_KEY_DOMAIN: &[u8] = b"ksp.raw_transaction_ingest.yellowstone_http.source_key.v1\0"; +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum RawTransactionIngestTrafficClass { + Nominal, + Repair, +} + +struct RawTransactionIngestFairTurnState { + active: bool, + last_granted: std::option::Option, + next_ticket: u64, + waiters: std::collections::VecDeque<(u64, RawTransactionIngestTrafficClass)>, +} + +struct RawTransactionIngestFairTurnGate { + max_waiters: usize, + notify: tokio::sync::Notify, + state: std::sync::Mutex, +} + +struct RawTransactionIngestFairWaiterGuard { + armed: bool, + gate: std::sync::Arc, + ticket: u64, +} + +struct RawTransactionIngestFairTurnGuard { + gate: std::sync::Arc, +} + +impl RawTransactionIngestFairTurnGate { + fn new(max_waiters: usize) -> Self { + return Self { + max_waiters, + notify: tokio::sync::Notify::new(), + state: std::sync::Mutex::new(RawTransactionIngestFairTurnState { + active: false, + last_granted: std::option::Option::None, + next_ticket: 1, + waiters: std::collections::VecDeque::new(), + }), + }; + } + + async fn acquire(self: std::sync::Arc, class: RawTransactionIngestTrafficClass) -> ksp_core_lib::Result { + let mut waiter = match RawTransactionIngestFairWaiterGuard::register(std::sync::Arc::clone(&self), class) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + loop { + let notify_gate = std::sync::Arc::clone(&self); + let notified = notify_gate.notify.notified(); + let granted = { + let mut state = match self.state.lock() { + std::result::Result::Ok(value) => value, + std::result::Result::Err(poisoned) => poisoned.into_inner(), + }; + if state.active || fair_turn_selected_ticket(&state) != std::option::Option::Some(waiter.ticket) { + false + } else { + let position = state.waiters.iter().position(|(ticket, _)| return *ticket == waiter.ticket); + let position = match position { + std::option::Option::Some(value) => value, + std::option::Option::None => return std::result::Result::Err(crate::runtime_error("source.fairness_waiter_missing")), + }; + let removed = state.waiters.remove(position); + let removed = match removed { + std::option::Option::Some(value) => value, + std::option::Option::None => return std::result::Result::Err(crate::runtime_error("source.fairness_waiter_missing")), + }; + if removed.1 != class { + return std::result::Result::Err(crate::runtime_error("source.fairness_waiter_class_mismatch")); + } + state.active = true; + state.last_granted = std::option::Option::Some(class); + true + } + }; + if granted { + waiter.armed = false; + std::mem::drop(notified); + std::mem::drop(notify_gate); + return std::result::Result::Ok(RawTransactionIngestFairTurnGuard { gate: self }); + } + notified.await; + } + } +} + +impl RawTransactionIngestFairWaiterGuard { + fn register(gate: std::sync::Arc, class: RawTransactionIngestTrafficClass) -> ksp_core_lib::Result { + let ticket = { + let mut state = match gate.state.lock() { + std::result::Result::Ok(value) => value, + std::result::Result::Err(poisoned) => poisoned.into_inner(), + }; + if gate.max_waiters == 0 || state.waiters.len() >= gate.max_waiters { + return std::result::Result::Err(crate::runtime_error("source.fairness_waiter_capacity_exceeded")); + } + let ticket = state.next_ticket; + state.next_ticket = match state.next_ticket.checked_add(1) { + std::option::Option::Some(value) => value, + std::option::Option::None => return std::result::Result::Err(crate::runtime_error("source.fairness_ticket_exhausted")), + }; + state.waiters.push_back((ticket, class)); + ticket + }; + gate.notify.notify_waiters(); + return std::result::Result::Ok(Self { armed: true, gate, ticket }); + } +} + +impl std::ops::Drop for RawTransactionIngestFairWaiterGuard { + fn drop(&mut self) { + if !self.armed { + return; + } + let mut state = match self.gate.state.lock() { + std::result::Result::Ok(value) => value, + std::result::Result::Err(poisoned) => poisoned.into_inner(), + }; + if let std::option::Option::Some(position) = state.waiters.iter().position(|(ticket, _)| return *ticket == self.ticket) { + let _removed = state.waiters.remove(position); + } + std::mem::drop(state); + self.gate.notify.notify_waiters(); + return; + } +} + +impl std::ops::Drop for RawTransactionIngestFairTurnGuard { + fn drop(&mut self) { + let mut state = match self.gate.state.lock() { + std::result::Result::Ok(value) => value, + std::result::Result::Err(poisoned) => poisoned.into_inner(), + }; + state.active = false; + std::mem::drop(state); + self.gate.notify.notify_waiters(); + return; + } +} + +fn fair_turn_selected_ticket(state: &RawTransactionIngestFairTurnState) -> std::option::Option { + let nominal = state.waiters.iter().find(|(_, class)| return *class == RawTransactionIngestTrafficClass::Nominal); + let repair = state.waiters.iter().find(|(_, class)| return *class == RawTransactionIngestTrafficClass::Repair); + let selected = match (nominal, repair) { + (std::option::Option::Some(nominal), std::option::Option::Some(repair)) => match state.last_granted { + std::option::Option::Some(RawTransactionIngestTrafficClass::Nominal) => repair, + std::option::Option::Some(RawTransactionIngestTrafficClass::Repair) | std::option::Option::None => nominal, + }, + (std::option::Option::Some(nominal), std::option::Option::None) => nominal, + (std::option::Option::None, std::option::Option::Some(repair)) => repair, + (std::option::Option::None, std::option::Option::None) => return std::option::Option::None, + }; + return std::option::Option::Some(selected.0); +} + +fn validate_repair_fairness_contract(admission_capacity: usize, persistence_concurrency: usize) -> ksp_core_lib::Result<()> { + if admission_capacity == 0 + || persistence_concurrency == 0 + || repair_block_fetch_limit(persistence_concurrency) == 0 + || MAX_RAW_TRANSACTION_INGEST_REPAIR_BURST != 1 + || !repair_fairness_catalog_is_complete() + { + return std::result::Result::Err(crate::runtime_error("runtime_resources.repair_fairness_invalid")); + } + return std::result::Result::Ok(()); +} + +fn repair_block_fetch_limit(existing_capacity: usize) -> usize { + return existing_capacity.min(4); +} + +fn repair_fairness_catalog_is_complete() -> bool { + let mut state = RawTransactionIngestFairTurnState { + active: false, + last_granted: std::option::Option::None, + next_ticket: 3, + waiters: std::collections::VecDeque::from([(1, RawTransactionIngestTrafficClass::Nominal), (2, RawTransactionIngestTrafficClass::Repair)]), + }; + if fair_turn_selected_ticket(&state) != std::option::Option::Some(1) { + return false; + } + state.last_granted = std::option::Option::Some(RawTransactionIngestTrafficClass::Nominal); + if fair_turn_selected_ticket(&state) != std::option::Option::Some(2) { + return false; + } + state.last_granted = std::option::Option::Some(RawTransactionIngestTrafficClass::Repair); + return fair_turn_selected_ticket(&state) == std::option::Option::Some(1); +} + #[derive(Clone, Copy, Debug, Eq, PartialEq)] enum RawTransactionIngestHttpDiscoveryStrategy { ClosedRange, @@ -2336,6 +2528,9 @@ impl crate::RawTransactionIngestRuntimeResources { { return std::result::Result::Err(error); } + if let std::result::Result::Err(error) = validate_repair_fairness_contract(settings.admission_queue_capacity(), settings.persistence_concurrency()) { + return std::result::Result::Err(error); + } let hydration_pending_budget = settings.admission_queue_capacity(); let hydration_in_flight_budget = settings.persistence_concurrency(); let inventory = std::sync::Arc::new(std::sync::Mutex::new(RawTransactionIngestSourceInventory::new(source_keys))); @@ -4263,6 +4458,7 @@ enum RawTransactionIngestSharedHydrationResult { } struct RawTransactionIngestGlobalHydrationRegistry { + hydration_fairness: std::sync::Arc, hydration_permits: std::sync::Arc, max_pending: usize, pending: std::sync::Mutex< @@ -4276,12 +4472,26 @@ struct RawTransactionIngestGlobalHydrationRegistry { impl RawTransactionIngestGlobalHydrationRegistry { fn new(max_pending: usize, max_in_flight: usize) -> Self { return Self { + hydration_fairness: std::sync::Arc::new(RawTransactionIngestFairTurnGate::new(max_pending)), hydration_permits: std::sync::Arc::new(tokio::sync::Semaphore::new(max_in_flight)), max_pending, pending: std::sync::Mutex::new(std::collections::BTreeMap::new()), }; } + async fn acquire_hydration_permit(&self, class: RawTransactionIngestTrafficClass) -> ksp_core_lib::Result { + let turn = match std::sync::Arc::clone(&self.hydration_fairness).acquire(class).await { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let permit = std::sync::Arc::clone(&self.hydration_permits).acquire_owned().await; + std::mem::drop(turn); + return match permit { + std::result::Result::Ok(value) => std::result::Result::Ok(value), + std::result::Result::Err(_) => std::result::Result::Err(crate::runtime_error("source.global_hydration_permit_closed")), + }; + } + fn subscribe_or_lead( &self, key: &RawTransactionIngestHydrationKey, @@ -4599,10 +4809,9 @@ async fn fetch_hydration_shared( }; if leader { let mut leader_guard = RawTransactionIngestHydrationLeaderGuard::new(std::sync::Arc::clone(&global_registry), key.clone()); - let permit = std::sync::Arc::clone(&global_registry.hydration_permits).acquire_owned().await; - let _permit = match permit { + let _permit = match global_registry.acquire_hydration_permit(RawTransactionIngestTrafficClass::Nominal).await { std::result::Result::Ok(value) => value, - std::result::Result::Err(_) => return std::result::Result::Err(crate::runtime_error("source.global_hydration_permit_closed")), + std::result::Result::Err(error) => return std::result::Result::Err(error), }; let fetched = fetch_hydration(http_pool, hydration_role, expected_network, key.clone(), commitment).await; return match fetched { diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs b/crates/ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs index 0ffd209..a74e0c6 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs @@ -1,7 +1,7 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs -// version: 33 +// version: 34 -//! External public, security, redaction and release-boundary hardening canaries through `v0.3.14-pre.009`. +//! External public, security, redaction and release-boundary hardening canaries through `v0.3.14-pre.010`. fn network(value: &'static str) -> std::option::Option { let result = ksp_store_lib::RawNetworkId::new(value); @@ -1234,3 +1234,32 @@ fn v0_3_14_pre_009_health_policy_is_present_future_coverage_gated_and_source_neu assert!(!snapshot.contains("WorkerHealth::Faulted"), "Faulted must remain a Worker lifecycle state rather than a new health enum variant"); return; } + +#[test] +fn v0_3_14_pre_010_repair_fairness_shares_existing_bounds_without_second_pipeline() { + let continuity = include_str!("../src/continuity.rs"); + let resources = include_str!("../src/runtime_resources.rs"); + let admission = include_str!("../src/admission.rs"); + let root = include_str!("../src/lib.rs"); + for required in [ + "MAX_RAW_TRANSACTION_INGEST_REPAIR_BURST: usize = 1", + "RawTransactionIngestTrafficClass", + "RawTransactionIngestFairTurnGate", + "RawTransactionIngestTrafficClass::Nominal", + "RawTransactionIngestTrafficClass::Repair", + "acquire_hydration_permit", + "validate_repair_fairness_contract", + "repair_block_fetch_limit", + ] { + assert!(resources.contains(required) || continuity.contains(required), "required pre.010 fairness guard missing: {required}"); + } + assert!(continuity.contains("MAX_RAW_TRANSACTION_INGEST_REPAIR_BLOCK_FETCH_IN_FLIGHT: usize = 4")); + assert!(continuity.contains("MAX_RAW_TRANSACTION_INGEST_REPAIR_DISCOVERY_WINDOW_SLOTS: u64 = 512")); + assert!(admission.contains("tokio::sync::mpsc::channel(capacity)")); + assert!(!resources.contains("tokio::sync::mpsc::channel("), "pre.010 created a second admission pipeline"); + assert_eq!(resources.matches("struct RawTransactionIngestGlobalHydrationRegistry").count(), 1); + assert!(!resources.contains("RepairHydrationRegistry")); + assert!(!root.contains("RawTransactionIngestTrafficClass")); + assert!(!root.contains("RawTransactionIngestFairTurnGate")); + return; +} diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs b/crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs index 0871acd..b935af1 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs @@ -1,7 +1,7 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs -// version: 28 +// version: 29 -//! Release-completeness canaries through the `v0.3.14-pre.009` multi-source health-policy tranche. +//! Release-completeness canaries through the `v0.3.14-pre.010` repair backpressure/fairness tranche. #[test] fn pre_010_production_module_inventory_is_exact() -> std::io::Result<()> { @@ -367,3 +367,22 @@ fn v0_3_14_pre_009_health_policy_canaries_are_present_without_public_source_iden assert!(!root.contains("pub use self::continuity::RawTransactionIngestTargetCoverage")); return; } + +#[test] +fn v0_3_14_pre_010_repair_fairness_canaries_are_present_without_public_or_pipeline_growth() { + let hardening = include_str!("hardening.rs"); + let resource_tests = include_str!("../unit_tests/runtime_resources.rs"); + let root = include_str!("../src/lib.rs"); + for required in [ + "v0_3_14_pre_010_fair_scheduler_alternates_nominal_and_repair_with_bounded_burst", + "v0_3_14_pre_010_fair_waiter_bound_is_checked_and_released_without_reserved_capacity", + "v0_3_14_pre_010_nominal_and_repair_hydration_share_one_capacity_one_permit_pool", + "v0_3_14_pre_010_repair_only_progresses_when_existing_capacity_is_one", + ] { + assert!(resource_tests.contains(required), "required pre.010 fairness canary missing: {required}"); + } + assert!(hardening.contains("v0_3_14_pre_010_repair_fairness_shares_existing_bounds_without_second_pipeline")); + assert!(!root.contains("RawTransactionIngestTrafficClass")); + assert!(!root.contains("RawTransactionIngestFairTurnGate")); + return; +} diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs index 33617a8..dc84e32 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs -// version: 31 +// version: 32 fn grpc_endpoint(cluster: &str) -> std::option::Option { return grpc_endpoint_with_identity(cluster, "yellowstone-fixture", "fixture-provider"); @@ -4142,3 +4142,82 @@ fn v0_3_14_pre_009_inventory_health_projection_tracks_reconciled_coverage() { assert!(reconciled.future_target_coverage()); return; } + +#[test] +fn v0_3_14_pre_010_fair_scheduler_alternates_nominal_and_repair_with_bounded_burst() { + assert_eq!(super::MAX_RAW_TRANSACTION_INGEST_REPAIR_BURST, 1); + let mut state = super::RawTransactionIngestFairTurnState { + active: false, + last_granted: std::option::Option::None, + next_ticket: 3, + waiters: std::collections::VecDeque::from([ + (1, super::RawTransactionIngestTrafficClass::Nominal), + (2, super::RawTransactionIngestTrafficClass::Repair), + ]), + }; + assert_eq!(super::fair_turn_selected_ticket(&state), std::option::Option::Some(1)); + state.last_granted = std::option::Option::Some(super::RawTransactionIngestTrafficClass::Nominal); + assert_eq!(super::fair_turn_selected_ticket(&state), std::option::Option::Some(2)); + state.last_granted = std::option::Option::Some(super::RawTransactionIngestTrafficClass::Repair); + assert_eq!(super::fair_turn_selected_ticket(&state), std::option::Option::Some(1)); + let _removed = state.waiters.pop_front(); + assert_eq!(super::fair_turn_selected_ticket(&state), std::option::Option::Some(2)); + return; +} + +#[test] +fn v0_3_14_pre_010_fair_waiter_bound_is_checked_and_released_without_reserved_capacity() { + assert!(super::validate_repair_fairness_contract(1, 1).is_ok()); + assert_eq!(super::repair_block_fetch_limit(1), 1); + assert_eq!(super::repair_block_fetch_limit(4), 4); + assert_eq!(super::repair_block_fetch_limit(64), 4); + let gate = std::sync::Arc::new(super::RawTransactionIngestFairTurnGate::new(1)); + let first = super::RawTransactionIngestFairWaiterGuard::register(std::sync::Arc::clone(&gate), super::RawTransactionIngestTrafficClass::Repair); + let first = match first { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + assert!(super::RawTransactionIngestFairWaiterGuard::register(std::sync::Arc::clone(&gate), super::RawTransactionIngestTrafficClass::Nominal,).is_err()); + std::mem::drop(first); + assert!(super::RawTransactionIngestFairWaiterGuard::register(gate, super::RawTransactionIngestTrafficClass::Nominal).is_ok()); + return; +} + +#[tokio::test] +async fn v0_3_14_pre_010_nominal_and_repair_hydration_share_one_capacity_one_permit_pool() { + let registry = std::sync::Arc::new(super::RawTransactionIngestGlobalHydrationRegistry::new(2, 1)); + let nominal = registry.acquire_hydration_permit(super::RawTransactionIngestTrafficClass::Nominal).await; + let nominal = match nominal { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + let repair_registry = std::sync::Arc::clone(®istry); + let repair = tokio::spawn(async move { + return repair_registry.acquire_hydration_permit(super::RawTransactionIngestTrafficClass::Repair).await; + }); + tokio::task::yield_now().await; + assert!(!repair.is_finished()); + std::mem::drop(nominal); + let repair = match repair.await { + std::result::Result::Ok(std::result::Result::Ok(value)) => value, + std::result::Result::Ok(std::result::Result::Err(_)) | std::result::Result::Err(_) => return, + }; + assert_eq!(registry.hydration_permits.available_permits(), 0); + std::mem::drop(repair); + assert_eq!(registry.hydration_permits.available_permits(), 1); + return; +} + +#[tokio::test] +async fn v0_3_14_pre_010_repair_only_progresses_when_existing_capacity_is_one() { + let registry = super::RawTransactionIngestGlobalHydrationRegistry::new(1, 1); + let permit = registry.acquire_hydration_permit(super::RawTransactionIngestTrafficClass::Repair).await; + let permit = match permit { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + assert_eq!(registry.hydration_permits.available_permits(), 0); + std::mem::drop(permit); + assert_eq!(registry.hydration_permits.available_permits(), 1); + return; +} diff --git a/deltas/0.3.14/pre.010.md b/deltas/0.3.14/pre.010.md new file mode 100644 index 0000000..4bfae50 --- /dev/null +++ b/deltas/0.3.14/pre.010.md @@ -0,0 +1,273 @@ + + + +# Delta `0.3.14-pre.010` — backpressure et fairness nominal / repair + +## Base requise + +```text +0.3.14-pre.009-fix.002 +workspace.package.version = 0.3.14-pre.9.fix.2 +deltas/0.3.14/pre.009-fix.002.md présent +``` + +## Gate de la base + +Le gate opérateur de `0.3.14-pre.009-fix.002` est validé avant ouverture de cette tranche : + +```text +cargo fmt --all : PASS +cargo fmt --all -- --check : PASS +audit Rust workspace rules : PASS +audit Markdown tables : PASS +cargo check --workspace : PASS +cargo clippy --workspace --all-targets --all-features -- -D warnings : PASS +cargo test -p ksp-worker-raw-transaction-ingest-lib --all-targets --all-features : PASS +``` + +Le gate Worker comprend notamment : + +```text +143 unit tests : PASS +cross_layer_completeness : 8 PASS +dependency_boundary : 19 PASS +hardening : 34 PASS +public_api : 20 PASS +release_completeness : 11 PASS +``` + +## Objectif + +Implémenter strictement la tranche `pre.010` du plan `035` : + +```text +conserver une admission Common RAW centrale unique +conserver un pipeline de persistence unique et ses limites existantes +faire partager au trafic nominal et repair le registre/coalescence global d'hydration existant +ordonner nominal/repair sans réserver une fraction fixe de capacité +permettre la progression lorsque la capacité existante vaut 1 +borner le burst repair lorsqu'un nominal attend +conserver un fanout logique getBlock <= 4 +conserver les fenêtres HTTP <= 512 slots +ne créer aucun second pool Store, aucune seconde queue d'admission et aucun second registre d'hydration +``` + +## Arbitre privé nominal / repair + +`runtime_resources.rs` introduit un arbitre privé source-neutral : + +```text +RawTransactionIngestTrafficClass + Nominal + Repair + +RawTransactionIngestFairTurnGate +``` + +Les waiters reçoivent des tickets monotones checked. La queue d'attente est bornée par la capacité déjà attribuée au registre global d'hydration ; aucun waiter pool illimité n'est ajouté. + +Lorsque les deux classes attendent : + +```text +premier arbitrage -> Nominal +après Nominal -> Repair +après Repair -> Nominal +``` + +Le burst repair maximal en présence de nominal est donc : + +```text +MAX_RAW_TRANSACTION_INGEST_REPAIR_BURST = 1 +``` + +Si une seule classe attend, elle progresse sans attendre artificiellement l'autre classe. Le mécanisme ne réserve donc jamais une capacité inexistante. + +## Hydration globale partagée + +Le `RawTransactionIngestGlobalHydrationRegistry` reste unique. + +Le même sémaphore : + +```text +hydration_permits +``` + +est maintenant précédé par l'arbitre nominal/repair. Le chemin live existant acquiert explicitement : + +```text +RawTransactionIngestTrafficClass::Nominal +``` + +La classe `Repair` utilise la même primitive d'acquisition et le même sémaphore ; aucun `RepairHydrationRegistry`, aucune seconde semaphore d'hydration et aucune seconde coalescence ne sont créés. + +L'arbitre gouverne l'ordre d'accès au permit ; le permit lui-même reste détenu pendant l'I/O comme auparavant. Les bornes globales existantes restent donc l'autorité de concurrence. + +## Faible capacité et fanout HTTP + +La validation run-local accepte explicitement : + +```text +admission_queue_capacity = 1 +persistence_concurrency = 1 +``` + +Le fanout bloc repair dérive uniquement de la capacité existante : + +```text +repair_block_fetch_limit(existing_capacity) = min(existing_capacity, 4) +``` + +Donc : + +```text +capacité 1 -> fanout 1 +capacité 4 -> fanout 4 +capacité 64 -> fanout 4 +``` + +Aucun worker ou permit n'est réservé à l'avance pour le repair. La borne historique reste également : + +```text +MAX_RAW_TRANSACTION_INGEST_REPAIR_ACTIVE_GAPS = 1 +MAX_RAW_TRANSACTION_INGEST_REPAIR_BLOCK_FETCH_IN_FLIGHT = 4 +MAX_RAW_TRANSACTION_INGEST_REPAIR_DISCOVERY_WINDOW_SLOTS = 512 +``` + +## Admission et persistence + +Cette tranche ne crée aucune nouvelle `mpsc::channel` dans `runtime_resources.rs`. + +Tout matériau Common RAW continue à utiliser : + +```text +RawTransactionAdmission +la Sender centrale existante +le pipeline de persistence existant +la même persistence_concurrency +``` + +La fairness ajoutée ne contourne donc ni backpressure admission ni convergence/persistence Store. + +## Compteurs et overflow + +Les tickets de fairness utilisent `checked_add`. Un épuisement de ticket ou un dépassement de la queue de waiters produit une erreur Worker stable ; aucun `saturating_add` n'est utilisé pour masquer un overflow. + +L'annulation d'un waiter retire son ticket de l'arbitre. La libération d'un turn réveille les waiters restants sans conserver de `std::sync::MutexGuard` à travers un `await`. + +## Tests ajoutés + +Les unit tests Worker couvrent : + +```text +alternance Nominal -> Repair -> Nominal lorsque les deux classes attendent +burst repair borné à 1 +queue de waiters bornée et libérée après annulation/drop +fanout bloc min(capacité, 4) +partage réel du même semaphore d'hydration entre Nominal et Repair +blocage Repair tant que le permit capacity=1 est occupé par Nominal +progression Repair seul avec capacity=1 +``` + +Les canaris `hardening` et `release_completeness` vérifient en plus : + +```text +une seule admission mpsc centrale +un seul RawTransactionIngestGlobalHydrationRegistry +absence de RepairHydrationRegistry +bornes 4 getBlock / 512 slots conservées +aucune fuite publique des classes ou de l'arbitre de fairness +``` + +## Hors périmètre inchangé + +```text +aucune nouvelle snapshot publique de gap/repair avant pre.011 +aucun changement de health policy pre.009 +aucun nouveau mécanisme replay/reconnect Worker +aucun respawn de source +aucun EARLY/shred +aucun Job Backfill depuis Worker +aucun nouveau provider ou SDK +aucune nouvelle dépendance +``` + +## Fichiers ajoutés + +```text +deltas/0.3.14/pre.010.md +``` + +## Fichiers modifiés + +```text +Cargo.toml +crates/ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs +crates/ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs +crates/ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs +crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs +``` + +## Fichiers supprimés + +```text +aucun +``` + +## Version Cargo + +Conformément à `VER-ID-009` : + +```text +header Cargo.toml : 581 -> 582 +workspace.package.version : 0.3.14-pre.9.fix.2 -> 0.3.14-pre.10 +``` + +Versions des fichiers modifiés : + +```text +runtime_resources.rs : 39 (bump depuis 38) +unit_tests/runtime_resources.rs : 32 (bump depuis 31) +tests/hardening.rs : 34 (bump depuis 33) +tests/release_completeness.rs : 29 (bump depuis 28) +``` + +## Validation exécutée dans l'environnement de préparation + +```text +python3 scripts/audit_rust_workspace_rules.py : PASS +python3 scripts/audit_markdown_tables.py README.md RULES.md ROADMAP.md CHANGELOG.md docs prompts crates deltas : PASS +scan des frontières Worker/Transport/Backfill/Store : PASS +scan de la crate-root historique : PASS +comparaison exacte pre.009-fix.002 -> pre.010 : PASS +``` + +Les gates Cargo ne sont pas déclarés PASS dans l'environnement de préparation lorsqu'ils ne peuvent pas y être exécutés. Ils restent obligatoires côté opérateur avant `pre.011`. + +## Décisions prises + +```text +fairness temporelle/alternée et non réservation de pourcentage +burst repair = 1 lorsque nominal attend +classe seule autorisée à progresser immédiatement +hydration nominal et repair sur le même registre + semaphore +fanout bloc dérivé de la capacité existante et plafonné à 4 +aucune nouvelle queue admission/persistence +``` + +## Questions ouvertes + +```text +aucune pour cette tranche +``` + +## Gate opérateur après application + +```bash +cargo fmt --all +cargo fmt --all -- --check +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 --all-targets --all-features +```