diff --git a/Cargo.toml b/Cargo.toml index 2cf765c..f47edec 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,12 @@ # file: Cargo.toml -# version: 195 +# version: 196 [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.3" +version = "0.2.7-pre.4" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" @@ -21,6 +21,7 @@ ed25519-dalek = { version = "^3.0", default-features = false } getrandom = { version = "^0.4", default-features = false } base64 = { version = "^0.23" } fs2 = { version = "^0.4" } +futures-util = { version = "^0.3", default-features = false } serde = { version = "^1.0" } serde_json = { version = "^1.0" } jsonschema = { version = "^0.50", default-features = false } @@ -31,6 +32,7 @@ tracing = { version = "^0.1", default-features = false } tracing-subscriber = { version = "^0.3", default-features = false } tracing-appender = { version = "^0.2", default-features = false } tokio = { version = "^1.53", default-features = false } +tokio-tungstenite = { version = "^0.30", default-features = false } tempfile = { version = "^3.27" } chrono = { version = "^0.4", default-features = false } tauri = { version = "^2.11" } diff --git a/crates/ksp-core-lib/tests/workspace_dependencies.rs b/crates/ksp-core-lib/tests/workspace_dependencies.rs index 308bc24..c155074 100644 --- a/crates/ksp-core-lib/tests/workspace_dependencies.rs +++ b/crates/ksp-core-lib/tests/workspace_dependencies.rs @@ -1,5 +1,5 @@ // file: crates/ksp-core-lib/tests/workspace_dependencies.rs -// version: 1 +// version: 2 //! Workspace-level dependency policy canaries owned by the foundational KSP test surface. @@ -61,8 +61,10 @@ fn transport_manifest_preserves_ksp_dependency_firewall() { } assert!(manifest.contains("ksp-core-lib")); assert!(manifest.contains("ksp-logging-lib")); + assert!(manifest.contains("futures-util = { workspace = true, features = [\"sink\", \"std\"] }")); assert!(manifest.contains("reqwest = { workspace = true, features = [\"rustls\"] }")); - assert!(manifest.contains("tokio = { workspace = true, features = [\"macros\", \"sync\", \"time\"] }")); + assert!(manifest.contains("tokio = { workspace = true, features = [\"macros\", \"rt\", \"sync\", \"time\"] }")); + assert!(manifest.contains("tokio-tungstenite = { workspace = true, features = [\"connect\", \"rustls-tls-webpki-roots\"] }")); assert!(manifest.contains("[dev-dependencies]")); - assert!(manifest.contains("tokio = { workspace = true, features = [\"rt\"] }")); + assert!(manifest.contains("tokio = { workspace = true, features = [\"net\", \"rt\"] }")); } diff --git a/crates/ksp-onchain-transport-lib/Cargo.toml b/crates/ksp-onchain-transport-lib/Cargo.toml index 6c154b8..84a3628 100644 --- a/crates/ksp-onchain-transport-lib/Cargo.toml +++ b/crates/ksp-onchain-transport-lib/Cargo.toml @@ -1,5 +1,5 @@ # file: crates/ksp-onchain-transport-lib/Cargo.toml -# version: 3 +# version: 4 [package] name = "ksp-onchain-transport-lib" @@ -10,13 +10,15 @@ repository.workspace = true [dependencies] ksp-core-lib = { path = "../ksp-core-lib" } ksp-logging-lib = { path = "../ksp-logging-lib" } +futures-util = { workspace = true, features = ["sink", "std"] } reqwest = { workspace = true, features = ["rustls"] } serde = { workspace = true, features = ["derive"] } serde_json.workspace = true -tokio = { workspace = true, features = ["macros", "sync", "time"] } +tokio = { workspace = true, features = ["macros", "rt", "sync", "time"] } +tokio-tungstenite = { workspace = true, features = ["connect", "rustls-tls-webpki-roots"] } [dev-dependencies] -tokio = { workspace = true, features = ["rt"] } +tokio = { workspace = true, features = ["net", "rt"] } [lints] workspace = true diff --git a/crates/ksp-onchain-transport-lib/README.md b/crates/ksp-onchain-transport-lib/README.md index db088c2..22d71d4 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` @@ -35,6 +35,7 @@ ksp-config-lib -> ksp-core-lib -> ksp-logging-lib -> reqwest / tokio / serde + -> tokio-tungstenite / futures-util ``` La direction inverse est interdite : @@ -85,6 +86,25 @@ total : 52 Les 14 méthodes historiques restent découvrables pour la compliance mais sont `Removed` et ne sont pas simulées comme appelables. +## Foundation WebSocket `0.2.7-pre.004` + +La première session physique WebSocket est matérialisée sans introduire de pool/scheduler automatique ni de registry de subscriptions anticipé. + +`WsSession::connect(WsEndpointSettings)` : + +- ouvre exactement une connexion physique pour un appel ; +- confie le socket à une tâche actor unique ; +- sérialise les commandes internes par un canal `mpsc` borné ; +- maintient une map bornée de requests JSON-RPC en attente ; +- publie `WsSessionSnapshot` via un état compact `watch` ; +- applique aux sockets les plafonds KSP de message, frame et write buffer ; +- ne projette jamais l'URL dans `Debug`, snapshot, erreurs KSP ou logs ; +- répond aux `Ping` reçus et tolère les `Pong`; le lifecycle complet `Close`/shutdown reste le gate `pre.005`. + +Le chemin JSON-RPC générique reste `pub(crate)` en `pre.004`. Il sert de primitive au futur moteur typed de subscriptions et **ne constitue pas une API publique raw provider-extension**. Le registry subscriptions, les IDs serveur, reconnect/resubscribe et backpressure par subscription restent respectivement dans les tranches prévues. + +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. + ## Résilience L'admission est calculée par couple endpoint/rôle. Le pool applique : @@ -148,4 +168,4 @@ Les deux sont `ignored` par défaut. Le smoke Transport appartient durablement - [`../../docs/validation/006-V0_2_3_HTTP_TRANSACTIONS.md`](../../docs/validation/006-V0_2_3_HTTP_TRANSACTIONS.md) — matrice finale validée `0.2.3` ; - [`../../docs/plans/011-V0_2_4_HTTP_BLOCKS_ECONOMICS_PLAN.md`](../../docs/plans/011-V0_2_4_HTTP_BLOCKS_ECONOMICS_PLAN.md) — plan Blocks/Economics et compliance HTTP finale ; - [`../../docs/validation/007-V0_2_4_HTTP_FINAL_COMPLIANCE.md`](../../docs/validation/007-V0_2_4_HTTP_FINAL_COMPLIANCE.md) — matrice finale validée `52/52 + 14/14` et audit `KSP-TRANSPORT-007` global ; -- [`../../config/std.transport.json`](../../config/std.transport.json) — configuration standard HTTP. +- [`../../config/std.transport.json`](../../config/std.transport.json) — configuration standard HTTP + WebSocket V2, avec lecture backward V1 HTTP-only. diff --git a/crates/ksp-onchain-transport-lib/USAGE.md b/crates/ksp-onchain-transport-lib/USAGE.md index 9c3d522..55689a5 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` @@ -65,7 +65,38 @@ let pool = match ksp_onchain_transport_lib::HttpTransportPool::new(resolved.into Le document standard peut contenir une URL provenant d'un `KSP_SECRET_*`. La valeur réelle est transmise au runtime, mais les projections sûres et `Debug` restent redacted. -## 3. Appels typés +## 3. Session physique WebSocket + +À partir de `0.2.7-pre.004`, un consumer peut créer explicitement une session physique : + +```rust +let ws_url = match ksp_onchain_transport_lib::WsEndpointUrl::parse("wss://api.devnet.solana.com") { + Ok(value) => value, + Err(error) => return Err(error), +}; +let endpoint = ksp_onchain_transport_lib::WsEndpointSettings::new( + "devnet_public", + true, + ksp_onchain_transport_lib::WsProviderName::new("solana-public"), + ksp_onchain_transport_lib::WsClusterName::new("devnet"), + ksp_onchain_transport_lib::WsProtocolKind::SolanaStandard, + ws_url, + ksp_onchain_transport_lib::WsSessionSettings::default(), +); +let session = match ksp_onchain_transport_lib::WsSession::connect(endpoint).await { + Ok(value) => value, + Err(error) => return Err(error), +}; +let snapshot = session.snapshot(); +``` + +Deux appels `WsSession::connect` avec le même endpoint créent volontairement deux connexions physiques distinctes. Il n'existe encore aucun pool de sessions automatique. + +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. + +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 Les wrappers typés se trouvent directement sur `HttpTransportPool`. @@ -125,7 +156,7 @@ let stake_minimum = pool.get_stake_minimum_delegation(&role, Some(&context)).awa `getBlock` possède également une forme bare-encoding legacy séparée et deprecated. Les valeurs Economics restent celles du runtime : le consumer ne doit pas supposer localement un taux d'inflation ou un minimum de délégation constant. -## 4. Exécution JSON-RPC standard générique +## 5. Exécution JSON-RPC standard générique Une méthode courante auditée peut être appelée via son descriptor : @@ -139,7 +170,7 @@ Cette API retourne un `serde_json::Value`. Elle reste utile pour les extensions Avant exécution, `ensure_runtime_supported()` est appliqué. Une méthode historique `Removed` retourne `ERROR_CODE_METHOD_REMOVED` au lieu d'émettre un appel réseau fictif. -## 5. Sélection et admission sans exécuter la requête +## 6. Sélection et admission sans exécuter la requête Pour inspecter le routing : @@ -154,13 +185,13 @@ Dans le même bloc, `acquire_for_method()` réserve réellement la capacité RPS `HttpRequestPermit` détient la capacité de concurrence jusqu'à sa destruction. Aucun verrou synchrone n'est conservé pendant l'attente réseau. -## 6. Snapshots runtime +## 7. Snapshots runtime `HttpTransportPool::snapshot()` fournit une vue sûre des endpoints/rôles : disponibilité, limites, requêtes en vol, cooldown restant et compteurs runtime. Les URLs d'endpoint n'y apparaissent jamais. -## 7. Retry et write submissions +## 8. Retry et write submissions La policy de retry est portée par la metadata des méthodes et `evaluate_transport_retry()`. @@ -168,7 +199,7 @@ Les reads/simulations classés `RetrySafe` peuvent être réessayés dans le bud Pour une opération `WriteSubmission / NeverAfterDispatch`, un timeout ou autre résultat ambigu après dispatch arrête la resoumission automatique. Le consumer métier ne doit pas contourner cette protection avec une boucle de retry externe aveugle. -## 8. Logging +## 9. Logging Les événements Transport utilisent le target : @@ -180,7 +211,7 @@ Ne jamais journaliser l'URL complète, un token provider, un body massif, une tr La configuration standard route les événements `info` de Transport vers un fichier dédié. Pour une investigation temporaire, élever uniquement ce target/sink à `debug` ou `trace`, puis revenir à `info` avant clôture du développement. -## 9. Smokes Devnet opt-in +## 10. Smokes Devnet opt-in Le smoke **Transport pur** construit ses settings programmatiquement et exerce un sous-ensemble représentatif d'Accounts/Tokens/Cluster, trois reads Transactions, puis des reads Blocks/Economics de la release stable `0.2.4` : diff --git a/crates/ksp-onchain-transport-lib/src/error.rs b/crates/ksp-onchain-transport-lib/src/error.rs index 7712654..0c50135 100644 --- a/crates/ksp-onchain-transport-lib/src/error.rs +++ b/crates/ksp-onchain-transport-lib/src/error.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/src/error.rs -// version: 3 +// version: 4 /// Error code used when no logical endpoint can satisfy a request. pub const ERROR_CODE_ENDPOINT_SELECTION_FAILED: ksp_core_lib::ErrorCode = ksp_core_lib::ErrorCode::new("onchain_transport", "endpoint_selection_failed"); @@ -27,3 +27,11 @@ pub const ERROR_CODE_RATE_LIMITED: ksp_core_lib::ErrorCode = ksp_core_lib::Error pub const ERROR_CODE_RPC_APPLICATION_ERROR: ksp_core_lib::ErrorCode = ksp_core_lib::ErrorCode::new("onchain_transport", "rpc_application_error"); /// Error code used when a transport deadline expires. pub const ERROR_CODE_TIMEOUT: ksp_core_lib::ErrorCode = ksp_core_lib::ErrorCode::new("onchain_transport", "timeout"); +/// Error code used when a bounded WebSocket runtime queue or pending-request capacity is exhausted. +pub const ERROR_CODE_WS_BACKPRESSURE_OVERFLOW: ksp_core_lib::ErrorCode = ksp_core_lib::ErrorCode::new("onchain_transport", "ws_backpressure_overflow"); +/// Error code used when a physical WebSocket connection or handshake fails. +pub const ERROR_CODE_WS_CONNECTION_FAILED: ksp_core_lib::ErrorCode = ksp_core_lib::ErrorCode::new("onchain_transport", "ws_connection_failed"); +/// Error code used when WebSocket wire data violates the KSP protocol contract. +pub const ERROR_CODE_WS_PROTOCOL_ERROR: ksp_core_lib::ErrorCode = ksp_core_lib::ErrorCode::new("onchain_transport", "ws_protocol_error"); +/// Error code used when a WebSocket session is no longer available to a caller. +pub const ERROR_CODE_WS_SESSION_CLOSED: ksp_core_lib::ErrorCode = ksp_core_lib::ErrorCode::new("onchain_transport", "ws_session_closed"); diff --git a/crates/ksp-onchain-transport-lib/src/lib.rs b/crates/ksp-onchain-transport-lib/src/lib.rs index 4a9acf6..508d16e 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: 22 +// version: 23 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -17,7 +17,8 @@ //! modern/legacy `getBlock`, positional inflation rewards, runtime-provided economics values and the final `KSP-TRANSPORT-007` compliance target. //! 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. Physical sockets and subscription execution are intentionally deferred to later prereleases. +//! 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. mod client; mod constants; @@ -37,6 +38,7 @@ mod rpc_tokens; mod rpc_transactions; mod settings; mod ws_lifecycle; +mod ws_session; mod ws_settings; /// Passive runtime availability reported for one logical HTTP endpoint. @@ -73,6 +75,14 @@ pub use self::error::ERROR_CODE_RATE_LIMITED; pub use self::error::ERROR_CODE_RPC_APPLICATION_ERROR; /// Error code used when a transport deadline expires. pub use self::error::ERROR_CODE_TIMEOUT; +/// Error code used when bounded WebSocket runtime capacity is exhausted. +pub use self::error::ERROR_CODE_WS_BACKPRESSURE_OVERFLOW; +/// Error code used when a physical WebSocket connection or handshake fails. +pub use self::error::ERROR_CODE_WS_CONNECTION_FAILED; +/// Error code used when WebSocket wire data violates protocol invariants. +pub use self::error::ERROR_CODE_WS_PROTOCOL_ERROR; +/// Error code used when a WebSocket session is no longer available. +pub use self::error::ERROR_CODE_WS_SESSION_CLOSED; /// JSON-RPC 2.0 error payload returned by a remote Solana endpoint. pub use self::json_rpc::JsonRpcErrorObject; /// Validated JSON-RPC 2.0 error response. @@ -309,6 +319,8 @@ pub use self::ws_lifecycle::WsSubscriptionKind; pub use self::ws_lifecycle::WsSubscriptionSnapshot; /// Observable lifecycle state of one logical WebSocket subscription. pub use self::ws_lifecycle::WsSubscriptionState; +/// Shareable handle for one explicitly created physical WebSocket session. +pub use self::ws_session::WsSession; /// Open cluster or network descriptor used by WebSocket endpoint settings. pub use self::ws_settings::WsClusterName; /// Runtime settings for one named WebSocket endpoint. diff --git a/crates/ksp-onchain-transport-lib/src/ws_lifecycle.rs b/crates/ksp-onchain-transport-lib/src/ws_lifecycle.rs index 0c3440f..274c0f4 100644 --- a/crates/ksp-onchain-transport-lib/src/ws_lifecycle.rs +++ b/crates/ksp-onchain-transport-lib/src/ws_lifecycle.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/src/ws_lifecycle.rs -// version: 2 +// version: 3 /// Stable local identity assigned to one physical WebSocket session. #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] @@ -183,7 +183,6 @@ pub struct WsSessionSnapshot { impl WsSessionSnapshot { /// Creates one safe session projection for Transport runtime internals. #[must_use] - #[cfg(test)] #[allow(clippy::too_many_arguments)] pub(crate) fn new( id: crate::WsSessionId, diff --git a/crates/ksp-onchain-transport-lib/src/ws_session.rs b/crates/ksp-onchain-transport-lib/src/ws_session.rs new file mode 100644 index 0000000..2f82d0c --- /dev/null +++ b/crates/ksp-onchain-transport-lib/src/ws_session.rs @@ -0,0 +1,574 @@ +// file: crates/ksp-onchain-transport-lib/src/ws_session.rs +// version: 1 + +use futures_util::SinkExt; // rust-rules: trait-import +use futures_util::StreamExt; // rust-rules: trait-import + +static NEXT_WS_SESSION_ID: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1); + +/// 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. +#[derive(Clone)] +pub struct WsSession { + id: crate::WsSessionId, + command_tx: tokio::sync::mpsc::Sender, + snapshot_rx: tokio::sync::watch::Receiver, + command_timeout: std::time::Duration, +} + +impl WsSession { + /// Opens one physical WebSocket connection for the supplied endpoint settings. + /// + /// Calling this function twice with the same endpoint creates two independent physical sessions. The function returns only after the WebSocket + /// handshake succeeds or the configured command timeout expires. + pub async fn connect(endpoint: crate::WsEndpointSettings) -> ksp_core_lib::Result { + let validation = crate::WsTransportSettings::new(std::vec![endpoint.clone()]).validate(); + if let std::result::Result::Err(error) = validation { + return std::result::Result::Err(error); + } + let id_result = next_session_id(); + let id = match id_result { + std::result::Result::Ok(id) => id, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let initial_snapshot = crate::WsSessionSnapshot::new( + id, + endpoint.name(), + endpoint.provider().clone(), + endpoint.cluster().clone(), + endpoint.protocol(), + crate::WsSessionState::Connecting, + 0, + 0, + 0, + std::vec::Vec::new(), + ); + let (command_tx, command_rx) = tokio::sync::mpsc::channel(endpoint.session().command_queue_capacity()); + 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(); + ksp_logging_lib::debug!( + target: crate::TRACING_TARGET, + session_id = id.get(), + endpoint_name = endpoint.name(), + provider = endpoint.provider().as_str(), + cluster = endpoint.cluster().as_str(), + 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 startup_wait = tokio::time::timeout(command_timeout, startup_rx).await; + 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(std::result::Result::Ok(std::result::Result::Err(error))) => { + join_handle.abort(); + return 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(_) => { + join_handle.abort(); + return std::result::Result::Err(ws_timeout_error(id, "WebSocket handshake exceeded the configured command timeout")); + }, + } + } + + /// Returns the stable local session identity. + #[must_use] + pub const fn id(&self) -> crate::WsSessionId { + return self.id; + } + + /// Returns the latest safe runtime snapshot published by the actor. + #[must_use] + pub fn snapshot(&self) -> crate::WsSessionSnapshot { + return self.snapshot_rx.borrow().clone(); + } + + /// Returns the latest observable physical-session state. + #[must_use] + pub fn state(&self) -> crate::WsSessionState { + return self.snapshot_rx.borrow().state(); + } + + /// 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 { + 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; + match send_wait { + std::result::Result::Ok(std::result::Result::Ok(())) => {}, + std::result::Result::Ok(std::result::Result::Err(_)) => { + return std::result::Result::Err(ws_session_closed_error(self.id, "WebSocket session command channel is closed")); + }, + std::result::Result::Err(_) => { + 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")) + }, + std::result::Result::Err(_) => { + std::result::Result::Err(ws_timeout_error(self.id, "WebSocket JSON-RPC request exceeded the configured command timeout")) + }, + }; + } +} + +impl std::fmt::Debug for WsSession { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + return formatter.debug_struct("WsSession").field("id", &self.id).field("snapshot", &self.snapshot()).finish(); + } +} + +enum WsSessionCommand { + ExecuteJsonRpc { + method: &'static str, + params: std::vec::Vec, + response_tx: tokio::sync::oneshot::Sender>, + }, +} + +struct PendingWsRequest { + method: &'static str, + deadline: tokio::time::Instant, + response_tx: tokio::sync::oneshot::Sender>, +} + +async fn run_ws_session_actor( + id: crate::WsSessionId, + endpoint: crate::WsEndpointSettings, + mut command_rx: tokio::sync::mpsc::Receiver, + snapshot_tx: tokio::sync::watch::Sender, + startup_tx: tokio::sync::oneshot::Sender>, +) { + let websocket_config = tokio_tungstenite::tungstenite::protocol::WebSocketConfig::default() + .write_buffer_size(0) + .max_write_buffer_size(endpoint.session().max_write_buffer_size_bytes()) + .max_message_size(std::option::Option::Some(endpoint.session().max_message_size_bytes())) + .max_frame_size(std::option::Option::Some(endpoint.session().max_frame_size_bytes())); + ksp_logging_lib::trace!( + target: crate::TRACING_TARGET, + session_id = id.get(), + endpoint_name = endpoint.name(), + provider = endpoint.provider().as_str(), + cluster = endpoint.cluster().as_str(), + "opening physical WebSocket connection" + ); + let connect_wait = tokio::time::timeout( + endpoint.session().command_timeout(), + tokio_tungstenite::connect_async_with_config(endpoint.url().as_str(), std::option::Option::Some(websocket_config), false), + ) + .await; + let (mut websocket, handshake_status) = match connect_wait { + std::result::Result::Ok(std::result::Result::Ok((websocket, response))) => (websocket, response.status().as_u16()), + std::result::Result::Ok(std::result::Result::Err(_)) => { + publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Failed, 0); + let error = ws_connection_error(id, &endpoint, "WebSocket connection or handshake failed"); + let _ = startup_tx.send(std::result::Result::Err(error)); + ksp_logging_lib::warn!( + target: crate::TRACING_TARGET, + session_id = id.get(), + endpoint_name = endpoint.name(), + provider = endpoint.provider().as_str(), + cluster = endpoint.cluster().as_str(), + "physical WebSocket connection failed" + ); + return; + }, + std::result::Result::Err(_) => { + publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Failed, 0); + let error = ws_timeout_error(id, "WebSocket connection handshake timed out"); + let _ = startup_tx.send(std::result::Result::Err(error)); + ksp_logging_lib::warn!( + target: crate::TRACING_TARGET, + session_id = id.get(), + endpoint_name = endpoint.name(), + "physical WebSocket connection handshake timed out" + ); + return; + }, + }; + publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Active, 0); + let _ = startup_tx.send(std::result::Result::Ok(())); + ksp_logging_lib::debug!( + target: crate::TRACING_TARGET, + session_id = id.get(), + endpoint_name = endpoint.name(), + handshake_status, + "physical WebSocket session is active" + ); + let mut next_request_id = 1_u64; + let mut pending = std::collections::BTreeMap::::new(); + loop { + let timeout_deadline = next_pending_deadline(&pending); + tokio::select! { + maybe_command = command_rx.recv() => { + let command = match maybe_command { + std::option::Option::Some(command) => command, + std::option::Option::None => { + 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; + 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; + } + 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; + } + publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Active, pending.len()); + }, + () = tokio::time::sleep_until(timeout_deadline) => { + expire_pending_requests(id, &mut pending); + publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Active, pending.len()); + }, + } + } +} + +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, + command: WsSessionCommand, +) -> bool +where + S: tokio::io::AsyncRead + tokio::io::AsyncWrite + std::marker::Unpin, +{ + 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") + .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(), + pending_request_count = pending.len(), + max_pending_requests = endpoint.session().max_pending_requests(), + "rejected WebSocket JSON-RPC request because pending capacity is exhausted" + ); + return true; + } + let request_id = *next_request_id; + let incremented = request_id.checked_add(1); + *next_request_id = match incremented { + std::option::Option::Some(value) => value, + std::option::Option::None => { + 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; + }, + }; + let request_result = crate::JsonRpcRequest::new(request_id, method, params); + let request = match request_result { + std::result::Result::Ok(request) => request, + std::result::Result::Err(error) => { + let _ = response_tx.send(std::result::Result::Err(error)); + return true; + }, + }; + let payload_result = request.to_json_string(); + let payload = match payload_result { + std::result::Result::Ok(payload) => payload, + std::result::Result::Err(error) => { + let _ = response_tx.send(std::result::Result::Err(error)); + return true; + }, + }; + let deadline = tokio::time::Instant::now() + endpoint.session().command_timeout(); + pending.insert(request_id, PendingWsRequest { method, deadline, response_tx }); + ksp_logging_lib::trace!( + target: crate::TRACING_TARGET, + session_id = id.get(), + request_id, + method, + 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; + } + return true; + }, + }; +} + +async fn handle_socket_message( + id: crate::WsSessionId, + endpoint: &crate::WsEndpointSettings, + maybe_message: std::option::Option>, + websocket: &mut tokio_tungstenite::WebSocketStream, + pending: &mut std::collections::BTreeMap, +) -> bool +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(_)) => { + ksp_logging_lib::warn!( + target: crate::TRACING_TARGET, + session_id = id.get(), + endpoint_name = endpoint.name(), + "physical WebSocket read failed" + ); + return false; + }, + std::option::Option::None => { + ksp_logging_lib::warn!( + target: crate::TRACING_TARGET, + session_id = id.get(), + endpoint_name = endpoint.name(), + "physical WebSocket stream ended" + ); + return false; + }, + }; + 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 + }, + 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::Pong(_) => { + ksp_logging_lib::trace!(target: crate::TRACING_TARGET, session_id = id.get(), "received WebSocket pong control frame"); + true + }, + tokio_tungstenite::tungstenite::Message::Close(_) => { + ksp_logging_lib::debug!(target: crate::TRACING_TARGET, session_id = id.get(), "remote peer closed physical WebSocket session"); + false + }, + tokio_tungstenite::tungstenite::Message::Frame(_) => { + ksp_logging_lib::trace!(target: crate::TRACING_TARGET, session_id = id.get(), "ignored internal WebSocket frame event"); + true + }, + }; +} + +fn handle_text_message(id: crate::WsSessionId, text: &str, pending: &mut std::collections::BTreeMap) -> bool { + 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; + }, + }; + 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; + }, + }; + if !object.contains_key("id") { + if object.contains_key("method") && object.contains_key("params") { + ksp_logging_lib::trace!( + target: crate::TRACING_TARGET, + session_id = id.get(), + "received WebSocket notification before subscription registry activation; safely ignored" + ); + return true; + } + ksp_logging_lib::warn!(target: crate::TRACING_TARGET, session_id = id.get(), "received structurally invalid WebSocket JSON-RPC message"); + return false; + } + 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; + }, + }; + let pending_request = match pending.remove(&response_id) { + std::option::Option::Some(pending_request) => pending_request, + std::option::Option::None => { + ksp_logging_lib::debug!( + target: crate::TRACING_TARGET, + session_id = id.get(), + response_id, + "ignored unknown or stale WebSocket JSON-RPC response id" + ); + return true; + }, + }; + let parsed = crate::parse_json_rpc_response_value(value, response_id); + let (result, keep_running) = 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) + }, + 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) + }, + }; + let _ = pending_request.response_tx.send(result); + ksp_logging_lib::trace!( + target: crate::TRACING_TARGET, + session_id = id.get(), + response_id, + method = pending_request.method, + pending_request_count = pending.len(), + "dispatched WebSocket JSON-RPC response to pending request" + ); + return keep_running; +} + +fn next_pending_deadline(pending: &std::collections::BTreeMap) -> tokio::time::Instant { + let mut earliest = std::option::Option::None::; + for request in pending.values() { + earliest = match earliest { + std::option::Option::Some(current) if current <= request.deadline => std::option::Option::Some(current), + _ => std::option::Option::Some(request.deadline), + }; + } + return match earliest { + std::option::Option::Some(deadline) => deadline, + std::option::Option::None => tokio::time::Instant::now() + std::time::Duration::from_secs(86_400), + }; +} + +fn expire_pending_requests(id: crate::WsSessionId, pending: &mut std::collections::BTreeMap) { + let now = tokio::time::Instant::now(); + let expired_ids = pending + .iter() + .filter_map(|(request_id, request)| if request.deadline <= now { std::option::Option::Some(*request_id) } else { std::option::Option::None }) + .collect::>(); + for request_id in expired_ids { + if let std::option::Option::Some(request) = pending.remove(&request_id) { + let error = ws_timeout_error(id, "WebSocket JSON-RPC request timed out while awaiting the remote response").with_context("method", request.method); + let _ = request.response_tx.send(std::result::Result::Err(error)); + ksp_logging_lib::debug!( + target: crate::TRACING_TARGET, + session_id = id.get(), + request_id, + method = request.method, + "expired pending WebSocket JSON-RPC request" + ); + } + } +} + +fn fail_all_pending( + pending: &mut std::collections::BTreeMap, + id: crate::WsSessionId, + code: ksp_core_lib::ErrorCode, + message: &'static str, +) { + let requests = std::mem::take(pending); + for (request_id, request) in requests { + let error = ksp_core_lib::Error::new(code, message) + .with_context("session_id", id.get().to_string()) + .with_context("request_id", request_id.to_string()) + .with_context("method", request.method); + let _ = request.response_tx.send(std::result::Result::Err(error)); + } +} + +fn publish_snapshot( + snapshot_tx: &tokio::sync::watch::Sender, + id: crate::WsSessionId, + endpoint: &crate::WsEndpointSettings, + state: crate::WsSessionState, + pending_request_count: usize, +) { + let snapshot = crate::WsSessionSnapshot::new( + id, + endpoint.name(), + endpoint.provider().clone(), + endpoint.cluster().clone(), + endpoint.protocol(), + state, + pending_request_count, + 0, + 0, + std::vec::Vec::new(), + ); + snapshot_tx.send_replace(snapshot); +} + +fn next_session_id() -> ksp_core_lib::Result { + let update_result = + NEXT_WS_SESSION_ID.fetch_update(std::sync::atomic::Ordering::Relaxed, std::sync::atomic::Ordering::Relaxed, |current| current.checked_add(1)); + let value = match update_result { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => { + return std::result::Result::Err(ksp_core_lib::Error::new(crate::ERROR_CODE_WS_PROTOCOL_ERROR, "WebSocket session identity space is exhausted")); + }, + }; + return match std::num::NonZeroU64::new(value) { + std::option::Option::Some(value) => std::result::Result::Ok(crate::WsSessionId::new(value)), + std::option::Option::None => { + std::result::Result::Err(ksp_core_lib::Error::new(crate::ERROR_CODE_WS_PROTOCOL_ERROR, "WebSocket session identity generator produced zero")) + }, + }; +} + +fn ws_connection_error(id: crate::WsSessionId, endpoint: &crate::WsEndpointSettings, message: &'static str) -> ksp_core_lib::Error { + return ksp_core_lib::Error::new(crate::ERROR_CODE_WS_CONNECTION_FAILED, message) + .with_context("session_id", id.get().to_string()) + .with_context("endpoint_name", endpoint.name()) + .with_context("provider", endpoint.provider().as_str()) + .with_context("cluster", endpoint.cluster().as_str()); +} + +fn ws_session_closed_error(id: crate::WsSessionId, message: &'static str) -> ksp_core_lib::Error { + return ksp_core_lib::Error::new(crate::ERROR_CODE_WS_SESSION_CLOSED, message).with_context("session_id", id.get().to_string()); +} + +fn ws_timeout_error(id: crate::WsSessionId, message: &'static str) -> ksp_core_lib::Error { + return ksp_core_lib::Error::new(crate::ERROR_CODE_TIMEOUT, message).with_context("session_id", id.get().to_string()); +} + +#[cfg(test)] +#[path = "../unit_tests/ws_session.rs"] +mod tests; diff --git a/crates/ksp-onchain-transport-lib/tests/public_api.rs b/crates/ksp-onchain-transport-lib/tests/public_api.rs index fe42a59..8749282 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: 25 +// version: 26 //! Integration tests for the public `ksp-onchain-transport-lib` consumer contract. @@ -549,3 +549,13 @@ fn public_v0_2_7_pre_002_websocket_settings_and_lifecycle_contracts_are_availabl assert_eq!(ksp_onchain_transport_lib::WsSubscriptionKind::Slot.as_str(), "slot"); assert_eq!(ksp_onchain_transport_lib::WsSubscriptionState::Requested, ksp_onchain_transport_lib::WsSubscriptionState::Requested); } + +#[test] +fn public_v0_2_7_pre_004_physical_websocket_session_contract_is_available_from_crate_root() { + let type_name = std::any::type_name::(); + assert!(type_name.ends_with("WsSession")); + assert_eq!(ksp_onchain_transport_lib::ERROR_CODE_WS_BACKPRESSURE_OVERFLOW.code(), "ws_backpressure_overflow"); + assert_eq!(ksp_onchain_transport_lib::ERROR_CODE_WS_CONNECTION_FAILED.code(), "ws_connection_failed"); + 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"); +} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs b/crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs new file mode 100644 index 0000000..8bfa2ac --- /dev/null +++ b/crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs @@ -0,0 +1,151 @@ +// file: crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs +// version: 1 + +use futures_util::SinkExt; // rust-rules: trait-import +use futures_util::StreamExt; // rust-rules: trait-import + +fn local_endpoint(url: &str) -> crate::WsEndpointSettings { + return crate::WsEndpointSettings::new( + "local_ws", + true, + crate::WsProviderName::new("local-fixture"), + crate::WsClusterName::new("local"), + crate::WsProtocolKind::SolanaStandard, + crate::WsEndpointUrl::parse(url).expect("local test WebSocket URL must parse"), + crate::WsSessionSettings::default(), + ); +} + +async fn bind_local_listener() -> (tokio::net::TcpListener, std::string::String) { + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.expect("local listener must bind"); + let address = listener.local_addr().expect("local listener must expose address"); + return (listener, format!("ws://{address}")); +} + +async fn read_request(websocket: &mut tokio_tungstenite::WebSocketStream) -> serde_json::Value { + let message = websocket.next().await.expect("request message must exist").expect("request message must decode"); + let text = message.to_text().expect("request must be text"); + return serde_json::from_str(text).expect("request must contain JSON"); +} + +async fn send_result(websocket: &mut tokio_tungstenite::WebSocketStream, request: &serde_json::Value, result: serde_json::Value) { + let id = request.get("id").and_then(serde_json::Value::as_u64).expect("request id must be numeric"); + let response = serde_json::json!({"jsonrpc":"2.0","id":id,"result":result}); + websocket.send(tokio_tungstenite::tungstenite::Message::Text(response.to_string().into())).await.expect("local response must send"); +} + +#[tokio::test(flavor = "current_thread")] +async fn websocket_session_connects_and_round_trips_internal_json_rpc() { + 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"); + let request = read_request(&mut websocket).await; + assert_eq!(request.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("getVersion")); + send_result(&mut websocket, &request, serde_json::json!({"solana-core":"fixture"})).await; + }); + let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed"); + assert_eq!(session.state(), crate::WsSessionState::Active); + assert_eq!(session.snapshot().pending_request_count(), 0); + let result = session.execute_json_rpc("getVersion", std::vec::Vec::new()).await.expect("fixture JSON-RPC call must succeed"); + assert_eq!(result.get("solana-core").and_then(serde_json::Value::as_str), std::option::Option::Some("fixture")); + tokio::task::yield_now().await; + assert_eq!(session.snapshot().pending_request_count(), 0); + server.await.expect("local server task must complete"); +} + +#[tokio::test(flavor = "current_thread")] +async fn websocket_sessions_are_physical_and_distinct_even_for_the_same_url() { + let (listener, url) = bind_local_listener().await; + let server = tokio::spawn(async move { + let first = listener.accept().await.expect("first local client must connect"); + let first_ws = tokio_tungstenite::accept_async(first.0).await.expect("first handshake must succeed"); + let second = listener.accept().await.expect("second local client must connect"); + let second_ws = tokio_tungstenite::accept_async(second.0).await.expect("second handshake must succeed"); + return (first_ws, second_ws); + }); + let first = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("first client must connect"); + let second = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("second client must connect"); + assert_ne!(first.id(), second.id()); + assert_eq!(first.snapshot().endpoint_name(), second.snapshot().endpoint_name()); + let _ = server.await.expect("local server task must complete"); +} + +#[tokio::test(flavor = "current_thread")] +async fn websocket_pending_map_dispatches_out_of_order_responses_by_request_id() { + 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"); + let first = read_request(&mut websocket).await; + let second = read_request(&mut websocket).await; + send_result(&mut websocket, &second, serde_json::json!("second")).await; + send_result(&mut websocket, &first, serde_json::json!("first")).await; + }); + let session = crate::WsSession::connect(local_endpoint(url.as_str())).await.expect("client handshake must succeed"); + let first_session = session.clone(); + let second_session = session.clone(); + let (first, second) = tokio::join!( + first_session.execute_json_rpc("firstMethod", std::vec::Vec::new()), + second_session.execute_json_rpc("secondMethod", std::vec::Vec::new()) + ); + assert_eq!(first.expect("first response must dispatch"), serde_json::json!("first")); + assert_eq!(second.expect("second response must dispatch"), serde_json::json!("second")); + tokio::task::yield_now().await; + assert_eq!(session.snapshot().pending_request_count(), 0); + server.await.expect("local server task must complete"); +} + +#[tokio::test(flavor = "current_thread")] +async fn websocket_rpc_application_error_does_not_fail_the_physical_session() { + 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"); + let request = read_request(&mut websocket).await; + let id = request.get("id").and_then(serde_json::Value::as_u64).expect("request id must be numeric"); + let response = serde_json::json!({"jsonrpc":"2.0","id":id,"error":{"code":-32602,"message":"fixture invalid params"}}); + websocket.send(tokio_tungstenite::tungstenite::Message::Text(response.to_string().into())).await.expect("local error response must 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"); + let error = session.execute_json_rpc("badMethod", std::vec::Vec::new()).await.expect_err("RPC application error must surface"); + assert_eq!(error.code(), crate::ERROR_CODE_RPC_APPLICATION_ERROR); + assert_eq!(session.state(), crate::WsSessionState::Active); + let result = session.execute_json_rpc("goodMethod", std::vec::Vec::new()).await.expect("session must remain usable"); + assert_eq!(result, serde_json::json!(true)); + server.await.expect("local server task must complete"); +} + +#[tokio::test(flavor = "current_thread")] +async fn websocket_connection_errors_do_not_echo_url_credentials() { + let defaults = crate::WsSessionSettings::default(); + let short_session = crate::WsSessionSettings::new( + std::time::Duration::from_millis(150), + defaults.close_timeout(), + defaults.reconnect().clone(), + defaults.resubscribe(), + defaults.command_queue_capacity(), + defaults.notification_queue_capacity(), + defaults.max_active_subscriptions(), + defaults.max_pending_requests(), + defaults.max_message_size_bytes(), + defaults.max_frame_size_bytes(), + defaults.max_write_buffer_size_bytes(), + ); + let endpoint = crate::WsEndpointSettings::new( + "redaction_fixture", + true, + crate::WsProviderName::new("fixture-provider"), + crate::WsClusterName::new("local"), + crate::WsProtocolKind::SolanaStandard, + crate::WsEndpointUrl::parse("ws://user:password@127.0.0.1:1/private?api-key=SECRET-CANARY").expect("test URL must parse"), + short_session, + ); + let error = crate::WsSession::connect(endpoint).await.expect_err("unreachable local endpoint must fail"); + let rendered = format!("{error:?}"); + assert!(!rendered.contains("SECRET-CANARY")); + assert!(!rendered.contains("password")); + assert!(!rendered.contains("private?")); +} diff --git a/deltas/0.2.7/pre.001-fix.001.md b/deltas/0.2.7/pre.001-fix.001.md index 53ba3b6..7c18336 100644 --- a/deltas/0.2.7/pre.001-fix.001.md +++ b/deltas/0.2.7/pre.001-fix.001.md @@ -56,18 +56,18 @@ La documentation Solana reste l'autorité pour la **surface publique annoncée** Le plan et la compliance suivent désormais explicitement les SIMDs pertinents : -| SIMD | Statut | Décision KSP principale | -|------------------------------------------------|-----------|----------------------------------------------------------------------------------------------------------------------| -| `0118` Partitioned Epoch Rewards Distribution | Activated | réutiliser les DTOs bloc/rewards lossless déjà acquis par HTTP | -| `0291` Commission Rate in Basis Points | Review | ne pas dériver localement une représentation de commission depuis une autre | -| `0296` Larger Transaction Size | Review | proposition jusqu'à 4096 octets ; aucune limite WS dérivée en dur de l'ancienne taille transaction de 1232 octets | -| `0298` Bank Hash in Block Footer | Idea | aucun `bankHash` spéculatif | -| `0301` parent bank hash | PR fermé | aucun `parentBankHash` spéculatif ; PR non mergée | -| `0307` Add Block Footer | Review | aucun `footer` spéculatif ; réaudit lorsque l'upstream l'expose | -| `0326` Alpenglow | Review | ne pas figer les sémantiques TowerBFT des flux unstable | -| `0337` Alpenglow Fast Leader Handover Markers | Review | surveiller l'impact futur sur shape/taille des blocs | -| `0384` Alpenglow migration | Review | ne pas supposer une séquence exhaustive de notifications commitment/optimistic confirmation | -| `0385` Transaction V1 | Review | conserver versions transaction et `maxSupportedTransactionVersion` génériques | +| SIMD | Statut | Décision KSP principale | +|-----------------------------------------------|-----------|-------------------------------------------------------------------------------------------------------------------| +| `0118` Partitioned Epoch Rewards Distribution | Activated | réutiliser les DTOs bloc/rewards lossless déjà acquis par HTTP | +| `0291` Commission Rate in Basis Points | Review | ne pas dériver localement une représentation de commission depuis une autre | +| `0296` Larger Transaction Size | Review | proposition jusqu'à 4096 octets ; aucune limite WS dérivée en dur de l'ancienne taille transaction de 1232 octets | +| `0298` Bank Hash in Block Footer | Idea | aucun `bankHash` spéculatif | +| `0301` parent bank hash | PR fermé | aucun `parentBankHash` spéculatif ; PR non mergée | +| `0307` Add Block Footer | Review | aucun `footer` spéculatif ; réaudit lorsque l'upstream l'expose | +| `0326` Alpenglow | Review | ne pas figer les sémantiques TowerBFT des flux unstable | +| `0337` Alpenglow Fast Leader Handover Markers | Review | surveiller l'impact futur sur shape/taille des blocs | +| `0384` Alpenglow migration | Review | ne pas supposer une séquence exhaustive de notifications commitment/optimistic confirmation | +| `0385` Transaction V1 | Review | conserver versions transaction et `maxSupportedTransactionVersion` génériques | Aucun de ces SIMDs n'ajoute, dans la baseline Agave `v4.2.1`, une dixième famille WebSocket standard. diff --git a/deltas/0.2.7/pre.004.md b/deltas/0.2.7/pre.004.md new file mode 100644 index 0000000..cdc2b31 --- /dev/null +++ b/deltas/0.2.7/pre.004.md @@ -0,0 +1,240 @@ + + + +# Delta `0.2.7-pre.004` — runtime WebSocket physique + actor JSON-RPC + +## 1. Base requise + +```text +0.2.7-pre.003 appliquée +workspace.package.version = 0.2.7-pre.3 +``` + +Le checkpoint opérateur reçu avant cette tranche est vert : `cargo fmt --all`, audit Python, `cargo check --workspace`, `cargo clippy --workspace --all-targets`, tests Transport, tests Config et `cargo test --workspace`. + +## 2. Signal technique + +Cette prerelease non-fix modifie dépendances, runtime Rust et tests. Conformément au workflow KSP : + +```text +livraison = 0.2.7-pre.004 +workspace.package.version = 0.2.7-pre.4 +commit = v0.2.7-pre.004 +``` + +Aucun tag prerelease. + +## 3. Dépendances WebSocket matérialisées + +Le réaudit du 22 août 2026 confirme les versions retenues depuis `pre.001` : + +```text +tokio-tungstenite 0.30.0 +futures-util 0.3.34 +``` + +Le root déclare sans features consumer : + +```toml +tokio-tungstenite = { version = "^0.30", default-features = false } +futures-util = { version = "^0.3", default-features = false } +``` + +Transport active seulement : + +```text +tokio-tungstenite : connect + rustls-tls-webpki-roots +futures-util : sink + std +tokio : macros + rt + sync + time +``` + +Le fixture serveur local ajoute `tokio/net` côté dev. + +Aucune dépendance Config, Store, Program, Wallet ou `tracing` direct n'est introduite. + +## 4. `WsSession` physique + +Nouvelle surface publique : + +```text +WsSession::connect(WsEndpointSettings) +WsSession::id() +WsSession::state() +WsSession::snapshot() +``` + +Un appel de `connect` crée exactement une connexion physique. Deux appels avec le même endpoint créent deux sockets indépendants ; aucun singleton, pool ou scheduler automatique n'est ajouté. + +Le caller ne reçoit jamais le socket brut. + +## 5. Actor propriétaire du socket + +Une tâche actor unique possède : + +```text +WebSocketStream +compteur JSON-RPC request id +map pending requests +bounded command receiver +publication WsSessionSnapshot +``` + +Le handle communique avec l'actor par `tokio::sync::mpsc` borné selon `command_queue_capacity`. + +Les snapshots sont publiés via `tokio::sync::watch` et conservent seulement les metadata sûres prévues en `pre.002`. + +## 6. Handshake et `WebSocketConfig` + +`WsSession::connect` attend le handshake sous `command_timeout` et configure explicitement : + +```text +write_buffer_size = 0 +max_write_buffer_size = WsSessionSettings.max_write_buffer_size_bytes +max_message_size = WsSessionSettings.max_message_size_bytes +max_frame_size = WsSessionSettings.max_frame_size_bytes +``` + +Le `write_buffer_size = 0` évite de rendre la validité de la configuration KSP dépendante du buffer par défaut interne de Tungstenite et garantit que le plafond configuré reste strictement supérieur au target buffer. + +Les tests oversized et les recalibrages éventuels restent le gate `pre.005`. + +## 7. Pending JSON-RPC + +La primitive interne actor : + +```text +execute_json_rpc(method, params) +``` + +reste **`pub(crate)`**. Elle n'est volontairement pas exposée comme API raw provider-extension publique. + +Comportement : + +- ID numérique KSP monotone par session ; +- sérialisation via `JsonRpcRequest` existant ; +- map `BTreeMap` bornée par `max_pending_requests` ; +- deadline par request issue de `command_timeout` ; +- dispatch des réponses par `id`, y compris si elles arrivent hors ordre ; +- erreurs JSON-RPC applicatives renvoyées au caller concerné sans teardown de la connexion ; +- ID réponse inconnu/stale ignoré avec diagnostic sûr ; +- JSON structurellement invalide classé erreur protocole session. + +Cette primitive sera consommée par le moteur de subscriptions à partir de `pre.006`. + +## 8. Lifecycle limité à la tranche + +`pre.004` matérialise : + +```text +Connecting -> Active +connection/read/write failure -> Failed +last handle dropped -> cleanup best-effort -> Closed +``` + +Le reconnect/resubscribe reste `pre.007`. + +Le shutdown async public, les budgets de Close et les fixtures peer hostile restent `pre.005`. + +Ping reçu est répondu par Pong afin de conserver l'interopérabilité du socket. Aucun heartbeat applicatif périodique n'est ajouté. + +## 9. Erreurs + +Nouveaux codes publics : + +```text +ws_backpressure_overflow +ws_connection_failed +ws_protocol_error +ws_session_closed +``` + +Les erreurs de connexion WebSocket ne conservent volontairement pas la source Tungstenite brute : celle-ci pourrait contenir une request/URI ou d'autres détails provider. Les erreurs KSP exposent uniquement `session_id`, endpoint logique, provider et cluster. + +## 10. Logging / tracing + +Toutes les émissions passent exclusivement par `ksp-logging-lib` et réutilisent : + +```text +crates/ksp-onchain-transport-lib/src/constants.rs +TRACING_TARGET = "ksp-onchain-transport-lib" +``` + +Répartition principale : + +```text +trace -> ouverture socket, send/dispatch JSON-RPC, Ping/Pong, notification prématurée ignorée +debug -> actor start, handshake actif, réponse stale, timeout pending, remote close +warn -> handshake/read/write failure, malformed wire, capacité pending épuisée +``` + +Ne sont jamais loggés : URL complète, credentials, query token, payload JSON-RPC complet ou notification brute. + +## 11. Serveur local déterministe + +Nouveaux tests runtime sans Internet : + +```text +handshake local + round-trip JSON-RPC +deux sessions physiques distinctes sur la même URL +deux requests concurrentes + réponses inversées +application error sans teardown de session +connection error sans fuite URL/credential +``` + +Le serveur utilise `tokio::net::TcpListener` + `tokio_tungstenite::accept_async`. + +## 12. Canaris workspace/public API + +Le canari workspace dependencies est synchronisé avec les nouvelles dépendances/features et continue de vérifier le firewall Transport. + +Le canari public API vérifie la disponibilité de `WsSession` et les quatre nouveaux codes d'erreur. + +## 13. Documentation synchronisée + +Mis à jour : + +```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 conformément à la politique de série prerelease. + +## 14. Validation de préparation + +Exécuté dans le sandbox : + +```text +python3 scripts/audit_rust_workspace_rules.py +General Rust rule audit: clean +Rust export completeness audit: 0 candidate(s) +KSP workspace Rust rule audit: clean + +inspection absence tracing direct OK +inspection TRACING_TARGET OK +inspection dépendances workspace/member OK +inspection Markdown tables OK +``` + +Cargo n'est pas disponible dans le sandbox de génération ; aucun résultat Cargo local n'est revendiqué. + +## 15. 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.004 +``` + +La tranche suivante est `0.2.7-pre.005` : adversarial limits frame/message/request, control frames, cancellation, close/shutdown explicite et peer hostile. 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 ea56292..75efee2 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.003`.** Les gates `pre.001` et `pre.002` sont clos. `pre.003` matérialise `std.transport` V2 HTTP + WebSocket, la lecture backward V1 HTTP-only et l'adapter `Config -> WsTransportSettings`. Le socket physique reste différé à `pre.004`. +> **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. ## 1. Objet et base vérifiée @@ -366,7 +366,7 @@ Audit au 2026-08-22 : | `tokio-websockets` | `0.13.3` | alternative viable, non retenue | strict/minimal et performant, mais exige davantage d'assemblage/features et n'apporte pas de besoin fonctionnel supérieur démontré pour cette foundation | | `fastwebsockets` | `0.10.0` | non retenue | plus bas niveau ; peut déléguer davantage de compliance au caller, inutile pour la première foundation KSP | -Landing prévu, **pas dans `pre.001`** : +Landing matérialisé par **`pre.004`** : ```toml # root [workspace.dependencies] @@ -379,7 +379,7 @@ futures-util = { workspace = true, features = ["std", "sink"] } tokio = { workspace = true, features = ["macros", "rt", "sync", "time"] } ``` -Le serveur de test local pourra activer `tokio/net` en dev si KSP utilise directement `TcpListener`. +Le serveur de test local active `tokio/net` en dev et utilise directement `TcpListener`. `url` n'est pas retenu a priori : l'endpoint KSP peut être validé et passé comme chaîne/request sans ajouter la feature uniquement par habitude bot3. `handshake`/`stream` sont déjà requis transitivement par `connect`/TLS et ne doivent pas être listés sans nécessité directe. @@ -486,6 +486,8 @@ Une tâche actor possède exclusivement : Le handle public `WsSession` communique avec cet actor par canal bounded. Aucun caller ne split/manipule directement le socket. +`pre.004` matérialise cette foundation : `WsSession::connect` ouvre une connexion physique explicite, le socket reste exclusivement dans l'actor, les commandes passent par `mpsc` borné, les réponses JSON-RPC sont dispatchées par ID KSP dans une map bornée et les snapshots sûrs sont publiés par `watch`. La primitive JSON-RPC reste `pub(crate)` afin de ne pas créer une API publique raw provider-extension avant les wrappers typed. Le registry des subscriptions et le mapping remote/local restent `pre.006`. + ### 9.4 Identités ```text @@ -868,7 +870,7 @@ L'inventaire officiel n'impose que 9 familles de subscriptions, mais le lifecycl pre.001 audit interne/externe + matrice 18 méthodes + bot3 + dependencies + threat model + plan/sizing 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 deps tokio-tungstenite/futures-util + actor physique + handshake/read/write + pending JSON-RPC + serveur local +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.006 registry subscriptions + IDs locaux + generic subscribe/unsubscribe engine + channels typed bounded pre.007 reconnect borné + resubscribe déterministe + continuity gap + races unsubscribe/reconnect diff --git a/docs/validation/010-V0_2_7_ONCHAIN_WEBSOCKET.md b/docs/validation/010-V0_2_7_ONCHAIN_WEBSOCKET.md index 527d204..1aa003b 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.003`.** Les settings/lifecycle `pre.002` et la composition Config V2 `pre.003` sont matérialisés. Les preuves runtime socket/subscription restent ouvertes. +> **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. ## 1. Baseline normative @@ -189,39 +189,39 @@ Les votes observés sont gossip/pre-consensus ; aucune garantie d'entrée dans l ## 6. Lifecycle compliance initiale -| Contrat | Décision `pre.001` | Gate cible | -|-------------------------------|-------------------------------------------------------------------|----------------------| -| plusieurs sessions / même URL | création physique explicite ; aucun singleton/pool automatique | `pre.004`, `pre.012` | -| plusieurs subs / session | registry actor par session | `pre.006` | -| ID public subscription | local KSP stable | `pre.006` | -| ID serveur | éphémère interne et remappé | `pre.006`/`pre.007` | -| session states | Disconnected/Connecting/Active/Reconnecting/Closing/Closed/Failed | **Done `pre.002`** | -| subscription states | Requested/Active/Resubscribing/Cancelling/Closed/Failed | **Done `pre.002`** | -| reconnect | physique uniquement, budget/backoff finis | `pre.007` | -| resubscribe | policy `Never` ou `ActiveSubscriptions`, ordre local déterministe | `pre.007` | -| continuity | gap observable, aucune promesse lossless | `pre.007` | -| backpressure | queue par sub bounded ; overflow => fail local explicite | `pre.008` | -| shutdown | explicite, bounded, annule reconnect et subscriptions | `pre.005` | -| keepalive | pas de ping applicatif périodique sans besoin démontré | `pre.005` | +| Contrat | Décision `pre.001` | Gate cible | +|-------------------------------|-------------------------------------------------------------------|--------------------------------------------| +| plusieurs sessions / même URL | création physique explicite ; aucun singleton/pool automatique | **Done `pre.004`**, canary final `pre.012` | +| plusieurs subs / session | registry actor par session | `pre.006` | +| ID public subscription | local KSP stable | `pre.006` | +| ID serveur | éphémère interne et remappé | `pre.006`/`pre.007` | +| session states | Disconnected/Connecting/Active/Reconnecting/Closing/Closed/Failed | **Done `pre.002`** | +| subscription states | Requested/Active/Resubscribing/Cancelling/Closed/Failed | **Done `pre.002`** | +| reconnect | physique uniquement, budget/backoff finis | `pre.007` | +| resubscribe | policy `Never` ou `ActiveSubscriptions`, ordre local déterministe | `pre.007` | +| continuity | gap observable, aucune promesse lossless | `pre.007` | +| backpressure | queue par sub bounded ; overflow => fail local explicite | `pre.008` | +| shutdown | explicite, bounded, annule reconnect et subscriptions | `pre.005` | +| keepalive | pas de ping applicatif périodique sans besoin démontré | `pre.005` | ## 7. Threat/security compliance initiale -| Invariant | Preuve attendue | Statut | -|---------------------------------------------------|---------------------------------------------|--------------------| -| URL/credentials absents de `Debug` | unit tests URL wrapper | **Done `pre.002`** | -| URL/credentials absents des erreurs | validation URL + future connection errors | Partial `pre.002` | -| URL/credentials absents des logs | safe fields `pre.002`, capture actor future | Partial `pre.002` | -| snapshots sans URL/raw payload | unit shape + public contract | **Done `pre.002`** | -| frame/message finis | settings bornés, enforcement socket futur | Partial `pre.002` | -| JSON borné indirectement par message | oversized + malformed fixture | Planned | -| queues notifications bornées | capacité settings, channel futur | Partial `pre.002` | -| pending RPC borné + timeout | settings bornés, runtime futur | Partial `pre.002` | -| reconnect loop bornée | repeated disconnect fixture | Planned | -| unsubscribe pendant reconnect ne resubscribe pas | race fixture | Planned | -| signature terminale ne resubscribe pas | terminal fixture | Planned | -| shutdown ne bloque pas | peer hostile/no close ack fixture | Planned | -| no Store/Program/Wallet/Config dep dans Transport | cargo tree + source canary | Planned | -| no direct `tracing` dans Transport | workspace audit + source audit | **Done `pre.002`** | +| Invariant | Preuve attendue | Statut | +|---------------------------------------------------|----------------------------------------------|------------------------------------------| +| URL/credentials absents de `Debug` | unit tests URL wrapper | **Done `pre.002`** | +| URL/credentials absents des erreurs | validation URL + connection errors safe | **Done through `pre.004`** | +| URL/credentials absents des logs | actor logs only safe endpoint metadata | Partial `pre.004`, capture finale future | +| snapshots sans URL/raw payload | unit shape + public contract | **Done `pre.002`** | +| frame/message finis | settings bornés + `WebSocketConfig` raccordé | Partial `pre.004`, adversarial `pre.005` | +| JSON borné indirectement par message | oversized + malformed fixture | Planned | +| queues notifications bornées | capacité settings, channel futur | Partial `pre.002` | +| pending RPC borné + timeout | map actor bornée + timeout request | **Done `pre.004`** | +| reconnect loop bornée | repeated disconnect fixture | Planned | +| unsubscribe pendant reconnect ne resubscribe pas | race fixture | Planned | +| signature terminale ne resubscribe pas | terminal fixture | Planned | +| shutdown ne bloque pas | peer hostile/no close ack fixture | Planned | +| no Store/Program/Wallet/Config dep dans Transport | cargo tree + source canary | Planned | +| no direct `tracing` dans Transport | workspace audit + source audit | **Done `pre.002`** | ## 8. Dependency compliance initiale @@ -249,7 +249,7 @@ https://docs.rs/crate/futures-util/0.3.34 Alternatives auditées mais non retenues : `tokio-websockets 0.13.3`, `fastwebsockets 0.10.0`. -Aucune de ces dependencies n'est ajoutée par `pre.001`; le graphe Cargo stable ne change pas dans ce gate hors signal de version workspace. +`pre.004` matérialise `tokio-tungstenite` et `futures-util` dans `[workspace.dependencies]` sans features consumer au root. Transport active seulement `connect`, `rustls-tls-webpki-roots`, `std` et `sink`; `tokio/net` est ajouté côté dev fixture local. ## 9. Config compliance initiale @@ -330,6 +330,34 @@ Logging : `TRACING_TARGET` reste défini dans `constants.rs`; les nouveaux évé Les defaults de taille/queue sont des **policies KSP locales**, pas des limites Solana. Leur enforcement réel et leurs tests oversized/slow-consumer restent attendus dans les tranches socket/backpressure. +## 9.2 Checkpoint runtime physique `pre.004` + +Surface matérialisée : + +```text +WsSession::connect(endpoint) public, une connexion physique par appel +actor socket propriétaire exclusif du WebSocket +command queue tokio mpsc bounded +pending JSON-RPC BTreeMap bounded par max_pending_requests +request timeout command_timeout +response dispatch par id numérique KSP +state snapshot watch + WsSessionSnapshot +raw JSON-RPC public non, primitive pub(crate) +reconnect/resubscribe non, gates futurs +``` + +Fixtures déterministes `pre.004` : + +```text +handshake local + JSON-RPC round-trip +deux sessions physiques distinctes sur la même URL +deux requests concurrentes + réponses inversées +erreur RPC applicative sans teardown de session +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`. + ## 10. Validation du gate `pre.001` Exécuté dans le sandbox :