diff --git a/Cargo.toml b/Cargo.toml index 9fc181b..2fd3f07 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,12 @@ # file: Cargo.toml -# version: 197 +# version: 198 [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.7-pre.4.fix.1" +version = "0.2.7-pre.5" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-onchain-transport-lib/README.md b/crates/ksp-onchain-transport-lib/README.md index 22d71d4..500b2d6 100644 --- a/crates/ksp-onchain-transport-lib/README.md +++ b/crates/ksp-onchain-transport-lib/README.md @@ -1,5 +1,5 @@ - + # `ksp-onchain-transport-lib` @@ -105,6 +105,14 @@ Le chemin JSON-RPC générique reste `pub(crate)` en `pre.004`. Il sert de primi Les tests déterministes utilisent un serveur WebSocket local et prouvent le handshake, le round-trip JSON-RPC, le dispatch de réponses hors ordre, l'isolation des erreurs RPC applicatives, deux sessions physiques distinctes sur la même URL et la redaction des erreurs de connexion. +### Durcissement `0.2.7-pre.005` + +La session physique dispose désormais de `WsSession::close().await`. Le signal de shutdown est distinct de la command queue, passe l'état en `Closing`, annule les requests JSON-RPC en attente, envoie un Close WebSocket best-effort sous `close_timeout`, puis publie `Closed`. Un peer qui ne répond pas au Close ne peut donc pas bloquer indéfiniment le shutdown. + +Les limites `max_message_size`, `max_frame_size`, `max_write_buffer_size` et `max_pending_requests` sont couvertes par des fixtures adversariales locales. Les requests outbound qui dépassent les bornes message/frame sont rejetées avant écriture ; les frames/messages inbound surdimensionnés sont rejetés par Tungstenite avant parse JSON. Les timeouts pending libèrent leur capacité sans faire tomber une session encore saine. + +Ping/Pong/Close sont traités comme control frames : le Pong automatique Tungstenite est flushé, un Close distant propre mène à `Closed`, tandis qu'une erreur I/O ou une violation de protocole mène à `Failed`. Aucun heartbeat applicatif périodique n'est ajouté. + ## Résilience L'admission est calculée par couple endpoint/rôle. Le pool applique : diff --git a/crates/ksp-onchain-transport-lib/USAGE.md b/crates/ksp-onchain-transport-lib/USAGE.md index 55689a5..c0fa7f5 100644 --- a/crates/ksp-onchain-transport-lib/USAGE.md +++ b/crates/ksp-onchain-transport-lib/USAGE.md @@ -1,5 +1,5 @@ - + # Utilisation de `ksp-onchain-transport-lib` @@ -94,6 +94,20 @@ Deux appels `WsSession::connect` avec le même endpoint créent volontairement d Le socket brut et la primitive JSON-RPC générique ne sont pas publics. Les wrappers `*Subscribe` typed et leurs handles seront ajoutés au-dessus de l'actor ; `pre.004` ne doit donc pas être utilisé comme escape hatch provider-specific. +### Fermeture explicite + +À partir de `0.2.7-pre.005`, fermer explicitement la session est la voie normale de shutdown : + +```rust +let session = ksp_onchain_transport_lib::WsSession::connect(endpoint).await?; +// ... utilisation future des subscriptions typed ... +session.close().await?; +``` + +`close()` agit sur toute la session physique, y compris les clones du handle. Il annule les requests en attente, publie `Closing`, tente le Close WebSocket dans le budget configuré, puis publie `Closed`. Une session `Closed` refuse les nouvelles requests internes. + +Les limites de taille et de capacité sont des policies KSP configurables par `WsSessionSettings`; elles ne doivent pas être interprétées comme des limites protocolaires Solana officielles. + Le snapshot expose seulement l'identité locale, les metadata logiques de l'endpoint, l'état et les compteurs sûrs. L'URL n'est jamais projetée. Le shutdown async explicite arrive en `pre.005`; la disparition de tous les handles déclenche seulement le cleanup actor best-effort de cette foundation. ## 4. Appels typés diff --git a/crates/ksp-onchain-transport-lib/src/lib.rs b/crates/ksp-onchain-transport-lib/src/lib.rs index 508d16e..9d3b7bf 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: 23 +// version: 24 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -18,7 +18,8 @@ //! The candidate surface therefore exposes typed wrappers for all 52 current audited Solana HTTP methods while retaining 14 removed historical descriptors. //! `0.2.7-pre.002` adds the provider-neutral WebSocket settings foundation, redacted endpoint URLs, explicit protocol-family discrimination, local session/subscription //! identities, observable lifecycle states and safe snapshots. `0.2.7-pre.004` adds the first physical WebSocket runtime: bounded handshake, one actor-owned socket, -//! bounded command/pending JSON-RPC flow, safe session snapshots and deterministic local-server fixtures. Subscription registration remains deferred. +//! bounded command/pending JSON-RPC flow, safe session snapshots and deterministic local-server fixtures. `0.2.7-pre.005` adds explicit bounded shutdown, +//! request/frame/message adversarial limits, control-frame handling and pending-request cancellation/timeout cleanup. Subscription registration remains deferred. mod client; mod constants; diff --git a/crates/ksp-onchain-transport-lib/src/ws_session.rs b/crates/ksp-onchain-transport-lib/src/ws_session.rs index 1be46c1..2887eab 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: 2 +// version: 3 use futures_util::SinkExt; // rust-rules: trait-import use futures_util::StreamExt; // rust-rules: trait-import @@ -9,13 +9,15 @@ static NEXT_WS_SESSION_ID: std::sync::atomic::AtomicU64 = std::sync::atomic::Ato /// 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 -/// interaction is serialized through a bounded command queue. +/// interaction is serialized through bounded channels. #[derive(Clone)] pub struct WsSession { id: crate::WsSessionId, command_tx: tokio::sync::mpsc::Sender, + shutdown_tx: tokio::sync::watch::Sender>, snapshot_rx: tokio::sync::watch::Receiver, command_timeout: std::time::Duration, + close_timeout: std::time::Duration, } impl WsSession { @@ -46,9 +48,11 @@ impl WsSession { std::vec::Vec::new(), ); let (command_tx, command_rx) = tokio::sync::mpsc::channel(endpoint.session().command_queue_capacity()); + let (shutdown_tx, shutdown_rx) = tokio::sync::watch::channel(std::option::Option::None::); let (snapshot_tx, snapshot_rx) = tokio::sync::watch::channel(initial_snapshot); let (startup_tx, startup_rx) = tokio::sync::oneshot::channel(); let command_timeout = endpoint.session().command_timeout(); + let close_timeout = endpoint.session().close_timeout(); ksp_logging_lib::debug!( target: crate::TRACING_TARGET, session_id = id.get(), @@ -58,25 +62,25 @@ impl WsSession { protocol = endpoint.protocol().as_str(), "starting physical WebSocket session actor" ); - let join_handle = tokio::spawn(run_ws_session_actor(id, endpoint, command_rx, snapshot_tx, startup_tx)); + let join_handle = tokio::spawn(run_ws_session_actor(id, endpoint, command_rx, shutdown_rx, snapshot_tx, startup_tx)); let startup_wait = tokio::time::timeout(command_timeout, startup_rx).await; - match startup_wait { + return match startup_wait { std::result::Result::Ok(std::result::Result::Ok(std::result::Result::Ok(()))) => { - return std::result::Result::Ok(Self { id, command_tx, snapshot_rx, command_timeout }); + std::result::Result::Ok(Self { id, command_tx, shutdown_tx, snapshot_rx, command_timeout, close_timeout }) }, std::result::Result::Ok(std::result::Result::Ok(std::result::Result::Err(error))) => { join_handle.abort(); - return std::result::Result::Err(error); + std::result::Result::Err(error) }, std::result::Result::Ok(std::result::Result::Err(_)) => { join_handle.abort(); - return std::result::Result::Err(ws_session_closed_error(id, "WebSocket session actor ended during startup")); + std::result::Result::Err(ws_session_closed_error(id, "WebSocket session actor ended during startup")) }, std::result::Result::Err(_) => { join_handle.abort(); - return std::result::Result::Err(ws_timeout_error(id, "WebSocket handshake exceeded the configured command timeout")); + std::result::Result::Err(ws_timeout_error(id, "WebSocket handshake exceeded the configured command timeout")) }, - } + }; } /// Returns the stable local session identity. @@ -97,11 +101,57 @@ impl WsSession { return self.snapshot_rx.borrow().state(); } + /// Explicitly closes this physical session under the configured close timeout. + /// + /// Shutdown is session-wide: calling this method through any clone requests the actor to enter `Closing`, cancels pending JSON-RPC requests, sends a + /// best-effort WebSocket Close frame and waits for the safe lifecycle snapshot to reach `Closed`. The shutdown signal is independent from the bounded + /// command queue so a saturated command queue cannot prevent close from being requested. + pub async fn close(&self) -> ksp_core_lib::Result<()> { + match self.state() { + crate::WsSessionState::Closed => return std::result::Result::Ok(()), + crate::WsSessionState::Failed => { + return std::result::Result::Err(ws_session_closed_error(self.id, "WebSocket session is already in terminal failed state")); + }, + _ => {}, + } + let deadline = tokio::time::Instant::now() + self.close_timeout; + self.shutdown_tx.send_replace(std::option::Option::Some(deadline)); + let mut snapshot_rx = self.snapshot_rx.clone(); + loop { + match snapshot_rx.borrow().state() { + crate::WsSessionState::Closed => return std::result::Result::Ok(()), + crate::WsSessionState::Failed => { + return std::result::Result::Err(ws_session_closed_error(self.id, "WebSocket session failed while explicit shutdown was in progress")); + }, + _ => {}, + } + let changed = tokio::time::timeout_at(deadline, snapshot_rx.changed()).await; + match changed { + std::result::Result::Ok(std::result::Result::Ok(())) => {}, + std::result::Result::Ok(std::result::Result::Err(_)) => { + return match snapshot_rx.borrow().state() { + crate::WsSessionState::Closed => std::result::Result::Ok(()), + _ => std::result::Result::Err(ws_session_closed_error(self.id, "WebSocket session actor ended before publishing Closed state")), + }; + }, + std::result::Result::Err(_) => { + if snapshot_rx.borrow().state() == crate::WsSessionState::Closed { + return std::result::Result::Ok(()); + } + return std::result::Result::Err(ws_timeout_error(self.id, "WebSocket explicit shutdown exceeded the configured close timeout")); + }, + } + } + } + /// Executes one internal JSON-RPC request through the actor-owned socket. /// /// This remains crate-private in `0.2.7`; standard subscriptions consume it without exposing a public raw provider-extension escape hatch. #[allow(dead_code)] // Staged in pre.004 and consumed by the subscription engine starting in pre.006. pub(crate) async fn execute_json_rpc(&self, method: &'static str, params: std::vec::Vec) -> ksp_core_lib::Result { + if self.state() != crate::WsSessionState::Active { + return std::result::Result::Err(ws_session_closed_error(self.id, "WebSocket session is not active")); + } let (response_tx, response_rx) = tokio::sync::oneshot::channel(); let command = WsSessionCommand::ExecuteJsonRpc { method, params, response_tx }; let send_wait = tokio::time::timeout(self.command_timeout, self.command_tx.send(command)).await; @@ -114,14 +164,10 @@ impl WsSession { return std::result::Result::Err(ws_timeout_error(self.id, "WebSocket session command queue remained unavailable until timeout")); }, } - let response_wait = tokio::time::timeout(self.command_timeout, response_rx).await; - return match response_wait { - std::result::Result::Ok(std::result::Result::Ok(result)) => result, - std::result::Result::Ok(std::result::Result::Err(_)) => { - std::result::Result::Err(ws_session_closed_error(self.id, "WebSocket session ended before the JSON-RPC response was delivered")) - }, + return match response_rx.await { + std::result::Result::Ok(result) => result, std::result::Result::Err(_) => { - std::result::Result::Err(ws_timeout_error(self.id, "WebSocket JSON-RPC request exceeded the configured command timeout")) + std::result::Result::Err(ws_session_closed_error(self.id, "WebSocket session ended before the JSON-RPC response was delivered")) }, }; } @@ -147,10 +193,18 @@ struct PendingWsRequest { response_tx: tokio::sync::oneshot::Sender>, } +enum WsActorIoOutcome { + Continue, + RemoteClosed, + ShutdownRequested { deadline: tokio::time::Instant }, + Failed { code: ksp_core_lib::ErrorCode, pending_message: &'static str }, +} + async fn run_ws_session_actor( id: crate::WsSessionId, endpoint: crate::WsEndpointSettings, mut command_rx: tokio::sync::mpsc::Receiver, + mut shutdown_rx: tokio::sync::watch::Receiver>, snapshot_tx: tokio::sync::watch::Sender, startup_tx: tokio::sync::oneshot::Sender>, ) { @@ -213,35 +267,81 @@ async fn run_ws_session_actor( let mut next_request_id = 1_u64; let mut pending = std::collections::BTreeMap::::new(); loop { + prune_cancelled_pending(id, &mut pending); let timeout_deadline = next_pending_deadline(&pending); tokio::select! { + biased; + shutdown_changed = shutdown_rx.changed() => { + let deadline = resolve_shutdown_deadline(&shutdown_rx, shutdown_changed, endpoint.session().close_timeout()); + close_session_actor(id, &endpoint, &snapshot_tx, &mut websocket, &mut pending, deadline).await; + return; + }, maybe_command = command_rx.recv() => { let command = match maybe_command { std::option::Option::Some(command) => command, std::option::Option::None => { + let deadline = tokio::time::Instant::now() + endpoint.session().close_timeout(); ksp_logging_lib::trace!(target: crate::TRACING_TARGET, session_id = id.get(), "all WebSocket session handles dropped; closing actor"); - let _ = websocket.send(tokio_tungstenite::tungstenite::Message::Close(std::option::Option::None)).await; + close_session_actor(id, &endpoint, &snapshot_tx, &mut websocket, &mut pending, deadline).await; + return; + }, + }; + let command_outcome = handle_session_command( + id, + &endpoint, + &mut websocket, + &mut pending, + &mut next_request_id, + &mut shutdown_rx, + command, + ) + .await; + match command_outcome { + WsActorIoOutcome::Continue => { + publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Active, pending.len()); + }, + WsActorIoOutcome::RemoteClosed => { publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Closed, pending.len()); fail_all_pending(&mut pending, id, crate::ERROR_CODE_WS_SESSION_CLOSED, "WebSocket session closed before pending response delivery"); return; }, - }; - let command_outcome = handle_session_command(id, &endpoint, &mut websocket, &mut pending, &mut next_request_id, command).await; - if !command_outcome { - publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Failed, pending.len()); - fail_all_pending(&mut pending, id, crate::ERROR_CODE_WS_CONNECTION_FAILED, "WebSocket connection failed while writing a request"); - return; + WsActorIoOutcome::ShutdownRequested { deadline } => { + close_session_actor(id, &endpoint, &snapshot_tx, &mut websocket, &mut pending, deadline).await; + return; + }, + WsActorIoOutcome::Failed { code, pending_message } => { + publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Failed, pending.len()); + fail_all_pending(&mut pending, id, code, pending_message); + return; + }, } - publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Active, pending.len()); }, maybe_message = websocket.next() => { - let keep_running = handle_socket_message(id, &endpoint, maybe_message, &mut websocket, &mut pending).await; - if !keep_running { - publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Failed, pending.len()); - fail_all_pending(&mut pending, id, crate::ERROR_CODE_WS_CONNECTION_FAILED, "WebSocket connection ended before pending response delivery"); - return; + let socket_outcome = handle_socket_message(id, &endpoint, maybe_message, &mut websocket, &mut pending, &mut shutdown_rx).await; + match socket_outcome { + WsActorIoOutcome::Continue => { + publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Active, pending.len()); + }, + WsActorIoOutcome::RemoteClosed => { + publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Closed, pending.len()); + fail_all_pending( + &mut pending, + id, + crate::ERROR_CODE_WS_SESSION_CLOSED, + "Remote peer closed WebSocket session before pending response delivery", + ); + return; + }, + WsActorIoOutcome::ShutdownRequested { deadline } => { + close_session_actor(id, &endpoint, &snapshot_tx, &mut websocket, &mut pending, deadline).await; + return; + }, + WsActorIoOutcome::Failed { code, pending_message } => { + publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Failed, pending.len()); + fail_all_pending(&mut pending, id, code, pending_message); + return; + }, } - publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Active, pending.len()); }, () = tokio::time::sleep_until(timeout_deadline) => { expire_pending_requests(id, &mut pending); @@ -251,18 +351,20 @@ async fn run_ws_session_actor( } } +#[allow(clippy::too_many_arguments)] async fn handle_session_command( id: crate::WsSessionId, endpoint: &crate::WsEndpointSettings, websocket: &mut tokio_tungstenite::WebSocketStream, pending: &mut std::collections::BTreeMap, next_request_id: &mut u64, + shutdown_rx: &mut tokio::sync::watch::Receiver>, command: WsSessionCommand, -) -> bool +) -> WsActorIoOutcome where S: tokio::io::AsyncRead + tokio::io::AsyncWrite + std::marker::Unpin, { - match command { + return match command { WsSessionCommand::ExecuteJsonRpc { method, params, response_tx } => { if pending.len() >= endpoint.session().max_pending_requests() { let error = ksp_core_lib::Error::new(crate::ERROR_CODE_WS_BACKPRESSURE_OVERFLOW, "WebSocket pending JSON-RPC request capacity is exhausted") @@ -276,7 +378,7 @@ where max_pending_requests = endpoint.session().max_pending_requests(), "rejected WebSocket JSON-RPC request because pending capacity is exhausted" ); - return true; + return WsActorIoOutcome::Continue; } let request_id = *next_request_id; let incremented = request_id.checked_add(1); @@ -286,7 +388,7 @@ where let error = ksp_core_lib::Error::new(crate::ERROR_CODE_WS_PROTOCOL_ERROR, "WebSocket JSON-RPC request identifier space is exhausted") .with_context("session_id", id.get().to_string()); let _ = response_tx.send(std::result::Result::Err(error)); - return true; + return WsActorIoOutcome::Continue; }, }; let request_result = crate::JsonRpcRequest::new(request_id, method, params); @@ -294,7 +396,7 @@ where std::result::Result::Ok(request) => request, std::result::Result::Err(error) => { let _ = response_tx.send(std::result::Result::Err(error)); - return true; + return WsActorIoOutcome::Continue; }, }; let payload_result = request.to_json_string(); @@ -302,11 +404,30 @@ where std::result::Result::Ok(payload) => payload, std::result::Result::Err(error) => { let _ = response_tx.send(std::result::Result::Err(error)); - return true; + return WsActorIoOutcome::Continue; }, }; - let deadline = tokio::time::Instant::now() + endpoint.session().command_timeout(); - pending.insert(request_id, PendingWsRequest { method, deadline, response_tx }); + let payload_size = payload.len(); + if payload_size > endpoint.session().max_message_size_bytes() || payload_size > endpoint.session().max_frame_size_bytes() { + let error = + ksp_core_lib::Error::new(crate::ERROR_CODE_INVALID_RPC_PARAMETERS, "WebSocket JSON-RPC request exceeds the configured outbound size bound") + .with_context("session_id", id.get().to_string()) + .with_context("method", method) + .with_context("payload_size_bytes", payload_size.to_string()) + .with_context("max_message_size_bytes", endpoint.session().max_message_size_bytes().to_string()) + .with_context("max_frame_size_bytes", endpoint.session().max_frame_size_bytes().to_string()); + let _ = response_tx.send(std::result::Result::Err(error)); + ksp_logging_lib::debug!( + target: crate::TRACING_TARGET, + session_id = id.get(), + method, + payload_size_bytes = payload_size, + max_message_size_bytes = endpoint.session().max_message_size_bytes(), + max_frame_size_bytes = endpoint.session().max_frame_size_bytes(), + "rejected oversized WebSocket JSON-RPC request before socket write" + ); + return WsActorIoOutcome::Continue; + } ksp_logging_lib::trace!( target: crate::TRACING_TARGET, session_id = id.get(), @@ -315,20 +436,65 @@ where pending_request_count = pending.len(), "sending WebSocket JSON-RPC request" ); - let send_result = websocket.send(tokio_tungstenite::tungstenite::Message::Text(payload.into())).await; - if send_result.is_err() { - ksp_logging_lib::warn!( - target: crate::TRACING_TARGET, - session_id = id.get(), - request_id, - method, - "WebSocket JSON-RPC request write failed" - ); - return false; + let message = tokio_tungstenite::tungstenite::Message::Text(payload.into()); + let send_result = tokio::select! { + biased; + shutdown_changed = shutdown_rx.changed() => { + let deadline = resolve_shutdown_deadline(shutdown_rx, shutdown_changed, endpoint.session().close_timeout()); + let error = ws_session_closed_error(id, "WebSocket JSON-RPC request cancelled by session shutdown").with_context("method", method); + let _ = response_tx.send(std::result::Result::Err(error)); + return WsActorIoOutcome::ShutdownRequested { deadline }; + }, + send_result = websocket.send(message) => send_result, + () = tokio::time::sleep(endpoint.session().command_timeout()) => { + let error = ws_timeout_error(id, "WebSocket JSON-RPC request write exceeded the configured command timeout").with_context("method", method); + let _ = response_tx.send(std::result::Result::Err(error)); + ksp_logging_lib::warn!( + target: crate::TRACING_TARGET, + session_id = id.get(), + request_id, + method, + "WebSocket JSON-RPC request write timed out" + ); + return WsActorIoOutcome::Failed { + code: crate::ERROR_CODE_WS_CONNECTION_FAILED, + pending_message: "WebSocket connection failed while a request write timed out", + }; + }, + }; + match send_result { + std::result::Result::Ok(()) => {}, + std::result::Result::Err(tokio_tungstenite::tungstenite::Error::WriteBufferFull(_)) + | std::result::Result::Err(tokio_tungstenite::tungstenite::Error::Capacity(_)) => { + let error = ksp_core_lib::Error::new(crate::ERROR_CODE_WS_BACKPRESSURE_OVERFLOW, "WebSocket write capacity is exhausted") + .with_context("session_id", id.get().to_string()) + .with_context("method", method); + let _ = response_tx.send(std::result::Result::Err(error)); + ksp_logging_lib::warn!( + target: crate::TRACING_TARGET, + session_id = id.get(), + request_id, + method, + "WebSocket request write capacity is exhausted" + ); + return WsActorIoOutcome::Continue; + }, + std::result::Result::Err(_) => { + let error = + ws_connection_error(id, endpoint, "WebSocket connection failed while writing a JSON-RPC request").with_context("method", method); + let _ = response_tx.send(std::result::Result::Err(error)); + ksp_logging_lib::warn!(target: crate::TRACING_TARGET, session_id = id.get(), request_id, method, "WebSocket JSON-RPC request write failed"); + return WsActorIoOutcome::Failed { + code: crate::ERROR_CODE_WS_CONNECTION_FAILED, + pending_message: "WebSocket connection failed while writing a request", + }; + }, } - return true; + let deadline = tokio::time::Instant::now() + endpoint.session().command_timeout(); + pending.insert(request_id, PendingWsRequest { method, deadline, response_tx }); + WsActorIoOutcome::Continue }, - } + }; } async fn handle_socket_message( @@ -337,70 +503,103 @@ async fn handle_socket_message( maybe_message: std::option::Option>, websocket: &mut tokio_tungstenite::WebSocketStream, pending: &mut std::collections::BTreeMap, -) -> bool + shutdown_rx: &mut tokio::sync::watch::Receiver>, +) -> WsActorIoOutcome where S: tokio::io::AsyncRead + tokio::io::AsyncWrite + std::marker::Unpin, { let message = match maybe_message { std::option::Option::Some(std::result::Result::Ok(message)) => message, - std::option::Option::Some(std::result::Result::Err(_)) => { + std::option::Option::Some(std::result::Result::Err(error)) => { + let code = classify_websocket_read_error(&error); ksp_logging_lib::warn!( target: crate::TRACING_TARGET, session_id = id.get(), endpoint_name = endpoint.name(), + protocol_error = code == crate::ERROR_CODE_WS_PROTOCOL_ERROR, "physical WebSocket read failed" ); - return false; + return WsActorIoOutcome::Failed { code, pending_message: "WebSocket read failed before pending response delivery" }; }, std::option::Option::None => { ksp_logging_lib::warn!( target: crate::TRACING_TARGET, session_id = id.get(), endpoint_name = endpoint.name(), - "physical WebSocket stream ended" + "physical WebSocket stream ended without a Close frame" ); - return false; + return WsActorIoOutcome::Failed { + code: crate::ERROR_CODE_WS_CONNECTION_FAILED, + pending_message: "WebSocket connection ended before pending response delivery", + }; }, }; return match message { tokio_tungstenite::tungstenite::Message::Text(text) => handle_text_message(id, text.as_str(), pending), tokio_tungstenite::tungstenite::Message::Binary(_) => { ksp_logging_lib::warn!(target: crate::TRACING_TARGET, session_id = id.get(), "received unexpected binary WebSocket message"); - false + WsActorIoOutcome::Failed { + code: crate::ERROR_CODE_WS_PROTOCOL_ERROR, + pending_message: "WebSocket wire data violated the expected JSON text protocol", + } }, - tokio_tungstenite::tungstenite::Message::Ping(payload) => { - ksp_logging_lib::trace!(target: crate::TRACING_TARGET, session_id = id.get(), "received WebSocket ping control frame"); - return websocket.send(tokio_tungstenite::tungstenite::Message::Pong(payload)).await.is_ok(); + tokio_tungstenite::tungstenite::Message::Ping(_) => { + ksp_logging_lib::trace!(target: crate::TRACING_TARGET, session_id = id.get(), "received WebSocket ping control frame; flushing automatic Pong"); + let flush_result = tokio::select! { + biased; + shutdown_changed = shutdown_rx.changed() => { + let deadline = resolve_shutdown_deadline(shutdown_rx, shutdown_changed, endpoint.session().close_timeout()); + return WsActorIoOutcome::ShutdownRequested { deadline }; + }, + flush_result = websocket.flush() => flush_result, + () = tokio::time::sleep(endpoint.session().command_timeout()) => { + return WsActorIoOutcome::Failed { + code: crate::ERROR_CODE_WS_CONNECTION_FAILED, + pending_message: "WebSocket connection failed while flushing automatic Pong", + }; + }, + }; + return match flush_result { + std::result::Result::Ok(()) => WsActorIoOutcome::Continue, + std::result::Result::Err(_) => WsActorIoOutcome::Failed { + code: crate::ERROR_CODE_WS_CONNECTION_FAILED, + pending_message: "WebSocket connection failed while flushing automatic Pong", + }, + }; }, tokio_tungstenite::tungstenite::Message::Pong(_) => { ksp_logging_lib::trace!(target: crate::TRACING_TARGET, session_id = id.get(), "received WebSocket pong control frame"); - true + WsActorIoOutcome::Continue }, tokio_tungstenite::tungstenite::Message::Close(_) => { ksp_logging_lib::debug!(target: crate::TRACING_TARGET, session_id = id.get(), "remote peer closed physical WebSocket session"); - false + let _ = tokio::time::timeout(endpoint.session().close_timeout(), websocket.flush()).await; + WsActorIoOutcome::RemoteClosed }, tokio_tungstenite::tungstenite::Message::Frame(_) => { ksp_logging_lib::trace!(target: crate::TRACING_TARGET, session_id = id.get(), "ignored internal WebSocket frame event"); - true + WsActorIoOutcome::Continue }, }; } -fn handle_text_message(id: crate::WsSessionId, text: &str, pending: &mut std::collections::BTreeMap) -> bool { +fn handle_text_message(id: crate::WsSessionId, text: &str, pending: &mut std::collections::BTreeMap) -> WsActorIoOutcome { let decoded = serde_json::from_str::(text); let value = match decoded { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => { ksp_logging_lib::warn!(target: crate::TRACING_TARGET, session_id = id.get(), "received malformed JSON WebSocket message"); - return false; + return WsActorIoOutcome::Failed { code: crate::ERROR_CODE_WS_PROTOCOL_ERROR, pending_message: "WebSocket JSON payload was malformed" }; }, }; let object = match value.as_object() { std::option::Option::Some(object) => object, std::option::Option::None => { ksp_logging_lib::warn!(target: crate::TRACING_TARGET, session_id = id.get(), "received non-object JSON WebSocket message"); - return false; + return WsActorIoOutcome::Failed { + code: crate::ERROR_CODE_WS_PROTOCOL_ERROR, + pending_message: "WebSocket JSON-RPC payload was not an object", + }; }, }; if !object.contains_key("id") { @@ -410,16 +609,22 @@ fn handle_text_message(id: crate::WsSessionId, text: &str, pending: &mut std::co session_id = id.get(), "received WebSocket notification before subscription registry activation; safely ignored" ); - return true; + return WsActorIoOutcome::Continue; } ksp_logging_lib::warn!(target: crate::TRACING_TARGET, session_id = id.get(), "received structurally invalid WebSocket JSON-RPC message"); - return false; + return WsActorIoOutcome::Failed { + code: crate::ERROR_CODE_WS_PROTOCOL_ERROR, + pending_message: "WebSocket JSON-RPC payload violated structural invariants", + }; } let response_id = match object.get("id").and_then(serde_json::Value::as_u64) { std::option::Option::Some(response_id) => response_id, std::option::Option::None => { ksp_logging_lib::warn!(target: crate::TRACING_TARGET, session_id = id.get(), "received WebSocket JSON-RPC response with invalid id"); - return false; + return WsActorIoOutcome::Failed { + code: crate::ERROR_CODE_WS_PROTOCOL_ERROR, + pending_message: "WebSocket JSON-RPC response contained an invalid id", + }; }, }; let pending_request = match pending.remove(&response_id) { @@ -431,24 +636,30 @@ fn handle_text_message(id: crate::WsSessionId, text: &str, pending: &mut std::co response_id, "ignored unknown or stale WebSocket JSON-RPC response id" ); - return true; + return WsActorIoOutcome::Continue; }, }; let parsed = crate::parse_json_rpc_response_value(value, response_id); - let (result, keep_running) = match parsed { + let (result, outcome) = match parsed { std::result::Result::Ok(response) => { let result = match response.into_result() { std::result::Result::Ok(value) => std::result::Result::Ok(value), std::result::Result::Err(error) => std::result::Result::Err(error.with_context("method", pending_request.method)), }; - (result, true) + (result, WsActorIoOutcome::Continue) }, std::result::Result::Err(_) => { let error = ksp_core_lib::Error::new(crate::ERROR_CODE_WS_PROTOCOL_ERROR, "WebSocket JSON-RPC response violates protocol invariants") .with_context("session_id", id.get().to_string()) .with_context("request_id", response_id.to_string()) .with_context("method", pending_request.method); - (std::result::Result::Err(error), false) + ( + std::result::Result::Err(error), + WsActorIoOutcome::Failed { + code: crate::ERROR_CODE_WS_PROTOCOL_ERROR, + pending_message: "WebSocket JSON-RPC response violated protocol invariants", + }, + ) }, }; let _ = pending_request.response_tx.send(result); @@ -460,7 +671,69 @@ fn handle_text_message(id: crate::WsSessionId, text: &str, pending: &mut std::co pending_request_count = pending.len(), "dispatched WebSocket JSON-RPC response to pending request" ); - return keep_running; + return outcome; +} + +async fn close_session_actor( + id: crate::WsSessionId, + endpoint: &crate::WsEndpointSettings, + snapshot_tx: &tokio::sync::watch::Sender, + websocket: &mut tokio_tungstenite::WebSocketStream, + pending: &mut std::collections::BTreeMap, + deadline: tokio::time::Instant, +) where + S: tokio::io::AsyncRead + tokio::io::AsyncWrite + std::marker::Unpin, +{ + publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Closing, pending.len()); + fail_all_pending(pending, id, crate::ERROR_CODE_WS_SESSION_CLOSED, "WebSocket session shutdown cancelled the pending request"); + publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Closing, 0); + ksp_logging_lib::debug!( + target: crate::TRACING_TARGET, + session_id = id.get(), + endpoint_name = endpoint.name(), + "closing physical WebSocket session" + ); + let now = tokio::time::Instant::now(); + let remaining = deadline.saturating_duration_since(now); + let io_deadline = now + (remaining / 2); + let close_wait = tokio::time::timeout_at(io_deadline, websocket.send(tokio_tungstenite::tungstenite::Message::Close(std::option::Option::None))).await; + match close_wait { + std::result::Result::Ok(std::result::Result::Ok(())) => { + ksp_logging_lib::trace!(target: crate::TRACING_TARGET, session_id = id.get(), "sent WebSocket Close control frame"); + }, + std::result::Result::Ok(std::result::Result::Err(_)) => { + ksp_logging_lib::debug!(target: crate::TRACING_TARGET, session_id = id.get(), "best-effort WebSocket Close frame could not be sent"); + }, + std::result::Result::Err(_) => { + ksp_logging_lib::debug!(target: crate::TRACING_TARGET, session_id = id.get(), "best-effort WebSocket Close frame reached close deadline"); + }, + } + publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Closed, 0); + ksp_logging_lib::debug!(target: crate::TRACING_TARGET, session_id = id.get(), endpoint_name = endpoint.name(), "physical WebSocket session is closed"); +} + +fn classify_websocket_read_error(error: &tokio_tungstenite::tungstenite::Error) -> ksp_core_lib::ErrorCode { + return match error { + tokio_tungstenite::tungstenite::Error::Capacity(_) + | tokio_tungstenite::tungstenite::Error::Protocol(_) + | tokio_tungstenite::tungstenite::Error::Utf8(_) + | tokio_tungstenite::tungstenite::Error::AttackAttempt => crate::ERROR_CODE_WS_PROTOCOL_ERROR, + _ => crate::ERROR_CODE_WS_CONNECTION_FAILED, + }; +} + +fn resolve_shutdown_deadline( + shutdown_rx: &tokio::sync::watch::Receiver>, + shutdown_changed: std::result::Result<(), tokio::sync::watch::error::RecvError>, + close_timeout: std::time::Duration, +) -> tokio::time::Instant { + if shutdown_changed.is_ok() { + let requested = shutdown_rx.borrow().to_owned(); + if let std::option::Option::Some(deadline) = requested { + return deadline; + } + } + return tokio::time::Instant::now() + close_timeout; } fn next_pending_deadline(pending: &std::collections::BTreeMap) -> tokio::time::Instant { @@ -500,6 +773,26 @@ fn expire_pending_requests(id: crate::WsSessionId, pending: &mut std::collection } } +fn prune_cancelled_pending(id: crate::WsSessionId, pending: &mut std::collections::BTreeMap) { + let mut cancelled_ids = std::vec::Vec::new(); + for (request_id, request) in pending.iter() { + if request.response_tx.is_closed() { + cancelled_ids.push(*request_id); + } + } + for request_id in cancelled_ids { + if let std::option::Option::Some(request) = pending.remove(&request_id) { + ksp_logging_lib::trace!( + target: crate::TRACING_TARGET, + session_id = id.get(), + request_id, + method = request.method, + "removed abandoned WebSocket JSON-RPC request after caller cancellation" + ); + } + } +} + fn fail_all_pending( pending: &mut std::collections::BTreeMap, id: crate::WsSessionId, diff --git a/crates/ksp-onchain-transport-lib/tests/public_api.rs b/crates/ksp-onchain-transport-lib/tests/public_api.rs index 8749282..752f5f1 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: 26 +// version: 27 //! Integration tests for the public `ksp-onchain-transport-lib` consumer contract. @@ -559,3 +559,10 @@ fn public_v0_2_7_pre_004_physical_websocket_session_contract_is_available_from_c assert_eq!(ksp_onchain_transport_lib::ERROR_CODE_WS_PROTOCOL_ERROR.code(), "ws_protocol_error"); assert_eq!(ksp_onchain_transport_lib::ERROR_CODE_WS_SESSION_CLOSED.code(), "ws_session_closed"); } + +#[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; + assert_eq!(ksp_onchain_transport_lib::WsSessionState::Closing, ksp_onchain_transport_lib::WsSessionState::Closing); + assert_eq!(ksp_onchain_transport_lib::WsSessionState::Closed, ksp_onchain_transport_lib::WsSessionState::Closed); +} 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 8bfa2ac..0c16021 100644 --- a/crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs +++ b/crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs @@ -1,10 +1,14 @@ // file: crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs -// version: 1 +// version: 2 use futures_util::SinkExt; // rust-rules: trait-import use futures_util::StreamExt; // rust-rules: trait-import fn local_endpoint(url: &str) -> crate::WsEndpointSettings { + return local_endpoint_with_session(url, crate::WsSessionSettings::default()); +} + +fn local_endpoint_with_session(url: &str, session: crate::WsSessionSettings) -> crate::WsEndpointSettings { return crate::WsEndpointSettings::new( "local_ws", true, @@ -12,7 +16,31 @@ fn local_endpoint(url: &str) -> crate::WsEndpointSettings { crate::WsClusterName::new("local"), crate::WsProtocolKind::SolanaStandard, crate::WsEndpointUrl::parse(url).expect("local test WebSocket URL must parse"), - crate::WsSessionSettings::default(), + session, + ); +} + +fn session_settings( + command_timeout: std::time::Duration, + close_timeout: std::time::Duration, + max_pending_requests: usize, + max_message_size_bytes: usize, + max_frame_size_bytes: usize, + max_write_buffer_size_bytes: usize, +) -> crate::WsSessionSettings { + let defaults = crate::WsSessionSettings::default(); + return crate::WsSessionSettings::new( + command_timeout, + close_timeout, + defaults.reconnect().clone(), + defaults.resubscribe(), + defaults.command_queue_capacity(), + defaults.notification_queue_capacity(), + defaults.max_active_subscriptions(), + max_pending_requests, + max_message_size_bytes, + max_frame_size_bytes, + max_write_buffer_size_bytes, ); } @@ -34,6 +62,28 @@ async fn send_result(websocket: &mut tokio_tungstenite::WebSocketStream assert_eq!(payload.as_ref(), b"ping-canary"), + other => panic!("expected Pong frame, got {other:?}"), + } + let _ = pong_tx.send(()); + let request = read_request(&mut websocket).await; + send_result(&mut websocket, &request, serde_json::json!(true)).await; + }); + let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed"); + tokio::time::timeout(std::time::Duration::from_secs(1), pong_rx) + .await + .expect("automatic Pong must remain bounded") + .expect("fixture server must observe Pong"); + assert_eq!(session.state(), crate::WsSessionState::Active); + let result = session.execute_json_rpc("afterPing", std::vec::Vec::new()).await.expect("session must remain usable after Ping/Pong"); + assert_eq!(result, serde_json::json!(true)); + server.await.expect("local server task must complete"); +} + +#[tokio::test(flavor = "current_thread")] +async fn websocket_remote_close_transitions_to_closed_instead_of_failed() { + let (listener, url) = bind_local_listener().await; + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.expect("local server must accept client"); + let mut websocket = tokio_tungstenite::accept_async(stream).await.expect("local WebSocket handshake must succeed"); + websocket.send(tokio_tungstenite::tungstenite::Message::Close(std::option::Option::None)).await.expect("fixture Close must send"); + let _ = tokio::time::timeout(std::time::Duration::from_millis(200), websocket.next()).await; + }); + let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed"); + wait_for_state(&session, crate::WsSessionState::Closed).await; + assert!(session.close().await.is_ok()); + server.await.expect("local server task must complete"); +} + +#[tokio::test(flavor = "current_thread")] +async fn websocket_explicit_close_cancels_pending_and_finishes_with_hostile_peer() { + let (listener, url) = bind_local_listener().await; + let (request_seen_tx, request_seen_rx) = tokio::sync::oneshot::channel(); + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.expect("local server must accept client"); + let mut websocket = tokio_tungstenite::accept_async(stream).await.expect("local WebSocket handshake must succeed"); + let _ = read_request(&mut websocket).await; + let _ = request_seen_tx.send(()); + tokio::time::sleep(std::time::Duration::from_millis(500)).await; + }); + let settings = session_settings(std::time::Duration::from_secs(1), std::time::Duration::from_millis(150), 8, 64 * 1024, 64 * 1024, 64 * 1024); + let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed"); + let request_session = session.clone(); + let pending = tokio::spawn(async move { + return request_session.execute_json_rpc("neverRespond", std::vec::Vec::new()).await; + }); + request_seen_rx.await.expect("hostile server must observe request"); + wait_for_pending_count(&session, 1).await; + let close_result = tokio::time::timeout(std::time::Duration::from_millis(300), session.close()).await; + close_result.expect("explicit close must remain bounded").expect("explicit close must finish cleanly"); + assert_eq!(session.state(), crate::WsSessionState::Closed); + assert_eq!(session.snapshot().pending_request_count(), 0); + let pending_error = pending.await.expect("pending request task must join").expect_err("shutdown must cancel pending request"); + assert_eq!(pending_error.code(), crate::ERROR_CODE_WS_SESSION_CLOSED); + server.await.expect("hostile server task must complete"); +} + +#[tokio::test(flavor = "current_thread")] +async fn websocket_repeated_connect_close_cycles_are_bounded() { + const CYCLES: usize = 8; + let (listener, url) = bind_local_listener().await; + let server = tokio::spawn(async move { + for _ in 0..CYCLES { + let (stream, _) = listener.accept().await.expect("local server must accept client"); + let mut websocket = tokio_tungstenite::accept_async(stream).await.expect("local WebSocket handshake must succeed"); + let message = websocket.next().await.expect("client Close must arrive").expect("client Close must decode"); + assert!(matches!(message, tokio_tungstenite::tungstenite::Message::Close(_))); + let _ = websocket.flush().await; + } + }); + let settings = session_settings(std::time::Duration::from_secs(1), std::time::Duration::from_millis(200), 8, 64 * 1024, 64 * 1024, 64 * 1024); + for _ in 0..CYCLES { + let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings.clone())).await.expect("client must connect"); + session.close().await.expect("client close must succeed"); + assert_eq!(session.state(), crate::WsSessionState::Closed); + } + tokio::time::timeout(std::time::Duration::from_secs(2), server) + .await + .expect("repeated close server must remain bounded") + .expect("repeated close server task must join"); +} diff --git a/deltas/0.2.7/pre.005.md b/deltas/0.2.7/pre.005.md new file mode 100644 index 0000000..3e005dd --- /dev/null +++ b/deltas/0.2.7/pre.005.md @@ -0,0 +1,172 @@ + + + +# Delta `0.2.7-pre.005` — limits adversariales + control frames + shutdown borné + +## 1. Base requise + +```text +0.2.7-pre.004-fix.001 appliqué +workspace.package.version = 0.2.7-pre.4.fix.1 +``` + +Le checkpoint opérateur reçu avant cette tranche est vert : `cargo fmt --all`, audit Python, `cargo check --workspace`, `cargo clippy --workspace --all-targets` et `cargo test --workspace`. Les 261 tests unitaires Transport de `pre.004-fix.001` passent également. + +## 2. Signal technique + +Cette prerelease modifie le runtime Rust, les tests et la documentation. Conformément aux règles KSP : + +```text +livraison = 0.2.7-pre.005 +workspace.package.version = 0.2.7-pre.5 +commit = v0.2.7-pre.005 +``` + +Aucun tag prerelease. + +## 3. Shutdown explicite `WsSession::close()` + +Nouvelle surface publique : + +```text +WsSession::close().await +``` + +Le signal de shutdown est transporté par un canal `watch` distinct de la command queue. Une command queue saturée ne peut donc pas empêcher la demande de fermeture. + +Lifecycle : + +```text +Active -> Closing -> Closed +remote Close propre -> Closed +I/O ou violation protocolaire -> Failed +``` + +`close()` : + +1. publie la demande de shutdown avec une deadline issue de `close_timeout` ; +2. l'actor passe `Closing` ; +3. toutes les requests JSON-RPC pending reçoivent `ws_session_closed` ; +4. l'actor tente un Close WebSocket best-effort sous la deadline ; +5. l'actor publie `Closed` ; +6. le caller attend `Closed` sans pouvoir rester bloqué indéfiniment. + +Le reconnect reste absent jusqu'à `pre.007`. + +## 4. Control frames + +Ping/Pong/Close sont maintenant testés explicitement. + +Pour Ping, Tungstenite met automatiquement le Pong correspondant en file lors de la lecture ; KSP force le `flush` de cette réponse automatique. Aucun Pong duplicatif et aucun heartbeat applicatif périodique ne sont ajoutés. + +Un Close distant propre est distingué d'une rupture I/O : il termine la session en `Closed`, alors qu'une rupture physique ou une violation de protocole termine en `Failed`. + +## 5. Bornes frame/message/write + +Les limites configurées depuis `pre.002` et raccordées à Tungstenite depuis `pre.004` sont désormais couvertes par des fixtures adversariales : + +```text +max_message_size +max_frame_size +max_write_buffer_size +``` + +Un payload JSON-RPC outbound dépassant `max_message_size` ou `max_frame_size` est rejeté avant écriture avec `invalid_rpc_parameters`, sans teardown d'une session saine. + +Un frame/message inbound surdimensionné est rejeté par Tungstenite avant que KSP tente le parse JSON ; l'anomalie est classée protocolaire et la session devient `Failed`. + +Aucune limite legacy Solana de type 1232 bytes n'est introduite. + +## 6. Pending requests et cancellation + +Le runtime `pre.005` ajoute les preuves suivantes : + +- `max_pending_requests` rejette uniquement la request excédentaire avec `ws_backpressure_overflow` ; +- une request pending silencieuse expire selon `command_timeout`, reçoit `timeout` et libère immédiatement sa capacité ; +- une fermeture explicite annule toutes les pending requests avec `ws_session_closed` ; +- un receiver caller abandonné est purgé du pending map au prochain cycle actor, et reste de toute façon borné par la deadline de request. + +La cancellation de subscriptions et les races unsubscribe/reconnect restent dans `pre.006`/`pre.007`. + +## 7. Socket I/O borné et shutdown prioritaire + +Les writes JSON-RPC et le flush Pong surveillent le signal shutdown pendant leur attente. La fermeture peut donc préempter une opération socket qui serait autrement bloquante. + +Un write JSON-RPC dépassant `command_timeout` provoque une erreur caller bornée et termine la session physique, car l'état de l'écriture devient impropre à une poursuite sûre. + +## 8. Fixtures adversariales locales + +Les tests Transport ajoutent notamment : + +```text +pending capacity = 1 avec deuxième request rejetée isolément +oversized outbound request rejetée avant write +oversized inbound frame/message rejeté avant JSON +pending request timeout + purge +Ping -> Pong + session toujours Active +remote Close -> Closed +close explicite + pending cancellation + peer hostile +8 cycles connect/close bornés +``` + +Toutes les fixtures restent locales et déterministes ; aucun réseau Solana réel n'est requis. + +## 9. Logging / sécurité + +Toutes les émissions passent exclusivement par `ksp-logging-lib` et réutilisent : + +```text +TRACING_TARGET = "ksp-onchain-transport-lib" +``` + +Les nouveaux diagnostics ne loggent que session ID, endpoint logique, méthode statique, compteurs et tailles numériques. Les URLs, credentials, query tokens et payloads JSON restent absents. + +Les erreurs Tungstenite brutes ne sont pas projetées dans les erreurs KSP afin d'éviter une fuite de données provider. + +## 10. Public API et documentation + +Le canari public API vérifie désormais `WsSession::close` et les états `Closing`/`Closed`. + +Synchronisés : + +```text +crates/ksp-onchain-transport-lib/README.md +crates/ksp-onchain-transport-lib/USAGE.md +docs/plans/014-V0_2_7_ONCHAIN_WEBSOCKET_PLAN.md +docs/validation/010-V0_2_7_ONCHAIN_WEBSOCKET.md +``` + +`ROADMAP.md` et `CHANGELOG.md` restent inchangés pendant la série prerelease. + +## 11. Validation de préparation + +Le sandbox de génération doit au minimum vérifier : + +```text +python3 scripts/audit_rust_workspace_rules.py +absence de tracing direct +présence du TRACING_TARGET existant +absence de question-mark operator dans le runtime modifié +overlay exact du delta +``` + +Cargo n'est pas disponible dans le sandbox de génération ; aucun résultat Cargo local n'est revendiqué. + +## 12. Gates opérateur avant commit + +```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 --workspace +``` + +Si le checkpoint est vert : + +```text +commit = v0.2.7-pre.005 +``` + +La tranche suivante est `0.2.7-pre.006` : registry subscriptions, IDs locaux stables, moteur generic subscribe/unsubscribe et channels typed bornés. diff --git a/docs/plans/014-V0_2_7_ONCHAIN_WEBSOCKET_PLAN.md b/docs/plans/014-V0_2_7_ONCHAIN_WEBSOCKET_PLAN.md index 75efee2..8f15312 100644 --- a/docs/plans/014-V0_2_7_ONCHAIN_WEBSOCKET_PLAN.md +++ b/docs/plans/014-V0_2_7_ONCHAIN_WEBSOCKET_PLAN.md @@ -1,9 +1,9 @@ - + # Plan `0.2.7` — WebSocket Solana standard -> **Statut : actif, `0.2.7-pre.004`.** `pre.003` est validé opérateur. `pre.004` matérialise les dépendances WebSocket, la première session physique actor-owned, le handshake/read/write JSON-RPC, la map bornée de requests en attente et les fixtures serveur local. Les subscriptions typed restent différées. +> **Statut : actif, `0.2.7-pre.005`.** `pre.004-fix.001` est validé opérateur. `pre.005` durcit la session physique avec limites adversariales, control frames, purge des pending requests et shutdown explicite borné. Les subscriptions typed restent différées. ## 1. Objet et base vérifiée @@ -666,6 +666,10 @@ Si la session est déjà reconnecting, elle ne se reconnecte jamais seulement po `Drop` peut déclencher un signal best-effort, mais ne remplace pas l'API async explicite pour les garanties de lifecycle. +Implémentation `pre.005` : le signal de shutdown est hors de la command queue afin qu'une saturation de celle-ci ne puisse pas empêcher la fermeture. `close()` utilise `close_timeout` comme budget caller, l'actor annule d'abord les pending requests, tente le Close WebSocket sous ce budget, puis publie `Closed`. Un Close distant propre mène également à `Closed`; une erreur I/O ou protocolaire reste terminale `Failed`. + +Les Ping reçus s'appuient sur le Pong automatiquement mis en file par Tungstenite et forcent son flush. Aucun second Pong applicatif ni heartbeat périodique n'est ajouté. + ## 14. Erreurs et anomalies de protocole Réutiliser lorsque possible les codes HTTP/JSON-RPC génériques existants : @@ -871,7 +875,7 @@ pre.001 audit interne/externe + matrice 18 méthodes + bot3 + dependencies + th pre.002 DONE — settings WS Transport + URL redaction + IDs/states/snapshots + tests de settings pre.003 DONE — std.transport V2 HTTP+WS + backward V1 + discriminateur WS + schema/fixtures + Config -> WsTransportSettings pre.004 DONE — deps tokio-tungstenite/futures-util + actor physique + handshake/read/write + pending JSON-RPC + serveur local -pre.005 limites frame/message/request + control frames + cancellation/close/shutdown + adversarial socket tests +pre.005 DONE — limites frame/message/request + control frames + cancellation/close/shutdown + adversarial socket tests pre.006 registry subscriptions + IDs locaux + generic subscribe/unsubscribe engine + channels typed bounded pre.007 reconnect borné + resubscribe déterministe + continuity gap + races unsubscribe/reconnect pre.008 backpressure per-sub + overflow/limits + leak/lifecycle adversarial tests diff --git a/docs/validation/010-V0_2_7_ONCHAIN_WEBSOCKET.md b/docs/validation/010-V0_2_7_ONCHAIN_WEBSOCKET.md index 1aa003b..9c82f27 100644 --- a/docs/validation/010-V0_2_7_ONCHAIN_WEBSOCKET.md +++ b/docs/validation/010-V0_2_7_ONCHAIN_WEBSOCKET.md @@ -1,9 +1,9 @@ - + # Validation `0.2.7` — WebSocket Solana standard -> **Statut : matrice active, `0.2.7-pre.004`.** Settings/lifecycle `pre.002`, Config V2 `pre.003` et première session physique actor-owned `pre.004` sont matérialisés. Registry subscriptions, shutdown complet, reconnect et wrappers typed restent ouverts. +> **Statut : matrice active, `0.2.7-pre.005`.** Settings/lifecycle `pre.002`, Config V2 `pre.003`, session physique `pre.004` et durcissement limits/control/shutdown `pre.005` sont matérialisés. Registry subscriptions, reconnect et wrappers typed restent ouverts. ## 1. Baseline normative @@ -358,6 +358,25 @@ erreur de connexion sans URL/credential dans Error Debug Les limites `max_message_size`, `max_frame_size` et `max_write_buffer_size` sont raccordées à `tungstenite::WebSocketConfig`. Les fixtures oversized et le shutdown/control-frame hostile restent explicitement `pre.005`. +## 9.3 Checkpoint limits/control/shutdown `pre.005` + +Gates déterministes ajoutés : + +```text +max_pending_requests saturé -> seule la request excédentaire est rejetée +outbound JSON-RPC > borne message/frame -> rejet avant socket write, session Active +inbound frame/message oversized -> rejet Tungstenite avant parse JSON, session Failed +pending request silencieuse -> timeout + purge de capacité, session Active +Ping distant -> Pong automatique flushé, session Active +Close distant propre -> Closed et non Failed +WsSession::close() -> Closing -> Closed +close() annule les pending requests avec ws_session_closed +peer hostile qui ne répond pas au Close -> shutdown borné +cycles connect/close répétés -> terminaison bornée +``` + +Le signal de shutdown est indépendant de la command queue et les opérations socket longues de l'actor surveillent ce signal. Aucun reconnect/resubscribe n'est activé par cette tranche. + ## 10. Validation du gate `pre.001` Exécuté dans le sandbox :