v0.2.8-pre.007
This commit is contained in:
@@ -1,11 +1,13 @@
|
||||
// file: crates/ksp-onchain-transport-lib/src/ws_session.rs
|
||||
// version: 13
|
||||
// version: 14
|
||||
|
||||
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);
|
||||
|
||||
const HELIUS_WS_HEARTBEAT_INTERVAL: std::time::Duration = std::time::Duration::from_secs(60);
|
||||
|
||||
type WsPhysicalStream = tokio_tungstenite::WebSocketStream<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>>;
|
||||
|
||||
/// Shareable handle for one explicitly created physical WebSocket session.
|
||||
@@ -412,6 +414,8 @@ async fn run_ws_session_actor(
|
||||
let mut pending = std::collections::BTreeMap::<u64, PendingWsRequest>::new();
|
||||
let mut subscriptions = std::collections::BTreeMap::<u64, crate::WsSubscriptionRuntime>::new();
|
||||
let mut remote_to_local = std::collections::BTreeMap::<u64, crate::WsSubscriptionId>::new();
|
||||
let heartbeat_enabled = helius_heartbeat_enabled(endpoint.protocol());
|
||||
let mut heartbeat_deadline = next_helius_heartbeat_deadline();
|
||||
loop {
|
||||
prune_cancelled_pending(id, &mut pending, &mut subscriptions, &mut remote_to_local);
|
||||
let timeout_deadline = next_pending_deadline(&pending);
|
||||
@@ -471,6 +475,13 @@ async fn run_ws_session_actor(
|
||||
)
|
||||
.await
|
||||
},
|
||||
() = tokio::time::sleep_until(heartbeat_deadline), if heartbeat_enabled => {
|
||||
let heartbeat = send_helius_heartbeat(id, &endpoint, &mut websocket, &mut shutdown_rx).await;
|
||||
if matches!(&heartbeat, WsActorIoOutcome::Continue) {
|
||||
heartbeat_deadline = next_helius_heartbeat_deadline();
|
||||
}
|
||||
heartbeat
|
||||
},
|
||||
() = tokio::time::sleep_until(timeout_deadline) => {
|
||||
expire_pending_requests(id, &mut pending, &mut subscriptions, &mut remote_to_local);
|
||||
WsActorIoOutcome::Continue
|
||||
@@ -550,7 +561,12 @@ async fn run_ws_session_actor(
|
||||
)
|
||||
.await;
|
||||
match recovery {
|
||||
WsReconnectOutcome::Connected { websocket: replacement } => websocket = *replacement,
|
||||
WsReconnectOutcome::Connected { websocket: replacement } => {
|
||||
websocket = *replacement;
|
||||
if heartbeat_enabled {
|
||||
heartbeat_deadline = next_helius_heartbeat_deadline();
|
||||
}
|
||||
},
|
||||
WsReconnectOutcome::ShutdownRequested { deadline } => {
|
||||
finish_disconnected_shutdown(
|
||||
id,
|
||||
@@ -611,7 +627,12 @@ async fn run_ws_session_actor(
|
||||
)
|
||||
.await;
|
||||
match recovery {
|
||||
WsReconnectOutcome::Connected { websocket: replacement } => websocket = *replacement,
|
||||
WsReconnectOutcome::Connected { websocket: replacement } => {
|
||||
websocket = *replacement;
|
||||
if heartbeat_enabled {
|
||||
heartbeat_deadline = next_helius_heartbeat_deadline();
|
||||
}
|
||||
},
|
||||
WsReconnectOutcome::ShutdownRequested { deadline } => {
|
||||
finish_disconnected_shutdown(
|
||||
id,
|
||||
@@ -658,6 +679,69 @@ async fn run_ws_session_actor(
|
||||
}
|
||||
}
|
||||
|
||||
const fn helius_heartbeat_enabled(protocol: crate::WsProtocolKind) -> bool {
|
||||
return matches!(protocol, crate::WsProtocolKind::HeliusLaserStream);
|
||||
}
|
||||
|
||||
fn next_helius_heartbeat_deadline() -> tokio::time::Instant {
|
||||
return tokio::time::Instant::now() + HELIUS_WS_HEARTBEAT_INTERVAL;
|
||||
}
|
||||
|
||||
async fn send_helius_heartbeat<S>(
|
||||
id: crate::WsSessionId,
|
||||
endpoint: &crate::WsEndpointSettings,
|
||||
websocket: &mut tokio_tungstenite::WebSocketStream<S>,
|
||||
shutdown_rx: &mut tokio::sync::watch::Receiver<std::option::Option<tokio::time::Instant>>,
|
||||
) -> WsActorIoOutcome
|
||||
where
|
||||
S: tokio::io::AsyncRead + tokio::io::AsyncWrite + std::marker::Unpin,
|
||||
{
|
||||
let message = tokio_tungstenite::tungstenite::Message::Ping(std::vec::Vec::new().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());
|
||||
return WsActorIoOutcome::ShutdownRequested { deadline };
|
||||
},
|
||||
send_result = websocket.send(message) => send_result,
|
||||
() = tokio::time::sleep(endpoint.session().command_timeout()) => {
|
||||
ksp_logging_lib::warn!(
|
||||
target: crate::TRACING_TARGET,
|
||||
session_id = id.get(),
|
||||
endpoint_name = endpoint.name(),
|
||||
"Helius WebSocket heartbeat Ping write timed out"
|
||||
);
|
||||
return WsActorIoOutcome::Failed {
|
||||
code: crate::ERROR_CODE_WS_CONNECTION_FAILED,
|
||||
pending_message: "WebSocket connection failed while writing Helius heartbeat Ping",
|
||||
};
|
||||
},
|
||||
};
|
||||
return match send_result {
|
||||
std::result::Result::Ok(()) => {
|
||||
ksp_logging_lib::trace!(
|
||||
target: crate::TRACING_TARGET,
|
||||
session_id = id.get(),
|
||||
endpoint_name = endpoint.name(),
|
||||
"sent Helius WebSocket heartbeat Ping control frame"
|
||||
);
|
||||
WsActorIoOutcome::Continue
|
||||
},
|
||||
std::result::Result::Err(_) => {
|
||||
ksp_logging_lib::warn!(
|
||||
target: crate::TRACING_TARGET,
|
||||
session_id = id.get(),
|
||||
endpoint_name = endpoint.name(),
|
||||
"Helius WebSocket heartbeat Ping write failed"
|
||||
);
|
||||
WsActorIoOutcome::Failed {
|
||||
code: crate::ERROR_CODE_WS_CONNECTION_FAILED,
|
||||
pending_message: "WebSocket connection failed while writing Helius heartbeat Ping",
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
fn websocket_config(endpoint: &crate::WsEndpointSettings) -> tokio_tungstenite::tungstenite::protocol::WebSocketConfig {
|
||||
return tokio_tungstenite::tungstenite::protocol::WebSocketConfig::default()
|
||||
.write_buffer_size(0)
|
||||
|
||||
Reference in New Issue
Block a user