0.3.15-pre.015-fix.001
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
<!-- file: crates/ksp-onchain-transport-lib/README.md -->
|
||||
<!-- version: 25 -->
|
||||
<!-- version: 26 -->
|
||||
|
||||
# `ksp-onchain-transport-lib`
|
||||
|
||||
@@ -239,7 +239,7 @@ YellowstoneGrpcSubscribeSnapshot
|
||||
|
||||
Les sept unary standards retenus sont `SubscribeReplayInfo`, `Ping`, `GetLatestBlockhash`, `GetBlockHeight`, `GetSlot`, `IsBlockhashValid` et `GetVersion`. `Subscribe` couvre accounts, slots, transactions, transaction status, blocks, block metadata et entries, avec `commitment`, `accounts_data_slice`, `ping` et `from_slot`. `SubscribeDeshred` reste hors scope de la foundation standard.
|
||||
|
||||
Le stream bidirectionnel est borné : request/update queues, tailles inbound/outbound, half-close, close timeout et reconnect budget sont explicites. La queue d'updates Yellowstone reste bornée mais sa saturation locale n'est plus un fault : l'acteur attend asynchronement qu'une place soit rendue et cesse donc de poller le stream Tonic, ce qui propage la backpressure via gRPC/HTTP/2 sans drop silencieux ni queue non bornée. Le shutdown préempte cette attente. La queue de mutations de requête reste au contraire synchrone/fail-fast : un caller qui la sature reçoit toujours `grpc_backpressure_overflow`.
|
||||
Le stream bidirectionnel est borné : request/update queues, tailles inbound/outbound, half-close, close timeout et reconnect budget sont explicites. La queue d'updates Yellowstone reste bornée mais sa saturation locale n'est plus un fault : l'acteur attend asynchronement qu'une place soit rendue et cesse donc de poller le stream Tonic, ce qui propage la backpressure via gRPC/HTTP/2 sans drop silencieux ni queue non bornée. Le shutdown préempte cette attente. La queue de mutations de requête reste au contraire synchrone/fail-fast : un caller qui la sature reçoit toujours `grpc_backpressure_overflow`. Pendant un shutdown local explicite, un `Status` gRPC reçu après le half-close client est traité comme une terminaison `Closed` de cette fermeture déjà engagée ; le même `Status` reçu pendant une session active conserve son comportement normal de reconnect/fault et n'est pas masqué.
|
||||
|
||||
Après reconnect, KSP rejoue la dernière requête complète acceptée et avance prudemment `from_slot` selon le dernier slot observé et `SubscribeReplayInfo.first_available`. Les compteurs de gap et duplicate sont de l’observabilité ; ils ne constituent jamais une garantie exactly-once ou lossless.
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/src/grpc_stream.rs
|
||||
// version: 6
|
||||
// version: 7
|
||||
|
||||
use tonic_prost::prost::Message; // rust-rules: trait-import
|
||||
|
||||
@@ -1012,16 +1012,16 @@ async fn finish_client_half_close(
|
||||
},
|
||||
std::result::Result::Ok(std::result::Result::Err(status)) => {
|
||||
let code = status.code().to_string();
|
||||
ksp_logging_lib::warn!(
|
||||
ksp_logging_lib::debug!(
|
||||
target: crate::TRACING_TARGET,
|
||||
endpoint_name,
|
||||
provider = provider.as_str(),
|
||||
cluster = cluster.as_str(),
|
||||
grpc_code = code.as_str(),
|
||||
"Yellowstone subscribe endpoint returned a status during graceful shutdown"
|
||||
"Yellowstone subscribe endpoint returned a status after local graceful shutdown request; treating session as closed"
|
||||
);
|
||||
snapshot.state = crate::YellowstoneGrpcSubscribeState::Failed;
|
||||
snapshot.terminal_error_code = std::option::Option::Some(crate::ERROR_CODE_GRPC_STATUS);
|
||||
snapshot.state = crate::YellowstoneGrpcSubscribeState::Closed;
|
||||
snapshot.terminal_error_code = std::option::Option::None;
|
||||
snapshot_tx.send_replace(*snapshot);
|
||||
return;
|
||||
},
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/tests/release_completeness.rs
|
||||
// version: 48
|
||||
// version: 49
|
||||
|
||||
//! Release-level completeness canaries for staged HTTP and WebSocket Transport coverage.
|
||||
|
||||
@@ -1434,3 +1434,25 @@ enum ReconnectOutcome",
|
||||
assert!(!stream_source.contains("unbounded_channel"));
|
||||
assert!(stream_tests.contains("yellowstone_slow_receiver_applies_bounded_backpressure_without_terminal_overflow"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn release_v0_3_15_pre_015_fix_001_local_close_status_is_closed_without_weakening_active_status_failure() {
|
||||
let stream_source = include_str!("../src/grpc_stream.rs");
|
||||
let stream_tests = include_str!("../unit_tests/grpc_stream.rs");
|
||||
let close_start = stream_source.find("async fn finish_client_half_close(").expect("graceful half-close helper must remain present");
|
||||
let close_tail = &stream_source[close_start..];
|
||||
let close_end = close_tail.find("\nfn stream_status_error(").expect("graceful half-close helper boundary must remain present");
|
||||
let close_helper = &close_tail[..close_end];
|
||||
let actor_start = stream_source.find("async fn run_subscribe_actor(").expect("Subscribe actor must remain present");
|
||||
let actor_tail = &stream_source[actor_start..];
|
||||
let actor_end = actor_tail.find("\nenum ReconnectOutcome").expect("Subscribe actor boundary must remain present");
|
||||
let actor = &actor_tail[..actor_end];
|
||||
assert!(close_helper.contains("returned a status after local graceful shutdown request; treating session as closed"));
|
||||
assert!(close_helper.contains("snapshot.state = crate::YellowstoneGrpcSubscribeState::Closed"));
|
||||
assert!(close_helper.contains("snapshot.terminal_error_code = std::option::Option::None"));
|
||||
assert!(!close_helper.contains("snapshot.terminal_error_code = std::option::Option::Some(crate::ERROR_CODE_GRPC_STATUS)"));
|
||||
assert!(actor.contains("stream_status_error(\"Subscribe\", status"));
|
||||
assert!(actor.contains("fail_actor(crate::ERROR_CODE_GRPC_STATUS"));
|
||||
assert!(stream_tests.contains("yellowstone_explicit_close_accepts_remote_status_after_local_half_close"));
|
||||
assert!(stream_tests.contains("yellowstone_stream_remote_status_is_safe_and_terminal"));
|
||||
}
|
||||
|
||||
@@ -1,10 +1,11 @@
|
||||
// file: crates/ksp-onchain-transport-lib/unit_tests/grpc_stream.rs
|
||||
// version: 7
|
||||
// version: 8
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
enum FixtureMode {
|
||||
RoundTrip,
|
||||
ClientHalfClose,
|
||||
ClientHalfCloseStatus,
|
||||
HostileClose,
|
||||
Flood,
|
||||
RemoteStatus,
|
||||
@@ -86,6 +87,13 @@ impl yellowstone_grpc_proto::geyser::geyser_server::Geyser for FixtureGeyser {
|
||||
half_close_seen.store(true, std::sync::atomic::Ordering::SeqCst);
|
||||
}
|
||||
},
|
||||
FixtureMode::ClientHalfCloseStatus => {
|
||||
let half_close = inbound.message().await;
|
||||
if matches!(half_close, std::result::Result::Ok(std::option::Option::None)) {
|
||||
half_close_seen.store(true, std::sync::atomic::Ordering::SeqCst);
|
||||
}
|
||||
let _ = outbound_tx.send(std::result::Result::Err(tonic::Status::unknown("GRPC-CLOSE-STATUS-SECRET-CANARY"))).await;
|
||||
},
|
||||
FixtureMode::HostileClose => {
|
||||
let half_close = inbound.message().await;
|
||||
if matches!(half_close, std::result::Result::Ok(std::option::Option::None)) {
|
||||
@@ -465,6 +473,18 @@ async fn yellowstone_explicit_close_half_closes_request_stream_before_deadline()
|
||||
server.stop().await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn yellowstone_explicit_close_accepts_remote_status_after_local_half_close() {
|
||||
let server = FixtureServer::start(FixtureMode::ClientHalfCloseStatus).await;
|
||||
let defaults = crate::YellowstoneGrpcSessionSettings::default();
|
||||
let settings = fixture_settings(server.endpoint_url.as_str(), 8, 8, defaults.max_inbound_message_size_bytes(), defaults.max_outbound_message_size_bytes());
|
||||
let channel = crate::YellowstoneGrpcChannel::connect(&settings).await.expect("fixture channel must connect");
|
||||
let session = channel.open_standard_subscribe(initial_request()).await.expect("fixture Subscribe stream must open");
|
||||
session.close().await.expect("remote status after local half-close must remain a successful graceful close");
|
||||
assert!(server.half_close_seen.load(std::sync::atomic::Ordering::SeqCst));
|
||||
server.stop().await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn yellowstone_hostile_server_shutdown_is_bounded_by_close_timeout() {
|
||||
let server = FixtureServer::start(FixtureMode::HostileClose).await;
|
||||
|
||||
Reference in New Issue
Block a user