diff --git a/Cargo.toml b/Cargo.toml index 28a383a..5011b3c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,12 @@ # file: Cargo.toml -# version: 570 +# version: 571 [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.3" +version = "0.3.14-pre.4" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-onchain-transport-lib/USAGE.md b/crates/ksp-onchain-transport-lib/USAGE.md index d679997..3ade543 100644 --- a/crates/ksp-onchain-transport-lib/USAGE.md +++ b/crates/ksp-onchain-transport-lib/USAGE.md @@ -1,5 +1,5 @@ - + # Utilisation de `ksp-onchain-transport-lib` @@ -387,9 +387,11 @@ let snapshot = stream.snapshot(); let closed = stream.close().await; ``` -`try_update()` remplace dynamiquement la requête complète tant que la session est `Active`. Une mutation pendant `Reconnecting` est refusée pour éviter une application ambiguë. Le snapshot expose reconnects, replay attempts, gaps, duplicates, dernier `from_slot` demandé et dernier slot observé, sans endpoint ni payload arbitraire. +`try_update()` remplace dynamiquement la requête complète tant que la session est `Active`. Une mutation pendant `Reconnecting` est refusée pour éviter une application ambiguë. Le snapshot expose reconnects, replay attempts, replay deliveries conservatrices, reprises dont la couverture reste non prouvée, gaps de rétention, duplicates, dernier `from_slot` demandé et dernier slot observé, sans endpoint ni payload arbitraire. -Le reconnect réutilise la dernière requête acceptée et peut avancer `from_slot`, mais le consumer doit traiter cette reprise comme best-effort. KSP ne promet ni exactly-once, ni replay historique complet, ni absence de fork/equivocation entre nœuds. +Une `replay_delivery_count` n'augmente que lorsque le premier flux post-reconnect redélivre exactement le slot de reprise demandé. Elle prouve une livraison au bord de replay, pas la complétude de l'intervalle. Dès que le premier update slot-bearing atteint ou dépasse cette borne, `replay_coverage_unproven_count` augmente aussi : Transport ne possède pas de preuve générique que tous les updates correspondant aux filtres ont été livrés entre la borne et la reprise live. Le compteur augmente également si un nouveau reconnect survient avant tout matériau slot-bearing. Cette valeur signifie explicitement « couverture non prouvée » ; elle ne prouve pas qu'un événement filtré existait ou a été perdu. + +Le reconnect réutilise la dernière requête acceptée et peut avancer `from_slot`, mais le consumer doit traiter cette reprise comme best-effort. `replay_attempt_count`, `replay_delivery_count` et une preuve de coverage sont des notions distinctes. Transport ne publie actuellement aucun compteur `replay_covered` générique, car ni l'acceptation de `from_slot` ni une livraison ponctuelle ne démontrent à elles seules une couverture historique complète. KSP ne promet ni exactly-once, ni replay historique complet, ni absence de fork/equivocation entre nœuds. ## 5. Appels typés diff --git a/crates/ksp-onchain-transport-lib/src/grpc_stream.rs b/crates/ksp-onchain-transport-lib/src/grpc_stream.rs index 683eab3..b2b5afd 100644 --- a/crates/ksp-onchain-transport-lib/src/grpc_stream.rs +++ b/crates/ksp-onchain-transport-lib/src/grpc_stream.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/src/grpc_stream.rs -// version: 4 +// version: 5 use tonic_prost::prost::Message; // rust-rules: trait-import @@ -34,6 +34,8 @@ pub struct YellowstoneGrpcSubscribeSnapshot { continuity_gap_count: u64, duplicate_update_count: u64, replay_attempt_count: u64, + replay_delivery_count: u64, + replay_coverage_unproven_count: u64, last_requested_from_slot: std::option::Option, last_observed_slot: std::option::Option, terminal_error_code: std::option::Option, @@ -84,6 +86,8 @@ impl YellowstoneGrpcSubscribeSnapshot { continuity_gap_count: 0, duplicate_update_count: 0, replay_attempt_count: 0, + replay_delivery_count: 0, + replay_coverage_unproven_count: 0, last_requested_from_slot: initial_from_slot, last_observed_slot: std::option::Option::None, terminal_error_code: std::option::Option::None, @@ -122,6 +126,25 @@ impl YellowstoneGrpcSubscribeSnapshot { return self.replay_attempt_count; } + /// Returns the number of replay-bearing reconnects that delivered the requested replay boundary slot again. + /// + /// This is conservative delivery evidence only. It does not prove that every matching update in the replay interval was delivered and must never be + /// interpreted as `replay_covered`. + #[must_use] + pub const fn replay_delivery_count(self) -> u64 { + return self.replay_delivery_count; + } + + /// Returns the number of successful replay-bearing reconnects whose target coverage remains unproven. + /// + /// The counter advances when the first post-reconnect slot-bearing update reaches or passes the requested replay boundary, because generic Transport + /// cannot prove from that delivery alone that every matching update in the replay interval was delivered. It also advances when another reconnect starts + /// before any slot-bearing replay material arrives. This is an explicit lack of coverage proof, not proof that a filtered event actually existed or was lost. + #[must_use] + pub const fn replay_coverage_unproven_count(self) -> u64 { + return self.replay_coverage_unproven_count; + } + /// Returns the most recent effective `from_slot` sent by KSP, including any clamp to `SubscribeReplayInfo.first_available`. #[must_use] pub const fn last_requested_from_slot(self) -> std::option::Option { @@ -506,17 +529,44 @@ enum UpdateIdentity { } struct ContinuityTracker { + pending_replay_from_slot: std::option::Option, recent_order: std::collections::VecDeque, recent_set: std::collections::HashSet, } impl ContinuityTracker { fn new() -> Self { - return Self { recent_order: std::collections::VecDeque::new(), recent_set: std::collections::HashSet::new() }; + return Self { + pending_replay_from_slot: std::option::Option::None, + recent_order: std::collections::VecDeque::new(), + recent_set: std::collections::HashSet::new(), + }; + } + + fn begin_replay(&mut self, from_slot: std::option::Option) { + self.pending_replay_from_slot = from_slot; + return; + } + + fn abandon_pending_replay(&mut self, snapshot: &mut crate::YellowstoneGrpcSubscribeSnapshot) -> bool { + if self.pending_replay_from_slot.take().is_none() { + return false; + } + snapshot.replay_coverage_unproven_count = snapshot.replay_coverage_unproven_count.saturating_add(1); + return true; } fn observe(&mut self, update: &crate::YellowstoneSubscribeUpdate, snapshot: &mut crate::YellowstoneGrpcSubscribeSnapshot) { if let std::option::Option::Some(slot) = update_slot(update) { + if let std::option::Option::Some(requested) = self.pending_replay_from_slot + && slot >= requested + { + if slot == requested { + snapshot.replay_delivery_count = snapshot.replay_delivery_count.saturating_add(1); + } + snapshot.replay_coverage_unproven_count = snapshot.replay_coverage_unproven_count.saturating_add(1); + self.pending_replay_from_slot = std::option::Option::None; + } snapshot.last_observed_slot = std::option::Option::Some(match snapshot.last_observed_slot { std::option::Option::Some(previous) => std::cmp::max(previous, slot), std::option::Option::None => slot, @@ -620,6 +670,7 @@ async fn run_subscribe_actor( &mut shutdown_rx, &mut snapshot, &snapshot_tx, + &mut tracker, ) .await { @@ -662,6 +713,7 @@ async fn run_subscribe_actor( &mut shutdown_rx, &mut snapshot, &snapshot_tx, + &mut tracker, ) .await { @@ -701,8 +753,12 @@ async fn reconnect_subscribe_stream( shutdown_rx: &mut tokio::sync::watch::Receiver>, snapshot: &mut crate::YellowstoneGrpcSubscribeSnapshot, snapshot_tx: &tokio::sync::watch::Sender, + tracker: &mut ContinuityTracker, ) -> ReconnectOutcome { clear_request_sender(request_state); + if tracker.abandon_pending_replay(snapshot) { + snapshot_tx.send_replace(*snapshot); + } snapshot.state = crate::YellowstoneGrpcSubscribeState::Reconnecting; snapshot.terminal_error_code = std::option::Option::None; snapshot_tx.send_replace(*snapshot); @@ -785,6 +841,7 @@ async fn reconnect_subscribe_stream( return ReconnectOutcome::Exhausted(error); } snapshot.reconnect_count = snapshot.reconnect_count.saturating_add(1); + tracker.begin_replay(effective_from_slot); snapshot.state = crate::YellowstoneGrpcSubscribeState::Active; snapshot.terminal_error_code = std::option::Option::None; snapshot_tx.send_replace(*snapshot); diff --git a/crates/ksp-onchain-transport-lib/tests/public_api.rs b/crates/ksp-onchain-transport-lib/tests/public_api.rs index 687766b..41c91b3 100644 --- a/crates/ksp-onchain-transport-lib/tests/public_api.rs +++ b/crates/ksp-onchain-transport-lib/tests/public_api.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/tests/public_api.rs -// version: 53 +// version: 54 //! Integration tests for the public `ksp-onchain-transport-lib` consumer contract. @@ -1032,6 +1032,8 @@ fn public_v0_2_9_pre_010_yellowstone_reconnect_snapshot_is_available_from_crate_ let _gap_count = ksp_onchain_transport_lib::YellowstoneGrpcSubscribeSnapshot::continuity_gap_count; let _duplicate_count = ksp_onchain_transport_lib::YellowstoneGrpcSubscribeSnapshot::duplicate_update_count; let _replay_count = ksp_onchain_transport_lib::YellowstoneGrpcSubscribeSnapshot::replay_attempt_count; + let _delivery_count = ksp_onchain_transport_lib::YellowstoneGrpcSubscribeSnapshot::replay_delivery_count; + let _coverage_unproven_count = ksp_onchain_transport_lib::YellowstoneGrpcSubscribeSnapshot::replay_coverage_unproven_count; let _requested = ksp_onchain_transport_lib::YellowstoneGrpcSubscribeSnapshot::last_requested_from_slot; let _observed = ksp_onchain_transport_lib::YellowstoneGrpcSubscribeSnapshot::last_observed_slot; let _terminal = ksp_onchain_transport_lib::YellowstoneGrpcSubscribeSnapshot::terminal_error_code; diff --git a/crates/ksp-onchain-transport-lib/tests/release_completeness.rs b/crates/ksp-onchain-transport-lib/tests/release_completeness.rs index 27c6539..4cdefc4 100644 --- a/crates/ksp-onchain-transport-lib/tests/release_completeness.rs +++ b/crates/ksp-onchain-transport-lib/tests/release_completeness.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/tests/release_completeness.rs -// version: 45 +// version: 46 //! Release-level completeness canaries for staged HTTP and WebSocket Transport coverage. @@ -1374,3 +1374,28 @@ fn release_v0_3_14_pre_003_ws_snapshot_source_reuses_actor_watch_without_second_ assert!(!facade_source.contains("tokio_tungstenite::connect_async")); return; } + +#[test] +fn release_v0_3_14_pre_004_native_replay_evidence_never_conflates_attempt_delivery_and_coverage() { + let stream_source = include_str!("../src/grpc_stream.rs"); + let stream_tests = include_str!("../unit_tests/grpc_stream.rs"); + for required in [ + "replay_attempt_count", + "replay_delivery_count", + "replay_coverage_unproven_count", + "pending_replay_from_slot", + "begin_replay", + "abandon_pending_replay", + "first_available > requested", + "This is conservative delivery evidence only", + "explicit lack of coverage proof", + ] { + assert!(stream_source.contains(required), "missing pre.004 native replay evidence token: {required}"); + } + assert!(stream_tests.contains("v0_3_14_pre_004_replay_acceptance_without_boundary_delivery_never_claims_coverage")); + assert!(stream_tests.contains("yellowstone_reconnect_replays_from_last_observed_slot_and_counts_duplicate_identity")); + assert!(stream_tests.contains("yellowstone_replay_info_proves_and_clamps_retention_gap_without_lossless_claim")); + assert!(!stream_source.contains("replay_covered_count")); + assert!(!stream_source.contains("replay_coverage_proven_count")); + return; +} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/grpc_stream.rs b/crates/ksp-onchain-transport-lib/unit_tests/grpc_stream.rs index 03ce308..a657dc6 100644 --- a/crates/ksp-onchain-transport-lib/unit_tests/grpc_stream.rs +++ b/crates/ksp-onchain-transport-lib/unit_tests/grpc_stream.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/unit_tests/grpc_stream.rs -// version: 5 +// version: 6 #[derive(Clone, Copy)] enum FixtureMode { @@ -13,6 +13,7 @@ enum FixtureMode { Idle, ReconnectReplay, ReplayGap, + ReplayAcceptedWithoutBoundary, ReconnectExhausted, } @@ -153,6 +154,22 @@ impl yellowstone_grpc_proto::geyser::geyser_server::Geyser for FixtureGeyser { } } }, + FixtureMode::ReplayAcceptedWithoutBoundary => { + if subscribe_call == 1 { + assert_eq!(initial.from_slot, std::option::Option::None); + let _ = outbound_tx.send(std::result::Result::Ok(slot_update(800))).await; + } else { + assert_eq!(subscribe_call, 2); + assert_eq!(initial.from_slot, std::option::Option::Some(800)); + if outbound_tx.send(std::result::Result::Ok(slot_update(805))).await.is_err() { + return; + } + let half_close = inbound.message().await; + if matches!(half_close, std::result::Result::Ok(std::option::Option::None)) { + half_close_seen.store(true, std::sync::atomic::Ordering::SeqCst); + } + } + }, FixtureMode::ReconnectExhausted => { assert_eq!(subscribe_call, 1); let _ = outbound_tx.send(std::result::Result::Ok(slot_update(700))).await; @@ -191,7 +208,7 @@ impl yellowstone_grpc_proto::geyser::geyser_server::Geyser for FixtureGeyser { return std::result::Result::Err(error); } let first_available = match self.mode { - FixtureMode::ReconnectReplay | FixtureMode::ReconnectExhausted => std::option::Option::Some(400), + FixtureMode::ReconnectReplay | FixtureMode::ReconnectExhausted | FixtureMode::ReplayAcceptedWithoutBoundary => std::option::Option::Some(400), FixtureMode::ReplayGap => std::option::Option::Some(505), _ => return std::result::Result::Err(tonic::Status::unimplemented("replay info is outside this fixture mode")), }; @@ -595,6 +612,8 @@ async fn yellowstone_reconnect_replays_from_last_observed_slot_and_counts_duplic assert_eq!(snapshot.state(), crate::YellowstoneGrpcSubscribeState::Active); assert_eq!(snapshot.reconnect_count(), 1); assert_eq!(snapshot.replay_attempt_count(), 1); + assert_eq!(snapshot.replay_delivery_count(), 1); + assert_eq!(snapshot.replay_coverage_unproven_count(), 1); assert_eq!(snapshot.continuity_gap_count(), 0); assert_eq!(snapshot.duplicate_update_count(), 1); assert_eq!(snapshot.last_requested_from_slot(), std::option::Option::Some(500)); @@ -627,6 +646,8 @@ async fn yellowstone_replay_info_proves_and_clamps_retention_gap_without_lossles let snapshot = session.snapshot(); assert_eq!(snapshot.reconnect_count(), 1); assert_eq!(snapshot.replay_attempt_count(), 1); + assert_eq!(snapshot.replay_delivery_count(), 1); + assert_eq!(snapshot.replay_coverage_unproven_count(), 1); assert_eq!(snapshot.continuity_gap_count(), 1); assert_eq!(snapshot.duplicate_update_count(), 0); assert_eq!(snapshot.last_requested_from_slot(), std::option::Option::Some(505)); @@ -635,6 +656,42 @@ async fn yellowstone_replay_info_proves_and_clamps_retention_gap_without_lossles server.stop().await; } +#[tokio::test(flavor = "current_thread")] +async fn v0_3_14_pre_004_replay_acceptance_without_boundary_delivery_never_claims_coverage() { + let server = FixtureServer::start(FixtureMode::ReplayAcceptedWithoutBoundary).await; + let defaults = crate::YellowstoneGrpcSessionSettings::default(); + let settings = fixture_settings_with_reconnect( + server.endpoint_url.as_str(), + 8, + 8, + defaults.max_inbound_message_size_bytes(), + defaults.max_outbound_message_size_bytes(), + crate::YellowstoneGrpcReconnectSettings::new(3, std::time::Duration::from_millis(5), std::time::Duration::from_millis(20)), + ); + let channel = crate::YellowstoneGrpcChannel::connect(&settings).await.expect("fixture channel must connect"); + let mut session = channel.open_standard_subscribe(initial_request()).await.expect("fixture Subscribe stream must open"); + let first = session.next_update().await.expect("first slot must decode").expect("first slot must be present"); + match first { + crate::YellowstoneSubscribeUpdate::Slot(value) => assert_eq!(value.slot(), 800), + _ => panic!("fixture must return first Slot"), + } + let resumed = session.next_update().await.expect("post-reconnect slot must decode").expect("post-reconnect slot must be present"); + match resumed { + crate::YellowstoneSubscribeUpdate::Slot(value) => assert_eq!(value.slot(), 805), + _ => panic!("fixture must return post-reconnect Slot"), + } + let snapshot = session.snapshot(); + assert_eq!(snapshot.reconnect_count(), 1); + assert_eq!(snapshot.replay_attempt_count(), 1); + assert_eq!(snapshot.replay_delivery_count(), 0); + assert_eq!(snapshot.replay_coverage_unproven_count(), 1); + assert_eq!(snapshot.continuity_gap_count(), 0); + assert_eq!(snapshot.last_requested_from_slot(), std::option::Option::Some(800)); + assert_eq!(snapshot.last_observed_slot(), std::option::Option::Some(805)); + session.close().await.expect("coverage-unproven fixture must close cleanly"); + server.stop().await; +} + #[tokio::test(flavor = "current_thread")] async fn yellowstone_reconnect_budget_exhaustion_is_terminal_and_safe() { let server = FixtureServer::start(FixtureMode::ReconnectExhausted).await; @@ -656,6 +713,8 @@ async fn yellowstone_reconnect_budget_exhaustion_is_terminal_and_safe() { assert_eq!(snapshot.state(), crate::YellowstoneGrpcSubscribeState::Failed); assert_eq!(snapshot.reconnect_count(), 0); assert_eq!(snapshot.replay_attempt_count(), 2); + assert_eq!(snapshot.replay_delivery_count(), 0); + assert_eq!(snapshot.replay_coverage_unproven_count(), 0); assert_eq!(snapshot.terminal_error_code(), std::option::Option::Some(crate::ERROR_CODE_GRPC_CHANNEL_FAILED)); let rendered = format!("{error:?} {session:?}"); assert!(!rendered.contains("GRPC-RECONNECT-SECRET-CANARY")); 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 c0fa5b8..d257ed1 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: 30 +// version: 31 use sha2::Digest; // rust-rules: trait-import @@ -3635,6 +3635,8 @@ struct RawTransactionIngestProcessingFrontierReporter { source_state: std::option::Option, source_reconnect_total: u64, source_replay_attempt_total: u64, + source_replay_delivery_total: u64, + source_replay_coverage_unproven_total: u64, source_continuity_gap_total: u64, source_overflow_total: u64, websocket_incident_anchor: std::option::Option, @@ -3648,6 +3650,8 @@ impl RawTransactionIngestProcessingFrontierReporter { source_state: std::option::Option::None, source_reconnect_total: 0, source_replay_attempt_total: 0, + source_replay_delivery_total: 0, + source_replay_coverage_unproven_total: 0, source_continuity_gap_total: 0, source_overflow_total: 0, websocket_incident_anchor: std::option::Option::None, @@ -3689,6 +3693,8 @@ impl RawTransactionIngestProcessingFrontierReporter { map_yellowstone_source_state(snapshot.state()), snapshot.reconnect_count(), snapshot.replay_attempt_count(), + snapshot.replay_delivery_count(), + snapshot.replay_coverage_unproven_count(), snapshot.continuity_gap_count(), ); } @@ -3777,23 +3783,33 @@ impl RawTransactionIngestProcessingFrontierReporter { state: crate::RawTransactionIngestSourceState, reconnect_total: u64, replay_attempt_total: u64, + replay_delivery_total: u64, + replay_coverage_unproven_total: u64, continuity_gap_total: u64, ) -> ksp_core_lib::Result<()> { if reconnect_total < self.source_reconnect_total || replay_attempt_total < self.source_replay_attempt_total + || replay_delivery_total < self.source_replay_delivery_total + || replay_coverage_unproven_total < self.source_replay_coverage_unproven_total || continuity_gap_total < self.source_continuity_gap_total { return std::result::Result::Err(crate::runtime_error("source.continuity_counter_regression")); } let gap_increased = continuity_gap_total > self.source_continuity_gap_total; + let replay_coverage_unproven_increased = replay_coverage_unproven_total > self.source_replay_coverage_unproven_total; self.source_state = std::option::Option::Some(state); self.source_reconnect_total = reconnect_total; self.source_replay_attempt_total = replay_attempt_total; + self.source_replay_delivery_total = replay_delivery_total; + self.source_replay_coverage_unproven_total = replay_coverage_unproven_total; self.source_continuity_gap_total = continuity_gap_total; self.publish(); if gap_increased { return std::result::Result::Err(crate::runtime_error("source.continuity_gap_proven")); } + if replay_coverage_unproven_increased { + return std::result::Result::Err(crate::runtime_error("source.replay_coverage_unproven")); + } return std::result::Result::Ok(()); } 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 eea25b0..62d9599 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: 27 +// version: 28 -//! External public, security, redaction and release-boundary hardening canaries through `v0.3.14-pre.003`. +//! External public, security, redaction and release-boundary hardening canaries through `v0.3.14-pre.004`. fn network(value: &'static str) -> std::option::Option { let result = ksp_store_lib::RawNetworkId::new(value); @@ -1024,3 +1024,33 @@ fn v0_3_14_pre_003_websocket_continuity_observation_reuses_transport_snapshot_so assert!(!worker.contains("set_from_slot(")); return; } + +#[test] +fn v0_3_14_pre_004_yellowstone_replay_evidence_is_transport_owned_and_fail_closed_without_coverage_proof() { + let worker = include_str!("../src/runtime_resources.rs"); + let transport = include_str!("../../ksp-onchain-transport-lib/src/grpc_stream.rs"); + for required in [ + "snapshot.replay_delivery_count()", + "snapshot.replay_coverage_unproven_count()", + "source_replay_delivery_total", + "source_replay_coverage_unproven_total", + "source.replay_coverage_unproven", + ] { + assert!(worker.contains(required), "required pre.004 Worker replay-evidence guard missing: {required}"); + } + for required in [ + "replay_delivery_count", + "replay_coverage_unproven_count", + "pending_replay_from_slot", + "tracker.begin_replay(effective_from_slot)", + "first_available > requested", + ] { + assert!(transport.contains(required), "required pre.004 Transport replay-evidence guard missing: {required}"); + } + for forbidden in ["subscribe_replay_info(", "set_from_slot(", "YellowstoneReplayInfo", "last_requested_from_slot()"] { + assert!(!worker.contains(forbidden), "pre.004 Worker took replay ownership: {forbidden}"); + } + assert!(!transport.contains("replay_covered_count")); + assert!(!transport.contains("replay_coverage_proven_count")); + 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 06e6d6e..4f3e8d4 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: 22 +// version: 23 -//! Release-completeness canaries through the `v0.3.14-pre.003` WebSocket continuity-observation tranche. +//! Release-completeness canaries through the `v0.3.14-pre.004` Yellowstone native replay-evidence tranche. #[test] fn pre_010_production_module_inventory_is_exact() -> std::io::Result<()> { @@ -248,3 +248,18 @@ fn v0_3_13_pre_012_cross_layer_completeness_security_inventory_is_exact() { assert!(public_api.contains("v0_3_13_pre_012_completeness_closure_adds_no_public_implementation_surface")); return; } + +#[test] +fn v0_3_14_pre_004_native_replay_evidence_canaries_are_present_without_public_replay_ownership() { + let hardening = include_str!("hardening.rs"); + let resources = include_str!("../src/runtime_resources.rs"); + let resource_tests = include_str!("../unit_tests/runtime_resources.rs"); + let root = include_str!("../src/lib.rs"); + assert!(hardening.contains("v0_3_14_pre_004_yellowstone_replay_evidence_is_transport_owned_and_fail_closed_without_coverage_proof")); + assert!(resource_tests.contains("v0_3_14_pre_004_replay_delivery_is_not_coverage_and_unproven_coverage_fails_closed")); + assert!(resources.contains("source.replay_coverage_unproven")); + assert!(!root.contains("ReplayInfo")); + assert!(!root.contains("from_slot")); + assert!(!root.contains("set_from_slot")); + 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 b007f75..620f0dc 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: 24 +// version: 25 fn grpc_endpoint(cluster: &str) -> std::option::Option { return grpc_endpoint_with_identity(cluster, "yellowstone-fixture", "fixture-provider"); @@ -2659,7 +2659,7 @@ fn pre_008_reconnect_replay_and_proven_gap_are_distinct_monotone_and_frontier_pr let (sender, receiver) = tokio::sync::watch::channel(crate::RawTransactionIngestProcessingFrontierProjection::empty()); let mut reporter = super::RawTransactionIngestProcessingFrontierReporter::new(sender); assert!(reporter.observe_pending(42).is_ok()); - assert!(reporter.observe_source_continuity(crate::RawTransactionIngestSourceState::Reconnecting, 0, 1, 0).is_ok()); + assert!(reporter.observe_source_continuity(crate::RawTransactionIngestSourceState::Reconnecting, 0, 1, 0, 0, 0).is_ok()); let replaying = *receiver.borrow(); assert_eq!(replaying.hydration_pending(), 1); assert_eq!(replaying.oldest_pending_slot(), std::option::Option::Some(42)); @@ -2667,8 +2667,8 @@ fn pre_008_reconnect_replay_and_proven_gap_are_distinct_monotone_and_frontier_pr assert_eq!(replaying.source_reconnect_total(), 0); assert_eq!(replaying.source_replay_attempt_total(), 1); assert_eq!(replaying.source_continuity_gap_total(), 0); - assert!(reporter.observe_source_continuity(crate::RawTransactionIngestSourceState::Active, 1, 1, 0).is_ok()); - let gap = reporter.observe_source_continuity(crate::RawTransactionIngestSourceState::Reconnecting, 1, 2, 1); + assert!(reporter.observe_source_continuity(crate::RawTransactionIngestSourceState::Active, 1, 1, 1, 0, 0).is_ok()); + let gap = reporter.observe_source_continuity(crate::RawTransactionIngestSourceState::Reconnecting, 1, 2, 1, 0, 1); assert!(gap.is_err()); let gap = match gap { std::result::Result::Ok(()) => return, @@ -2684,7 +2684,37 @@ fn pre_008_reconnect_replay_and_proven_gap_are_distinct_monotone_and_frontier_pr assert_eq!(proven.source_reconnect_total(), 1); assert_eq!(proven.source_replay_attempt_total(), 2); assert_eq!(proven.source_continuity_gap_total(), 1); - let regression = reporter.observe_source_continuity(crate::RawTransactionIngestSourceState::Active, 0, 2, 1); + let regression = reporter.observe_source_continuity(crate::RawTransactionIngestSourceState::Active, 0, 2, 1, 0, 1); + assert!(regression.is_err()); + let regression = match regression { + std::result::Result::Ok(()) => return, + std::result::Result::Err(value) => value, + }; + assert!(regression.context().iter().any(|context| { + return context.value() == "source.continuity_counter_regression"; + })); + return; +} + +#[test] +fn v0_3_14_pre_004_replay_delivery_is_not_coverage_and_unproven_coverage_fails_closed() { + let (sender, receiver) = tokio::sync::watch::channel(crate::RawTransactionIngestProcessingFrontierProjection::empty()); + let mut reporter = super::RawTransactionIngestProcessingFrontierReporter::new(sender); + assert!(reporter.observe_source_continuity(crate::RawTransactionIngestSourceState::Active, 1, 1, 0, 0, 0).is_ok()); + let attempted = *receiver.borrow(); + assert_eq!(attempted.source_reconnect_total(), 1); + assert_eq!(attempted.source_replay_attempt_total(), 1); + assert_eq!(attempted.source_continuity_gap_total(), 0); + let unproven = reporter.observe_source_continuity(crate::RawTransactionIngestSourceState::Active, 1, 1, 1, 1, 0); + assert!(unproven.is_err()); + let unproven = match unproven { + std::result::Result::Ok(()) => return, + std::result::Result::Err(value) => value, + }; + assert!(unproven.context().iter().any(|context| { + return context.value() == "source.replay_coverage_unproven"; + })); + let regression = reporter.observe_source_continuity(crate::RawTransactionIngestSourceState::Active, 1, 1, 0, 1, 0); assert!(regression.is_err()); let regression = match regression { std::result::Result::Ok(()) => return, diff --git a/deltas/0.3.14/pre.004.md b/deltas/0.3.14/pre.004.md new file mode 100644 index 0000000..452ba60 --- /dev/null +++ b/deltas/0.3.14/pre.004.md @@ -0,0 +1,367 @@ + + + +# Delta `0.3.14-pre.004` — preuve conservative de replay natif Yellowstone + +## Base requise + +```text +0.3.14-pre.003 +workspace.package.version = 0.3.14-pre.3 +deltas/0.3.14/pre.003.md présent +``` + +## Gate de la base + +Le gate opérateur de `0.3.14-pre.003` 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-onchain-transport-lib --all-targets --all-features : PASS +cargo test -p ksp-worker-raw-transaction-ingest-lib --all-targets --all-features : PASS +``` + +Le gate Transport comprend notamment : + +```text +390 unit tests : PASS +public_api : 53 PASS +release_completeness : 45 PASS +smokes live : ignorés par défaut conformément au contrat existant +``` + +Le gate Worker comprend notamment : + +```text +119 unit tests : PASS +cross_layer_completeness : 8 PASS +dependency_boundary : 19 PASS +hardening : 28 PASS +public_api : 20 PASS +release_completeness : 5 PASS +``` + +## Objectif + +Implémenter strictement la tranche `pre.004` du plan `035` : + +```text +conserver l'ownership Yellowstone from_slot / SubscribeReplayInfo dans Transport +distinguer replay attempt de replay delivery +distinguer toute delivery ponctuelle d'une preuve de coverage +conserver first_available comme preuve de limite de rétention uniquement +rendre observable une reprise acceptée dont la couverture reste non prouvée +rester fail-closed côté Worker tant que les tranches de réconciliation ne sont pas présentes +ne créer aucun replay Worker, aucun repair HTTP et aucune nouvelle source +``` + +## Transport — distinction attempt / delivery / coverage non prouvée + +`YellowstoneGrpcSubscribeSnapshot` conserve `replay_attempt_count` et ajoute deux compteurs sûrs : + +```text +replay_delivery_count +replay_coverage_unproven_count +``` + +La sémantique est volontairement conservative. + +`replay_delivery_count` n'augmente que lorsqu'une reconnexion replay-bearing redélivre exactement le slot de reprise effectivement demandé par Transport. Cette observation prouve uniquement qu'un matériau slot-bearing a été livré au bord demandé. Elle ne prouve pas que tous les updates correspondant aux filtres ont été rejoués sur l'intervalle. + +`replay_coverage_unproven_count` augmente lorsqu'une reconnexion replay-bearing a réussi et que : + +```text +le premier update slot-bearing post-reconnect atteint ou dépasse la borne demandée +ou +un nouveau reconnect commence avant tout matériau slot-bearing replayé +``` + +Même une redélivrance exacte de la borne fait donc progresser `replay_delivery_count` **et** `replay_coverage_unproven_count` : la première prouve une delivery ponctuelle, la seconde conserve explicitement l'absence de preuve de couverture complète. + +Cette valeur signifie explicitement **coverage non prouvée**. Elle ne prouve pas qu'un événement filtré existait, qu'il a été perdu, ni que le provider a nécessairement violé son contrat. + +Aucun compteur générique `replay_covered` / `replay_coverage_proven` n'est ajouté. Le replay natif Transport ne possède pas aujourd'hui une preuve générique suffisante pour faire cette affirmation. + +## Transport — tracker de replay + +Le `ContinuityTracker` existant conserve maintenant une seule borne privée : + +```text +pending_replay_from_slot +``` + +Lorsqu'un reconnect physique Yellowstone réussit avec un `effective_from_slot`, Transport ouvre cette preuve pending. + +Sur le premier update slot-bearing : + +```text +slot == requested -> delivery conservatrice prouvée + coverage globale toujours non prouvée +slot > requested -> aucune delivery de borne prouvée + coverage globale non prouvée +slot < requested -> aucune conclusion positive ; la preuve reste pending +``` + +Si le stream se coupe de nouveau avant résolution de la borne pending, la coverage de cette tentative est classée non prouvée avant d'ouvrir la tentative suivante. + +Les `Ping`/`Pong` ne ferment jamais cette preuve parce qu'ils ne portent aucun slot. + +## Transport — rétention `SubscribeReplayInfo` + +Le contrat antérieur reste inchangé : + +```text +first_available > requested +``` + +incrémente `continuity_gap_count`, clamp le `from_slot` effectif dans Transport et signifie uniquement que l'endpoint ne peut pas rejouer avant `first_available`. + +Le clamp ne devient pas une preuve de perte d'un événement filtré et ne devient pas une preuve de coverage complète. + +## Cas replay accepté mais couverture non prouvée + +Une fixture dédiée couvre le cas adversarial prévu au plan : + +```text +avant coupure : slot 800 observé +reconnect : from_slot = 800 accepté +premier update après reconnect : slot 805 +``` + +Résultat attendu : + +```text +reconnect_count = 1 +replay_attempt_count = 1 +replay_delivery_count = 0 +replay_coverage_unproven_count = 1 +continuity_gap_count = 0 +``` + +Le résultat ne prétend ni qu'un update manquait réellement entre 800 et 805, ni que le replay est complet. Il démontre seulement que KSP ne possède pas la preuve conservative de delivery requise pour assimiler la reprise à une couverture saine. + +La fixture de replay redélivrant `slot 500` conserve au contraire : + +```text +replay_attempt_count = 1 +replay_delivery_count = 1 +replay_coverage_unproven_count = 1 +``` + +Cette delivery ne devient toujours pas un `replay_covered`. + +## Worker — consommation des preuves Transport + +`RawTransactionIngestProcessingFrontierReporter` observe désormais de manière monotone : + +```text +replay_attempt_total +replay_delivery_total +replay_coverage_unproven_total +continuity_gap_total +``` + +Le Worker ne calcule aucun `from_slot`, n'appelle jamais `SubscribeReplayInfo` et ne construit aucune requête replay. + +La phase `replay_attempt` peut être observée avant toute delivery et ne ferme aucune obligation de coverage. Une delivery de borne fait ensuite progresser la preuve de delivery mais conserve simultanément `coverage_unproven`. + +Une progression de `replay_coverage_unproven_total` sans gap de rétention produit actuellement : + +```text +source.replay_coverage_unproven +``` + +et conserve la terminalité fail-closed de la source. Cette terminalité est volontairement temporaire jusqu'aux tranches de coverage redondante et de réconciliation du plan `035` ; elle évite qu'un replay accepté mais non prouvé soit assimilé silencieusement à une continuité saine. + +Si `continuity_gap_total` augmente simultanément, la preuve de rétention plus forte conserve la priorité et produit le contrat existant : + +```text +source.continuity_gap_proven +``` + +Les nouveaux compteurs privés sont soumis au même rejet de régression que les compteurs reconnect/replay existants. + +## Surface publique + +Transport étend uniquement le snapshot Yellowstone déjà public avec : + +```text +YellowstoneGrpcSubscribeSnapshot::replay_delivery_count() +YellowstoneGrpcSubscribeSnapshot::replay_coverage_unproven_count() +``` + +Aucun type provider, payload, endpoint, credential ou proto upstream n'est exposé. + +La façade publique Worker n'est pas étendue. Les preuves supplémentaires restent privées au runtime Worker dans cette tranche. + +## Documentation d'utilisation + +`ksp-onchain-transport-lib/USAGE.md` explicite désormais que : + +```text +attempt != delivery != coverage +une delivery ponctuelle n'est pas une preuve de replay_covered +coverage_unproven est une absence explicite de preuve, pas une preuve de perte +``` + +Le document reste version-neutral conformément aux règles workspace. + +## Tests ajoutés/étendus + +Transport couvre : + +```text +replay avec redélivrance exacte de la borne -> delivery = 1 +rétention first_available -> gap = 1 tout en conservant delivery distincte +reconnects exhaustés -> attempts sans fausse delivery +replay accepté puis progression au-delà de la borne -> coverage_unproven = 1 +absence de compteur générique replay_covered/replay_coverage_proven +accessibilité crate-root des deux getters snapshot +``` + +Worker couvre : + +```text +attempt seul != delivery +delivery de borne != coverage +coverage_unproven devient fail-closed +régression delivery/coverage_unproven rejetée +absence persistante de set_from_slot / SubscribeReplayInfo / YellowstoneReplayInfo dans Worker +``` + +## Fichiers ajoutés + +```text +deltas/0.3.14/pre.004.md +``` + +## Fichiers modifiés + +```text +Cargo.toml +crates/ksp-onchain-transport-lib/USAGE.md +crates/ksp-onchain-transport-lib/src/grpc_stream.rs +crates/ksp-onchain-transport-lib/tests/public_api.rs +crates/ksp-onchain-transport-lib/tests/release_completeness.rs +crates/ksp-onchain-transport-lib/unit_tests/grpc_stream.rs +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`, cette nouvelle prerelease synchronise la version workspace : + +```text +header Cargo.toml : 570 -> 571 +workspace.package.version : 0.3.14-pre.3 -> 0.3.14-pre.4 +``` + +Versions des fichiers modifiés : + +```text +ksp-onchain-transport-lib/USAGE.md : 25 -> 26 +ksp-onchain-transport-lib/src/grpc_stream.rs : 4 -> 5 +ksp-onchain-transport-lib/tests/public_api.rs : 53 -> 54 +ksp-onchain-transport-lib/tests/release_completeness.rs : 45 -> 46 +ksp-onchain-transport-lib/unit_tests/grpc_stream.rs : 5 -> 6 +ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs : 30 -> 31 +ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs : 27 -> 28 +ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs : 22 -> 23 +ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs : 24 -> 25 +``` + +## Frontières préservées + +```text +from_slot reste Transport-owned +SubscribeReplayInfo reste Transport-owned +aucun set_from_slot dans Worker +aucun replay loop Worker +aucun retry réseau Worker +aucun repair HTTP +aucune nouvelle source +aucun nouveau provider +aucune nouvelle dépendance +aucune nouvelle feature +aucun changement Config +aucun changement Store +aucun changement Job Backfill +aucun backend Store physique depuis Worker +aucun replay_covered générique inventé par Transport +aucun élargissement de la façade publique Worker +``` + +## Validations exécutées + +Dans le sandbox de préparation : + +```text +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 +scan exact de possession replay dans Worker +scan des nouveaux compteurs attempt/delivery/coverage_unproven +comparaison exacte 0.3.14-pre.003 -> 0.3.14-pre.004 +contrôle des versions de fichier modifiées +contrôle du contenu de l'archive delta +unzip -t de l'archive delta +réapplication de l'archive sur la base puis comparaison byte-exacte +``` + +## Validations non exécutées + +Le sandbox de préparation ne fournit pas le toolchain Cargo/Rust. Les gates suivants restent donc explicitement à exécuter côté opérateur : + +```text +cargo fmt --all +cargo fmt --all -- --check +cargo check --workspace +cargo clippy --workspace --all-targets --all-features -- -D warnings +cargo test -p ksp-onchain-transport-lib --all-targets --all-features +cargo test -p ksp-worker-raw-transaction-ingest-lib --all-targets --all-features +``` + +## Décisions prises + +```text +replay attempt ne prouve aucune delivery +delivery au slot demandé ne prouve aucune coverage complète +absence de delivery conservative avant progression/reconnect suivant devient coverage_unproven +first_available reste une preuve de rétention et non une preuve d'événement perdu +aucun replay_covered n'est revendiqué par Transport générique +le Worker reste fail-closed sur coverage_unproven jusqu'aux tranches de réconciliation +``` + +## Questions ouvertes + +```text +aucune pour pre.004 +``` + +## Tranche suivante + +`pre.005` reste dédiée à la coverage redondante : relations exact/superset conservatrices, epochs/ranges de coverage et refus explicite des équivalences cross-family opportunistes. + +## 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-onchain-transport-lib --all-targets --all-features +cargo test -p ksp-worker-raw-transaction-ingest-lib --all-targets --all-features +```