diff --git a/Cargo.toml b/Cargo.toml index a0a6b93..9f1aca3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,12 @@ # file: Cargo.toml -# version: 250 +# version: 251 [workspace] resolver = "3" members = ["crates/ksp-app-config-desk", "crates/ksp-app-wallet-desk", "crates/ksp-config-lib", "crates/ksp-core-lib", "crates/ksp-logging-lib", "crates/ksp-onchain-transport-lib", "crates/ksp-wallet-lib"] [workspace.package] -version = "0.2.9-pre.9.fix.1" +version = "0.2.9-pre.10" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-onchain-transport-lib/src/grpc_stream.rs b/crates/ksp-onchain-transport-lib/src/grpc_stream.rs index 6ff6508..49fe019 100644 --- a/crates/ksp-onchain-transport-lib/src/grpc_stream.rs +++ b/crates/ksp-onchain-transport-lib/src/grpc_stream.rs @@ -1,9 +1,10 @@ // file: crates/ksp-onchain-transport-lib/src/grpc_stream.rs -// version: 1 +// version: 2 use tonic_prost::prost::Message; // rust-rules: trait-import const AUTO_SUBSCRIBE_PING_ID: i32 = 1; +const MAX_RECENT_UPDATE_IDENTITIES: usize = 512; const PATH_SUBSCRIBE: &str = "/geyser.Geyser/Subscribe"; /// Safe lifecycle state of one standard Yellowstone bidirectional subscribe session. @@ -11,25 +12,112 @@ const PATH_SUBSCRIBE: &str = "/geyser.Geyser/Subscribe"; pub enum YellowstoneGrpcSubscribeState { /// The bidirectional stream is active and accepts request mutations. Active, + /// The current stream ended and KSP is inside its bounded reconnect policy. + Reconnecting, /// KSP has started a bounded graceful half-close. Closing, - /// The stream ended normally, including a server half-close. + /// The stream ended normally, including an explicit close or a server half-close when reconnect is disabled. Closed, - /// The stream terminated because a transport, protocol, decoding or backpressure error occurred. + /// The stream terminated because a transport, protocol, decoding, backpressure or exhausted-reconnect error occurred. Failed, } +/// Safe continuity and lifecycle snapshot for one standard Yellowstone Subscribe session. +/// +/// Gap and duplicate counters are deliberately conservative. A continuity gap is counted only when `SubscribeReplayInfo` proves that the requested replay +/// slot is older than the endpoint's first retained slot. This proves unavailable replay coverage, not that a matching filtered update necessarily existed +/// or was lost. Duplicate updates are observed by bounded identity matching and are still delivered to the caller; KSP does not claim exactly-once delivery. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct YellowstoneGrpcSubscribeSnapshot { + state: crate::YellowstoneGrpcSubscribeState, + reconnect_count: u64, + continuity_gap_count: u64, + duplicate_update_count: u64, + replay_attempt_count: u64, + last_requested_from_slot: std::option::Option, + last_observed_slot: std::option::Option, + terminal_error_code: std::option::Option, +} + +impl YellowstoneGrpcSubscribeSnapshot { + const fn new(initial_from_slot: std::option::Option) -> Self { + return Self { + state: crate::YellowstoneGrpcSubscribeState::Active, + reconnect_count: 0, + continuity_gap_count: 0, + duplicate_update_count: 0, + replay_attempt_count: 0, + last_requested_from_slot: initial_from_slot, + last_observed_slot: std::option::Option::None, + terminal_error_code: std::option::Option::None, + }; + } + + /// Returns the current safe lifecycle state. + #[must_use] + pub const fn state(self) -> crate::YellowstoneGrpcSubscribeState { + return self.state; + } + + /// Returns the number of successful automatic stream reconnections. + #[must_use] + pub const fn reconnect_count(self) -> u64 { + return self.reconnect_count; + } + + /// Returns the number of replay discontinuities proven by `SubscribeReplayInfo` retention bounds. + #[must_use] + pub const fn continuity_gap_count(self) -> u64 { + return self.continuity_gap_count; + } + + /// Returns the number of bounded update identities observed more than once. + /// + /// Duplicate observations are not suppressed and this counter is not an exactly-once guarantee. + #[must_use] + pub const fn duplicate_update_count(self) -> u64 { + return self.duplicate_update_count; + } + + /// Returns the number of automatic reconnect attempts that requested a replay slot. + #[must_use] + pub const fn replay_attempt_count(self) -> u64 { + return self.replay_attempt_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 { + return self.last_requested_from_slot; + } + + /// Returns the highest slot observed from slot-bearing standard updates. + #[must_use] + pub const fn last_observed_slot(self) -> std::option::Option { + return self.last_observed_slot; + } + + /// Returns the safe terminal KSP error code when the session failed. + #[must_use] + pub const fn terminal_error_code(self) -> std::option::Option { + return self.terminal_error_code; + } +} + /// Standard Solana Yellowstone bidirectional `Subscribe` session layered on one KSP-owned physical gRPC channel. /// /// Request mutations and decoded updates use bounded Tokio channels. The raw Tonic stream and upstream protobuf messages remain private to Transport. +/// When the remote stream ends, KSP can reopen it with the bounded reconnect policy from [`crate::YellowstoneGrpcSessionSettings`]. The latest accepted +/// complete request is resubmitted deterministically and `from_slot` is advanced to at least the highest observed slot. Mutations are rejected while the +/// session is reconnecting so no request can be ambiguously applied to an old or replacement stream. pub struct SolanaYellowstoneGrpcSubscribeSession { endpoint_name: std::string::String, provider: crate::YellowstoneGrpcProviderName, cluster: crate::YellowstoneGrpcClusterName, - request_tx: std::option::Option>, + request_state: std::sync::Arc>, update_rx: tokio::sync::mpsc::Receiver>, shutdown_tx: tokio::sync::watch::Sender>, - state_rx: tokio::sync::watch::Receiver, + snapshot_rx: tokio::sync::watch::Receiver, task: tokio::task::JoinHandle<()>, close_timeout: std::time::Duration, max_outbound_message_size_bytes: usize, @@ -48,7 +136,7 @@ impl SolanaYellowstoneGrpcSubscribeSession { return &self.provider; } - /// Returns the open cluster descriptor. + /// Returns the open cluster descriptor without exposing endpoint credentials. #[must_use] pub const fn cluster(&self) -> &crate::YellowstoneGrpcClusterName { return &self.cluster; @@ -57,18 +145,25 @@ impl SolanaYellowstoneGrpcSubscribeSession { /// Returns the current safe stream lifecycle state. #[must_use] pub fn state(&self) -> crate::YellowstoneGrpcSubscribeState { - return self.state_rx.borrow().public_state(); + return self.snapshot_rx.borrow().state(); + } + + /// Returns the current safe reconnect/replay observability snapshot. + #[must_use] + pub fn snapshot(&self) -> crate::YellowstoneGrpcSubscribeSnapshot { + return *self.snapshot_rx.borrow(); } /// Queues one complete standard Yellowstone request mutation without waiting for network dispatch. /// /// The mutation is rejected synchronously when local validation fails, the encoded request exceeds the configured outbound bound, the bounded request - /// queue is full, or the session is no longer active. Accepted mutations preserve bounded-channel admission order. + /// queue is full, or the session is not currently active. In particular, mutations are rejected during reconnect so the request state used to reopen the + /// stream cannot race with a caller mutation. pub fn try_update(&self, request: &crate::YellowstoneSubscribeRequest) -> ksp_core_lib::Result<()> { if self.state() != crate::YellowstoneGrpcSubscribeState::Active { return std::result::Result::Err(subscribe_session_error( crate::ERROR_CODE_GRPC_SESSION_CLOSED, - "Yellowstone subscribe session is not active", + "Yellowstone subscribe session is not active for request mutation", self.endpoint_name.as_str(), self.provider.as_str(), self.cluster.as_str(), @@ -87,12 +182,24 @@ impl SolanaYellowstoneGrpcSubscribeSession { self.cluster.as_str(), )); } - let sender = match self.request_tx.as_ref() { + let mut shared = match self.request_state.lock() { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => { + return std::result::Result::Err(subscribe_session_error( + crate::ERROR_CODE_GRPC_SESSION_CLOSED, + "Yellowstone subscribe request state is unavailable", + self.endpoint_name.as_str(), + self.provider.as_str(), + self.cluster.as_str(), + )); + }, + }; + let sender = match shared.sender.as_ref() { std::option::Option::Some(value) => value, std::option::Option::None => { return std::result::Result::Err(subscribe_session_error( crate::ERROR_CODE_GRPC_SESSION_CLOSED, - "Yellowstone subscribe request channel is closed", + "Yellowstone subscribe request channel is unavailable during lifecycle transition", self.endpoint_name.as_str(), self.provider.as_str(), self.cluster.as_str(), @@ -100,7 +207,10 @@ impl SolanaYellowstoneGrpcSubscribeSession { }, }; return match sender.try_send(wire) { - std::result::Result::Ok(()) => std::result::Result::Ok(()), + std::result::Result::Ok(()) => { + shared.latest_request = request.clone(); + std::result::Result::Ok(()) + }, std::result::Result::Err(tokio::sync::mpsc::error::TrySendError::Full(_)) => std::result::Result::Err(subscribe_session_error( crate::ERROR_CODE_GRPC_BACKPRESSURE_OVERFLOW, "Yellowstone subscribe request queue is full", @@ -120,26 +230,33 @@ impl SolanaYellowstoneGrpcSubscribeSession { /// Receives the next decoded standard Yellowstone update. /// - /// A normal server half-close returns `Ok(None)`. Remote statuses, malformed updates and overflow terminate the session and surface a safe KSP error. + /// Transient remote stream loss is hidden while KSP performs a bounded reconnect. Terminal remote statuses after the reconnect budget, malformed updates + /// and overflow surface a safe KSP error. A normal close returns `Ok(None)`. pub async fn next_update(&mut self) -> ksp_core_lib::Result> { return match self.update_rx.recv().await { std::option::Option::Some(std::result::Result::Ok(update)) => std::result::Result::Ok(std::option::Option::Some(update)), std::option::Option::Some(std::result::Result::Err(error)) => { - self.request_tx.take(); + clear_request_sender(&self.request_state); std::result::Result::Err(error) }, std::option::Option::None => { - self.request_tx.take(); - match *self.state_rx.borrow() { - SubscribeActorState::Closed => std::result::Result::Ok(std::option::Option::None), - SubscribeActorState::Failed(code) => std::result::Result::Err(subscribe_session_error( - code, - terminal_message(code), - self.endpoint_name.as_str(), - self.provider.as_str(), - self.cluster.as_str(), - )), - SubscribeActorState::Active | SubscribeActorState::Closing => std::result::Result::Err(subscribe_session_error( + clear_request_sender(&self.request_state); + let snapshot = self.snapshot(); + match snapshot.state() { + crate::YellowstoneGrpcSubscribeState::Closed => std::result::Result::Ok(std::option::Option::None), + crate::YellowstoneGrpcSubscribeState::Failed => { + let code = snapshot.terminal_error_code().unwrap_or(crate::ERROR_CODE_GRPC_SESSION_CLOSED); + std::result::Result::Err(subscribe_session_error( + code, + terminal_message(code), + self.endpoint_name.as_str(), + self.provider.as_str(), + self.cluster.as_str(), + )) + }, + crate::YellowstoneGrpcSubscribeState::Active + | crate::YellowstoneGrpcSubscribeState::Reconnecting + | crate::YellowstoneGrpcSubscribeState::Closing => std::result::Result::Err(subscribe_session_error( crate::ERROR_CODE_GRPC_SESSION_CLOSED, "Yellowstone subscribe update channel ended before a terminal state was published", self.endpoint_name.as_str(), @@ -151,11 +268,13 @@ impl SolanaYellowstoneGrpcSubscribeSession { }; } - /// Gracefully half-closes the client request side and waits up to the configured close timeout for the server side to finish. + /// Gracefully closes the logical session and prevents any further reconnect attempt. pub async fn close(mut self) -> ksp_core_lib::Result<()> { - let terminal_before_close = *self.state_rx.borrow(); - self.request_tx.take(); - if terminal_before_close == SubscribeActorState::Active { + let terminal_before_close = self.snapshot(); + clear_request_sender(&self.request_state); + if terminal_before_close.state() == crate::YellowstoneGrpcSubscribeState::Active + || terminal_before_close.state() == crate::YellowstoneGrpcSubscribeState::Reconnecting + { let deadline = tokio::time::Instant::now() + self.close_timeout; self.shutdown_tx.send_replace(std::option::Option::Some(deadline)); } @@ -184,16 +303,22 @@ impl SolanaYellowstoneGrpcSubscribeSession { )); }, } - return match *self.state_rx.borrow() { - SubscribeActorState::Closed => std::result::Result::Ok(()), - SubscribeActorState::Failed(code) => std::result::Result::Err(subscribe_session_error( - code, - terminal_message(code), - self.endpoint_name.as_str(), - self.provider.as_str(), - self.cluster.as_str(), - )), - SubscribeActorState::Active | SubscribeActorState::Closing => std::result::Result::Err(subscribe_session_error( + let snapshot = self.snapshot(); + return match snapshot.state() { + crate::YellowstoneGrpcSubscribeState::Closed => std::result::Result::Ok(()), + crate::YellowstoneGrpcSubscribeState::Failed => { + let code = snapshot.terminal_error_code().unwrap_or(crate::ERROR_CODE_GRPC_SESSION_CLOSED); + std::result::Result::Err(subscribe_session_error( + code, + terminal_message(code), + self.endpoint_name.as_str(), + self.provider.as_str(), + self.cluster.as_str(), + )) + }, + crate::YellowstoneGrpcSubscribeState::Active + | crate::YellowstoneGrpcSubscribeState::Reconnecting + | crate::YellowstoneGrpcSubscribeState::Closing => std::result::Result::Err(subscribe_session_error( crate::ERROR_CODE_GRPC_SESSION_CLOSED, "Yellowstone subscribe actor ended without a terminal lifecycle state", self.endpoint_name.as_str(), @@ -206,13 +331,17 @@ impl SolanaYellowstoneGrpcSubscribeSession { impl std::fmt::Debug for SolanaYellowstoneGrpcSubscribeSession { fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let request_queue_capacity = match self.request_state.lock() { + std::result::Result::Ok(value) => value.sender.as_ref().map(|sender| return sender.capacity()), + std::result::Result::Err(_) => std::option::Option::None, + }; return formatter .debug_struct("SolanaYellowstoneGrpcSubscribeSession") .field("endpoint_name", &self.endpoint_name) .field("provider", &self.provider) .field("cluster", &self.cluster) - .field("state", &self.state()) - .field("request_queue_capacity", &self.request_tx.as_ref().map(|sender| return sender.capacity())) + .field("snapshot", &self.snapshot()) + .field("request_queue_capacity", &request_queue_capacity) .field("update_queue_capacity", &self.update_rx.capacity()) .finish(); } @@ -220,8 +349,9 @@ impl std::fmt::Debug for SolanaYellowstoneGrpcSubscribeSession { impl Drop for SolanaYellowstoneGrpcSubscribeSession { fn drop(&mut self) { - self.request_tx.take(); - if self.state() == crate::YellowstoneGrpcSubscribeState::Active { + clear_request_sender(&self.request_state); + let state = self.state(); + if state == crate::YellowstoneGrpcSubscribeState::Active || state == crate::YellowstoneGrpcSubscribeState::Reconnecting { self.shutdown_tx.send_replace(std::option::Option::Some(tokio::time::Instant::now() + self.close_timeout)); } } @@ -237,139 +367,70 @@ pub(crate) async fn open_yellowstone_subscribe_session( cluster: crate::YellowstoneGrpcClusterName, initial_request: crate::YellowstoneSubscribeRequest, ) -> ksp_core_lib::Result { - let initial_wire = match crate::yellowstone_subscribe_request_to_wire(&initial_request) { + let initial_wire = match checked_request_wire( + &initial_request, + settings.max_outbound_message_size_bytes(), + endpoint_name.as_str(), + provider.as_str(), + cluster.as_str(), + "initial Yellowstone subscribe request exceeds the configured outbound message bound", + ) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; - if initial_wire.encoded_len() > settings.max_outbound_message_size_bytes() { - return std::result::Result::Err(subscribe_session_error( - crate::ERROR_CODE_GRPC_BACKPRESSURE_OVERFLOW, - "initial Yellowstone subscribe request exceeds the configured outbound message bound", - endpoint_name.as_str(), - provider.as_str(), - cluster.as_str(), - )); - } - let (request_tx, request_rx) = tokio::sync::mpsc::channel(settings.request_channel_capacity()); - if request_tx.try_send(initial_wire).is_err() { - return std::result::Result::Err(subscribe_session_error( - crate::ERROR_CODE_GRPC_BACKPRESSURE_OVERFLOW, - "initial Yellowstone subscribe request cannot enter the bounded request queue", - endpoint_name.as_str(), - provider.as_str(), - cluster.as_str(), - )); - } - let path = match PATH_SUBSCRIBE.parse::() { - std::result::Result::Ok(value) => value, - std::result::Result::Err(_) => { - return std::result::Result::Err(subscribe_session_error( - crate::ERROR_CODE_GRPC_CHANNEL_FAILED, - "internal Yellowstone Subscribe method path is invalid", - endpoint_name.as_str(), - provider.as_str(), - cluster.as_str(), - )); - }, - }; - let mut grpc = tonic::client::Grpc::new(channel) - .max_decoding_message_size(settings.max_inbound_message_size_bytes()) - .max_encoding_message_size(settings.max_outbound_message_size_bytes()); - let open_future = async { - if grpc.ready().await.is_err() { - return std::result::Result::Err(subscribe_session_error( - crate::ERROR_CODE_GRPC_CHANNEL_FAILED, - "Yellowstone gRPC channel is not ready for Subscribe dispatch", - endpoint_name.as_str(), - provider.as_str(), - cluster.as_str(), - )); - } - let mut request = tonic::Request::new(MpscStream::new(request_rx)); - for entry in &metadata { - if let std::result::Result::Err(error) = entry.append_to(request.metadata_mut()) { - return std::result::Result::Err(error); - } - } - let response = match grpc - .streaming( - request, - path, - tonic_prost::ProstCodec::::default(), - ) + let (incoming, request_tx) = + match open_physical_subscribe_stream(channel.clone(), &metadata, &settings, endpoint_name.as_str(), provider.as_str(), cluster.as_str(), initial_wire) .await { std::result::Result::Ok(value) => value, - std::result::Result::Err(status) => { - return std::result::Result::Err(stream_status_error("SubscribeOpen", status, endpoint_name.as_str(), provider.as_str(), cluster.as_str())); - }, + std::result::Result::Err(error) => return std::result::Result::Err(error), }; - return std::result::Result::Ok(response.into_inner()); - }; - let incoming = match tokio::time::timeout(settings.connect_timeout(), open_future).await { - std::result::Result::Ok(std::result::Result::Ok(value)) => value, - std::result::Result::Ok(std::result::Result::Err(error)) => return std::result::Result::Err(error), - std::result::Result::Err(_) => { - return std::result::Result::Err(subscribe_session_error( - crate::ERROR_CODE_TIMEOUT, - "Yellowstone Subscribe stream opening exceeded the configured connection timeout", - endpoint_name.as_str(), - provider.as_str(), - cluster.as_str(), - )); - }, - }; + let request_state = std::sync::Arc::new(std::sync::Mutex::new(SharedRequestState { + sender: std::option::Option::Some(request_tx), + latest_request: initial_request.clone(), + })); let (update_tx, update_rx) = tokio::sync::mpsc::channel(settings.update_channel_capacity()); let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(std::option::Option::None::); - let (state_tx, state_rx) = tokio::sync::watch::channel(SubscribeActorState::Active); - let actor_request_tx = request_tx.clone(); + let initial_snapshot = crate::YellowstoneGrpcSubscribeSnapshot::new(initial_request.from_slot()); + let (snapshot_tx, snapshot_rx) = tokio::sync::watch::channel(initial_snapshot); + let actor_request_state = request_state.clone(); let actor_endpoint_name = endpoint_name.clone(); let actor_provider = provider.clone(); let actor_cluster = cluster.clone(); + let actor_settings = settings.clone(); let close_timeout = settings.close_timeout(); let max_outbound_message_size_bytes = settings.max_outbound_message_size_bytes(); let task = tokio::spawn(run_subscribe_actor( actor_endpoint_name, actor_provider, actor_cluster, + channel, + metadata, + actor_settings, incoming, - actor_request_tx, + actor_request_state, update_tx, shutdown_rx, - state_tx, - max_outbound_message_size_bytes, + snapshot_tx, + initial_snapshot, )); return std::result::Result::Ok(crate::SolanaYellowstoneGrpcSubscribeSession { endpoint_name, provider, cluster, - request_tx: std::option::Option::Some(request_tx), + request_state, update_rx, shutdown_tx, - state_rx, + snapshot_rx, task, close_timeout, max_outbound_message_size_bytes, }); } -#[derive(Clone, Copy, Debug, Eq, PartialEq)] -enum SubscribeActorState { - Active, - Closing, - Closed, - Failed(ksp_core_lib::ErrorCode), -} - -impl SubscribeActorState { - const fn public_state(self) -> crate::YellowstoneGrpcSubscribeState { - return match self { - Self::Active => crate::YellowstoneGrpcSubscribeState::Active, - Self::Closing => crate::YellowstoneGrpcSubscribeState::Closing, - Self::Closed => crate::YellowstoneGrpcSubscribeState::Closed, - Self::Failed(_) => crate::YellowstoneGrpcSubscribeState::Failed, - }; - } +struct SharedRequestState { + sender: std::option::Option>, + latest_request: crate::YellowstoneSubscribeRequest, } struct MpscStream { @@ -390,26 +451,76 @@ impl futures_util::Stream for MpscStream { } } +#[derive(Clone, Eq, Hash, PartialEq)] +enum UpdateIdentity { + Account { slot: u64, pubkey: ksp_core_lib::Pubkey, write_version: u64 }, + Slot { slot: u64, status: crate::YellowstoneSlotStatus }, + Transaction { slot: u64, signature: crate::YellowstoneTransactionSignature }, + TransactionStatus { slot: u64, signature: crate::YellowstoneTransactionSignature }, + Block { slot: u64, blockhash: std::string::String }, + BlockMeta { slot: u64, blockhash: std::string::String }, + Entry { slot: u64, index: u64, hash: crate::YellowstoneHashBytes }, +} + +struct ContinuityTracker { + 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() }; + } + + fn observe(&mut self, update: &crate::YellowstoneSubscribeUpdate, snapshot: &mut crate::YellowstoneGrpcSubscribeSnapshot) { + if let std::option::Option::Some(slot) = update_slot(update) { + 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, + }); + } + let identity = match update_identity(update) { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + if snapshot.reconnect_count > 0 && self.recent_set.contains(&identity) { + snapshot.duplicate_update_count = snapshot.duplicate_update_count.saturating_add(1); + } + if self.recent_set.insert(identity.clone()) { + self.recent_order.push_back(identity); + if self.recent_order.len() > MAX_RECENT_UPDATE_IDENTITIES { + if let std::option::Option::Some(oldest) = self.recent_order.pop_front() { + self.recent_set.remove(&oldest); + } + } + } + } +} + #[allow(clippy::too_many_arguments)] async fn run_subscribe_actor( endpoint_name: std::string::String, provider: crate::YellowstoneGrpcProviderName, cluster: crate::YellowstoneGrpcClusterName, + channel: tonic::transport::Channel, + metadata: std::vec::Vec, + settings: crate::YellowstoneGrpcSessionSettings, mut incoming: tonic::Streaming, - request_tx: tokio::sync::mpsc::Sender, + request_state: std::sync::Arc>, update_tx: tokio::sync::mpsc::Sender>, mut shutdown_rx: tokio::sync::watch::Receiver>, - state_tx: tokio::sync::watch::Sender, - max_outbound_message_size_bytes: usize, + snapshot_tx: tokio::sync::watch::Sender, + mut snapshot: crate::YellowstoneGrpcSubscribeSnapshot, ) { - let mut request_tx = std::option::Option::Some(request_tx); + let mut tracker = ContinuityTracker::new(); loop { tokio::select! { shutdown_changed = shutdown_rx.changed() => { let deadline = shutdown_deadline(&shutdown_rx, shutdown_changed); - state_tx.send_replace(SubscribeActorState::Closing); - request_tx.take(); - finish_client_half_close(&endpoint_name, &provider, &cluster, &mut incoming, deadline, &state_tx).await; + snapshot.state = crate::YellowstoneGrpcSubscribeState::Closing; + snapshot_tx.send_replace(snapshot); + clear_request_sender(&request_state); + finish_client_half_close(&endpoint_name, &provider, &cluster, &mut incoming, deadline, &mut snapshot, &snapshot_tx).await; return; } message = incoming.message() => { @@ -419,31 +530,18 @@ async fn run_subscribe_actor( std::result::Result::Ok(value) => value, std::result::Result::Err(error) => { let _ = update_tx.try_send(std::result::Result::Err(error)); - state_tx.send_replace(SubscribeActorState::Failed(crate::ERROR_CODE_INVALID_RESPONSE)); - request_tx.take(); + fail_actor(crate::ERROR_CODE_INVALID_RESPONSE, &request_state, &mut snapshot, &snapshot_tx); return; }, }; if matches!(&update, crate::YellowstoneSubscribeUpdate::Ping(_)) { - let ping_wire = ping_request_wire(); - if ping_wire.encoded_len() > max_outbound_message_size_bytes { - state_tx.send_replace(SubscribeActorState::Failed(crate::ERROR_CODE_GRPC_BACKPRESSURE_OVERFLOW)); - request_tx.take(); - return; - } - let sender = match request_tx.as_ref() { - std::option::Option::Some(value) => value, - std::option::Option::None => { - state_tx.send_replace(SubscribeActorState::Failed(crate::ERROR_CODE_GRPC_SESSION_CLOSED)); - return; - }, - }; - if sender.try_send(ping_wire).is_err() { - state_tx.send_replace(SubscribeActorState::Failed(crate::ERROR_CODE_GRPC_BACKPRESSURE_OVERFLOW)); - request_tx.take(); + if let std::result::Result::Err(code) = send_automatic_ping(&request_state, settings.max_outbound_message_size_bytes()) { + fail_actor(code, &request_state, &mut snapshot, &snapshot_tx); return; } } + tracker.observe(&update, &mut snapshot); + snapshot_tx.send_replace(snapshot); match update_tx.try_send(std::result::Result::Ok(update)) { std::result::Result::Ok(()) => {}, std::result::Result::Err(tokio::sync::mpsc::error::TrySendError::Full(_)) => { @@ -454,28 +552,87 @@ async fn run_subscribe_actor( cluster = cluster.as_str(), "Yellowstone subscribe update queue overflowed" ); - state_tx.send_replace(SubscribeActorState::Failed(crate::ERROR_CODE_GRPC_BACKPRESSURE_OVERFLOW)); - request_tx.take(); + fail_actor(crate::ERROR_CODE_GRPC_BACKPRESSURE_OVERFLOW, &request_state, &mut snapshot, &snapshot_tx); return; }, std::result::Result::Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) => { - state_tx.send_replace(SubscribeActorState::Closed); - request_tx.take(); + close_actor(&request_state, &mut snapshot, &snapshot_tx); return; }, } }, std::result::Result::Ok(std::option::Option::None) => { - state_tx.send_replace(SubscribeActorState::Closed); - request_tx.take(); - return; + if settings.reconnect().max_retries() == 0 { + close_actor(&request_state, &mut snapshot, &snapshot_tx); + return; + } + match reconnect_subscribe_stream( + &endpoint_name, + &provider, + &cluster, + channel.clone(), + &metadata, + &settings, + &request_state, + &mut shutdown_rx, + &mut snapshot, + &snapshot_tx, + ) + .await + { + ReconnectOutcome::Connected(value) => incoming = value, + ReconnectOutcome::Shutdown => { + close_actor(&request_state, &mut snapshot, &snapshot_tx); + return; + }, + ReconnectOutcome::Exhausted(error) => { + let _ = update_tx.try_send(std::result::Result::Err(error)); + fail_actor(crate::ERROR_CODE_GRPC_CHANNEL_FAILED, &request_state, &mut snapshot, &snapshot_tx); + return; + }, + } }, std::result::Result::Err(status) => { - let error = stream_status_error("Subscribe", status, endpoint_name.as_str(), provider.as_str(), cluster.as_str()); - let _ = update_tx.try_send(std::result::Result::Err(error)); - state_tx.send_replace(SubscribeActorState::Failed(crate::ERROR_CODE_GRPC_STATUS)); - request_tx.take(); - return; + if settings.reconnect().max_retries() == 0 { + let error = stream_status_error("Subscribe", status, endpoint_name.as_str(), provider.as_str(), cluster.as_str()); + let _ = update_tx.try_send(std::result::Result::Err(error)); + fail_actor(crate::ERROR_CODE_GRPC_STATUS, &request_state, &mut snapshot, &snapshot_tx); + return; + } + let grpc_code = status.code().to_string(); + ksp_logging_lib::warn!( + target: crate::TRACING_TARGET, + endpoint_name = endpoint_name.as_str(), + provider = provider.as_str(), + cluster = cluster.as_str(), + grpc_code = grpc_code.as_str(), + "Yellowstone subscribe stream ended with remote status; starting bounded reconnect" + ); + match reconnect_subscribe_stream( + &endpoint_name, + &provider, + &cluster, + channel.clone(), + &metadata, + &settings, + &request_state, + &mut shutdown_rx, + &mut snapshot, + &snapshot_tx, + ) + .await + { + ReconnectOutcome::Connected(value) => incoming = value, + ReconnectOutcome::Shutdown => { + close_actor(&request_state, &mut snapshot, &snapshot_tx); + return; + }, + ReconnectOutcome::Exhausted(error) => { + let _ = update_tx.try_send(std::result::Result::Err(error)); + fail_actor(crate::ERROR_CODE_GRPC_CHANNEL_FAILED, &request_state, &mut snapshot, &snapshot_tx); + return; + }, + } }, } } @@ -483,20 +640,266 @@ async fn run_subscribe_actor( } } +enum ReconnectOutcome { + Connected(tonic::Streaming), + Shutdown, + Exhausted(ksp_core_lib::Error), +} + +#[allow(clippy::too_many_arguments)] +async fn reconnect_subscribe_stream( + endpoint_name: &str, + provider: &crate::YellowstoneGrpcProviderName, + cluster: &crate::YellowstoneGrpcClusterName, + channel: tonic::transport::Channel, + metadata: &[crate::YellowstoneGrpcMetadataEntry], + settings: &crate::YellowstoneGrpcSessionSettings, + request_state: &std::sync::Arc>, + shutdown_rx: &mut tokio::sync::watch::Receiver>, + snapshot: &mut crate::YellowstoneGrpcSubscribeSnapshot, + snapshot_tx: &tokio::sync::watch::Sender, +) -> ReconnectOutcome { + clear_request_sender(request_state); + snapshot.state = crate::YellowstoneGrpcSubscribeState::Reconnecting; + snapshot.terminal_error_code = std::option::Option::None; + snapshot_tx.send_replace(*snapshot); + let mut backoff = settings.reconnect().initial_backoff(); + let mut gap_recorded = false; + for attempt in 0..settings.reconnect().max_retries() { + let sleep = tokio::time::sleep(backoff); + tokio::pin!(sleep); + tokio::select! { + _ = &mut sleep => {}, + changed = shutdown_rx.changed() => { + let _ = changed; + return ReconnectOutcome::Shutdown; + } + } + let latest_request = match clone_latest_request(request_state) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return ReconnectOutcome::Exhausted(error), + }; + let resume_slot = max_optional_slot(latest_request.from_slot(), snapshot.last_observed_slot); + let mut reconnect_request = latest_request; + let effective_from_slot = match resume_slot { + std::option::Option::Some(requested) => { + let first_available = replay_first_available(channel.clone(), metadata, settings.clone()).await; + match first_available { + std::option::Option::Some(first_available) if first_available > requested => { + if !gap_recorded { + snapshot.continuity_gap_count = snapshot.continuity_gap_count.saturating_add(1); + gap_recorded = true; + ksp_logging_lib::warn!( + target: crate::TRACING_TARGET, + endpoint_name, + provider = provider.as_str(), + cluster = cluster.as_str(), + requested_from_slot = requested, + first_available, + "Yellowstone replay retention proves unavailable replay coverage; clamping reconnect from_slot" + ); + } + std::option::Option::Some(first_available) + }, + _ => std::option::Option::Some(requested), + } + }, + std::option::Option::None => std::option::Option::None, + }; + reconnect_request.set_from_slot(effective_from_slot); + snapshot.last_requested_from_slot = effective_from_slot; + if effective_from_slot.is_some() { + snapshot.replay_attempt_count = snapshot.replay_attempt_count.saturating_add(1); + } + snapshot_tx.send_replace(*snapshot); + let wire = match checked_request_wire( + &reconnect_request, + settings.max_outbound_message_size_bytes(), + endpoint_name, + provider.as_str(), + cluster.as_str(), + "Yellowstone reconnect subscribe request exceeds the configured outbound message bound", + ) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return ReconnectOutcome::Exhausted(error), + }; + let open = open_physical_subscribe_stream(channel.clone(), metadata, settings, endpoint_name, provider.as_str(), cluster.as_str(), wire); + tokio::pin!(open); + let result = tokio::select! { + value = &mut open => std::option::Option::Some(value), + changed = shutdown_rx.changed() => { + let _ = changed; + std::option::Option::None + } + }; + let result = match result { + std::option::Option::Some(value) => value, + std::option::Option::None => return ReconnectOutcome::Shutdown, + }; + match result { + std::result::Result::Ok((incoming, request_tx)) => { + if let std::result::Result::Err(error) = replace_request_sender(request_state, request_tx) { + return ReconnectOutcome::Exhausted(error); + } + snapshot.reconnect_count = snapshot.reconnect_count.saturating_add(1); + snapshot.state = crate::YellowstoneGrpcSubscribeState::Active; + snapshot.terminal_error_code = std::option::Option::None; + snapshot_tx.send_replace(*snapshot); + ksp_logging_lib::info!( + target: crate::TRACING_TARGET, + endpoint_name, + provider = provider.as_str(), + cluster = cluster.as_str(), + reconnect_attempt = attempt.saturating_add(1), + replay = effective_from_slot.is_some(), + "reopened Yellowstone subscribe stream within bounded reconnect policy" + ); + return ReconnectOutcome::Connected(incoming); + }, + std::result::Result::Err(error) => { + ksp_logging_lib::warn!( + target: crate::TRACING_TARGET, + endpoint_name, + provider = provider.as_str(), + cluster = cluster.as_str(), + reconnect_attempt = attempt.saturating_add(1), + error_code = error.code().code(), + "Yellowstone subscribe reconnect attempt failed" + ); + }, + } + backoff = std::cmp::min(backoff.saturating_mul(2), settings.reconnect().max_backoff()); + } + return ReconnectOutcome::Exhausted(subscribe_session_error( + crate::ERROR_CODE_GRPC_CHANNEL_FAILED, + "Yellowstone subscribe reconnect budget is exhausted", + endpoint_name, + provider.as_str(), + cluster.as_str(), + )); +} + +async fn replay_first_available( + channel: tonic::transport::Channel, + metadata: &[crate::YellowstoneGrpcMetadataEntry], + settings: crate::YellowstoneGrpcSessionSettings, +) -> std::option::Option { + let unary = crate::SolanaYellowstoneGrpcUnaryClient::new(channel, metadata.to_vec(), settings); + return match unary.subscribe_replay_info().await { + std::result::Result::Ok(info) => info.first_available(), + std::result::Result::Err(error) => { + ksp_logging_lib::warn!( + target: crate::TRACING_TARGET, + error_code = error.code().code(), + "Yellowstone SubscribeReplayInfo is unavailable during reconnect; continuing without inferred retention clamp" + ); + std::option::Option::None + }, + }; +} + +#[allow(clippy::too_many_arguments)] +async fn open_physical_subscribe_stream( + channel: tonic::transport::Channel, + metadata: &[crate::YellowstoneGrpcMetadataEntry], + settings: &crate::YellowstoneGrpcSessionSettings, + endpoint_name: &str, + provider: &str, + cluster: &str, + initial_wire: yellowstone_grpc_proto::geyser::SubscribeRequest, +) -> ksp_core_lib::Result<( + tonic::Streaming, + tokio::sync::mpsc::Sender, +)> { + let (request_tx, request_rx) = tokio::sync::mpsc::channel(settings.request_channel_capacity()); + if request_tx.try_send(initial_wire).is_err() { + return std::result::Result::Err(subscribe_session_error( + crate::ERROR_CODE_GRPC_BACKPRESSURE_OVERFLOW, + "initial Yellowstone subscribe request cannot enter the bounded request queue", + endpoint_name, + provider, + cluster, + )); + } + let path = match PATH_SUBSCRIBE.parse::() { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => { + return std::result::Result::Err(subscribe_session_error( + crate::ERROR_CODE_GRPC_CHANNEL_FAILED, + "internal Yellowstone Subscribe method path is invalid", + endpoint_name, + provider, + cluster, + )); + }, + }; + let mut grpc = tonic::client::Grpc::new(channel) + .max_decoding_message_size(settings.max_inbound_message_size_bytes()) + .max_encoding_message_size(settings.max_outbound_message_size_bytes()); + let open_future = async { + if grpc.ready().await.is_err() { + return std::result::Result::Err(subscribe_session_error( + crate::ERROR_CODE_GRPC_CHANNEL_FAILED, + "Yellowstone gRPC channel is not ready for Subscribe dispatch", + endpoint_name, + provider, + cluster, + )); + } + let mut request = tonic::Request::new(MpscStream::new(request_rx)); + for entry in metadata { + if let std::result::Result::Err(error) = entry.append_to(request.metadata_mut()) { + return std::result::Result::Err(error); + } + } + let response = match grpc + .streaming( + request, + path, + tonic_prost::ProstCodec::::default(), + ) + .await + { + std::result::Result::Ok(value) => value, + std::result::Result::Err(status) => { + return std::result::Result::Err(stream_status_error("SubscribeOpen", status, endpoint_name, provider, cluster)); + }, + }; + return std::result::Result::Ok(response.into_inner()); + }; + let incoming = match tokio::time::timeout(settings.connect_timeout(), open_future).await { + std::result::Result::Ok(std::result::Result::Ok(value)) => value, + std::result::Result::Ok(std::result::Result::Err(error)) => return std::result::Result::Err(error), + std::result::Result::Err(_) => { + return std::result::Result::Err(subscribe_session_error( + crate::ERROR_CODE_TIMEOUT, + "Yellowstone Subscribe stream opening exceeded the configured connection timeout", + endpoint_name, + provider, + cluster, + )); + }, + }; + return std::result::Result::Ok((incoming, request_tx)); +} + async fn finish_client_half_close( endpoint_name: &str, provider: &crate::YellowstoneGrpcProviderName, cluster: &crate::YellowstoneGrpcClusterName, incoming: &mut tonic::Streaming, deadline: tokio::time::Instant, - state_tx: &tokio::sync::watch::Sender, + snapshot: &mut crate::YellowstoneGrpcSubscribeSnapshot, + snapshot_tx: &tokio::sync::watch::Sender, ) { loop { let message = tokio::time::timeout_at(deadline, incoming.message()).await; match message { std::result::Result::Ok(std::result::Result::Ok(std::option::Option::Some(_))) => {}, std::result::Result::Ok(std::result::Result::Ok(std::option::Option::None)) => { - state_tx.send_replace(SubscribeActorState::Closed); + snapshot.state = crate::YellowstoneGrpcSubscribeState::Closed; + snapshot.terminal_error_code = std::option::Option::None; + snapshot_tx.send_replace(*snapshot); return; }, std::result::Result::Ok(std::result::Result::Err(status)) => { @@ -509,17 +912,173 @@ async fn finish_client_half_close( grpc_code = code.as_str(), "Yellowstone subscribe endpoint returned a status during graceful shutdown" ); - state_tx.send_replace(SubscribeActorState::Failed(crate::ERROR_CODE_GRPC_STATUS)); + snapshot.state = crate::YellowstoneGrpcSubscribeState::Failed; + snapshot.terminal_error_code = std::option::Option::Some(crate::ERROR_CODE_GRPC_STATUS); + snapshot_tx.send_replace(*snapshot); return; }, std::result::Result::Err(_) => { - state_tx.send_replace(SubscribeActorState::Failed(crate::ERROR_CODE_TIMEOUT)); + snapshot.state = crate::YellowstoneGrpcSubscribeState::Failed; + snapshot.terminal_error_code = std::option::Option::Some(crate::ERROR_CODE_TIMEOUT); + snapshot_tx.send_replace(*snapshot); return; }, } } } +fn checked_request_wire( + request: &crate::YellowstoneSubscribeRequest, + max_outbound_message_size_bytes: usize, + endpoint_name: &str, + provider: &str, + cluster: &str, + oversized_message: &'static str, +) -> ksp_core_lib::Result { + let wire = match crate::yellowstone_subscribe_request_to_wire(request) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + if wire.encoded_len() > max_outbound_message_size_bytes { + return std::result::Result::Err(subscribe_session_error( + crate::ERROR_CODE_GRPC_BACKPRESSURE_OVERFLOW, + oversized_message, + endpoint_name, + provider, + cluster, + )); + } + return std::result::Result::Ok(wire); +} + +fn clone_latest_request(request_state: &std::sync::Arc>) -> ksp_core_lib::Result { + return match request_state.lock() { + std::result::Result::Ok(value) => std::result::Result::Ok(value.latest_request.clone()), + std::result::Result::Err(_) => std::result::Result::Err(ksp_core_lib::Error::new( + crate::ERROR_CODE_GRPC_SESSION_CLOSED, + "Yellowstone subscribe request state is unavailable during reconnect", + )), + }; +} + +fn replace_request_sender( + request_state: &std::sync::Arc>, + sender: tokio::sync::mpsc::Sender, +) -> ksp_core_lib::Result<()> { + return match request_state.lock() { + std::result::Result::Ok(mut value) => { + value.sender = std::option::Option::Some(sender); + std::result::Result::Ok(()) + }, + std::result::Result::Err(_) => std::result::Result::Err(ksp_core_lib::Error::new( + crate::ERROR_CODE_GRPC_SESSION_CLOSED, + "Yellowstone subscribe request state is unavailable after reconnect", + )), + }; +} + +fn clear_request_sender(request_state: &std::sync::Arc>) { + if let std::result::Result::Ok(mut value) = request_state.lock() { + value.sender.take(); + } +} + +fn send_automatic_ping( + request_state: &std::sync::Arc>, + max_outbound_message_size_bytes: usize, +) -> std::result::Result<(), ksp_core_lib::ErrorCode> { + let ping_wire = ping_request_wire(); + if ping_wire.encoded_len() > max_outbound_message_size_bytes { + return std::result::Result::Err(crate::ERROR_CODE_GRPC_BACKPRESSURE_OVERFLOW); + } + let shared = match request_state.lock() { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(crate::ERROR_CODE_GRPC_SESSION_CLOSED), + }; + let sender = match shared.sender.as_ref() { + std::option::Option::Some(value) => value, + std::option::Option::None => return std::result::Result::Err(crate::ERROR_CODE_GRPC_SESSION_CLOSED), + }; + return match sender.try_send(ping_wire) { + std::result::Result::Ok(()) => std::result::Result::Ok(()), + std::result::Result::Err(_) => std::result::Result::Err(crate::ERROR_CODE_GRPC_BACKPRESSURE_OVERFLOW), + }; +} + +fn close_actor( + request_state: &std::sync::Arc>, + snapshot: &mut crate::YellowstoneGrpcSubscribeSnapshot, + snapshot_tx: &tokio::sync::watch::Sender, +) { + clear_request_sender(request_state); + snapshot.state = crate::YellowstoneGrpcSubscribeState::Closed; + snapshot.terminal_error_code = std::option::Option::None; + snapshot_tx.send_replace(*snapshot); +} + +fn fail_actor( + code: ksp_core_lib::ErrorCode, + request_state: &std::sync::Arc>, + snapshot: &mut crate::YellowstoneGrpcSubscribeSnapshot, + snapshot_tx: &tokio::sync::watch::Sender, +) { + clear_request_sender(request_state); + snapshot.state = crate::YellowstoneGrpcSubscribeState::Failed; + snapshot.terminal_error_code = std::option::Option::Some(code); + snapshot_tx.send_replace(*snapshot); +} + +fn update_slot(update: &crate::YellowstoneSubscribeUpdate) -> std::option::Option { + return match update { + crate::YellowstoneSubscribeUpdate::Account(value) => std::option::Option::Some(value.slot()), + crate::YellowstoneSubscribeUpdate::Slot(value) => std::option::Option::Some(value.slot()), + crate::YellowstoneSubscribeUpdate::Transaction(value) => std::option::Option::Some(value.slot()), + crate::YellowstoneSubscribeUpdate::TransactionStatus(value) => std::option::Option::Some(value.slot()), + crate::YellowstoneSubscribeUpdate::Block(value) => std::option::Option::Some(value.slot()), + crate::YellowstoneSubscribeUpdate::BlockMeta(value) => std::option::Option::Some(value.slot()), + crate::YellowstoneSubscribeUpdate::Entry(value) => std::option::Option::Some(value.entry().slot()), + crate::YellowstoneSubscribeUpdate::Ping(_) | crate::YellowstoneSubscribeUpdate::Pong(_) => std::option::Option::None, + }; +} + +fn update_identity(update: &crate::YellowstoneSubscribeUpdate) -> std::option::Option { + return match update { + crate::YellowstoneSubscribeUpdate::Account(value) => std::option::Option::Some(UpdateIdentity::Account { + slot: value.slot(), + pubkey: *value.account().pubkey(), + write_version: value.account().write_version(), + }), + crate::YellowstoneSubscribeUpdate::Slot(value) => std::option::Option::Some(UpdateIdentity::Slot { slot: value.slot(), status: value.status() }), + crate::YellowstoneSubscribeUpdate::Transaction(value) => { + std::option::Option::Some(UpdateIdentity::Transaction { slot: value.slot(), signature: value.transaction().signature() }) + }, + crate::YellowstoneSubscribeUpdate::TransactionStatus(value) => { + std::option::Option::Some(UpdateIdentity::TransactionStatus { slot: value.slot(), signature: value.signature() }) + }, + crate::YellowstoneSubscribeUpdate::Block(value) => { + std::option::Option::Some(UpdateIdentity::Block { slot: value.slot(), blockhash: value.blockhash().to_owned() }) + }, + crate::YellowstoneSubscribeUpdate::BlockMeta(value) => { + std::option::Option::Some(UpdateIdentity::BlockMeta { slot: value.slot(), blockhash: value.blockhash().to_owned() }) + }, + crate::YellowstoneSubscribeUpdate::Entry(value) => { + let entry = value.entry(); + std::option::Option::Some(UpdateIdentity::Entry { slot: entry.slot(), index: entry.index(), hash: entry.hash() }) + }, + crate::YellowstoneSubscribeUpdate::Ping(_) | crate::YellowstoneSubscribeUpdate::Pong(_) => std::option::Option::None, + }; +} + +fn max_optional_slot(left: std::option::Option, right: std::option::Option) -> std::option::Option { + return match (left, right) { + (std::option::Option::Some(left), std::option::Option::Some(right)) => std::option::Option::Some(std::cmp::max(left, right)), + (std::option::Option::Some(value), std::option::Option::None) | (std::option::Option::None, std::option::Option::Some(value)) => { + std::option::Option::Some(value) + }, + (std::option::Option::None, std::option::Option::None) => std::option::Option::None, + }; +} + fn shutdown_deadline( shutdown_rx: &tokio::sync::watch::Receiver>, shutdown_changed: std::result::Result<(), tokio::sync::watch::error::RecvError>, @@ -570,6 +1129,9 @@ fn terminal_message(code: ksp_core_lib::ErrorCode) -> &'static str { if code == crate::ERROR_CODE_GRPC_STATUS { return "Yellowstone subscribe session terminated after a remote gRPC status"; } + if code == crate::ERROR_CODE_GRPC_CHANNEL_FAILED { + return "Yellowstone subscribe session exhausted its bounded reconnect policy"; + } if code == crate::ERROR_CODE_INVALID_RESPONSE { return "Yellowstone subscribe session terminated after an invalid update"; } diff --git a/crates/ksp-onchain-transport-lib/src/lib.rs b/crates/ksp-onchain-transport-lib/src/lib.rs index 5bb88dc..2470882 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: 42 +// version: 43 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -140,6 +140,8 @@ pub use self::grpc_settings::YellowstoneGrpcSessionSettings; pub use self::grpc_settings::YellowstoneGrpcTransportSettings; /// Standard Yellowstone bidirectional Subscribe session. pub use self::grpc_stream::SolanaYellowstoneGrpcSubscribeSession; +/// Safe Yellowstone reconnect/replay continuity snapshot. +pub use self::grpc_stream::YellowstoneGrpcSubscribeSnapshot; /// Safe Yellowstone bidirectional Subscribe lifecycle state. pub use self::grpc_stream::YellowstoneGrpcSubscribeState; /// One validated standard Yellowstone account predicate. diff --git a/crates/ksp-onchain-transport-lib/tests/public_api.rs b/crates/ksp-onchain-transport-lib/tests/public_api.rs index f3aac48..b0381b1 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: 47 +// version: 48 //! Integration tests for the public `ksp-onchain-transport-lib` consumer contract. @@ -993,3 +993,22 @@ fn public_v0_2_9_pre_009_yellowstone_bidi_session_contract_is_available_from_cra assert_eq!(ksp_onchain_transport_lib::ERROR_CODE_GRPC_SESSION_CLOSED.domain(), "onchain_transport"); assert_eq!(ksp_onchain_transport_lib::ERROR_CODE_GRPC_SESSION_CLOSED.code(), "grpc_session_closed"); } + +#[test] +fn public_v0_2_9_pre_010_yellowstone_reconnect_snapshot_is_available_from_crate_root() { + fn assert_copy() { + let _ = std::marker::PhantomData::; + return; + } + assert_copy::(); + let _snapshot = std::any::type_name::(); + let _state = ksp_onchain_transport_lib::YellowstoneGrpcSubscribeState::Reconnecting; + let _session_snapshot = ksp_onchain_transport_lib::SolanaYellowstoneGrpcSubscribeSession::snapshot; + let _reconnect_count = ksp_onchain_transport_lib::YellowstoneGrpcSubscribeSnapshot::reconnect_count; + 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 _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 83779e4..8b4a379 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: 40 +// version: 41 //! Release-level completeness canaries for staged HTTP and WebSocket Transport coverage. @@ -1264,7 +1264,7 @@ fn release_v0_2_9_pre_008_standard_blocks_contract_remains_complete() { } #[test] -fn release_v0_2_9_pre_009_opens_one_bounded_standard_bidi_session_without_reconnect_or_provider_coupling() { +fn release_v0_2_9_pre_009_bounded_standard_bidi_session_contract_remains_complete() { let stream_source = include_str!("../src/grpc_stream.rs"); let subscribe_source = include_str!("../src/grpc_subscribe.rs"); let channel_source = include_str!("../src/grpc_channel.rs"); @@ -1301,10 +1301,46 @@ fn release_v0_2_9_pre_009_opens_one_bounded_standard_bidi_session_without_reconn assert!(!stream_source.contains("PublicNode")); assert!(!stream_source.contains("OrbitFlare")); assert!(!stream_source.contains("Helius")); - assert!(!stream_source.contains("reconnect")); assert!(!crate_root.contains("pub use tonic")); assert!(!crate_root.contains("pub use yellowstone_grpc_proto")); let _session = std::any::type_name::(); let _state = ksp_onchain_transport_lib::YellowstoneGrpcSubscribeState::Active; let _update = std::any::type_name::(); } + +#[test] +fn release_v0_2_9_pre_010_adds_bounded_reconnect_replay_and_conservative_continuity_observability() { + let stream_source = include_str!("../src/grpc_stream.rs"); + let crate_root = include_str!("../src/lib.rs"); + for required in [ + "YellowstoneGrpcSubscribeState::Reconnecting", + "YellowstoneGrpcSubscribeSnapshot", + "reconnect_subscribe_stream", + "replay_first_available", + "SubscribeReplayInfo", + "last_requested_from_slot", + "last_observed_slot", + "reconnect_count", + "replay_attempt_count", + "continuity_gap_count", + "duplicate_update_count", + "MAX_RECENT_UPDATE_IDENTITIES", + "first_available > requested", + "max_optional_slot", + "reconnect budget is exhausted", + "request mutation", + ] { + assert!(stream_source.contains(required), "missing pre.010 reconnect/replay contract token: {required}"); + } + assert!(stream_source.contains("still delivered to the caller")); + assert!(stream_source.contains("does not claim exactly-once delivery")); + assert!(!stream_source.contains("unbounded_channel")); + assert!(!stream_source.contains("SubscribeDeshred")); + assert!(!stream_source.contains("PublicNode")); + assert!(!stream_source.contains("OrbitFlare")); + assert!(!stream_source.contains("Helius")); + assert!(!crate_root.contains("pub use tonic")); + assert!(!crate_root.contains("pub use yellowstone_grpc_proto")); + let _snapshot = std::any::type_name::(); + let _state = ksp_onchain_transport_lib::YellowstoneGrpcSubscribeState::Reconnecting; +} 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 4685914..016bb8d 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: 1 +// version: 2 #[derive(Clone, Copy)] enum FixtureMode { @@ -11,12 +11,16 @@ enum FixtureMode { Malformed, Oversized, Idle, + ReconnectReplay, + ReplayGap, + ReconnectExhausted, } #[derive(Clone)] struct FixtureGeyser { mode: FixtureMode, half_close_seen: std::sync::Arc, + subscribe_calls: std::sync::Arc, } #[allow(clippy::implicit_return)] // tonic::async_trait generates async wrapper tails outside the authored fixture bodies. @@ -32,6 +36,10 @@ impl yellowstone_grpc_proto::geyser::geyser_server::Geyser for FixtureGeyser { if let std::result::Result::Err(error) = verify_fixture_metadata(request.metadata()) { return std::result::Result::Err(error); } + let subscribe_call = self.subscribe_calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst).saturating_add(1); + if matches!(self.mode, FixtureMode::ReconnectExhausted) && subscribe_call > 1 { + return std::result::Result::Err(tonic::Status::unavailable("GRPC-RECONNECT-SECRET-CANARY")); + } let mut inbound = request.into_inner(); let (outbound_tx, outbound_rx) = tokio::sync::mpsc::channel(64); let mode = self.mode; @@ -111,6 +119,45 @@ impl yellowstone_grpc_proto::geyser::geyser_server::Geyser for FixtureGeyser { FixtureMode::Idle => { let _ = inbound.message().await; }, + FixtureMode::ReconnectReplay => { + if subscribe_call == 1 { + assert_eq!(initial.from_slot, std::option::Option::None); + let _ = outbound_tx.send(std::result::Result::Ok(slot_update(500))).await; + } else { + assert_eq!(subscribe_call, 2); + assert_eq!(initial.from_slot, std::option::Option::Some(500)); + if outbound_tx.send(std::result::Result::Ok(slot_update(500))).await.is_err() { + return; + } + if outbound_tx.send(std::result::Result::Ok(slot_update(501))).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::ReplayGap => { + if subscribe_call == 1 { + assert_eq!(initial.from_slot, std::option::Option::None); + let _ = outbound_tx.send(std::result::Result::Ok(slot_update(500))).await; + } else { + assert_eq!(subscribe_call, 2); + assert_eq!(initial.from_slot, std::option::Option::Some(505)); + if outbound_tx.send(std::result::Result::Ok(slot_update(505))).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; + }, } }); return std::result::Result::Ok(tonic::Response::new(super::MpscStream::new(outbound_rx))); @@ -125,9 +172,17 @@ impl yellowstone_grpc_proto::geyser::geyser_server::Geyser for FixtureGeyser { async fn subscribe_replay_info( &self, - _request: tonic::Request, + request: tonic::Request, ) -> std::result::Result, tonic::Status> { - return std::result::Result::Err(tonic::Status::unimplemented("unary is outside this fixture")); + if let std::result::Result::Err(error) = verify_fixture_metadata(request.metadata()) { + return std::result::Result::Err(error); + } + let first_available = match self.mode { + FixtureMode::ReconnectReplay | FixtureMode::ReconnectExhausted => 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")), + }; + return std::result::Result::Ok(tonic::Response::new(yellowstone_grpc_proto::geyser::SubscribeReplayInfoResponse { first_available })); } async fn ping( @@ -176,6 +231,7 @@ impl yellowstone_grpc_proto::geyser::geyser_server::Geyser for FixtureGeyser { struct FixtureServer { endpoint_url: std::string::String, half_close_seen: std::sync::Arc, + subscribe_calls: std::sync::Arc, shutdown: std::option::Option>, task: tokio::task::JoinHandle<()>, } @@ -187,9 +243,15 @@ impl FixtureServer { let local_address = incoming.local_addr().expect("fixture gRPC listener must expose local address"); let (shutdown, shutdown_receiver) = tokio::sync::oneshot::channel(); let half_close_seen = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)); + let subscribe_calls = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); let service_half_close_seen = half_close_seen.clone(); + let service_subscribe_calls = subscribe_calls.clone(); let task = tokio::spawn(async move { - let service = yellowstone_grpc_proto::geyser::geyser_server::GeyserServer::new(FixtureGeyser { mode, half_close_seen: service_half_close_seen }); + let service = yellowstone_grpc_proto::geyser::geyser_server::GeyserServer::new(FixtureGeyser { + mode, + half_close_seen: service_half_close_seen, + subscribe_calls: service_subscribe_calls, + }); let result = tonic::transport::Server::builder() .serve_with_incoming_shutdown(service, incoming, async move { let _ = shutdown_receiver.await; @@ -200,6 +262,7 @@ impl FixtureServer { return Self { endpoint_url: format!("http://{local_address}"), half_close_seen, + subscribe_calls, shutdown: std::option::Option::Some(shutdown), task, }; @@ -229,13 +292,32 @@ fn fixture_settings( update_capacity: usize, max_inbound_message_size_bytes: usize, max_outbound_message_size_bytes: usize, +) -> crate::YellowstoneGrpcEndpointSettings { + let defaults = crate::YellowstoneGrpcSessionSettings::default(); + return fixture_settings_with_reconnect( + url, + request_capacity, + update_capacity, + max_inbound_message_size_bytes, + max_outbound_message_size_bytes, + crate::YellowstoneGrpcReconnectSettings::new(0, defaults.reconnect().initial_backoff(), defaults.reconnect().max_backoff()), + ); +} + +fn fixture_settings_with_reconnect( + url: &str, + request_capacity: usize, + update_capacity: usize, + max_inbound_message_size_bytes: usize, + max_outbound_message_size_bytes: usize, + reconnect: crate::YellowstoneGrpcReconnectSettings, ) -> crate::YellowstoneGrpcEndpointSettings { let defaults = crate::YellowstoneGrpcSessionSettings::default(); let session = crate::YellowstoneGrpcSessionSettings::new( defaults.connect_timeout(), defaults.unary_timeout(), std::time::Duration::from_millis(250), - defaults.reconnect().clone(), + reconnect, request_capacity, update_capacity, max_inbound_message_size_bytes, @@ -466,3 +548,138 @@ async fn yellowstone_outbound_mutation_size_is_rejected_before_queue_dispatch() session.close().await.expect("idle fixture must observe client half-close"); server.stop().await; } + +#[tokio::test(flavor = "current_thread")] +async fn yellowstone_reconnect_replays_from_last_observed_slot_and_counts_duplicate_identity() { + let server = FixtureServer::start(FixtureMode::ReconnectReplay).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(), 500), + _ => panic!("fixture must return first Slot"), + } + let duplicate = session.next_update().await.expect("replayed duplicate must decode").expect("replayed duplicate must be present"); + match duplicate { + crate::YellowstoneSubscribeUpdate::Slot(value) => assert_eq!(value.slot(), 500), + _ => panic!("fixture must return replayed Slot"), + } + let resumed = session.next_update().await.expect("resumed slot must decode").expect("resumed slot must be present"); + match resumed { + crate::YellowstoneSubscribeUpdate::Slot(value) => assert_eq!(value.slot(), 501), + _ => panic!("fixture must return resumed Slot"), + } + let snapshot = session.snapshot(); + assert_eq!(snapshot.state(), crate::YellowstoneGrpcSubscribeState::Active); + assert_eq!(snapshot.reconnect_count(), 1); + assert_eq!(snapshot.replay_attempt_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)); + assert_eq!(snapshot.last_observed_slot(), std::option::Option::Some(501)); + session.close().await.expect("reconnected stream must close cleanly"); + assert_eq!(server.subscribe_calls.load(std::sync::atomic::Ordering::SeqCst), 2); + server.stop().await; +} + +#[tokio::test(flavor = "current_thread")] +async fn yellowstone_replay_info_proves_and_clamps_retention_gap_without_lossless_claim() { + let server = FixtureServer::start(FixtureMode::ReplayGap).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 _ = session.next_update().await.expect("first slot must decode").expect("first slot must be present"); + let resumed = session.next_update().await.expect("clamped replay must decode").expect("clamped replay must be present"); + match resumed { + crate::YellowstoneSubscribeUpdate::Slot(value) => assert_eq!(value.slot(), 505), + _ => panic!("fixture must return clamped Slot"), + } + let snapshot = session.snapshot(); + assert_eq!(snapshot.reconnect_count(), 1); + assert_eq!(snapshot.replay_attempt_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)); + assert_eq!(snapshot.last_observed_slot(), std::option::Option::Some(505)); + session.close().await.expect("gap 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; + 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(2, std::time::Duration::from_millis(5), std::time::Duration::from_millis(10)), + ); + 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 _ = session.next_update().await.expect("first slot must decode").expect("first slot must be present"); + let error = session.next_update().await.expect_err("failed reconnect attempts must exhaust the bounded budget"); + assert_eq!(error.code(), crate::ERROR_CODE_GRPC_CHANNEL_FAILED); + let snapshot = session.snapshot(); + assert_eq!(snapshot.state(), crate::YellowstoneGrpcSubscribeState::Failed); + assert_eq!(snapshot.reconnect_count(), 0); + assert_eq!(snapshot.replay_attempt_count(), 2); + 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")); + assert!(!rendered.contains(server.endpoint_url.as_str())); + assert_eq!(server.subscribe_calls.load(std::sync::atomic::Ordering::SeqCst), 3); + let _ = session.close().await; + server.stop().await; +} + +#[tokio::test(flavor = "current_thread")] +async fn yellowstone_shutdown_interrupts_reconnect_backoff_and_mutation_is_rejected_during_reconnect() { + let server = FixtureServer::start(FixtureMode::ReconnectReplay).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(500), std::time::Duration::from_millis(500)), + ); + 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 _ = session.next_update().await.expect("first slot must decode").expect("first slot must be present"); + tokio::time::timeout(std::time::Duration::from_secs(1), async { + loop { + if session.state() == crate::YellowstoneGrpcSubscribeState::Reconnecting { + return; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("fixture must enter reconnecting state"); + let mutation_error = session.try_update(&mutated_request()).expect_err("request mutation must be rejected during reconnect"); + assert_eq!(mutation_error.code(), crate::ERROR_CODE_GRPC_SESSION_CLOSED); + session.close().await.expect("shutdown must interrupt reconnect backoff without opening another stream"); + assert_eq!(server.subscribe_calls.load(std::sync::atomic::Ordering::SeqCst), 1); + server.stop().await; +} diff --git a/deltas/0.2.9/pre.010.md b/deltas/0.2.9/pre.010.md new file mode 100644 index 0000000..50e8113 --- /dev/null +++ b/deltas/0.2.9/pre.010.md @@ -0,0 +1,94 @@ + + + +# Delta `0.2.9-pre.010` — reconnect borné + replay prudent + continuité observable + +## Base + +```text +0.2.9-pre.009-fix.001 +gate opérateur final : fmt/audit/check/Clippy/workspace PASS sans warning +Transport : 379 unit + 48 public API + 42 release-completeness + 4 doctests +``` + +## Changements runtime + +- ajoute l'état public `YellowstoneGrpcSubscribeState::Reconnecting` ; +- ajoute `YellowstoneGrpcSubscribeSnapshot` avec compteurs sûrs de reconnect/replay/continuité ; +- active le `YellowstoneGrpcReconnectSettings` déjà introduit en `pre.002` pour le stream `Subscribe` ; +- rouvre le stream avec un backoff exponentiel borné et un nombre d'essais borné ; +- conserve le dernier `YellowstoneSubscribeRequest` complet accepté ; +- après perte du stream, resoumet ce request déterministement avec un `from_slot` au moins égal au plus haut slot déjà observé ; +- consulte `SubscribeReplayInfo` avant les tentatives de replay lorsque `from_slot` est disponible ; +- si `first_available > from_slot`, incrémente `continuity_gap_count` et clamp le replay à `first_available` ; +- maintient un cache borné de 512 identités d'updates pour observer les duplicates autour d'un replay ; +- les duplicates restent livrés : aucun contrat exactly-once n'est inventé ; +- rejette les mutations request pendant `Reconnecting` pour éviter une application ambiguë entre deux streams physiques ; +- un shutdown pendant le backoff interrompt immédiatement la boucle et interdit une nouvelle ouverture ; +- un épuisement du budget devient terminal avec un code KSP sûr `grpc_channel_failed`. + +## Sémantique de continuité + +```text +from_slot de reconnect = max(from_slot explicite du dernier request, last_observed_slot) +ReplayInfo = information de rétention ; pas preuve de replay complet +continuity gap = couverture de replay indisponible, comptée seulement si first_available > slot demandé +slot manquant entre deux updates filtrés = jamais assimilé automatiquement à un gap +duplicate = identité KSP bornée déjà observée ; update toujours livrée +exactly-once = non garanti +lossless = non garanti +ordre global sans gap = non garanti +``` + +Identités retenues pour l'observation bornée : + +```text +Account slot + pubkey + write_version +Slot slot + status +Transaction slot + signature +TransactionStatus slot + signature +Block slot + blockhash +BlockMeta slot + blockhash +Entry slot + index + hash +Ping/Pong hors déduplication +``` + +La couverture d'une éventuelle divergence de node reste volontairement limitée : deux Block/BlockMeta de même slot mais de blockhash différent ne sont pas classés duplicate. KSP ne généralise pas cette preuve aux familles qui ne transportent pas de blockhash. + +## Fixture locale + +| Cas | Preuve | +|------------------------------------------------|-------------| +| reconnect après server half-close | PASS source | +| `from_slot = last_observed_slot` | PASS source | +| ReplayInfo disponible | PASS source | +| replay duplicate observable | PASS source | +| couverture replay indisponible prouvée + clamp | PASS source | +| budget reconnect épuisé | PASS source | +| mutation pendant reconnect rejetée | PASS source | +| shutdown pendant backoff | PASS source | +| endpoint/Status secret non réémis | PASS source | +| anciens gates bidi avec reconnect `0` | PASS source | + +## Frontières + +```text +OUT pre.010 : Config V3 / PublicNode / provider facade +OUT 0.2.9 : SubscribeDeshred +aucune nouvelle dépendance / feature Cargo +aucun yellowstone-grpc-client runtime +``` + +## Validation candidate + +```bash +cargo fmt --all +python3 scripts/audit_rust_workspace_rules.py +cargo check --workspace +cargo clippy --workspace --all-targets +cargo test -p ksp-onchain-transport-lib +cargo test -p ksp-core-lib --test workspace_dependencies +cargo test --workspace +``` + +Le `cargo tree` n'est pas requis : aucune dépendance ni feature n'a changé. diff --git a/docs/plans/016-V0_2_9_YELLOWSTONE_GRPC_PLAN.md b/docs/plans/016-V0_2_9_YELLOWSTONE_GRPC_PLAN.md index 2667c48..e78ac92 100644 --- a/docs/plans/016-V0_2_9_YELLOWSTONE_GRPC_PLAN.md +++ b/docs/plans/016-V0_2_9_YELLOWSTONE_GRPC_PLAN.md @@ -1,9 +1,9 @@ - + # Plan `0.2.9` — moteur Yellowstone gRPC + standard Solana + PublicNode -> **Statut : `0.2.9-pre.009` est fonctionnellement verte sur gate opérateur (Transport 379 unit + 48 public API + 42 completeness + 4 doctests, workspace PASS), mais Clippy signale deux warnings d’hygiène. `0.2.9-pre.009-fix.001` est candidate pour les supprimer sans modifier le protocole ni le lifecycle ; reconnect/replay restent `pre.010`.** +> **Statut : `0.2.9-pre.009-fix.001` est fermé sur gate opérateur intégralement vert : fmt/audit/check/Clippy sans warning, Transport 379 unit + 48 public API + 42 completeness + 4 doctests, dependency canary et workspace PASS. `0.2.9-pre.010` est candidate pour le reconnect/replay KSP-owned, les gaps de couverture replay prouvés par ReplayInfo et les duplicates observables sans promesse lossless.** ## 1. Objet, base et état d'ouverture @@ -934,11 +934,11 @@ pre.007 DONE — standard Solana : Transactions + transaction_status pre.008 DONE — standard Solana : Blocks + block_meta + entry budget : 15–20 min ; gate final fix.001 : fmt/audit/check/Clippy/workspace PASS sans warning + Transport 370/47/41/4 -pre.009 FIX.001 CANDIDATE — moteur partagé : bidi mutation + Ping/Pong + half-close + backpressure + shutdown - budget : 15–20 min ; gate fonctionnel 379/48/42/4 + workspace PASS ; fix de deux warnings Clippy +pre.009 DONE — moteur partagé : bidi mutation + Ping/Pong + half-close + backpressure + shutdown + budget : 15–20 min ; gate final fix.001 : fmt/audit/check/Clippy/workspace PASS sans warning + Transport 379/48/42/4 -pre.010 moteur partagé : reconnect/resubscribe + from_slot/ReplayInfo + gaps/duplicates - budget : 15–20 min ; preuve : reconnect local déterministe + aucune promesse lossless +pre.010 CANDIDATE — moteur partagé : reconnect/resubscribe + from_slot/ReplayInfo + gaps/duplicates + budget : 15–20 min ; preuve source : reconnect local déterministe + ReplayInfo clamp + duplicate observable + aucune promesse lossless pre.011 Config V3 + séparation protocol/provider + profils PublicNode Mainnet/Testnet budget : 15–20 min ; preuve : V1/V2 backward + schema/mapping/redaction + Config -> Transport @@ -1387,4 +1387,34 @@ Le fix applique deux corrections sans `allow` : Aucun changement n’est apporté aux neuf variantes wire, à la mutation request, à Ping/Pong, au half-close, au shutdown, à la backpressure, aux dépendances ou aux frontières de `pre.010`. -**Gate fix requis :** fmt/audit/check/Clippy sans warning + Transport 379/48/42/4 + dependency canary + workspace. +**Gate fix final :** fmt/audit/check/Clippy sans warning + Transport 379/48/42/4 + dependency canary + workspace **PASS le 2026-08-24**. + + +## 28. `pre.010` — reconnect KSP-owned, reprise `from_slot` et continuité prudente + +Le proto courant conserve `SubscribeRequest.from_slot` et le unary `SubscribeReplayInfo.first_available`. Le changelog upstream relu le `2026-08-24` confirme qu’un replay `from_slot` peut avoir des défauts spécifiques de famille : le correctif du `2026-07-22` concernait précisément des subscriptions Blocks acceptées mais reprenant live avec un state gap. KSP ne transforme donc pas la présence de `from_slot` en garantie lossless. + +Politique retenue : + +```text +stream perdu -> état Reconnecting +budget -> max_retries borné déjà dans YellowstoneGrpcReconnectSettings +backoff -> exponentiel, initial/max bornés +request de resubscribe -> dernier request complet accepté +resume slot -> max(from_slot explicite, highest observed slot) +ReplayInfo first_available > resume -> couverture replay indisponible prouvée + compteur + clamp à first_available +ReplayInfo indisponible -> reconnect poursuit sans inférer de gap +mutation pendant reconnect -> rejet explicite, pas de coalescing ambigu +duplicate replay -> identité bornée comptée mais update toujours livrée +shutdown pendant backoff -> interruption, aucune nouvelle ouverture +``` + +`YellowstoneGrpcSubscribeSnapshot` expose `reconnect_count`, `continuity_gap_count`, `duplicate_update_count`, `replay_attempt_count`, `last_requested_from_slot`, `last_observed_slot`, l’état lifecycle et le terminal error code sûr. Le cache d’identité est borné à 512 entrées et ne contient pas les payloads arbitraires complets. + +La détection de gap est volontairement **plus stricte** qu’un simple saut de numéro de slot : avec des filtres Accounts/Transactions/Blocks, l’absence d’update sur un slot intermédiaire peut être normale. Seule une borne de rétention `first_available` supérieure au replay demandé constitue ici une preuve objective que cette plage n'est plus rejouable ; elle ne prouve pas qu'un update correspondant aux filtres existait ou a été perdu. + +De même, KSP ne supprime pas les duplicates : il les observe autour du replay, les compte et continue de les livrer. Cela maintient explicitement les non-promesses `exactly-once`, `lossless` et `ordre global sans gap`. + +Couverture de divergence : les identités Block/BlockMeta incluent le blockhash ; deux histoires de même slot avec blockhash différent ne sont donc pas prises pour un duplicate. Aucune généralisation d’equivocation n’est annoncée pour les familles sans preuve de blockhash. + +**Gate candidat :** fixture locale reconnect/replay, duplicate, gap de couverture replay, budget épuisé, mutation rejetée pendant reconnect et shutdown pendant backoff ; audit statique clean. Aucun changement de dépendance ou feature. diff --git a/docs/validation/012-V0_2_9_YELLOWSTONE_GRPC.md b/docs/validation/012-V0_2_9_YELLOWSTONE_GRPC.md index 459404f..7dda050 100644 --- a/docs/validation/012-V0_2_9_YELLOWSTONE_GRPC.md +++ b/docs/validation/012-V0_2_9_YELLOWSTONE_GRPC.md @@ -1,9 +1,9 @@ - + # Validation `0.2.9` — moteur Yellowstone + standard Solana + PublicNode -> **Statut : `pre.009` est fonctionnellement verte : Transport 379 unit + 48 public API + 42 completeness + 4 doctests et workspace PASS. Clippy signale deux warnings d’hygiène ; `pre.009-fix.001` est candidate pour les supprimer sans changement de protocole/lifecycle.** +> **Statut : `pre.009-fix.001` est fermé : fmt/audit/check/Clippy sans warning, Transport 379 unit + 48 public API + 42 completeness + 4 doctests, dependency canary et workspace PASS. `pre.010` est candidate reconnect/replay avec observabilité prudente des gaps/duplicates.** ## 1. Autorités du gate @@ -287,25 +287,25 @@ DONE pre.003 source/tests channel/client Debug sans URL, metadata value ni raw ## 11. Lifecycle / backpressure / replay -| Cas | Attendu | État | -|--------------------------------------------|------------------------------------------|-----------------------------------| -| stream open | session bornée | CANDIDATE pre.009 | -| request mutation | ordre déterministe | CANDIDATE pre.009 | -| server Ping -> client request ping -> Pong | explicite | CANDIDATE pre.009 | -| server half-close | terminal normal avant reconnect | CANDIDATE pre.009 | -| client close | half-close + cleanup borné | CANDIDATE pre.009 | -| receiver drop | Drop session => cleanup best-effort | CANDIDATE pre.009 | -| slow subscription | pas de queue infinie | CANDIDATE pre.009 | -| inbound oversized | limite Tonic avant payload KSP | CANDIDATE pre.009 | -| outbound oversized | encoded_len avant queue/write | CANDIDATE pre.009 | -| reconnect budget | borné | TODO | -| resubscribe order | déterministe | TODO | -| `from_slot` | utilisé sans promesse lossless | TODO | -| ReplayInfo | informatif | DONE unary ; usage reconnect TODO | -| duplicates | observables | TODO | -| gaps | observables | TODO | -| divergent node history | couverture documentée | TODO | -| shutdown during reconnect | aucune nouvelle connexion après shutdown | TODO | +| Cas | Attendu | État | +|--------------------------------------------|-----------------------------------------------------------------|----------------------------------------| +| stream open | session bornée | DONE pre.009 / gate fix.001 PASS | +| request mutation | ordre déterministe | DONE pre.009 | +| server Ping -> client request ping -> Pong | explicite | DONE pre.009 | +| server half-close | reconnect si budget > 0 ; terminal normal sinon | CANDIDATE pre.010 | +| client close | half-close + cleanup borné | DONE pre.009 | +| receiver drop | Drop session => cleanup best-effort | DONE pre.009 | +| slow subscription | pas de queue infinie | DONE pre.009 | +| inbound oversized | limite Tonic avant payload KSP | DONE pre.009 | +| outbound oversized | encoded_len avant queue/write | DONE pre.009 | +| reconnect budget | borné | CANDIDATE pre.010 | +| resubscribe order | dernier request complet accepté | CANDIDATE pre.010 | +| `from_slot` | max(explicite, highest observed), sans promesse lossless | CANDIDATE pre.010 | +| ReplayInfo | informatif ; clamp seulement si gap de couverture replay prouvé | DONE unary + CANDIDATE usage reconnect | +| duplicates | identités bornées observables, jamais supprimées | CANDIDATE pre.010 | +| gaps | seulement si `first_available > requested_from_slot` | CANDIDATE pre.010 | +| divergent node history | Block/BlockMeta hash-aware ; pas de claim global | CANDIDATE couverture documentée | +| shutdown during reconnect | aucune nouvelle connexion après shutdown | CANDIDATE pre.010 | Claims interdits sans nouvelle preuve : @@ -411,8 +411,8 @@ pre.005 DONE standard: accounts + slots 15 pre.006 DONE structure: namespace privé HTTP `http_*` 15–20 min ; gate PASS pre.007 DONE standard: transactions + transaction_status 15–20 min ; gate PASS pre.008 DONE standard: blocks + block_meta + entry 15–20 min ; gate final fix.001 PASS sans warning -pre.009 FIX.001 CANDIDATE moteur: bidi/backpressure/half-close/shutdown 15–20 min ; gate fonctionnel vert, 2 warnings Clippy -pre.010 TODO moteur: reconnect/replay/gap/duplicate 15–20 min +pre.009 DONE moteur: bidi/backpressure/half-close/shutdown 15–20 min ; gate final fix.001 PASS sans warning, 379/48/42/4 +pre.010 CANDIDATE moteur: reconnect/replay/gap/duplicate 15–20 min ; source/fixture candidate, gate opérateur à exécuter pre.011 TODO Config V3 + protocol/provider + profils PublicNode 15–20 min pre.012 TODO PublicNode live + compliance + docs/prompt 0.2.10 15–20 min rel.001 TODO stable @@ -918,7 +918,7 @@ Le flux sortant utilise la même `mpsc` bornée que celle fournie à `tonic::cli `pre.009` ne tente aucun reconnect. Un half-close serveur est terminal normal ; `Status`, decode invalide et overflow sont terminaux en erreur. Le budget de reconnect, l'ordre de resubscribe, `from_slot`/ReplayInfo, gaps, duplicates et node divergence restent exclusivement `pre.010`. -**Verdict `pre.009` :** contrat fonctionnel validé ; fix Clippy-only requis avant fermeture. +**Verdict `pre.009` :** fermé après `pre.009-fix.001`, gate opérateur final intégralement vert et sans warning. ## 28. Gate `pre.009-fix.001` — taille de l’enum update + canari `Send` @@ -937,6 +937,45 @@ Le flux sortant utilise la même `mpsc` bornée que celle fournie à `tonic::cli | `large_enum_variant` | Transaction ~664 B, Block ~280 B | `Transaction(Box)` | | `extra_unused_type_parameters` | helper `assert_send()` n’utilise pas T | `PhantomData` matérialise l’usage compile-time | -Le boxing ne modifie ni le protobuf ni `YellowstoneTransactionUpdate` : il réduit uniquement la taille du discriminant public transporté dans la queue bidi. Le canari `Send` conserve la même contrainte et ne construit aucune session réseau. Aucun `allow` Clippy n’est ajouté. +Le boxing ne modifie ni le protobuf ni `YellowstoneTransactionUpdate` : il réduit uniquement la taille du discriminant public transporté dans la queue bidi. Le canari `Send` conserve la même contrainte et ne construit aucune session réseau. Aucun `allow` Clippy n’est ajouté. Le gate opérateur final du `2026-08-24` repasse intégralement vert, Clippy sans warning, avec les compteurs 379/48/42/4 inchangés. **Verdict fix candidate :** correction minimale prête ; fermeture de `pre.009` après gate opérateur sans warning. + + +## 29. Candidate `pre.010` — reconnect/replay et continuité observable + +Réaduit upstream au `2026-08-24` : + +```text +SubscribeRequest.from_slot toujours présent +SubscribeReplayInfoResponse.first_available? toujours présent +changelog 2026-07-22 fix replay Blocks/state gap +changelog 2026-06-15 upstream autoreconnect traite l'equivocation par quarantaine/blockhash +``` + +KSP conserve sa stratégie B : aucune importation de la sémantique `yellowstone-grpc-client` autoreconnect. Le stream KSP rouvre lui-même `/geyser.Geyser/Subscribe` sur le channel N1 existant. + +| Preuve source `pre.010` | État | +|-------------------------------------------------------------|--------------------| +| état public `Reconnecting` | SOURCE | +| snapshot reconnect/replay/gap/duplicate | SOURCE | +| budget `max_retries` | SOURCE | +| backoff exponentiel initial/max | SOURCE | +| dernier request complet resoumis | SOURCE | +| resume `max(explicit, last_observed)` | SOURCE | +| ReplayInfo consulté avant replay | SOURCE | +| `first_available > requested` => replay non couvert + clamp | SOURCE+TEST | +| duplicate identité bornée, update livrée | SOURCE+TEST | +| mutation refusée pendant `Reconnecting` | SOURCE+TEST | +| shutdown interrompt le backoff | SOURCE+TEST | +| épuisement budget => terminal sûr | SOURCE+TEST | +| cache identité borné à 512 | SOURCE | +| exactement-once/lossless | EXPLICIT NON-CLAIM | +| Config V3 / PublicNode | OUT pre.010 | +| `SubscribeDeshred` | OUT 0.2.9 | + +La détection de gap n'utilise pas les sauts entre slots d'updates filtrés : un filtre peut légitimement ne produire aucun message pendant plusieurs slots. Le seul gap compté par cette tranche est une plage de replay objectivement devenue indisponible d'après `SubscribeReplayInfo.first_available` ; cela ne prouve pas qu'un update correspondant aux filtres existait ou a été perdu. + +Le cache de duplicates retient des identités minimales Account/Slot/Transaction/TransactionStatus/Block/BlockMeta/Entry. Il ne retient pas les payloads complets et ne supprime jamais un update. Pour Block/BlockMeta, le blockhash fait partie de l'identité afin de ne pas classer deux histoires divergentes de même slot comme un simple duplicate ; aucune détection globale d'equivocation n'est prétendue. + +**Gate candidate attendu :** audit statique clean ; après exécution opérateur, Transport attendu autour de 383 unit / 49 public API / 43 release-completeness / 4 doctests. Aucune dépendance/feature changée, donc pas de `cargo tree` supplémentaire.