// file: crates/ksp-onchain-transport-lib/unit_tests/grpc_stream.rs // version: 3 #[derive(Clone, Copy)] enum FixtureMode { RoundTrip, ClientHalfClose, HostileClose, Flood, RemoteStatus, Malformed, Oversized, Idle, ReconnectReplay, ReplayGap, ReconnectExhausted, } #[derive(Clone)] struct FixtureGeyser { mode: FixtureMode, half_close_seen: std::sync::Arc, subscribe_calls: std::sync::Arc, } #[allow(clippy::implicit_return)] // tonic::async_trait generates async wrapper tails outside the authored fixture bodies. #[tonic::async_trait] impl yellowstone_grpc_proto::geyser::geyser_server::Geyser for FixtureGeyser { type SubscribeStream = super::MpscStream>; async fn subscribe( &self, request: tonic::Request>, ) -> std::result::Result, tonic::Status> { if let std::result::Result::Err(error) = verify_fixture_metadata(request.metadata()) { return std::result::Result::Err(error); } let subscribe_call = self.subscribe_calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst).saturating_add(1); if matches!(self.mode, FixtureMode::ReconnectExhausted) && subscribe_call > 1 { return std::result::Result::Err(tonic::Status::unavailable("GRPC-RECONNECT-SECRET-CANARY")); } let mut inbound = request.into_inner(); let (outbound_tx, outbound_rx) = tokio::sync::mpsc::channel(64); let mode = self.mode; let half_close_seen = self.half_close_seen.clone(); tokio::spawn(async move { let initial = inbound.message().await; let initial = match initial { std::result::Result::Ok(std::option::Option::Some(value)) => value, _ => return, }; match mode { FixtureMode::RoundTrip => { assert!(initial.slots.contains_key("initial")); if outbound_tx.send(std::result::Result::Ok(ping_update())).await.is_err() { return; } let ping = inbound.message().await; let ping = match ping { std::result::Result::Ok(std::option::Option::Some(value)) => value, _ => return, }; assert_eq!(ping.ping.map(|value| value.id), std::option::Option::Some(1)); assert!(ping.accounts.is_empty()); assert!(ping.slots.is_empty()); assert!(ping.transactions.is_empty()); assert!(ping.transactions_status.is_empty()); assert!(ping.blocks.is_empty()); assert!(ping.blocks_meta.is_empty()); assert!(ping.entry.is_empty()); if outbound_tx.send(std::result::Result::Ok(pong_update(1))).await.is_err() { return; } let mutation = inbound.message().await; let mutation = match mutation { std::result::Result::Ok(std::option::Option::Some(value)) => value, _ => return, }; assert!(mutation.accounts.contains_key("mutated")); let _ = outbound_tx.send(std::result::Result::Ok(slot_update(901))).await; }, FixtureMode::ClientHalfClose => { 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); } }, FixtureMode::HostileClose => { 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); } tokio::time::sleep(std::time::Duration::from_secs(5)).await; drop(outbound_tx); }, FixtureMode::Flood => { for slot in 0_u64..32_u64 { if outbound_tx.send(std::result::Result::Ok(slot_update(1_000 + slot))).await.is_err() { return; } } }, FixtureMode::RemoteStatus => { let _ = outbound_tx.send(std::result::Result::Err(tonic::Status::permission_denied("GRPC-STREAM-REMOTE-SECRET-CANARY"))).await; }, FixtureMode::Malformed => { let _ = outbound_tx .send(std::result::Result::Ok(yellowstone_grpc_proto::geyser::SubscribeUpdate { filters: std::vec::Vec::new(), update_oneof: std::option::Option::None, created_at: std::option::Option::None, })) .await; }, FixtureMode::Oversized => { let _ = outbound_tx.send(std::result::Result::Ok(slot_update(77))).await; }, FixtureMode::Idle => { let _ = inbound.message().await; }, FixtureMode::ReconnectReplay => { if subscribe_call == 1 { assert_eq!(initial.from_slot, std::option::Option::None); let _ = outbound_tx.send(std::result::Result::Ok(slot_update(500))).await; } else { assert_eq!(subscribe_call, 2); assert_eq!(initial.from_slot, std::option::Option::Some(500)); if outbound_tx.send(std::result::Result::Ok(slot_update(500))).await.is_err() { return; } if outbound_tx.send(std::result::Result::Ok(slot_update(501))).await.is_err() { return; } 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); } } }, FixtureMode::ReplayGap => { if subscribe_call == 1 { assert_eq!(initial.from_slot, std::option::Option::None); let _ = outbound_tx.send(std::result::Result::Ok(slot_update(500))).await; } else { assert_eq!(subscribe_call, 2); assert_eq!(initial.from_slot, std::option::Option::Some(505)); if outbound_tx.send(std::result::Result::Ok(slot_update(505))).await.is_err() { return; } 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); } } }, FixtureMode::ReconnectExhausted => { assert_eq!(subscribe_call, 1); let _ = outbound_tx.send(std::result::Result::Ok(slot_update(700))).await; }, } }); return std::result::Result::Ok(tonic::Response::new(super::MpscStream::new(outbound_rx))); } type SubscribeDeshredStream = futures_util::stream::Empty>; async fn subscribe_deshred( &self, _request: tonic::Request>, ) -> std::result::Result, tonic::Status> { return std::result::Result::Err(tonic::Status::unimplemented("deshred is outside KSP 0.2.9")); } async fn subscribe_replay_info( &self, request: tonic::Request, ) -> std::result::Result, tonic::Status> { if let std::result::Result::Err(error) = verify_fixture_metadata(request.metadata()) { return std::result::Result::Err(error); } let first_available = match self.mode { FixtureMode::ReconnectReplay | FixtureMode::ReconnectExhausted => std::option::Option::Some(400), FixtureMode::ReplayGap => std::option::Option::Some(505), _ => return std::result::Result::Err(tonic::Status::unimplemented("replay info is outside this fixture mode")), }; return std::result::Result::Ok(tonic::Response::new(yellowstone_grpc_proto::geyser::SubscribeReplayInfoResponse { first_available })); } async fn ping( &self, _request: tonic::Request, ) -> std::result::Result, tonic::Status> { return std::result::Result::Err(tonic::Status::unimplemented("unary is outside this fixture")); } async fn get_latest_blockhash( &self, _request: tonic::Request, ) -> std::result::Result, tonic::Status> { return std::result::Result::Err(tonic::Status::unimplemented("unary is outside this fixture")); } async fn get_block_height( &self, _request: tonic::Request, ) -> std::result::Result, tonic::Status> { return std::result::Result::Err(tonic::Status::unimplemented("unary is outside this fixture")); } async fn get_slot( &self, _request: tonic::Request, ) -> std::result::Result, tonic::Status> { return std::result::Result::Err(tonic::Status::unimplemented("unary is outside this fixture")); } async fn is_blockhash_valid( &self, _request: tonic::Request, ) -> std::result::Result, tonic::Status> { return std::result::Result::Err(tonic::Status::unimplemented("unary is outside this fixture")); } async fn get_version( &self, _request: tonic::Request, ) -> std::result::Result, tonic::Status> { return std::result::Result::Err(tonic::Status::unimplemented("unary is outside this fixture")); } } struct FixtureServer { endpoint_url: std::string::String, half_close_seen: std::sync::Arc, subscribe_calls: std::sync::Arc, shutdown: std::option::Option>, task: tokio::task::JoinHandle<()>, } impl FixtureServer { async fn start(mode: FixtureMode) -> Self { let bind_address: std::net::SocketAddr = "127.0.0.1:0".parse().expect("fixture bind address must parse"); let incoming = tonic::transport::server::TcpIncoming::bind(bind_address).expect("fixture gRPC listener must bind"); let local_address = incoming.local_addr().expect("fixture gRPC listener must expose local address"); let (shutdown, shutdown_receiver) = tokio::sync::oneshot::channel(); let half_close_seen = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)); let subscribe_calls = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)); let service_half_close_seen = half_close_seen.clone(); let service_subscribe_calls = subscribe_calls.clone(); let task = tokio::spawn(async move { let service = yellowstone_grpc_proto::geyser::geyser_server::GeyserServer::new(FixtureGeyser { mode, half_close_seen: service_half_close_seen, subscribe_calls: service_subscribe_calls, }); let result = tonic::transport::Server::builder() .serve_with_incoming_shutdown(service, incoming, async move { let _ = shutdown_receiver.await; }) .await; assert!(result.is_ok()); }); return Self { endpoint_url: format!("http://{local_address}"), half_close_seen, subscribe_calls, shutdown: std::option::Option::Some(shutdown), task, }; } async fn stop(mut self) { if let std::option::Option::Some(shutdown) = self.shutdown.take() { let _ = shutdown.send(()); } let result = self.task.await; assert!(result.is_ok()); } } fn verify_fixture_metadata(metadata: &tonic::metadata::MetadataMap) -> std::result::Result<(), tonic::Status> { let public = metadata.get("x-ksp-public").and_then(|value| return value.to_str().ok()); let secret = metadata.get("x-ksp-token").and_then(|value| return value.to_str().ok()); if public != std::option::Option::Some("fixture-public") || secret != std::option::Option::Some("GRPC-STREAM-SECRET-CANARY") { return std::result::Result::Err(tonic::Status::unauthenticated("fixture metadata mismatch")); } return std::result::Result::Ok(()); } fn fixture_settings( url: &str, request_capacity: usize, update_capacity: usize, max_inbound_message_size_bytes: usize, max_outbound_message_size_bytes: usize, ) -> crate::YellowstoneGrpcEndpointSettings { let defaults = crate::YellowstoneGrpcSessionSettings::default(); return fixture_settings_with_reconnect( url, request_capacity, update_capacity, max_inbound_message_size_bytes, max_outbound_message_size_bytes, crate::YellowstoneGrpcReconnectSettings::new(0, defaults.reconnect().initial_backoff(), defaults.reconnect().max_backoff()), ); } fn fixture_settings_with_reconnect( url: &str, request_capacity: usize, update_capacity: usize, max_inbound_message_size_bytes: usize, max_outbound_message_size_bytes: usize, reconnect: crate::YellowstoneGrpcReconnectSettings, ) -> crate::YellowstoneGrpcEndpointSettings { let defaults = crate::YellowstoneGrpcSessionSettings::default(); let session = crate::YellowstoneGrpcSessionSettings::new( defaults.connect_timeout(), defaults.unary_timeout(), std::time::Duration::from_millis(250), reconnect, request_capacity, update_capacity, max_inbound_message_size_bytes, max_outbound_message_size_bytes, ); let metadata = std::vec![ crate::YellowstoneGrpcMetadataEntry::public("x-ksp-public", "fixture-public").expect("fixture public metadata must be valid"), crate::YellowstoneGrpcMetadataEntry::secret("x-ksp-token", "GRPC-STREAM-SECRET-CANARY").expect("fixture secret metadata must be valid"), ]; return crate::YellowstoneGrpcEndpointSettings::new( "fixture-stream", true, crate::YellowstoneGrpcProviderName::new("fixture-provider"), crate::YellowstoneGrpcClusterName::new("devnet"), crate::YellowstoneGrpcEndpointUrl::parse(url).expect("fixture URL must parse"), session, ) .with_metadata(metadata) .expect("fixture metadata settings must validate"); } fn initial_request() -> crate::YellowstoneSubscribeRequest { let mut request = crate::YellowstoneSubscribeRequest::new(); request .insert_slot_filter( crate::YellowstoneSubscribeFilterName::new("initial").expect("fixture filter name must validate"), crate::YellowstoneSubscribeSlotFilter::new(), ) .expect("fixture slot filter must insert"); return request; } fn mutated_request() -> crate::YellowstoneSubscribeRequest { let mut request = crate::YellowstoneSubscribeRequest::new(); request .insert_account_filter( crate::YellowstoneSubscribeFilterName::new("mutated").expect("fixture filter name must validate"), crate::YellowstoneSubscribeAccountFilter::new(), ) .expect("fixture account filter must insert"); return request; } fn ping_update() -> yellowstone_grpc_proto::geyser::SubscribeUpdate { return yellowstone_grpc_proto::geyser::SubscribeUpdate { filters: std::vec::Vec::new(), update_oneof: std::option::Option::Some(yellowstone_grpc_proto::geyser::subscribe_update::UpdateOneof::Ping( yellowstone_grpc_proto::geyser::SubscribeUpdatePing {}, )), created_at: std::option::Option::None, }; } fn pong_update(id: i32) -> yellowstone_grpc_proto::geyser::SubscribeUpdate { return yellowstone_grpc_proto::geyser::SubscribeUpdate { filters: std::vec::Vec::new(), update_oneof: std::option::Option::Some(yellowstone_grpc_proto::geyser::subscribe_update::UpdateOneof::Pong( yellowstone_grpc_proto::geyser::SubscribeUpdatePong { id }, )), created_at: std::option::Option::None, }; } fn slot_update(slot: u64) -> yellowstone_grpc_proto::geyser::SubscribeUpdate { return yellowstone_grpc_proto::geyser::SubscribeUpdate { filters: std::vec!["initial".to_owned()], update_oneof: std::option::Option::Some(yellowstone_grpc_proto::geyser::subscribe_update::UpdateOneof::Slot( yellowstone_grpc_proto::geyser::SubscribeUpdateSlot { slot, parent: std::option::Option::Some(slot.saturating_sub(1)), status: yellowstone_grpc_proto::geyser::SlotStatus::SlotProcessed as i32, dead_error: std::option::Option::None, }, )), created_at: std::option::Option::None, }; } #[tokio::test(flavor = "current_thread")] async fn yellowstone_bidi_round_trip_mutates_request_replies_to_ping_and_observes_server_half_close() { let server = FixtureServer::start(FixtureMode::RoundTrip).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 mut session = channel.open_standard_subscribe(initial_request()).await.expect("fixture Subscribe stream must open"); assert_eq!(session.state(), crate::YellowstoneGrpcSubscribeState::Active); let ping = session.next_update().await.expect("Ping update must decode").expect("Ping update must be present"); assert!(matches!(ping, crate::YellowstoneSubscribeUpdate::Ping(_))); let pong = session.next_update().await.expect("Pong update must decode").expect("Pong update must be present"); match pong { crate::YellowstoneSubscribeUpdate::Pong(value) => assert_eq!(value.id(), 1), _ => panic!("fixture must return Pong"), } session.try_update(&mutated_request()).expect("request mutation must enter bounded queue"); let slot = session.next_update().await.expect("slot update must decode").expect("slot update must be present"); match slot { crate::YellowstoneSubscribeUpdate::Slot(value) => assert_eq!(value.slot(), 901), _ => panic!("fixture must return Slot"), } assert!(session.next_update().await.expect("server half-close must be normal").is_none()); assert_eq!(session.state(), crate::YellowstoneGrpcSubscribeState::Closed); session.close().await.expect("already half-closed fixture must close cleanly"); server.stop().await; } #[tokio::test(flavor = "current_thread")] async fn yellowstone_explicit_close_half_closes_request_stream_before_deadline() { let server = FixtureServer::start(FixtureMode::ClientHalfClose).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("graceful half-close must complete"); 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; 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"); let error = session.close().await.expect_err("hostile server must hit the bounded graceful close deadline"); assert_eq!(error.code(), crate::ERROR_CODE_TIMEOUT); assert!(server.half_close_seen.load(std::sync::atomic::Ordering::SeqCst)); server.stop().await; } #[tokio::test(flavor = "current_thread")] async fn yellowstone_session_drop_best_effort_half_closes_request_stream() { let server = FixtureServer::start(FixtureMode::ClientHalfClose).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"); drop(session); let observed = tokio::time::timeout(std::time::Duration::from_secs(1), async { loop { if server.half_close_seen.load(std::sync::atomic::Ordering::SeqCst) { return true; } tokio::task::yield_now().await; } }) .await .expect("fixture must observe dropped-session half-close before timeout"); assert!(observed); server.stop().await; } #[tokio::test(flavor = "current_thread")] async fn yellowstone_slow_receiver_overflow_is_terminal_and_observable() { let server = FixtureServer::start(FixtureMode::Flood).await; let defaults = crate::YellowstoneGrpcSessionSettings::default(); let settings = fixture_settings(server.endpoint_url.as_str(), 8, 2, 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 mut session = channel.open_standard_subscribe(initial_request()).await.expect("fixture Subscribe stream must open"); tokio::time::sleep(std::time::Duration::from_millis(20)).await; assert!(session.next_update().await.expect("first queued update must remain readable").is_some()); assert!(session.next_update().await.expect("second queued update must remain readable").is_some()); let error = session.next_update().await.expect_err("overflow must become terminal after bounded queued updates drain"); assert_eq!(error.code(), crate::ERROR_CODE_GRPC_BACKPRESSURE_OVERFLOW); let rendered = format!("{error:?} {session:?}"); assert!(!rendered.contains("GRPC-STREAM-SECRET-CANARY")); assert!(!rendered.contains(server.endpoint_url.as_str())); let _ = session.close().await; server.stop().await; } #[tokio::test(flavor = "current_thread")] async fn yellowstone_stream_remote_status_is_safe_and_terminal() { let server = FixtureServer::start(FixtureMode::RemoteStatus).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 mut session = channel.open_standard_subscribe(initial_request()).await.expect("fixture Subscribe stream must open"); let error = session.next_update().await.expect_err("remote Status must fail stream"); assert_eq!(error.code(), crate::ERROR_CODE_GRPC_STATUS); let rendered = format!("{error:?} {session:?}"); assert!(!rendered.contains("GRPC-STREAM-REMOTE-SECRET-CANARY")); assert!(!rendered.contains("GRPC-STREAM-SECRET-CANARY")); assert!(!rendered.contains(server.endpoint_url.as_str())); let _ = session.close().await; server.stop().await; } #[tokio::test(flavor = "current_thread")] async fn yellowstone_malformed_update_is_rejected_without_raw_wire_escape() { let server = FixtureServer::start(FixtureMode::Malformed).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 mut session = channel.open_standard_subscribe(initial_request()).await.expect("fixture Subscribe stream must open"); let error = session.next_update().await.expect_err("missing update oneof must be rejected"); assert_eq!(error.code(), crate::ERROR_CODE_INVALID_RESPONSE); let _ = session.close().await; server.stop().await; } #[tokio::test(flavor = "current_thread")] async fn yellowstone_inbound_message_limit_is_enforced_by_tonic_stream_decoder() { let server = FixtureServer::start(FixtureMode::Oversized).await; let defaults = crate::YellowstoneGrpcSessionSettings::default(); let settings = fixture_settings(server.endpoint_url.as_str(), 8, 8, 1, defaults.max_outbound_message_size_bytes()); let channel = crate::YellowstoneGrpcChannel::connect(&settings).await.expect("fixture channel must connect"); let mut session = channel.open_standard_subscribe(initial_request()).await.expect("fixture Subscribe stream must open"); let error = session.next_update().await.expect_err("message above configured inbound bound must fail decoding"); assert_eq!(error.code(), crate::ERROR_CODE_GRPC_STATUS); let rendered = format!("{error:?} {session:?}"); assert!(!rendered.contains(server.endpoint_url.as_str())); let _ = session.close().await; server.stop().await; } #[tokio::test(flavor = "current_thread")] async fn yellowstone_outbound_mutation_size_is_rejected_before_queue_dispatch() { let server = FixtureServer::start(FixtureMode::Idle).await; let settings = fixture_settings(server.endpoint_url.as_str(), 8, 8, 64 * 1024 * 1024, 1); let channel = crate::YellowstoneGrpcChannel::connect(&settings).await.expect("fixture channel must connect"); let session = channel .open_standard_subscribe(crate::YellowstoneSubscribeRequest::new()) .await .expect("empty initial request must fit one-byte fixture bound"); let error = session.try_update(&mutated_request()).expect_err("encoded mutation above configured bound must be rejected before dispatch"); assert_eq!(error.code(), crate::ERROR_CODE_GRPC_BACKPRESSURE_OVERFLOW); session.close().await.expect("idle fixture must observe client half-close"); server.stop().await; } #[tokio::test(flavor = "current_thread")] async fn yellowstone_reconnect_replays_from_last_observed_slot_and_counts_duplicate_identity() { let server = FixtureServer::start(FixtureMode::ReconnectReplay).await; let defaults = crate::YellowstoneGrpcSessionSettings::default(); let settings = fixture_settings_with_reconnect( server.endpoint_url.as_str(), 8, 8, defaults.max_inbound_message_size_bytes(), defaults.max_outbound_message_size_bytes(), crate::YellowstoneGrpcReconnectSettings::new(3, std::time::Duration::from_millis(5), std::time::Duration::from_millis(20)), ); let channel = crate::YellowstoneGrpcChannel::connect(&settings).await.expect("fixture channel must connect"); let mut session = channel.open_standard_subscribe(initial_request()).await.expect("fixture Subscribe stream must open"); let first = session.next_update().await.expect("first slot must decode").expect("first slot must be present"); match first { crate::YellowstoneSubscribeUpdate::Slot(value) => assert_eq!(value.slot(), 500), _ => panic!("fixture must return first Slot"), } let duplicate = session.next_update().await.expect("replayed duplicate must decode").expect("replayed duplicate must be present"); match duplicate { crate::YellowstoneSubscribeUpdate::Slot(value) => assert_eq!(value.slot(), 500), _ => panic!("fixture must return replayed Slot"), } let resumed = session.next_update().await.expect("resumed slot must decode").expect("resumed slot must be present"); match resumed { crate::YellowstoneSubscribeUpdate::Slot(value) => assert_eq!(value.slot(), 501), _ => panic!("fixture must return resumed Slot"), } let snapshot = session.snapshot(); assert_eq!(snapshot.state(), crate::YellowstoneGrpcSubscribeState::Active); assert_eq!(snapshot.reconnect_count(), 1); assert_eq!(snapshot.replay_attempt_count(), 1); assert_eq!(snapshot.continuity_gap_count(), 0); assert_eq!(snapshot.duplicate_update_count(), 1); assert_eq!(snapshot.last_requested_from_slot(), std::option::Option::Some(500)); assert_eq!(snapshot.last_observed_slot(), std::option::Option::Some(501)); session.close().await.expect("reconnected stream must close cleanly"); assert_eq!(server.subscribe_calls.load(std::sync::atomic::Ordering::SeqCst), 2); server.stop().await; } #[tokio::test(flavor = "current_thread")] async fn yellowstone_replay_info_proves_and_clamps_retention_gap_without_lossless_claim() { let server = FixtureServer::start(FixtureMode::ReplayGap).await; let defaults = crate::YellowstoneGrpcSessionSettings::default(); let settings = fixture_settings_with_reconnect( server.endpoint_url.as_str(), 8, 8, defaults.max_inbound_message_size_bytes(), defaults.max_outbound_message_size_bytes(), crate::YellowstoneGrpcReconnectSettings::new(3, std::time::Duration::from_millis(5), std::time::Duration::from_millis(20)), ); let channel = crate::YellowstoneGrpcChannel::connect(&settings).await.expect("fixture channel must connect"); let mut session = channel.open_standard_subscribe(initial_request()).await.expect("fixture Subscribe stream must open"); let _ = session.next_update().await.expect("first slot must decode").expect("first slot must be present"); let resumed = session.next_update().await.expect("clamped replay must decode").expect("clamped replay must be present"); match resumed { crate::YellowstoneSubscribeUpdate::Slot(value) => assert_eq!(value.slot(), 505), _ => panic!("fixture must return clamped Slot"), } let snapshot = session.snapshot(); assert_eq!(snapshot.reconnect_count(), 1); assert_eq!(snapshot.replay_attempt_count(), 1); assert_eq!(snapshot.continuity_gap_count(), 1); assert_eq!(snapshot.duplicate_update_count(), 0); assert_eq!(snapshot.last_requested_from_slot(), std::option::Option::Some(505)); assert_eq!(snapshot.last_observed_slot(), std::option::Option::Some(505)); session.close().await.expect("gap fixture must close cleanly"); server.stop().await; } #[tokio::test(flavor = "current_thread")] async fn yellowstone_reconnect_budget_exhaustion_is_terminal_and_safe() { let server = FixtureServer::start(FixtureMode::ReconnectExhausted).await; let defaults = crate::YellowstoneGrpcSessionSettings::default(); let settings = fixture_settings_with_reconnect( server.endpoint_url.as_str(), 8, 8, defaults.max_inbound_message_size_bytes(), defaults.max_outbound_message_size_bytes(), crate::YellowstoneGrpcReconnectSettings::new(2, std::time::Duration::from_millis(5), std::time::Duration::from_millis(10)), ); let channel = crate::YellowstoneGrpcChannel::connect(&settings).await.expect("fixture channel must connect"); let mut session = channel.open_standard_subscribe(initial_request()).await.expect("fixture Subscribe stream must open"); let _ = session.next_update().await.expect("first slot must decode").expect("first slot must be present"); let error = session.next_update().await.expect_err("failed reconnect attempts must exhaust the bounded budget"); assert_eq!(error.code(), crate::ERROR_CODE_GRPC_CHANNEL_FAILED); let snapshot = session.snapshot(); assert_eq!(snapshot.state(), crate::YellowstoneGrpcSubscribeState::Failed); assert_eq!(snapshot.reconnect_count(), 0); assert_eq!(snapshot.replay_attempt_count(), 2); assert_eq!(snapshot.terminal_error_code(), std::option::Option::Some(crate::ERROR_CODE_GRPC_CHANNEL_FAILED)); let rendered = format!("{error:?} {session:?}"); assert!(!rendered.contains("GRPC-RECONNECT-SECRET-CANARY")); assert!(!rendered.contains(server.endpoint_url.as_str())); assert_eq!(server.subscribe_calls.load(std::sync::atomic::Ordering::SeqCst), 3); let _ = session.close().await; server.stop().await; } #[tokio::test(flavor = "current_thread")] async fn yellowstone_shutdown_interrupts_reconnect_backoff_and_mutation_is_rejected_during_reconnect() { let server = FixtureServer::start(FixtureMode::ReconnectReplay).await; let defaults = crate::YellowstoneGrpcSessionSettings::default(); let settings = fixture_settings_with_reconnect( server.endpoint_url.as_str(), 8, 8, defaults.max_inbound_message_size_bytes(), defaults.max_outbound_message_size_bytes(), crate::YellowstoneGrpcReconnectSettings::new(3, std::time::Duration::from_millis(500), std::time::Duration::from_millis(500)), ); let channel = crate::YellowstoneGrpcChannel::connect(&settings).await.expect("fixture channel must connect"); let mut session = channel.open_standard_subscribe(initial_request()).await.expect("fixture Subscribe stream must open"); let _ = session.next_update().await.expect("first slot must decode").expect("first slot must be present"); tokio::time::timeout(std::time::Duration::from_secs(1), async { loop { if session.state() == crate::YellowstoneGrpcSubscribeState::Reconnecting { return; } tokio::task::yield_now().await; } }) .await .expect("fixture must enter reconnecting state"); let mutation_error = session.try_update(&mutated_request()).expect_err("request mutation must be rejected during reconnect"); assert_eq!(mutation_error.code(), crate::ERROR_CODE_GRPC_SESSION_CLOSED); session.close().await.expect("shutdown must interrupt reconnect backoff without opening another stream"); assert_eq!(server.subscribe_calls.load(std::sync::atomic::Ordering::SeqCst), 1); server.stop().await; }