v0.2.7-pre.008
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
<!-- file: crates/ksp-onchain-transport-lib/README.md -->
|
||||
<!-- version: 13 -->
|
||||
<!-- version: 14 -->
|
||||
|
||||
# `ksp-onchain-transport-lib`
|
||||
|
||||
@@ -101,7 +101,7 @@ La première session physique WebSocket est matérialisée sans introduire de po
|
||||
- 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 a été durci ensuite en `pre.005`.
|
||||
|
||||
Le chemin JSON-RPC générique reste `pub(crate)`. Il sert de primitive au moteur typed de subscriptions et **ne constitue pas une API publique raw provider-extension**. Le registry subscriptions et le mapping remote/local sont matérialisés en `pre.006`; reconnect/resubscribe est actif depuis `pre.007`, tandis que le durcissement backpressure per-sub reste dans la tranche suivante.
|
||||
Le chemin JSON-RPC générique reste `pub(crate)`. Il sert de primitive au moteur typed de subscriptions et **ne constitue pas une API publique raw provider-extension**. Le registry subscriptions et le mapping remote/local sont matérialisés en `pre.006`, reconnect/resubscribe est actif depuis `pre.007` et le backpressure borné par subscription est effectif depuis `pre.008`.
|
||||
|
||||
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.
|
||||
|
||||
@@ -136,7 +136,15 @@ Avec `WsResubscribePolicy::Never`, la session physique peut se reconnecter mais
|
||||
|
||||
Le compteur de continuity gaps est un signal d'observabilité, pas une garantie de livraison. Transport n'ajoute aucun backfill HTTP et ne promet aucune continuité lossless pendant l'intervalle de déconnexion.
|
||||
|
||||
Le durcissement adversarial final de backpressure/leaks reste `pre.008`.
|
||||
### Backpressure et libération de capacité `0.2.7-pre.008`
|
||||
|
||||
Chaque subscription dispose de sa propre queue typed bornée par `notification_queue_capacity`. Le runtime ne droppe jamais silencieusement une notification lorsque cette queue est pleine : il incrémente `WsSessionSnapshot::overflow_count()`, fait passer uniquement le handle lent à `Failed`, publie `ERROR_CODE_WS_BACKPRESSURE_OVERFLOW` via `WsSubscription::terminal_error_code()` et programme un `*Unsubscribe` distant best-effort. Les autres subscriptions et la session physique restent utilisables.
|
||||
|
||||
Les autres terminaisons en échec publient également un code KSP sûr sur le handle : erreur protocolaire, timeout, erreur RPC applicative ou perte physique terminale. Les fermetures normales et les unsubscriptions réussis conservent `terminal_error_code() == None`. Aucun payload distant, remote subscription ID ou endpoint URL n'est projeté dans cette cause.
|
||||
|
||||
`max_active_subscriptions` reste une limite d'admission distincte du compteur d'overflow de notifications : un rejet de création ne l'incrémente pas. Lorsqu'une subscription est fermée, échoue ou que son receiver est abandonné puis détecté sur la notification suivante, son entrée runtime et son binding distant sont nettoyés et la capacité locale redevient réutilisable.
|
||||
|
||||
Les fixtures adversariales prouvent l'isolation d'un consumer lent, la survie d'une subscription saine, le cleanup distant best-effort, la réutilisation de capacité après unsubscribe ou abandon du receiver et la conservation du compteur d'overflow à travers les snapshots. Aucune promesse de livraison lossless n'est ajoutée.
|
||||
|
||||
## Résilience
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
<!-- file: crates/ksp-onchain-transport-lib/USAGE.md -->
|
||||
<!-- version: 13 -->
|
||||
<!-- version: 14 -->
|
||||
|
||||
# Utilisation de `ksp-onchain-transport-lib`
|
||||
|
||||
@@ -106,6 +106,16 @@ Depuis `0.2.7-pre.007`, les settings de session contrôlent réellement le recon
|
||||
|
||||
`unsubscribe().await` peut être appelé pendant `Reconnecting` ou `Resubscribing`. La cancellation locale gagne et le handle ne redevient jamais `Active`. Un ACK distant tardif est nettoyé best-effort par l'actor. Aucun backfill HTTP n'est déclenché automatiquement ; le consumer doit traiter `continuity_gap_count` comme un signal de réconciliation éventuelle.
|
||||
|
||||
### Backpressure par subscription
|
||||
|
||||
Depuis `0.2.7-pre.008`, `WsSessionSettings::notification_queue_capacity()` borne réellement la queue de chaque `WsSubscription<T>`. Le consumer doit donc drainer `recv()` selon son débit métier. Une queue pleine ne bloque pas l'actor et n'affecte pas les autres subscriptions : la subscription lente devient terminale avec `state() == Failed` et `terminal_error_code() == Some(ERROR_CODE_WS_BACKPRESSURE_OVERFLOW)`, tandis que `WsSessionSnapshot::overflow_count()` est incrémenté.
|
||||
|
||||
Un échec terminal non lié à l'overflow expose lui aussi un `ErrorCode` KSP sûr via `terminal_error_code()`. Une fermeture normale conserve `None`. Cette projection ne contient ni payload de notification, ni remote subscription ID, ni URL d'endpoint.
|
||||
|
||||
`max_active_subscriptions` borne séparément le nombre d'entrées logiques enregistrées. Son rejet utilise le même domaine d'erreur de capacité mais n'incrémente pas `overflow_count`, réservé aux queues de notifications saturées. Une subscription fermée ou nettoyée après abandon de son receiver libère sa capacité locale ; l'actor tente aussi de supprimer son binding distant sans rendre ce cleanup bloquant.
|
||||
|
||||
Le consumer doit traiter `overflow_count` et `continuity_gap_count` comme deux signaux distincts : le premier indique une perte locale par saturation d'un consumer, le second une interruption de continuité liée à une reconnexion. Aucun des deux n'implique un replay ou un backfill automatique.
|
||||
|
||||
### Fermeture explicite
|
||||
|
||||
À partir de `0.2.7-pre.005`, fermer explicitement la session est la voie normale de shutdown :
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/src/lib.rs
|
||||
// version: 24
|
||||
// version: 25
|
||||
|
||||
#![warn(missing_docs)]
|
||||
#![deny(unreachable_pub)]
|
||||
@@ -16,10 +16,13 @@
|
||||
//! modern/legacy `getTransaction` coverage. `0.2.4` completes the HTTP surface with all ten Blocks and five Economics wrappers, including complete
|
||||
//! 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. `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. `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.
|
||||
//! `0.2.7-pre.002` adds the provider-neutral WebSocket settings foundation, redacted endpoint URLs, explicit protocol-family discrimination, local session and
|
||||
//! subscription identities, observable lifecycle states and safe snapshots. `0.2.7-pre.004` adds the first physical WebSocket runtime with one
|
||||
//! actor-owned socket,
|
||||
//! bounded handshake, command/pending JSON-RPC flow and deterministic local-server fixtures. `0.2.7-pre.005` adds explicit bounded shutdown, adversarial
|
||||
//! request/frame/message limits and control-frame handling. `0.2.7-pre.006` adds the typed subscription registry with stable local IDs and internal remote-ID
|
||||
//! routing. `0.2.7-pre.007` adds finite reconnect, deterministic resubscribe and continuity-gap tracking. `0.2.7-pre.008` makes per-subscription notification
|
||||
//! backpressure terminal and observable, preserves safe terminal error codes, performs best-effort remote cleanup and proves bounded capacity reuse.
|
||||
|
||||
mod client;
|
||||
mod constants;
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/src/ws_lifecycle.rs
|
||||
// version: 4
|
||||
// version: 5
|
||||
|
||||
/// Stable local identity assigned to one physical WebSocket session.
|
||||
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
|
||||
@@ -170,13 +170,20 @@ pub struct WsSubscriptionSnapshot {
|
||||
kind: crate::WsSubscriptionKind,
|
||||
state: crate::WsSubscriptionState,
|
||||
remote_bound: bool,
|
||||
terminal_error_code: std::option::Option<ksp_core_lib::ErrorCode>,
|
||||
}
|
||||
|
||||
impl WsSubscriptionSnapshot {
|
||||
/// Creates one safe subscription lifecycle projection for Transport runtime internals.
|
||||
#[must_use]
|
||||
pub(crate) const fn new(id: crate::WsSubscriptionId, kind: crate::WsSubscriptionKind, state: crate::WsSubscriptionState, remote_bound: bool) -> Self {
|
||||
return Self { id, kind, state, remote_bound };
|
||||
pub(crate) const fn new(
|
||||
id: crate::WsSubscriptionId,
|
||||
kind: crate::WsSubscriptionKind,
|
||||
state: crate::WsSubscriptionState,
|
||||
remote_bound: bool,
|
||||
terminal_error_code: std::option::Option<ksp_core_lib::ErrorCode>,
|
||||
) -> Self {
|
||||
return Self { id, kind, state, remote_bound, terminal_error_code };
|
||||
}
|
||||
|
||||
/// Returns the stable local subscription identity.
|
||||
@@ -204,6 +211,12 @@ impl WsSubscriptionSnapshot {
|
||||
pub const fn remote_bound(&self) -> bool {
|
||||
return self.remote_bound;
|
||||
}
|
||||
|
||||
/// Returns the safe terminal error code when this subscription ended because of a failure.
|
||||
#[must_use]
|
||||
pub const fn terminal_error_code(&self) -> std::option::Option<ksp_core_lib::ErrorCode> {
|
||||
return self.terminal_error_code;
|
||||
}
|
||||
}
|
||||
|
||||
/// Safe runtime snapshot for one physical WebSocket session.
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/src/ws_session.rs
|
||||
// version: 7
|
||||
// version: 8
|
||||
|
||||
use futures_util::SinkExt; // rust-rules: trait-import
|
||||
use futures_util::StreamExt; // rust-rules: trait-import
|
||||
@@ -275,10 +275,16 @@ enum WsActorIoOutcome {
|
||||
Continue,
|
||||
RemoteClosed,
|
||||
ShutdownRequested { deadline: tokio::time::Instant },
|
||||
StaleSubscribeAck { kind: crate::WsSubscriptionKind, remote_id: u64 },
|
||||
BestEffortUnsubscribe { kind: crate::WsSubscriptionKind, remote_id: u64 },
|
||||
Failed { code: ksp_core_lib::ErrorCode, pending_message: &'static str },
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Default)]
|
||||
struct WsRuntimeCounters {
|
||||
continuity_gap_count: u64,
|
||||
overflow_count: u64,
|
||||
}
|
||||
|
||||
enum WsReconnectOutcome {
|
||||
Connected { websocket: std::boxed::Box<WsPhysicalStream> },
|
||||
ShutdownRequested { deadline: tokio::time::Instant },
|
||||
@@ -323,7 +329,7 @@ async fn run_ws_session_actor(
|
||||
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, 0, &std::collections::BTreeMap::new());
|
||||
publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Failed, 0, WsRuntimeCounters::default(), &std::collections::BTreeMap::new());
|
||||
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!(
|
||||
@@ -337,7 +343,7 @@ async fn run_ws_session_actor(
|
||||
return;
|
||||
},
|
||||
std::result::Result::Err(_) => {
|
||||
publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Failed, 0, 0, &std::collections::BTreeMap::new());
|
||||
publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Failed, 0, WsRuntimeCounters::default(), &std::collections::BTreeMap::new());
|
||||
let error = ws_timeout_error(id, "WebSocket connection handshake timed out");
|
||||
let _ = startup_tx.send(std::result::Result::Err(error));
|
||||
ksp_logging_lib::warn!(
|
||||
@@ -350,7 +356,16 @@ async fn run_ws_session_actor(
|
||||
},
|
||||
};
|
||||
let mut continuity_gap_count = 0_u64;
|
||||
publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Active, 0, continuity_gap_count, &std::collections::BTreeMap::new());
|
||||
let mut overflow_count = 0_u64;
|
||||
publish_snapshot(
|
||||
&snapshot_tx,
|
||||
id,
|
||||
&endpoint,
|
||||
crate::WsSessionState::Active,
|
||||
0,
|
||||
WsRuntimeCounters { continuity_gap_count, overflow_count },
|
||||
&std::collections::BTreeMap::new(),
|
||||
);
|
||||
let _ = startup_tx.send(std::result::Result::Ok(()));
|
||||
ksp_logging_lib::debug!(
|
||||
target: crate::TRACING_TARGET,
|
||||
@@ -388,6 +403,7 @@ async fn run_ws_session_actor(
|
||||
&mut subscriptions,
|
||||
&mut remote_to_local,
|
||||
continuity_gap_count,
|
||||
overflow_count,
|
||||
deadline,
|
||||
)
|
||||
.await;
|
||||
@@ -417,6 +433,7 @@ async fn run_ws_session_actor(
|
||||
&mut pending,
|
||||
&mut subscriptions,
|
||||
&mut remote_to_local,
|
||||
&mut overflow_count,
|
||||
&mut shutdown_rx,
|
||||
)
|
||||
.await
|
||||
@@ -428,10 +445,18 @@ async fn run_ws_session_actor(
|
||||
};
|
||||
match outcome {
|
||||
WsActorIoOutcome::Continue => {
|
||||
publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Active, pending.len(), continuity_gap_count, &subscriptions);
|
||||
publish_snapshot(
|
||||
&snapshot_tx,
|
||||
id,
|
||||
&endpoint,
|
||||
crate::WsSessionState::Active,
|
||||
pending.len(),
|
||||
WsRuntimeCounters { continuity_gap_count, overflow_count },
|
||||
&subscriptions,
|
||||
);
|
||||
},
|
||||
WsActorIoOutcome::StaleSubscribeAck { kind, remote_id } => {
|
||||
let cleanup = best_effort_stale_unsubscribe(id, &endpoint, &mut websocket, &mut shutdown_rx, kind, remote_id).await;
|
||||
WsActorIoOutcome::BestEffortUnsubscribe { kind, remote_id } => {
|
||||
let cleanup = best_effort_remote_unsubscribe(id, &endpoint, &mut websocket, &mut shutdown_rx, kind, remote_id).await;
|
||||
if let WsActorIoOutcome::ShutdownRequested { deadline } = cleanup {
|
||||
close_session_actor(
|
||||
id,
|
||||
@@ -442,12 +467,21 @@ async fn run_ws_session_actor(
|
||||
&mut subscriptions,
|
||||
&mut remote_to_local,
|
||||
continuity_gap_count,
|
||||
overflow_count,
|
||||
deadline,
|
||||
)
|
||||
.await;
|
||||
return;
|
||||
}
|
||||
publish_snapshot(&snapshot_tx, id, &endpoint, crate::WsSessionState::Active, pending.len(), continuity_gap_count, &subscriptions);
|
||||
publish_snapshot(
|
||||
&snapshot_tx,
|
||||
id,
|
||||
&endpoint,
|
||||
crate::WsSessionState::Active,
|
||||
pending.len(),
|
||||
WsRuntimeCounters { continuity_gap_count, overflow_count },
|
||||
&subscriptions,
|
||||
);
|
||||
},
|
||||
WsActorIoOutcome::ShutdownRequested { deadline } => {
|
||||
close_session_actor(
|
||||
@@ -459,6 +493,7 @@ async fn run_ws_session_actor(
|
||||
&mut subscriptions,
|
||||
&mut remote_to_local,
|
||||
continuity_gap_count,
|
||||
overflow_count,
|
||||
deadline,
|
||||
)
|
||||
.await;
|
||||
@@ -476,6 +511,7 @@ async fn run_ws_session_actor(
|
||||
&mut subscriptions,
|
||||
&mut remote_to_local,
|
||||
&mut continuity_gap_count,
|
||||
&mut overflow_count,
|
||||
crate::ERROR_CODE_WS_CONNECTION_FAILED,
|
||||
"Remote peer closed WebSocket session before pending response delivery",
|
||||
)
|
||||
@@ -491,16 +527,35 @@ async fn run_ws_session_actor(
|
||||
&mut subscriptions,
|
||||
&mut remote_to_local,
|
||||
continuity_gap_count,
|
||||
overflow_count,
|
||||
deadline,
|
||||
);
|
||||
return;
|
||||
},
|
||||
WsReconnectOutcome::HandlesDropped => {
|
||||
finish_disconnected_close(id, &endpoint, &snapshot_tx, &mut pending, &mut subscriptions, &mut remote_to_local, continuity_gap_count);
|
||||
finish_disconnected_close(
|
||||
id,
|
||||
&endpoint,
|
||||
&snapshot_tx,
|
||||
&mut pending,
|
||||
&mut subscriptions,
|
||||
&mut remote_to_local,
|
||||
continuity_gap_count,
|
||||
overflow_count,
|
||||
);
|
||||
return;
|
||||
},
|
||||
WsReconnectOutcome::Exhausted => {
|
||||
finish_reconnect_exhaustion(id, &endpoint, &snapshot_tx, &mut pending, &mut subscriptions, &mut remote_to_local, continuity_gap_count);
|
||||
finish_reconnect_exhaustion(
|
||||
id,
|
||||
&endpoint,
|
||||
&snapshot_tx,
|
||||
&mut pending,
|
||||
&mut subscriptions,
|
||||
&mut remote_to_local,
|
||||
continuity_gap_count,
|
||||
overflow_count,
|
||||
);
|
||||
return;
|
||||
},
|
||||
}
|
||||
@@ -517,6 +572,7 @@ async fn run_ws_session_actor(
|
||||
&mut subscriptions,
|
||||
&mut remote_to_local,
|
||||
&mut continuity_gap_count,
|
||||
&mut overflow_count,
|
||||
code,
|
||||
pending_message,
|
||||
)
|
||||
@@ -532,16 +588,35 @@ async fn run_ws_session_actor(
|
||||
&mut subscriptions,
|
||||
&mut remote_to_local,
|
||||
continuity_gap_count,
|
||||
overflow_count,
|
||||
deadline,
|
||||
);
|
||||
return;
|
||||
},
|
||||
WsReconnectOutcome::HandlesDropped => {
|
||||
finish_disconnected_close(id, &endpoint, &snapshot_tx, &mut pending, &mut subscriptions, &mut remote_to_local, continuity_gap_count);
|
||||
finish_disconnected_close(
|
||||
id,
|
||||
&endpoint,
|
||||
&snapshot_tx,
|
||||
&mut pending,
|
||||
&mut subscriptions,
|
||||
&mut remote_to_local,
|
||||
continuity_gap_count,
|
||||
overflow_count,
|
||||
);
|
||||
return;
|
||||
},
|
||||
WsReconnectOutcome::Exhausted => {
|
||||
finish_reconnect_exhaustion(id, &endpoint, &snapshot_tx, &mut pending, &mut subscriptions, &mut remote_to_local, continuity_gap_count);
|
||||
finish_reconnect_exhaustion(
|
||||
id,
|
||||
&endpoint,
|
||||
&snapshot_tx,
|
||||
&mut pending,
|
||||
&mut subscriptions,
|
||||
&mut remote_to_local,
|
||||
continuity_gap_count,
|
||||
overflow_count,
|
||||
);
|
||||
return;
|
||||
},
|
||||
}
|
||||
@@ -570,11 +645,12 @@ async fn recover_websocket_session(
|
||||
subscriptions: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionRuntime>,
|
||||
remote_to_local: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionId>,
|
||||
continuity_gap_count: &mut u64,
|
||||
overflow_count: &mut u64,
|
||||
code: ksp_core_lib::ErrorCode,
|
||||
pending_message: &'static str,
|
||||
) -> WsReconnectOutcome {
|
||||
fail_pending_for_reconnect(pending, id, code, pending_message, subscriptions, remote_to_local);
|
||||
prepare_subscriptions_for_reconnect(endpoint.session().resubscribe(), subscriptions, remote_to_local);
|
||||
prepare_subscriptions_for_reconnect(endpoint.session().resubscribe(), subscriptions, remote_to_local, code);
|
||||
*continuity_gap_count = (*continuity_gap_count).saturating_add(1);
|
||||
let max_retries = endpoint.session().reconnect().max_retries();
|
||||
if max_retries == 0 {
|
||||
@@ -582,7 +658,15 @@ async fn recover_websocket_session(
|
||||
}
|
||||
let mut attempt = 1_u32;
|
||||
while attempt <= max_retries {
|
||||
publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Reconnecting { attempt }, pending.len(), *continuity_gap_count, subscriptions);
|
||||
publish_snapshot(
|
||||
snapshot_tx,
|
||||
id,
|
||||
endpoint,
|
||||
crate::WsSessionState::Reconnecting { attempt },
|
||||
pending.len(),
|
||||
WsRuntimeCounters { continuity_gap_count: *continuity_gap_count, overflow_count: *overflow_count },
|
||||
subscriptions,
|
||||
);
|
||||
ksp_logging_lib::warn!(
|
||||
target: crate::TRACING_TARGET,
|
||||
session_id = id.get(),
|
||||
@@ -631,12 +715,21 @@ async fn recover_websocket_session(
|
||||
subscriptions,
|
||||
remote_to_local,
|
||||
*continuity_gap_count,
|
||||
overflow_count,
|
||||
attempt,
|
||||
)
|
||||
.await;
|
||||
match restore {
|
||||
WsActorIoOutcome::Continue => {
|
||||
publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Active, pending.len(), *continuity_gap_count, subscriptions);
|
||||
publish_snapshot(
|
||||
snapshot_tx,
|
||||
id,
|
||||
endpoint,
|
||||
crate::WsSessionState::Active,
|
||||
pending.len(),
|
||||
WsRuntimeCounters { continuity_gap_count: *continuity_gap_count, overflow_count: *overflow_count },
|
||||
subscriptions,
|
||||
);
|
||||
ksp_logging_lib::debug!(
|
||||
target: crate::TRACING_TARGET,
|
||||
session_id = id.get(),
|
||||
@@ -657,13 +750,13 @@ async fn recover_websocket_session(
|
||||
subscriptions,
|
||||
remote_to_local,
|
||||
);
|
||||
prepare_subscriptions_for_reconnect(endpoint.session().resubscribe(), subscriptions, remote_to_local);
|
||||
prepare_subscriptions_for_reconnect(endpoint.session().resubscribe(), subscriptions, remote_to_local, crate::ERROR_CODE_WS_CONNECTION_FAILED);
|
||||
},
|
||||
WsActorIoOutcome::Failed { code: restore_code, pending_message: restore_message } => {
|
||||
fail_pending_for_reconnect(pending, id, restore_code, restore_message, subscriptions, remote_to_local);
|
||||
prepare_subscriptions_for_reconnect(endpoint.session().resubscribe(), subscriptions, remote_to_local);
|
||||
prepare_subscriptions_for_reconnect(endpoint.session().resubscribe(), subscriptions, remote_to_local, restore_code);
|
||||
},
|
||||
WsActorIoOutcome::StaleSubscribeAck { .. } => {},
|
||||
WsActorIoOutcome::BestEffortUnsubscribe { .. } => {},
|
||||
}
|
||||
attempt = attempt.saturating_add(1);
|
||||
}
|
||||
@@ -683,6 +776,7 @@ async fn restore_subscriptions_after_reconnect(
|
||||
subscriptions: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionRuntime>,
|
||||
remote_to_local: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionId>,
|
||||
continuity_gap_count: u64,
|
||||
overflow_count: &mut u64,
|
||||
attempt: u32,
|
||||
) -> WsActorIoOutcome {
|
||||
let restore_ids = subscriptions
|
||||
@@ -700,13 +794,14 @@ async fn restore_subscriptions_after_reconnect(
|
||||
let (request_id, deadline) = match write_result {
|
||||
std::result::Result::Ok(result) => result,
|
||||
std::result::Result::Err((error, WsActorIoOutcome::Continue)) => {
|
||||
close_local_subscription(subscriptions, remote_to_local, subscription_id, crate::WsSubscriptionState::Failed);
|
||||
let error_code = error.code();
|
||||
fail_local_subscription(subscriptions, remote_to_local, subscription_id, error_code);
|
||||
ksp_logging_lib::warn!(
|
||||
target: crate::TRACING_TARGET,
|
||||
session_id = id.get(),
|
||||
subscription_id = subscription_id.get(),
|
||||
subscription_kind = kind.as_str(),
|
||||
error_code = error.code().code(),
|
||||
error_code = error_code.code(),
|
||||
"logical WebSocket resubscribe could not be queued"
|
||||
);
|
||||
continue;
|
||||
@@ -752,7 +847,7 @@ async fn restore_subscriptions_after_reconnect(
|
||||
}
|
||||
},
|
||||
maybe_message = websocket.next() => {
|
||||
handle_socket_message(id, endpoint, maybe_message, websocket, pending, subscriptions, remote_to_local, shutdown_rx).await
|
||||
handle_socket_message(id, endpoint, maybe_message, websocket, pending, subscriptions, remote_to_local, overflow_count, shutdown_rx).await
|
||||
},
|
||||
() = tokio::time::sleep_until(timeout_deadline) => {
|
||||
expire_pending_requests(id, pending, subscriptions, remote_to_local);
|
||||
@@ -761,15 +856,23 @@ async fn restore_subscriptions_after_reconnect(
|
||||
};
|
||||
match outcome {
|
||||
WsActorIoOutcome::Continue => {},
|
||||
WsActorIoOutcome::StaleSubscribeAck { kind: stale_kind, remote_id } => {
|
||||
let cleanup = best_effort_stale_unsubscribe(id, endpoint, websocket, shutdown_rx, stale_kind, remote_id).await;
|
||||
WsActorIoOutcome::BestEffortUnsubscribe { kind: stale_kind, remote_id } => {
|
||||
let cleanup = best_effort_remote_unsubscribe(id, endpoint, websocket, shutdown_rx, stale_kind, remote_id).await;
|
||||
if let WsActorIoOutcome::ShutdownRequested { deadline } = cleanup {
|
||||
return WsActorIoOutcome::ShutdownRequested { deadline };
|
||||
}
|
||||
},
|
||||
_ => return outcome,
|
||||
}
|
||||
publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Reconnecting { attempt }, pending.len(), continuity_gap_count, subscriptions);
|
||||
publish_snapshot(
|
||||
snapshot_tx,
|
||||
id,
|
||||
endpoint,
|
||||
crate::WsSessionState::Reconnecting { attempt },
|
||||
pending.len(),
|
||||
WsRuntimeCounters { continuity_gap_count, overflow_count: *overflow_count },
|
||||
subscriptions,
|
||||
);
|
||||
}
|
||||
}
|
||||
return WsActorIoOutcome::Continue;
|
||||
@@ -791,7 +894,7 @@ async fn handle_reconnecting_command_with_socket(
|
||||
});
|
||||
close_local_subscription(subscriptions, remote_to_local, subscription_id, crate::WsSubscriptionState::Closed);
|
||||
if let std::option::Option::Some((kind, remote_id)) = cleanup {
|
||||
let cleanup_outcome = best_effort_stale_unsubscribe(id, endpoint, websocket, shutdown_rx, kind, remote_id).await;
|
||||
let cleanup_outcome = best_effort_remote_unsubscribe(id, endpoint, websocket, shutdown_rx, kind, remote_id).await;
|
||||
if let WsActorIoOutcome::ShutdownRequested { deadline } = cleanup_outcome {
|
||||
let _ = response_tx.send(std::result::Result::Ok(false));
|
||||
return WsActorIoOutcome::ShutdownRequested { deadline };
|
||||
@@ -943,9 +1046,11 @@ fn prepare_subscriptions_for_reconnect(
|
||||
policy: crate::WsResubscribePolicy,
|
||||
subscriptions: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionRuntime>,
|
||||
remote_to_local: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionId>,
|
||||
failure_code: ksp_core_lib::ErrorCode,
|
||||
) {
|
||||
remote_to_local.clear();
|
||||
let mut terminal = std::vec::Vec::new();
|
||||
let mut failed = std::vec::Vec::new();
|
||||
let mut already_terminal = std::vec::Vec::new();
|
||||
for runtime in subscriptions.values_mut() {
|
||||
runtime.remote_id = std::option::Option::None;
|
||||
match runtime.state {
|
||||
@@ -953,19 +1058,18 @@ fn prepare_subscriptions_for_reconnect(
|
||||
if policy == crate::WsResubscribePolicy::ActiveSubscriptions {
|
||||
runtime.set_state(crate::WsSubscriptionState::Resubscribing);
|
||||
} else {
|
||||
runtime.set_state(crate::WsSubscriptionState::Failed);
|
||||
terminal.push(runtime.id);
|
||||
failed.push(runtime.id);
|
||||
}
|
||||
},
|
||||
crate::WsSubscriptionState::Requested | crate::WsSubscriptionState::Cancelling => {
|
||||
runtime.set_state(crate::WsSubscriptionState::Failed);
|
||||
terminal.push(runtime.id);
|
||||
},
|
||||
crate::WsSubscriptionState::Closed | crate::WsSubscriptionState::Failed => terminal.push(runtime.id),
|
||||
crate::WsSubscriptionState::Requested | crate::WsSubscriptionState::Cancelling => failed.push(runtime.id),
|
||||
crate::WsSubscriptionState::Closed | crate::WsSubscriptionState::Failed => already_terminal.push(runtime.id),
|
||||
}
|
||||
}
|
||||
for subscription_id in terminal {
|
||||
close_local_subscription(subscriptions, remote_to_local, subscription_id, crate::WsSubscriptionState::Failed);
|
||||
for subscription_id in failed {
|
||||
fail_local_subscription(subscriptions, remote_to_local, subscription_id, failure_code);
|
||||
}
|
||||
for subscription_id in already_terminal {
|
||||
subscriptions.remove(&subscription_id.get());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -998,7 +1102,7 @@ fn fail_pending_for_reconnect(
|
||||
let _ = response_tx.send(std::result::Result::Err(error));
|
||||
},
|
||||
PendingWsResponse::Subscribe { subscription_id, response_tx } => {
|
||||
close_local_subscription(subscriptions, remote_to_local, subscription_id, crate::WsSubscriptionState::Failed);
|
||||
fail_local_subscription(subscriptions, remote_to_local, subscription_id, code);
|
||||
let _ = response_tx.send(std::result::Result::Err(error));
|
||||
},
|
||||
PendingWsResponse::Resubscribe { .. } => {},
|
||||
@@ -1010,7 +1114,7 @@ fn fail_pending_for_reconnect(
|
||||
}
|
||||
}
|
||||
|
||||
async fn best_effort_stale_unsubscribe(
|
||||
async fn best_effort_remote_unsubscribe(
|
||||
id: crate::WsSessionId,
|
||||
endpoint: &crate::WsEndpointSettings,
|
||||
websocket: &mut WsPhysicalStream,
|
||||
@@ -1041,7 +1145,7 @@ async fn best_effort_stale_unsubscribe(
|
||||
target: crate::TRACING_TARGET,
|
||||
session_id = id.get(),
|
||||
subscription_kind = kind.as_str(),
|
||||
"sent best-effort unsubscribe for stale remote subscription acknowledgement"
|
||||
"sent best-effort remote WebSocket unsubscribe"
|
||||
);
|
||||
}
|
||||
return WsActorIoOutcome::Continue;
|
||||
@@ -1055,12 +1159,21 @@ fn finish_disconnected_shutdown(
|
||||
subscriptions: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionRuntime>,
|
||||
remote_to_local: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionId>,
|
||||
continuity_gap_count: u64,
|
||||
overflow_count: u64,
|
||||
_deadline: tokio::time::Instant,
|
||||
) {
|
||||
publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Closing, pending.len(), continuity_gap_count, subscriptions);
|
||||
publish_snapshot(
|
||||
snapshot_tx,
|
||||
id,
|
||||
endpoint,
|
||||
crate::WsSessionState::Closing,
|
||||
pending.len(),
|
||||
WsRuntimeCounters { continuity_gap_count, overflow_count },
|
||||
subscriptions,
|
||||
);
|
||||
fail_all_pending(pending, id, crate::ERROR_CODE_WS_SESSION_CLOSED, "WebSocket session shutdown cancelled pending request");
|
||||
terminate_all_subscriptions(subscriptions, remote_to_local, crate::WsSubscriptionState::Closed);
|
||||
publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Closed, 0, continuity_gap_count, subscriptions);
|
||||
publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Closed, 0, WsRuntimeCounters { continuity_gap_count, overflow_count }, subscriptions);
|
||||
}
|
||||
|
||||
fn finish_disconnected_close(
|
||||
@@ -1071,10 +1184,11 @@ fn finish_disconnected_close(
|
||||
subscriptions: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionRuntime>,
|
||||
remote_to_local: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionId>,
|
||||
continuity_gap_count: u64,
|
||||
overflow_count: u64,
|
||||
) {
|
||||
fail_all_pending(pending, id, crate::ERROR_CODE_WS_SESSION_CLOSED, "WebSocket session handles were dropped during reconnect");
|
||||
terminate_all_subscriptions(subscriptions, remote_to_local, crate::WsSubscriptionState::Closed);
|
||||
publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Closed, 0, continuity_gap_count, subscriptions);
|
||||
publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Closed, 0, WsRuntimeCounters { continuity_gap_count, overflow_count }, subscriptions);
|
||||
}
|
||||
|
||||
fn finish_reconnect_exhaustion(
|
||||
@@ -1085,10 +1199,11 @@ fn finish_reconnect_exhaustion(
|
||||
subscriptions: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionRuntime>,
|
||||
remote_to_local: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionId>,
|
||||
continuity_gap_count: u64,
|
||||
overflow_count: u64,
|
||||
) {
|
||||
fail_all_pending(pending, id, crate::ERROR_CODE_WS_CONNECTION_FAILED, "WebSocket reconnect budget was exhausted");
|
||||
terminate_all_subscriptions(subscriptions, remote_to_local, crate::WsSubscriptionState::Failed);
|
||||
publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Failed, 0, continuity_gap_count, subscriptions);
|
||||
fail_all_subscriptions(subscriptions, remote_to_local, crate::ERROR_CODE_WS_CONNECTION_FAILED);
|
||||
publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Failed, 0, WsRuntimeCounters { continuity_gap_count, overflow_count }, subscriptions);
|
||||
ksp_logging_lib::warn!(
|
||||
target: crate::TRACING_TARGET,
|
||||
session_id = id.get(),
|
||||
@@ -1154,6 +1269,7 @@ where
|
||||
},
|
||||
};
|
||||
let (state_tx, _) = tokio::sync::watch::channel(crate::WsSubscriptionState::Requested);
|
||||
let (terminal_error_tx, _) = tokio::sync::watch::channel(std::option::Option::None::<ksp_core_lib::ErrorCode>);
|
||||
let resubscribe_params = params.clone();
|
||||
subscriptions.insert(
|
||||
subscription_id.get(),
|
||||
@@ -1164,6 +1280,7 @@ where
|
||||
params: resubscribe_params,
|
||||
remote_id: std::option::Option::None,
|
||||
state_tx,
|
||||
terminal_error_tx,
|
||||
dispatcher,
|
||||
},
|
||||
);
|
||||
@@ -1182,9 +1299,8 @@ where
|
||||
WsActorIoOutcome::Continue
|
||||
},
|
||||
std::result::Result::Err((error, outcome)) => {
|
||||
if let std::option::Option::Some(mut runtime) = subscriptions.remove(&subscription_id.get()) {
|
||||
runtime.set_state(crate::WsSubscriptionState::Failed);
|
||||
}
|
||||
let error_code = error.code();
|
||||
fail_local_subscription(subscriptions, remote_to_local, subscription_id, error_code);
|
||||
let _ = response_tx.send(std::result::Result::Err(error));
|
||||
outcome
|
||||
},
|
||||
@@ -1367,6 +1483,7 @@ async fn handle_socket_message<S>(
|
||||
pending: &mut std::collections::BTreeMap<u64, PendingWsRequest>,
|
||||
subscriptions: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionRuntime>,
|
||||
remote_to_local: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionId>,
|
||||
overflow_count: &mut u64,
|
||||
shutdown_rx: &mut tokio::sync::watch::Receiver<std::option::Option<tokio::time::Instant>>,
|
||||
) -> WsActorIoOutcome
|
||||
where
|
||||
@@ -1399,7 +1516,7 @@ where
|
||||
},
|
||||
};
|
||||
return match message {
|
||||
tokio_tungstenite::tungstenite::Message::Text(text) => handle_text_message(id, text.as_str(), pending, subscriptions, remote_to_local),
|
||||
tokio_tungstenite::tungstenite::Message::Text(text) => handle_text_message(id, text.as_str(), pending, subscriptions, remote_to_local, overflow_count),
|
||||
tokio_tungstenite::tungstenite::Message::Binary(_) => {
|
||||
ksp_logging_lib::warn!(target: crate::TRACING_TARGET, session_id = id.get(), "received unexpected binary WebSocket message");
|
||||
WsActorIoOutcome::Failed {
|
||||
@@ -1453,6 +1570,7 @@ fn handle_text_message(
|
||||
pending: &mut std::collections::BTreeMap<u64, PendingWsRequest>,
|
||||
subscriptions: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionRuntime>,
|
||||
remote_to_local: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionId>,
|
||||
overflow_count: &mut u64,
|
||||
) -> WsActorIoOutcome {
|
||||
let decoded = serde_json::from_str::<serde_json::Value>(text);
|
||||
let value = match decoded {
|
||||
@@ -1473,7 +1591,7 @@ fn handle_text_message(
|
||||
},
|
||||
};
|
||||
if !object.contains_key("id") {
|
||||
return handle_subscription_notification(id, object, subscriptions, remote_to_local);
|
||||
return handle_subscription_notification(id, object, subscriptions, remote_to_local, overflow_count);
|
||||
}
|
||||
let response_id = match object.get("id").and_then(serde_json::Value::as_u64) {
|
||||
std::option::Option::Some(response_id) => response_id,
|
||||
@@ -1538,7 +1656,7 @@ fn dispatch_pending_response(
|
||||
let remote_id = match value.as_u64() {
|
||||
std::option::Option::Some(remote_id) => remote_id,
|
||||
std::option::Option::None => {
|
||||
close_local_subscription(subscriptions, remote_to_local, subscription_id, crate::WsSubscriptionState::Failed);
|
||||
fail_local_subscription(subscriptions, remote_to_local, subscription_id, crate::ERROR_CODE_WS_PROTOCOL_ERROR);
|
||||
let error = ksp_core_lib::Error::new(
|
||||
crate::ERROR_CODE_WS_PROTOCOL_ERROR,
|
||||
"WebSocket subscribe response did not contain a numeric remote subscription id",
|
||||
@@ -1553,7 +1671,7 @@ fn dispatch_pending_response(
|
||||
},
|
||||
};
|
||||
if remote_to_local.contains_key(&remote_id) {
|
||||
close_local_subscription(subscriptions, remote_to_local, subscription_id, crate::WsSubscriptionState::Failed);
|
||||
fail_local_subscription(subscriptions, remote_to_local, subscription_id, crate::ERROR_CODE_WS_PROTOCOL_ERROR);
|
||||
let error = ksp_core_lib::Error::new(crate::ERROR_CODE_WS_PROTOCOL_ERROR, "WebSocket endpoint reused an active remote subscription id")
|
||||
.with_context("session_id", id.get().to_string())
|
||||
.with_context("subscription_id", subscription_id.get().to_string());
|
||||
@@ -1568,7 +1686,12 @@ fn dispatch_pending_response(
|
||||
runtime.remote_id = std::option::Option::Some(remote_id);
|
||||
runtime.set_state(crate::WsSubscriptionState::Active);
|
||||
remote_to_local.insert(remote_id, subscription_id);
|
||||
crate::WsSubscriptionRegistration::new(subscription_id, runtime.kind, runtime.state_tx.subscribe())
|
||||
crate::WsSubscriptionRegistration::new(
|
||||
subscription_id,
|
||||
runtime.kind,
|
||||
runtime.state_tx.subscribe(),
|
||||
runtime.terminal_error_tx.subscribe(),
|
||||
)
|
||||
},
|
||||
std::option::Option::None => {
|
||||
let error = ksp_core_lib::Error::new(
|
||||
@@ -1593,7 +1716,8 @@ fn dispatch_pending_response(
|
||||
);
|
||||
},
|
||||
std::result::Result::Err(error) => {
|
||||
close_local_subscription(subscriptions, remote_to_local, subscription_id, crate::WsSubscriptionState::Failed);
|
||||
let error_code = error.code();
|
||||
fail_local_subscription(subscriptions, remote_to_local, subscription_id, error_code);
|
||||
let _ = response_tx.send(std::result::Result::Err(error.with_context("method", pending_request.method)));
|
||||
},
|
||||
},
|
||||
@@ -1602,7 +1726,7 @@ fn dispatch_pending_response(
|
||||
let remote_id = match value.as_u64() {
|
||||
std::option::Option::Some(remote_id) => remote_id,
|
||||
std::option::Option::None => {
|
||||
close_local_subscription(subscriptions, remote_to_local, subscription_id, crate::WsSubscriptionState::Failed);
|
||||
fail_local_subscription(subscriptions, remote_to_local, subscription_id, crate::ERROR_CODE_WS_PROTOCOL_ERROR);
|
||||
return WsActorIoOutcome::Failed {
|
||||
code: crate::ERROR_CODE_WS_PROTOCOL_ERROR,
|
||||
pending_message: "WebSocket resubscribe response violated protocol invariants",
|
||||
@@ -1610,7 +1734,7 @@ fn dispatch_pending_response(
|
||||
},
|
||||
};
|
||||
if remote_to_local.contains_key(&remote_id) {
|
||||
close_local_subscription(subscriptions, remote_to_local, subscription_id, crate::WsSubscriptionState::Failed);
|
||||
fail_local_subscription(subscriptions, remote_to_local, subscription_id, crate::ERROR_CODE_WS_PROTOCOL_ERROR);
|
||||
return WsActorIoOutcome::Failed {
|
||||
code: crate::ERROR_CODE_WS_PROTOCOL_ERROR,
|
||||
pending_message: "WebSocket endpoint reused an active remote subscription id during resubscribe",
|
||||
@@ -1626,7 +1750,7 @@ fn dispatch_pending_response(
|
||||
subscription_kind = kind.as_str(),
|
||||
"late WebSocket resubscribe acknowledgement lost to local cancellation"
|
||||
);
|
||||
return WsActorIoOutcome::StaleSubscribeAck { kind, remote_id };
|
||||
return WsActorIoOutcome::BestEffortUnsubscribe { kind, remote_id };
|
||||
}
|
||||
if let std::option::Option::Some(runtime) = subscriptions.get_mut(&subscription_id.get()) {
|
||||
runtime.remote_id = std::option::Option::Some(remote_id);
|
||||
@@ -1641,13 +1765,15 @@ fn dispatch_pending_response(
|
||||
"logical WebSocket subscription restored with a new remote binding"
|
||||
);
|
||||
},
|
||||
std::result::Result::Err(_) => {
|
||||
close_local_subscription(subscriptions, remote_to_local, subscription_id, crate::WsSubscriptionState::Failed);
|
||||
std::result::Result::Err(error) => {
|
||||
let error_code = error.code();
|
||||
fail_local_subscription(subscriptions, remote_to_local, subscription_id, error_code);
|
||||
ksp_logging_lib::warn!(
|
||||
target: crate::TRACING_TARGET,
|
||||
session_id = id.get(),
|
||||
subscription_id = subscription_id.get(),
|
||||
subscription_kind = kind.as_str(),
|
||||
error_code = error_code.code(),
|
||||
"remote WebSocket resubscribe returned an application error"
|
||||
);
|
||||
},
|
||||
@@ -1700,6 +1826,7 @@ fn handle_subscription_notification(
|
||||
object: &serde_json::Map<std::string::String, serde_json::Value>,
|
||||
subscriptions: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionRuntime>,
|
||||
remote_to_local: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionId>,
|
||||
overflow_count: &mut u64,
|
||||
) -> WsActorIoOutcome {
|
||||
if object.get("jsonrpc").and_then(serde_json::Value::as_str) != std::option::Option::Some("2.0") {
|
||||
return WsActorIoOutcome::Failed {
|
||||
@@ -1767,38 +1894,50 @@ fn handle_subscription_notification(
|
||||
subscription_kind = subscription_kind.as_str(),
|
||||
"WebSocket notification method mismatched the registered subscription family"
|
||||
);
|
||||
close_local_subscription(subscriptions, remote_to_local, local_id, crate::WsSubscriptionState::Failed);
|
||||
return WsActorIoOutcome::Continue;
|
||||
fail_local_subscription(subscriptions, remote_to_local, local_id, crate::ERROR_CODE_WS_PROTOCOL_ERROR);
|
||||
return WsActorIoOutcome::BestEffortUnsubscribe { kind: subscription_kind, remote_id };
|
||||
}
|
||||
let subscription_kind = runtime.kind;
|
||||
let dispatch = (runtime.dispatcher)(result);
|
||||
match dispatch {
|
||||
crate::WsNotificationDispatchOutcome::Delivered => {},
|
||||
return match dispatch {
|
||||
crate::WsNotificationDispatchOutcome::Delivered => WsActorIoOutcome::Continue,
|
||||
crate::WsNotificationDispatchOutcome::ReceiverClosed => {
|
||||
close_local_subscription(subscriptions, remote_to_local, local_id, crate::WsSubscriptionState::Closed);
|
||||
ksp_logging_lib::debug!(
|
||||
target: crate::TRACING_TARGET,
|
||||
session_id = id.get(),
|
||||
subscription_id = local_id.get(),
|
||||
subscription_kind = subscription_kind.as_str(),
|
||||
"logical WebSocket subscription receiver was dropped; scheduling remote cleanup"
|
||||
);
|
||||
WsActorIoOutcome::BestEffortUnsubscribe { kind: subscription_kind, remote_id }
|
||||
},
|
||||
crate::WsNotificationDispatchOutcome::QueueFull => {
|
||||
*overflow_count = (*overflow_count).saturating_add(1);
|
||||
ksp_logging_lib::warn!(
|
||||
target: crate::TRACING_TARGET,
|
||||
session_id = id.get(),
|
||||
subscription_id = local_id.get(),
|
||||
subscription_kind = subscription_kind.as_str(),
|
||||
"bounded WebSocket subscription notification queue is full"
|
||||
overflow_count = *overflow_count,
|
||||
"bounded WebSocket subscription notification queue overflowed; failing only the slow subscription"
|
||||
);
|
||||
close_local_subscription(subscriptions, remote_to_local, local_id, crate::WsSubscriptionState::Failed);
|
||||
fail_local_subscription(subscriptions, remote_to_local, local_id, crate::ERROR_CODE_WS_BACKPRESSURE_OVERFLOW);
|
||||
WsActorIoOutcome::BestEffortUnsubscribe { kind: subscription_kind, remote_id }
|
||||
},
|
||||
crate::WsNotificationDispatchOutcome::DecodeFailed => {
|
||||
crate::WsNotificationDispatchOutcome::DecodeFailed { code } => {
|
||||
ksp_logging_lib::warn!(
|
||||
target: crate::TRACING_TARGET,
|
||||
session_id = id.get(),
|
||||
subscription_id = local_id.get(),
|
||||
subscription_kind = subscription_kind.as_str(),
|
||||
error_code = code.code(),
|
||||
"typed WebSocket notification decoding failed for one subscription"
|
||||
);
|
||||
close_local_subscription(subscriptions, remote_to_local, local_id, crate::WsSubscriptionState::Failed);
|
||||
fail_local_subscription(subscriptions, remote_to_local, local_id, code);
|
||||
WsActorIoOutcome::BestEffortUnsubscribe { kind: subscription_kind, remote_id }
|
||||
},
|
||||
}
|
||||
return WsActorIoOutcome::Continue;
|
||||
};
|
||||
}
|
||||
|
||||
async fn close_session_actor<S>(
|
||||
@@ -1810,14 +1949,23 @@ async fn close_session_actor<S>(
|
||||
subscriptions: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionRuntime>,
|
||||
remote_to_local: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionId>,
|
||||
continuity_gap_count: u64,
|
||||
overflow_count: u64,
|
||||
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(), continuity_gap_count, subscriptions);
|
||||
publish_snapshot(
|
||||
snapshot_tx,
|
||||
id,
|
||||
endpoint,
|
||||
crate::WsSessionState::Closing,
|
||||
pending.len(),
|
||||
WsRuntimeCounters { continuity_gap_count, overflow_count },
|
||||
subscriptions,
|
||||
);
|
||||
fail_all_pending(pending, id, crate::ERROR_CODE_WS_SESSION_CLOSED, "WebSocket session shutdown cancelled the pending request");
|
||||
terminate_all_subscriptions(subscriptions, remote_to_local, crate::WsSubscriptionState::Closed);
|
||||
publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Closing, 0, continuity_gap_count, subscriptions);
|
||||
publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Closing, 0, WsRuntimeCounters { continuity_gap_count, overflow_count }, subscriptions);
|
||||
ksp_logging_lib::debug!(
|
||||
target: crate::TRACING_TARGET,
|
||||
session_id = id.get(),
|
||||
@@ -1839,7 +1987,7 @@ async fn close_session_actor<S>(
|
||||
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, continuity_gap_count, subscriptions);
|
||||
publish_snapshot(snapshot_tx, id, endpoint, crate::WsSessionState::Closed, 0, WsRuntimeCounters { continuity_gap_count, overflow_count }, subscriptions);
|
||||
ksp_logging_lib::debug!(target: crate::TRACING_TARGET, session_id = id.get(), endpoint_name = endpoint.name(), "physical WebSocket session is closed");
|
||||
}
|
||||
|
||||
@@ -1921,11 +2069,11 @@ fn expire_pending_response(
|
||||
let _ = response_tx.send(std::result::Result::Err(error));
|
||||
},
|
||||
PendingWsResponse::Subscribe { subscription_id, response_tx } => {
|
||||
close_local_subscription(subscriptions, remote_to_local, subscription_id, crate::WsSubscriptionState::Failed);
|
||||
fail_local_subscription(subscriptions, remote_to_local, subscription_id, crate::ERROR_CODE_TIMEOUT);
|
||||
let _ = response_tx.send(std::result::Result::Err(error));
|
||||
},
|
||||
PendingWsResponse::Resubscribe { subscription_id, .. } => {
|
||||
close_local_subscription(subscriptions, remote_to_local, subscription_id, crate::WsSubscriptionState::Failed);
|
||||
fail_local_subscription(subscriptions, remote_to_local, subscription_id, crate::ERROR_CODE_TIMEOUT);
|
||||
},
|
||||
PendingWsResponse::Unsubscribe { subscription_id, response_tx } => {
|
||||
close_local_subscription(subscriptions, remote_to_local, subscription_id, crate::WsSubscriptionState::Closed);
|
||||
@@ -2026,6 +2174,32 @@ fn terminate_all_subscriptions(
|
||||
remote_to_local.clear();
|
||||
}
|
||||
|
||||
fn fail_all_subscriptions(
|
||||
subscriptions: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionRuntime>,
|
||||
remote_to_local: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionId>,
|
||||
error_code: ksp_core_lib::ErrorCode,
|
||||
) {
|
||||
for runtime in subscriptions.values_mut() {
|
||||
runtime.remote_id = std::option::Option::None;
|
||||
runtime.fail_with_code(error_code);
|
||||
}
|
||||
remote_to_local.clear();
|
||||
}
|
||||
|
||||
fn fail_local_subscription(
|
||||
subscriptions: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionRuntime>,
|
||||
remote_to_local: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionId>,
|
||||
subscription_id: crate::WsSubscriptionId,
|
||||
error_code: ksp_core_lib::ErrorCode,
|
||||
) {
|
||||
if let std::option::Option::Some(mut runtime) = subscriptions.remove(&subscription_id.get()) {
|
||||
if let std::option::Option::Some(remote_id) = runtime.remote_id.take() {
|
||||
remote_to_local.remove(&remote_id);
|
||||
}
|
||||
runtime.fail_with_code(error_code);
|
||||
}
|
||||
}
|
||||
|
||||
fn close_local_subscription(
|
||||
subscriptions: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionRuntime>,
|
||||
remote_to_local: &mut std::collections::BTreeMap<u64, crate::WsSubscriptionId>,
|
||||
@@ -2046,7 +2220,7 @@ fn publish_snapshot(
|
||||
endpoint: &crate::WsEndpointSettings,
|
||||
state: crate::WsSessionState,
|
||||
pending_request_count: usize,
|
||||
continuity_gap_count: u64,
|
||||
counters: WsRuntimeCounters,
|
||||
subscriptions: &std::collections::BTreeMap<u64, crate::WsSubscriptionRuntime>,
|
||||
) {
|
||||
let subscription_snapshots = subscriptions.values().map(crate::WsSubscriptionRuntime::snapshot).collect();
|
||||
@@ -2058,8 +2232,8 @@ fn publish_snapshot(
|
||||
endpoint.protocol(),
|
||||
state,
|
||||
pending_request_count,
|
||||
continuity_gap_count,
|
||||
0,
|
||||
counters.continuity_gap_count,
|
||||
counters.overflow_count,
|
||||
subscription_snapshots,
|
||||
);
|
||||
snapshot_tx.send_replace(snapshot);
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/src/ws_subscription.rs
|
||||
// version: 2
|
||||
// version: 3
|
||||
|
||||
/// Typed handle for one logical Solana WebSocket subscription.
|
||||
///
|
||||
@@ -11,6 +11,7 @@ pub struct WsSubscription<T> {
|
||||
kind: crate::WsSubscriptionKind,
|
||||
notification_rx: tokio::sync::mpsc::Receiver<ksp_core_lib::Result<T>>,
|
||||
state_rx: tokio::sync::watch::Receiver<crate::WsSubscriptionState>,
|
||||
terminal_error_rx: tokio::sync::watch::Receiver<std::option::Option<ksp_core_lib::ErrorCode>>,
|
||||
command_tx: tokio::sync::mpsc::Sender<crate::WsSessionCommand>,
|
||||
command_timeout: std::time::Duration,
|
||||
}
|
||||
@@ -30,6 +31,7 @@ impl<T> WsSubscription<T> {
|
||||
kind: registration.kind,
|
||||
notification_rx,
|
||||
state_rx: registration.state_rx,
|
||||
terminal_error_rx: registration.terminal_error_rx,
|
||||
command_tx,
|
||||
command_timeout,
|
||||
};
|
||||
@@ -53,6 +55,14 @@ impl<T> WsSubscription<T> {
|
||||
return *self.state_rx.borrow();
|
||||
}
|
||||
|
||||
/// Returns the safe terminal error code when this logical subscription failed.
|
||||
///
|
||||
/// Successful local cancellation and normal completion use `None`. The value never contains remote payloads or endpoint credentials.
|
||||
#[must_use]
|
||||
pub fn terminal_error_code(&self) -> std::option::Option<ksp_core_lib::ErrorCode> {
|
||||
return *self.terminal_error_rx.borrow();
|
||||
}
|
||||
|
||||
/// Receives the next typed notification or terminal typed-decoding error.
|
||||
///
|
||||
/// The underlying queue is bounded by `WsSessionSettings::notification_queue_capacity`. `None` means the actor closed this logical subscription and no
|
||||
@@ -116,7 +126,7 @@ pub(crate) enum WsNotificationDispatchOutcome {
|
||||
Delivered,
|
||||
ReceiverClosed,
|
||||
QueueFull,
|
||||
DecodeFailed,
|
||||
DecodeFailed { code: ksp_core_lib::ErrorCode },
|
||||
}
|
||||
|
||||
/// Type-erased actor-owned dispatcher for one heterogeneous typed notification channel.
|
||||
@@ -130,16 +140,24 @@ where
|
||||
{
|
||||
let (notification_tx, notification_rx) = tokio::sync::mpsc::channel(capacity);
|
||||
let dispatcher = move |value: serde_json::Value| -> WsNotificationDispatchOutcome {
|
||||
let decoded = decoder(value);
|
||||
return match decoded {
|
||||
std::result::Result::Ok(notification) => match notification_tx.try_send(std::result::Result::Ok(notification)) {
|
||||
std::result::Result::Ok(()) => WsNotificationDispatchOutcome::Delivered,
|
||||
std::result::Result::Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) => WsNotificationDispatchOutcome::ReceiverClosed,
|
||||
std::result::Result::Err(tokio::sync::mpsc::error::TrySendError::Full(_)) => WsNotificationDispatchOutcome::QueueFull,
|
||||
let permit = match notification_tx.try_reserve() {
|
||||
std::result::Result::Ok(permit) => permit,
|
||||
std::result::Result::Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) => {
|
||||
return WsNotificationDispatchOutcome::ReceiverClosed;
|
||||
},
|
||||
std::result::Result::Err(tokio::sync::mpsc::error::TrySendError::Full(_)) => {
|
||||
return WsNotificationDispatchOutcome::QueueFull;
|
||||
},
|
||||
};
|
||||
return match decoder(value) {
|
||||
std::result::Result::Ok(notification) => {
|
||||
permit.send(std::result::Result::Ok(notification));
|
||||
WsNotificationDispatchOutcome::Delivered
|
||||
},
|
||||
std::result::Result::Err(error) => {
|
||||
let _ = notification_tx.try_send(std::result::Result::Err(error));
|
||||
WsNotificationDispatchOutcome::DecodeFailed
|
||||
let code = error.code();
|
||||
permit.send(std::result::Result::Err(error));
|
||||
WsNotificationDispatchOutcome::DecodeFailed { code }
|
||||
},
|
||||
};
|
||||
};
|
||||
@@ -151,6 +169,7 @@ pub(crate) struct WsSubscriptionRegistration {
|
||||
id: crate::WsSubscriptionId,
|
||||
kind: crate::WsSubscriptionKind,
|
||||
state_rx: tokio::sync::watch::Receiver<crate::WsSubscriptionState>,
|
||||
terminal_error_rx: tokio::sync::watch::Receiver<std::option::Option<ksp_core_lib::ErrorCode>>,
|
||||
}
|
||||
|
||||
impl WsSubscriptionRegistration {
|
||||
@@ -159,8 +178,9 @@ impl WsSubscriptionRegistration {
|
||||
id: crate::WsSubscriptionId,
|
||||
kind: crate::WsSubscriptionKind,
|
||||
state_rx: tokio::sync::watch::Receiver<crate::WsSubscriptionState>,
|
||||
terminal_error_rx: tokio::sync::watch::Receiver<std::option::Option<ksp_core_lib::ErrorCode>>,
|
||||
) -> Self {
|
||||
return Self { id, kind, state_rx };
|
||||
return Self { id, kind, state_rx, terminal_error_rx };
|
||||
}
|
||||
}
|
||||
|
||||
@@ -178,6 +198,8 @@ pub(crate) struct WsSubscriptionRuntime {
|
||||
pub(crate) remote_id: std::option::Option<u64>,
|
||||
/// Lifecycle publisher observed by the public typed handle.
|
||||
pub(crate) state_tx: tokio::sync::watch::Sender<crate::WsSubscriptionState>,
|
||||
/// Safe terminal failure code publisher observed by the public typed handle.
|
||||
pub(crate) terminal_error_tx: tokio::sync::watch::Sender<std::option::Option<ksp_core_lib::ErrorCode>>,
|
||||
/// Type-erased dispatcher into the bounded typed notification channel.
|
||||
pub(crate) dispatcher: WsNotificationDispatcher,
|
||||
}
|
||||
@@ -185,7 +207,7 @@ pub(crate) struct WsSubscriptionRuntime {
|
||||
impl WsSubscriptionRuntime {
|
||||
/// Builds the safe session-snapshot projection for this runtime entry.
|
||||
pub(crate) fn snapshot(&self) -> crate::WsSubscriptionSnapshot {
|
||||
return crate::WsSubscriptionSnapshot::new(self.id, self.kind, self.state, self.remote_id.is_some());
|
||||
return crate::WsSubscriptionSnapshot::new(self.id, self.kind, self.state, self.remote_id.is_some(), *self.terminal_error_tx.borrow());
|
||||
}
|
||||
|
||||
/// Updates the runtime state and publishes it to the typed handle.
|
||||
@@ -193,6 +215,12 @@ impl WsSubscriptionRuntime {
|
||||
self.state = state;
|
||||
self.state_tx.send_replace(state);
|
||||
}
|
||||
|
||||
/// Publishes a terminal failure code before moving the logical subscription to `Failed`.
|
||||
pub(crate) fn fail_with_code(&mut self, code: ksp_core_lib::ErrorCode) {
|
||||
self.terminal_error_tx.send_replace(std::option::Option::Some(code));
|
||||
self.set_state(crate::WsSubscriptionState::Failed);
|
||||
}
|
||||
}
|
||||
|
||||
fn subscription_closed_error(session_id: crate::WsSessionId, subscription_id: crate::WsSubscriptionId, message: &'static str) -> ksp_core_lib::Error {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/tests/public_api.rs
|
||||
// version: 28
|
||||
// version: 29
|
||||
|
||||
//! Integration tests for the public `ksp-onchain-transport-lib` consumer contract.
|
||||
|
||||
@@ -577,3 +577,11 @@ fn public_v0_2_7_pre_006_typed_websocket_subscription_handle_is_available_from_c
|
||||
let _recv = ksp_onchain_transport_lib::WsSubscription::<serde_json::Value>::recv;
|
||||
let _unsubscribe = ksp_onchain_transport_lib::WsSubscription::<serde_json::Value>::unsubscribe;
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn public_v0_2_7_pre_008_backpressure_observability_contract_is_available_from_crate_root() {
|
||||
let _terminal_error = ksp_onchain_transport_lib::WsSubscription::<serde_json::Value>::terminal_error_code;
|
||||
let _snapshot_terminal_error = ksp_onchain_transport_lib::WsSubscriptionSnapshot::terminal_error_code;
|
||||
let _overflow_count = ksp_onchain_transport_lib::WsSessionSnapshot::overflow_count;
|
||||
assert_eq!(ksp_onchain_transport_lib::ERROR_CODE_WS_BACKPRESSURE_OVERFLOW.code(), "ws_backpressure_overflow");
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/unit_tests/ws_lifecycle.rs
|
||||
// version: 2
|
||||
// version: 3
|
||||
|
||||
fn non_zero(value: u64) -> std::num::NonZeroU64 {
|
||||
return std::num::NonZeroU64::new(value).expect("test ID must be non-zero");
|
||||
@@ -45,6 +45,7 @@ fn websocket_snapshots_expose_safe_metadata_without_remote_ids_or_urls() {
|
||||
crate::WsSubscriptionKind::Slot,
|
||||
crate::WsSubscriptionState::Active,
|
||||
true,
|
||||
std::option::Option::None,
|
||||
);
|
||||
let snapshot = crate::WsSessionSnapshot::new(
|
||||
crate::WsSessionId::new(non_zero(3)),
|
||||
@@ -64,11 +65,27 @@ fn websocket_snapshots_expose_safe_metadata_without_remote_ids_or_urls() {
|
||||
assert_eq!(snapshot.continuity_gap_count(), 1);
|
||||
assert_eq!(snapshot.subscription_count(), 1);
|
||||
assert!(snapshot.subscriptions()[0].remote_bound());
|
||||
assert_eq!(snapshot.subscriptions()[0].terminal_error_code(), std::option::Option::None);
|
||||
assert_eq!(snapshot.overflow_count(), 0);
|
||||
let rendered = format!("{snapshot:?}");
|
||||
assert!(!rendered.contains("wss://"));
|
||||
assert!(!rendered.contains("remote_subscription_id"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn websocket_subscription_snapshot_preserves_only_safe_terminal_error_code() {
|
||||
let snapshot = crate::WsSubscriptionSnapshot::new(
|
||||
crate::WsSubscriptionId::new(non_zero(10)),
|
||||
crate::WsSubscriptionKind::Logs,
|
||||
crate::WsSubscriptionState::Failed,
|
||||
false,
|
||||
std::option::Option::Some(crate::ERROR_CODE_WS_BACKPRESSURE_OVERFLOW),
|
||||
);
|
||||
assert_eq!(snapshot.state(), crate::WsSubscriptionState::Failed);
|
||||
assert_eq!(snapshot.terminal_error_code(), std::option::Option::Some(crate::ERROR_CODE_WS_BACKPRESSURE_OVERFLOW));
|
||||
assert!(!format!("{snapshot:?}").contains("remote_subscription_id"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn websocket_subscription_kinds_map_exact_standard_method_triplets() {
|
||||
let cases = [
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/unit_tests/ws_session.rs
|
||||
// version: 6
|
||||
// version: 7
|
||||
|
||||
use futures_util::SinkExt; // rust-rules: trait-import
|
||||
use futures_util::StreamExt; // rust-rules: trait-import
|
||||
@@ -61,6 +61,27 @@ fn reconnect_session_settings(max_retries: u32, backoff: std::time::Duration, re
|
||||
);
|
||||
}
|
||||
|
||||
fn subscription_session_settings(
|
||||
notification_queue_capacity: usize,
|
||||
max_active_subscriptions: usize,
|
||||
reconnect: crate::WsReconnectSettings,
|
||||
) -> crate::WsSessionSettings {
|
||||
let defaults = crate::WsSessionSettings::default();
|
||||
return crate::WsSessionSettings::new(
|
||||
std::time::Duration::from_millis(250),
|
||||
std::time::Duration::from_millis(200),
|
||||
reconnect,
|
||||
crate::WsResubscribePolicy::ActiveSubscriptions,
|
||||
defaults.command_queue_capacity(),
|
||||
notification_queue_capacity,
|
||||
max_active_subscriptions,
|
||||
defaults.max_pending_requests(),
|
||||
defaults.max_message_size_bytes(),
|
||||
defaults.max_frame_size_bytes(),
|
||||
defaults.max_write_buffer_size_bytes(),
|
||||
);
|
||||
}
|
||||
|
||||
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");
|
||||
@@ -112,6 +133,39 @@ async fn wait_for_gap_count(session: &crate::WsSession, expected: u64) {
|
||||
}
|
||||
}
|
||||
|
||||
async fn wait_for_overflow_count(session: &crate::WsSession, expected: u64) {
|
||||
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(1);
|
||||
loop {
|
||||
if session.snapshot().overflow_count() == expected {
|
||||
return;
|
||||
}
|
||||
assert!(tokio::time::Instant::now() < deadline, "session did not reach overflow count {expected}");
|
||||
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
|
||||
}
|
||||
}
|
||||
|
||||
async fn wait_for_subscription_count(session: &crate::WsSession, expected: usize) {
|
||||
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(1);
|
||||
loop {
|
||||
if session.snapshot().subscription_count() == expected {
|
||||
return;
|
||||
}
|
||||
assert!(tokio::time::Instant::now() < deadline, "session did not reach subscription count {expected}");
|
||||
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
|
||||
}
|
||||
}
|
||||
|
||||
async fn wait_for_close_frame(websocket: &mut tokio_tungstenite::WebSocketStream<tokio::net::TcpStream>) {
|
||||
loop {
|
||||
let message = websocket.next().await;
|
||||
match message {
|
||||
std::option::Option::Some(std::result::Result::Ok(tokio_tungstenite::tungstenite::Message::Close(_))) => return,
|
||||
std::option::Option::Some(std::result::Result::Ok(_)) => {},
|
||||
std::option::Option::Some(std::result::Result::Err(_)) | std::option::Option::None => return,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn websocket_session_connects_and_round_trips_internal_json_rpc() {
|
||||
let (listener, url) = bind_local_listener().await;
|
||||
@@ -623,6 +677,7 @@ async fn websocket_resubscribe_never_keeps_session_but_fails_logical_subscriptio
|
||||
.expect("root subscribe must succeed");
|
||||
wait_for_gap_count(&session, 1).await;
|
||||
assert_eq!(subscription.state(), crate::WsSubscriptionState::Failed);
|
||||
assert_eq!(subscription.terminal_error_code(), std::option::Option::Some(crate::ERROR_CODE_WS_CONNECTION_FAILED));
|
||||
assert_eq!(session.snapshot().subscription_count(), 0);
|
||||
let result = session.execute_json_rpc("afterNeverPolicy", std::vec::Vec::new()).await.expect("physical session must reconnect without resubscribe");
|
||||
assert_eq!(result, serde_json::json!(true));
|
||||
@@ -719,6 +774,7 @@ async fn websocket_resubscribe_application_error_fails_only_one_subscription() {
|
||||
.expect("logs notification must decode");
|
||||
assert_eq!(notification, serde_json::json!({"restored": true}));
|
||||
assert_eq!(account.state(), crate::WsSubscriptionState::Failed);
|
||||
assert_eq!(account.terminal_error_code(), std::option::Option::Some(crate::ERROR_CODE_RPC_APPLICATION_ERROR));
|
||||
assert_eq!(logs.state(), crate::WsSubscriptionState::Active);
|
||||
assert_eq!(session.state(), crate::WsSessionState::Active);
|
||||
assert_eq!(session.snapshot().continuity_gap_count(), 1);
|
||||
@@ -867,6 +923,9 @@ async fn websocket_notification_method_mismatch_fails_only_the_logical_subscript
|
||||
let subscribe = read_request(&mut websocket).await;
|
||||
send_result(&mut websocket, &subscribe, serde_json::json!(5)).await;
|
||||
send_notification(&mut websocket, "rootNotification", 5, serde_json::json!(1)).await;
|
||||
let cleanup = read_request(&mut websocket).await;
|
||||
assert_eq!(cleanup.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("slotUnsubscribe"));
|
||||
assert_eq!(cleanup.get("params"), std::option::Option::Some(&serde_json::json!([5])));
|
||||
let request = read_request(&mut websocket).await;
|
||||
send_result(&mut websocket, &request, serde_json::json!(true)).await;
|
||||
});
|
||||
@@ -878,6 +937,7 @@ async fn websocket_notification_method_mismatch_fails_only_the_logical_subscript
|
||||
.await
|
||||
.expect("slot subscribe must succeed");
|
||||
wait_for_subscription_state(&subscription, crate::WsSubscriptionState::Failed).await;
|
||||
assert_eq!(subscription.terminal_error_code(), std::option::Option::Some(crate::ERROR_CODE_WS_PROTOCOL_ERROR));
|
||||
assert!(subscription.recv().await.is_none());
|
||||
assert_eq!(session.state(), crate::WsSessionState::Active);
|
||||
let result = session.execute_json_rpc("afterMismatch", std::vec::Vec::new()).await.expect("physical session must remain usable");
|
||||
@@ -894,6 +954,9 @@ async fn websocket_typed_notification_decode_failure_fails_only_one_subscription
|
||||
let subscribe = read_request(&mut websocket).await;
|
||||
send_result(&mut websocket, &subscribe, serde_json::json!(8)).await;
|
||||
send_notification(&mut websocket, "rootNotification", 8, serde_json::json!("not-a-slot")).await;
|
||||
let cleanup = read_request(&mut websocket).await;
|
||||
assert_eq!(cleanup.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("rootUnsubscribe"));
|
||||
assert_eq!(cleanup.get("params"), std::option::Option::Some(&serde_json::json!([8])));
|
||||
let request = read_request(&mut websocket).await;
|
||||
send_result(&mut websocket, &request, serde_json::json!(true)).await;
|
||||
});
|
||||
@@ -912,8 +975,162 @@ async fn websocket_typed_notification_decode_failure_fails_only_one_subscription
|
||||
let error = subscription.recv().await.expect("decode error must be delivered").expect_err("fixture payload must fail typed decoder");
|
||||
assert_eq!(error.code(), crate::ERROR_CODE_INVALID_RESPONSE);
|
||||
wait_for_subscription_state(&subscription, crate::WsSubscriptionState::Failed).await;
|
||||
assert_eq!(subscription.terminal_error_code(), std::option::Option::Some(crate::ERROR_CODE_INVALID_RESPONSE));
|
||||
assert_eq!(session.state(), crate::WsSessionState::Active);
|
||||
let result = session.execute_json_rpc("afterDecodeFailure", std::vec::Vec::new()).await.expect("physical 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_notification_queue_overflow_fails_only_slow_subscription_and_cleans_remote_binding() {
|
||||
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 slow_subscribe = read_request(&mut websocket).await;
|
||||
assert_eq!(slow_subscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("slotSubscribe"));
|
||||
send_result(&mut websocket, &slow_subscribe, serde_json::json!(41)).await;
|
||||
let healthy_subscribe = read_request(&mut websocket).await;
|
||||
assert_eq!(healthy_subscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("rootSubscribe"));
|
||||
send_result(&mut websocket, &healthy_subscribe, serde_json::json!(42)).await;
|
||||
send_notification(&mut websocket, "slotNotification", 41, serde_json::json!(1)).await;
|
||||
send_notification(&mut websocket, "slotNotification", 41, serde_json::json!(2)).await;
|
||||
let cleanup = read_request(&mut websocket).await;
|
||||
assert_eq!(cleanup.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("slotUnsubscribe"));
|
||||
assert_eq!(cleanup.get("params"), std::option::Option::Some(&serde_json::json!([41])));
|
||||
send_notification(&mut websocket, "rootNotification", 42, serde_json::json!(99)).await;
|
||||
wait_for_close_frame(&mut websocket).await;
|
||||
});
|
||||
let reconnect = crate::WsReconnectSettings::new(0, std::time::Duration::from_millis(10), std::time::Duration::from_millis(10));
|
||||
let settings = subscription_session_settings(1, 2, reconnect);
|
||||
let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed");
|
||||
let mut slow = session
|
||||
.subscribe_typed(crate::WsSubscriptionKind::Slot, std::vec::Vec::new(), |value| {
|
||||
return std::result::Result::Ok(value);
|
||||
})
|
||||
.await
|
||||
.expect("slow subscription must register");
|
||||
let mut healthy = session
|
||||
.subscribe_typed(crate::WsSubscriptionKind::Root, std::vec::Vec::new(), |value| {
|
||||
return std::result::Result::Ok(value);
|
||||
})
|
||||
.await
|
||||
.expect("healthy subscription must register");
|
||||
wait_for_subscription_state(&slow, crate::WsSubscriptionState::Failed).await;
|
||||
wait_for_overflow_count(&session, 1).await;
|
||||
assert_eq!(slow.terminal_error_code(), std::option::Option::Some(crate::ERROR_CODE_WS_BACKPRESSURE_OVERFLOW));
|
||||
assert_eq!(session.state(), crate::WsSessionState::Active);
|
||||
assert_eq!(session.snapshot().subscription_count(), 1);
|
||||
let first_slow = slow.recv().await.expect("first queued notification must remain observable").expect("first queued notification must decode");
|
||||
assert_eq!(first_slow, serde_json::json!(1));
|
||||
assert!(slow.recv().await.is_none());
|
||||
let healthy_value = healthy.recv().await.expect("healthy notification must remain available").expect("healthy notification must decode");
|
||||
assert_eq!(healthy_value, serde_json::json!(99));
|
||||
assert_eq!(healthy.state(), crate::WsSubscriptionState::Active);
|
||||
assert_eq!(healthy.terminal_error_code(), std::option::Option::None);
|
||||
session.close().await.expect("session close must remain bounded");
|
||||
server.await.expect("local server task must complete");
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn websocket_active_subscription_limit_rejects_excess_without_leaking_capacity() {
|
||||
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_subscribe = read_request(&mut websocket).await;
|
||||
assert_eq!(first_subscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("slotSubscribe"));
|
||||
send_result(&mut websocket, &first_subscribe, serde_json::json!(11)).await;
|
||||
let first_unsubscribe = read_request(&mut websocket).await;
|
||||
assert_eq!(first_unsubscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("slotUnsubscribe"));
|
||||
assert_eq!(first_unsubscribe.get("params"), std::option::Option::Some(&serde_json::json!([11])));
|
||||
send_result(&mut websocket, &first_unsubscribe, serde_json::json!(true)).await;
|
||||
let replacement_subscribe = read_request(&mut websocket).await;
|
||||
assert_eq!(replacement_subscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("rootSubscribe"));
|
||||
send_result(&mut websocket, &replacement_subscribe, serde_json::json!(12)).await;
|
||||
let replacement_unsubscribe = read_request(&mut websocket).await;
|
||||
assert_eq!(replacement_unsubscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("rootUnsubscribe"));
|
||||
assert_eq!(replacement_unsubscribe.get("params"), std::option::Option::Some(&serde_json::json!([12])));
|
||||
send_result(&mut websocket, &replacement_unsubscribe, serde_json::json!(true)).await;
|
||||
wait_for_close_frame(&mut websocket).await;
|
||||
});
|
||||
let reconnect = crate::WsReconnectSettings::new(0, std::time::Duration::from_millis(10), std::time::Duration::from_millis(10));
|
||||
let settings = subscription_session_settings(1, 1, reconnect);
|
||||
let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed");
|
||||
let mut first = session
|
||||
.subscribe_typed(crate::WsSubscriptionKind::Slot, std::vec::Vec::new(), |value| {
|
||||
return std::result::Result::Ok(value);
|
||||
})
|
||||
.await
|
||||
.expect("first subscription must register");
|
||||
let excess = session
|
||||
.subscribe_typed(crate::WsSubscriptionKind::Root, std::vec::Vec::new(), |value| {
|
||||
return std::result::Result::Ok(value);
|
||||
})
|
||||
.await
|
||||
.expect_err("subscription above configured active limit must be rejected");
|
||||
assert_eq!(excess.code(), crate::ERROR_CODE_WS_BACKPRESSURE_OVERFLOW);
|
||||
assert_eq!(session.snapshot().overflow_count(), 0);
|
||||
assert_eq!(session.snapshot().subscription_count(), 1);
|
||||
assert!(first.unsubscribe().await.expect("first unsubscribe must succeed"));
|
||||
assert_eq!(first.state(), crate::WsSubscriptionState::Closed);
|
||||
assert_eq!(first.terminal_error_code(), std::option::Option::None);
|
||||
wait_for_subscription_count(&session, 0).await;
|
||||
let mut replacement = session
|
||||
.subscribe_typed(crate::WsSubscriptionKind::Root, std::vec::Vec::new(), |value| {
|
||||
return std::result::Result::Ok(value);
|
||||
})
|
||||
.await
|
||||
.expect("capacity must become reusable after terminal cleanup");
|
||||
assert_eq!(replacement.id().get(), 2);
|
||||
assert!(replacement.unsubscribe().await.expect("replacement unsubscribe must succeed"));
|
||||
wait_for_subscription_count(&session, 0).await;
|
||||
session.close().await.expect("session close must remain bounded");
|
||||
server.await.expect("local server task must complete");
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn websocket_dropped_notification_receiver_triggers_remote_cleanup_and_releases_capacity() {
|
||||
let (listener, url) = bind_local_listener().await;
|
||||
let (trigger_tx, trigger_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 first_subscribe = read_request(&mut websocket).await;
|
||||
send_result(&mut websocket, &first_subscribe, serde_json::json!(71)).await;
|
||||
trigger_rx.await.expect("client must signal receiver drop");
|
||||
send_notification(&mut websocket, "slotNotification", 71, serde_json::json!(1)).await;
|
||||
let cleanup = read_request(&mut websocket).await;
|
||||
assert_eq!(cleanup.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("slotUnsubscribe"));
|
||||
assert_eq!(cleanup.get("params"), std::option::Option::Some(&serde_json::json!([71])));
|
||||
let replacement_subscribe = read_request(&mut websocket).await;
|
||||
assert_eq!(replacement_subscribe.get("method").and_then(serde_json::Value::as_str), std::option::Option::Some("rootSubscribe"));
|
||||
send_result(&mut websocket, &replacement_subscribe, serde_json::json!(72)).await;
|
||||
let replacement_unsubscribe = read_request(&mut websocket).await;
|
||||
send_result(&mut websocket, &replacement_unsubscribe, serde_json::json!(true)).await;
|
||||
wait_for_close_frame(&mut websocket).await;
|
||||
});
|
||||
let reconnect = crate::WsReconnectSettings::new(0, std::time::Duration::from_millis(10), std::time::Duration::from_millis(10));
|
||||
let settings = subscription_session_settings(1, 1, reconnect);
|
||||
let session = crate::WsSession::connect(local_endpoint_with_session(url.as_str(), settings)).await.expect("client handshake must succeed");
|
||||
let subscription = session
|
||||
.subscribe_typed(crate::WsSubscriptionKind::Slot, std::vec::Vec::new(), |value| {
|
||||
return std::result::Result::Ok(value);
|
||||
})
|
||||
.await
|
||||
.expect("first subscription must register");
|
||||
drop(subscription);
|
||||
trigger_tx.send(()).expect("receiver-drop trigger must send");
|
||||
wait_for_subscription_count(&session, 0).await;
|
||||
assert_eq!(session.snapshot().overflow_count(), 0);
|
||||
let mut replacement = session
|
||||
.subscribe_typed(crate::WsSubscriptionKind::Root, std::vec::Vec::new(), |value| {
|
||||
return std::result::Result::Ok(value);
|
||||
})
|
||||
.await
|
||||
.expect("capacity must be reusable after dropped receiver cleanup");
|
||||
assert!(replacement.unsubscribe().await.expect("replacement unsubscribe must succeed"));
|
||||
session.close().await.expect("session close must remain bounded");
|
||||
server.await.expect("local server task must complete");
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user