v0.2.7-pre.007-fix.001

This commit is contained in:
2026-08-22 21:00:29 +02:00
parent b67fa89f44
commit 4c540d67a7
4 changed files with 160 additions and 16 deletions

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-onchain-transport-lib/src/ws_session.rs
// version: 6
// version: 7
use futures_util::SinkExt; // rust-rules: trait-import
use futures_util::StreamExt; // rust-rules: trait-import
@@ -280,7 +280,7 @@ enum WsActorIoOutcome {
}
enum WsReconnectOutcome {
Connected { websocket: WsPhysicalStream },
Connected { websocket: std::boxed::Box<WsPhysicalStream> },
ShutdownRequested { deadline: tokio::time::Instant },
HandlesDropped,
Exhausted,
@@ -293,7 +293,7 @@ enum WsReconnectControlOutcome {
}
enum WsConnectAttemptOutcome {
Connected { websocket: WsPhysicalStream, handshake_status: u16 },
Connected { websocket: std::boxed::Box<WsPhysicalStream>, handshake_status: u16 },
Retry,
ShutdownRequested { deadline: tokio::time::Instant },
HandlesDropped,
@@ -481,7 +481,7 @@ async fn run_ws_session_actor(
)
.await;
match recovery {
WsReconnectOutcome::Connected { websocket: replacement } => websocket = replacement,
WsReconnectOutcome::Connected { websocket: replacement } => websocket = *replacement,
WsReconnectOutcome::ShutdownRequested { deadline } => {
finish_disconnected_shutdown(
id,
@@ -522,7 +522,7 @@ async fn run_ws_session_actor(
)
.await;
match recovery {
WsReconnectOutcome::Connected { websocket: replacement } => websocket = replacement,
WsReconnectOutcome::Connected { websocket: replacement } => websocket = *replacement,
WsReconnectOutcome::ShutdownRequested { deadline } => {
finish_disconnected_shutdown(
id,
@@ -610,7 +610,7 @@ async fn recover_websocket_session(
handshake_status,
"replacement physical WebSocket connection established"
);
websocket
*websocket
},
WsConnectAttemptOutcome::Retry => {
attempt = attempt.saturating_add(1);
@@ -645,7 +645,7 @@ async fn recover_websocket_session(
continuity_gap_count = *continuity_gap_count,
"physical WebSocket reconnect completed and retry budget reset"
);
return WsReconnectOutcome::Connected { websocket };
return WsReconnectOutcome::Connected { websocket: std::boxed::Box::new(websocket) };
},
WsActorIoOutcome::ShutdownRequested { deadline } => return WsReconnectOutcome::ShutdownRequested { deadline },
WsActorIoOutcome::RemoteClosed => {
@@ -872,7 +872,7 @@ async fn connect_replacement_websocket(
result = &mut connect => {
return match result {
std::result::Result::Ok((websocket, response)) => WsConnectAttemptOutcome::Connected {
websocket,
websocket: std::boxed::Box::new(websocket),
handshake_status: response.status().as_u16(),
},
std::result::Result::Err(_) => {

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs
// version: 5
// version: 6
use futures_util::SinkExt; // rust-rules: trait-import
use futures_util::StreamExt; // rust-rules: trait-import
@@ -280,17 +280,22 @@ async fn websocket_oversized_outbound_request_is_rejected_before_socket_write()
}
#[tokio::test(flavor = "current_thread")]
async fn websocket_oversized_inbound_frame_fails_before_json_decode() {
async fn websocket_oversized_inbound_frame_triggers_reconnect_before_json_decode() {
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 (stream, _) = listener.accept().await.expect("local server must accept initial client");
let mut websocket = tokio_tungstenite::accept_async(stream).await.expect("initial local WebSocket handshake must succeed");
websocket.send(tokio_tungstenite::tungstenite::Message::Text("X".repeat(512).into())).await.expect("oversized fixture message must send");
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
let (replacement_stream, _) = listener.accept().await.expect("local server must accept replacement client");
let mut replacement = tokio_tungstenite::accept_async(replacement_stream).await.expect("replacement local WebSocket handshake must succeed");
let _ = replacement.next().await;
let _ = replacement.flush().await;
});
let settings = session_settings(std::time::Duration::from_secs(1), std::time::Duration::from_millis(200), 8, 128, 64, 1024);
let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed");
wait_for_state(&session, crate::WsSessionState::Failed).await;
wait_for_gap_count(&session, 1).await;
assert_eq!(session.state(), crate::WsSessionState::Active);
session.close().await.expect("recovered session must close cleanly");
server.await.expect("local server task must complete");
}