diff --git a/crates/zakura-network/src/zakura/block_sync/service.rs b/crates/zakura-network/src/zakura/block_sync/service.rs index 3f7b655c8f..64d5ccb690 100644 --- a/crates/zakura-network/src/zakura/block_sync/service.rs +++ b/crates/zakura-network/src/zakura/block_sync/service.rs @@ -1,9 +1,9 @@ use super::{config::*, events::*, peer_registry::SessionAdmission, wire::*, *}; use crate::zakura::{ - handle_pipe_exit, spawn_supervised_pipe, FramedRecv, FramedSend, OrderedSendError, - OrderedSessionDemand, OrderedStreamOpening, OrderedStreamPolicy, Peer, PeerStreamSession, - Service, ServicePeerSnapshot, SinkReject, Stream, StreamMode, ZakuraBlockSyncCandidateState, - ZakuraConnId, ZakuraPeerId, FRAME_HEADER_BYTES, + handle_pipe_exit, spawn_supervised_pipe, FramedRecv, FramedSend, OrderedSendError, Peer, + PeerStreamSession, Service, ServicePeerSnapshot, SessionDemand, SessionOpening, SessionPolicy, + SinkReject, Stream, StreamMode, ZakuraBlockSyncCandidateState, ZakuraConnId, ZakuraPeerId, + FRAME_HEADER_BYTES, }; use std::{ sync::atomic::{AtomicU64, Ordering}, @@ -26,7 +26,7 @@ const BLOCK_SYNC_SERVICE_STREAMS: [Stream; 1] = [Stream { version: ZAKURA_BLOCK_SYNC_STREAM_VERSION, frame_cap: MAX_BS_FRAME_BYTES, capability: ZAKURA_CAP_BLOCK_SYNC, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }]; /// Service-declared streams for native block sync. @@ -451,28 +451,28 @@ impl Service for BlockSyncService { block_sync_streams() } - fn ordered_stream_policy(&self, _kind: u16) -> OrderedStreamPolicy { - OrderedStreamPolicy { - opening: OrderedStreamOpening::EitherSide, + fn session_policy(&self) -> SessionPolicy { + SessionPolicy { + opening: SessionOpening::EitherSide, reopen: true, } } - fn ordered_session_demand( + fn session_demand( &self, conn_id: ZakuraConnId, peer: &ZakuraPeerId, _negotiated: u64, direction: ServicePeerDirection, - ) -> OrderedSessionDemand { + ) -> SessionDemand { if let Some(deadline) = self.peer_park_deadline(peer) { - return OrderedSessionDemand::RetryAt(deadline); + return SessionDemand::RetryAt(deadline); } let mut peer_snapshot = self.inner.peer_snapshot.clone(); peer_snapshot.borrow_and_update(); if !self.peer_slots_free(direction) { - return OrderedSessionDemand::WaitForChange(Box::pin(async move { + return SessionDemand::WaitForChange(Box::pin(async move { if peer_snapshot.changed().await.is_err() { std::future::pending::<()>().await; } @@ -495,7 +495,7 @@ impl Service for BlockSyncService { .is_empty() { let mut service_demand = self.service_demand.clone(); - return OrderedSessionDemand::WaitForChange(Box::pin(async move { + return SessionDemand::WaitForChange(Box::pin(async move { if let Some(demand) = service_demand.as_mut() { tokio::select! { changed = candidates.changed() => { @@ -516,7 +516,7 @@ impl Service for BlockSyncService { } } - OrderedSessionDemand::OpenNow + SessionDemand::OpenNow } fn wants_peer( @@ -760,8 +760,8 @@ impl Service for BlockSyncService { return false; }; matches!( - self.ordered_session_demand(conn_id, peer, ZAKURA_CAP_BLOCK_SYNC, direction), - OrderedSessionDemand::OpenNow + self.session_demand(conn_id, peer, ZAKURA_CAP_BLOCK_SYNC, direction), + SessionDemand::OpenNow ) } diff --git a/crates/zakura-network/src/zakura/block_sync/tests.rs b/crates/zakura-network/src/zakura/block_sync/tests.rs index 462da5ca72..0c9fa96346 100644 --- a/crates/zakura-network/src/zakura/block_sync/tests.rs +++ b/crates/zakura-network/src/zakura/block_sync/tests.rs @@ -32,8 +32,8 @@ use crate::zakura::{ framed_channel, testkit::{await_until, TraceCapture, TraceValue}, trace::BlockBodySource, - FramedRecv, FramedSend, OrderedSessionDemand, Peer, Service, ServicePeerSnapshot, - ServiceRegistry, StreamMode, ZakuraBlockSyncCandidateState, + FramedRecv, FramedSend, Peer, Service, ServicePeerSnapshot, ServiceRegistry, SessionDemand, + StreamMode, ZakuraBlockSyncCandidateState, }; use zakura_chain::{ fmt::HexDebug, @@ -6232,7 +6232,7 @@ fn block_sync_stream_declares_kind_capability_version_and_frame_cap() { assert_eq!(stream.kind, ZAKURA_STREAM_BLOCK_SYNC); assert_eq!(stream.version, ZAKURA_BLOCK_SYNC_STREAM_VERSION); assert_eq!(stream.capability, ZAKURA_CAP_BLOCK_SYNC); - assert_eq!(stream.mode, StreamMode::Ordered); + assert_eq!(stream.mode, StreamMode::Persistent); assert_eq!(stream.frame_cap, MAX_BS_FRAME_BYTES); } @@ -6255,14 +6255,14 @@ async fn service_registry_routes_block_sync_by_exact_capability_and_version() { .is_none()); assert_eq!( registry - .ordered_streams_for_negotiated(ZAKURA_CAP_BLOCK_SYNC) + .persistent_streams_for_negotiated(ZAKURA_CAP_BLOCK_SYNC) .iter() .map(|stream| stream.kind) .collect::>(), vec![ZAKURA_STREAM_BLOCK_SYNC] ); - assert!(registry.ordered_streams_for_negotiated(0).is_empty()); - assert!(registry.wants_ordered_stream( + assert!(registry.persistent_streams_for_negotiated(0).is_empty()); + assert!(registry.wants_session( ZAKURA_STREAM_BLOCK_SYNC, ZAKURA_CAP_BLOCK_SYNC, &peer, @@ -15311,24 +15311,24 @@ async fn parked_connection_cleanup_allows_a_fresh_connection_after_cooldown() { handle.park_session_for_test(&peer, old_conn_id, Duration::ZERO); assert!(matches!( - service.ordered_session_demand( + service.session_demand( old_conn_id, &peer, ZAKURA_CAP_BLOCK_SYNC, ServicePeerDirection::Outbound, ), - OrderedSessionDemand::WaitForChange(_), + SessionDemand::WaitForChange(_), )); service.remove_peer(&peer, old_conn_id); assert!(matches!( - service.ordered_session_demand( + service.session_demand( new_conn_id, &peer, ZAKURA_CAP_BLOCK_SYNC, ServicePeerDirection::Outbound, ), - OrderedSessionDemand::OpenNow + SessionDemand::OpenNow )); reactor_task.abort(); } @@ -15353,13 +15353,13 @@ async fn same_connection_block_sync_session_waits_at_tip_then_reopens_for_new_wo let conn_id = 17; handle.park_session_for_test(&peer, conn_id, Duration::ZERO); - let demand = service.ordered_session_demand( + let demand = service.session_demand( conn_id, &peer, ZAKURA_CAP_BLOCK_SYNC, ServicePeerDirection::Outbound, ); - let OrderedSessionDemand::WaitForChange(changed) = demand else { + let SessionDemand::WaitForChange(changed) = demand else { panic!("a locally parked session must stay absent while block sync is at tip"); }; @@ -15375,13 +15375,13 @@ async fn same_connection_block_sync_session_waits_at_tip_then_reopens_for_new_wo .expect("new block work wakes the parked session demand"); assert!(matches!( - service.ordered_session_demand( + service.session_demand( conn_id, &peer, ZAKURA_CAP_BLOCK_SYNC, ServicePeerDirection::Outbound, ), - OrderedSessionDemand::OpenNow, + SessionDemand::OpenNow, )); reactor_task.abort(); } @@ -15420,7 +15420,7 @@ async fn serving_only_coordinator_demand_keeps_block_session_available_during_fa let conn_id = 18; handle.park_session_for_test(&peer, conn_id, Duration::ZERO); - let OrderedSessionDemand::WaitForChange(changed) = service.ordered_session_demand( + let SessionDemand::WaitForChange(changed) = service.session_demand( conn_id, &peer, ZAKURA_CAP_BLOCK_SYNC, @@ -15440,13 +15440,13 @@ async fn serving_only_coordinator_demand_keeps_block_session_available_during_fa .await .expect("fallback service demand wakes the parked ordered session"); assert!(matches!( - service.ordered_session_demand( + service.session_demand( conn_id, &peer, ZAKURA_CAP_BLOCK_SYNC, ServicePeerDirection::Outbound, ), - OrderedSessionDemand::OpenNow, + SessionDemand::OpenNow, )); reactor_task.abort(); } diff --git a/crates/zakura-network/src/zakura/discovery/service.rs b/crates/zakura-network/src/zakura/discovery/service.rs index 5c4c3cf1e8..0625c66771 100644 --- a/crates/zakura-network/src/zakura/discovery/service.rs +++ b/crates/zakura-network/src/zakura/discovery/service.rs @@ -24,9 +24,9 @@ use tokio_util::sync::CancellationToken; use crate::zakura::{ handle_pipe_exit, spawn_supervised_peer_task, spawn_supervised_pipe, BlockSyncHandle, CloseCause, Event, Flow, Frame, FramedRecv, FramedSend, HeaderSyncHandle, OrderedSendError, - OrderedSessionDemand, OrderedStreamOpening, OrderedStreamPolicy, Peer, PeerStreamSession, Pipe, - Service, ServiceAdmissionDecision, ServicePeerDirection, SinkReject, Stream, StreamMode, - ZakuraConnId, ZakuraPeerId, LOCAL_MAX_CONTROL_FRAME_BYTES, ZAKURA_CAP_DISCOVERY, + Peer, PeerStreamSession, Pipe, Service, ServiceAdmissionDecision, ServicePeerDirection, + SessionDemand, SessionOpening, SessionPolicy, SinkReject, Stream, StreamMode, ZakuraConnId, + ZakuraPeerId, LOCAL_MAX_CONTROL_FRAME_BYTES, ZAKURA_CAP_DISCOVERY, }; #[cfg(test)] @@ -47,7 +47,7 @@ const DISCOVERY_SERVICE_STREAMS: [Stream; 1] = [Stream { version: ZAKURA_DISCOVERY_STREAM_VERSION, frame_cap: LOCAL_MAX_CONTROL_FRAME_BYTES, capability: ZAKURA_CAP_DISCOVERY, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }]; /// Service-declared streams for native discovery. @@ -305,20 +305,20 @@ impl Service for DiscoveryService { discovery_streams() } - fn ordered_stream_policy(&self, _kind: u16) -> OrderedStreamPolicy { - OrderedStreamPolicy { - opening: OrderedStreamOpening::InitiatorOnly, + fn session_policy(&self) -> SessionPolicy { + SessionPolicy { + opening: SessionOpening::InitiatorOnly, reopen: true, } } - fn ordered_session_demand( + fn session_demand( &self, conn_id: ZakuraConnId, peer: &ZakuraPeerId, _negotiated: u64, direction: ServicePeerDirection, - ) -> OrderedSessionDemand { + ) -> SessionDemand { if self .session_states .lock() @@ -326,7 +326,7 @@ impl Service for DiscoveryService { .get(&(peer.clone(), conn_id)) .is_some_and(|record| record.state == DiscoverySessionState::Retired) { - return OrderedSessionDemand::Retire; + return SessionDemand::Retire; } let mut peers = self.handle.subscribe_peer_snapshot(); @@ -336,14 +336,14 @@ impl Service for DiscoveryService { ServicePeerDirection::Outbound => snapshot.outbound_slots_free, }; if slots_free == 0 { - return OrderedSessionDemand::WaitForChange(Box::pin(async move { + return SessionDemand::WaitForChange(Box::pin(async move { if peers.changed().await.is_err() { std::future::pending::<()>().await; } })); } - OrderedSessionDemand::OpenNow + SessionDemand::OpenNow } fn wants_peer( @@ -1970,13 +1970,13 @@ mod tests { ); assert_eq!(record.state, DiscoverySessionState::Retired); assert!(matches!( - service.ordered_session_demand( + service.session_demand( 0, &peer_id, ZAKURA_CAP_DISCOVERY, ServicePeerDirection::Inbound, ), - OrderedSessionDemand::Retire, + SessionDemand::Retire, )); Ok(()) } @@ -2313,13 +2313,13 @@ mod tests { // A finished discovery-only exchange retires the session so the handler // cannot reopen the stream while the connection tears down. assert!(matches!( - service.ordered_session_demand( + service.session_demand( 0, &peer_id, ZAKURA_CAP_DISCOVERY | ZAKURA_CAP_BLOCK_SYNC, ServicePeerDirection::Inbound, ), - OrderedSessionDemand::Retire, + SessionDemand::Retire, )); wait_for_discovery_inbound_peers(&handle, 0).await; @@ -2557,13 +2557,13 @@ mod tests { // The refresh loop still owns this session. // The handler must keep the stream eligible. assert!(!matches!( - service.ordered_session_demand( + service.session_demand( 0, &peer_id, ZAKURA_CAP_DISCOVERY | ZAKURA_CAP_HEADER_SYNC, ServicePeerDirection::Inbound, ), - OrderedSessionDemand::Retire, + SessionDemand::Retire, )); connection_cancel.cancel(); diff --git a/crates/zakura-network/src/zakura/handler.rs b/crates/zakura-network/src/zakura/handler.rs index b08ecfa181..e3c2b1cc04 100644 --- a/crates/zakura-network/src/zakura/handler.rs +++ b/crates/zakura-network/src/zakura/handler.rs @@ -1,8 +1,8 @@ //! Zakura P2P v2 endpoint, protocol handler, and bounded connection serving. -mod ordered_pair; +mod service_session; mod trace; -use ordered_pair::{spawn_ordered_pair, PendingOrderedPairs, PreparedOrderedStream, SetupIo}; +use service_session::{spawn_service_session, PendingSessions, PreparedStream, SetupIo}; use std::{ collections::{HashMap, HashSet}, @@ -46,7 +46,7 @@ use zakura_chain::{ use self::trace::ZakuraConnTrace; use super::discovery::{self, native_dial_supervised, spawn_native_bootstrap_dialer, RedialPolicy}; use super::trace::{reject_reason_label, ZakuraTrace}; -use super::transport::{worker_framed_channel, FramedWorkerRecv, QueuedFrame}; +use super::transport::{worker_framed_channel, FramedWorkerRecv, SessionLayout}; #[cfg(any(test, feature = "zakura-testkit"))] use crate::zakura::drive_header_sync_actions; #[cfg(any(test, feature = "zakura-testkit"))] @@ -59,11 +59,11 @@ use crate::{ AuthenticatedPeerRegistration, BlockSyncAction, BlockSyncFrontiers, BlockSyncHandle, BlockSyncService, BlockSyncStartup, BoxRunFuture, Clock, CloseCause, Frame, FramedRecv, FramedSend, FullStateFrontiers, HeaderSyncPassthroughService, HeaderSyncService, - HeaderSyncStartup, OrderedSessionDemand, OrderedStreamOpening, OrderedStreamPair, - OrderedStreamPolicy, Peer, RealClock, Service, ServicePeerDirection, ServiceRegistry, - ServiceStream, SinkReject, Stream, StreamMode, StreamPrelude, ZakuraAcceptedLimits, - ZakuraBlockSyncConfig, ZakuraConnId, ZakuraControlAck, ZakuraControlHello, - ZakuraControlRole, ZakuraControlValidation, ZakuraHandshakeConfig, ZakuraHandshakePath, + HeaderSyncStartup, Peer, RealClock, Service, ServicePeerDirection, ServiceRegistry, + ServiceStream, SessionDemand, SessionOpening, SessionPolicy, SinkReject, Stream, + StreamMode, StreamPrelude, StreamWritePolicy, ZakuraAcceptedLimits, ZakuraBlockSyncConfig, + ZakuraConnId, ZakuraControlAck, ZakuraControlHello, ZakuraControlRole, + ZakuraControlValidation, ZakuraHandshakeConfig, ZakuraHandshakePath, ZakuraHeaderSyncConfig, ZakuraInitialLimits, ZakuraLimits, ZakuraPeerId, ZakuraPeerSupervisor, ZakuraProtocolError, ZakuraRejectReason, ZakuraServiceId, ZakuraUpgradeDialStart, CONTROL_ACK_MAGIC, CONTROL_HELLO_MAGIC, CONTROL_VERSION, @@ -170,9 +170,6 @@ const STREAM_WORKER_DRAIN_TIMEOUT: Duration = Duration::from_secs(1); const ORDERED_STREAM_REOPEN_BACKOFF: Duration = Duration::from_millis(250); const ORDERED_STREAM_REOPEN_BACKOFF_CAP: Duration = Duration::from_secs(8); const OUTBOUND_STREAM_WRITE_TIMEOUT: Duration = Duration::from_secs(10); -// A paused sibling can make shared-credit updates take longer than ten seconds -// on a slow link even while complete blocks keep arriving. -const PAIRED_DATA_WRITE_TIMEOUT: Duration = Duration::from_secs(32); const OUTBOUND_REQUEST_RESPONSE_TIMEOUT: Duration = Duration::from_secs(30); // Mirrors the legacy gossip compatibility protocol. Compile-time assertions // below keep this transport-side budget validator pinned to the codec constants. @@ -1548,7 +1545,7 @@ struct IncomingStreamSetup { stream_id: u64, stream: Stream, prelude: StreamPrelude, - pair: Option<(OrderedStreamPair, u64)>, + session: Option<(SessionLayout, u64)>, } struct StreamAdmission<'a> { @@ -1607,7 +1604,8 @@ struct StreamWorkerContext { message_payload_limits: &'static [(u16, usize)], message_types: Option<&'static [u16]>, queue_depths: Option<(usize, usize)>, - session_resources: Option>, + write_policy: StreamWritePolicy, + session_resources: Option>, outbound_frame_cap: u32, message_bucket: SharedMessageBucket, connection_token: CancellationToken, @@ -1616,14 +1614,11 @@ struct StreamWorkerContext { freshness_tx: watch::Sender, } -struct AdmittedOrderedSession { +struct AdmittedSession { kind: u16, - version: u16, session_id: u64, - recv: FramedRecv, - send: FramedSend, cancel_token: CancellationToken, - companion: Option, + streams: Vec, } struct ServiceStreamRole { @@ -1633,62 +1628,53 @@ struct ServiceStreamRole { send: FramedSend, } -impl AdmittedOrderedSession { +impl AdmittedSession { fn into_service_streams(self) -> HashMap { - let mut streams = HashMap::new(); - if let Some(role) = self.companion { - streams.insert( - role.kind, - ServiceStream::new( - self.session_id, - role.version, - role.recv, - role.send, - self.cancel_token.clone(), - ), - ); - } - streams.insert( - self.kind, - ServiceStream::new( - self.session_id, - self.version, - self.recv, - self.send, - self.cancel_token, - ), - ); - streams + self.streams + .into_iter() + .map(|role| { + ( + role.kind, + ServiceStream::new( + self.session_id, + role.version, + role.recv, + role.send, + self.cancel_token.clone(), + ), + ) + }) + .collect() } } #[derive(Copy, Clone, Debug)] -struct OrderedSessionExit { +struct SessionExit { stream: Stream, session_id: u64, opened_locally: bool, } #[derive(Copy, Clone, Debug, Eq, PartialEq)] -enum OrderedSessionWaitReason { +enum SessionWaitReason { Demand, Transport, } #[derive(Copy, Clone, Debug, Default, Eq, PartialEq)] -enum OrderedSessionReopenState { +enum SessionReopenState { #[default] Idle, - Waiting(OrderedSessionWaitReason), + Waiting(SessionWaitReason), Retired, } -type OrderedSessionWait = futures::stream::Once>; -type OrderedSessionWaits = StreamMap; +type SessionWait = futures::stream::Once>; +type SessionWaits = StreamMap; /// Connection-owned lifecycle for one negotiated ordered stream kind. #[derive(Debug)] -struct OrderedSessionState { +struct ServiceSessionState { stream: Stream, local_session_id: Option, remote_session_id: Option, @@ -1696,17 +1682,17 @@ struct OrderedSessionState { // service can immediately park the stream again, and resetting here would // turn persistent refusal into a base-delay reopen loop. reopen_attempts: u32, - reopen_state: OrderedSessionReopenState, + reopen_state: SessionReopenState, } -impl OrderedSessionState { +impl ServiceSessionState { fn new(stream: Stream) -> Self { Self { stream, local_session_id: None, remote_session_id: None, reopen_attempts: 0, - reopen_state: OrderedSessionReopenState::Idle, + reopen_state: SessionReopenState::Idle, } } @@ -1728,17 +1714,17 @@ impl OrderedSessionState { true } - fn schedule_demand(&mut self, waits: &mut OrderedSessionWaits, demand: OrderedSessionDemand) { + fn schedule_demand(&mut self, waits: &mut SessionWaits, demand: SessionDemand) { match demand { - OrderedSessionDemand::OpenNow => { + SessionDemand::OpenNow => { self.cancel_wait(waits); self.install_wait( waits, - OrderedSessionWaitReason::Demand, + SessionWaitReason::Demand, Box::pin(future::ready(())), ); } - OrderedSessionDemand::RetryAt(at) => { + SessionDemand::RetryAt(at) => { self.install_demand_wait( waits, Box::pin(async move { @@ -1746,18 +1732,18 @@ impl OrderedSessionState { }), ); } - OrderedSessionDemand::WaitForChange(changed) => { + SessionDemand::WaitForChange(changed) => { self.install_demand_wait(waits, changed); } - OrderedSessionDemand::Retire => { + SessionDemand::Retire => { waits.remove(&self.stream.kind); - self.reopen_state = OrderedSessionReopenState::Retired; + self.reopen_state = SessionReopenState::Retired; } } } - fn schedule_transport_backoff(&mut self, waits: &mut OrderedSessionWaits) -> Option { - if self.reopen_state != OrderedSessionReopenState::Idle { + fn schedule_transport_backoff(&mut self, waits: &mut SessionWaits) -> Option { + if self.reopen_state != SessionReopenState::Idle { return None; } debug_assert!( @@ -1769,48 +1755,43 @@ impl OrderedSessionState { self.reopen_attempts = self.reopen_attempts.saturating_add(1); self.install_wait( waits, - OrderedSessionWaitReason::Transport, + SessionWaitReason::Transport, Box::pin(tokio::time::sleep(delay)), ); Some(delay) } - fn finish_wait(&mut self, waits: &mut OrderedSessionWaits) { + fn finish_wait(&mut self, waits: &mut SessionWaits) { let removed = waits.remove(&self.stream.kind); debug_assert!( removed.is_some(), "yielded ordered session wait remains keyed until it is handled" ); - if matches!(self.reopen_state, OrderedSessionReopenState::Waiting(_)) { - self.reopen_state = OrderedSessionReopenState::Idle; + if matches!(self.reopen_state, SessionReopenState::Waiting(_)) { + self.reopen_state = SessionReopenState::Idle; } } - fn cancel_wait(&mut self, waits: &mut OrderedSessionWaits) { + fn cancel_wait(&mut self, waits: &mut SessionWaits) { waits.remove(&self.stream.kind); - if matches!(self.reopen_state, OrderedSessionReopenState::Waiting(_)) { - self.reopen_state = OrderedSessionReopenState::Idle; + if matches!(self.reopen_state, SessionReopenState::Waiting(_)) { + self.reopen_state = SessionReopenState::Idle; } } - fn install_demand_wait( - &mut self, - waits: &mut OrderedSessionWaits, - wait: BoxRunFuture<'static, ()>, - ) { - if self.reopen_state == OrderedSessionReopenState::Waiting(OrderedSessionWaitReason::Demand) - { + fn install_demand_wait(&mut self, waits: &mut SessionWaits, wait: BoxRunFuture<'static, ()>) { + if self.reopen_state == SessionReopenState::Waiting(SessionWaitReason::Demand) { return; } self.cancel_wait(waits); - self.install_wait(waits, OrderedSessionWaitReason::Demand, wait); + self.install_wait(waits, SessionWaitReason::Demand, wait); } fn install_wait( &mut self, - waits: &mut OrderedSessionWaits, - reason: OrderedSessionWaitReason, + waits: &mut SessionWaits, + reason: SessionWaitReason, wait: BoxRunFuture<'static, ()>, ) { let replaced = waits.insert(self.stream.kind, futures::stream::once(wait)); @@ -1818,7 +1799,7 @@ impl OrderedSessionState { replaced.is_none(), "ordered session wait is cancelled before replacement" ); - self.reopen_state = OrderedSessionReopenState::Waiting(reason); + self.reopen_state = SessionReopenState::Waiting(reason); } } @@ -1827,8 +1808,8 @@ impl OrderedSessionState { /// The initiator opens ordinary ordered streams. Block sync is symmetric, so /// either side may open it and simultaneous offers use the connection's /// deterministic collision tiebreak. -fn may_open_ordered_stream(policy: OrderedStreamPolicy, is_initiator: bool) -> bool { - is_initiator || policy.opening == OrderedStreamOpening::EitherSide +fn may_open_ordered_stream(policy: SessionPolicy, is_initiator: bool) -> bool { + is_initiator || policy.opening == SessionOpening::EitherSide } /// Whether this endpoint proactively opens an ordered stream. @@ -1838,13 +1819,13 @@ fn may_open_ordered_stream(policy: OrderedStreamPolicy, is_initiator: bool) -> b /// accept either-side streams, preserving compatibility with peers that race /// an offer during upgrade. fn opens_ordered_stream_locally( - policy: OrderedStreamPolicy, + policy: SessionPolicy, is_initiator: bool, i_open_collision_winner: bool, ) -> bool { match policy.opening { - OrderedStreamOpening::InitiatorOnly => is_initiator, - OrderedStreamOpening::EitherSide => i_open_collision_winner, + SessionOpening::InitiatorOnly => is_initiator, + SessionOpening::EitherSide => i_open_collision_winner, } } @@ -1852,13 +1833,13 @@ fn opens_ordered_stream_locally( /// sibling services remain healthy. The endpoint designated by the transport /// policy keeps offering a replacement with bounded backoff. fn should_reopen_ordered_session( - exited: OrderedSessionExit, - policy: OrderedStreamPolicy, + exited: SessionExit, + policy: SessionPolicy, is_initiator: bool, i_open_collision_winner: bool, connection_cancelled: bool, ) -> bool { - exited.stream.mode == StreamMode::Ordered + exited.stream.mode == StreamMode::Persistent && policy.reopen && opens_ordered_stream_locally(policy, is_initiator, i_open_collision_winner) && !connection_cancelled @@ -2379,20 +2360,25 @@ impl ZakuraProtocolHandler { let accepted_capabilities = context.accepted_capabilities; let stream_sem = Arc::new(Semaphore::new(usize::from(limits.max_open_streams))); let mut workers = JoinSet::new(); - let (ordered_session_exit_tx, mut ordered_session_exit_rx) = mpsc::unbounded_channel(); - let mut ordered_session_waits = OrderedSessionWaits::new(); - let mut pending_pairs = PendingOrderedPairs::default(); + let (session_exit_tx, mut session_exit_rx) = mpsc::unbounded_channel(); + let mut session_waits = SessionWaits::new(); + let mut pending_sessions = PendingSessions::default(); let mut incoming_setup: Option>> = None; let mut open_limiter = TokenBucket::new(limits.stream_open_rate_per_second); let mut message_buckets = MessageRateBuckets::new(); let (freshness_tx, freshness_rx) = watch::channel(Instant::now()); let negotiated_ordered_streams = self .registry - .ordered_streams_for_negotiated(accepted_capabilities); - let mut ordered_sessions: HashMap = negotiated_ordered_streams + .persistent_streams_for_negotiated(accepted_capabilities); + let mut service_sessions: HashMap = negotiated_ordered_streams .iter() .copied() - .map(|stream| (stream.kind, OrderedSessionState::new(stream))) + .filter(|stream| { + self.registry + .session_layout(*stream) + .is_some_and(|layout| layout.primary() == *stream) + }) + .map(|stream| (stream.kind, ServiceSessionState::new(stream))) .collect(); // The dialer opens ordinary ordered streams. For symmetric block sync, // the mirror-stable node-id tiebreak designates one proactive opener, @@ -2404,12 +2390,12 @@ impl ZakuraProtocolHandler { for stream in negotiated_ordered_streams.iter().copied() { if self .registry - .ordered_stream_pair(stream) - .is_some_and(|pair| stream != pair.data) + .session_layout(stream) + .is_some_and(|layout| stream != layout.primary()) { continue; } - let policy = self.registry.ordered_stream_policy(stream.kind); + let policy = self.registry.session_policy(stream.kind); if !opens_ordered_stream_locally( policy, context.is_initiator, @@ -2418,14 +2404,14 @@ impl ZakuraProtocolHandler { continue; } - match self.registry.ordered_session_demand( + match self.registry.session_demand( stream.kind, conn_id, accepted_capabilities, &peer_id, context.direction, ) { - OrderedSessionDemand::OpenNow => ordered_streams.push(stream), + SessionDemand::OpenNow => ordered_streams.push(stream), demand => deferred_ordered_streams.push((stream, demand)), } } @@ -2433,10 +2419,14 @@ impl ZakuraProtocolHandler { .registry .request_response_streams_for_negotiated(accepted_capabilities) .len(); - // Deferred services do not open streams here; each active pair needs both roles. + // Count every required stream in the sessions that local demand will open. let opening_stream_count: usize = ordered_streams .iter() - .map(|stream| self.registry.ordered_stream_pair(*stream).map_or(1, |_| 2)) + .map(|stream| { + self.registry + .session_layout(*stream) + .map_or(1, |layout| layout.streams.len()) + }) .sum(); if opening_stream_count > usize::from(limits.max_open_streams) { debug!( @@ -2485,7 +2475,7 @@ impl ZakuraProtocolHandler { let mut opened_capabilities = 0; for stream in ordered_streams { let admitted = match self - .open_ordered_service_stream( + .open_service_session( &connection, stream, &mut workers, @@ -2499,18 +2489,18 @@ impl ZakuraProtocolHandler { conn.clone(), peer_id.clone(), context.direction, - ordered_session_exit_tx.clone(), + session_exit_tx.clone(), ) .await { Ok(admitted) => admitted, - Err(ZakuraHandlerError::OrderedSessionFull) => { + Err(ZakuraHandlerError::SessionFull) => { // Demand is advisory: another connection can reserve the // last slot before this open. Retry only this service. - ordered_sessions + service_sessions .get_mut(&stream.kind) .expect("selected stream has negotiated session state") - .schedule_transport_backoff(&mut ordered_session_waits); + .schedule_transport_backoff(&mut session_waits); continue; } Err(error) => { @@ -2529,7 +2519,7 @@ impl ZakuraProtocolHandler { } }; opened_capabilities |= stream.capability; - ordered_sessions + service_sessions .get_mut(&admitted.kind) .expect("opened ordered stream was selected from negotiated session state") .local_session_id = Some(admitted.session_id); @@ -2563,39 +2553,39 @@ impl ZakuraProtocolHandler { ?demand, "deferring ordered service session according to reactor demand" ); - ordered_sessions + service_sessions .get_mut(&stream.kind) .expect("deferred ordered stream has negotiated session state") - .schedule_demand(&mut ordered_session_waits, demand); + .schedule_demand(&mut session_waits, demand); } } loop { - let pair_deadline = pending_pairs.deadline(); + let session_deadline = pending_sessions.deadline(); tokio::select! { biased; _ = connection_token.cancelled() => break, _ = async { - match pair_deadline { + match session_deadline { Some(deadline) => tokio::time::sleep_until(deadline).await, None => future::pending().await, } - } => pending_pairs.expire(Instant::now()), + } => pending_sessions.expire(Instant::now()), _ = freshness_reaper(freshness_rx.clone(), limits.idle_timeout), if run_freshness_reaper => { connection.close(VarInt::from_u32(ZAKURA_CLOSE_NEUTRAL), b"idle"); close_cause.record("idle_timeout"); break; } - Some(exited) = ordered_session_exit_rx.recv() => { + Some(exited) = session_exit_rx.recv() => { // Stop tracking the dead generation, or a legitimate reopen of the // same kind would look like a duplicate stream and kill the whole // connection, taking block sync and gossip down with header sync. - let Some(session) = ordered_sessions.get_mut(&exited.stream.kind) else { + let Some(session) = service_sessions.get_mut(&exited.stream.kind) else { continue; }; let removed_active_session = session.remove_active_session(exited.opened_locally, exited.session_id); - let policy = self.registry.ordered_stream_policy(exited.stream.kind); + let policy = self.registry.session_policy(exited.stream.kind); if removed_active_session && !session.has_active_session() && should_reopen_ordered_session( @@ -2607,7 +2597,7 @@ impl ZakuraProtocolHandler { ) { if let Some(delay) = - session.schedule_transport_backoff(&mut ordered_session_waits) + session.schedule_transport_backoff(&mut session_waits) { debug!( stream_kind = exited.stream.kind, @@ -2619,13 +2609,13 @@ impl ZakuraProtocolHandler { } } } - Some((kind, ())) = ordered_session_waits.next(), if !ordered_session_waits.is_empty() => { - let Some(session) = ordered_sessions.get_mut(&kind) else { + Some((kind, ())) = session_waits.next(), if !session_waits.is_empty() => { + let Some(session) = service_sessions.get_mut(&kind) else { continue; }; - session.finish_wait(&mut ordered_session_waits); + session.finish_wait(&mut session_waits); let stream = session.stream; - let policy = self.registry.ordered_stream_policy(stream.kind); + let policy = self.registry.session_policy(stream.kind); if connection_token.is_cancelled() || !opens_ordered_stream_locally( policy, @@ -2637,7 +2627,7 @@ impl ZakuraProtocolHandler { continue; } - let demand = self.registry.ordered_session_demand( + let demand = self.registry.session_demand( stream.kind, conn_id, accepted_capabilities, @@ -2645,15 +2635,15 @@ impl ZakuraProtocolHandler { context.direction, ); match demand { - OrderedSessionDemand::OpenNow => {} + SessionDemand::OpenNow => {} demand => { - session.schedule_demand(&mut ordered_session_waits, demand); + session.schedule_demand(&mut session_waits, demand); continue; } } match self - .open_ordered_service_stream( + .open_service_session( &connection, stream, &mut workers, @@ -2667,15 +2657,15 @@ impl ZakuraProtocolHandler { conn.clone(), peer_id.clone(), context.direction, - ordered_session_exit_tx.clone(), + session_exit_tx.clone(), ) .await { Ok(admitted) => { - let session = ordered_sessions + let session = service_sessions .get_mut(&admitted.kind) .expect("opened ordered stream has negotiated session state"); - session.cancel_wait(&mut ordered_session_waits); + session.cancel_wait(&mut session_waits); session.local_session_id = Some(admitted.session_id); let service_streams = admitted.into_service_streams(); let admitted_capabilities = self.registry.add_escalated_peer( @@ -2693,10 +2683,10 @@ impl ZakuraProtocolHandler { cleanup_guard.add_admitted_capabilities(admitted_capabilities); } Err(error) => { - let delay = ordered_sessions + let delay = service_sessions .get_mut(&stream.kind) .expect("failed ordered stream open has negotiated session state") - .schedule_transport_backoff(&mut ordered_session_waits); + .schedule_transport_backoff(&mut session_waits); debug!( ?error, stream_kind = stream.kind, @@ -2733,14 +2723,14 @@ impl ZakuraProtocolHandler { }; if let Some(admitted) = self.finish_bi_stream_setup( setup, &mut admission, per_stream_queue_depth, - ordered_session_exit_tx.clone(), &mut pending_pairs, + session_exit_tx.clone(), &mut pending_sessions, ) { let kind = admitted.kind; // A stream kind we never negotiated is a protocol // fault: tear the connection down (unchanged // strictness). - if !ordered_sessions.contains_key(&kind) { + if !service_sessions.contains_key(&kind) { debug!( stream_kind = kind, "closing peer after unexpected ordered stream" @@ -2750,7 +2740,7 @@ impl ZakuraProtocolHandler { continue; } let is_collision = - ordered_sessions[&kind].local_session_id.is_some(); + service_sessions[&kind].local_session_id.is_some(); // The deterministic winner keeps its own stream // and parks the peer's. The loser falls through // and adopts the peer's stream. @@ -2765,23 +2755,12 @@ impl ZakuraProtocolHandler { // The old remote generation remains tracked until its // workers have finished. - if ordered_sessions[&kind].remote_session_id.is_some() { - if admitted.companion.is_some() { - // QUIC can deliver a replacement before both old - // role workers finish. Park this offer locally; - // the opener retries through its normal backoff. - admitted.cancel_token.cancel(); - continue; - } - debug!( - stream_kind = kind, - "closing peer after duplicate ordered stream" - ); - close_cause.record("duplicate_stream"); - connection_token.cancel(); + if service_sessions[&kind].remote_session_id.is_some() { + // A replacement can arrive before the retiring workers finish. + admitted.cancel_token.cancel(); continue; } - ordered_sessions + service_sessions .get_mut(&kind) .expect("accepted ordered stream has negotiated session state") .remote_session_id = Some(admitted.session_id); @@ -2790,28 +2769,14 @@ impl ZakuraProtocolHandler { // demand. For a collision we lost, demand is // implied because we already opened the same kind. let demand = (!is_collision).then(|| { - if admitted.companion.is_some() { - self.registry.reserved_ordered_session_demand( - kind, - conn_id, - accepted_capabilities, - &peer_id, - context.direction, - ) - } else { - self.registry.ordered_session_demand( - kind, - conn_id, - accepted_capabilities, - &peer_id, - context.direction, - ) - } + self.registry.reserved_session_demand( + kind, conn_id, accepted_capabilities, &peer_id, context.direction, + ) }); if demand .as_ref() .is_some_and(|demand| { - !matches!(demand, OrderedSessionDemand::OpenNow) + !matches!(demand, SessionDemand::OpenNow) }) { metrics::counter!( @@ -2827,14 +2792,14 @@ impl ZakuraProtocolHandler { // We are not adopting it after all, so this kind // is free again -- both for a later re-offer by // the peer and for the demand re-check below. - let session = ordered_sessions + let session = service_sessions .get_mut(&kind) .expect("parked ordered stream has negotiated session state"); session.remove_active_session(false, admitted.session_id); // If this side may open the stream, re-check its // local demand. Otherwise the entitled remote // opener observes the parked stream and retries. - let policy = self.registry.ordered_stream_policy(kind); + let policy = self.registry.session_policy(kind); if policy.reopen && opens_ordered_stream_locally( policy, @@ -2843,7 +2808,7 @@ impl ZakuraProtocolHandler { ) { session.schedule_demand( - &mut ordered_session_waits, + &mut session_waits, demand.expect("non-open demand exists because this branch checked it"), ); } @@ -2851,10 +2816,10 @@ impl ZakuraProtocolHandler { continue; } - ordered_sessions + service_sessions .get_mut(&kind) .expect("adopted ordered stream has negotiated session state") - .cancel_wait(&mut ordered_session_waits); + .cancel_wait(&mut session_waits); let service_streams = admitted.into_service_streams(); // We keep the full accepted capability context so // discovery can make cross-service ownership @@ -2970,7 +2935,7 @@ impl ZakuraProtocolHandler { connection_token.cancel(); drop(incoming_setup); - drop(pending_pairs); + drop(pending_sessions); while let Some(joined) = timeout(STREAM_WORKER_DRAIN_TIMEOUT, workers.join_next()) .await .ok() @@ -2993,7 +2958,7 @@ impl ZakuraProtocolHandler { } #[allow(clippy::too_many_arguments)] - async fn open_ordered_service_stream( + async fn open_service_session( &self, connection: &Connection, stream: Stream, @@ -3008,71 +2973,49 @@ impl ZakuraProtocolHandler { conn: ZakuraConnTrace, peer_id: ZakuraPeerId, direction: ServicePeerDirection, - ordered_session_exit_tx: mpsc::UnboundedSender, - ) -> Result { - let pair = self.registry.ordered_stream_pair(stream); - let resources = if pair.is_some() { - self.registry - .service_for_kind(stream.kind) - .expect("a selected stream has an owning service") - .reserve_ordered_session(direction) - .map_err(|_| ZakuraHandlerError::OrderedSessionFull)? - } else { - None - }; - let pair_id = pair.map(|_| random_stream_session_seed()); - let mut primary = self - .prepare_ordered_stream( - connection, - stream, - pair_id, - stream_sem, - message_buckets, - limits, - connection_token.clone(), - close_cause.clone(), - freshness_tx.clone(), - conn.clone(), - peer_id.clone(), - ) - .await?; - primary.set_session_resources(resources.clone()); - if let Some(pair) = pair { - if pair.data != stream { - return Err(ZakuraHandlerError::InvalidOrderedPair); - } - let mut requests = self + session_exit_tx: mpsc::UnboundedSender, + ) -> Result { + let layout = self + .registry + .session_layout(stream) + .expect("only negotiated persistent streams open service sessions"); + if layout.primary() != stream { + return Err(ZakuraHandlerError::InvalidServiceSession); + } + let resources = self + .registry + .service_for_kind(stream.kind) + .expect("a selected stream has an owning service") + .reserve_session(direction) + .map_err(|_| ZakuraHandlerError::SessionFull)?; + let wire_id = layout.is_multi_stream().then(random_stream_session_seed); + let mut prepared = Vec::with_capacity(layout.streams.len()); + for member in layout.streams.iter().copied() { + let mut member = self .prepare_ordered_stream( connection, - pair.requests, - pair_id, + member, + wire_id, stream_sem, message_buckets, limits, - connection_token, - close_cause, - freshness_tx, - conn, - peer_id, + connection_token.clone(), + close_cause.clone(), + freshness_tx.clone(), + conn.clone(), + peer_id.clone(), ) .await?; - requests.set_session_resources(resources); - Ok(spawn_ordered_pair( - workers, - primary, - requests, - per_stream_queue_depth, - true, - ordered_session_exit_tx, - )) - } else { - Ok(primary.spawn_single( - workers, - per_stream_queue_depth, - true, - ordered_session_exit_tx, - )) + member.set_session_resources(resources.clone()); + prepared.push(member); } + Ok(spawn_service_session( + workers, + prepared, + per_stream_queue_depth, + true, + session_exit_tx, + )) } /// Start at most one peer-paced setup read per connection. The returned @@ -3168,9 +3111,9 @@ impl ZakuraProtocolHandler { return None; } - if stream.mode == StreamMode::Ordered + if stream.mode == StreamMode::Persistent && registry - .ordered_streams_for_negotiated(accepted_capabilities) + .persistent_streams_for_negotiated(accepted_capabilities) .iter() .all(|selected| *selected != stream) { @@ -3213,11 +3156,8 @@ impl ZakuraProtocolHandler { return None; } - if stream.mode == StreamMode::Ordered - && !may_open_ordered_stream( - registry.ordered_stream_policy(stream.kind), - !is_initiator, - ) + if stream.mode == StreamMode::Persistent + && !may_open_ordered_stream(registry.session_policy(stream.kind), !is_initiator) { debug!( stream_kind = stream.kind, @@ -3237,23 +3177,19 @@ impl ZakuraProtocolHandler { .increment(1); conn.trace_stream("accepted", stream_id, Some(stream_kind)); - let pair = registry.ordered_stream_pair(stream); - let pair = if let Some(pair) = pair { - if prelude.request_id.is_some() { - let _ = send.reset(VarInt::from_u32(ZAKURA_CLOSE_BAD_PRELUDE)); - let _ = recv.stop(VarInt::from_u32(ZAKURA_CLOSE_BAD_PRELUDE)); - return None; - } + let session = if let Some(layout) = registry.session_layout(stream) { let mut bytes = [0u8; 8]; - if !matches!( - timeout(limits.prelude_timeout, recv.read_exact(&mut bytes)).await, - Ok(Ok(())) - ) { + if layout.is_multi_stream() + && !matches!( + timeout(limits.prelude_timeout, recv.read_exact(&mut bytes)).await, + Ok(Ok(())) + ) + { let _ = send.reset(VarInt::from_u32(ZAKURA_CLOSE_BAD_PRELUDE)); let _ = recv.stop(VarInt::from_u32(ZAKURA_CLOSE_BAD_PRELUDE)); return None; } - Some((pair, u64::from_le_bytes(bytes))) + Some((layout, u64::from_le_bytes(bytes))) } else { None }; @@ -3263,7 +3199,7 @@ impl ZakuraProtocolHandler { stream_id, stream, prelude, - pair, + session, }) })) } @@ -3275,20 +3211,20 @@ impl ZakuraProtocolHandler { incoming: IncomingStreamSetup, admission: &mut StreamAdmission<'_>, per_stream_queue_depth: usize, - ordered_session_exit_tx: mpsc::UnboundedSender, - pending_pairs: &mut PendingOrderedPairs, - ) -> Option { + session_exit_tx: mpsc::UnboundedSender, + pending_sessions: &mut PendingSessions, + ) -> Option { let IncomingStreamSetup { io, permit, stream_id, stream, prelude, - pair, + session, } = incoming; let (mut send, mut recv) = io.take(); - let resources = if let Some((pair, _)) = pair { - match pending_pairs.reserve_or_share(pair, &self.registry, admission.direction) { + let resources = if let Some((layout, _)) = &session { + match pending_sessions.reserve_or_share(layout, &self.registry, admission.direction) { Ok(resources) => resources, Err(_) => { let _ = send.reset(VarInt::from_u32(ZAKURA_CLOSE_RESOURCE)); @@ -3301,7 +3237,9 @@ impl ZakuraProtocolHandler { }; let message_bucket = message_bucket_for( admission.message_buckets, - pair.map_or(prelude.stream_kind, |(pair, _)| pair.data.kind), + session + .as_ref() + .map_or(prelude.stream_kind, |(layout, _)| layout.primary().kind), admission.limits.message_rate_per_second, RealClock, ); @@ -3317,6 +3255,7 @@ impl ZakuraProtocolHandler { message_payload_limits: self.registry.message_payload_limits(stream), message_types: self.registry.message_types(stream), queue_depths: self.registry.stream_queue_depths(stream), + write_policy: self.registry.stream_write_policy(stream), session_resources: resources, outbound_frame_cap: peer_accepted_frame_cap( &admission.limits, @@ -3330,48 +3269,33 @@ impl ZakuraProtocolHandler { freshness_tx: admission.freshness_tx.clone(), }; - if let Some((pair, pair_id)) = pair { - let prepared = PreparedOrderedStream::new(send, recv, stream, prelude, context); - return match pending_pairs.insert(pair, pair_id, prepared) { - Ok(Some((data, requests))) => Some(spawn_ordered_pair( + if let Some((layout, wire_id)) = session { + let prepared = PreparedStream::new(send, recv, stream, prelude, context); + return match pending_sessions.insert(&layout, wire_id, prepared) { + Ok(Some(streams)) => Some(spawn_service_session( admission.workers, - data, - requests, + streams, per_stream_queue_depth, false, - ordered_session_exit_tx, + session_exit_tx, )), Ok(None) => None, Err(error) => { - debug!(?error, "rejecting mismatched ordered stream pair"); - admission.close_cause.record("invalid_ordered_pair"); + debug!(?error, "rejecting mismatched service session"); + admission.close_cause.record("invalid_service_session"); admission.connection_token.cancel(); None } }; } - if stream.mode == StreamMode::RequestResponse { - admission.workers.spawn(request_stream_worker( - send, - recv, - prelude, - context, - self.registry.clone(), - )); - None - } else { - Some(spawn_persistent_stream_worker( - admission.workers, - send, - recv, - stream, - prelude, - context, - per_stream_queue_depth, - false, - ordered_session_exit_tx, - )) - } + admission.workers.spawn(request_stream_worker( + send, + recv, + prelude, + context, + self.registry.clone(), + )); + None } async fn register_and_serve( @@ -3385,7 +3309,7 @@ impl ZakuraProtocolHandler { // Count each symmetric session regardless of which endpoint opens it. let ordered_stream_count = self .registry - .ordered_streams_for_escalation( + .persistent_streams_for_escalation( context.accepted_capabilities, &peer_id, context.direction, @@ -4147,6 +4071,7 @@ async fn run_native_initiator_handshake( } #[allow(clippy::too_many_arguments)] +#[cfg(test)] fn spawn_persistent_stream_worker( workers: &mut JoinSet<()>, send: SendStream, @@ -4156,42 +4081,15 @@ fn spawn_persistent_stream_worker( context: StreamWorkerContext, queue_depth: usize, opened_locally: bool, - ordered_session_exit_tx: mpsc::UnboundedSender, -) -> AdmittedOrderedSession { - let (inbound_depth, outbound_depth) = - bounded_stream_queue_depths(queue_depth, context.queue_depths); - let (to_service_tx, to_service_rx) = mpsc::channel(inbound_depth); - let (from_service_tx, from_service_rx) = worker_framed_channel(outbound_depth); - let admitted = AdmittedOrderedSession { - kind: prelude.stream_kind, - version: prelude.stream_version, - session_id: context.stream_id, - recv: FramedRecv::new(to_service_rx), - send: from_service_tx, - cancel_token: context.stream_token.clone(), - companion: None, - }; - - let exit = OrderedSessionExit { - stream, - session_id: admitted.session_id, + session_exit_tx: mpsc::UnboundedSender, +) -> AdmittedSession { + spawn_service_session( + workers, + vec![PreparedStream::new(send, recv, stream, prelude, context)], + queue_depth, opened_locally, - }; - workers.spawn(async move { - persistent_stream_worker( - send, - recv, - prelude, - context, - to_service_tx, - from_service_rx, - inbound_depth, - ) - .await; - let _ = ordered_session_exit_tx.send(exit); - }); - - admitted + session_exit_tx, + ) } fn bounded_stream_queue_depths( @@ -4206,17 +4104,11 @@ fn bounded_stream_queue_depths( }) } -#[derive(Copy, Clone, Eq, PartialEq)] -enum OrderedWritePolicy { - Standalone, - PairData, - PairRequests, -} - #[derive(Debug, Error)] #[error("Zakura outbound frame write timed out")] struct OrderedFrameWriteTimeout; +#[cfg(test)] async fn persistent_stream_worker( send: SendStream, recv: RecvStream, @@ -4234,7 +4126,6 @@ async fn persistent_stream_worker( inbound_tx, outbound_rx, queue_depth_limit, - OrderedWritePolicy::Standalone, None, ) .await; @@ -4249,7 +4140,6 @@ async fn persistent_stream_worker_with_policy( inbound_tx: mpsc::Sender, outbound_rx: FramedWorkerRecv, queue_depth_limit: usize, - write_policy: OrderedWritePolicy, remote_close: Option, ) { let context = Arc::new(context); @@ -4262,10 +4152,9 @@ async fn persistent_stream_worker_with_policy( let reader_context = Arc::clone(&context); let reader_remote_close = remote_close.clone(); let reader = tokio_util::task::AbortOnDropHandle::new(tokio::spawn(async move { - // Pair teardown must interrupt a blocked writer without waiting for it + // Session teardown must interrupt a blocked writer without waiting for it // to poll the reader's terminal event. Preserve the close cause first. - let _cancel_pair_on_exit = (write_policy != OrderedWritePolicy::Standalone) - .then(|| reader_context.stream_token.clone().drop_guard()); + let _cancel_session_on_exit = reader_context.stream_token.clone().drop_guard(); let mut recv = recv; loop { let frame = tokio::select! { @@ -4371,32 +4260,21 @@ async fn persistent_stream_worker_with_policy( } => { match outbound { Some(queued_frame) => { - // Standalone streams finish the current frame on local - // cancellation. Paired streams can interrupt the write: - // teardown resets both roles before either can be reused. + // Cancellation resets the stream before a replacement can write. let result = tokio::select! { biased; _ = context.connection_token.cancelled() => break, - _ = context.stream_token.cancelled(), - if write_policy != OrderedWritePolicy::Standalone => break, - result = async { - if write_policy == OrderedWritePolicy::Standalone { - write_queued_ordered_frame(&mut send, queued_frame, - context.limits, context.outbound_frame_cap).await - } else { - queued_frame.write_with(|frame| write_ordered_frame_with_policy( - &mut send, frame, context.limits, - context.outbound_frame_cap, write_policy, - )).await - } - } => result, + _ = context.stream_token.cancelled() => break, + result = queued_frame.write_with(|frame| write_ordered_frame_with_policy( + &mut send, frame, context.limits, + context.outbound_frame_cap, context.write_policy, + )) => result, }; if let Err(error) = result { - if write_policy == OrderedWritePolicy::PairData - && error.is::() + if error.is::() { debug!(stream_kind, stream_id = context.stream_id, - "retiring Zakura stream pair after data write timeout"); + "retiring Zakura service session after stream write timeout"); break; } if ordered_stream_write_was_stopped(&error) { @@ -4450,12 +4328,9 @@ async fn persistent_stream_worker_with_policy( } } - if write_policy != OrderedWritePolicy::Standalone { - // A cancelled pair cannot leave a partial frame followed by a graceful - // FIN. Reset it before any replacement session may write. - let _ = send.reset(VarInt::from_u32(ZAKURA_CLOSE_NEUTRAL)); - context.stream_token.cancel(); - } + // Never leave a partial frame followed by a graceful FIN. + let _ = send.reset(VarInt::from_u32(ZAKURA_CLOSE_NEUTRAL)); + context.stream_token.cancel(); reader.abort(); // Keep the stream permit until the reader has actually dropped its buffers. let _ = reader.await; @@ -4791,6 +4666,7 @@ async fn write_control_payload( Ok(()) } +#[cfg(test)] async fn write_ordered_frame( send: &mut SendStream, frame: Frame, @@ -4802,7 +4678,7 @@ async fn write_ordered_frame( frame, limits, max_frame_bytes, - OrderedWritePolicy::Standalone, + StreamWritePolicy::Timeout(OUTBOUND_STREAM_WRITE_TIMEOUT), ) .await } @@ -4812,7 +4688,7 @@ async fn write_ordered_frame_with_policy( frame: Frame, limits: ZakuraConnectionLimits, max_frame_bytes: u32, - write_policy: OrderedWritePolicy, + write_policy: StreamWritePolicy, ) -> Result<(), BoxError> { // Mirror `write_response_frame`: a persistent ordered-stream frame whose // payload exceeds the peer's negotiated `max_message_bytes` would be @@ -4829,35 +4705,17 @@ async fn write_ordered_frame_with_policy( .into()); } let frame = frame.encode(max_frame_bytes)?; - if write_policy == OrderedWritePolicy::PairRequests { - // Serving backpressure can outlast a write deadline while downloads - // progress. The service's liveness policy or either reader ending cancels the pair. - send.write_all(&frame).await?; - } else { - let write_timeout = if write_policy == OrderedWritePolicy::PairData { - PAIRED_DATA_WRITE_TIMEOUT - } else { - OUTBOUND_STREAM_WRITE_TIMEOUT - }; - timeout(write_timeout, send.write_all(&frame)) - .await - .map_err(|_| -> BoxError { Box::new(OrderedFrameWriteTimeout) })??; + match write_policy { + StreamWritePolicy::UntilCancelled => send.write_all(&frame).await?, + StreamWritePolicy::Timeout(duration) => { + timeout(duration, send.write_all(&frame)) + .await + .map_err(|_| -> BoxError { Box::new(OrderedFrameWriteTimeout) })?? + } } Ok(()) } -/// Retain the queued frame's ownership through its QUIC write or cancellation. -async fn write_queued_ordered_frame( - send: &mut SendStream, - queued_frame: QueuedFrame, - limits: ZakuraConnectionLimits, - max_frame_bytes: u32, -) -> Result<(), BoxError> { - queued_frame - .write_with(|frame| write_ordered_frame(send, frame, limits, max_frame_bytes)) - .await -} - #[allow(clippy::too_many_arguments)] async fn write_outbound_request_frame( connection: &Connection, @@ -5710,8 +5568,8 @@ pub enum ZakuraHandlerError { #[error("invalid message type {0} for this stream role")] InvalidMessageType(u16), /// Two ordered stream roles failed to name one complete session. - #[error("invalid Zakura ordered stream pair")] - InvalidOrderedPair, + #[error("invalid Zakura service session")] + InvalidServiceSession, /// A bounded read/write timed out. #[error("Zakura {0} timed out")] Timeout(&'static str), @@ -5748,7 +5606,7 @@ pub enum ZakuraHandlerError { ResourceLimit(&'static str), /// Another connection reserved the last service session slot. #[error("ordered service session capacity is full")] - OrderedSessionFull, + SessionFull, /// The peer exceeded its per-kind inbound message rate. #[error("Zakura message rate exceeded")] RateLimited, @@ -6586,7 +6444,7 @@ mod tests { version: 1, frame_cap: 1024, capability: ZAKURA_CAP_LEGACY_GOSSIP, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }; let service = GenerationGuardedRecordingService::new(vec![stream]); let registry = Arc::new( @@ -6701,7 +6559,7 @@ mod tests { version: 1, frame_cap: 1024, capability: ZAKURA_CAP_LEGACY_GOSSIP, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }; let service = GenerationGuardedRecordingService::new(vec![stream]); let registry = Arc::new( @@ -7171,7 +7029,7 @@ mod tests { version: 1, frame_cap: 64 * 1024, capability: 1 << 16, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }; let _guard = zakura_test::init(); @@ -7303,8 +7161,7 @@ mod tests { } #[tokio::test] - async fn header_ordered_session_waits_for_coordinator_capability_epoch() -> Result<(), BoxError> - { + async fn header_session_waits_for_coordinator_capability_epoch() -> Result<(), BoxError> { use zakura_node_services::sync_lifecycle::{ BlockServiceDemand, HeaderServiceDemand, LifecycleEpoch, SyncServiceDemand, }; @@ -7328,7 +7185,7 @@ mod tests { ZAKURA_CAP_HEADER_SYNC, ServicePeerDirection::Outbound, )); - let OrderedSessionDemand::WaitForChange(changed) = service.ordered_session_demand( + let SessionDemand::WaitForChange(changed) = service.session_demand( test_conn_id(), &peer, ZAKURA_CAP_HEADER_SYNC, @@ -7791,15 +7648,15 @@ mod tests { version: ZAKURA_HEADER_SYNC_STREAM_VERSION, frame_cap: 1, capability: ZAKURA_CAP_HEADER_SYNC, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }; - let exit = OrderedSessionExit { + let exit = SessionExit { stream: initiator_opened, session_id: 1, opened_locally: true, }; - let initiator_policy = OrderedStreamPolicy { - opening: OrderedStreamOpening::InitiatorOnly, + let initiator_policy = SessionPolicy { + opening: SessionOpening::InitiatorOnly, reopen: true, }; assert!(should_reopen_ordered_session( @@ -7825,7 +7682,7 @@ mod tests { )); assert!(!should_reopen_ordered_session( exit, - OrderedStreamPolicy::default(), + SessionPolicy::default(), true, false, false, @@ -7834,7 +7691,7 @@ mod tests { // The authorized side replaces an accepted stream too; replacement does // not depend on which physical generation just exited. assert!(should_reopen_ordered_session( - OrderedSessionExit { + SessionExit { opened_locally: false, ..exit }, @@ -7844,15 +7701,15 @@ mod tests { false, )); - let either_peer = OrderedSessionExit { + let either_peer = SessionExit { stream: Stream { kind: ZAKURA_STREAM_BLOCK_SYNC, ..initiator_opened }, ..exit }; - let either_policy = OrderedStreamPolicy { - opening: OrderedStreamOpening::EitherSide, + let either_policy = SessionPolicy { + opening: SessionOpening::EitherSide, reopen: true, }; assert!(should_reopen_ordered_session( @@ -7872,7 +7729,7 @@ mod tests { assert!(opens_ordered_stream_locally(either_policy, false, true)); assert!(!opens_ordered_stream_locally(either_policy, true, false)); - let request_response = OrderedSessionExit { + let request_response = SessionExit { stream: Stream { kind: 102, mode: StreamMode::RequestResponse, @@ -7891,17 +7748,17 @@ mod tests { #[test] fn either_side_session_has_one_proactive_opener_across_connection_roles() { - let policy = OrderedStreamPolicy { - opening: OrderedStreamOpening::EitherSide, + let policy = SessionPolicy { + opening: SessionOpening::EitherSide, reopen: true, }; - let exit = OrderedSessionExit { + let exit = SessionExit { stream: Stream { kind: ZAKURA_STREAM_BLOCK_SYNC, version: ZAKURA_BLOCK_SYNC_STREAM_VERSION, frame_cap: 1, capability: ZAKURA_CAP_BLOCK_SYNC, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }, session_id: 1, opened_locally: false, @@ -7936,9 +7793,9 @@ mod tests { version: ZAKURA_HEADER_SYNC_STREAM_VERSION, frame_cap: 1, capability: ZAKURA_CAP_HEADER_SYNC, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }; - let mut session = OrderedSessionState::new(stream); + let mut session = ServiceSessionState::new(stream); session.remote_session_id = Some(2); assert!(!session.remove_active_session(false, 1)); @@ -7953,33 +7810,33 @@ mod tests { version: ZAKURA_HEADER_SYNC_STREAM_VERSION, frame_cap: 1, capability: ZAKURA_CAP_HEADER_SYNC, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, } } #[tokio::test(start_paused = true)] - async fn ordered_session_waits_deduplicate_demand() { + async fn session_waits_deduplicate_demand() { let stream = test_ordered_stream(); - let mut session = OrderedSessionState::new(stream); - let mut waits = OrderedSessionWaits::new(); + let mut session = ServiceSessionState::new(stream); + let mut waits = SessionWaits::new(); let (change_tx, change_rx) = watch::channel(()); session.schedule_demand( &mut waits, - OrderedSessionDemand::WaitForChange(Box::pin(async move { + SessionDemand::WaitForChange(Box::pin(async move { let mut change_rx = change_rx; let _ = change_rx.changed().await; })), ); session.schedule_demand( &mut waits, - OrderedSessionDemand::RetryAt(std::time::Instant::now() + Duration::from_secs(60)), + SessionDemand::RetryAt(std::time::Instant::now() + Duration::from_secs(60)), ); assert_eq!(waits.len(), 1); assert_eq!( session.reopen_state, - OrderedSessionReopenState::Waiting(OrderedSessionWaitReason::Demand) + SessionReopenState::Waiting(SessionWaitReason::Demand) ); assert_eq!( change_tx.receiver_count(), @@ -7989,10 +7846,10 @@ mod tests { } #[tokio::test(start_paused = true)] - async fn ordered_session_demand_replaces_transport_backoff() { + async fn session_demand_replaces_transport_backoff() { let stream = test_ordered_stream(); - let mut session = OrderedSessionState::new(stream); - let mut waits = OrderedSessionWaits::new(); + let mut session = ServiceSessionState::new(stream); + let mut waits = SessionWaits::new(); assert_eq!( session.schedule_transport_backoff(&mut waits), @@ -8000,26 +7857,26 @@ mod tests { ); session.schedule_demand( &mut waits, - OrderedSessionDemand::RetryAt(std::time::Instant::now() + Duration::from_secs(60)), + SessionDemand::RetryAt(std::time::Instant::now() + Duration::from_secs(60)), ); assert_eq!(session.reopen_attempts, 1); assert_eq!(waits.len(), 1); assert_eq!( session.reopen_state, - OrderedSessionReopenState::Waiting(OrderedSessionWaitReason::Demand) + SessionReopenState::Waiting(SessionWaitReason::Demand) ); } #[tokio::test(start_paused = true)] async fn ordered_session_exit_keeps_existing_demand_wait() { let stream = test_ordered_stream(); - let mut session = OrderedSessionState::new(stream); - let mut waits = OrderedSessionWaits::new(); + let mut session = ServiceSessionState::new(stream); + let mut waits = SessionWaits::new(); session.schedule_demand( &mut waits, - OrderedSessionDemand::RetryAt(std::time::Instant::now() + Duration::from_secs(60)), + SessionDemand::RetryAt(std::time::Instant::now() + Duration::from_secs(60)), ); assert_eq!(session.schedule_transport_backoff(&mut waits), None); @@ -8027,24 +7884,24 @@ mod tests { assert_eq!(waits.len(), 1); assert_eq!( session.reopen_state, - OrderedSessionReopenState::Waiting(OrderedSessionWaitReason::Demand) + SessionReopenState::Waiting(SessionWaitReason::Demand) ); } #[tokio::test] async fn ordered_session_retirement_cancels_and_blocks_transport_waits() { let stream = test_ordered_stream(); - let mut session = OrderedSessionState::new(stream); - let mut waits = OrderedSessionWaits::new(); + let mut session = ServiceSessionState::new(stream); + let mut waits = SessionWaits::new(); session.schedule_demand( &mut waits, - OrderedSessionDemand::RetryAt(std::time::Instant::now() + Duration::from_secs(60)), + SessionDemand::RetryAt(std::time::Instant::now() + Duration::from_secs(60)), ); - session.schedule_demand(&mut waits, OrderedSessionDemand::Retire); + session.schedule_demand(&mut waits, SessionDemand::Retire); assert!(waits.is_empty()); - assert_eq!(session.reopen_state, OrderedSessionReopenState::Retired); + assert_eq!(session.reopen_state, SessionReopenState::Retired); assert_eq!(session.schedule_transport_backoff(&mut waits), None); assert_eq!(session.reopen_attempts, 0); } @@ -8052,26 +7909,26 @@ mod tests { #[tokio::test] async fn ordered_session_adoption_drops_pending_wait() { let stream = test_ordered_stream(); - let mut session = OrderedSessionState::new(stream); - let mut waits = OrderedSessionWaits::new(); + let mut session = ServiceSessionState::new(stream); + let mut waits = SessionWaits::new(); - session.schedule_demand(&mut waits, OrderedSessionDemand::OpenNow); + session.schedule_demand(&mut waits, SessionDemand::OpenNow); session.cancel_wait(&mut waits); assert!(waits.is_empty()); - assert_eq!(session.reopen_state, OrderedSessionReopenState::Idle); + assert_eq!(session.reopen_state, SessionReopenState::Idle); } #[tokio::test] - async fn ordered_session_demand_change_yields_exactly_one_reopen() { + async fn session_demand_change_yields_exactly_one_reopen() { let stream = test_ordered_stream(); - let mut session = OrderedSessionState::new(stream); - let mut waits = OrderedSessionWaits::new(); + let mut session = ServiceSessionState::new(stream); + let mut waits = SessionWaits::new(); let (change_tx, mut change_rx) = watch::channel(()); session.schedule_demand( &mut waits, - OrderedSessionDemand::WaitForChange(Box::pin(async move { + SessionDemand::WaitForChange(Box::pin(async move { let _ = change_rx.changed().await; })), ); @@ -8090,14 +7947,14 @@ mod tests { } #[tokio::test(start_paused = true)] - async fn ordered_session_demand_deadline_yields_exactly_one_reopen() { + async fn session_demand_deadline_yields_exactly_one_reopen() { let stream = test_ordered_stream(); - let mut session = OrderedSessionState::new(stream); - let mut waits = OrderedSessionWaits::new(); + let mut session = ServiceSessionState::new(stream); + let mut waits = SessionWaits::new(); session.schedule_demand( &mut waits, - OrderedSessionDemand::RetryAt(std::time::Instant::now() + Duration::from_secs(60)), + SessionDemand::RetryAt(std::time::Instant::now() + Duration::from_secs(60)), ); tokio::time::advance(Duration::from_secs(60)).await; let (kind, ()) = waits @@ -8113,13 +7970,13 @@ mod tests { #[test] fn ordered_session_connection_teardown_drops_pending_waits() { let stream = test_ordered_stream(); - let mut session = OrderedSessionState::new(stream); - let mut waits = OrderedSessionWaits::new(); + let mut session = ServiceSessionState::new(stream); + let mut waits = SessionWaits::new(); let (change_tx, mut change_rx) = watch::channel(()); session.schedule_demand( &mut waits, - OrderedSessionDemand::WaitForChange(Box::pin(async move { + SessionDemand::WaitForChange(Box::pin(async move { let _ = change_rx.changed().await; })), ); @@ -8196,7 +8053,7 @@ mod tests { version: ZAKURA_BLOCK_SYNC_STREAM_VERSION, frame_cap: 2_000_009, capability: ZAKURA_CAP_BLOCK_SYNC, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }; let encoded = frame.encode(stream.frame_cap)?; // Opening a QUIC stream becomes visible to the receiver after the first bytes. @@ -8217,6 +8074,7 @@ mod tests { message_payload_limits: &[], message_types: None, queue_depths: None, + write_policy: StreamWritePolicy::Timeout(OUTBOUND_STREAM_WRITE_TIMEOUT), session_resources: None, outbound_frame_cap: stream.frame_cap, message_bucket: Arc::new(std::sync::Mutex::new(TokenBucket::new(128))), @@ -8409,7 +8267,7 @@ mod tests { version: ZAKURA_STREAM_VERSION_1, frame_cap: 2_000_009, capability: ZAKURA_CAP_BLOCK_SYNC, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }; let context = StreamWorkerContext { conn: ZakuraConnTrace::without_peer(1), @@ -8421,6 +8279,7 @@ mod tests { message_payload_limits: &[], message_types: None, queue_depths: None, + write_policy: StreamWritePolicy::Timeout(OUTBOUND_STREAM_WRITE_TIMEOUT), session_resources: None, outbound_frame_cap: application_frame_cap(&limits, stream), message_bucket: Arc::new(std::sync::Mutex::new(TokenBucket::new(128))), @@ -8437,7 +8296,7 @@ mod tests { max_frame_bytes: inbound_frame_cap_for_stream(&limits, stream), }; let mut workers = JoinSet::new(); - let (ordered_session_exit_tx, mut ordered_session_exit_rx) = mpsc::unbounded_channel(); + let (session_exit_tx, mut session_exit_rx) = mpsc::unbounded_channel(); let admitted = spawn_persistent_stream_worker( &mut workers, server_send, @@ -8447,7 +8306,7 @@ mod tests { context, 1, true, - ordered_session_exit_tx, + session_exit_tx, ); let response = Frame { @@ -8455,7 +8314,7 @@ mod tests { flags: 0, payload: vec![7; 1024 * 1024], }; - admitted.send.try_send(response.clone()).unwrap(); + admitted.streams[0].send.try_send(response.clone()).unwrap(); // Reading the header proves the write has started. The much smaller // QUIC windows keep the remaining payload blocked until we drain it. let mut header = [0; FRAME_HEADER_BYTES]; @@ -8465,21 +8324,16 @@ mod tests { &response.encode(stream.frame_cap)?[..FRAME_HEADER_BYTES] ); admitted.cancel_token.cancel(); + let mut payload = vec![0; response.payload.len()]; assert!( - timeout(Duration::from_millis(50), ordered_session_exit_rx.recv()) - .await + timeout(Duration::from_secs(2), client_recv.read_exact(&mut payload)) + .await? .is_err(), - "stream cancellation must finish the current frame before reporting exit" - ); - let mut payload = vec![0; response.payload.len()]; - timeout(Duration::from_secs(2), client_recv.read_exact(&mut payload)).await??; - assert_eq!( - payload, response.payload, - "the peer must receive a complete frame" + "session cancellation resets an unfinished frame" ); // The exit must be reported, or the connection loop never prunes the dead // generation and never reopens the stream. - let exited = timeout(Duration::from_secs(1), ordered_session_exit_rx.recv()) + let exited = timeout(Duration::from_secs(1), session_exit_rx.recv()) .await .expect("stream cancellation reports worker exit") .expect("exit channel stays open"); @@ -8620,7 +8474,7 @@ mod tests { version: ZAKURA_STREAM_VERSION_1, frame_cap: LOCAL_MAX_CONTROL_FRAME_BYTES, capability: ZAKURA_CAP_LEGACY_GOSSIP, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }; let context = StreamWorkerContext { conn: ZakuraConnTrace::without_peer(1), @@ -8632,6 +8486,7 @@ mod tests { message_payload_limits: &[], message_types: None, queue_depths: None, + write_policy: StreamWritePolicy::Timeout(OUTBOUND_STREAM_WRITE_TIMEOUT), session_resources: None, outbound_frame_cap: application_frame_cap(&limits, stream), message_bucket: Arc::new(std::sync::Mutex::new(TokenBucket::new( @@ -8764,7 +8619,7 @@ mod tests { version: ZAKURA_STREAM_VERSION_1, frame_cap: LOCAL_MAX_CONTROL_FRAME_BYTES, capability: ZAKURA_CAP_DISCOVERY, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }; let frame_cap = application_frame_cap(&limits, stream); @@ -9098,7 +8953,7 @@ mod tests { version: ZAKURA_STREAM_VERSION_1, frame_cap: LOCAL_MAX_CONTROL_FRAME_BYTES, capability: ZAKURA_CAP_LEGACY_GOSSIP, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }; let limits = ZakuraConnectionLimits { @@ -9358,15 +9213,15 @@ mod tests { close_cause: CloseCause::new(), freshness_tx, }; - let (ordered_session_exit_tx, _ordered_session_exit_rx) = mpsc::unbounded_channel(); + let (session_exit_tx, _session_exit_rx) = mpsc::unbounded_channel(); let admitted = handler .admit_bi_stream( server_send, server_recv, &mut admission, 16, - ordered_session_exit_tx, - &mut PendingOrderedPairs::default(), + session_exit_tx, + &mut PendingSessions::default(), ) .await; @@ -9490,46 +9345,53 @@ mod tests { #[test] fn supported_stream_accepts_registered_kinds_at_declared_version_only() { - let registry = ServiceRegistry::new(vec![Arc::new(DeclaredStreamService { - streams: vec![ - Stream { - kind: LEGACY_GOSSIP_STREAM_KIND, - version: ZAKURA_STREAM_VERSION_1, - frame_cap: 1024, - capability: ZAKURA_CAP_LEGACY_GOSSIP, - mode: StreamMode::Ordered, - }, - Stream { - kind: LEGACY_REQUEST_STREAM_KIND, - version: ZAKURA_STREAM_VERSION_1, - frame_cap: 1024, - capability: ZAKURA_CAP_LEGACY_GOSSIP, - mode: StreamMode::RequestResponse, - }, - Stream { - kind: DISCOVERY_STREAM_KIND, - version: ZAKURA_STREAM_VERSION_1, - frame_cap: 1024, - capability: ZAKURA_CAP_DISCOVERY, - mode: StreamMode::Ordered, - }, - Stream { - kind: HEADER_SYNC_STREAM_KIND, - version: ZAKURA_HEADER_SYNC_STREAM_VERSION, - frame_cap: 1024, - capability: ZAKURA_CAP_HEADER_SYNC, - mode: StreamMode::Ordered, - }, - Stream { - kind: ZAKURA_STREAM_BLOCK_SYNC, - version: ZAKURA_STREAM_VERSION_1, - frame_cap: MAX_BS_FRAME_BYTES, - capability: crate::zakura::ZAKURA_CAP_BLOCK_SYNC, - mode: StreamMode::Ordered, - }, - ], - }) as Arc]) - .expect("test registry declares unique stream kinds"); + let streams = vec![ + Stream { + kind: LEGACY_GOSSIP_STREAM_KIND, + version: ZAKURA_STREAM_VERSION_1, + frame_cap: 1024, + capability: ZAKURA_CAP_LEGACY_GOSSIP, + mode: StreamMode::Persistent, + }, + Stream { + kind: LEGACY_REQUEST_STREAM_KIND, + version: ZAKURA_STREAM_VERSION_1, + frame_cap: 1024, + capability: ZAKURA_CAP_LEGACY_GOSSIP, + mode: StreamMode::RequestResponse, + }, + Stream { + kind: DISCOVERY_STREAM_KIND, + version: ZAKURA_STREAM_VERSION_1, + frame_cap: 1024, + capability: ZAKURA_CAP_DISCOVERY, + mode: StreamMode::Persistent, + }, + Stream { + kind: HEADER_SYNC_STREAM_KIND, + version: ZAKURA_HEADER_SYNC_STREAM_VERSION, + frame_cap: 1024, + capability: ZAKURA_CAP_HEADER_SYNC, + mode: StreamMode::Persistent, + }, + Stream { + kind: ZAKURA_STREAM_BLOCK_SYNC, + version: ZAKURA_STREAM_VERSION_1, + frame_cap: MAX_BS_FRAME_BYTES, + capability: crate::zakura::ZAKURA_CAP_BLOCK_SYNC, + mode: StreamMode::Persistent, + }, + ]; + let services = streams + .into_iter() + .map(|stream| { + Arc::new(DeclaredStreamService { + streams: vec![stream], + }) as Arc + }) + .collect(); + let registry = + ServiceRegistry::new(services).expect("test registry declares unique stream kinds"); for (kind, version) in [ (LEGACY_GOSSIP_STREAM_KIND, ZAKURA_STREAM_VERSION_1), @@ -9565,7 +9427,9 @@ mod tests { "the predecessor header-sync stream version is rejected" ); assert!( - registry.ordered_streams_for_negotiated(1 << 4).is_empty(), + registry + .persistent_streams_for_negotiated(1 << 4) + .is_empty(), "the retired predecessor capability opens no header-sync stream" ); diff --git a/crates/zakura-network/src/zakura/handler/ordered_pair.rs b/crates/zakura-network/src/zakura/handler/ordered_pair.rs deleted file mode 100644 index b92244963b..0000000000 --- a/crates/zakura-network/src/zakura/handler/ordered_pair.rs +++ /dev/null @@ -1,426 +0,0 @@ -//! Set up and retire two persistent QUIC streams as one service session. -//! -//! A pair has a request stream and a data/control stream. In block sync's paired -//! layout, `GetBlocks` uses the request stream; `Status`, blocks, and ending -//! messages use the data/control stream. Each stream has its own `SendStream` -//! and `RecvStream`: those are the two directions of one stream. -//! -//! The service registry defines the roles and their allowed messages. Each role -//! has its own queues and worker. A blocked request reader or writer does not stop -//! the data/control worker from being polled. Both streams share cancellation and -//! are retired together. -//! -//! Setup protects each stream with `SetupIo`, retains its metadata in -//! `PreparedOrderedStream`, and starts both workers together once the pair is -//! complete. Incoming roles can arrive in either order; `PendingOrderedPairs` -//! matches them by the opener's pair ID under a bounded setup deadline. - -use super::*; -use crate::zakura::OrderedStreamPair; - -/// Own both directions of one QUIC stream until setup hands them to a new owner. -/// -/// While the handles are present, dropping this guard resets sending and stops -/// receiving. This also cleans up failed or cancelled setup, including a partly -/// written prelude (the header identifying the stream's role and version). -/// `take` transfers ownership and leaves `None`, disabling this guard's cleanup. -pub(super) struct SetupIo(Option<(SendStream, RecvStream)>); - -impl SetupIo { - pub(super) fn new(send: SendStream, recv: RecvStream) -> Self { - Self(Some((send, recv))) - } - - /// Borrow the handles for setup I/O while retaining responsibility for cleanup. - pub(super) fn streams(&mut self) -> (&mut SendStream, &mut RecvStream) { - let (send, recv) = self - .0 - .as_mut() - .expect("setup owns both stream halves until handoff"); - (send, recv) - } - - /// Hand both handles to the caller, which becomes responsible for their lifetime. - pub(super) fn take(mut self) -> (SendStream, RecvStream) { - self.0 - .take() - .expect("setup owns both stream halves until handoff") - } -} - -impl Drop for SetupIo { - fn drop(&mut self) { - if let Some((send, recv)) = &mut self.0 { - let _ = send.reset(VarInt::from_u32(ZAKURA_CLOSE_RESOURCE)); - let _ = recv.stop(VarInt::from_u32(ZAKURA_CLOSE_RESOURCE)); - } - } -} - -/// One stream's I/O, setup metadata, and resources, before its worker starts. -/// -/// Retaining this value keeps its transport stream permit and any service -/// reservation charged, including while it waits for the other role. Dropping -/// it stops the stream through `SetupIo` and releases its resource ownership. -pub(super) struct PreparedOrderedStream { - io: SetupIo, - stream: Stream, - prelude: StreamPrelude, - context: StreamWorkerContext, -} - -impl PreparedOrderedStream { - /// Attach the pair's shared service reservation before starting its workers. - pub(super) fn set_session_resources( - &mut self, - resources: Option>, - ) { - self.context.session_resources = resources; - } - - pub(super) fn new( - send: SendStream, - recv: RecvStream, - stream: Stream, - prelude: StreamPrelude, - context: StreamWorkerContext, - ) -> Self { - Self { - io: SetupIo(Some((send, recv))), - stream, - prelude, - context, - } - } - - /// Start a standalone ordered service that does not declare a companion role. - pub(super) fn spawn_single( - self, - workers: &mut JoinSet<()>, - queue_depth: usize, - opened_locally: bool, - exits: mpsc::UnboundedSender, - ) -> AdmittedOrderedSession { - let (send, recv) = self.io.take(); - spawn_persistent_stream_worker( - workers, - send, - recv, - self.stream, - self.prelude, - self.context, - queue_depth, - opened_locally, - exits, - ) - } -} - -/// The first incoming role, held until its companion arrives or setup expires. -struct PendingPair { - // Shared wire ID chosen by the opener, scoped to this connection and opener. - id: u64, - first: PreparedOrderedStream, - // Starts with the first role; subsequent arrivals cannot extend it. - deadline: Instant, -} - -/// Match incoming request and data/control streams on this connection. -/// -/// Retain at most one incomplete offer per negotiated pair. Keying by the data -/// role's kind, rather than the peer's chosen pair ID, prevents changing IDs from -/// creating extra pending offers. Each waiting stream owns its transport permit -/// and shares the service reservation that will move into the running pair. -#[derive(Default)] -pub(super) struct PendingOrderedPairs { - pairs: HashMap, - // One entry per negotiated pair, so repeated failures cannot grow this map. - retry_after: HashMap, -} - -impl PendingOrderedPairs { - /// Reserve service capacity for the first role, or share the waiting role's - /// reservation. A service that limits session slots charges the pair once; - /// each role also holds its own transport stream permit. Recently expired - /// offers must wait out their cooldown before reserving again. `insert` - /// checks the pair ID and role. - pub(super) fn reserve_or_share( - &self, - pair: OrderedStreamPair, - registry: &ServiceRegistry, - direction: ServicePeerDirection, - ) -> Result< - Option>, - crate::zakura::OrderedSessionFull, - > { - if self - .retry_after - .get(&pair.data.kind) - .is_some_and(|retry| *retry > Instant::now()) - { - return Err(crate::zakura::OrderedSessionFull); - } - match self.pairs.get(&pair.data.kind) { - Some(pending) => Ok(pending.first.context.session_resources.clone()), - None => registry - .service_for_kind(pair.data.kind) - .expect("a selected pair has an owning service") - .reserve_ordered_session(direction), - } - } - - /// Earliest setup deadline, used to wake the connection loop even if no more - /// stream bytes arrive from the peer. - pub(super) fn deadline(&self) -> Option { - self.pairs.values().map(|pair| pair.deadline).min() - } - - /// Drop expired offers, stopping their streams and releasing their ownership. - /// A short retry delay gives other connections a chance to use that capacity. - pub(super) fn expire(&mut self, now: Instant) { - self.retry_after.retain(|_, retry| *retry > now); - for (kind, pair) in self.pairs.extract_if(|_, pair| pair.deadline <= now) { - // Give waiting connections a setup interval to claim the released - // capacity before this peer can reserve another incomplete pair. - self.retry_after - .insert(kind, now + pair.first.context.limits.prelude_timeout); - } - } - - /// Retain the first role with `Ok(None)`, or return a complete pair as - /// `Ok(Some((data, requests)))`, regardless of which role arrived first. - /// - /// Both roles must match the negotiated pair, carry the same nonzero wire ID, - /// and arrive before the first role's deadline. Two copies of one role cannot - /// form a pair. The connection handler validates negotiation before calling. - pub(super) fn insert( - &mut self, - pair: OrderedStreamPair, - id: u64, - incoming: PreparedOrderedStream, - ) -> Result, ZakuraHandlerError> { - if id == 0 || (incoming.stream != pair.data && incoming.stream != pair.requests) { - return Err(ZakuraHandlerError::InvalidOrderedPair); - } - // Taking the first role out also makes it drop if matching the second - // role fails below, so a rejected offer cannot keep its old stream alive. - let Some(pending) = self.pairs.remove(&pair.data.kind) else { - let deadline = Instant::now() + incoming.context.limits.prelude_timeout; - self.pairs.insert( - pair.data.kind, - PendingPair { - id, - first: incoming, - deadline, - }, - ); - return Ok(None); - }; - if pending.id != id - || pending.first.stream == incoming.stream - || Instant::now() >= pending.deadline - { - return Err(ZakuraHandlerError::InvalidOrderedPair); - } - if incoming.stream == pair.data { - Ok(Some((incoming, pending.first))) - } else { - Ok(Some((pending.first, incoming))) - } - } -} - -/// Start a complete pair with separate queues and independently polled workers. -/// -/// Request writes may wait for serving capacity while data writes retain a -/// deadline. Both workers share cancellation: when either ends, the other must -/// stop too. Publish one session exit only after both workers have finished, so -/// the connection handler observes the pair's lifetime as a single unit. -pub(super) fn spawn_ordered_pair( - workers: &mut JoinSet<()>, - mut data: PreparedOrderedStream, - mut requests: PreparedOrderedStream, - queue_depth: usize, - opened_locally: bool, - exits: mpsc::UnboundedSender, -) -> AdmittedOrderedSession { - // Tell the service setup is complete. The workers and application senders - // keep the shared service slot owned through teardown. - if let Some(resources) = &data.context.session_resources { - resources.admitted(); - } - let cancel = data.context.connection_token.child_token(); - // Let the service distinguish peer closure from locally initiated cancellation. - let remote_close = CancellationToken::new(); - data.context.stream_token = cancel.clone(); - requests.context.stream_token = cancel.clone(); - let (inbound_depth, outbound_depth) = - bounded_stream_queue_depths(queue_depth, data.context.queue_depths); - let (data_tx, data_rx) = mpsc::channel(inbound_depth); - let (data_send, data_out) = worker_framed_channel(outbound_depth); - // One raw queued request plus one reader-held frame. The service owns its - // decoded waiting request separately from these transport bounds. - let (request_tx, request_rx) = mpsc::channel(1); - let (request_send, request_out) = worker_framed_channel(1); - // Expose data/control as the primary stream and requests as its companion. - // The local data-stream ID identifies this session in lifecycle events; - // the wire pair ID was only needed to match the two roles during setup. - let admitted = AdmittedOrderedSession { - kind: data.stream.kind, - version: data.stream.version, - session_id: data.context.stream_id, - recv: FramedRecv::new(data_rx).with_remote_close(remote_close.clone()), - send: data_send.with_session_resources(data.context.session_resources.clone()), - cancel_token: cancel.clone(), - companion: Some(ServiceStreamRole { - kind: requests.stream.kind, - version: requests.stream.version, - recv: FramedRecv::new(request_rx).with_remote_close(remote_close.clone()), - send: request_send.with_session_resources(requests.context.session_resources.clone()), - }), - }; - let exit = OrderedSessionExit { - stream: data.stream, - session_id: admitted.session_id, - opened_locally, - }; - workers.spawn(async move { - // Aborting this supervising task must cancel the pair as well. - let _cancel_on_exit = cancel.clone().drop_guard(); - let (data_send, data_recv) = data.io.take(); - let (request_send, request_recv) = requests.io.take(); - // Poll both workers concurrently. Awaiting one before starting the other - // would make request backpressure obstruct data/control traffic again. - tokio::join!( - async { - persistent_stream_worker_with_policy( - data_send, - data_recv, - data.prelude, - data.context, - data_tx, - data_out, - inbound_depth, - OrderedWritePolicy::PairData, - Some(remote_close.clone()), - ) - .await; - cancel.cancel(); - }, - async { - persistent_stream_worker_with_policy( - request_send, - request_recv, - requests.prelude, - requests.context, - request_tx, - request_out, - 1, - OrderedWritePolicy::PairRequests, - Some(remote_close.clone()), - ) - .await; - cancel.cancel(); - }, - ); - let _ = exits.send(exit); - }); - admitted -} - -impl ZakuraProtocolHandler { - /// Complete local setup of one outbound stream without starting its worker. - /// - /// Reserve a transport stream slot, open a bidirectional stream, and write its - /// prelude under bounded waits. For a pair, the caller supplies the same - /// nonzero `pair_id` for both roles and attaches their shared service resources - /// afterward. The returned value retains the stream slot until handoff or drop. - #[allow(clippy::too_many_arguments)] - pub(super) async fn prepare_ordered_stream( - &self, - connection: &Connection, - stream: Stream, - pair_id: Option, - stream_sem: &Arc, - message_buckets: &mut MessageRateBuckets, - limits: ZakuraConnectionLimits, - connection_token: CancellationToken, - close_cause: CloseCause, - freshness_tx: watch::Sender, - conn: ZakuraConnTrace, - peer_id: ZakuraPeerId, - ) -> Result { - let stream_id = self.next_stream_id.fetch_add(1, Ordering::Relaxed); - let permit = stream_sem - .clone() - .try_acquire_owned() - .map_err(|_| ZakuraHandlerError::ResourceLimit("ordered stream permit"))?; - let io = timeout(OUTBOUND_STREAM_WRITE_TIMEOUT, connection.open_bi()) - .await - .map_err(|_| ZakuraHandlerError::Timeout("open ordered service stream"))??; - let mut io = SetupIo(Some(io)); - let prelude = StreamPrelude { - magic: STREAM_PRELUDE_MAGIC, - stream_kind: stream.kind, - stream_version: stream.version, - request_id: None, - max_frame_bytes: inbound_frame_cap_for_stream(&limits, stream), - }; - let mut bytes = prelude.encode()?; - // The pair ID follows the ordinary prelude. It matches persistent roles; - // it is separate from the prelude's per-request `request_id` field. - if let Some(id) = pair_id { - bytes.extend_from_slice(&id.to_le_bytes()); - } - timeout( - OUTBOUND_STREAM_WRITE_TIMEOUT, - io.0.as_mut() - .expect("setup retains its stream halves") - .0 - .write_all(&bytes), - ) - .await - .map_err(|_| ZakuraHandlerError::Timeout("ordered stream prelude write"))??; - // Both roles spend the same message-rate budget, so splitting a service - // across two streams does not double its allowance. - let bucket_kind = self - .registry - .ordered_stream_pair(stream) - .map_or(stream.kind, |pair| pair.data.kind); - let message_bucket = message_bucket_for( - message_buckets, - bucket_kind, - limits.message_rate_per_second, - RealClock, - ); - let context = StreamWorkerContext { - conn: conn.clone(), - peer_id, - stream_id, - _permit: permit, - limits, - inbound_frame_cap: prelude.max_frame_bytes, - message_payload_limits: self.registry.message_payload_limits(stream), - message_types: self.registry.message_types(stream), - queue_depths: self.registry.stream_queue_depths(stream), - session_resources: None, - outbound_frame_cap: application_frame_cap(&limits, stream), - message_bucket, - stream_token: connection_token.child_token(), - connection_token, - close_cause, - freshness_tx, - }; - metrics::counter!("zakura.p2p.stream.accepted", "stream_kind" => stream_kind_label(stream.kind)).increment(1); - conn.trace_stream("accepted", stream_id, Some(stream_kind_label(stream.kind))); - Ok(PreparedOrderedStream { - io, - stream, - prelude, - context, - }) - } -} - -#[cfg(test)] -mod tests; diff --git a/crates/zakura-network/src/zakura/handler/service_session.rs b/crates/zakura-network/src/zakura/handler/service_session.rs new file mode 100644 index 0000000000..4ef03490e0 --- /dev/null +++ b/crates/zakura-network/src/zakura/handler/service_session.rs @@ -0,0 +1,342 @@ +//! Set up every required persistent stream before admitting a service session. +//! +//! Each member has independent queues and workers. The session shares a wire +//! identifier, admission reservation, message budget, and cancellation scope. +//! Single-stream sessions retain their existing prelude without a wire identifier. + +use super::*; +use crate::zakura::transport::SessionLayout; + +/// Own both directions of one QUIC stream until setup hands them to a new owner. +/// +/// While the handles are present, dropping this guard resets sending and stops +/// receiving. This also cleans up failed or cancelled setup, including a partly +/// written prelude (the header identifying the stream's role and version). +/// `take` transfers ownership and leaves `None`, disabling this guard's cleanup. +pub(super) struct SetupIo(Option<(SendStream, RecvStream)>); + +impl SetupIo { + pub(super) fn new(send: SendStream, recv: RecvStream) -> Self { + Self(Some((send, recv))) + } + + /// Borrow the handles for setup I/O while retaining responsibility for cleanup. + pub(super) fn streams(&mut self) -> (&mut SendStream, &mut RecvStream) { + let (send, recv) = self + .0 + .as_mut() + .expect("setup owns both stream halves until handoff"); + (send, recv) + } + + /// Hand both handles to the caller, which becomes responsible for their lifetime. + pub(super) fn take(mut self) -> (SendStream, RecvStream) { + self.0 + .take() + .expect("setup owns both stream halves until handoff") + } +} + +impl Drop for SetupIo { + fn drop(&mut self) { + if let Some((send, recv)) = &mut self.0 { + let _ = send.reset(VarInt::from_u32(ZAKURA_CLOSE_RESOURCE)); + let _ = recv.stop(VarInt::from_u32(ZAKURA_CLOSE_RESOURCE)); + } + } +} + +/// One stream's I/O, setup metadata, and resources, before its worker starts. +/// +/// Retaining this value keeps its transport stream permit and any service +/// reservation charged, including while it waits for the remaining members. Dropping +/// it stops the stream through `SetupIo` and releases its resource ownership. +pub(super) struct PreparedStream { + io: SetupIo, + stream: Stream, + prelude: StreamPrelude, + context: StreamWorkerContext, +} + +impl PreparedStream { + /// Attach the session's shared service reservation before starting its workers. + pub(super) fn set_session_resources( + &mut self, + resources: Option>, + ) { + self.context.session_resources = resources; + } + + pub(super) fn new( + send: SendStream, + recv: RecvStream, + stream: Stream, + prelude: StreamPrelude, + context: StreamWorkerContext, + ) -> Self { + Self { + io: SetupIo(Some((send, recv))), + stream, + prelude, + context, + } + } +} + +/// Incomplete setup retains every arrived stream under the first arrival's deadline. +struct PendingSession { + id: u64, + streams: Vec, + deadline: Instant, +} + +/// At most one incomplete offer per service session on this connection. +#[derive(Default)] +pub(super) struct PendingSessions { + sessions: HashMap, + retry_after: HashMap, +} + +impl PendingSessions { + /// Charge the service once, then share its reservation across all members. + pub(super) fn reserve_or_share( + &self, + layout: &SessionLayout, + registry: &ServiceRegistry, + direction: ServicePeerDirection, + ) -> Result>, crate::zakura::SessionFull> { + let kind = layout.primary().kind; + if self + .retry_after + .get(&kind) + .is_some_and(|retry| *retry > Instant::now()) + { + return Err(crate::zakura::SessionFull); + } + match self.sessions.get(&kind) { + Some(pending) => Ok(pending.streams[0].context.session_resources.clone()), + None => registry + .service_for_kind(kind) + .expect("a selected session has an owning service") + .reserve_session(direction), + } + } + + pub(super) fn deadline(&self) -> Option { + self.sessions.values().map(|session| session.deadline).min() + } + + /// Expiry releases every arrived member and briefly defers another offer. + pub(super) fn expire(&mut self, now: Instant) { + self.retry_after.retain(|_, retry| *retry > now); + for (kind, session) in self + .sessions + .extract_if(|_, session| session.deadline <= now) + { + self.retry_after.insert( + kind, + now + session.streams[0].context.limits.prelude_timeout, + ); + } + } + + /// Admit the complete layout in kind order, regardless of arrival order. + pub(super) fn insert( + &mut self, + layout: &SessionLayout, + id: u64, + incoming: PreparedStream, + ) -> Result>, ZakuraHandlerError> { + let kind = layout.primary().kind; + // Remove first so every invalid continuation releases the existing offer. + let pending = self.sessions.remove(&kind); + if (layout.is_multi_stream() && id == 0) || !layout.streams.contains(&incoming.stream) { + return Err(ZakuraHandlerError::InvalidServiceSession); + } + let mut pending = pending.unwrap_or_else(|| PendingSession { + id, + streams: Vec::with_capacity(layout.streams.len()), + deadline: Instant::now() + incoming.context.limits.prelude_timeout, + }); + if pending.id != id + || Instant::now() >= pending.deadline + || pending.streams.iter().any(|s| { + s.stream.kind == incoming.stream.kind || !layout.streams.contains(&s.stream) + }) + { + return Err(ZakuraHandlerError::InvalidServiceSession); + } + pending.streams.push(incoming); + if pending.streams.len() == layout.streams.len() { + pending.streams.sort_unstable_by_key(|s| s.stream.kind); + Ok(Some(pending.streams)) + } else { + self.sessions.insert(kind, pending); + Ok(None) + } + } +} + +/// Run all members independently and report exit only after every worker finishes. +pub(super) fn spawn_service_session( + workers: &mut JoinSet<()>, + streams: Vec, + queue_depth: usize, + opened_locally: bool, + exits: mpsc::UnboundedSender, +) -> AdmittedSession { + let primary = streams + .first() + .expect("a validated session has at least one stream"); + if let Some(resources) = &primary.context.session_resources { + resources.admitted(); + } + let cancel = primary.context.connection_token.child_token(); + let remote_close = CancellationToken::new(); + let mut admitted = AdmittedSession { + kind: primary.stream.kind, + session_id: primary.context.stream_id, + cancel_token: cancel.clone(), + streams: Vec::with_capacity(streams.len()), + }; + let exit = SessionExit { + stream: primary.stream, + session_id: admitted.session_id, + opened_locally, + }; + let mut running = futures::stream::FuturesUnordered::new(); + for mut prepared in streams { + prepared.context.stream_token = cancel.clone(); + let (inbound_depth, outbound_depth) = + bounded_stream_queue_depths(queue_depth, prepared.context.queue_depths); + let (inbound_tx, inbound_rx) = mpsc::channel(inbound_depth); + let (sender, outbound_rx) = worker_framed_channel(outbound_depth); + admitted.streams.push(ServiceStreamRole { + kind: prepared.stream.kind, + version: prepared.stream.version, + recv: FramedRecv::new(inbound_rx).with_remote_close(remote_close.clone()), + send: sender.with_session_resources(prepared.context.session_resources.clone()), + }); + let remote_close = remote_close.clone(); + running.push(async move { + let (send, recv) = prepared.io.take(); + persistent_stream_worker_with_policy( + send, + recv, + prepared.prelude, + prepared.context, + inbound_tx, + outbound_rx, + inbound_depth, + Some(remote_close), + ) + .await; + }); + } + workers.spawn(async move { + let _cancel_on_exit = cancel.clone().drop_guard(); + while running.next().await.is_some() { + cancel.cancel(); + } + let _ = exits.send(exit); + }); + admitted +} + +impl ZakuraProtocolHandler { + /// Complete local setup of one outbound stream without starting its worker. + /// + /// Reserve a transport stream slot, open a bidirectional stream, and write its + /// prelude under bounded waits. For a multi-stream session, the caller supplies the same + /// nonzero `session_id` for every member and attaches their shared service resources + /// afterward. The returned value retains the stream slot until handoff or drop. + #[allow(clippy::too_many_arguments)] + pub(super) async fn prepare_ordered_stream( + &self, + connection: &Connection, + stream: Stream, + session_id: Option, + stream_sem: &Arc, + message_buckets: &mut MessageRateBuckets, + limits: ZakuraConnectionLimits, + connection_token: CancellationToken, + close_cause: CloseCause, + freshness_tx: watch::Sender, + conn: ZakuraConnTrace, + peer_id: ZakuraPeerId, + ) -> Result { + let stream_id = self.next_stream_id.fetch_add(1, Ordering::Relaxed); + let permit = stream_sem + .clone() + .try_acquire_owned() + .map_err(|_| ZakuraHandlerError::ResourceLimit("ordered stream permit"))?; + let io = timeout(OUTBOUND_STREAM_WRITE_TIMEOUT, connection.open_bi()) + .await + .map_err(|_| ZakuraHandlerError::Timeout("open ordered service stream"))??; + let mut io = SetupIo(Some(io)); + let prelude = StreamPrelude { + magic: STREAM_PRELUDE_MAGIC, + stream_kind: stream.kind, + stream_version: stream.version, + request_id: None, + max_frame_bytes: inbound_frame_cap_for_stream(&limits, stream), + }; + let mut bytes = prelude.encode()?; + // The session ID follows the ordinary prelude. It matches persistent roles; + // it is separate from the prelude's per-request `request_id` field. + if let Some(id) = session_id { + bytes.extend_from_slice(&id.to_le_bytes()); + } + timeout( + OUTBOUND_STREAM_WRITE_TIMEOUT, + io.0.as_mut() + .expect("setup retains its stream halves") + .0 + .write_all(&bytes), + ) + .await + .map_err(|_| ZakuraHandlerError::Timeout("ordered stream prelude write"))??; + // All members spend the same message-rate budget, so splitting a service + // across several streams does not double its allowance. + let bucket_kind = self + .registry + .session_layout(stream) + .map_or(stream.kind, |layout| layout.primary().kind); + let message_bucket = message_bucket_for( + message_buckets, + bucket_kind, + limits.message_rate_per_second, + RealClock, + ); + let context = StreamWorkerContext { + conn: conn.clone(), + peer_id, + stream_id, + _permit: permit, + limits, + inbound_frame_cap: prelude.max_frame_bytes, + message_payload_limits: self.registry.message_payload_limits(stream), + message_types: self.registry.message_types(stream), + queue_depths: self.registry.stream_queue_depths(stream), + write_policy: self.registry.stream_write_policy(stream), + session_resources: None, + outbound_frame_cap: application_frame_cap(&limits, stream), + message_bucket, + stream_token: connection_token.child_token(), + connection_token, + close_cause, + freshness_tx, + }; + metrics::counter!("zakura.p2p.stream.accepted", "stream_kind" => stream_kind_label(stream.kind)).increment(1); + conn.trace_stream("accepted", stream_id, Some(stream_kind_label(stream.kind))); + Ok(PreparedStream { + io, + stream, + prelude, + context, + }) + } +} + +#[cfg(test)] +mod tests; diff --git a/crates/zakura-network/src/zakura/handler/ordered_pair/tests.rs b/crates/zakura-network/src/zakura/handler/service_session/tests.rs similarity index 79% rename from crates/zakura-network/src/zakura/handler/ordered_pair/tests.rs rename to crates/zakura-network/src/zakura/handler/service_session/tests.rs index 6053b23260..2e9e0a951a 100644 --- a/crates/zakura-network/src/zakura/handler/ordered_pair/tests.rs +++ b/crates/zakura-network/src/zakura/handler/service_session/tests.rs @@ -7,24 +7,28 @@ const DATA: Stream = Stream { version: 1, frame_cap: 2 * 1024 * 1024, capability: 1 << 16, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }; const REQUESTS: Stream = Stream { kind: 65, ..DATA }; +const EVENTS: Stream = Stream { kind: 67, ..DATA }; +const ONE_SHOT: Stream = Stream { + kind: 68, + mode: StreamMode::RequestResponse, + ..DATA +}; const SIBLING: Stream = Stream { kind: 66, capability: 1 << 17, ..DATA }; -const PAIR: OrderedStreamPair = OrderedStreamPair { - data: DATA, - requests: REQUESTS, -}; const ALPN: &[u8] = b"/zakura/test/ordered-pair/1"; +const TEST_DATA_WRITE_TIMEOUT: Duration = Duration::from_secs(32); const TEST_TIMEOUT: Duration = Duration::from_secs(30); #[derive(Debug)] -struct PairService { - opening: OrderedStreamOpening, +struct SessionService { + streams: &'static [Stream], + opening: SessionOpening, fail_first_reservation: std::sync::atomic::AtomicBool, capacity: Option>, sessions: mpsc::Sender, @@ -35,7 +39,7 @@ struct SessionSlot { _permit: OwnedSemaphorePermit, } -impl crate::zakura::OrderedSessionResources for SessionSlot { +impl crate::zakura::SessionResources for SessionSlot { fn admitted(&self) {} } @@ -64,23 +68,20 @@ impl Service for SiblingService { fn remove_peer(&self, _: &ZakuraPeerId, _: ZakuraConnId) {} } -impl Service for PairService { - fn reserve_ordered_session( +impl Service for SessionService { + fn reserve_session( &self, _: ServicePeerDirection, - ) -> Result< - Option>, - crate::zakura::OrderedSessionFull, - > { + ) -> Result>, crate::zakura::SessionFull> { // Inject the outcome of another connection taking the final slot after // this connection's advisory OpenNow check. A retry can succeed. if self.fail_first_reservation.swap(false, Ordering::SeqCst) { - Err(crate::zakura::OrderedSessionFull) + Err(crate::zakura::SessionFull) } else if let Some(capacity) = &self.capacity { let permit = capacity .clone() .try_acquire_owned() - .map_err(|_| crate::zakura::OrderedSessionFull)?; + .map_err(|_| crate::zakura::SessionFull)?; Ok(Some(Arc::new(SessionSlot { _permit: permit }))) } else { Ok(None) @@ -90,16 +91,23 @@ impl Service for PairService { "test-pair" } fn streams(&self) -> &[Stream] { - &[DATA, REQUESTS] + self.streams + } + fn stream_write_policy(&self, stream: Stream) -> StreamWritePolicy { + if stream == REQUESTS { + StreamWritePolicy::UntilCancelled + } else { + StreamWritePolicy::Timeout(TEST_DATA_WRITE_TIMEOUT) + } } - fn ordered_stream_pair(&self, stream: Stream) -> Option { - [DATA, REQUESTS].contains(&stream).then_some(PAIR) + fn stream_queue_depths(&self, stream: Stream) -> Option<(usize, usize)> { + Some(if stream == EVENTS { (3, 3) } else { (1, 1) }) } - fn stream_queue_depths(&self, _: Stream) -> Option<(usize, usize)> { - Some((1, 1)) + fn as_request_response(&self) -> Option<&dyn crate::zakura::RequestResponseService> { + Some(self) } - fn ordered_stream_policy(&self, _: u16) -> OrderedStreamPolicy { - OrderedStreamPolicy { + fn session_policy(&self) -> SessionPolicy { + SessionPolicy { opening: self.opening, reopen: true, } @@ -113,6 +121,20 @@ impl Service for PairService { fn remove_peer(&self, _: &ZakuraPeerId, _: ZakuraConnId) {} } +impl crate::zakura::RequestResponseService for SessionService { + fn request_frame<'a>( + &'a self, + _: ZakuraPeerId, + _: u16, + _: u64, + _: u32, + _: u32, + frame: Frame, + ) -> BoxRunFuture<'a, Result, SinkReject>> { + Box::pin(async move { Ok(vec![frame]) }) + } +} + struct Session { id: u64, conn_id: u64, @@ -176,6 +198,21 @@ impl Fixture { fail_first_reservation: bool, max_open_streams: Option, retired_sibling: bool, + ) -> Result { + Self::start_with_streams( + fail_first_reservation, + max_open_streams, + retired_sibling, + &[DATA, REQUESTS], + ) + .await + } + + async fn start_with_streams( + fail_first_reservation: bool, + max_open_streams: Option, + retired_sibling: bool, + streams: &'static [Stream], ) -> Result { let mut local = ZakuraLocalLimits::from_config(&Config::default()); if let Some(max_open_streams) = max_open_streams { @@ -199,9 +236,10 @@ impl Fixture { local.clone(), Arc::new( ServiceRegistry::new(vec![ - Arc::new(PairService { + Arc::new(SessionService { + streams, capacity: None, - opening: OrderedStreamOpening::EitherSide, + opening: SessionOpening::EitherSide, sessions, fail_first_reservation: std::sync::atomic::AtomicBool::new( fail_reservation, @@ -319,7 +357,7 @@ async fn paired_data_timeout_preserves_sibling_and_reopens_pair() -> Result<(), } Ok::<_, BoxError>(()) })); - timeout(PAIRED_DATA_WRITE_TIMEOUT + Duration::from_secs(10), async { + timeout(TEST_DATA_WRITE_TIMEOUT + Duration::from_secs(10), async { loop { let ping = frame(1, 17, 64); sibling_send.send(ping.clone()).await?; @@ -333,7 +371,7 @@ async fn paired_data_timeout_preserves_sibling_and_reopens_pair() -> Result<(), }) .await .expect("the paired data writer retires its session at the write deadline")?; - assert!(started.elapsed() >= PAIRED_DATA_WRITE_TIMEOUT); + assert!(started.elapsed() >= TEST_DATA_WRITE_TIMEOUT); timeout(TEST_TIMEOUT, server.cancel.cancelled()).await?; assert!(!client.connection_cancel.is_cancelled()); assert!(!server.connection_cancel.is_cancelled()); @@ -500,19 +538,89 @@ struct RawFixture { connection: Connection, serving: AbortOnDropHandle>, shutdown: CancellationToken, - service: Arc, + service: Arc, sessions: mpsc::Receiver, siblings: mpsc::Receiver, } +#[tokio::test] +async fn three_stream_session_reopens_at_the_exact_stream_limit() -> Result<(), BoxError> { + let mut fixture = + Fixture::start_with_streams(false, Some(3), true, &[DATA, REQUESTS, EVENTS]).await?; + let mut previous_id = None; + for _ in 0..2 { + let mut client = timeout(TEST_TIMEOUT, fixture.client_sessions.recv()) + .await? + .ok_or("missing client session")?; + let mut server = timeout(TEST_TIMEOUT, fixture.server_sessions.recv()) + .await? + .ok_or("missing server session")?; + let mut client_streams = Vec::new(); + let mut server_streams = Vec::new(); + for stream in [DATA, REQUESTS, EVENTS] { + client_streams.push(client.take_stream_with_session_id(stream.kind).unwrap()); + server_streams.push(server.take_stream_with_session_id(stream.kind).unwrap()); + } + let client_id = client_streams[0].0; + let server_id = server_streams[0].0; + assert_ne!(previous_id, Some(client_id)); + previous_id = Some(client_id); + for ((id, _, send), (remote_id, recv, _)) in + client_streams.iter().zip(server_streams.iter_mut()) + { + assert_eq!(*id, client_id); + assert_eq!(*remote_id, server_id); + let message = frame(1, 21, 8); + timeout(TEST_TIMEOUT, send.send(message.clone())).await??; + assert_eq!(timeout(TEST_TIMEOUT, recv.recv()).await?, Some(message)); + } + client.service_cancel_token().cancel(); + timeout(TEST_TIMEOUT, server.service_cancel_token().cancelled()).await?; + assert!(!client.cancel_token().is_cancelled()); + } + fixture.close().await +} + +#[tokio::test] +async fn invalid_third_member_releases_the_entire_pending_session() -> Result<(), BoxError> { + for (stream, id) in [(REQUESTS, 9), (EVENTS, 10), (EVENTS, 0)] { + let mut fixture = + RawFixture::start_with_streams(1, Duration::from_secs(3), &[DATA, REQUESTS, EVENTS]) + .await?; + let _data = fixture.offer(DATA, Some(9)).await?; + let _requests = fixture.offer(REQUESTS, Some(9)).await?; + fixture.wait_for_slots(0, TEST_TIMEOUT).await?; + let _invalid = fixture.offer(stream, Some(id)).await?; + fixture.wait_for_slots(1, Duration::from_secs(1)).await?; + timeout(Duration::from_secs(1), async { + while !fixture.serving.is_finished() { + tokio::time::sleep(Duration::from_millis(5)).await; + } + }) + .await?; + assert!(fixture.sessions.try_recv().is_err()); + fixture.close().await?; + } + Ok(()) +} + impl RawFixture { async fn start(slots: usize, setup_timeout: Duration) -> Result { + Self::start_with_streams(slots, setup_timeout, &[DATA, REQUESTS]).await + } + + async fn start_with_streams( + slots: usize, + setup_timeout: Duration, + streams: &'static [Stream], + ) -> Result { let (router, client, connection, remote) = raw_connection().await?; let local = ZakuraLocalLimits::from_config(&Config::default()); let (sessions_tx, sessions) = mpsc::channel(2); let (siblings_tx, siblings) = mpsc::channel(1); - let service = Arc::new(PairService { - opening: OrderedStreamOpening::InitiatorOnly, + let service = Arc::new(SessionService { + streams, + opening: SessionOpening::InitiatorOnly, fail_first_reservation: std::sync::atomic::AtomicBool::new(false), capacity: Some(Arc::new(Semaphore::new(slots))), sessions: sessions_tx, @@ -614,6 +722,116 @@ impl RawFixture { } } +#[tokio::test] +async fn three_stream_session_waits_for_every_member_and_leaves_requests_independent( +) -> Result<(), BoxError> { + let mut fixture = RawFixture::start_with_streams( + 1, + Duration::from_secs(3), + &[EVENTS, ONE_SHOT, REQUESTS, DATA], + ) + .await?; + let (mut sibling_send, _sibling_recv) = fixture.offer(SIBLING, None).await?; + let mut sibling = timeout(TEST_TIMEOUT, fixture.siblings.recv()) + .await? + .ok_or("missing sibling")?; + let (mut sibling_recv, _sibling_sender) = sibling.take_stream(SIBLING.kind).unwrap(); + + let mut events = fixture.offer(EVENTS, Some(42)).await?; + let mut requests = fixture.offer(REQUESTS, Some(42)).await?; + let event = frame(3, 11, 8); + let request = frame(1, 12, 8); + events.0.write_all(&event.encode(EVENTS.frame_cap)?).await?; + requests + .0 + .write_all(&request.encode(REQUESTS.frame_cap)?) + .await?; + fixture.wait_for_slots(0, TEST_TIMEOUT).await?; + assert!( + timeout(Duration::from_millis(100), fixture.sessions.recv()) + .await + .is_err(), + "two of three required streams cannot start the service" + ); + + let _data = fixture.offer(DATA, Some(42)).await?; + let mut peer = timeout(TEST_TIMEOUT, fixture.sessions.recv()) + .await? + .ok_or("missing session")?; + let cancel = peer.service_cancel_token(); + let (data_id, mut data_recv, _data_send) = peer.take_stream_with_session_id(DATA.kind).unwrap(); + let (request_id, mut request_recv, _request_send) = + peer.take_stream_with_session_id(REQUESTS.kind).unwrap(); + let (event_id, mut event_recv, event_send) = + peer.take_stream_with_session_id(EVENTS.kind).unwrap(); + assert_eq!((data_id, data_id), (request_id, event_id)); + assert_eq!( + event_send.max_capacity(), + 3, + "the service controls each stream's queue" + ); + assert_eq!(timeout(TEST_TIMEOUT, event_recv.recv()).await?, Some(event)); + assert_eq!( + timeout(TEST_TIMEOUT, request_recv.recv()).await?, + Some(request) + ); + assert!(peer.take_stream(ONE_SHOT.kind).is_none()); + + // A per-request stream carries no session identifier and completes independently. + let (mut send, mut recv) = fixture.connection.open_bi().await?; + let mut bytes = StreamPrelude { + magic: STREAM_PRELUDE_MAGIC, + stream_kind: ONE_SHOT.kind, + stream_version: ONE_SHOT.version, + request_id: Some(99), + max_frame_bytes: ONE_SHOT.frame_cap, + } + .encode()?; + let echo = frame(4, 13, 8).encode(ONE_SHOT.frame_cap)?; + bytes.extend_from_slice(&echo); + timeout(TEST_TIMEOUT, send.write_all(&bytes)).await??; + send.finish()?; + assert_eq!(timeout(TEST_TIMEOUT, recv.read_to_end(1024)).await??, echo); + assert!(!cancel.is_cancelled()); + + events.0.reset(0u32.into())?; + events.1.stop(0u32.into())?; + timeout(TEST_TIMEOUT, cancel.cancelled()).await?; + assert!(timeout(TEST_TIMEOUT, data_recv.recv()).await?.is_none()); + assert!(timeout(TEST_TIMEOUT, request_recv.recv()).await?.is_none()); + assert!(timeout(TEST_TIMEOUT, event_recv.recv()).await?.is_none()); + assert!(!peer.cancel_token().is_cancelled()); + let ping = frame(2, 14, 8); + sibling_send + .write_all(&ping.encode(SIBLING.frame_cap)?) + .await?; + assert_eq!( + timeout(TEST_TIMEOUT, sibling_recv.recv()).await?, + Some(ping) + ); + fixture.close().await +} + +#[tokio::test] +async fn incomplete_three_stream_session_releases_all_arrivals_on_expiry() -> Result<(), BoxError> { + let mut fixture = + RawFixture::start_with_streams(1, Duration::from_millis(300), &[DATA, REQUESTS, EVENTS]) + .await?; + let (_request_send, mut request_recv) = fixture.offer(REQUESTS, Some(1)).await?; + let (_event_send, mut event_recv) = fixture.offer(EVENTS, Some(1)).await?; + fixture.wait_for_slots(0, TEST_TIMEOUT).await?; + fixture.wait_for_slots(1, TEST_TIMEOUT).await?; + assert!(fixture.sessions.try_recv().is_err()); + assert!(timeout(TEST_TIMEOUT, request_recv.read_exact(&mut [0; 1])) + .await? + .is_err()); + assert!(timeout(TEST_TIMEOUT, event_recv.read_exact(&mut [0; 1])) + .await? + .is_err()); + assert!(fixture.connection.close_reason().is_none()); + fixture.close().await +} + #[tokio::test] async fn withheld_pair_id_does_not_reserve_service_capacity() -> Result<(), BoxError> { let fixture = RawFixture::start(1, Duration::from_secs(2)).await?; @@ -622,7 +840,7 @@ async fn withheld_pair_id_does_not_reserve_service_capacity() -> Result<(), BoxE assert_eq!(fixture.capacity().available_permits(), 1); let outbound = fixture .service - .reserve_ordered_session(ServicePeerDirection::Outbound)?; + .reserve_session(ServicePeerDirection::Outbound)?; drop(outbound); send.write_all(&1u64.to_le_bytes()).await?; fixture.wait_for_slots(0, TEST_TIMEOUT).await?; @@ -660,7 +878,7 @@ async fn expired_pair_cannot_reclaim_capacity_before_an_outgoing_session() -> Re } let outbound = fixture .service - .reserve_ordered_session(ServicePeerDirection::Outbound)?; + .reserve_session(ServicePeerDirection::Outbound)?; assert_eq!(fixture.capacity().available_permits(), 0); drop(outbound); fixture.close().await @@ -763,6 +981,7 @@ fn raw_worker_context(client: &Endpoint, slots: Arc) -> StreamWorkerC message_payload_limits: &[], message_types: None, queue_depths: None, + write_policy: StreamWritePolicy::UntilCancelled, session_resources: None, outbound_frame_cap: DATA.frame_cap, message_bucket: Arc::new(std::sync::Mutex::new(TokenBucket::new(128))), @@ -806,7 +1025,6 @@ async fn paired_request_reader_close_interrupts_a_blocked_write() -> Result<(), inbound_tx, outbound_rx, 1, - OrderedWritePolicy::PairRequests, Some(remote_close.clone()), ))); assert_eq!( @@ -938,9 +1156,10 @@ async fn incomplete_pairs_expire_and_mismatched_roles_release_stream_permits( Network::Mainnet, ZakuraHandshakeConfig::for_network(&Network::Mainnet), local.clone(), - Arc::new(ServiceRegistry::new(vec![Arc::new(PairService { + Arc::new(ServiceRegistry::new(vec![Arc::new(SessionService { + streams: &[DATA, REQUESTS], capacity: None, - opening: OrderedStreamOpening::EitherSide, + opening: SessionOpening::EitherSide, fail_first_reservation: std::sync::atomic::AtomicBool::new(false), sessions, })])?), @@ -950,7 +1169,7 @@ async fn incomplete_pairs_expire_and_mismatched_roles_release_stream_permits( let peer = ZakuraPeerId::new(client.node_id().as_bytes().to_vec())?; let (freshness, _freshness_rx) = watch::channel(Instant::now()); let cancel = CancellationToken::new(); - let mut pending = PendingOrderedPairs::default(); + let mut pending = PendingSessions::default(); let mut workers = JoinSet::new(); let mut buckets = MessageRateBuckets::new(); let mut open_limiter = TokenBucket::new(100); @@ -1005,7 +1224,11 @@ async fn incomplete_pairs_expire_and_mismatched_roles_release_stream_permits( "incomplete setup is retired locally" ); assert!(pending - .reserve_or_share(PAIR, &handler.registry, ServicePeerDirection::Inbound) + .reserve_or_share( + &handler.registry.session_layout(DATA).unwrap(), + &handler.registry, + ServicePeerDirection::Inbound + ) .is_err()); tokio::time::sleep(limits.prelude_timeout).await; } @@ -1032,9 +1255,10 @@ async fn ineligible_pair_opener_is_rejected_before_service_reservation() -> Resu let _guard = zakura_test::init(); let (router, client, connection, remote) = raw_connection().await?; let (sessions, _sessions_rx) = mpsc::channel(1); - let service = Arc::new(PairService { + let service = Arc::new(SessionService { + streams: &[DATA, REQUESTS], capacity: None, - opening: OrderedStreamOpening::InitiatorOnly, + opening: SessionOpening::InitiatorOnly, fail_first_reservation: std::sync::atomic::AtomicBool::new(true), sessions, }); @@ -1051,7 +1275,7 @@ async fn ineligible_pair_opener_is_rejected_before_service_reservation() -> Resu let mut workers = JoinSet::new(); let mut open_limiter = TokenBucket::new(100); let mut buckets = MessageRateBuckets::new(); - let mut pending = PendingOrderedPairs::default(); + let mut pending = PendingSessions::default(); let (exits, _exit_rx) = mpsc::unbounded_channel(); let mut admission = StreamAdmission { is_initiator: true, diff --git a/crates/zakura-network/src/zakura/handler/tests/connection.rs b/crates/zakura-network/src/zakura/handler/tests/connection.rs index add2f04ed8..5fac19dafd 100644 --- a/crates/zakura-network/src/zakura/handler/tests/connection.rs +++ b/crates/zakura-network/src/zakura/handler/tests/connection.rs @@ -11,9 +11,9 @@ impl ZakuraProtocolHandler { recv: RecvStream, admission: &mut StreamAdmission<'_>, queue_depth: usize, - exits: mpsc::UnboundedSender, - pending: &mut PendingOrderedPairs, - ) -> Option { + exits: mpsc::UnboundedSender, + pending: &mut PendingSessions, + ) -> Option { let incoming = self.begin_bi_stream_setup(send, recv, admission)?.await?; self.finish_bi_stream_setup(incoming, admission, queue_depth, exits, pending) } diff --git a/crates/zakura-network/src/zakura/header_sync/service.rs b/crates/zakura-network/src/zakura/header_sync/service.rs index 96e22957fb..378d25e9fc 100644 --- a/crates/zakura-network/src/zakura/header_sync/service.rs +++ b/crates/zakura-network/src/zakura/header_sync/service.rs @@ -18,8 +18,8 @@ use super::{events::*, pipe::run_peer, wire::*, FRAME_HEADER_BYTES}; use crate::zakura::ZakuraSupervisorHandle; use crate::zakura::{ handle_pipe_exit, spawn_supervised_pipe, BoxRunFuture, CloseCause, Frame, FramedRecv, - FramedSend, OrderedSendError, OrderedSessionDemand, OrderedStreamOpening, OrderedStreamPolicy, - Peer, PeerStreamSession, Service, ServicePeerDirection, Sink, SinkReject, Stream, StreamMode, + FramedSend, OrderedSendError, Peer, PeerStreamSession, Service, ServicePeerDirection, + SessionDemand, SessionOpening, SessionPolicy, Sink, SinkReject, Stream, StreamMode, ZakuraConnId, ZakuraPeerId, ZAKURA_CAP_HEADER_SYNC, }; @@ -35,7 +35,7 @@ const HEADER_SYNC_SERVICE_STREAMS: [Stream; 1] = [Stream { version: ZAKURA_HEADER_SYNC_STREAM_VERSION, frame_cap: HEADER_SYNC_FRAME_CAP, capability: ZAKURA_CAP_HEADER_SYNC, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }]; /// The sole stream declaration for native header sync. @@ -58,7 +58,7 @@ mod stream_tests { version: 8, frame_cap: HEADER_SYNC_FRAME_CAP, capability: 1 << 5, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, } ); } @@ -482,12 +482,12 @@ impl HeaderSyncService { self } - fn coordinator_demand(&self) -> Option { + fn coordinator_demand(&self) -> Option { let mut service_demand = self.service_demand.clone()?; if service_demand.borrow().header.is_enabled() { return None; } - Some(OrderedSessionDemand::WaitForChange(Box::pin(async move { + Some(SessionDemand::WaitForChange(Box::pin(async move { loop { if service_demand.changed().await.is_err() { std::future::pending::<()>().await; @@ -509,20 +509,20 @@ impl Service for HeaderSyncService { header_sync_streams() } - fn ordered_stream_policy(&self, _kind: u16) -> OrderedStreamPolicy { - OrderedStreamPolicy { - opening: OrderedStreamOpening::InitiatorOnly, + fn session_policy(&self) -> SessionPolicy { + SessionPolicy { + opening: SessionOpening::InitiatorOnly, reopen: true, } } - fn ordered_session_demand( + fn session_demand( &self, _conn_id: ZakuraConnId, peer: &ZakuraPeerId, _negotiated: u64, direction: ServicePeerDirection, - ) -> OrderedSessionDemand { + ) -> SessionDemand { if let Some(demand) = self.coordinator_demand() { return demand; } @@ -533,7 +533,7 @@ impl Service for HeaderSyncService { ServicePeerDirection::Outbound => snapshot.outbound_slots_free, }; if slots_free == 0 { - return OrderedSessionDemand::WaitForChange(Box::pin(async move { + return SessionDemand::WaitForChange(Box::pin(async move { if peers.changed().await.is_err() { std::future::pending::<()>().await; } @@ -541,7 +541,7 @@ impl Service for HeaderSyncService { } let Some(node_id) = header_peer_node_id(peer) else { - return OrderedSessionDemand::Retire; + return SessionDemand::Retire; }; if self .header_sync @@ -551,12 +551,12 @@ impl Service for HeaderSyncService { { // The reactor publishes the backed-off set without per-node deadlines. // The service re-offers each target after one conservative backoff window. - return OrderedSessionDemand::RetryAt( + return SessionDemand::RetryAt( std::time::Instant::now() + HEADER_SYNC_ADVISORY_BACKOFF, ); } - OrderedSessionDemand::OpenNow + SessionDemand::OpenNow } fn wants_peer( @@ -802,19 +802,19 @@ impl Service for HeaderSyncPassthroughService { header_sync_streams() } - fn ordered_stream_policy(&self, kind: u16) -> OrderedStreamPolicy { - self.inner.ordered_stream_policy(kind) + fn session_policy(&self) -> SessionPolicy { + self.inner.session_policy() } - fn ordered_session_demand( + fn session_demand( &self, conn_id: ZakuraConnId, peer: &ZakuraPeerId, negotiated: u64, direction: ServicePeerDirection, - ) -> OrderedSessionDemand { + ) -> SessionDemand { self.inner - .ordered_session_demand(conn_id, peer, negotiated, direction) + .session_demand(conn_id, peer, negotiated, direction) } fn wants_peer( diff --git a/crates/zakura-network/src/zakura/legacy_gossip.rs b/crates/zakura-network/src/zakura/legacy_gossip.rs index 819450773d..e178ac42ae 100644 --- a/crates/zakura-network/src/zakura/legacy_gossip.rs +++ b/crates/zakura-network/src/zakura/legacy_gossip.rs @@ -40,10 +40,10 @@ use crate::{ use super::trace::BlockBodySource; use super::{ - spawn_supervised_peer_task, BoxRunFuture, Frame, FramedSend, OrderedSendError, - OrderedSessionDemand, OrderedStreamOpening, OrderedStreamPolicy, Peer, RequestResponseService, - Service as ZakuraService, ServicePeerDirection, SinkReject, Stream, StreamMode, ZakuraConnId, - ZakuraPeerHandle, ZakuraPeerId, ZakuraSupervisorHandle, ZakuraTrace, FRAME_HEADER_BYTES, + spawn_supervised_peer_task, BoxRunFuture, Frame, FramedSend, OrderedSendError, Peer, + RequestResponseService, Service as ZakuraService, ServicePeerDirection, SessionDemand, + SessionOpening, SessionPolicy, SinkReject, Stream, StreamMode, ZakuraConnId, ZakuraPeerHandle, + ZakuraPeerId, ZakuraSupervisorHandle, ZakuraTrace, FRAME_HEADER_BYTES, LOCAL_MAX_CONTROL_FRAME_BYTES, ZAKURA_CAP_LEGACY_GOSSIP, }; @@ -145,7 +145,7 @@ const LEGACY_GOSSIP_SERVICE_STREAMS: [Stream; 2] = [ version: LEGACY_GOSSIP_VERSION, frame_cap: LOCAL_MAX_CONTROL_FRAME_BYTES, capability: ZAKURA_CAP_LEGACY_GOSSIP, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }, Stream { kind: ZAKURA_STREAM_LEGACY_REQUESTS, @@ -2544,9 +2544,9 @@ impl ZakuraService for LegacyGossipSink { legacy_gossip_streams() } - fn ordered_stream_policy(&self, _kind: u16) -> OrderedStreamPolicy { - OrderedStreamPolicy { - opening: OrderedStreamOpening::InitiatorOnly, + fn session_policy(&self) -> SessionPolicy { + SessionPolicy { + opening: SessionOpening::InitiatorOnly, reopen: true, } } @@ -2560,17 +2560,17 @@ impl ZakuraService for LegacyGossipSink { self.outbound.owns_connection(peer, conn_id) } - fn ordered_session_demand( + fn session_demand( &self, conn_id: ZakuraConnId, peer: &ZakuraPeerId, _negotiated: u64, _direction: ServicePeerDirection, - ) -> OrderedSessionDemand { + ) -> SessionDemand { if self.outbound.is_retired(peer, conn_id) { - return OrderedSessionDemand::Retire; + return SessionDemand::Retire; } - OrderedSessionDemand::OpenNow + SessionDemand::OpenNow } fn add_peer(&self, mut peer: Peer) { @@ -4700,8 +4700,8 @@ mod tests { "a retired gossip stream must not hold a reopen-gap claim" ); assert!(matches!( - sink.ordered_session_demand(conn_id, &peer_id, 0, ServicePeerDirection::Outbound), - OrderedSessionDemand::Retire, + sink.session_demand(conn_id, &peer_id, 0, ServicePeerDirection::Outbound), + SessionDemand::Retire, )); let (send, _rx) = framed_channel(1); assert!( @@ -4763,8 +4763,8 @@ mod tests { "a reset churn counter must keep the reopen-gap claim" ); assert!(matches!( - sink.ordered_session_demand(conn_id, &peer_id, 0, ServicePeerDirection::Outbound), - OrderedSessionDemand::OpenNow, + sink.session_demand(conn_id, &peer_id, 0, ServicePeerDirection::Outbound), + SessionDemand::OpenNow, )); } diff --git a/crates/zakura-network/src/zakura/testkit/cluster.rs b/crates/zakura-network/src/zakura/testkit/cluster.rs index 4c4a02b63b..d29fa2a228 100644 --- a/crates/zakura-network/src/zakura/testkit/cluster.rs +++ b/crates/zakura-network/src/zakura/testkit/cluster.rs @@ -495,7 +495,7 @@ mod tests { version: 1, frame_cap: CUSTOM_FRAME_CAP_BYTES, capability: CUSTOM_FRAME_CAP_CAPABILITY, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, }]; #[derive(Debug, Default)] diff --git a/crates/zakura-network/src/zakura/testkit/recorder.rs b/crates/zakura-network/src/zakura/testkit/recorder.rs index 1c018ee76e..bb0ce6fc88 100644 --- a/crates/zakura-network/src/zakura/testkit/recorder.rs +++ b/crates/zakura-network/src/zakura/testkit/recorder.rs @@ -141,7 +141,7 @@ impl Service for InboundRecorder { for stream in self .streams() .iter() - .filter(|stream| matches!(stream.mode, crate::zakura::StreamMode::Ordered)) + .filter(|stream| matches!(stream.mode, crate::zakura::StreamMode::Persistent)) { let Some((session_id, mut recv, _send)) = peer.take_stream_with_session_id(stream.kind) else { diff --git a/crates/zakura-network/src/zakura/transport/io.rs b/crates/zakura-network/src/zakura/transport/io.rs index 947d86c68f..5288ef657d 100644 --- a/crates/zakura-network/src/zakura/transport/io.rs +++ b/crates/zakura-network/src/zakura/transport/io.rs @@ -80,7 +80,7 @@ impl FramedRecv { #[derive(Clone, Debug)] pub struct FramedSend { sender: FramedSender, - session_resources: Option>, + session_resources: Option>, } #[derive(Clone, Debug)] @@ -108,7 +108,7 @@ impl FramedSend { /// Keep service admission charged while application senders still own the session. pub(crate) fn with_session_resources( mut self, - resources: Option>, + resources: Option>, ) -> Self { self.session_resources = resources; self diff --git a/crates/zakura-network/src/zakura/transport/mod.rs b/crates/zakura-network/src/zakura/transport/mod.rs index 0cae40d11c..1ad62e9d56 100644 --- a/crates/zakura-network/src/zakura/transport/mod.rs +++ b/crates/zakura-network/src/zakura/transport/mod.rs @@ -29,11 +29,12 @@ pub(crate) use pipe::{ handle_pipe_exit, spawn_supervised_peer_task, spawn_supervised_pipe, CloseCause, Edge, Flow, Node, NodeKind, Pipe, PipeCx, PipeShape, }; +pub(crate) use registry::SessionLayout; pub use registry::{RegistryError, ServiceRegistry}; pub(crate) use service::ServiceStream; pub use service::{ - BoxRunFuture, OrderedSessionDemand, OrderedSessionFull, OrderedSessionResources, - OrderedStreamOpening, OrderedStreamPair, OrderedStreamPolicy, Peer, RequestResponseService, - Service, Sink, SinkReject, Source, Stream, StreamMode, + BoxRunFuture, Peer, RequestResponseService, Service, SessionDemand, SessionFull, + SessionOpening, SessionPolicy, SessionResources, Sink, SinkReject, Source, Stream, StreamMode, + StreamWritePolicy, }; pub use session::{OrderedSendError, PeerStreamSession}; diff --git a/crates/zakura-network/src/zakura/transport/registry.rs b/crates/zakura-network/src/zakura/transport/registry.rs index e3f51bb124..7c6e4c6b6f 100644 --- a/crates/zakura-network/src/zakura/transport/registry.rs +++ b/crates/zakura-network/src/zakura/transport/registry.rs @@ -8,20 +8,20 @@ use std::{ use thiserror::Error; use super::{ - Frame, OrderedSessionDemand, OrderedStreamPair, OrderedStreamPolicy, Peer, Service, SinkReject, - Stream, StreamMode, + Frame, Peer, Service, SessionDemand, SessionPolicy, SinkReject, Stream, StreamMode, + StreamWritePolicy, }; use crate::zakura::{ServicePeerDirection, ZakuraConnId, ZakuraPeerId}; /// Errors returned while building a [`ServiceRegistry`]. #[derive(Debug, Error)] pub enum RegistryError { - /// A paired session has missing, mismatched, or inconsistently declared roles. - #[error("service {service} declared an invalid ordered stream pair for kind {kind}")] - InvalidOrderedPair { - /// Service declaring the pair. + /// A service declared incompatible persistent session layouts. + #[error("service {service} declared an invalid service session for kind {kind}")] + InvalidSessionLayout { + /// Service declaring the session. service: &'static str, - /// Stream with an inconsistent pair declaration. + /// Stream with an inconsistent session declaration. kind: u16, }, /// Two services declared the same stream kind. @@ -69,6 +69,23 @@ pub enum RegistryError { }, } +/// A validated set of persistent streams admitted as one service session. +#[derive(Clone, Debug, Eq, PartialEq)] +pub(crate) struct SessionLayout { + pub(crate) streams: Arc<[Stream]>, +} + +impl SessionLayout { + /// The lowest stream kind supplies the session identity and layout version. + pub(crate) fn primary(&self) -> Stream { + self.streams[0] + } + + pub(crate) fn is_multi_stream(&self) -> bool { + self.streams.len() > 1 + } +} + /// Registry of Zakura protocol services. #[derive(Clone, Debug, Default)] pub struct ServiceRegistry { @@ -76,6 +93,7 @@ pub struct ServiceRegistry { by_kind: HashMap, by_capability: HashMap>, supported_capabilities: u64, + session_layouts: HashMap<(u16, u16), SessionLayout>, } impl ServiceRegistry { @@ -84,32 +102,13 @@ impl ServiceRegistry { let mut by_kind: HashMap = HashMap::new(); let mut by_capability: HashMap> = HashMap::new(); let mut supported_capabilities = 0; + let mut session_layouts = HashMap::new(); for (index, service) in services.iter().enumerate() { let mut service_capabilities = HashSet::new(); let mut service_streams = HashSet::new(); for stream in service.streams() { - if let Some(pair) = service.ordered_stream_pair(*stream) { - let valid = pair.data != pair.requests - && pair.data.kind != pair.requests.kind - && pair.data.mode == StreamMode::Ordered - && pair.requests.mode == StreamMode::Ordered - && pair.data.capability == pair.requests.capability - && (*stream == pair.data || *stream == pair.requests) - && service.streams().contains(&pair.data) - && service.streams().contains(&pair.requests) - && service.ordered_stream_pair(pair.data) == Some(pair) - && service.ordered_stream_pair(pair.requests) == Some(pair) - && service.ordered_stream_policy(pair.data.kind) - == service.ordered_stream_policy(pair.requests.kind); - if !valid { - return Err(RegistryError::InvalidOrderedPair { - service: service.name(), - kind: stream.kind, - }); - } - } // Each stream must map to exactly one capability bit, otherwise // `supported_capabilities` (the OR below) and per-bit // `services_for_capability` lookups disagree. @@ -145,6 +144,37 @@ impl ServiceRegistry { service_capabilities.insert(stream.capability); } + let mut layouts: HashMap> = HashMap::new(); + for stream in service + .streams() + .iter() + .filter(|s| s.mode == StreamMode::Persistent) + { + layouts.entry(stream.capability).or_default().push(*stream); + } + let mut primary_kind = None; + for mut streams in layouts.into_values() { + streams.sort_unstable_by_key(|stream| stream.kind); + let primary = streams[0]; + let invalid = primary_kind.is_some_and(|kind| kind != primary.kind) + || streams + .windows(2) + .any(|roles| roles[0].kind == roles[1].kind); + if invalid { + return Err(RegistryError::InvalidSessionLayout { + service: service.name(), + kind: primary.kind, + }); + } + primary_kind = Some(primary.kind); + let layout = SessionLayout { + streams: streams.into(), + }; + for stream in layout.streams.iter() { + session_layouts.insert((stream.kind, stream.version), layout.clone()); + } + } + for capability in service_capabilities { by_capability.entry(capability).or_default().push(index); } @@ -155,6 +185,7 @@ impl ServiceRegistry { by_kind, by_capability, supported_capabilities, + session_layouts, }) } @@ -260,34 +291,50 @@ impl ServiceRegistry { self.supported_capabilities } - /// Return this exact stream version's validated pair declaration. - pub fn ordered_stream_pair(&self, stream: Stream) -> Option { - self.service_for_kind(stream.kind)? - .ordered_stream_pair(stream) + /// Return the complete persistent layout containing this exact stream. + pub(crate) fn session_layout(&self, stream: Stream) -> Option { + self.session_layouts + .get(&(stream.kind, stream.version)) + .filter(|layout| layout.streams.contains(&stream)) + .cloned() + } + + pub(crate) fn stream_write_policy(&self, stream: Stream) -> StreamWritePolicy { + self.service_for_kind(stream.kind) + .expect("a registered stream has an owning service") + .stream_write_policy(stream) } - /// Ordered streams negotiated with a peer, in registry service order. - pub fn ordered_streams_for_negotiated(&self, negotiated: u64) -> Vec { + fn selected_session_streams(&self, service: &dyn Service, negotiated: u64) -> Vec { + service + .streams() + .iter() + .filter(|stream| { + stream.mode == StreamMode::Persistent && negotiated & stream.capability != 0 + }) + .filter_map(|stream| self.session_layout(*stream)) + .max_by_key(|layout| layout.primary().version) + .map_or_else(Vec::new, |layout| layout.streams.to_vec()) + } + + /// Persistent streams negotiated with a peer, in registry service order. + pub fn persistent_streams_for_negotiated(&self, negotiated: u64) -> Vec { let mut streams = Vec::new(); for service in self.services_for_negotiated(negotiated) { - streams.extend(selected_streams( - service.as_ref(), - negotiated, - StreamMode::Ordered, - )); + streams.extend(self.selected_session_streams(service.as_ref(), negotiated)); } streams } - /// Ordered streams that should be lazily escalated for this peer now. + /// Persistent streams that should be lazily escalated for this peer now. /// /// The connection loop applies its per-kind opening policy to each returned /// stream. This demand check narrows the negotiated capabilities to services /// that currently have local interest and room; the owning reactor still /// makes the final admission decision after the typed session arrives. - pub fn ordered_streams_for_escalation( + pub fn persistent_streams_for_escalation( &self, negotiated: u64, peer_id: &ZakuraPeerId, @@ -300,18 +347,14 @@ impl ServiceRegistry { continue; } - streams.extend(selected_streams( - service.as_ref(), - negotiated, - StreamMode::Ordered, - )); + streams.extend(self.selected_session_streams(service.as_ref(), negotiated)); } streams } /// Return true when the service owning `kind` still wants this peer. - pub fn wants_ordered_stream( + pub fn wants_session( &self, kind: u16, negotiated: u64, @@ -326,41 +369,41 @@ impl ServiceRegistry { } /// Return the owning service's static ordered-stream policy. - pub fn ordered_stream_policy(&self, kind: u16) -> OrderedStreamPolicy { + pub fn session_policy(&self, kind: u16) -> SessionPolicy { self.service_for_kind(kind) - .map(|service| service.ordered_stream_policy(kind)) + .map(|service| service.session_policy()) .unwrap_or_default() } /// Return the owning service's current demand for an absent ordered session. - pub fn ordered_session_demand( + pub fn session_demand( &self, kind: u16, conn_id: ZakuraConnId, negotiated: u64, peer_id: &ZakuraPeerId, direction: ServicePeerDirection, - ) -> OrderedSessionDemand { + ) -> SessionDemand { let Some(service) = self.service_for_kind(kind) else { - return OrderedSessionDemand::Retire; + return SessionDemand::Retire; }; - service.ordered_session_demand(conn_id, peer_id, negotiated, direction) + service.session_demand(conn_id, peer_id, negotiated, direction) } - /// Recheck demand after a complete pair has reserved its service capacity. - pub(crate) fn reserved_ordered_session_demand( + /// Recheck demand after a complete session has reserved its service capacity. + pub(crate) fn reserved_session_demand( &self, kind: u16, conn_id: ZakuraConnId, negotiated: u64, peer_id: &ZakuraPeerId, direction: ServicePeerDirection, - ) -> OrderedSessionDemand { + ) -> SessionDemand { let Some(service) = self.service_for_kind(kind) else { - return OrderedSessionDemand::Retire; + return SessionDemand::Retire; }; - service.reserved_ordered_session_demand(conn_id, peer_id, negotiated, direction) + service.reserved_session_demand(conn_id, peer_id, negotiated, direction) } /// Request/response streams negotiated with a peer, in registry service order. @@ -555,22 +598,7 @@ fn selected_streams(service: &dyn Service, negotiated: u64, mode: StreamMode) -> selected.push(stream); } } - // TODO: Choose complete pair alternatives before selecting per-kind versions - // when adding a second paired protocol version. For example, data/request - // alternatives 3/1 and 2/4 select 3/4, and this filter rejects both roles. - // Block sync's paired layout has only one alternative, and peers can only - // negotiate locally registered versions, so this case is unreachable for - // the current layout. Defer the selection change until multiple pair - // alternatives are needed, and implement it before registering them. selected - .iter() - .copied() - .filter(|stream| { - service.ordered_stream_pair(*stream).is_none_or(|pair| { - selected.contains(&pair.data) && selected.contains(&pair.requests) - }) - }) - .collect() } #[cfg(test)] @@ -590,7 +618,6 @@ mod tests { added: Mutex>, added_streams: Mutex>>, removed: Mutex>, - pairs: Vec, } impl TestService { @@ -602,7 +629,6 @@ mod tests { added: Mutex::new(Vec::new()), added_streams: Mutex::new(Vec::new()), removed: Mutex::new(Vec::new()), - pairs: Vec::new(), }) } @@ -623,13 +649,6 @@ mod tests { &self.streams } - fn ordered_stream_pair(&self, stream: Stream) -> Option { - self.pairs - .iter() - .copied() - .find(|pair| pair.data == stream || pair.requests == stream) - } - fn wants_peer( &self, _peer: &ZakuraPeerId, @@ -692,7 +711,7 @@ mod tests { version: 1, frame_cap: 1024, capability, - mode: StreamMode::Ordered, + mode: StreamMode::Persistent, } } @@ -718,67 +737,44 @@ mod tests { assert!(registry.is_supported_stream(5, 7)); assert!(registry.is_supported_stream(5, 8)); assert_eq!( - registry.ordered_streams_for_negotiated(0b0001), + registry.persistent_streams_for_negotiated(0b0001), vec![versioned_stream(5, 7, 0b0001)] ); assert_eq!( - registry.ordered_streams_for_negotiated(0b0011), + registry.persistent_streams_for_negotiated(0b0011), vec![versioned_stream(5, 8, 0b0010)] ); } #[test] - fn paired_versions_are_selected_together_and_cannot_leave_an_orphan_role() { - let legacy = versioned_stream(6, 2, 1); - let pair = OrderedStreamPair { - data: versioned_stream(6, 3, 2), - requests: versioned_stream(7, 1, 2), - }; - let later = versioned_stream(6, 4, 4); - let mut service = TestService::new("pair", vec![legacy, pair.data, pair.requests, later]); - Arc::get_mut(&mut service).unwrap().pairs.push(pair); - let registry = ServiceRegistry::new(vec![service]).unwrap(); - assert_eq!(registry.ordered_streams_for_negotiated(1), vec![legacy]); - assert_eq!( - registry.ordered_streams_for_negotiated(3), - vec![pair.data, pair.requests] + fn session_versions_are_selected_as_complete_layouts() { + let legacy = versioned_stream(6, 1, 1); + let older = [versioned_stream(6, 2, 2), versioned_stream(7, 4, 2)]; + let newer = [ + versioned_stream(6, 3, 4), + versioned_stream(7, 1, 4), + versioned_stream(8, 1, 4), + ]; + let service = TestService::new( + "session", + [vec![legacy], older.to_vec(), newer.to_vec()].concat(), ); - assert_eq!(registry.ordered_streams_for_negotiated(7), vec![later]); + let registry = ServiceRegistry::new(vec![service]).unwrap(); + assert_eq!(registry.persistent_streams_for_negotiated(1), vec![legacy]); + assert_eq!(registry.persistent_streams_for_negotiated(3), older); + // Per-kind selection would incorrectly choose stream 7 version 4. + assert_eq!(registry.persistent_streams_for_negotiated(7), newer); } #[test] - fn invalid_pair_declarations_are_rejected_before_handshake() { - let data = versioned_stream(6, 3, 2); - let requests = versioned_stream(7, 1, 2); - for pair in [ - OrderedStreamPair { - data, - requests: data, - }, - OrderedStreamPair { - data, - requests: Stream { - capability: 4, - ..requests - }, - }, - OrderedStreamPair { - data, - requests: Stream { - mode: StreamMode::RequestResponse, - ..requests - }, - }, - OrderedStreamPair { - data, - requests: versioned_stream(8, 1, 2), - }, + fn session_layouts_require_one_version_per_kind_and_a_stable_primary_kind() { + for streams in [ + vec![versioned_stream(6, 1, 1), versioned_stream(6, 2, 1)], + vec![versioned_stream(6, 1, 1), versioned_stream(7, 1, 2)], ] { - let mut service = TestService::new("pair", vec![data, requests]); - Arc::get_mut(&mut service).unwrap().pairs.push(pair); assert!(matches!( - ServiceRegistry::new(vec![service]), - Err(RegistryError::InvalidOrderedPair { .. }) + ServiceRegistry::new(vec![TestService::new("invalid", streams)]), + Err(RegistryError::InvalidSessionLayout { .. }) )); } } @@ -805,7 +801,10 @@ mod tests { #[test] fn registry_builds_kind_and_capability_lookups() { - let header = TestService::new("header", vec![stream(5, 0b0001), stream(6, 0b0010)]); + let header = TestService::new( + "header", + vec![stream(5, 0b0001), versioned_stream(5, 2, 0b0010)], + ); let gossip = TestService::new("gossip", vec![stream(2, 0b0100)]); let registry = ServiceRegistry::new(vec![header.clone(), gossip.clone()]) @@ -886,7 +885,10 @@ mod tests { #[test] fn supported_capabilities_are_or_of_declared_streams() { - let header = TestService::new("header", vec![stream(5, 0b0001), stream(6, 0b0010)]); + let header = TestService::new( + "header", + vec![stream(5, 0b0001), versioned_stream(5, 2, 0b0010)], + ); let gossip = TestService::new("gossip", vec![stream(2, 0b0100)]); let registry = ServiceRegistry::new(vec![header, gossip]) @@ -897,7 +899,10 @@ mod tests { #[test] fn services_for_negotiated_matches_any_bit_once_in_registration_order() { - let header = TestService::new("header", vec![stream(5, 0b0001), stream(6, 0b0010)]); + let header = TestService::new( + "header", + vec![stream(5, 0b0001), versioned_stream(5, 2, 0b0010)], + ); let gossip = TestService::new("gossip", vec![stream(2, 0b0100)]); let discovery = TestService::new("discovery", vec![stream(4, 0b1000)]); @@ -920,7 +925,7 @@ mod tests { "multi-capability", vec![ stream(5, 0b0001), - stream(6, 0b0010), + versioned_stream(5, 2, 0b0010), request_response_one, request_response_two, ], @@ -930,12 +935,12 @@ mod tests { let peer = ZakuraPeerId::new(vec![8; 32]).expect("32-byte test peer id is valid"); let ordered_kinds: Vec<_> = registry - .ordered_streams_for_negotiated(0b0001) + .persistent_streams_for_negotiated(0b0001) .iter() .map(|stream| stream.kind) .collect(); let escalated_kinds: Vec<_> = registry - .ordered_streams_for_escalation(0b0001, &peer, ServicePeerDirection::Outbound) + .persistent_streams_for_escalation(0b0001, &peer, ServicePeerDirection::Outbound) .iter() .map(|stream| stream.kind) .collect(); @@ -1058,7 +1063,7 @@ mod tests { } #[test] - fn ordered_streams_for_escalation_filters_services_without_demand() { + fn persistent_streams_for_escalation_filters_services_without_demand() { let header = TestService::new("header", vec![stream(5, 0b0001)]); let discovery = TestService::new("discovery", vec![stream(4, 0b0010)]); let registry = ServiceRegistry::new(vec![header.clone(), discovery.clone()]) @@ -1067,8 +1072,11 @@ mod tests { header.set_wants(false); - let streams = - registry.ordered_streams_for_escalation(0b0011, &peer, ServicePeerDirection::Outbound); + let streams = registry.persistent_streams_for_escalation( + 0b0011, + &peer, + ServicePeerDirection::Outbound, + ); let stream_kinds: Vec<_> = streams.iter().map(|stream| stream.kind).collect(); assert_eq!(stream_kinds, [4]); diff --git a/crates/zakura-network/src/zakura/transport/service.rs b/crates/zakura-network/src/zakura/transport/service.rs index b44cc1ae71..f828a642a7 100644 --- a/crates/zakura-network/src/zakura/transport/service.rs +++ b/crates/zakura-network/src/zakura/transport/service.rs @@ -1,6 +1,13 @@ //! Zakura protocol service trait surface. -use std::{collections::HashMap, fmt, future::Future, net::IpAddr, pin::Pin, time::Instant}; +use std::{ + collections::HashMap, + fmt, + future::Future, + net::IpAddr, + pin::Pin, + time::{Duration, Instant}, +}; use thiserror::Error; use tokio_util::sync::CancellationToken; @@ -20,41 +27,41 @@ pub type BoxRunFuture<'a, T> = Pin + Send + 'a>>; #[derive(Copy, Clone, Debug, Eq, PartialEq)] pub enum StreamMode { /// A long-lived ordered stream between connected peers. - Ordered, + Persistent, /// A short-lived request/response stream opened per request. RequestResponse, } -/// Which endpoint may proactively open an ordered service stream. +/// Which endpoint may proactively open a service session. #[derive(Copy, Clone, Debug, Eq, PartialEq)] -pub enum OrderedStreamOpening { +pub enum SessionOpening { /// Only the endpoint that initiated the authenticated connection opens the stream. InitiatorOnly, /// Either endpoint may open the stream; simultaneous opens use the transport tiebreak. EitherSide, } -/// Static transport policy for one ordered service stream. +/// Static transport policy for one persistent service session. #[derive(Copy, Clone, Debug, Eq, PartialEq)] -pub struct OrderedStreamPolicy { +pub struct SessionPolicy { /// Which endpoint may proactively open the stream. - pub opening: OrderedStreamOpening, + pub opening: SessionOpening, /// Whether a locally ended session may be re-admitted on the same connection. pub reopen: bool, } -impl Default for OrderedStreamPolicy { +impl Default for SessionPolicy { fn default() -> Self { Self { - opening: OrderedStreamOpening::InitiatorOnly, + opening: SessionOpening::InitiatorOnly, reopen: false, } } } -/// A service's current decision for an absent ordered session. -pub enum OrderedSessionDemand { - /// Open and admit the ordered stream now. +/// A service's current decision for an absent service session. +pub enum SessionDemand { + /// Open and admit the complete session now. OpenNow, /// Re-check demand at this instant. RetryAt(Instant), @@ -64,7 +71,7 @@ pub enum OrderedSessionDemand { Retire, } -impl fmt::Debug for OrderedSessionDemand { +impl fmt::Debug for SessionDemand { fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { match self { Self::OpenNow => formatter.write_str("OpenNow"), @@ -90,34 +97,29 @@ pub struct Stream { pub mode: StreamMode, } -/// Two persistent ordered streams selected and retired as one service session. -/// -/// Both roles belong to the same service and capability. The data stream is the -/// session's transport identity; the request stream shares its cancellation and -/// message budget. Each role carries the same nonzero eight-byte pair identifier -/// immediately after its ordinary prelude, scoped to its connection and opener. +/// A persistent stream's write deadline within its service session. #[derive(Copy, Clone, Debug, Eq, PartialEq)] -pub struct OrderedStreamPair { - /// Stream carrying responses and control messages, with bounded writes. - pub data: Stream, - /// Stream whose writes may wait while the session remains valid. - pub requests: Stream, +pub enum StreamWritePolicy { + /// Retire the session if a complete frame cannot be written within this time. + Timeout(Duration), + /// Let the service's progress policy or session cancellation end the wait. + UntilCancelled, } -/// A service slot held from paired-stream setup through the last transport and +/// A service slot held from session setup through the last transport and /// application sender owner. The service can release its setup allowance once -/// both roles are ready, while retaining its session allowance through teardown. -pub trait OrderedSessionResources: fmt::Debug + Send + Sync { - /// Both roles have completed setup. Called once for a successfully built pair. +/// all members are ready, while retaining its session allowance through teardown. +pub trait SessionResources: fmt::Debug + Send + Sync { + /// All required streams have completed setup. Called once per session. fn admitted(&self); } /// The service has no capacity for another establishing or retiring session. #[derive(Debug, Error)] -#[error("ordered service session capacity is full")] -pub struct OrderedSessionFull; +#[error("service session capacity is full")] +pub struct SessionFull; -/// Transport state for one ordered service stream. +/// Transport state for one persistent stream within a service session. #[derive(Debug)] pub(crate) struct ServiceStream { pub(crate) session_id: u64, @@ -312,14 +314,14 @@ impl Peer { } } - /// Take ownership of a stream pair for `kind`. + /// Take ownership of a stream's receive and send handles for `kind`. pub fn take_stream(&mut self, kind: u16) -> Option<(FramedRecv, FramedSend)> { self.streams .remove(&kind) .map(|stream| (stream.recv, stream.send)) } - /// Take ownership of a stream pair and its owning ordered-stream generation. + /// Take ownership of a stream's receive and send handles and its owning session identity. pub fn take_stream_with_session_id( &mut self, kind: u16, @@ -329,7 +331,7 @@ impl Peer { .map(|stream| (stream.session_id, stream.recv, stream.send)) } - /// Take ownership of a stream pair, its version, and ordered-stream generation. + /// Take ownership of a stream's receive and send handles, its version, and session identity. pub fn take_versioned_stream_with_session_id( &mut self, kind: u16, @@ -391,7 +393,16 @@ pub trait Service: fmt::Debug + Send + Sync + 'static { /// Stable service name for logs and diagnostics. fn name(&self) -> &'static str; - /// Streams this service owns. + /// Stream types this service owns. + /// + /// Persistent streams with the same capability form one complete session + /// layout. The transport admits every member together and retires the + /// session when any member ends. Request/response streams open per request + /// and do not participate in session setup. + /// + /// Alternative layouts use distinct capabilities and the same lowest stream + /// kind. That stream's version ranks complete layouts during negotiation. + /// Advance its version whenever the session layout changes. fn streams(&self) -> &[Stream]; /// Payload size limits for this stream, as `(message_type, maximum_bytes)` pairs. @@ -415,63 +426,63 @@ pub trait Service: fmt::Debug + Send + Sync + 'static { None } - /// Reserve service capacity before starting a stream pair. The returned - /// owner lives through incomplete setup and both workers' eventual teardown. - fn reserve_ordered_session( + /// Reserve service capacity before starting a persistent session. The returned + /// owner lives through incomplete setup and all workers' eventual teardown. + fn reserve_session( &self, _direction: ServicePeerDirection, - ) -> Result>, OrderedSessionFull> { + ) -> Result>, SessionFull> { Ok(None) } - /// Return the complete pair containing `stream`, if this version uses one. + /// Control writes independently for each persistent stream. /// - /// Both declarations must be present in [`Service::streams`] with the same - /// capability and opening policy. Ordinary ordered streams return `None`. - fn ordered_stream_pair(&self, _stream: Stream) -> Option { - None + /// Services that allow indefinite writes must enforce their own bounded + /// progress policy. Cancelling the session always interrupts a pending write. + fn stream_write_policy(&self, _stream: Stream) -> StreamWritePolicy { + StreamWritePolicy::Timeout(Duration::from_secs(10)) } - /// Return the transport-owned opening and re-admission policy for `kind`. + /// Return the opening and re-admission policy for the whole service session. /// /// The default preserves the legacy one-shot initiator-opens behavior. - fn ordered_stream_policy(&self, _kind: u16) -> OrderedStreamPolicy { - OrderedStreamPolicy::default() + fn session_policy(&self) -> SessionPolicy { + SessionPolicy::default() } - /// Return this service's current demand for an absent ordered session. + /// Return this service's current demand for an absent service session. /// - /// Services that opt into [`OrderedStreamPolicy::reopen`] should override + /// Services that opt into [`SessionPolicy::reopen`] should override /// this method so local cooldowns, capacity, and usefulness remain - /// reactor-owned. [`OrderedSessionDemand::WaitForChange`] avoids periodic + /// reactor-owned. [`SessionDemand::WaitForChange`] avoids periodic /// transport polling while the service is full or has no useful work. - fn ordered_session_demand( + fn session_demand( &self, _conn_id: ZakuraConnId, peer: &ZakuraPeerId, negotiated: u64, direction: ServicePeerDirection, - ) -> OrderedSessionDemand { + ) -> SessionDemand { if self.wants_peer(peer, negotiated, direction) { - OrderedSessionDemand::OpenNow + SessionDemand::OpenNow } else { - OrderedSessionDemand::Retire + SessionDemand::Retire } } - /// Recheck demand for a complete pair that already owns its setup reservation. + /// Recheck demand for a complete session that already owns its setup reservation. /// /// Services with session reservations must retain cooldown and usefulness /// checks here without requiring capacity for a second reservation. The /// default preserves ordinary demand checks for services without reservations. - fn reserved_ordered_session_demand( + fn reserved_session_demand( &self, conn_id: ZakuraConnId, peer: &ZakuraPeerId, negotiated: u64, direction: ServicePeerDirection, - ) -> OrderedSessionDemand { - self.ordered_session_demand(conn_id, peer, negotiated, direction) + ) -> SessionDemand { + self.session_demand(conn_id, peer, negotiated, direction) } /// Return whether this service currently wants a new session for `peer`. diff --git a/docs/changelog/unreleased/956.md b/docs/changelog/unreleased/956.md new file mode 100644 index 0000000000..0d17e912c2 --- /dev/null +++ b/docs/changelog/unreleased/956.md @@ -0,0 +1,5 @@ +## Changed + +- Retire only the affected native P2P service session when a persistent stream + write times out, preserving unrelated services on the connection + ([#956](https://github.com/zakura-core/zakura/pull/956)). diff --git a/docs/design/service-sessions.md b/docs/design/service-sessions.md new file mode 100644 index 0000000000..1315abef29 --- /dev/null +++ b/docs/design/service-sessions.md @@ -0,0 +1,122 @@ +# Service sessions + +A service declares its stream types through `Service::streams()`. The transport +groups persistent streams with the same capability into one session layout. +The transport admits the complete layout through one `Service::add_peer()` call. +Each stream has independent readers, writers, and queues. + +Request/response stream types remain in the declaration, but the transport opens +their streams per request. They do not participate in persistent session setup. +The service controls their application lifecycle. + +## Declaring a session + +For example, a service can declare data, requests, and events as three persistent +streams with the same capability: + +```rust +const DATA: Stream = Stream { + kind: 64, + version: 1, + frame_cap: 1024 * 1024, + capability: 1 << 16, + mode: StreamMode::Persistent, +}; +const REQUESTS: Stream = Stream { kind: 65, ..DATA }; +const EVENTS: Stream = Stream { kind: 66, ..DATA }; +const LOOKUP: Stream = Stream { + kind: 67, + mode: StreamMode::RequestResponse, + ..DATA +}; + +// Inside impl Service: +fn streams(&self) -> &[Stream] { + &[DATA, REQUESTS, EVENTS, LOOKUP] +} + +fn session_policy(&self) -> SessionPolicy { + SessionPolicy { + opening: SessionOpening::EitherSide, + reopen: true, + } +} +``` + +These identifiers illustrate the API. A production protocol must allocate its +own stream kinds and capability bit. + +The service does not declare membership a second time. The transport waits for +data, requests, and events before handing their receive/send handles to the +service. It does not wait for a lookup request. + +The service uses `message_types()`, `message_payload_limits()`, and +`stream_queue_depths()` to specify each stream's traffic and bounds. The protocol +defines message assignments; peers do not negotiate individual message types. +The service routes outgoing messages to the appropriate sender. + +`stream_write_policy()` sets each persistent stream's write deadline. The default +is ten seconds. A service can choose another duration or `UntilCancelled`. +A service that chooses `UntilCancelled` must enforce its own progress deadline. + +## Negotiating complete layouts + +Each capability identifies a complete persistent layout for its service. +Alternative layouts use different capabilities. Each alternative retains the +same lowest stream kind, whose version ranks the layouts. The registry selects +the highest mutually supported version of that primary stream and includes every +member of its layout. It never mixes members from different alternatives. + +For example, primary/request versions `3/1` and `2/4` select `3/1` when both +capabilities are available. The request stream's higher version in the older +layout does not override that choice. + +Adding a required stream changes the protocol layout. Allocate a new capability +and advance the primary stream's version. Keep the primary kind stable. +Changing a message assignment also requires a compatible protocol transition. +The transport cannot make an old peer understand a new layout or message. + +## Setup and retirement + +Multi-stream sessions append the same nonzero eight-byte session identifier to +each ordinary stream prelude. This retains #943's two-stream setup encoding. +Single-stream sessions retain their existing prelude without an extra identifier. + +The transport holds at most one incomplete session per service and connection. +The first complete member starts the setup deadline. Later members cannot extend +that deadline. Duplicate members, mismatched identifiers, and invalid declarations +cannot complete a session. Expiry releases every arrived member and defers another +offer through the existing cooldown. + +`reserve_session()` charges service capacity once during setup. +`SessionResources::admitted()` signals complete setup. +The workers and application senders retain the shared resource owner until they +finish or drop it. Every member also consumes a transport stream slot. + +Every persistent member shares a local session identity, cancellation token, and +message-rate budget. Closing any member retires the session. Cancellation resets +unfinished writes before a replacement can send frames. A write deadline retires +the session without closing unrelated services on the connection. Protocol +violations can still close the connection. + +The transport reports session exit after every worker and reader finishes. +Reopening follows the service's policy and demand. Ephemeral request completion +does not cancel the persistent session. + +Setup readiness does not impose ordering across streams. For example, a request +can arrive before a status message on another stream. The service must handle +that ordering or perform an application handshake before processing requests. + +## Migrating a pair consumer + +Remove `OrderedStreamPair` and `ordered_stream_pair()`. Declare the persistent +members with one capability in `streams()`. Use the `Session*` policy, demand, +and resource APIs. The transport supplies all declared members together. + +Move role-specific queue limits and write deadlines into the service hooks. +For the block-sync activation following #943, the service must declare the +one-slot request queue, the request write policy, and the 32-second data write +deadline. The transport no longer assigns those policies by role name. + +Production block sync in #943 remains a single-stream protocol. This change does +not activate the later block-sync layout.