From 6721a4142aa2dc25c0076a665c765754a4caa200 Mon Sep 17 00:00:00 2001 From: SinuS Von SifriduS Date: Fri, 11 Sep 2026 21:54:19 +0200 Subject: [PATCH] v0.3.14-pre.003 --- Cargo.toml | 4 +- crates/ksp-onchain-transport-lib/src/lib.rs | 4 +- .../src/ws_protocol_session.rs | 14 +- .../src/ws_session.rs | 45 ++- .../tests/public_api.rs | 13 +- .../tests/release_completeness.rs | 22 +- .../unit_tests/ws_session.rs | 39 +- .../src/continuity.rs | 81 ++++- .../src/lib.rs | 4 +- .../src/runtime_resources.rs | 204 ++++++++++- .../tests/hardening.rs | 16 +- .../tests/release_completeness.rs | 5 +- .../unit_tests/continuity.rs | 38 +- .../unit_tests/runtime_resources.rs | 72 +++- deltas/0.3.14/pre.003.md | 333 ++++++++++++++++++ 15 files changed, 872 insertions(+), 22 deletions(-) create mode 100644 deltas/0.3.14/pre.003.md diff --git a/Cargo.toml b/Cargo.toml index 555b196..28a383a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,12 @@ # file: Cargo.toml -# version: 569 +# version: 570 [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.2.fix.2" +version = "0.3.14-pre.3" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-onchain-transport-lib/src/lib.rs b/crates/ksp-onchain-transport-lib/src/lib.rs index acb0cdc..1152d03 100644 --- a/crates/ksp-onchain-transport-lib/src/lib.rs +++ b/crates/ksp-onchain-transport-lib/src/lib.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/src/lib.rs -// version: 47 +// version: 48 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -555,6 +555,8 @@ pub use self::ws_protocol_session::HeliusLaserStreamWsSession; pub use self::ws_protocol_session::SolanaStandardWsSession; /// Shareable compatibility handle for one explicitly created standard Solana physical WebSocket session. pub use self::ws_session::WsSession; +/// Cloneable latest-value observer for one safe physical WebSocket session snapshot. +pub use self::ws_session::WsSessionSnapshotSource; /// Open cluster or network descriptor used by WebSocket endpoint settings. pub use self::ws_settings::WsClusterName; /// Runtime settings for one named WebSocket endpoint. diff --git a/crates/ksp-onchain-transport-lib/src/ws_protocol_session.rs b/crates/ksp-onchain-transport-lib/src/ws_protocol_session.rs index 5ba2982..89ef975 100644 --- a/crates/ksp-onchain-transport-lib/src/ws_protocol_session.rs +++ b/crates/ksp-onchain-transport-lib/src/ws_protocol_session.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/src/ws_protocol_session.rs -// version: 6 +// version: 7 /// Typed facade for one standard Solana WebSocket physical session. /// @@ -42,6 +42,12 @@ impl SolanaStandardWsSession { return self.inner.snapshot(); } + /// Returns a cloneable latest-value observer for safe physical-session snapshots from the shared actor. + #[must_use] + pub fn snapshot_source(&self) -> crate::WsSessionSnapshotSource { + return self.inner.snapshot_source(); + } + /// Returns the latest observable physical-session state. #[must_use] pub fn state(&self) -> crate::WsSessionState { @@ -117,6 +123,12 @@ impl HeliusLaserStreamWsSession { return self.inner.snapshot(); } + /// Returns a cloneable latest-value observer for safe physical-session snapshots from the shared actor. + #[must_use] + pub fn snapshot_source(&self) -> crate::WsSessionSnapshotSource { + return self.inner.snapshot_source(); + } + /// Returns the latest observable physical-session state. #[must_use] pub fn state(&self) -> crate::WsSessionState { diff --git a/crates/ksp-onchain-transport-lib/src/ws_session.rs b/crates/ksp-onchain-transport-lib/src/ws_session.rs index 0164485..beb0f5b 100644 --- a/crates/ksp-onchain-transport-lib/src/ws_session.rs +++ b/crates/ksp-onchain-transport-lib/src/ws_session.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/src/ws_session.rs -// version: 14 +// version: 15 use futures_util::SinkExt; // rust-rules: trait-import use futures_util::StreamExt; // rust-rules: trait-import @@ -10,6 +10,43 @@ const HELIUS_WS_HEARTBEAT_INTERVAL: std::time::Duration = std::time::Duration::f type WsPhysicalStream = tokio_tungstenite::WebSocketStream>; +/// Cloneable latest-value observer for one physical WebSocket session snapshot. +/// +/// The observer exposes only the already-safe [`crate::WsSessionSnapshot`] projection and keeps the internal Tokio watch channel private. Cloning it does +/// not create another socket, actor, reconnect loop or subscription registry. +#[derive(Clone)] +pub struct WsSessionSnapshotSource { + receiver: tokio::sync::watch::Receiver, +} + +impl crate::WsSessionSnapshotSource { + fn new(receiver: tokio::sync::watch::Receiver) -> Self { + return Self { receiver }; + } + + /// Returns the current safe physical-session snapshot without waiting for another actor transition. + #[must_use] + pub fn current(&self) -> crate::WsSessionSnapshot { + return (*self.receiver.borrow()).clone(); + } + + /// Waits for one newer safe physical-session snapshot. + /// + /// `None` means the owning WebSocket actor dropped the latest-value publisher and no further snapshot can arrive. + pub async fn wait_for_change(&mut self) -> std::option::Option { + return match self.receiver.changed().await { + std::result::Result::Ok(()) => std::option::Option::Some((*self.receiver.borrow_and_update()).clone()), + std::result::Result::Err(_) => std::option::Option::None, + }; + } +} + +impl std::fmt::Debug for crate::WsSessionSnapshotSource { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter.debug_struct("WsSessionSnapshotSource").field("current", &self.current()).finish(); + } +} + /// Shareable handle for one explicitly created physical WebSocket session. /// /// The handle never exposes the sensitive endpoint URL or the underlying socket. All socket I/O is owned by one internal actor task and all caller @@ -124,6 +161,12 @@ impl WsSession { return self.snapshot_rx.borrow().clone(); } + /// Returns a cloneable latest-value observer for safe physical-session snapshots. + #[must_use] + pub fn snapshot_source(&self) -> crate::WsSessionSnapshotSource { + return crate::WsSessionSnapshotSource::new(self.snapshot_rx.clone()); + } + /// Returns the latest observable physical-session state. #[must_use] pub fn state(&self) -> crate::WsSessionState { diff --git a/crates/ksp-onchain-transport-lib/tests/public_api.rs b/crates/ksp-onchain-transport-lib/tests/public_api.rs index 89dd645..687766b 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: 52 +// version: 53 //! Integration tests for the public `ksp-onchain-transport-lib` consumer contract. @@ -569,6 +569,17 @@ fn public_v0_2_7_pre_004_physical_websocket_session_contract_is_available_from_c assert_eq!(ksp_onchain_transport_lib::ERROR_CODE_WS_SESSION_CLOSED.code(), "ws_session_closed"); } +#[test] +fn public_v0_3_14_pre_003_websocket_snapshot_source_is_available_from_crate_root() { + let _physical = ksp_onchain_transport_lib::WsSession::snapshot_source; + let _standard = ksp_onchain_transport_lib::SolanaStandardWsSession::snapshot_source; + let _helius = ksp_onchain_transport_lib::HeliusLaserStreamWsSession::snapshot_source; + let _current = ksp_onchain_transport_lib::WsSessionSnapshotSource::current; + let _wait = ksp_onchain_transport_lib::WsSessionSnapshotSource::wait_for_change; + assert!(std::any::type_name::().ends_with("WsSessionSnapshotSource")); + return; +} + #[test] fn public_v0_2_7_pre_005_bounded_websocket_close_contract_is_available_from_crate_root() { let _close = ksp_onchain_transport_lib::WsSession::close; diff --git a/crates/ksp-onchain-transport-lib/tests/release_completeness.rs b/crates/ksp-onchain-transport-lib/tests/release_completeness.rs index 5e0e6c6..27c6539 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: 44 +// version: 45 //! Release-level completeness canaries for staged HTTP and WebSocket Transport coverage. @@ -1354,3 +1354,23 @@ fn release_v0_3_13_pre_002_yellowstone_subscribe_identity_remains_opaque_and_dep assert!(!root.contains("pub use yellowstone_grpc_proto")); return; } + +#[test] +fn release_v0_3_14_pre_003_ws_snapshot_source_reuses_actor_watch_without_second_runtime() { + let session_source = include_str!("../src/ws_session.rs"); + let facade_source = include_str!("../src/ws_protocol_session.rs"); + let session_tests = include_str!("../unit_tests/ws_session.rs"); + assert!(session_source.contains("pub struct WsSessionSnapshotSource")); + assert!(session_source.contains("self.snapshot_rx.clone()")); + assert!(facade_source.matches("pub fn snapshot_source").count() >= 2); + for required in [ + "v0_3_14_pre_003_snapshot_source_tracks_latest_value_without_owning_runtime", + "websocket_reconnect_resubscribes_in_local_id_order_and_remaps_remote_ids", + "websocket_notification_queue_overflow_fails_only_slow_subscription_and_cleans_remote_binding", + ] { + assert!(session_tests.contains(required), "missing pre.003 WebSocket continuity test: {required}"); + } + assert!(!session_source.contains("WsSessionSnapshotSource {\n socket:")); + assert!(!facade_source.contains("tokio_tungstenite::connect_async")); + return; +} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs b/crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs index a48c3d8..d6648f8 100644 --- a/crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs +++ b/crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs -// version: 11 +// version: 12 use futures_util::SinkExt; // rust-rules: trait-import use futures_util::StreamExt; // rust-rules: trait-import @@ -1355,3 +1355,40 @@ async fn websocket_dropped_notification_receiver_triggers_remote_cleanup_and_rel session.close().await.expect("session close must remain bounded"); server.await.expect("local server task must complete"); } + +#[tokio::test(flavor = "current_thread")] +async fn v0_3_14_pre_003_snapshot_source_tracks_latest_value_without_owning_runtime() { + let session_id = crate::WsSessionId::new(std::num::NonZeroU64::new(1).expect("test session ID must be non-zero")); + let initial = crate::WsSessionSnapshot::new( + session_id, + "local_ws", + crate::WsProviderName::new("local-fixture"), + crate::WsClusterName::new("local"), + crate::WsProtocolKind::SolanaStandard, + crate::WsSessionState::Active, + 0, + 0, + 0, + std::vec::Vec::new(), + ); + let reconnecting = crate::WsSessionSnapshot::new( + session_id, + "local_ws", + crate::WsProviderName::new("local-fixture"), + crate::WsClusterName::new("local"), + crate::WsProtocolKind::SolanaStandard, + crate::WsSessionState::Reconnecting { attempt: 1 }, + 0, + 1, + 0, + std::vec::Vec::new(), + ); + let (sender, receiver) = tokio::sync::watch::channel(initial.clone()); + let mut source = crate::WsSessionSnapshotSource::new(receiver); + assert_eq!(source.current(), initial); + sender.send_replace(reconnecting.clone()); + assert_eq!(source.wait_for_change().await, std::option::Option::Some(reconnecting)); + drop(sender); + assert_eq!(source.wait_for_change().await, std::option::Option::None); + return; +} diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/src/continuity.rs b/crates/ksp-worker-raw-transaction-ingest-lib/src/continuity.rs index 4c1918d..9d8838b 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/src/continuity.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/src/continuity.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/src/continuity.rs -// version: 3 +// version: 4 /// Maximum number of simultaneously retained non-repaired run-local gaps. const MAX_RAW_TRANSACTION_INGEST_OPEN_REPAIR_GAPS: usize = 64; @@ -150,6 +150,85 @@ impl RawTransactionIngestGapRange { } } +/// Private run-local anchor for one WebSocket continuity incident. +/// +/// The anchor never derives slots from wall-clock time. Its inclusive start is the latest slot actually observed by that source before Transport reported a +/// reconnect or notification overflow; the end remains absent until the first post-incident source slot is observed. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub(crate) struct RawTransactionIngestWebSocketIncidentAnchor { + end_slot: std::option::Option, + overflow_total: u64, + reconnect_total: u64, + saw_overflow: bool, + saw_reconnect: bool, + start_slot: u64, +} + +impl crate::RawTransactionIngestWebSocketIncidentAnchor { + /// Creates one open incident anchor from monotone Transport counters and one actually observed source slot. + pub(crate) fn new(start_slot: u64, reconnect_total: u64, overflow_total: u64, saw_reconnect: bool, saw_overflow: bool) -> ksp_core_lib::Result { + if !saw_reconnect && !saw_overflow { + return std::result::Result::Err(crate::runtime_error("continuity.websocket_incident_reason_missing")); + } + return std::result::Result::Ok(Self { + end_slot: std::option::Option::None, + overflow_total, + reconnect_total, + saw_overflow, + saw_reconnect, + start_slot, + }); + } + + /// Extends one still-open incident with newer monotone Transport counters without changing its earliest start slot. + pub(crate) fn extend(&mut self, reconnect_total: u64, overflow_total: u64, saw_reconnect: bool, saw_overflow: bool) -> ksp_core_lib::Result<()> { + if self.end_slot.is_some() { + return std::result::Result::Err(crate::runtime_error("continuity.websocket_incident_already_closed")); + } + if reconnect_total < self.reconnect_total || overflow_total < self.overflow_total { + return std::result::Result::Err(crate::runtime_error("continuity.websocket_incident_counter_regression")); + } + self.reconnect_total = reconnect_total; + self.overflow_total = overflow_total; + self.saw_reconnect = self.saw_reconnect || saw_reconnect; + self.saw_overflow = self.saw_overflow || saw_overflow; + return std::result::Result::Ok(()); + } + + /// Closes the inclusive incident range at the first post-incident source slot. + pub(crate) fn close_at(&mut self, end_slot: u64) -> ksp_core_lib::Result<()> { + let range = match RawTransactionIngestGapRange::new(self.start_slot, end_slot) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + self.end_slot = std::option::Option::Some(range.end_slot()); + return std::result::Result::Ok(()); + } + + /// Returns the first source slot that safely anchors the incident, inclusively. + #[cfg(test)] + pub(crate) const fn start_slot(self) -> u64 { + return self.start_slot; + } + + /// Returns the first post-incident source slot when the incident has become range-bounded. + pub(crate) const fn end_slot(self) -> std::option::Option { + return self.end_slot; + } + + /// Returns whether at least one physical WebSocket reconnect contributed to this incident. + #[cfg(test)] + pub(crate) const fn saw_reconnect(self) -> bool { + return self.saw_reconnect; + } + + /// Returns whether at least one Transport notification overflow contributed to this incident. + #[cfg(test)] + pub(crate) const fn saw_overflow(self) -> bool { + return self.saw_overflow; + } +} + #[derive(Clone, Copy, Debug, Eq, PartialEq, Ord, PartialOrd)] struct RawTransactionIngestGapId(u64); 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 3b3d802..bb0f069 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: 28 +// version: 29 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -117,6 +117,8 @@ pub(crate) use self::continuity::RawTransactionIngestContinuityCapabilityDescrip pub(crate) use self::continuity::RawTransactionIngestContinuityContracts; /// Private provider-neutral coverage scope used by run-local continuity proof contracts. pub(crate) use self::continuity::RawTransactionIngestCoverageScope; +/// Private run-local WebSocket incident anchor built only from observed source slots and safe Transport counters. +pub(crate) use self::continuity::RawTransactionIngestWebSocketIncidentAnchor; /// Creates one terminal content-conflict error without copying conflicting material into diagnostics. pub(crate) use self::error::content_conflict_error; /// Creates one terminal counter-exhaustion error without exposing runtime material. 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 43cb8da..c0fa5b8 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: 29 +// version: 30 use sha2::Digest; // rust-rules: trait-import @@ -1055,12 +1055,17 @@ impl crate::RawTransactionIngestHeliusTransactionSource { return std::result::Result::Err(source_transport_error(error.code())); }, }; + let mut session_snapshot_source = session.snapshot_source(); let hydration = self.hydration_context(); let mut coordinator = RawTransactionIngestHydrationCoordinator::with_global_registry(global_hydration_registry, hydration_pending_limit, hydration_in_flight_limit); let mut processing_frontier = RawTransactionIngestProcessingFrontierReporter::new(processing_frontier_sender); - processing_frontier.set_source_state(crate::RawTransactionIngestSourceState::Active); + if let std::result::Result::Err(error) = processing_frontier.observe_websocket_session_snapshot(session_snapshot_source.current()) { + let _closed = session.close().await; + return std::result::Result::Err(error); + } let mut fault = std::option::Option::None; + let mut bounded_websocket_incident = false; loop { if *stop_receiver.borrow() { break; @@ -1069,7 +1074,11 @@ impl crate::RawTransactionIngestHeliusTransactionSource { fault = std::option::Option::Some(error); break; } - let can_receive = coordinator.can_receive(); + if bounded_websocket_incident && coordinator.pending_signal_count == 0 && coordinator.tasks.is_empty() { + fault = std::option::Option::Some(crate::runtime_error("source.continuity_gap_proven")); + break; + } + let can_receive = coordinator.can_receive() && !bounded_websocket_incident; if !can_receive && coordinator.tasks.is_empty() { fault = std::option::Option::Some(crate::runtime_error("source.hydration_stalled")); break; @@ -1079,6 +1088,19 @@ impl crate::RawTransactionIngestHeliusTransactionSource { _ = stop_receiver.changed() => { break; } + source_snapshot = session_snapshot_source.wait_for_change() => { + let source_snapshot = match source_snapshot { + std::option::Option::Some(value) => value, + std::option::Option::None => { + fault = std::option::Option::Some(crate::runtime_error("source.websocket_snapshot_closed")); + break; + }, + }; + if let std::result::Result::Err(error) = processing_frontier.observe_websocket_session_snapshot(source_snapshot) { + fault = std::option::Option::Some(error); + break; + } + } joined = coordinator.tasks.join_next(), if !coordinator.tasks.is_empty() => { let joined = match joined { std::option::Option::Some(value) => value, @@ -1143,10 +1165,18 @@ impl crate::RawTransactionIngestHeliusTransactionSource { break; }, }; + let incident_bounded = match processing_frontier.observe_websocket_post_incident_slot(signal.slot) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => { + fault = std::option::Option::Some(error); + break; + }, + }; if let std::result::Result::Err(error) = coordinator.queue_signal(&hydration, signal, received_at, &mut processing_frontier) { fault = std::option::Option::Some(error); break; } + bounded_websocket_incident = bounded_websocket_incident || incident_bounded; } } } @@ -1541,8 +1571,12 @@ impl crate::RawTransactionIngestStandardBlockSource { return std::result::Result::Err(source_transport_error(error.code())); }, }; + let mut session_snapshot_source = session.snapshot_source(); let mut processing_frontier = RawTransactionIngestProcessingFrontierReporter::new(processing_frontier_sender); - processing_frontier.set_source_state(crate::RawTransactionIngestSourceState::Active); + if let std::result::Result::Err(error) = processing_frontier.observe_websocket_session_snapshot(session_snapshot_source.current()) { + let _closed = session.close().await; + return std::result::Result::Err(error); + } let mut fault = std::option::Option::None; 'source: loop { if *stop_receiver.borrow() { @@ -1553,6 +1587,20 @@ impl crate::RawTransactionIngestStandardBlockSource { _ = stop_receiver.changed() => { break; } + source_snapshot = session_snapshot_source.wait_for_change() => { + let source_snapshot = match source_snapshot { + std::option::Option::Some(value) => value, + std::option::Option::None => { + fault = std::option::Option::Some(crate::runtime_error("source.websocket_snapshot_closed")); + break; + }, + }; + if let std::result::Result::Err(error) = processing_frontier.observe_websocket_session_snapshot(source_snapshot) { + fault = std::option::Option::Some(error); + break; + } + continue; + } value = subscription.recv() => value, }; let notification = match notification { @@ -1584,11 +1632,22 @@ impl crate::RawTransactionIngestStandardBlockSource { break; }, }; + let incident_bounded = match processing_frontier.observe_websocket_post_incident_slot(slot) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => { + fault = std::option::Option::Some(error); + break; + }, + }; if ingresses.is_empty() { if let std::result::Result::Err(error) = processing_frontier.observe_settled(slot) { fault = std::option::Option::Some(error); break; } + if incident_bounded { + fault = std::option::Option::Some(crate::runtime_error("source.continuity_gap_proven")); + break; + } continue; } for ingress in ingresses { @@ -1611,6 +1670,10 @@ impl crate::RawTransactionIngestStandardBlockSource { fault = std::option::Option::Some(error); break; } + if incident_bounded { + fault = std::option::Option::Some(crate::runtime_error("source.continuity_gap_proven")); + break; + } } processing_frontier.discard_all_pending(); processing_frontier.set_source_state(crate::RawTransactionIngestSourceState::Closing); @@ -1777,12 +1840,17 @@ impl crate::RawTransactionIngestStandardLogsSource { return std::result::Result::Err(source_transport_error(error.code())); }, }; + let mut session_snapshot_source = session.snapshot_source(); let hydration = self.hydration_context(); let mut coordinator = RawTransactionIngestHydrationCoordinator::with_global_registry(global_hydration_registry, hydration_pending_limit, hydration_in_flight_limit); let mut processing_frontier = RawTransactionIngestProcessingFrontierReporter::new(processing_frontier_sender); - processing_frontier.set_source_state(crate::RawTransactionIngestSourceState::Active); + if let std::result::Result::Err(error) = processing_frontier.observe_websocket_session_snapshot(session_snapshot_source.current()) { + let _closed = session.close().await; + return std::result::Result::Err(error); + } let mut fault = std::option::Option::None; + let mut bounded_websocket_incident = false; loop { if *stop_receiver.borrow() { break; @@ -1791,7 +1859,11 @@ impl crate::RawTransactionIngestStandardLogsSource { fault = std::option::Option::Some(error); break; } - let can_receive = coordinator.can_receive(); + if bounded_websocket_incident && coordinator.pending_signal_count == 0 && coordinator.tasks.is_empty() { + fault = std::option::Option::Some(crate::runtime_error("source.continuity_gap_proven")); + break; + } + let can_receive = coordinator.can_receive() && !bounded_websocket_incident; if !can_receive && coordinator.tasks.is_empty() { fault = std::option::Option::Some(crate::runtime_error("source.hydration_stalled")); break; @@ -1801,6 +1873,19 @@ impl crate::RawTransactionIngestStandardLogsSource { _ = stop_receiver.changed() => { break; } + source_snapshot = session_snapshot_source.wait_for_change() => { + let source_snapshot = match source_snapshot { + std::option::Option::Some(value) => value, + std::option::Option::None => { + fault = std::option::Option::Some(crate::runtime_error("source.websocket_snapshot_closed")); + break; + }, + }; + if let std::result::Result::Err(error) = processing_frontier.observe_websocket_session_snapshot(source_snapshot) { + fault = std::option::Option::Some(error); + break; + } + } joined = coordinator.tasks.join_next(), if !coordinator.tasks.is_empty() => { let joined = match joined { std::option::Option::Some(value) => value, @@ -1858,10 +1943,18 @@ impl crate::RawTransactionIngestStandardLogsSource { break; }, }; + let incident_bounded = match processing_frontier.observe_websocket_post_incident_slot(signal.slot) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => { + fault = std::option::Option::Some(error); + break; + }, + }; if let std::result::Result::Err(error) = coordinator.queue_signal(&hydration, signal, received_at, &mut processing_frontier) { fault = std::option::Option::Some(error); break; } + bounded_websocket_incident = bounded_websocket_incident || incident_bounded; } } } @@ -3364,6 +3457,18 @@ fn map_yellowstone_source_state(state: ksp_onchain_transport_lib::YellowstoneGrp }; } +fn map_websocket_source_state(state: ksp_onchain_transport_lib::WsSessionState) -> crate::RawTransactionIngestSourceState { + return match state { + ksp_onchain_transport_lib::WsSessionState::Disconnected + | ksp_onchain_transport_lib::WsSessionState::Connecting + | ksp_onchain_transport_lib::WsSessionState::Reconnecting { .. } => crate::RawTransactionIngestSourceState::Reconnecting, + ksp_onchain_transport_lib::WsSessionState::Active => crate::RawTransactionIngestSourceState::Active, + ksp_onchain_transport_lib::WsSessionState::Closing => crate::RawTransactionIngestSourceState::Closing, + ksp_onchain_transport_lib::WsSessionState::Closed => crate::RawTransactionIngestSourceState::Closed, + ksp_onchain_transport_lib::WsSessionState::Failed => crate::RawTransactionIngestSourceState::Failed, + }; +} + #[derive(Clone, Copy, Debug, Eq, PartialEq)] struct RawTransactionIngestProcessingSlotState { pending: usize, @@ -3447,6 +3552,10 @@ impl RawTransactionIngestProcessingFrontier { return; } + fn highest_observed_slot(&self) -> std::option::Option { + return self.slots.keys().next_back().copied(); + } + fn projection(&self) -> crate::RawTransactionIngestProcessingFrontierProjection { let oldest_pending_slot = self.slots.iter().find_map(|(slot, state)| { if state.pending == 0 { @@ -3527,6 +3636,8 @@ struct RawTransactionIngestProcessingFrontierReporter { source_reconnect_total: u64, source_replay_attempt_total: u64, source_continuity_gap_total: u64, + source_overflow_total: u64, + websocket_incident_anchor: std::option::Option, } impl RawTransactionIngestProcessingFrontierReporter { @@ -3538,6 +3649,8 @@ impl RawTransactionIngestProcessingFrontierReporter { source_reconnect_total: 0, source_replay_attempt_total: 0, source_continuity_gap_total: 0, + source_overflow_total: 0, + websocket_incident_anchor: std::option::Option::None, }; } @@ -3580,6 +3693,85 @@ impl RawTransactionIngestProcessingFrontierReporter { ); } + fn observe_websocket_session_snapshot(&mut self, snapshot: ksp_onchain_transport_lib::WsSessionSnapshot) -> ksp_core_lib::Result<()> { + return self.observe_websocket_continuity(map_websocket_source_state(snapshot.state()), snapshot.continuity_gap_count(), snapshot.overflow_count()); + } + + fn observe_websocket_continuity( + &mut self, + state: crate::RawTransactionIngestSourceState, + reconnect_total: u64, + overflow_total: u64, + ) -> ksp_core_lib::Result<()> { + if reconnect_total < self.source_reconnect_total || overflow_total < self.source_overflow_total { + return std::result::Result::Err(crate::runtime_error("source.websocket_continuity_counter_regression")); + } + let continuity_gap_total = match reconnect_total.checked_add(overflow_total) { + std::option::Option::Some(value) => value, + std::option::Option::None => return std::result::Result::Err(crate::counter_exhausted_error("source.websocket_continuity_gap_total")), + }; + if continuity_gap_total < self.source_continuity_gap_total { + return std::result::Result::Err(crate::runtime_error("source.continuity_counter_regression")); + } + let saw_reconnect = reconnect_total > self.source_reconnect_total; + let saw_overflow = overflow_total > self.source_overflow_total; + self.source_reconnect_total = reconnect_total; + self.source_replay_attempt_total = 0; + self.source_continuity_gap_total = continuity_gap_total; + self.source_overflow_total = overflow_total; + if saw_reconnect || saw_overflow { + let start_slot = match self.frontier.highest_observed_slot() { + std::option::Option::Some(value) => value, + std::option::Option::None => { + self.publish(); + return std::result::Result::Err(crate::runtime_error("source.websocket_incident_unbounded")); + }, + }; + match self.websocket_incident_anchor.as_mut() { + std::option::Option::Some(anchor) => { + if let std::result::Result::Err(error) = anchor.extend(reconnect_total, overflow_total, saw_reconnect, saw_overflow) { + self.publish(); + return std::result::Result::Err(error); + } + }, + std::option::Option::None => { + let anchor = + match crate::RawTransactionIngestWebSocketIncidentAnchor::new(start_slot, reconnect_total, overflow_total, saw_reconnect, saw_overflow) + { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => { + self.publish(); + return std::result::Result::Err(error); + }, + }; + self.websocket_incident_anchor = std::option::Option::Some(anchor); + }, + } + } + self.source_state = if self.websocket_incident_anchor.is_some() && state == crate::RawTransactionIngestSourceState::Active { + std::option::Option::Some(crate::RawTransactionIngestSourceState::Reconnecting) + } else { + std::option::Option::Some(state) + }; + self.publish(); + return std::result::Result::Ok(()); + } + + fn observe_websocket_post_incident_slot(&mut self, slot: u64) -> ksp_core_lib::Result { + let anchor = match self.websocket_incident_anchor.as_mut() { + std::option::Option::Some(value) => value, + std::option::Option::None => return std::result::Result::Ok(false), + }; + if anchor.end_slot().is_some() { + return std::result::Result::Ok(false); + } + if let std::result::Result::Err(error) = anchor.close_at(slot) { + return std::result::Result::Err(error); + } + self.publish(); + return std::result::Result::Ok(true); + } + fn observe_source_continuity( &mut self, state: crate::RawTransactionIngestSourceState, 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 09df283..eea25b0 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: 26 +// version: 27 -//! External public, security, redaction and release-boundary hardening canaries through `v0.3.14-pre.002`. +//! External public, security, redaction and release-boundary hardening canaries through `v0.3.14-pre.003`. fn network(value: &'static str) -> std::option::Option { let result = ksp_store_lib::RawNetworkId::new(value); @@ -1012,3 +1012,15 @@ fn v0_3_14_pre_002_continuity_contracts_are_private_bounded_and_io_free() { assert!(root.contains("pub(crate) use self::continuity::RawTransactionIngestContinuityContracts;")); return; } + +#[test] +fn v0_3_14_pre_003_websocket_continuity_observation_reuses_transport_snapshot_sources_only() { + let worker = include_str!("../src/runtime_resources.rs"); + assert!(worker.matches("session.snapshot_source()").count() >= 3); + assert!(worker.contains("observe_websocket_session_snapshot")); + assert!(worker.contains("observe_websocket_post_incident_slot")); + assert!(!worker.contains("tokio_tungstenite::connect_async")); + assert!(!worker.contains("WebSocketStream<")); + assert!(!worker.contains("set_from_slot(")); + 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 58bb712..06e6d6e 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: 21 +// version: 22 -//! Release-completeness canaries through the `v0.3.14-pre.002` continuity-contract tranche. +//! Release-completeness canaries through the `v0.3.14-pre.003` WebSocket continuity-observation tranche. #[test] fn pre_010_production_module_inventory_is_exact() -> std::io::Result<()> { @@ -144,6 +144,7 @@ fn pre_010_external_hardening_suite_is_present_and_scoped() { "v0_3_13_pre_010_multi_source_health_is_conservative_counted_and_redacted", "v0_3_13_pre_011_shutdown_races_are_bounded_joined_atomic_and_counter_safe", "v0_3_14_pre_002_continuity_contracts_are_private_bounded_and_io_free", + "v0_3_14_pre_003_websocket_continuity_observation_reuses_transport_snapshot_sources_only", ] { assert!(hardening.contains(required), "required pre.010 hardening canary missing: {required}"); } diff --git a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/continuity.rs b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/continuity.rs index cab4b41..b809e48 100644 --- a/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/continuity.rs +++ b/crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/continuity.rs @@ -1,5 +1,5 @@ // file: crates/ksp-worker-raw-transaction-ingest-lib/unit_tests/continuity.rs -// version: 2 +// version: 3 fn network() -> std::option::Option { return match ksp_store_lib::RawNetworkId::new("mainnet") { @@ -239,3 +239,39 @@ fn pre_002_gap_ledger_enforces_open_gap_bound_and_next_id_monotonicity() { assert!(stale_next.validate_invariants().is_err()); return; } + +#[test] +fn pre_003_websocket_incident_anchor_is_inclusive_monotone_and_bounded() { + let mut anchor = match crate::RawTransactionIngestWebSocketIncidentAnchor::new(100, 1, 0, true, false) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("unexpected anchor failure: {error}"), + }; + assert_eq!(anchor.start_slot(), 100); + assert_eq!(anchor.end_slot(), std::option::Option::None); + assert!(anchor.saw_reconnect()); + assert!(!anchor.saw_overflow()); + assert!(anchor.extend(2, 1, true, true).is_ok()); + assert!(anchor.saw_reconnect()); + assert!(anchor.saw_overflow()); + assert!(anchor.extend(1, 1, false, false).is_err()); + assert!(anchor.close_at(100 + super::MAX_RAW_TRANSACTION_INGEST_REPAIR_RANGE_SLOTS - 1).is_ok()); + assert_eq!(anchor.end_slot(), std::option::Option::Some(100 + super::MAX_RAW_TRANSACTION_INGEST_REPAIR_RANGE_SLOTS - 1)); + assert!(anchor.extend(3, 1, true, false).is_err()); + return; +} + +#[test] +fn pre_003_websocket_incident_anchor_rejects_missing_reason_reversal_and_oversize() { + assert!(crate::RawTransactionIngestWebSocketIncidentAnchor::new(100, 0, 0, false, false).is_err()); + let mut reversed = match crate::RawTransactionIngestWebSocketIncidentAnchor::new(100, 0, 1, false, true) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("unexpected overflow anchor failure: {error}"), + }; + assert!(reversed.close_at(99).is_err()); + let mut oversized = match crate::RawTransactionIngestWebSocketIncidentAnchor::new(100, 1, 0, true, false) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("unexpected oversize anchor fixture failure: {error}"), + }; + assert!(oversized.close_at(100 + super::MAX_RAW_TRANSACTION_INGEST_REPAIR_RANGE_SLOTS).is_err()); + 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 7650126..b007f75 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: 23 +// version: 24 fn grpc_endpoint(cluster: &str) -> std::option::Option { return grpc_endpoint_with_identity(cluster, "yellowstone-fixture", "fixture-provider"); @@ -3347,3 +3347,73 @@ async fn v0_3_13_pre_011_stop_racing_ready_source_fault_preserves_fault_and_join assert_eq!(sibling_active.load(std::sync::atomic::Ordering::Acquire), 0); return; } + +#[test] +fn v0_3_14_pre_003_websocket_reconnect_is_anchored_before_terminal_gap_projection() { + let (sender, receiver) = tokio::sync::watch::channel(crate::RawTransactionIngestProcessingFrontierProjection::empty()); + let mut reporter = super::RawTransactionIngestProcessingFrontierReporter::new(sender); + assert!(reporter.observe_settled(100).is_ok()); + assert!(reporter.observe_websocket_continuity(crate::RawTransactionIngestSourceState::Active, 0, 0).is_ok()); + assert!(reporter.observe_websocket_continuity(crate::RawTransactionIngestSourceState::Reconnecting, 1, 0).is_ok()); + let anchor = match reporter.websocket_incident_anchor { + std::option::Option::Some(value) => value, + std::option::Option::None => panic!("reconnect incident anchor missing"), + }; + assert_eq!(anchor.start_slot(), 100); + assert_eq!(anchor.end_slot(), std::option::Option::None); + assert!(anchor.saw_reconnect()); + assert!(!anchor.saw_overflow()); + assert!(reporter.observe_websocket_continuity(crate::RawTransactionIngestSourceState::Active, 1, 0).is_ok()); + let bounded = reporter.observe_websocket_post_incident_slot(105); + assert!(matches!(bounded, std::result::Result::Ok(true))); + let anchor = match reporter.websocket_incident_anchor { + std::option::Option::Some(value) => value, + std::option::Option::None => panic!("bounded reconnect incident anchor missing"), + }; + assert_eq!(anchor.end_slot(), std::option::Option::Some(105)); + let projection = *receiver.borrow(); + assert_eq!(projection.source_reconnect_total(), 1); + assert_eq!(projection.source_replay_attempt_total(), 0); + assert_eq!(projection.source_continuity_gap_total(), 1); + return; +} + +#[test] +fn v0_3_14_pre_003_websocket_overflow_and_reconnect_storm_share_earliest_anchor() { + let (sender, receiver) = tokio::sync::watch::channel(crate::RawTransactionIngestProcessingFrontierProjection::empty()); + let mut reporter = super::RawTransactionIngestProcessingFrontierReporter::new(sender); + assert!(reporter.observe_settled(200).is_ok()); + assert!(reporter.observe_websocket_continuity(crate::RawTransactionIngestSourceState::Active, 0, 1).is_ok()); + assert!(reporter.observe_websocket_continuity(crate::RawTransactionIngestSourceState::Reconnecting, 1, 2).is_ok()); + let anchor = match reporter.websocket_incident_anchor { + std::option::Option::Some(value) => value, + std::option::Option::None => panic!("overflow incident anchor missing"), + }; + assert_eq!(anchor.start_slot(), 200); + assert_eq!(anchor.end_slot(), std::option::Option::None); + assert!(anchor.saw_reconnect()); + assert!(anchor.saw_overflow()); + let projection = *receiver.borrow(); + assert_eq!(projection.source_reconnect_total(), 1); + assert_eq!(projection.source_continuity_gap_total(), 3); + let bounded = reporter.observe_websocket_post_incident_slot(201); + assert!(matches!(bounded, std::result::Result::Ok(true))); + return; +} + +#[test] +fn v0_3_14_pre_003_websocket_incident_without_observed_slot_is_unbounded_and_counter_regression_fails() { + let (sender, _receiver) = tokio::sync::watch::channel(crate::RawTransactionIngestProcessingFrontierProjection::empty()); + let mut reporter = super::RawTransactionIngestProcessingFrontierReporter::new(sender); + let unbounded = reporter.observe_websocket_continuity(crate::RawTransactionIngestSourceState::Reconnecting, 1, 0); + assert!(unbounded.is_err()); + let unbounded = match unbounded { + std::result::Result::Ok(()) => return, + std::result::Result::Err(value) => value, + }; + assert!(unbounded.context().iter().any(|context| { + return context.value() == "source.websocket_incident_unbounded"; + })); + assert!(reporter.observe_websocket_continuity(crate::RawTransactionIngestSourceState::Active, 0, 0).is_err()); + return; +} diff --git a/deltas/0.3.14/pre.003.md b/deltas/0.3.14/pre.003.md new file mode 100644 index 0000000..320176b --- /dev/null +++ b/deltas/0.3.14/pre.003.md @@ -0,0 +1,333 @@ + + + +# Delta `0.3.14-pre.003` — observabilité latest-value de continuité WebSocket + +## Base requise + +```text +0.3.14-pre.002-fix.002 +workspace.package.version = 0.3.14-pre.2.fix.2 +deltas/0.3.14/pre.002-fix.002.md présent +``` + +## Gate de la base + +Le gate opérateur de `0.3.14-pre.002-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 +114 unit tests : PASS +cross_layer_completeness : 8 PASS +dependency_boundary : 19 PASS +hardening : 27 PASS +public_api : 20 PASS +release_completeness : 5 PASS +``` + +## Objectif + +Implémenter strictement la tranche `pre.003` du plan `035` : + +```text +ajouter dans Transport une source latest-value sûre pour les snapshots WebSocket typés +la déléguer depuis les façades Standard et Helius existantes +observer les incidents reconnect/overflow dans les sources Standard Logs, Standard Block et Helius du Worker +ancrer un incident sur le dernier slot réellement observé par cette source +borner l'incident avec le premier slot post-incident réellement reçu +conserver une terminalité conservative tant que les tranches de replay/repair ultérieures ne sont pas présentes +ne créer aucun nouveau socket, actor, retry ou repair HTTP +``` + +Cette tranche ne remplace pas encore un gap WebSocket par une réparation. Elle rend l'incident explicitement observable et bornable en réutilisant exclusivement les snapshots déjà publiés par l'actor Transport existant. + +## Transport — source latest-value WebSocket + +`WsSession` possède déjà un `tokio::sync::watch::Receiver` interne utilisé pour `snapshot()` et `state()`. + +`pre.003` ajoute `WsSessionSnapshotSource`, un observateur cloneable qui encapsule uniquement un clone de ce receiver et expose : + +```text +current() -> dernier WsSessionSnapshot sûr +wait_for_change() -> prochaine valeur latest-value, ou None si l'actor disparaît +``` + +Le type ne contient ni URL, ni credentials, ni socket, ni payload de notification. Son `Debug` reste fondé sur la projection sûre existante. + +`WsSession::snapshot_source()` clone le receiver existant. Les façades : + +```text +SolanaStandardWsSession +HeliusLaserStreamWsSession +``` + +délèguent cette méthode au même `WsSession` physique. Aucun second runtime WebSocket n'est créé. + +## Worker — ancre d'incident WebSocket + +Le Worker introduit un contrat crate-private `RawTransactionIngestWebSocketIncidentAnchor`. + +Il conserve uniquement : + +```text +start_slot +end_slot optionnel +compteurs de continuity gap / overflow observés au moment de l'incident +motifs reconnect et/ou overflow +``` + +La borne de début respecte le plan `035` : + +```text +start_slot = dernier slot réellement observé par la source avant l'incident +``` + +Aucun timestamp, frontier dérivé d'une autre source ou `last_slot + 1` n'est utilisé comme preuve de début. + +La borne de fin est : + +```text +end_slot = premier slot réellement reçu par la même source après l'incident +``` + +Le range reste inclusif. Les validations réutilisent les bornes run-local introduites en `pre.002` et rejettent notamment inversion, dépassement de plage et compteur régressif. + +Si un reconnect/overflow est observé avant qu'un slot source n'ait été vu, le Worker ne fabrique pas d'historique : l'incident est `source.websocket_incident_unbounded`. + +## Worker — observation des snapshots + +`RawTransactionIngestProcessingFrontierReporter` observe maintenant les `WsSessionSnapshot` Transport et maintient : + +```text +continuity gap total WebSocket +notification overflow total +ancre d'incident WebSocket éventuelle +projection de lifecycle source +``` + +Une source redevenue `Active` reste projetée `Reconnecting` tant qu'un incident ouvert n'a pas reçu sa borne post-incident. Cela évite de publier transitoirement une continuité saine alors que le gap n'est pas encore borné. + +Aucune mutation de la politique de reconnect Transport n'est effectuée par le Worker. + +## Standard Logs et Helius transaction + +Les deux sources réutilisent maintenant `session.snapshot_source()` en parallèle de la réception métier et des tâches d'hydration existantes. + +Lorsqu'un incident est détecté : + +```text +le premier signal post-incident borne le range de manière inclusive +ce signal est quand même admis dans le coordinator d'hydration existant +aucune nouvelle notification n'est admise après cette borne dans pre.003 +les signaux/tâches déjà admis sont drainés +la source termine ensuite par source.continuity_gap_proven +``` + +Ce séquencement évite de perdre artificiellement le premier événement post-incident tout en restant fail-closed tant que les tranches de repair ne sont pas implémentées. + +## Standard Block + +La source Standard Block observe le même snapshot latest-value Transport. + +Le premier bloc post-incident : + +```text +borne l'incident +projette et envoie toutes ses ingresses via l'admission centrale existante +marque ensuite le slot settled +termine enfin par source.continuity_gap_proven +``` + +Un bloc sans transaction qualifiée peut quand même fermer la borne de slot et est settled avant la terminalité conservative. + +## Yellowstone + +Le chemin Yellowstone existant reste inchangé fonctionnellement dans cette tranche : + +```text +aucune mutation Worker de from_slot +aucun ownership Worker de SubscribeReplayInfo +aucune nouvelle tentative de replay +aucune modification du contrat Transport Yellowstone +``` + +La preuve de replay native et la distinction `replay_attempt` / couverture effective restent réservées à `pre.004`. + +## Tests ajoutés/étendus + +Transport couvre : + +```text +source latest-value courante +publication d'un changement +fermeture de la source lorsque l'actor publisher disparaît +accessibilité crate-root de WsSessionSnapshotSource +présence des délégations Standard/Helius +réutilisation du watch actor existant sans second runtime +``` + +Worker couvre : + +```text +ancre inclusive et bornée +extension monotone d'un incident reconnect/overflow +rejet d'un incident sans motif +rejet de range inversé ou surdimensionné +incident sans slot précédent -> erreur explicite +compteurs WebSocket régressifs -> erreur explicite +reconnect et overflow successifs -> même ancre la plus ancienne +absence de socket/retry/from_slot possédé par le Worker +``` + +## Fichiers ajoutés + +```text +deltas/0.3.14/pre.003.md +``` + +## Fichiers modifiés + +```text +Cargo.toml +crates/ksp-onchain-transport-lib/src/lib.rs +crates/ksp-onchain-transport-lib/src/ws_protocol_session.rs +crates/ksp-onchain-transport-lib/src/ws_session.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/ws_session.rs +crates/ksp-worker-raw-transaction-ingest-lib/src/continuity.rs +crates/ksp-worker-raw-transaction-ingest-lib/src/lib.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/continuity.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 : 569 -> 570 +workspace.package.version : 0.3.14-pre.2.fix.2 -> 0.3.14-pre.3 +``` + +Versions des fichiers modifiés : + +```text +ksp-onchain-transport-lib/src/lib.rs : 47 -> 48 +ksp-onchain-transport-lib/src/ws_protocol_session.rs : 6 -> 7 +ksp-onchain-transport-lib/src/ws_session.rs : 14 -> 15 +ksp-onchain-transport-lib/tests/public_api.rs : 52 -> 53 +ksp-onchain-transport-lib/tests/release_completeness.rs : 44 -> 45 +ksp-onchain-transport-lib/unit_tests/ws_session.rs : 11 -> 12 +ksp-worker-raw-transaction-ingest-lib/src/continuity.rs : 3 -> 4 +ksp-worker-raw-transaction-ingest-lib/src/lib.rs : 28 -> 29 +ksp-worker-raw-transaction-ingest-lib/src/runtime_resources.rs : 29 -> 30 +ksp-worker-raw-transaction-ingest-lib/tests/hardening.rs : 26 -> 27 +ksp-worker-raw-transaction-ingest-lib/tests/release_completeness.rs : 21 -> 22 +ksp-worker-raw-transaction-ingest-lib/unit_tests/continuity.rs : 2 -> 3 +ksp-worker-raw-transaction-ingest-lib/unit_tests/runtime_resources.rs : 23 -> 24 +``` + +## Frontières préservées + +```text +aucun repair HTTP +aucun nouveau socket WebSocket +aucun second actor WebSocket +aucun retry Worker +aucune nouvelle dépendance +aucune nouvelle feature +aucun changement Config +aucun changement Store +aucun changement Job Backfill +aucun backend Store physique depuis Worker +aucune mutation Worker de Yellowstone from_slot +aucun SubscribeReplayInfo possédé par Worker +aucun élargissement de la façade publique Worker +src/lib.rs Worker reste exempt de ReplayInfo/from_slot/repair/backfill +``` + +## 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 ReplayInfo/from_slot/repair/backfill dans Worker src/lib.rs +scan de possession WebSocket directe dans Worker +scan set_from_slot dans les sources de production Worker +comparaison exacte 0.3.14-pre.002-fix.002 -> 0.3.14-pre.003 +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 +l'observabilité WS réutilise le watch latest-value déjà possédé par Transport +le Worker observe mais ne pilote pas la reconnexion WebSocket +un gap WS commence au dernier slot source réellement observé, inclusivement +le premier slot source post-incident est la borne de fin inclusive +aucune histoire antérieure n'est inventée sans ancre source +le premier événement post-incident est admis avant la terminalité conservative +pre.003 reste fail-closed et ne prétend pas encore réparer le gap +``` + +## Questions ouvertes + +```text +aucune pour pre.003 +``` + +## Tranche suivante + +`pre.004` reste dédiée à la preuve de replay natif Yellowstone : qualification de `SubscribeReplayInfo`, preuve de couverture effective et conservation stricte de l'ownership `from_slot` dans Transport. + +## 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 +```