diff --git a/crates/bin/src/main.rs b/crates/bin/src/main.rs index 115426f6..8468c94d 100644 --- a/crates/bin/src/main.rs +++ b/crates/bin/src/main.rs @@ -217,6 +217,7 @@ fn main() -> Result<(), Box> { ); let identify = config.identify()?; let p2p_context = Context { + data_columns_consumer: None, gossip_producer: incoming_gossip_producer, gossip_consumer: outgoing_gossip_producer .cache_ref() diff --git a/crates/common/src/spine/messages.rs b/crates/common/src/spine/messages.rs index 5832dfad..10f6c240 100644 --- a/crates/common/src/spine/messages.rs +++ b/crates/common/src/spine/messages.rs @@ -8,8 +8,8 @@ use flux::timing::Nanos; use silver_beacon_state_data::SLOTS_PER_EPOCH; use crate::{ - DataKind, Enr, GossipTopic, Identify, MessageId, Origin, P2pStreamId, PeerId, StreamProtocol, - TCacheProducer, TCacheRead, TMultiProducer, + CacheFrameRef, DataKind, Enr, GossipTopic, Identify, MessageId, Origin, P2pStreamId, PeerId, + StreamProtocol, TCacheProducer, TCacheRead, TMultiProducer, column_util::columns_of, ssz_view::{ BLOCKS_BY_RANGE_REQ_SIZE, DC_BY_RANGE_REQ_MAX, @@ -711,6 +711,7 @@ pub enum RpcSeverity { #[allow(clippy::large_enum_variant)] pub enum P2pSend { Gossip(GossipMsgOut), + SegmentedGossip { peer_id: usize, frame: CacheFrameRef }, Identify(usize), Rpc(RpcOutbound), } @@ -719,6 +720,7 @@ impl P2pSend { pub fn peer_id(&self) -> usize { match self { P2pSend::Gossip(gossip_msg_out) => gossip_msg_out.peer_id, + P2pSend::SegmentedGossip { peer_id, .. } => *peer_id, P2pSend::Identify(peer) => *peer, P2pSend::Rpc(rpc_outbound) => rpc_outbound.peer_id(), } @@ -726,7 +728,7 @@ impl P2pSend { pub fn protocol(&self) -> StreamProtocol { match self { - P2pSend::Gossip(_) => StreamProtocol::GossipSub, + P2pSend::Gossip(_) | P2pSend::SegmentedGossip { .. } => StreamProtocol::GossipSub, P2pSend::Identify(_) => StreamProtocol::Identity, P2pSend::Rpc(rpc_outbound) => rpc_outbound.protocol(), } diff --git a/crates/e2e/src/stack.rs b/crates/e2e/src/stack.rs index 5844fc44..ae38028b 100644 --- a/crates/e2e/src/stack.rs +++ b/crates/e2e/src/stack.rs @@ -241,6 +241,7 @@ impl PublisherStack { .expect("cluster outbound random access"); let context = Context { + data_columns_consumer: None, gossip_producer: gossip_in_producer, gossip_consumer: gossip_out_ra_for_network, rpc_producer: rpc_in_producer, @@ -373,6 +374,7 @@ impl EchoStack { .expect("cluster outbound random access"); let context = Context { + data_columns_consumer: None, gossip_producer: gossip_in_producer, gossip_consumer: protobuf_ra_for_network, rpc_producer: rpc_in_producer, diff --git a/crates/network/benches/quic_basic.rs b/crates/network/benches/quic_basic.rs index 23f5d0eb..be6207ec 100644 --- a/crates/network/benches/quic_basic.rs +++ b/crates/network/benches/quic_basic.rs @@ -67,6 +67,7 @@ pub fn broadcast(c: &mut Criterion) { ); let context = Context { + data_columns_consumer: None, gossip_producer: gi_producer, gossip_consumer: go_consumer, rpc_producer: rpc_in, @@ -133,6 +134,7 @@ pub fn broadcast(c: &mut Criterion) { cluster_in.cache_ref().random_access("cluster_out", true).unwrap(); let context = Context { + data_columns_consumer: None, gossip_producer: gi_producer, gossip_consumer: go_consumer, rpc_producer: rpc_in, diff --git a/crates/network/benches/quic_pingpong.rs b/crates/network/benches/quic_pingpong.rs index 3309161e..79433bac 100644 --- a/crates/network/benches/quic_pingpong.rs +++ b/crates/network/benches/quic_pingpong.rs @@ -66,6 +66,7 @@ pub fn broadcast(c: &mut Criterion) { let p2p = P2p::new(keypair, server_endpoint, 1024, Default::default()); let context = Context { + data_columns_consumer: None, gossip_producer: gi_producer, gossip_consumer: go_consumer, rpc_producer: rpc_in, @@ -139,6 +140,7 @@ pub fn broadcast(c: &mut Criterion) { cluster_in.cache_ref().random_access("cluster_out", true).unwrap(); let context = Context { + data_columns_consumer: None, gossip_producer: gi_producer, gossip_consumer: go_consumer, rpc_producer: rpc_in, diff --git a/crates/network/src/lib.rs b/crates/network/src/lib.rs index 18d45cc3..076fe04b 100644 --- a/crates/network/src/lib.rs +++ b/crates/network/src/lib.rs @@ -41,6 +41,16 @@ silver_common::declare_counters! { GossipStallDisconnect, // RPC codecs currently retained in the network-tile-wide free list. RpcCodecPoolIdle, + CacheSegmentedAdmitted, + CacheSegmentedRejected, + CacheSegmentedCapacity, + CacheSegmentedSegments, + CacheSegmentedOwnerAllocations, + CacheSegmentedFrames, + // Reserved owner slots, including owners already held by Quinn. + CacheSegmentedOwners, + // Reserved ranges per recipient, including descriptors; not unique cache backing bytes. + CacheSegmentedRetainedBytes, } } diff --git a/crates/network/src/p2p/context.rs b/crates/network/src/p2p/context.rs index 46096f1c..61e50573 100644 --- a/crates/network/src/p2p/context.rs +++ b/crates/network/src/p2p/context.rs @@ -15,6 +15,7 @@ pub struct Context { pub cluster_nodes: Option, pub cluster_inbound_producer: TProducer, pub cluster_outbound_consumer: TRandomAccess, + pub data_columns_consumer: Option>, } impl Context { diff --git a/crates/network/src/p2p/mod.rs b/crates/network/src/p2p/mod.rs index 1060261b..0780100e 100644 --- a/crates/network/src/p2p/mod.rs +++ b/crates/network/src/p2p/mod.rs @@ -13,12 +13,13 @@ use buffa::{Message, MessageView}; pub use context::{ClusterNodes, Context}; use fxhash::{FxHashMap, FxHashSet}; use mio::{Poll, net::UdpSocket}; +use quic::SegmentedGossipLimits; pub(crate) use quic::{Peer, create_client_config}; pub use quic::{SendResult, create_endpoint, create_server_config}; use quinn_proto::{ConnectionHandle, DatagramEvent, Endpoint}; use silver_common::{ - ClusterMsgOut, GossipMsgOut, Identify, Keypair, P2pConnectionStats, P2pStreamId, PeerId, - ProtoIdentify, ProtoIdentifyView, RpcOutbound, RpcRequestOutbound, TCacheRead, + CacheFrameRef, ClusterMsgOut, GossipMsgOut, Identify, Keypair, P2pConnectionStats, P2pStreamId, + PeerId, ProtoIdentify, ProtoIdentifyView, RpcOutbound, RpcRequestOutbound, TCacheRead, }; use crate::{ @@ -112,6 +113,8 @@ pub struct P2p { keypair: Keypair, endpoint: Endpoint, peers: FxHashMap, + // Peer queues and Quinn owners must drop before their shared budgets. + segmented_limits: Option>, rpc_codec_pool: RpcCodecPool, banned: FxHashSet, timeout: Option, @@ -134,6 +137,7 @@ impl P2p { keypair, endpoint, peers: FxHashMap::default(), + segmented_limits: None, rpc_codec_pool: RpcCodecPool::default(), banned: FxHashSet::default(), timeout: Some(Duration::ZERO), @@ -358,6 +362,9 @@ impl P2p { }; NetworkCounters::P2pConnections.set(self.peers.len() as u64); + if let Some(limits) = &self.segmented_limits { + limits.publish_gauges(); + } did_work } @@ -371,6 +378,23 @@ impl P2p { } } + pub fn enqueue_segmented_gossip( + &mut self, + peer_id: usize, + frame: CacheFrameRef, + context: &mut Context, + ) -> SendResult { + match self.peers.get_mut(&ConnectionHandle(peer_id)) { + Some(peer) => peer.send_segmented_gossip( + frame, + context, + self.segmented_limits.get_or_insert_with(Box::default), + &mut self.rpc_codec_pool, + ), + None => SendResult::UnknownPeer, + } + } + pub fn enqueue_rpc_out(&mut self, msg: RpcOutbound, context: &mut Context) -> SendResult { match self.peers.get_mut(&ConnectionHandle(msg.peer_id())) { Some(peer) => { diff --git a/crates/network/src/p2p/quic/gossip_frame.rs b/crates/network/src/p2p/quic/gossip_frame.rs new file mode 100644 index 00000000..8ee15335 --- /dev/null +++ b/crates/network/src/p2p/quic/gossip_frame.rs @@ -0,0 +1,195 @@ +use std::{cell::Cell, ptr::NonNull, time::Instant}; + +use bytes::Bytes; +use silver_common::{ + AcquiredCacheFrame, AcquiredCacheSegment, AcquiredRange, CacheFrameView, TRead, +}; + +use super::{Leased, leased::OutboundLeaseWheel}; +use crate::{NetworkCounters, p2p::Context}; + +const MAX_RETAINED_BYTES: usize = 64 * 1024 * 1024; +const MAX_RETAINED_OWNERS: usize = 8 * 1024; + +#[derive(Debug)] +pub(crate) enum OutboundGossip { + Contiguous(Leased), + Segmented(SegmentedFrame), +} + +pub(crate) struct SegmentedGossipLimits { + frames: Cell, + owners: Cell, + retained_bytes: Cell, + max_frames: usize, + max_bytes: usize, + max_owners: usize, +} + +impl Default for SegmentedGossipLimits { + fn default() -> Self { + Self::new(128) + } +} + +impl SegmentedGossipLimits { + pub(crate) fn new(max_frames: usize) -> Self { + Self { + frames: Cell::new(0), + owners: Cell::new(0), + retained_bytes: Cell::new(0), + max_frames, + max_bytes: MAX_RETAINED_BYTES, + max_owners: MAX_RETAINED_OWNERS, + } + } + + pub(crate) fn acquire( + &self, + view: CacheFrameView, + context: &mut Context, + wheel: &OutboundLeaseWheel, + now: Instant, + ) -> Option { + let bytes = view.descriptor_len().checked_add(view.wire_len())?; + let owners = view.segment_count() + 1; + if bytes > self.max_bytes.saturating_sub(self.retained_bytes.get()) || + owners > self.max_owners.saturating_sub(self.owners.get()) || + self.frames.get() >= self.max_frames + { + NetworkCounters::CacheSegmentedCapacity.inc(); + return None; + } + self.frames.set(self.frames.get() + 1); + self.owners.set(self.owners.get() + owners); + self.retained_bytes.set(self.retained_bytes.get() + bytes); + let budget = FrameBudget { limits: NonNull::from(self), owners, bytes }; + let frame = view.acquire_segments( + &mut context.gossip_consumer, + context.data_columns_consumer.as_deref_mut(), + )?; + NetworkCounters::CacheSegmentedAdmitted.inc(); + NetworkCounters::CacheSegmentedSegments.add(frame.segment_count() as u64); + Some(SegmentedFrame { segments: wheel.leased(frame, now), budget }) + } + + pub(crate) fn publish_gauges(&self) { + NetworkCounters::CacheSegmentedFrames.set(self.frames.get() as u64); + NetworkCounters::CacheSegmentedOwners.set(self.owners.get() as u64); + NetworkCounters::CacheSegmentedRetainedBytes.set(self.retained_bytes.get() as u64); + } +} + +impl Drop for SegmentedGossipLimits { + fn drop(&mut self) { + debug_assert_eq!(self.frames.get(), 0, "limits dropped with active frames"); + debug_assert_eq!(self.owners.get(), 0, "limits dropped with active owners"); + debug_assert_eq!(self.retained_bytes.get(), 0, "limits dropped with retained bytes"); + } +} + +#[derive(Debug)] +pub(crate) struct SegmentedFrame { + segments: Leased, + budget: FrameBudget, +} + +impl SegmentedFrame { + pub(crate) fn wire_len(&self) -> usize { + self.segments.wire_len() + } + + pub(crate) fn into_writer(self) -> SegmentedWriter { + let remaining = self.wire_len(); + SegmentedWriter { frame: self, current: Bytes::new(), descriptor: Bytes::new(), remaining } + } +} + +#[derive(Debug)] +pub(crate) struct SegmentedWriter { + frame: SegmentedFrame, + current: Bytes, + descriptor: Bytes, + remaining: usize, +} + +impl SegmentedWriter { + pub(crate) fn chunk(&mut self) -> Option<&mut Bytes> { + if self.current.is_empty() { + let SegmentedFrame { segments, budget } = &mut self.frame; + self.current = match segments.take_next()? { + AcquiredCacheSegment::Framing(range) => { + if self.descriptor.is_empty() { + self.descriptor = budget.owner(segments.child(segments.descriptor_range())); + } + self.descriptor.slice(range) + } + AcquiredCacheSegment::Data(range) => budget.owner(segments.child(range)), + }; + } + Some(&mut self.current) + } + + pub(crate) fn written(&mut self, bytes: usize) -> bool { + assert!(bytes <= self.remaining); + self.remaining -= bytes; + self.remaining == 0 + } +} + +#[derive(Debug)] +struct FrameBudget { + limits: NonNull, + owners: usize, + bytes: usize, +} + +// Budgets and owners remain on NetworkTile. P2p's boxed limits outlive all +// peer queues, stream state, and Quinn-owned Bytes. +unsafe impl Send for FrameBudget {} + +impl FrameBudget { + fn owner(&mut self, data: Leased) -> Bytes { + assert!(self.owners > 0 && data.len() <= self.bytes); + self.owners -= 1; + self.bytes -= data.len(); + let owner = SegmentOwner { data, limits: self.limits }; + NetworkCounters::CacheSegmentedOwnerAllocations.inc(); + Bytes::from_owner(owner) + } +} + +impl Drop for FrameBudget { + fn drop(&mut self) { + let limits = unsafe { self.limits.as_ref() }; + limits.frames.set(limits.frames.get() - 1); + limits.owners.set(limits.owners.get() - self.owners); + limits.retained_bytes.set(limits.retained_bytes.get() - self.bytes); + } +} + +struct SegmentOwner { + data: Leased, + limits: NonNull, +} + +// Bytes requires Send. Creation and destruction remain on NetworkTile, +// whose boxed limits outlive every peer's Quinn connection. +unsafe impl Send for SegmentOwner {} + +impl AsRef<[u8]> for SegmentOwner { + fn as_ref(&self) -> &[u8] { + self.data.as_ref() + } +} + +impl Drop for SegmentOwner { + fn drop(&mut self) { + let limits = unsafe { self.limits.as_ref() }; + limits.owners.set(limits.owners.get() - 1); + limits.retained_bytes.set(limits.retained_bytes.get() - self.data.len()); + } +} + +#[cfg(test)] +mod tests; diff --git a/crates/network/src/p2p/quic/gossip_frame/tests.rs b/crates/network/src/p2p/quic/gossip_frame/tests.rs new file mode 100644 index 00000000..9610011f --- /dev/null +++ b/crates/network/src/p2p/quic/gossip_frame/tests.rs @@ -0,0 +1,526 @@ +use std::{ + alloc::{GlobalAlloc, Layout, System}, + io::Write, + mem, + net::SocketAddr, + time::Duration, +}; + +use quinn_proto::StreamId; +use silver_common::{ + AcquiredWithOffset, CacheFrameRef, CacheSegment, P2pStreamId, StreamProtocol, SubLayout, + SubReservationRef, TCache, TCacheProducer, TProducer, +}; + +use super::*; +use crate::p2p::streams::{ + AcquiredRpcOutbound, StreamError, StreamIo, gossip_out::GossipWriteState, +}; + +thread_local! { + static ALLOCATIONS: Cell = const { Cell::new(0) }; +} + +struct CountingAllocator; + +unsafe impl GlobalAlloc for CountingAllocator { + unsafe fn alloc(&self, layout: Layout) -> *mut u8 { + ALLOCATIONS.with(|count| count.set(count.get() + 1)); + unsafe { System.alloc(layout) } + } + unsafe fn alloc_zeroed(&self, layout: Layout) -> *mut u8 { + ALLOCATIONS.with(|count| count.set(count.get() + 1)); + unsafe { System.alloc_zeroed(layout) } + } + unsafe fn realloc(&self, ptr: *mut u8, layout: Layout, new_size: usize) -> *mut u8 { + ALLOCATIONS.with(|count| count.set(count.get() + 1)); + unsafe { System.realloc(ptr, layout, new_size) } + } + unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) { + unsafe { System.dealloc(ptr, layout) } + } +} + +#[global_allocator] +static ALLOCATOR: CountingAllocator = CountingAllocator; + +struct Harness { + context: Box, + limits: Box, + wheel: Box, + gossip: TProducer, + columns: TProducer, + now: Instant, +} + +impl Harness { + fn new() -> Self { + // Counter files initialise lazily on first touch; warm them so the + // zero-allocation baselines below measure only the frame path. + NetworkCounters::CacheSegmentedAdmitted.inc(); + let gossip = TCache::producer("", 1 << 18); + let columns = TCache::producer("", 1 << 18); + let rpc = TCache::producer("", 1 << 16); + let cluster = TCache::producer("", 1 << 16); + let now = Instant::now(); + Self { + context: Box::new(Context { + gossip_consumer: gossip.cache_ref().strict_random_access("", true).unwrap(), + data_columns_consumer: Some(Box::new( + columns.cache_ref().retained_random_access("").unwrap(), + )), + gossip_producer: TCache::producer("", 1 << 16), + rpc_consumer: rpc.cache_ref().random_access("", true).unwrap(), + rpc_producer: rpc, + identify: None, + cluster_nodes: None, + cluster_inbound_producer: TCache::producer("", 1 << 16), + cluster_outbound_consumer: cluster + .cache_ref() + .strict_random_access("", true) + .unwrap(), + }), + limits: Box::new(SegmentedGossipLimits::new(2)), + wheel: Box::new(OutboundLeaseWheel::new(now)), + gossip, + columns, + now, + } + } + + fn assembly(&mut self, accept: bool) -> (CacheFrameRef, SubReservationRef) { + let reference = self + .columns + .sub_reservation(SubLayout { parts: 2, first_len: 64, second_len: 8 }, b"", b"") + .unwrap(); + let pending = self + .columns + .view_sub_reservation(reference) + .unwrap() + .claim(0) + .unwrap() + .write(&[0xab; 64], &[0xcd; 8]) + .unwrap(); + if accept { + pending + .acquire(self.context.data_columns_consumer.as_deref_mut().unwrap()) + .unwrap() + .accept() + .unwrap(); + } + let frame = CacheFrameRef::write( + &mut self.gossip, + self.now + Duration::from_secs(1), + b"head", + [ + CacheSegment::Framing { offset: 0, length: 4 }, + CacheSegment::Shared { + reservation: reference, + part: 0, + second: false, + offset: 0, + length: 64, + }, + CacheSegment::Shared { + reservation: reference, + part: 0, + second: true, + offset: 0, + length: 8, + }, + ] + .into_iter(), + ) + .unwrap(); + (frame, reference) + } + + fn acquire(&mut self, frame: CacheFrameRef) -> Option { + let view = frame.acquire(&mut self.context.gossip_consumer, self.now).ok()?; + self.limits.acquire(view, &mut self.context, &self.wheel, self.now) + } +} + +impl SegmentedWriter { + fn take_chunk(&mut self) -> Bytes { + let chunk = mem::take(self.chunk().unwrap()); + self.written(chunk.len()); + chunk + } +} + +struct MockIo { + pending: Option, + budget: usize, + retained: Vec, + written: Vec, + fail: bool, +} + +impl MockIo { + fn new(message: OutboundGossip, budget: usize) -> Self { + Self { + pending: Some(message), + budget, + retained: Vec::with_capacity(256), + written: Vec::with_capacity(4096), + fail: false, + } + } +} + +impl StreamIo for MockIo { + fn cluster_next(&mut self) -> Option> { + None + } + + fn write_to_stream(&mut self, _: StreamId, data: &[u8]) -> Result { + let n = self.budget.min(data.len()); + self.written.extend_from_slice(&data[..n]); + Ok(n) + } + fn write_leased_to_stream( + &mut self, + id: StreamId, + data: Leased, + ) -> Result { + self.write_chunks(id, &mut [Bytes::from_owner(data)]) + } + fn write_chunks(&mut self, _: StreamId, chunks: &mut [Bytes]) -> Result { + if self.fail { + return Err(StreamError::StreamClosed); + } + let mut remaining = self.budget; + for chunk in chunks { + let n = remaining.min(chunk.len()); + if n != 0 { + let accepted = chunk.split_to(n); + self.written.extend_from_slice(&accepted); + self.retained.push(accepted); + remaining -= n; + } + } + Ok(self.budget - remaining) + } + fn read_from_stream(&mut self, _: StreamId, _: &mut [u8]) -> Result { + unreachable!() + } + fn close_write(&mut self, _: StreamId) -> Result<(), StreamError> { + Ok(()) + } + fn rpc_next(&mut self) -> Option { + None + } + fn gossip_next(&mut self) -> Option { + self.pending.take() + } + fn remote_addr(&self) -> SocketAddr { + "127.0.0.1:0".parse().unwrap() + } +} + +fn stream() -> P2pStreamId { + P2pStreamId::new(0, 4, StreamProtocol::GossipSub, false) +} + +#[test] +fn admission_is_atomic_and_limits_include_ack_owners() { + let mut h = Harness::new(); + let (invalid, _) = h.assembly(false); + assert!(h.acquire(invalid).is_none()); + assert_eq!(h.wheel.active_count(), 0); + assert_eq!(h.limits.owners.get(), 0); + assert_eq!(h.limits.retained_bytes.get(), 0); + assert_eq!(h.limits.frames.get(), 0); + + let (valid, _) = h.assembly(true); + let mut first = h.acquire(valid).unwrap().into_writer(); + let second = h.acquire(valid).unwrap(); + assert!(h.acquire(valid).is_none()); + drop(first.take_chunk()); + let retained = first.take_chunk(); + drop(first); + drop(second); + assert_eq!(h.limits.frames.get(), 0); + assert_eq!(h.limits.owners.get(), 1); + assert_eq!(h.limits.retained_bytes.get(), 64); + assert_eq!(h.wheel.active_count(), 1); + drop(retained); + assert_eq!(h.limits.retained_bytes.get(), 0); + assert_eq!(h.wheel.active_count(), 0); +} + +#[test] +fn segments_are_allocated_lazily_and_blocked_retries_survive_expiry() { + let mut h = Harness::new(); + let (reference, assembly) = h.assembly(true); + let warm = h.acquire(reference).unwrap(); + drop(warm); + let before = ALLOCATIONS.with(Cell::get); + let frame = h.acquire(reference).unwrap(); + assert_eq!(ALLOCATIONS.with(Cell::get) - before, 0); + let cell_ptr = assembly + .acquire(h.context.data_columns_consumer.as_deref_mut().unwrap()) + .unwrap() + .ranges(0) + .unwrap()[0] + .as_ref() + .as_ptr(); + let mut io = MockIo::new(OutboundGossip::Segmented(frame), 7); + let before = ALLOCATIONS.with(Cell::get); + let mut state = GossipWriteState::Idle.spin(&mut io, &stream()).unwrap(); + assert_eq!(ALLOCATIONS.with(Cell::get) - before, 2); + assert!(matches!(state, GossipWriteState::WritingSegments(_))); + assert_eq!(io.written[0], 76); + io.budget = 0; + let before = ALLOCATIONS.with(Cell::get); + for _ in 0..100 { + state = state.spin(&mut io, &stream()).unwrap(); + } + assert_eq!(ALLOCATIONS.with(Cell::get) - before, 0); + h.columns.view_sub_reservation(assembly).unwrap().close(); + h.context.data_columns_consumer.as_deref_mut().unwrap().advance_retention(h.columns.next_seq()); + assert!(h.acquire(reference).is_none()); + let mut filled = 0; + while let Some(mut reservation) = h.columns.reserve(8192, true) { + reservation.buffer().unwrap().fill(0xee); + reservation.increment_offset(8192); + filled += 1; + assert!(filled < 32); + } + assert!(filled > 0); + io.budget = 13; + let before = ALLOCATIONS.with(Cell::get); + for _ in 0..20 { + state = state.spin(&mut io, &stream()).unwrap(); + if matches!(state, GossipWriteState::Idle) { + break; + } + } + // Only the previously untouched proof creates a new owner. + assert_eq!(ALLOCATIONS.with(Cell::get) - before, 1); + assert!(matches!(state, GossipWriteState::Idle)); + assert_eq!(&io.written[..5], b"\x4chead"); + assert_eq!(&io.written[5..69], &[0xab; 64]); + assert_eq!(&io.written[69..], &[0xcd; 8]); + assert_eq!(io.retained[1].as_ptr(), cell_ptr); + assert_eq!(h.limits.frames.get(), 0); + assert!(h.columns.reserve(8192, true).is_none()); + assert!(h.wheel.expire(h.now + Duration::from_secs(11)).is_some()); + io.retained.clear(); + assert_eq!(h.limits.owners.get(), 0); + assert_eq!(h.wheel.active_count(), 0); + h.context.data_columns_consumer.as_deref_mut().unwrap().advance_retention(h.columns.next_seq()); + assert!(h.columns.reserve(8192, true).is_some()); +} + +#[test] +fn mid_frame_error_never_starts_the_next_frame() { + let mut h = Harness::new(); + let (reference, _) = h.assembly(true); + let first = h.acquire(reference).unwrap(); + let second = h.acquire(reference).unwrap(); + let mut io = MockIo::new(OutboundGossip::Segmented(first), 5); + let state = GossipWriteState::Idle.spin(&mut io, &stream()).unwrap(); + assert_eq!(io.written.len(), 10); + io.pending = Some(OutboundGossip::Segmented(second)); + io.fail = true; + assert!(state.spin(&mut io, &stream()).is_err()); + assert_eq!(io.written.len(), 10); + assert!(io.pending.is_some()); + drop(io); + assert_eq!(h.wheel.active_count(), 0); + assert_eq!(h.limits.frames.get(), 0); + assert_eq!(h.limits.owners.get(), 0); + assert_eq!(h.limits.retained_bytes.get(), 0); +} + +#[test] +fn fanout_has_independent_delivery_leases() { + let mut h = Harness::new(); + let (reference, _) = h.assembly(true); + let wheel = Box::new(OutboundLeaseWheel::new(h.now)); + let mut first = h.acquire(reference).unwrap().into_writer(); + let view = reference.acquire(&mut h.context.gossip_consumer, h.now).unwrap(); + let mut second = h.limits.acquire(view, &mut h.context, &wheel, h.now).unwrap().into_writer(); + drop(first.take_chunk()); + drop(second.take_chunk()); + let first_cell = first.take_chunk(); + let second_cell = second.take_chunk(); + assert_eq!(first_cell.as_ptr(), second_cell.as_ptr()); + drop(first_cell); + drop(second_cell); + drop(first); + assert_eq!(h.wheel.active_count(), 0); + assert!(wheel.active_count() > 0); + assert!(wheel.expire(h.now + Duration::from_secs(11)).is_some()); + drop(second); + assert_eq!(wheel.active_count(), 0); +} + +#[test] +fn byte_and_owner_limits_remain_charged_until_ack() { + let mut h = Harness::new(); + let (reference, _) = h.assembly(true); + let view = reference.acquire(&mut h.context.gossip_consumer, h.now).unwrap(); + h.limits.max_bytes = view.descriptor_len() + view.wire_len(); + h.limits.max_owners = view.segment_count() + 1; + drop(view); + let mut frame = h.acquire(reference).unwrap().into_writer(); + drop(frame.take_chunk()); + let ack_owner = frame.take_chunk(); + drop(frame); + assert!(h.acquire(reference).is_none()); + h.limits.max_bytes = MAX_RETAINED_BYTES; + assert!(h.acquire(reference).is_none()); + drop(ack_owner); + let frame = h.acquire(reference).unwrap(); + drop(frame); + assert_eq!(h.limits.owners.get(), 0); +} + +#[test] +fn adjacent_ranges_share_one_owner_without_gathering() { + let mut h = Harness::new(); + let mut reservation = h.gossip.reserve(64, true).unwrap(); + reservation.write_all(&[0xab; 64]).unwrap(); + let read = reservation.read(); + let reference = CacheFrameRef::write( + &mut h.gossip, + h.now + Duration::from_secs(1), + b"", + [CacheSegment::Gossip { read, offset: 2, length: 4 }, CacheSegment::Gossip { + read, + offset: 6, + length: 8, + }] + .into_iter(), + ) + .unwrap(); + let mut frame = h.acquire(reference).unwrap().into_writer(); + let before = ALLOCATIONS.with(Cell::get); + let chunk = frame.take_chunk(); + assert_eq!(ALLOCATIONS.with(Cell::get) - before, 1); + assert_eq!(chunk.as_ref(), &[0xab; 12]); + assert!(frame.chunk().is_none()); + drop(frame); + assert_eq!(h.limits.owners.get(), 1); + assert_eq!(h.limits.retained_bytes.get(), 12); + drop(chunk); + assert_eq!(h.limits.owners.get(), 0); +} + +#[test] +fn framing_fragments_share_one_lazy_owner_until_the_last_ack() { + let mut h = Harness::new(); + let reference = CacheFrameRef::write( + &mut h.gossip, + h.now + Duration::from_secs(1), + b"abcdef", + [ + CacheSegment::Framing { offset: 0, length: 2 }, + CacheSegment::Framing { offset: 4, length: 2 }, + CacheSegment::Framing { offset: 2, length: 2 }, + ] + .into_iter(), + ) + .unwrap(); + let before = ALLOCATIONS.with(Cell::get); + let before = ALLOCATIONS.with(Cell::get); + let mut writer = h.acquire(reference).unwrap().into_writer(); + assert_eq!(ALLOCATIONS.with(Cell::get) - before, 0); + let descriptor_len = writer.frame.segments.descriptor_len(); + let chunks = [writer.take_chunk(), writer.take_chunk(), writer.take_chunk()]; + assert_eq!(ALLOCATIONS.with(Cell::get) - before, 1); + assert!(writer.chunk().is_none()); + drop(writer); + assert_eq!(h.limits.frames.get(), 0); + assert_eq!(h.limits.owners.get(), 1); + assert_eq!(h.limits.retained_bytes.get(), descriptor_len); + let [first, second, third] = chunks; + assert_eq!(first.as_ref(), b"ab"); + assert_eq!(second.as_ref(), b"ef"); + assert_eq!(third.as_ref(), b"cd"); + drop(first); + drop(second); + assert_eq!(h.wheel.active_count(), 1); + drop(third); + assert_eq!(h.limits.owners.get(), 0); + assert_eq!(h.limits.retained_bytes.get(), 0); + assert_eq!(h.wheel.active_count(), 0); +} + +#[test] +fn dropping_before_the_prefix_allocates_no_owners_and_releases_the_lease() { + let mut h = Harness::new(); + let (reference, _) = h.assembly(true); + let frame = h.acquire(reference).unwrap(); + assert_eq!(h.wheel.active_count(), 1); + let mut io = MockIo::new(OutboundGossip::Segmented(frame), 0); + let before = ALLOCATIONS.with(Cell::get); + let state = GossipWriteState::Idle.spin(&mut io, &stream()).unwrap(); + assert!(matches!(state, GossipWriteState::WritingLength { written: 0, .. })); + assert!(h.wheel.expire(h.now + Duration::from_secs(11)).is_some()); + drop(state); + assert_eq!(ALLOCATIONS.with(Cell::get) - before, 0); + assert_eq!(h.limits.frames.get(), 0); + assert_eq!(h.limits.owners.get(), 0); + assert_eq!(h.limits.retained_bytes.get(), 0); + assert_eq!(h.wheel.active_count(), 0); +} + +#[test] +fn partial_length_prefix_and_mixed_frames_preserve_boundaries() { + let mut h = Harness::new(); + let reference = CacheFrameRef::write( + &mut h.gossip, + h.now + Duration::from_secs(1), + &[0xab; 300], + [CacheSegment::Framing { offset: 0, length: 300 }].into_iter(), + ) + .unwrap(); + let frame = h.acquire(reference).unwrap(); + let mut io = MockIo::new(OutboundGossip::Segmented(frame), 1); + let mut state = GossipWriteState::Idle.spin(&mut io, &stream()).unwrap(); + assert_eq!(io.written, [0xac]); + assert!(matches!(state, GossipWriteState::WritingLength { written: 1, .. })); + let mut write = h.gossip.reserve(3, true).unwrap(); + write.write_all(b"end").unwrap(); + let read = h.context.gossip_consumer.acquire_strict(write.read()).unwrap(); + io.pending = Some(OutboundGossip::Contiguous(h.wheel.leased(read, h.now))); + io.budget = 0; + state = state.spin(&mut io, &stream()).unwrap(); + assert_eq!(io.written, [0xac]); + io.budget = 17; + for _ in 0..30 { + state = state.spin(&mut io, &stream()).unwrap(); + if matches!(state, GossipWriteState::Idle) { + break; + } + } + assert!(matches!(state, GossipWriteState::Idle)); + assert_eq!(&io.written[..2], &[0xac, 0x02]); + assert_eq!(&io.written[2..302], &[0xab; 300]); + assert_eq!(&io.written[302..], b"\x03end"); +} + +#[test] +fn contiguous_baseline_uses_one_owner_per_write_attempt() { + let mut h = Harness::new(); + let mut reservation = h.gossip.reserve(64, true).unwrap(); + reservation.write_all(&[0xab; 64]).unwrap(); + let read = h.context.gossip_consumer.acquire_strict(reservation.read()).unwrap(); + let mut io = MockIo::new(OutboundGossip::Contiguous(h.wheel.leased(read, h.now)), 7); + let before = ALLOCATIONS.with(Cell::get); + let mut state = GossipWriteState::Idle.spin(&mut io, &stream()).unwrap(); + assert_eq!(ALLOCATIONS.with(Cell::get) - before, 1); + io.budget = 0; + let before = ALLOCATIONS.with(Cell::get); + for _ in 0..10 { + state = state.spin(&mut io, &stream()).unwrap(); + } + assert_eq!(ALLOCATIONS.with(Cell::get) - before, 10); + drop(state); + drop(io); + assert_eq!(h.wheel.active_count(), 0); +} diff --git a/crates/network/src/p2p/quic/leased.rs b/crates/network/src/p2p/quic/leased.rs index a522a63f..b16fd69e 100644 --- a/crates/network/src/p2p/quic/leased.rs +++ b/crates/network/src/p2p/quic/leased.rs @@ -1,16 +1,23 @@ use std::{ cell::Cell, - ops::Deref, + ops::{Deref, DerefMut}, ptr::NonNull, time::{Duration, Instant}, }; +use silver_common::cells::GOSSIP_DELIVERY_RETENTION; + /// End-to-end age after which outbound gossip delivery is stale, measured from /// enqueue until Quinn releases every owner after ACK or teardown. Expiry is /// rounded up to the next wheel tick, so detection occurs within one /// additional second. pub(crate) const GOSSIP_DELIVERY_TIMEOUT: Duration = Duration::from_secs(10); pub(crate) const OUTBOUND_LEASE_TICK: Duration = Duration::from_secs(1); +const _: () = assert!( + GOSSIP_DELIVERY_RETENTION.as_nanos() >= + GOSSIP_DELIVERY_TIMEOUT.saturating_add(OUTBOUND_LEASE_TICK).as_nanos(), + "delivery retention must cover timeout and timer-wheel rounding" +); const OUTBOUND_LEASE_BUCKETS: usize = 32; /// Per-peer timer wheel for outbound delivery leases. The allocation keeps @@ -182,6 +189,12 @@ impl Deref for Leased { } } +impl DerefMut for Leased { + fn deref_mut(&mut self) -> &mut Self::Target { + &mut self.value + } +} + impl> AsRef<[u8]> for Leased { fn as_ref(&self) -> &[u8] { self.value.as_ref() diff --git a/crates/network/src/p2p/quic/mod.rs b/crates/network/src/p2p/quic/mod.rs index e3fca037..6e689aa0 100644 --- a/crates/network/src/p2p/quic/mod.rs +++ b/crates/network/src/p2p/quic/mod.rs @@ -8,10 +8,12 @@ use silver_common::{Keypair, PeerId}; use super::tls; +mod gossip_frame; mod leased; mod peer; mod stream; +pub(crate) use gossip_frame::{OutboundGossip, SegmentedGossipLimits, SegmentedWriter}; pub(crate) use leased::Leased; #[cfg(test)] pub(crate) use leased::OutboundLeaseWheel; diff --git a/crates/network/src/p2p/quic/peer.rs b/crates/network/src/p2p/quic/peer.rs index a5e235f2..0f363966 100644 --- a/crates/network/src/p2p/quic/peer.rs +++ b/crates/network/src/p2p/quic/peer.rs @@ -13,7 +13,8 @@ use quinn_proto::{ VarInt, }; use silver_common::{ - P2pConnectionStats, P2pStreamId, PeerId, StreamProtocol, TRead, rpc_rate_limit::RpcRateLimitSet, + CacheFrameRef, P2pConnectionStats, P2pStreamId, PeerId, StreamProtocol, TRead, + rpc_rate_limit::RpcRateLimitSet, }; use crate::{ @@ -22,8 +23,7 @@ use crate::{ NetEvent, context::Context, quic::{ - SendResult, - leased::{Leased, OutboundLeaseWheel}, + Leased, OutboundGossip, SegmentedGossipLimits, SendResult, leased::OutboundLeaseWheel, stream::StreamIoImpl, }, streams::{ @@ -171,6 +171,36 @@ impl Peer { if self.check_outbound_delivery_timeout(now, rpc_codec_pool) { return SendResult::ConnectionClosing; } + let msg = OutboundGossip::Contiguous(self.outbound_lease_wheel.leased(msg, now)); + self.queue_gossip(msg) + } + + pub(crate) fn send_segmented_gossip( + &mut self, + frame: CacheFrameRef, + context: &mut Context, + limits: &SegmentedGossipLimits, + rpc_codec_pool: &mut RpcCodecPool, + ) -> SendResult { + if self.connection.is_closed() { + return SendResult::ConnectionClosing; + } + let now = Instant::now(); + if self.check_outbound_delivery_timeout(now, rpc_codec_pool) { + return SendResult::ConnectionClosing; + } + let acquired = frame + .acquire(&mut context.gossip_consumer, now) + .ok() + .and_then(|view| limits.acquire(view, context, &self.outbound_lease_wheel, now)); + let Some(frame) = acquired else { + crate::NetworkCounters::CacheSegmentedRejected.inc(); + return SendResult::MessageDropped; + }; + self.queue_gossip(OutboundGossip::Segmented(frame)) + } + + fn queue_gossip(&mut self, msg: OutboundGossip) -> SendResult { self.dirty = true; let stream_id = match self.outbound_gossip { Some(id) => id, @@ -182,7 +212,6 @@ impl Peer { None => return SendResult::StreamCreationError, }, }; - let msg = self.outbound_lease_wheel.leased(msg, now); if let Some(stream) = self.streams.get_mut(&stream_id) { if let OutboundBuffer::Gossip(buffer) = &mut stream.out_buffer { let dropped = buffer.add_msg(msg); @@ -1141,7 +1170,7 @@ impl Stream { pub(super) enum OutboundBuffer { Unset, - Gossip(OutBuffer>), + Gossip(OutBuffer), Rpc(OutBuffer), Cluster(OutBuffer>), } @@ -1229,7 +1258,7 @@ mod tests { use mio::{Poll, Token}; use quinn_proto::{DatagramEvent, Endpoint, EndpointConfig}; - use silver_common::{Enr, Keypair, TCache, TCacheProducer, TConsumer, TProducer}; + use silver_common::{CacheSegment, Enr, Keypair, TCache, TCacheProducer, TConsumer, TProducer}; use super::*; use crate::{ @@ -1405,6 +1434,7 @@ mod tests { Self { context: Context { + data_columns_consumer: None, gossip_producer: gossip_in_p, gossip_consumer: gossip_out_c, rpc_producer: rpc_in_p, @@ -2133,6 +2163,104 @@ mod tests { ); } + #[test] + fn segmented_rpc_crosses_quinn_and_releases_owners_after_ack() { + let mut client_h = PeerHarness::new(); + let mut server_h = PeerHarness::new(); + client_h.context.gossip_consumer = + client_h.gossip_out_producer.cache_ref().strict_random_access("", true).unwrap(); + let mut columns = TCache::producer("", 1 << 16); + client_h.context.data_columns_consumer = + Some(Box::new(columns.cache_ref().retained_random_access("").unwrap())); + let limits = Box::new(SegmentedGossipLimits::new(2)); + let mut pair = PeerPair::new(); + // RPC.subscriptions = [{ subscribe: true, topicID: "t" }]. + let payload = b"\x0a\x05\x08\x01\x12\x01t"; + let read = { + let mut write = columns.reserve(5, false).unwrap(); + write.write_all(&payload[2..]).unwrap(); + write.flush().unwrap(); + write.read() + }; + let frame = CacheFrameRef::write( + &mut client_h.gossip_out_producer, + Instant::now() + Duration::from_secs(1), + &payload[..2], + [ + CacheSegment::Framing { offset: 0, length: 2 }, + CacheSegment::DataColumns { read, offset: 0, length: 2 }, + CacheSegment::DataColumns { read, offset: 2, length: 3 }, + ] + .into_iter(), + ) + .unwrap(); + assert_eq!( + pair.client_peer.send_segmented_gossip( + frame, + &mut client_h.context, + &limits, + &mut client_h.rpc_codec_pool + ), + SendResult::Ok + ); + client_h + .context + .data_columns_consumer + .as_deref_mut() + .unwrap() + .advance_retention(columns.next_seq()); + wait_for(&mut pair, &mut client_h, &mut server_h, 200, |_, s| !s.received.is_empty()); + let wire: Vec<_> = server_h.received.values().flatten().copied().collect(); + assert_eq!(wire, payload); + let now = Instant::now(); + for i in 0..200 { + if pair.client_peer.outbound_lease_wheel.active_count() == 0 { + break; + } + pair.step( + now + Duration::from_millis(i), + &mut client_h, + &mut server_h, + &mut |_| {}, + &mut |_| {}, + ); + } + assert_eq!(pair.client_peer.outbound_lease_wheel.active_count(), 0); + } + + #[test] + fn segmented_queue_overflow_drops_only_whole_unstarted_frames() { + let mut h = PeerHarness::new(); + h.context.gossip_consumer = + h.gossip_out_producer.cache_ref().strict_random_access("", true).unwrap(); + let limits = Box::new(SegmentedGossipLimits::new(2)); + let mut pair = PeerPair::new(); + let frame = CacheFrameRef::write( + &mut h.gossip_out_producer, + Instant::now() + Duration::from_secs(1), + b"payload", + [CacheSegment::Framing { offset: 0, length: 7 }].into_iter(), + ) + .unwrap(); + let peer = &mut pair.client_peer; + assert_eq!( + peer.send_segmented_gossip(frame, &mut h.context, &limits, &mut h.rpc_codec_pool), + SendResult::Ok + ); + let id = peer.outbound_gossip.unwrap(); + peer.streams.get_mut(&id).unwrap().out_buffer = OutboundBuffer::Gossip(OutBuffer::new(1)); + assert_eq!(peer.outbound_lease_wheel.active_count(), 0); + for expected in [SendResult::Ok, SendResult::MessageDropped, SendResult::MessageDropped] { + assert_eq!( + peer.send_segmented_gossip(frame, &mut h.context, &limits, &mut h.rpc_codec_pool), + expected + ); + assert_eq!(peer.outbound_lease_wheel.active_count(), 1); + } + peer.clear_streams(&mut h.rpc_codec_pool); + assert_eq!(peer.outbound_lease_wheel.active_count(), 0); + } + #[test] fn bidirectional_data_transfer() { let mut client_h = PeerHarness::new(); diff --git a/crates/network/src/p2p/quic/stream.rs b/crates/network/src/p2p/quic/stream.rs index b1b1cecb..010f60dc 100644 --- a/crates/network/src/p2p/quic/stream.rs +++ b/crates/network/src/p2p/quic/stream.rs @@ -5,7 +5,7 @@ use quinn_proto::{Connection, StreamId, WriteError}; use silver_common::AcquiredWithOffset; use crate::p2p::{ - quic::{leased::Leased, peer::OutboundBuffer}, + quic::{OutboundGossip, leased::Leased, peer::OutboundBuffer}, streams::{AcquiredRpcOutbound, StreamError, StreamIo}, }; @@ -71,6 +71,14 @@ impl<'a> StreamIo for StreamIoImpl<'a> { Ok(offset) } + fn write_chunks(&mut self, id: StreamId, chunks: &mut [Bytes]) -> Result { + match self.connection.send_stream(id).write_chunks(chunks) { + Ok(wrote) => Ok(wrote.bytes), + Err(WriteError::Blocked) => Ok(0), + Err(e) => Err(e.into()), + } + } + fn close_write(&mut self, id: StreamId) -> Result<(), StreamError> { // Finish errors Stopped and Closed are no-ops. let _ = self.connection.send_stream(id).finish(); @@ -84,7 +92,7 @@ impl<'a> StreamIo for StreamIoImpl<'a> { } } - fn gossip_next(&mut self) -> Option> { + fn gossip_next(&mut self) -> Option { match self.outbound { OutboundBuffer::Gossip(out_buffer) => out_buffer.pop(), _ => None, diff --git a/crates/network/src/p2p/streams/gossip_in.rs b/crates/network/src/p2p/streams/gossip_in.rs index 86facf6d..82f39dd8 100644 --- a/crates/network/src/p2p/streams/gossip_in.rs +++ b/crates/network/src/p2p/streams/gossip_in.rs @@ -228,7 +228,7 @@ mod tests { None } - fn gossip_next(&mut self) -> Option> { + fn gossip_next(&mut self) -> Option { None } @@ -361,6 +361,87 @@ mod tests { assert_eq!(&frame[header..], b"cccccc"); } + #[test] + fn unfinished_body_survives_cache_pressure_and_releases_space() { + for complete in [false, true] { + const CAPACITY: usize = 1 << 17; + let mut producer = TCache::producer("", CAPACITY); + let mut consumer = producer.cache_ref().random_access("", true).unwrap(); + let slow_id = P2pStreamId::new(0, 4, StreamProtocol::GossipSub, true); + let fast_id = P2pStreamId::new(1, 4, StreamProtocol::GossipSub, true); + let header = size_of::(); + let now = Instant::now(); + let mut partial = vec![100u8]; + partial.extend_from_slice(&[0xaa; 10]); + let mut slow_io = MockIo { data: partial, pos: 0 }; + let slow = GossipReadState::default() + .spin(&mut slow_io, &mut producer, &slow_id, now, &mut |_| {}) + .unwrap(); + assert!(matches!(slow, GossipReadState::ReadingBody { remaining: 90, .. })); + + let mut wire = vec![100u8]; + wire.extend_from_slice(&[0xbb; 100]); + let mut fast_io = MockIo { data: wire, pos: 0 }; + let mut fast = GossipReadState::default(); + loop { + fast_io.pos = 0; + fast = fast + .spin(&mut fast_io, &mut producer, &fast_id, now, &mut |event| { + let NetEvent::Gossip { msg, .. } = event else { + panic!("expected gossip"); + }; + let acquired = consumer.acquire(msg); + assert_eq!(&acquired.buffer().unwrap().0[header..], &[0xbb; 100]); + }) + .unwrap(); + assert!(producer.next_seq() <= CAPACITY as u64); + if matches!(fast, GossipReadState::AllocBody { .. }) { + break; + } + } + let blocked_head = producer.next_seq(); + let GossipReadState::ReadingBody { reservation, .. } = &slow else { + panic!("expected unfinished body"); + }; + assert_eq!(&reservation.buffer().unwrap()[header..header + 10], &[0xaa; 10]); + + let later = now + GOSSIP_BODY_STALL_TIMEOUT + Duration::from_millis(1); + if complete { + slow_io.data.extend_from_slice(&[0xcc; 90]); + let mut received = None; + slow.spin(&mut slow_io, &mut producer, &slow_id, later, &mut |event| { + let NetEvent::Gossip { msg, .. } = event else { + panic!("expected gossip"); + }; + received = Some(msg); + }) + .unwrap(); + let bytes = producer.read_buffer(received.unwrap()).unwrap(); + assert_eq!(&bytes[..header], slow_id.as_ref()); + assert_eq!(&bytes[header..header + 10], &[0xaa; 10]); + assert_eq!(&bytes[header + 10..], &[0xcc; 90]); + } else { + assert!(matches!( + slow.spin(&mut slow_io, &mut producer, &slow_id, later, &mut |_| {}), + Err(StreamError::ReadStall) + )); + } + + let mut received = None; + let fast = fast + .spin(&mut fast_io, &mut producer, &fast_id, later, &mut |event| { + let NetEvent::Gossip { msg, .. } = event else { + panic!("expected gossip"); + }; + received = Some(msg); + }) + .unwrap(); + assert!(matches!(fast, GossipReadState::ReadingLength { read: 0, .. })); + assert!(producer.next_seq() > blocked_head); + assert_eq!(&producer.read_buffer(received.unwrap()).unwrap()[header..], &[0xbb; 100]); + } + } + #[test] fn frame_length_boundaries() { assert_eq!( diff --git a/crates/network/src/p2p/streams/gossip_out.rs b/crates/network/src/p2p/streams/gossip_out.rs index 7c1799d2..024127b6 100644 --- a/crates/network/src/p2p/streams/gossip_out.rs +++ b/crates/network/src/p2p/streams/gossip_out.rs @@ -1,9 +1,11 @@ +use std::slice; + use silver_common::{MAX_GOSSIP_FRAME_SIZE, P2pStreamId, TRead}; use crate::{ NetworkCounters, p2p::{ - quic::Leased, + quic::{Leased, OutboundGossip, SegmentedWriter}, streams::{StreamError, StreamIo}, }, }; @@ -16,7 +18,7 @@ pub(crate) enum GossipWriteState { buffer: [u8; 10], limit: usize, written: usize, - message: Leased, + message: OutboundGossip, }, /// Writing body. `offset`/`length` track progress into the current /// message; the handler provides body bytes via `send_data`. @@ -25,6 +27,7 @@ pub(crate) enum GossipWriteState { length: usize, message: Leased, }, + WritingSegments(SegmentedWriter), } enum Spin { @@ -56,7 +59,10 @@ impl GossipWriteState { Self::Idle => match io.gossip_next() { Some(message) => { let mut buffer = [0u8; 10]; - let len = message.len()?; + let len = match &message { + OutboundGossip::Contiguous(message) => message.len()?, + OutboundGossip::Segmented(frame) => frame.wire_len(), + }; if len > MAX_GOSSIP_FRAME_SIZE { return Err(StreamError::GossipFrameTooLarge); } @@ -73,10 +79,13 @@ impl GossipWriteState { let n = io.write_to_stream(p2p_id.stream_id(), &buffer[written..limit])?; written += n; if written == limit { - return Ok(Spin::Next(Self::Writing { - offset: 0, - length: message.len()?, - message, + return Ok(Spin::Next(match message { + OutboundGossip::Contiguous(message) => { + Self::Writing { offset: 0, length: message.len()?, message } + } + OutboundGossip::Segmented(frame) => { + Self::WritingSegments(frame.into_writer()) + } })); } Ok(Spin::Ok(Self::WritingLength { buffer, limit, written, message })) @@ -95,16 +104,31 @@ impl GossipWriteState { } Ok(Spin::Ok(Self::Writing { offset, length, message })) } + Self::WritingSegments(mut frame) => { + let chunk = frame.chunk().ok_or(StreamError::InvalidGossipFrame)?; + let n = io.write_chunks(p2p_id.stream_id(), slice::from_mut(chunk))?; + let chunk_complete = chunk.is_empty(); + if frame.written(n) { + Ok(Spin::Next(Self::Idle)) + } else if chunk_complete { + Ok(Spin::Next(Self::WritingSegments(frame))) + } else { + Ok(Spin::Ok(Self::WritingSegments(frame))) + } + } } } } #[cfg(test)] mod tests { - use std::{io::Write as _, net::SocketAddr, time::Instant}; + use std::{array, io::Write as _, net::SocketAddr, time::Instant}; + use bytes::Bytes; use quinn_proto::StreamId; - use silver_common::{StreamProtocol, TCache, TCacheProducer, TProducer, TRandomAccess}; + use silver_common::{ + AcquiredWithOffset, StreamProtocol, TCache, TCacheProducer, TProducer, TRandomAccess, + }; use super::*; use crate::p2p::{quic::OutboundLeaseWheel, streams::AcquiredRpcOutbound}; @@ -113,23 +137,29 @@ mod tests { /// most `budget` bytes per write call (0 = peer granting no credit). struct MockIo { pending: Option>, - retained: Vec>, + retained: Vec, + written: Vec, budget: usize, } impl StreamIo for MockIo { fn write_to_stream(&mut self, _id: StreamId, data: &[u8]) -> Result { - Ok(data.len().min(self.budget)) + let n = data.len().min(self.budget); + self.written.extend_from_slice(&data[..n]); + Ok(n) } fn write_leased_to_stream( &mut self, _id: StreamId, - data: Leased, + data: Leased, ) -> Result { - let n = data.as_ref().len().min(self.budget); + let mut data = Bytes::from_owner(data); + let n = data.len().min(self.budget); if n != 0 { - self.retained.push(data); + let accepted = data.split_to(n); + self.written.extend_from_slice(&accepted); + self.retained.push(accepted); } Ok(n) } @@ -146,8 +176,8 @@ mod tests { None } - fn gossip_next(&mut self) -> Option> { - self.pending.take() + fn gossip_next(&mut self) -> Option { + self.pending.take().map(OutboundGossip::Contiguous) } fn remote_addr(&self) -> SocketAddr { @@ -178,7 +208,12 @@ mod tests { let (_consumer, _producer, msg) = queued_msg("test_gossip_wstall"); let now = Instant::now(); let wheel = Box::new(OutboundLeaseWheel::new(now)); - let mut io = MockIo { pending: Some(wheel.leased(msg, now)), retained: vec![], budget: 0 }; + let mut io = MockIo { + pending: Some(wheel.leased(msg, now)), + retained: vec![], + written: vec![], + budget: 0, + }; let state = GossipWriteState::Idle.spin(&mut io, &p2p_id).expect("blocked write parks"); assert!(matches!(state, GossipWriteState::WritingLength { written: 0, .. })); @@ -193,8 +228,12 @@ mod tests { let (_consumer, _producer, msg) = queued_msg("test_gossip_wprogress"); let now = Instant::now(); let wheel = Box::new(OutboundLeaseWheel::new(now)); - let mut io = - MockIo { pending: Some(wheel.leased(msg, now)), retained: vec![], budget: usize::MAX }; + let mut io = MockIo { + pending: Some(wheel.leased(msg, now)), + retained: vec![], + written: vec![], + budget: usize::MAX, + }; let state = GossipWriteState::Idle.spin(&mut io, &p2p_id).expect("write completes"); assert!(matches!(state, GossipWriteState::Idle)); @@ -203,4 +242,83 @@ mod tests { io.retained.clear(); assert_eq!(wheel.active_count(), 0); } + + #[test] + fn strict_partial_write_resumes_at_offset_and_pins_bytes_until_last_ack_owner_drops() { + const CAPACITY: usize = 1 << 18; + const CHURN_BYTES: usize = 8 * 1024; + + let mut producer = TCache::producer("", CAPACITY); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + let payload: [u8; 513] = array::from_fn(|i| i as u8); + let mut reservation = producer.reserve(payload.len(), true).unwrap(); + reservation.write_all(&payload).unwrap(); + let read = reservation.read(); + let message = consumer.acquire_strict(read).unwrap(); + let payload_ptr = message.buffer().unwrap().0.as_ptr(); + let now = Instant::now(); + let wheel = Box::new(OutboundLeaseWheel::new(now)); + let mut io = MockIo { + pending: Some(wheel.leased(message, now)), + retained: vec![], + written: vec![], + budget: 17, + }; + let p2p_id = P2pStreamId::new(0, 4, StreamProtocol::GossipSub, false); + let mut state = GossipWriteState::Idle.spin(&mut io, &p2p_id).unwrap(); + assert!(matches!(state, GossipWriteState::Writing { offset: 17, .. })); + assert_eq!(io.written[..2], [0x81, 0x04]); + assert_eq!(&io.written[2..], &payload[..17]); + assert_eq!(io.retained[0].as_ptr(), payload_ptr); + assert_eq!(wheel.active_count(), 2); + + io.budget = 0; + state = state.spin(&mut io, &p2p_id).unwrap(); + assert!(matches!(state, GossipWriteState::Writing { offset: 17, .. })); + assert_eq!(io.written.len(), 2 + 17); + assert_eq!(io.retained.len(), 1); + assert_eq!(wheel.active_count(), 2, "blocked attempts must release their child owner"); + + let mut produced = 0; + while let Some(mut reservation) = producer.reserve(CHURN_BYTES, true) { + reservation.buffer().unwrap().fill(0xee); + reservation.increment_offset(CHURN_BYTES); + drop(consumer.acquire_strict(reservation.read()).unwrap()); + produced += 1; + assert!(produced <= CAPACITY / CHURN_BYTES, "overwrote a pinned send"); + } + assert!(produced > 0); + + io.budget = 31; + for _ in 0..payload.len().div_ceil(io.budget) { + state = state.spin(&mut io, &p2p_id).unwrap(); + if matches!(state, GossipWriteState::Idle) { + break; + } + } + assert!(matches!(state, GossipWriteState::Idle)); + assert_eq!(&io.written[..2], &[0x81, 0x04]); + assert_eq!(&io.written[2..], &payload); + let mut offset = 0; + for chunk in &io.retained { + assert_eq!(chunk.as_ref(), &payload[offset..offset + chunk.len()]); + assert_eq!(chunk.as_ptr(), payload_ptr.wrapping_add(offset)); + offset += chunk.len(); + } + assert_eq!(offset, payload.len()); + assert_eq!(wheel.active_count(), io.retained.len() as u64); + + let ack_held = io.retained[0].slice(3..); + let clone = ack_held.clone(); + io.retained.clear(); + assert_eq!(wheel.active_count(), 1); + assert!(producer.reserve(CHURN_BYTES, true).is_none()); + drop(ack_held); + assert_eq!(clone.as_ref(), &payload[3..17]); + assert_eq!(wheel.active_count(), 1); + assert!(producer.reserve(CHURN_BYTES, true).is_none()); + drop(clone); + assert_eq!(wheel.active_count(), 0); + assert!(producer.reserve(CHURN_BYTES, true).is_some()); + } } diff --git a/crates/network/src/p2p/streams/mod.rs b/crates/network/src/p2p/streams/mod.rs index 1a1ceb86..4f67ee38 100644 --- a/crates/network/src/p2p/streams/mod.rs +++ b/crates/network/src/p2p/streams/mod.rs @@ -1,11 +1,15 @@ use std::{array::TryFromSliceError, fmt, net::SocketAddr}; use buffa::DecodeError; +use bytes::Bytes; use quinn_proto::{FinishError, ReadError, ReadableError, StreamId, WriteError}; use silver_common::{AcquiredWithOffset, TCacheError, TRead}; use thiserror::Error; -use crate::p2p::{quic::Leased, streams::snappy::SnappyError}; +use crate::p2p::{ + quic::{Leased, OutboundGossip}, + streams::snappy::SnappyError, +}; mod cluster_in; mod cluster_out; @@ -61,10 +65,13 @@ pub trait StreamIo { id: StreamId, data: Leased, ) -> Result; + fn write_chunks(&mut self, _id: StreamId, _chunks: &mut [Bytes]) -> Result { + Err(StreamError::InvalidGossipFrame) + } fn read_from_stream(&mut self, id: StreamId, data: &mut [u8]) -> Result; fn close_write(&mut self, id: StreamId) -> Result<(), StreamError>; fn rpc_next(&mut self) -> Option; - fn gossip_next(&mut self) -> Option>; fn cluster_next(&mut self) -> Option>; + fn gossip_next(&mut self) -> Option; fn remote_addr(&self) -> SocketAddr; } diff --git a/crates/network/src/p2p/streams/negotiate.rs b/crates/network/src/p2p/streams/negotiate.rs index 3d310ce4..9970f8a4 100644 --- a/crates/network/src/p2p/streams/negotiate.rs +++ b/crates/network/src/p2p/streams/negotiate.rs @@ -293,7 +293,7 @@ mod tests { None } - fn gossip_next(&mut self) -> Option> { + fn gossip_next(&mut self) -> Option { None } fn cluster_next(&mut self) -> Option> { diff --git a/crates/network/src/p2p/streams/rpc/request_in.rs b/crates/network/src/p2p/streams/rpc/request_in.rs index c6695784..cb5eebf5 100644 --- a/crates/network/src/p2p/streams/rpc/request_in.rs +++ b/crates/network/src/p2p/streams/rpc/request_in.rs @@ -147,7 +147,7 @@ mod tests { use std::net::SocketAddr; use quinn_proto::StreamId; - use silver_common::{StreamProtocol, TCache, TCacheProducer, TRead}; + use silver_common::{StreamProtocol, TCache, TCacheProducer}; use super::*; use crate::p2p::streams::{ @@ -187,7 +187,7 @@ mod tests { None } - fn gossip_next(&mut self) -> Option> { + fn gossip_next(&mut self) -> Option { None } @@ -203,7 +203,7 @@ mod tests { Ok(data.as_ref().len()) } - fn cluster_next(&mut self) -> Option> { + fn cluster_next(&mut self) -> Option> { None } } diff --git a/crates/network/src/p2p/streams/rpc/response_in.rs b/crates/network/src/p2p/streams/rpc/response_in.rs index 6a0f760f..420f5108 100644 --- a/crates/network/src/p2p/streams/rpc/response_in.rs +++ b/crates/network/src/p2p/streams/rpc/response_in.rs @@ -326,7 +326,7 @@ mod tests { use std::net::SocketAddr; use quinn_proto::StreamId; - use silver_common::{StreamProtocol, TCache, TRead, ssz_view::DATA_COLUMN_SIDECAR_GLOAS_MIN}; + use silver_common::{StreamProtocol, TCache, ssz_view::DATA_COLUMN_SIDECAR_GLOAS_MIN}; use super::*; use crate::p2p::streams::{rpc::AcquiredRpcOutbound, snappy::SnappyEncoder}; @@ -362,7 +362,7 @@ mod tests { None } - fn gossip_next(&mut self) -> Option> { + fn gossip_next(&mut self) -> Option { None } @@ -378,7 +378,7 @@ mod tests { Ok(data.as_ref().len()) } - fn cluster_next(&mut self) -> Option> { + fn cluster_next(&mut self) -> Option> { None } } diff --git a/crates/network/src/tile.rs b/crates/network/src/tile.rs index 658b4cc9..0042ab2c 100644 --- a/crates/network/src/tile.rs +++ b/crates/network/src/tile.rs @@ -11,7 +11,7 @@ use quinn_proto::Transmit; use secp256k1::PublicKey; use silver_common::{ BeaconStateEvent, ClusterIn, ClusterMsgIn, ClusterMsgOut, GossipMsgIn, GossipMsgOut, P2pSend, - PeerControl, PeerEvent, PeerStats, RpcInbound, RpcOutbound, SilverSpine, + PeerControl, PeerEvent, PeerStats, RpcInbound, RpcOutbound, SilverSpine, cells::RetentionEvent, }; use silver_discovery::{DiscV5, Discovery, DiscoveryEvent}; @@ -102,6 +102,11 @@ impl NetworkTile { fn body(&mut self, adapter: &mut SpineAdapter) { // Consume peer control messages let now = Instant::now(); + if let Some(consumer) = &mut self.inner.context.data_columns_consumer { + adapter.consume(|event: RetentionEvent, _| { + consumer.advance_retention(event.retain_from); + }); + } adapter.consume(|peer_control: PeerControl, _producers| { self.handle_peer_control(peer_control, now); }); @@ -205,6 +210,12 @@ impl NetworkTile { tracing::debug!(peer=gossip_msg_out.peer_id, "send gossip"); self.inner.enqueue_gossip(gossip_msg_out) }, + P2pSend::SegmentedGossip { peer_id, frame } => { + gossips += 1; + self.inner.p2p_endpoint.enqueue_segmented_gossip( + peer_id, frame, &mut self.inner.context, + ) + } P2pSend::Identify(peer) => { self.inner.p2p_endpoint.enqueue_identify(peer) } @@ -327,6 +338,10 @@ where discovery_addr: SocketAddr, discovery: D, ) -> Result { + assert!( + context.data_columns_consumer.as_ref().is_none_or(|consumer| consumer.is_retained()), + "data columns consumer must have a fixed retention boundary" + ); let poll = Poll::new()?; let p2p_socket = Socket::new(p2p_addr, &poll, P2P_SOCKET_TOKEN)?; let disc_socket = Socket::new(discovery_addr, &poll, DISC_SOCKET_TOKEN)?;