From 3ff0b883bfc5ffad4ee41ab4f21dc6b3d7ac88cb Mon Sep 17 00:00:00 2001 From: vladimir-ea Date: Fri, 11 Sep 2026 14:48:32 +0100 Subject: [PATCH 1/4] partial payloads: tcache range acquire and sub-reservations Producer tracks min_allocation so uncommitted reservations block ring reuse. MultiProducer serialises allocation behind a claim, initialising slot headers before the head advances. Consumers gain pinned sub-slot ranges, fixed retention boundaries, per-part sub-reservations with claim/validate/accept states, and gossip frame descriptors over pinned segments. Gossip decompression computes the message id before the offset commits the reservation. Co-Authored-By: Claude Fable 5 --- .../common/src/spine/tcache/gossip_frame.rs | 299 ++++++++++++ .../src/spine/tcache/gossip_frame/acquired.rs | 141 ++++++ .../src/spine/tcache/gossip_frame/tests.rs | 445 ++++++++++++++++++ 3 files changed, 885 insertions(+) create mode 100644 crates/common/src/spine/tcache/gossip_frame.rs create mode 100644 crates/common/src/spine/tcache/gossip_frame/acquired.rs create mode 100644 crates/common/src/spine/tcache/gossip_frame/tests.rs diff --git a/crates/common/src/spine/tcache/gossip_frame.rs b/crates/common/src/spine/tcache/gossip_frame.rs new file mode 100644 index 00000000..ff30daec --- /dev/null +++ b/crates/common/src/spine/tcache/gossip_frame.rs @@ -0,0 +1,299 @@ +use std::{io::Write, ops::Range, time::Instant}; + +use super::{ + AcquiredRange, AcquiredRead, Producer, RandomAccessConsumer, SubReservationRef, TCacheProducer, + TCacheRead, +}; +use crate::MAX_GOSSIP_FRAME_SIZE; + +mod acquired; +pub use acquired::{AcquiredGossipFrame, AcquiredGossipSegment}; + +pub const MAX_GOSSIP_SEGMENTS: usize = 2 * 128 + 8; +const HEADER_BYTES: usize = 16; +const SEGMENT_BYTES: usize = 40; +const MAGIC: [u8; 8] = *b"SGFRAME1"; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum GossipFrameError { + InvalidDescriptor, + TooLarge, + CacheFull, + Expired, + Stale, +} + +#[derive(Clone, Copy, Debug)] +pub enum GossipSegment { + Framing { + offset: usize, + length: usize, + }, + Gossip { + read: TCacheRead, + offset: usize, + length: usize, + }, + DataColumns { + read: TCacheRead, + offset: usize, + length: usize, + }, + Shared { + reservation: SubReservationRef, + part: usize, + second: bool, + offset: usize, + length: usize, + }, +} + +// Only the builder constructs this handle. Encoded source identities originate +// from typed cache descriptors, never from network input. +// Pointer identities and trusted layout metadata make this an in-process +// format. +#[derive(Clone, Copy, Debug)] +pub struct GossipFrameRef { + descriptor: TCacheRead, + expires: Instant, +} + +impl GossipFrameRef { + pub fn write( + producer: &mut Producer, + expires: Instant, + framing: &[u8], + segments: impl ExactSizeIterator, + ) -> Result { + let count = segments.len(); + if count == 0 || count > MAX_GOSSIP_SEGMENTS || framing.len() > MAX_GOSSIP_FRAME_SIZE { + return Err(GossipFrameError::TooLarge); + } + let framing_start = HEADER_BYTES + count * SEGMENT_BYTES; + let descriptor_len = framing_start + framing.len(); + let cache = producer.cache_ref(); + let mut reservation = + producer.reserve(descriptor_len, false).ok_or(GossipFrameError::CacheFull)?; + let buffer = reservation.buffer().map_err(|_| GossipFrameError::Stale)?; + buffer[..8].copy_from_slice(&MAGIC); + buffer[8..12].copy_from_slice(&(count as u32).to_le_bytes()); + buffer[framing_start..].copy_from_slice(framing); + let mut total = 0usize; + let mut written = 0; + for (index, segment) in segments.enumerate() { + if index >= count { + return Err(GossipFrameError::InvalidDescriptor); + } + let (kind, read, metadata, offset, length) = match segment { + GossipSegment::Framing { offset, length } => { + if offset.checked_add(length).is_none_or(|end| end > framing.len()) { + return Err(GossipFrameError::InvalidDescriptor); + } + (0u64, None, 0, offset, length) + } + GossipSegment::Gossip { read, offset, length } => { + if read.tcache.cache != cache.cache { + return Err(GossipFrameError::InvalidDescriptor); + } + (1, Some(read), 0, offset, length) + } + GossipSegment::DataColumns { read, offset, length } => { + (2, Some(read), 0, offset, length) + } + GossipSegment::Shared { reservation, part, second, offset, length } => { + if part >= 128 { + return Err(GossipFrameError::InvalidDescriptor); + } + let metadata = ((reservation.header_bytes as u64) << 32) | + ((part as u64) << 1) | + u64::from(second); + (3, Some(reservation.read()), metadata, offset, length) + } + }; + if length == 0 || offset > u32::MAX as usize || length > u32::MAX as usize { + return Err(GossipFrameError::InvalidDescriptor); + } + total = total.checked_add(length).ok_or(GossipFrameError::TooLarge)?; + if total > MAX_GOSSIP_FRAME_SIZE { + return Err(GossipFrameError::TooLarge); + } + let start = HEADER_BYTES + index * SEGMENT_BYTES; + let entry = &mut buffer[start..start + SEGMENT_BYTES]; + entry[..8].copy_from_slice(&kind.to_le_bytes()); + entry[8..16] + .copy_from_slice(&read.map_or(0, |r| r.tcache.cache as usize as u64).to_le_bytes()); + entry[16..24].copy_from_slice(&read.map_or(0, |r| r.seq).to_le_bytes()); + entry[24..32].copy_from_slice(&metadata.to_le_bytes()); + entry[32..36].copy_from_slice(&(offset as u32).to_le_bytes()); + entry[36..40].copy_from_slice(&(length as u32).to_le_bytes()); + written += 1; + } + if written != count { + return Err(GossipFrameError::InvalidDescriptor); + } + buffer[12..16].copy_from_slice(&(total as u32).to_le_bytes()); + reservation.flush().map_err(|_| GossipFrameError::Stale)?; + Ok(Self { descriptor: reservation.read(), expires }) + } + + pub fn read(self) -> TCacheRead { + self.descriptor + } + + pub fn acquire( + self, + consumer: &mut RandomAccessConsumer, + now: Instant, + ) -> Result { + if now >= self.expires { + return Err(GossipFrameError::Expired); + } + if !consumer.is_strict() || consumer.cache.cache != self.descriptor.tcache.cache { + return Err(GossipFrameError::InvalidDescriptor); + } + let read = consumer.acquire_strict(self.descriptor).ok_or(GossipFrameError::Stale)?; + let buffer = read.buffer().map_err(|_| GossipFrameError::Stale)?.0; + if buffer.len() < HEADER_BYTES || buffer[..8] != MAGIC { + return Err(GossipFrameError::InvalidDescriptor); + } + let count = u32::from_le_bytes(buffer[8..12].try_into().unwrap()) as usize; + let wire_len = u32::from_le_bytes(buffer[12..16].try_into().unwrap()) as usize; + if count == 0 || count > MAX_GOSSIP_SEGMENTS { + return Err(GossipFrameError::InvalidDescriptor); + } + let framing_start = HEADER_BYTES + count * SEGMENT_BYTES; + if framing_start > buffer.len() || wire_len == 0 || wire_len > MAX_GOSSIP_FRAME_SIZE { + return Err(GossipFrameError::InvalidDescriptor); + } + let descriptor_len = buffer.len(); + let view = GossipFrameView { read, count, wire_len, framing_start, descriptor_len }; + let mut total = 0usize; + for segment in view.segments() { + if segment.length == 0 || + segment.kind > 3 || + (segment.kind == 0 && + segment + .offset + .checked_add(segment.length) + .is_none_or(|end| end > view.descriptor_len() - framing_start)) + { + return Err(GossipFrameError::InvalidDescriptor); + } + total = total.checked_add(segment.length).ok_or(GossipFrameError::TooLarge)?; + } + if total != wire_len { + return Err(GossipFrameError::InvalidDescriptor); + } + Ok(view) + } +} + +#[derive(Debug)] +pub struct GossipFrameView { + read: AcquiredRead, + count: usize, + wire_len: usize, + framing_start: usize, + descriptor_len: usize, +} + +impl GossipFrameView { + pub fn wire_len(&self) -> usize { + self.wire_len + } + + pub fn segment_count(&self) -> usize { + self.count + } + + pub fn descriptor_len(&self) -> usize { + self.descriptor_len + } + + pub fn descriptor_range(&self) -> AcquiredRange { + self.read.with_range(0, self.descriptor_len()).expect("acquired descriptor") + } + + pub fn acquire_segments( + self, + gossip: &mut RandomAccessConsumer, + columns: Option<&mut RandomAccessConsumer>, + ) -> Option { + AcquiredGossipFrame::new(self, gossip, columns) + } + + pub fn segments(&self) -> impl ExactSizeIterator + '_ { + self.read.buffer().expect("acquired descriptor").0[HEADER_BYTES..self.framing_start] + .chunks_exact(SEGMENT_BYTES) + .map(|entry| GossipFrameSegment::decode(entry, self.framing_start)) + } + + fn segment(&self, index: usize) -> GossipFrameSegment { + assert!(index < self.count); + let start = HEADER_BYTES + index * SEGMENT_BYTES; + let buffer = self.read.buffer().expect("acquired descriptor").0; + GossipFrameSegment::decode(&buffer[start..start + SEGMENT_BYTES], self.framing_start) + } +} + +pub struct GossipFrameSegment { + kind: u64, + cache: u64, + seq: u64, + metadata: u64, + offset: usize, + length: usize, + framing_start: usize, +} + +impl GossipFrameSegment { + fn decode(entry: &[u8], framing_start: usize) -> Self { + Self { + kind: u64::from_le_bytes(entry[..8].try_into().unwrap()), + cache: u64::from_le_bytes(entry[8..16].try_into().unwrap()), + seq: u64::from_le_bytes(entry[16..24].try_into().unwrap()), + metadata: u64::from_le_bytes(entry[24..32].try_into().unwrap()), + offset: u32::from_le_bytes(entry[32..36].try_into().unwrap()) as usize, + length: u32::from_le_bytes(entry[36..40].try_into().unwrap()) as usize, + framing_start, + } + } + + pub fn framing_range(&self) -> Option> { + (self.kind == 0).then(|| { + self.framing_start + self.offset..self.framing_start + self.offset + self.length + }) + } + + pub fn acquire( + &self, + gossip: &mut RandomAccessConsumer, + columns: Option<&mut RandomAccessConsumer>, + ) -> Option { + let consumer = match self.kind { + 1 => gossip, + 2 | 3 => columns?, + _ => return None, + }; + if !consumer.is_strict() || + consumer.cache.cache as usize as u64 != self.cache || + !self.seq.is_multiple_of(super::ALIGN as u64) + { + return None; + } + let read = TCacheRead { tcache: consumer.cache, seq: self.seq }; + if self.kind == 3 { + let reference = + SubReservationRef { read, header_bytes: (self.metadata >> 32) as usize }; + let acquired = reference.acquire(consumer).ok()?; + let [first, second] = acquired.ranges(((self.metadata as u32) >> 1) as usize)?; + let range = if self.metadata & 1 == 0 { first } else { second }; + range.slice(self.offset, self.length) + } else { + consumer.acquire_strict(read)?.with_range(self.offset, self.length) + } + } +} + +#[cfg(test)] +mod tests; diff --git a/crates/common/src/spine/tcache/gossip_frame/acquired.rs b/crates/common/src/spine/tcache/gossip_frame/acquired.rs new file mode 100644 index 00000000..2aa2d1e7 --- /dev/null +++ b/crates/common/src/spine/tcache/gossip_frame/acquired.rs @@ -0,0 +1,141 @@ +use std::{ + mem, + ops::{Deref, Range}, + ptr::NonNull, +}; + +use super::{ + AcquiredRange, AcquiredRead, GossipFrameSegment, GossipFrameView, RandomAccessConsumer, + SubReservationRef, TCacheRead, +}; + +pub enum AcquiredGossipSegment { + Framing(Range), + Data(AcquiredRange), +} + +#[derive(Debug)] +pub struct AcquiredGossipFrame { + view: GossipFrameView, + gossip: NonNull, + columns: Option>, + // Every non-framing descriptor in [next, acquired_end) owns one bucket + // count. The descriptor stays pinned until those counts are released. + next: usize, + acquired_end: usize, +} + +// As with AcquiredRead, consumers stay at stable addresses and outlive their +// reads. Acquisition, handoff, and drops remain on the consumer's tile. +unsafe impl Send for AcquiredGossipFrame {} + +impl AcquiredGossipFrame { + pub(super) fn new( + view: GossipFrameView, + gossip: &mut RandomAccessConsumer, + mut columns: Option<&mut RandomAccessConsumer>, + ) -> Option { + let mut frame = Self { + view, + gossip: NonNull::from(&mut *gossip), + columns: columns.as_deref_mut().map(NonNull::from), + next: 0, + acquired_end: 0, + }; + for segment in frame.view.segments() { + if segment.kind != 0 { + let range = segment.acquire(gossip, columns.as_deref_mut())?; + // No fallible work separates forgetting this owner and recording + // its count in the frame's acquired prefix. + mem::forget(range); + } + frame.acquired_end += 1; + } + Some(frame) + } + + pub fn take_next(&mut self) -> Option { + if self.next == self.acquired_end { + return None; + } + let segment = self.view.segment(self.next); + if let Some(range) = segment.framing_range() { + self.next += 1; + return Some(AcquiredGossipSegment::Framing(range)); + } + let mut range = self.take_range(&segment); + while self.next < self.acquired_end { + let next = self.view.segment(self.next); + if next.kind == 0 || + range.read.consumer != self.consumer(&next).as_ptr() || + range.read.seq() != next.seq || + range.offset + range.length != Self::offset(&next, range.read.read) + { + break; + } + let next = self.take_range(&next); + assert!(range.extend_contiguous(&next)); + } + Some(AcquiredGossipSegment::Data(range)) + } + + fn consumer(&self, segment: &GossipFrameSegment) -> NonNull { + if segment.kind == 1 { self.gossip } else { self.columns.expect("acquired column segment") } + } + + fn take_read(&mut self, segment: &GossipFrameSegment) -> AcquiredRead { + let consumer = self.consumer(segment); + // Each call transfers exactly one existing count. Nothing increments + // here, and frame cleanup excludes the transferred descriptor. + let read = AcquiredRead { + consumer: consumer.as_ptr(), + read: TCacheRead { tcache: unsafe { consumer.as_ref() }.cache, seq: segment.seq }, + acquired: self.view.read.acquired, + }; + self.next += 1; + read + } + + fn take_range(&mut self, segment: &GossipFrameSegment) -> AcquiredRange { + let read = self.take_read(segment); + let offset = Self::offset(segment, read.read); + AcquiredRange { read, offset, length: segment.length } + } + + fn offset(segment: &GossipFrameSegment, read: TCacheRead) -> usize { + if segment.kind != 3 { + return segment.offset; + } + let reference = SubReservationRef { read, header_bytes: (segment.metadata >> 32) as usize }; + // Admission validated this part and retains its count. Its immutable + // layout remains valid even after the reservation is closed. + let base = unsafe { + reference.acquired_offset( + ((segment.metadata as u32) >> 1) as usize, + segment.metadata & 1 != 0, + ) + }; + base + segment.offset + } +} + +impl Deref for AcquiredGossipFrame { + type Target = GossipFrameView; + + fn deref(&self) -> &Self::Target { + &self.view + } +} + +impl Drop for AcquiredGossipFrame { + fn drop(&mut self) { + while self.next < self.acquired_end { + let segment = self.view.segment(self.next); + if segment.kind == 0 { + self.next += 1; + } else { + drop(self.take_read(&segment)); + } + } + } +} diff --git a/crates/common/src/spine/tcache/gossip_frame/tests.rs b/crates/common/src/spine/tcache/gossip_frame/tests.rs new file mode 100644 index 00000000..9ed1cad5 --- /dev/null +++ b/crates/common/src/spine/tcache/gossip_frame/tests.rs @@ -0,0 +1,445 @@ +use std::{ + panic::{AssertUnwindSafe, catch_unwind}, + time::Duration, +}; + +use super::*; +use crate::{P2pSend, SubLayout, TCache}; + +fn write(producer: &mut Producer, bytes: &[u8]) -> TCacheRead { + let mut reservation = producer.reserve(bytes.len(), false).unwrap(); + reservation.write_all(bytes).unwrap(); + reservation.flush().unwrap(); + reservation.read() +} + +#[test] +fn copy_handle_round_trips_framing_and_source_ranges() { + fn is_copy() {} + is_copy::(); + is_copy::(); + let mut producer = TCache::producer("", 1 << 18); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + let source = write(&mut producer, b"0123456789"); + let now = Instant::now(); + let frame = GossipFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"ab--cd", + [ + GossipSegment::Framing { offset: 0, length: 2 }, + GossipSegment::Gossip { read: source, offset: 3, length: 4 }, + GossipSegment::Framing { offset: 4, length: 2 }, + ] + .into_iter(), + ) + .unwrap(); + let view = frame.acquire(&mut consumer, now).unwrap(); + assert_eq!(view.wire_len(), 8); + assert_eq!(view.segment_count(), 3); + let descriptor = view.descriptor_range(); + let mut wire = Vec::new(); + for segment in view.segments() { + if let Some(range) = segment.framing_range() { + wire.extend_from_slice(&descriptor.as_ref()[range]); + } else { + wire.extend_from_slice(segment.acquire(&mut consumer, None).unwrap().as_ref()); + } + } + assert_eq!(wire, b"ab3456cd"); + assert!(matches!( + frame.acquire(&mut consumer, now + Duration::from_secs(1)), + Err(GossipFrameError::Expired) + )); +} + +#[test] +fn shared_segments_expose_only_verified_subranges() { + let mut producer = TCache::producer("", 1 << 18); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + let mut columns = TCache::producer("", 1 << 18); + let mut reader = Box::new(columns.cache_ref().retained_random_access("").unwrap()); + let reference = columns + .sub_reservation(SubLayout { parts: 2, first_len: 4, second_len: 2 }, b"", b"") + .unwrap(); + let pending = columns + .view_sub_reservation(reference) + .unwrap() + .claim(0) + .unwrap() + .write(b"cell", b"pf") + .unwrap(); + let now = Instant::now(); + let frame = GossipFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"", + [ + GossipSegment::Shared { + reservation: reference, + part: 0, + second: false, + offset: 1, + length: 2, + }, + GossipSegment::Shared { + reservation: reference, + part: 0, + second: true, + offset: 0, + length: 2, + }, + ] + .into_iter(), + ) + .unwrap(); + let view = frame.acquire(&mut consumer, now).unwrap(); + assert!(view.segments().next().unwrap().acquire(&mut consumer, Some(&mut reader)).is_none()); + pending.acquire(&mut reader).unwrap().accept().unwrap(); + let ranges: Vec<_> = view + .segments() + .map(|segment| segment.acquire(&mut consumer, Some(&mut reader)).unwrap()) + .collect(); + assert_eq!(ranges[0].as_ref(), b"el"); + assert_eq!(ranges[1].as_ref(), b"pf"); + columns.view_sub_reservation(reference).unwrap().close(); + reader.advance_retention(columns.next_seq()); + assert!(view.segments().next().unwrap().acquire(&mut consumer, Some(&mut reader)).is_none()); + assert_eq!(ranges[0].as_ref(), b"el"); +} + +#[test] +fn descriptor_bounds_and_sources_are_checked() { + let mut producer = TCache::producer("", 1 << 18); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + let mut other = TCache::producer("", 1 << 18); + let mut other_reader = Box::new(other.cache_ref().retained_random_access("").unwrap()); + let other_read = write(&mut other, b"data"); + let expires = Instant::now() + Duration::from_secs(1); + for segment in [ + GossipSegment::Framing { offset: usize::MAX, length: 1 }, + GossipSegment::Framing { offset: 0, length: 0 }, + GossipSegment::Framing { offset: 1, length: 4 }, + GossipSegment::Gossip { read: other_read, offset: 0, length: 4 }, + ] { + assert!( + GossipFrameRef::write(&mut producer, expires, b"data", [segment].into_iter()).is_err() + ); + } + assert!(GossipFrameRef::write(&mut producer, expires, b"", [].into_iter()).is_err()); + assert!( + GossipFrameRef::write( + &mut producer, + expires, + b"x", + std::iter::repeat_n( + GossipSegment::Framing { offset: 0, length: 1 }, + MAX_GOSSIP_SEGMENTS + 1, + ) + ) + .is_err() + ); + let frame = GossipFrameRef::write( + &mut producer, + expires, + b"", + [GossipSegment::DataColumns { read: other_read, offset: 2, length: 4 }].into_iter(), + ) + .unwrap(); + let view = frame.acquire(&mut consumer, Instant::now()).unwrap(); + let segment = view.segments().next().unwrap(); + assert!(segment.acquire(&mut consumer, None).is_none()); + assert!(segment.acquire(&mut consumer, Some(&mut other_reader)).is_none()); + + for bytes in [&b"not a descriptor"[..], &b"SGFRAME1\xff\xff\xff\xff\x01\x00\x00\x00"[..]] { + let malformed = GossipFrameRef { descriptor: write(&mut producer, bytes), expires }; + assert!(matches!( + malformed.acquire(&mut consumer, Instant::now()), + Err(GossipFrameError::InvalidDescriptor) + )); + } +} + +#[test] +fn owned_range_slices_check_bounds() { + let mut producer = TCache::producer("", 1 << 16); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + let read = write(&mut producer, b"0123456789"); + let acquired = consumer.acquire_strict(read).unwrap(); + let range = acquired.with_range(2, 6).unwrap(); + assert_eq!(range.clone().slice(2, 3).unwrap().as_ref(), b"456"); + assert!(range.clone().slice(6, 0).unwrap().is_empty()); + assert!(range.clone().slice(6, 1).is_none()); + assert!(range.clone().slice(usize::MAX, 1).is_none()); + assert!(range.slice(1, usize::MAX).is_none()); + let mut first = acquired.with_range(0, 3).unwrap(); + let next = acquired.with_range(3, 4).unwrap(); + let gap = acquired.with_range(8, 1).unwrap(); + assert!(!first.extend_contiguous(&gap)); + assert!(first.extend_contiguous(&next)); + drop(next); + assert_eq!(first.as_ref(), b"0123456"); +} + +#[test] +fn failed_acquisition_releases_every_successful_prefix_without_touching_other_reads() { + let mut producer = TCache::producer("", 1 << 18); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + let mut columns = TCache::producer("", 1 << 18); + let mut reader = Box::new(columns.cache_ref().retained_random_access("").unwrap()); + let gossip = write(&mut producer, b"gossip"); + let column = write(&mut columns, b"column"); + let gossip_guard = consumer.acquire_strict(gossip).unwrap(); + let column_guard = reader.acquire_strict(column).unwrap(); + let shared = columns + .sub_reservation(SubLayout { parts: 2, first_len: 4, second_len: 2 }, b"", b"") + .unwrap(); + columns + .view_sub_reservation(shared) + .unwrap() + .claim(0) + .unwrap() + .write(b"cell", b"pf") + .unwrap() + .acquire(&mut reader) + .unwrap() + .accept() + .unwrap(); + let now = Instant::now(); + let segments = [ + GossipSegment::Framing { offset: 0, length: 1 }, + GossipSegment::Gossip { read: gossip, offset: 1, length: 3 }, + GossipSegment::DataColumns { read: column, offset: 0, length: 4 }, + GossipSegment::Shared { reservation: shared, part: 0, second: false, offset: 0, length: 4 }, + GossipSegment::Shared { reservation: shared, part: 0, second: true, offset: 0, length: 2 }, + GossipSegment::Gossip { read: gossip, offset: 1, length: 3 }, + GossipSegment::Framing { offset: 0, length: 1 }, + ]; + for failed in 0..segments.len() { + let mut descriptors = segments; + descriptors[failed] = GossipSegment::DataColumns { read: column, offset: 100, length: 1 }; + let frame = GossipFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"f", + descriptors.into_iter(), + ) + .unwrap(); + assert!( + frame + .acquire(&mut consumer, now) + .unwrap() + .acquire_segments(&mut consumer, Some(&mut reader)) + .is_none() + ); + assert_eq!(consumer.active_count(), 1, "failed descriptor {failed}"); + assert_eq!(reader.active_count(), 1, "failed descriptor {failed}"); + assert_eq!(gossip_guard.buffer().unwrap().0, b"gossip"); + assert_eq!(column_guard.buffer().unwrap().0, b"column"); + } + drop(gossip_guard); + drop(column_guard); + assert_eq!(consumer.active_count(), 0); + assert_eq!(reader.active_count(), 0); +} + +#[test] +fn handoff_transfers_counts_and_frame_drop_releases_only_the_remainder() { + let mut producer = TCache::producer("", 1 << 18); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + let mut columns = TCache::producer("", 1 << 18); + let mut reader = Box::new(columns.cache_ref().retained_random_access("").unwrap()); + let gossip = write(&mut producer, b"gossip"); + let column = write(&mut columns, b"column"); + let now = Instant::now(); + let reference = GossipFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"f", + [ + GossipSegment::Framing { offset: 0, length: 1 }, + GossipSegment::Gossip { read: gossip, offset: 0, length: 6 }, + GossipSegment::DataColumns { read: column, offset: 0, length: 6 }, + GossipSegment::DataColumns { read: column, offset: 0, length: 6 }, + GossipSegment::Framing { offset: 0, length: 1 }, + ] + .into_iter(), + ) + .unwrap(); + let mut frame = reference + .acquire(&mut consumer, now) + .unwrap() + .acquire_segments(&mut consumer, Some(&mut reader)) + .unwrap(); + assert_eq!(consumer.active_count(), 2); + assert_eq!(reader.active_count(), 2); + assert!(matches!(frame.take_next(), Some(AcquiredGossipSegment::Framing(_)))); + let Some(AcquiredGossipSegment::Data(gossip_range)) = frame.take_next() else { panic!() }; + let Some(AcquiredGossipSegment::Data(column_range)) = frame.take_next() else { panic!() }; + assert_eq!(consumer.active_count(), 2); + assert_eq!(reader.active_count(), 2); + reader.advance_retention(columns.next_seq()); + drop(frame); + assert_eq!(consumer.active_count(), 1); + assert_eq!(reader.active_count(), 1); + assert_eq!(gossip_range.as_ref(), b"gossip"); + assert_eq!(column_range.as_ref(), b"column"); + drop(gossip_range); + drop(column_range); + assert_eq!(consumer.active_count(), 0); + assert_eq!(reader.active_count(), 0); +} + +#[test] +fn shared_handoff_survives_closure_without_exposing_unverified_gaps() { + let mut producer = TCache::producer("", 1 << 18); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + let mut columns = TCache::producer("", 1 << 18); + let mut reader = Box::new(columns.cache_ref().retained_random_access("").unwrap()); + let shared = columns + .sub_reservation(SubLayout { parts: 3, first_len: 4, second_len: 2 }, b"hdr", b"mid") + .unwrap(); + for part in [0, 2] { + columns + .view_sub_reservation(shared) + .unwrap() + .claim(part) + .unwrap() + .write(&[part as u8; 4], &[part as u8 + 10; 2]) + .unwrap() + .acquire(&mut reader) + .unwrap() + .accept() + .unwrap(); + } + let now = Instant::now(); + let reference = GossipFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"", + [ + GossipSegment::Shared { + reservation: shared, + part: 0, + second: false, + offset: 1, + length: 3, + }, + GossipSegment::Shared { + reservation: shared, + part: 2, + second: false, + offset: 0, + length: 4, + }, + GossipSegment::Shared { + reservation: shared, + part: 0, + second: true, + offset: 0, + length: 2, + }, + GossipSegment::Shared { + reservation: shared, + part: 2, + second: true, + offset: 1, + length: 1, + }, + ] + .into_iter(), + ) + .unwrap(); + let mut frame = reference + .acquire(&mut consumer, now) + .unwrap() + .acquire_segments(&mut consumer, Some(&mut reader)) + .unwrap(); + assert_eq!(reader.active_count(), 4); + columns.view_sub_reservation(shared).unwrap().close(); + reader.advance_retention(columns.next_seq()); + assert!( + reference + .acquire(&mut consumer, now) + .unwrap() + .acquire_segments(&mut consumer, Some(&mut reader)) + .is_none() + ); + assert_eq!(reader.active_count(), 4); + for (index, expected) in [&[0; 3][..], &[2; 4], &[10; 2], &[12; 1]].into_iter().enumerate() { + let Some(AcquiredGossipSegment::Data(range)) = frame.take_next() else { panic!() }; + assert_eq!(reader.active_count(), 4 - index); + assert_eq!(range.as_ref(), expected); + drop(range); + assert_eq!(reader.active_count(), 3 - index); + } + assert!(frame.take_next().is_none()); + drop(frame); + assert_eq!(consumer.active_count(), 0); + assert_eq!(reader.active_count(), 0); +} + +#[test] +fn coalescing_transfers_one_pin_and_releases_redundant_pins() { + let mut producer = TCache::producer("", 1 << 18); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + let source = write(&mut producer, b"0123456789"); + let now = Instant::now(); + let reference = GossipFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"", + [ + GossipSegment::Gossip { read: source, offset: 1, length: 3 }, + GossipSegment::Gossip { read: source, offset: 4, length: 3 }, + GossipSegment::Gossip { read: source, offset: 7, length: 2 }, + ] + .into_iter(), + ) + .unwrap(); + let mut frame = reference + .acquire(&mut consumer, now) + .unwrap() + .acquire_segments(&mut consumer, None) + .unwrap(); + assert_eq!(consumer.active_count(), 4); + let Some(AcquiredGossipSegment::Data(range)) = frame.take_next() else { panic!() }; + assert_eq!(consumer.active_count(), 2); + assert!(frame.take_next().is_none()); + drop(frame); + assert_eq!(consumer.active_count(), 1); + assert_eq!(range.as_ref(), b"12345678"); + drop(range); + assert_eq!(consumer.active_count(), 0); +} + +#[test] +fn unwinding_releases_untransferred_pins_but_not_the_handed_off_read() { + let mut producer = TCache::producer("", 1 << 18); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + let source = write(&mut producer, b"data"); + let now = Instant::now(); + let reference = GossipFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"", + [GossipSegment::Gossip { read: source, offset: 0, length: 4 }; 3].into_iter(), + ) + .unwrap(); + let mut handed_off = None; + let result = catch_unwind(AssertUnwindSafe(|| { + let mut frame = reference + .acquire(&mut consumer, now) + .unwrap() + .acquire_segments(&mut consumer, None) + .unwrap(); + let Some(AcquiredGossipSegment::Data(range)) = frame.take_next() else { panic!() }; + handed_off = Some(range); + panic!("abort a partially handed-off frame"); + })); + assert!(result.is_err()); + assert_eq!(consumer.active_count(), 1); + assert_eq!(handed_off.as_ref().unwrap().as_ref(), b"data"); + drop(handed_off); + assert_eq!(consumer.active_count(), 0); +} From 2747006a555bf94f66d8f73d6fd8d38173a80e1f Mon Sep 17 00:00:00 2001 From: vladimir-ea Date: Fri, 11 Sep 2026 14:49:28 +0100 Subject: [PATCH 2/4] partial payloads: cell store and control ingress CellStore assembles data column sidecars cell-by-cell in a tcache via per-row sub-reservations, alongside verified full sidecars. Capacity is derived from the blob schedule and delivery retention. Slot rollover closes assemblies and emits a retention boundary for downstream retained consumers. Control gains CellIngress driving context admission, validation outcomes, and availability events. Co-Authored-By: Claude Fable 5 --- crates/control/src/lib.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/crates/control/src/lib.rs b/crates/control/src/lib.rs index 7f336888..0df2a7ae 100644 --- a/crates/control/src/lib.rs +++ b/crates/control/src/lib.rs @@ -1,5 +1,6 @@ pub mod cell_ingress; pub mod cluster; +pub mod cell_ingress; mod counters; pub mod sync_engine; mod tile; From 342e0472e773eb89a9cd1b605c3331e871a8a307 Mon Sep 17 00:00:00 2001 From: vladimir-ea Date: Fri, 11 Sep 2026 14:50:30 +0100 Subject: [PATCH 3/4] partial payloads: segmented gossip send path P2pSend::SegmentedGossip carries a GossipFrameRef per peer. The network tile acquires frame segments through strict consumers, budgets retained bytes and owners, and streams chunks to quinn zero-copy via Bytes owners that release pins on ack. Retention events advance the data columns consumer boundary. Delivery retention is asserted to cover the lease timeout. Co-Authored-By: Claude Fable 5 --- crates/bin/src/main.rs | 1 + crates/common/src/spine/messages.rs | 8 +- crates/e2e/src/stack.rs | 2 + crates/network/benches/quic_basic.rs | 2 + crates/network/benches/quic_pingpong.rs | 2 + crates/network/src/lib.rs | 10 + crates/network/src/p2p/context.rs | 1 + crates/network/src/p2p/mod.rs | 26 +- crates/network/src/p2p/quic/gossip_frame.rs | 195 +++++++ .../src/p2p/quic/gossip_frame/tests.rs | 511 ++++++++++++++++++ crates/network/src/p2p/quic/leased.rs | 15 +- crates/network/src/p2p/quic/mod.rs | 2 + crates/network/src/p2p/quic/peer.rs | 140 ++++- crates/network/src/p2p/quic/stream.rs | 12 +- crates/network/src/p2p/streams/gossip_in.rs | 85 ++- crates/network/src/p2p/streams/gossip_out.rs | 156 +++++- crates/network/src/p2p/streams/mod.rs | 13 +- crates/network/src/p2p/streams/negotiate.rs | 2 +- .../network/src/p2p/streams/rpc/request_in.rs | 4 +- .../src/p2p/streams/rpc/response_in.rs | 4 +- crates/network/src/tile.rs | 17 +- 21 files changed, 1165 insertions(+), 43 deletions(-) create mode 100644 crates/network/src/p2p/quic/gossip_frame.rs create mode 100644 crates/network/src/p2p/quic/gossip_frame/tests.rs 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..d348ac06 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, + DataKind, Enr, GossipFrameRef, 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: GossipFrameRef }, 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..de8acc25 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, + GossipSegmentedAdmitted, + GossipSegmentedRejected, + GossipSegmentedCapacity, + GossipSegmentedSegments, + GossipSegmentedOwnerAllocations, + GossipSegmentedFrames, + // Reserved owner slots, including owners already held by Quinn. + GossipSegmentedOwners, + // Reserved ranges per recipient, including descriptors; not unique cache backing bytes. + GossipSegmentedRetainedBytes, } } 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..ddc2a9aa 100644 --- a/crates/network/src/p2p/mod.rs +++ b/crates/network/src/p2p/mod.rs @@ -13,11 +13,12 @@ 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, + CacheFrameRef, ClusterMsgOut, GossipMsgOut, Identify, Keypair, P2pConnectionStats, P2pStreamId, PeerId, ProtoIdentify, ProtoIdentifyView, RpcOutbound, RpcRequestOutbound, TCacheRead, }; @@ -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: GossipFrameRef, + 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..db8a2887 --- /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::{ + AcquiredGossipFrame, AcquiredGossipSegment, AcquiredRange, GossipFrameView, 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: GossipFrameView, + 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::GossipSegmentedCapacity.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::GossipSegmentedAdmitted.inc(); + NetworkCounters::GossipSegmentedSegments.add(frame.segment_count() as u64); + Some(SegmentedFrame { segments: wheel.leased(frame, now), budget }) + } + + pub(crate) fn publish_gauges(&self) { + NetworkCounters::GossipSegmentedFrames.set(self.frames.get() as u64); + NetworkCounters::GossipSegmentedOwners.set(self.owners.get() as u64); + NetworkCounters::GossipSegmentedRetainedBytes.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()? { + AcquiredGossipSegment::Framing(range) => { + if self.descriptor.is_empty() { + self.descriptor = budget.owner(segments.child(segments.descriptor_range())); + } + self.descriptor.slice(range) + } + AcquiredGossipSegment::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::GossipSegmentedOwnerAllocations.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..a0541b4b --- /dev/null +++ b/crates/network/src/p2p/quic/gossip_frame/tests.rs @@ -0,0 +1,511 @@ +use std::{ + alloc::{GlobalAlloc, Layout, System}, + io::Write, + mem, + net::SocketAddr, + time::Duration, +}; + +use quinn_proto::StreamId; +use silver_common::{ + AcquiredWithOffset, GossipFrameRef, GossipSegment, 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 { + let gossip = TCache::producer("", 1 << 18); + let columns = TCache::producer("", 1 << 18); + let rpc = 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, + }), + limits: Box::new(SegmentedGossipLimits::new(2)), + wheel: Box::new(OutboundLeaseWheel::new(now)), + gossip, + columns, + now, + } + } + + fn assembly(&mut self, accept: bool) -> (GossipFrameRef, 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 = GossipFrameRef::write( + &mut self.gossip, + self.now + Duration::from_secs(1), + b"head", + [ + GossipSegment::Framing { offset: 0, length: 4 }, + GossipSegment::Shared { + reservation: reference, + part: 0, + second: false, + offset: 0, + length: 64, + }, + GossipSegment::Shared { + reservation: reference, + part: 0, + second: true, + offset: 0, + length: 8, + }, + ] + .into_iter(), + ) + .unwrap(); + (frame, reference) + } + + fn acquire(&mut self, frame: GossipFrameRef) -> 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 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 = GossipFrameRef::write( + &mut h.gossip, + h.now + Duration::from_secs(1), + b"", + [GossipSegment::Gossip { read, offset: 2, length: 4 }, GossipSegment::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 = GossipFrameRef::write( + &mut h.gossip, + h.now + Duration::from_secs(1), + b"abcdef", + [ + GossipSegment::Framing { offset: 0, length: 2 }, + GossipSegment::Framing { offset: 4, length: 2 }, + GossipSegment::Framing { offset: 2, length: 2 }, + ] + .into_iter(), + ) + .unwrap(); + 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 = GossipFrameRef::write( + &mut h.gossip, + h.now + Duration::from_secs(1), + &[0xab; 300], + [GossipSegment::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..141667a7 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, + GossipFrameRef, 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}, + 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: GossipFrameRef, + 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::GossipSegmentedRejected.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 = GossipFrameRef::write( + &mut client_h.gossip_out_producer, + Instant::now() + Duration::from_secs(1), + &payload[..2], + [ + GossipSegment::Framing { offset: 0, length: 2 }, + GossipSegment::DataColumns { read, offset: 0, length: 2 }, + GossipSegment::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 = GossipFrameRef::write( + &mut h.gossip_out_producer, + Instant::now() + Duration::from_secs(1), + b"payload", + [GossipSegment::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..089f9692 100644 --- a/crates/network/src/p2p/streams/gossip_in.rs +++ b/crates/network/src/p2p/streams/gossip_in.rs @@ -194,7 +194,7 @@ mod tests { use std::net::SocketAddr; use quinn_proto::StreamId; - use silver_common::{StreamProtocol, TCache, TRead}; + use silver_common::{StreamProtocol, TCache}; use super::*; use crate::p2p::streams::AcquiredRpcOutbound; @@ -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..b462a8d0 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 silver_common::{AcquiredWithOffset, TCacheError}; 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..6cb31156 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 } diff --git a/crates/network/src/p2p/streams/rpc/response_in.rs b/crates/network/src/p2p/streams/rpc/response_in.rs index 6a0f760f..40c8dde3 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 } 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)?; From cb374e5cf498fff6d6d6e1dc91315d1b60461d66 Mon Sep 17 00:00:00 2001 From: vladimir-ea Date: Tue, 15 Sep 2026 10:29:04 +0100 Subject: [PATCH 4/4] rebase snafu --- crates/common/src/spine/messages.rs | 4 +- .../common/src/spine/tcache/gossip_frame.rs | 299 ------------ .../src/spine/tcache/gossip_frame/acquired.rs | 141 ------ .../src/spine/tcache/gossip_frame/tests.rs | 445 ------------------ crates/control/src/lib.rs | 1 - crates/network/src/lib.rs | 16 +- crates/network/src/p2p/mod.rs | 6 +- crates/network/src/p2p/quic/gossip_frame.rs | 24 +- .../src/p2p/quic/gossip_frame/tests.rs | 45 +- crates/network/src/p2p/quic/peer.rs | 20 +- crates/network/src/p2p/streams/gossip_in.rs | 2 +- crates/network/src/p2p/streams/mod.rs | 2 +- .../network/src/p2p/streams/rpc/request_in.rs | 2 +- .../src/p2p/streams/rpc/response_in.rs | 2 +- 14 files changed, 69 insertions(+), 940 deletions(-) delete mode 100644 crates/common/src/spine/tcache/gossip_frame.rs delete mode 100644 crates/common/src/spine/tcache/gossip_frame/acquired.rs delete mode 100644 crates/common/src/spine/tcache/gossip_frame/tests.rs diff --git a/crates/common/src/spine/messages.rs b/crates/common/src/spine/messages.rs index d348ac06..10f6c240 100644 --- a/crates/common/src/spine/messages.rs +++ b/crates/common/src/spine/messages.rs @@ -8,7 +8,7 @@ use flux::timing::Nanos; use silver_beacon_state_data::SLOTS_PER_EPOCH; use crate::{ - DataKind, Enr, GossipFrameRef, GossipTopic, Identify, MessageId, Origin, P2pStreamId, PeerId, + CacheFrameRef, DataKind, Enr, GossipTopic, Identify, MessageId, Origin, P2pStreamId, PeerId, StreamProtocol, TCacheProducer, TCacheRead, TMultiProducer, column_util::columns_of, ssz_view::{ @@ -711,7 +711,7 @@ pub enum RpcSeverity { #[allow(clippy::large_enum_variant)] pub enum P2pSend { Gossip(GossipMsgOut), - SegmentedGossip { peer_id: usize, frame: GossipFrameRef }, + SegmentedGossip { peer_id: usize, frame: CacheFrameRef }, Identify(usize), Rpc(RpcOutbound), } diff --git a/crates/common/src/spine/tcache/gossip_frame.rs b/crates/common/src/spine/tcache/gossip_frame.rs deleted file mode 100644 index ff30daec..00000000 --- a/crates/common/src/spine/tcache/gossip_frame.rs +++ /dev/null @@ -1,299 +0,0 @@ -use std::{io::Write, ops::Range, time::Instant}; - -use super::{ - AcquiredRange, AcquiredRead, Producer, RandomAccessConsumer, SubReservationRef, TCacheProducer, - TCacheRead, -}; -use crate::MAX_GOSSIP_FRAME_SIZE; - -mod acquired; -pub use acquired::{AcquiredGossipFrame, AcquiredGossipSegment}; - -pub const MAX_GOSSIP_SEGMENTS: usize = 2 * 128 + 8; -const HEADER_BYTES: usize = 16; -const SEGMENT_BYTES: usize = 40; -const MAGIC: [u8; 8] = *b"SGFRAME1"; - -#[derive(Clone, Copy, Debug, PartialEq, Eq)] -pub enum GossipFrameError { - InvalidDescriptor, - TooLarge, - CacheFull, - Expired, - Stale, -} - -#[derive(Clone, Copy, Debug)] -pub enum GossipSegment { - Framing { - offset: usize, - length: usize, - }, - Gossip { - read: TCacheRead, - offset: usize, - length: usize, - }, - DataColumns { - read: TCacheRead, - offset: usize, - length: usize, - }, - Shared { - reservation: SubReservationRef, - part: usize, - second: bool, - offset: usize, - length: usize, - }, -} - -// Only the builder constructs this handle. Encoded source identities originate -// from typed cache descriptors, never from network input. -// Pointer identities and trusted layout metadata make this an in-process -// format. -#[derive(Clone, Copy, Debug)] -pub struct GossipFrameRef { - descriptor: TCacheRead, - expires: Instant, -} - -impl GossipFrameRef { - pub fn write( - producer: &mut Producer, - expires: Instant, - framing: &[u8], - segments: impl ExactSizeIterator, - ) -> Result { - let count = segments.len(); - if count == 0 || count > MAX_GOSSIP_SEGMENTS || framing.len() > MAX_GOSSIP_FRAME_SIZE { - return Err(GossipFrameError::TooLarge); - } - let framing_start = HEADER_BYTES + count * SEGMENT_BYTES; - let descriptor_len = framing_start + framing.len(); - let cache = producer.cache_ref(); - let mut reservation = - producer.reserve(descriptor_len, false).ok_or(GossipFrameError::CacheFull)?; - let buffer = reservation.buffer().map_err(|_| GossipFrameError::Stale)?; - buffer[..8].copy_from_slice(&MAGIC); - buffer[8..12].copy_from_slice(&(count as u32).to_le_bytes()); - buffer[framing_start..].copy_from_slice(framing); - let mut total = 0usize; - let mut written = 0; - for (index, segment) in segments.enumerate() { - if index >= count { - return Err(GossipFrameError::InvalidDescriptor); - } - let (kind, read, metadata, offset, length) = match segment { - GossipSegment::Framing { offset, length } => { - if offset.checked_add(length).is_none_or(|end| end > framing.len()) { - return Err(GossipFrameError::InvalidDescriptor); - } - (0u64, None, 0, offset, length) - } - GossipSegment::Gossip { read, offset, length } => { - if read.tcache.cache != cache.cache { - return Err(GossipFrameError::InvalidDescriptor); - } - (1, Some(read), 0, offset, length) - } - GossipSegment::DataColumns { read, offset, length } => { - (2, Some(read), 0, offset, length) - } - GossipSegment::Shared { reservation, part, second, offset, length } => { - if part >= 128 { - return Err(GossipFrameError::InvalidDescriptor); - } - let metadata = ((reservation.header_bytes as u64) << 32) | - ((part as u64) << 1) | - u64::from(second); - (3, Some(reservation.read()), metadata, offset, length) - } - }; - if length == 0 || offset > u32::MAX as usize || length > u32::MAX as usize { - return Err(GossipFrameError::InvalidDescriptor); - } - total = total.checked_add(length).ok_or(GossipFrameError::TooLarge)?; - if total > MAX_GOSSIP_FRAME_SIZE { - return Err(GossipFrameError::TooLarge); - } - let start = HEADER_BYTES + index * SEGMENT_BYTES; - let entry = &mut buffer[start..start + SEGMENT_BYTES]; - entry[..8].copy_from_slice(&kind.to_le_bytes()); - entry[8..16] - .copy_from_slice(&read.map_or(0, |r| r.tcache.cache as usize as u64).to_le_bytes()); - entry[16..24].copy_from_slice(&read.map_or(0, |r| r.seq).to_le_bytes()); - entry[24..32].copy_from_slice(&metadata.to_le_bytes()); - entry[32..36].copy_from_slice(&(offset as u32).to_le_bytes()); - entry[36..40].copy_from_slice(&(length as u32).to_le_bytes()); - written += 1; - } - if written != count { - return Err(GossipFrameError::InvalidDescriptor); - } - buffer[12..16].copy_from_slice(&(total as u32).to_le_bytes()); - reservation.flush().map_err(|_| GossipFrameError::Stale)?; - Ok(Self { descriptor: reservation.read(), expires }) - } - - pub fn read(self) -> TCacheRead { - self.descriptor - } - - pub fn acquire( - self, - consumer: &mut RandomAccessConsumer, - now: Instant, - ) -> Result { - if now >= self.expires { - return Err(GossipFrameError::Expired); - } - if !consumer.is_strict() || consumer.cache.cache != self.descriptor.tcache.cache { - return Err(GossipFrameError::InvalidDescriptor); - } - let read = consumer.acquire_strict(self.descriptor).ok_or(GossipFrameError::Stale)?; - let buffer = read.buffer().map_err(|_| GossipFrameError::Stale)?.0; - if buffer.len() < HEADER_BYTES || buffer[..8] != MAGIC { - return Err(GossipFrameError::InvalidDescriptor); - } - let count = u32::from_le_bytes(buffer[8..12].try_into().unwrap()) as usize; - let wire_len = u32::from_le_bytes(buffer[12..16].try_into().unwrap()) as usize; - if count == 0 || count > MAX_GOSSIP_SEGMENTS { - return Err(GossipFrameError::InvalidDescriptor); - } - let framing_start = HEADER_BYTES + count * SEGMENT_BYTES; - if framing_start > buffer.len() || wire_len == 0 || wire_len > MAX_GOSSIP_FRAME_SIZE { - return Err(GossipFrameError::InvalidDescriptor); - } - let descriptor_len = buffer.len(); - let view = GossipFrameView { read, count, wire_len, framing_start, descriptor_len }; - let mut total = 0usize; - for segment in view.segments() { - if segment.length == 0 || - segment.kind > 3 || - (segment.kind == 0 && - segment - .offset - .checked_add(segment.length) - .is_none_or(|end| end > view.descriptor_len() - framing_start)) - { - return Err(GossipFrameError::InvalidDescriptor); - } - total = total.checked_add(segment.length).ok_or(GossipFrameError::TooLarge)?; - } - if total != wire_len { - return Err(GossipFrameError::InvalidDescriptor); - } - Ok(view) - } -} - -#[derive(Debug)] -pub struct GossipFrameView { - read: AcquiredRead, - count: usize, - wire_len: usize, - framing_start: usize, - descriptor_len: usize, -} - -impl GossipFrameView { - pub fn wire_len(&self) -> usize { - self.wire_len - } - - pub fn segment_count(&self) -> usize { - self.count - } - - pub fn descriptor_len(&self) -> usize { - self.descriptor_len - } - - pub fn descriptor_range(&self) -> AcquiredRange { - self.read.with_range(0, self.descriptor_len()).expect("acquired descriptor") - } - - pub fn acquire_segments( - self, - gossip: &mut RandomAccessConsumer, - columns: Option<&mut RandomAccessConsumer>, - ) -> Option { - AcquiredGossipFrame::new(self, gossip, columns) - } - - pub fn segments(&self) -> impl ExactSizeIterator + '_ { - self.read.buffer().expect("acquired descriptor").0[HEADER_BYTES..self.framing_start] - .chunks_exact(SEGMENT_BYTES) - .map(|entry| GossipFrameSegment::decode(entry, self.framing_start)) - } - - fn segment(&self, index: usize) -> GossipFrameSegment { - assert!(index < self.count); - let start = HEADER_BYTES + index * SEGMENT_BYTES; - let buffer = self.read.buffer().expect("acquired descriptor").0; - GossipFrameSegment::decode(&buffer[start..start + SEGMENT_BYTES], self.framing_start) - } -} - -pub struct GossipFrameSegment { - kind: u64, - cache: u64, - seq: u64, - metadata: u64, - offset: usize, - length: usize, - framing_start: usize, -} - -impl GossipFrameSegment { - fn decode(entry: &[u8], framing_start: usize) -> Self { - Self { - kind: u64::from_le_bytes(entry[..8].try_into().unwrap()), - cache: u64::from_le_bytes(entry[8..16].try_into().unwrap()), - seq: u64::from_le_bytes(entry[16..24].try_into().unwrap()), - metadata: u64::from_le_bytes(entry[24..32].try_into().unwrap()), - offset: u32::from_le_bytes(entry[32..36].try_into().unwrap()) as usize, - length: u32::from_le_bytes(entry[36..40].try_into().unwrap()) as usize, - framing_start, - } - } - - pub fn framing_range(&self) -> Option> { - (self.kind == 0).then(|| { - self.framing_start + self.offset..self.framing_start + self.offset + self.length - }) - } - - pub fn acquire( - &self, - gossip: &mut RandomAccessConsumer, - columns: Option<&mut RandomAccessConsumer>, - ) -> Option { - let consumer = match self.kind { - 1 => gossip, - 2 | 3 => columns?, - _ => return None, - }; - if !consumer.is_strict() || - consumer.cache.cache as usize as u64 != self.cache || - !self.seq.is_multiple_of(super::ALIGN as u64) - { - return None; - } - let read = TCacheRead { tcache: consumer.cache, seq: self.seq }; - if self.kind == 3 { - let reference = - SubReservationRef { read, header_bytes: (self.metadata >> 32) as usize }; - let acquired = reference.acquire(consumer).ok()?; - let [first, second] = acquired.ranges(((self.metadata as u32) >> 1) as usize)?; - let range = if self.metadata & 1 == 0 { first } else { second }; - range.slice(self.offset, self.length) - } else { - consumer.acquire_strict(read)?.with_range(self.offset, self.length) - } - } -} - -#[cfg(test)] -mod tests; diff --git a/crates/common/src/spine/tcache/gossip_frame/acquired.rs b/crates/common/src/spine/tcache/gossip_frame/acquired.rs deleted file mode 100644 index 2aa2d1e7..00000000 --- a/crates/common/src/spine/tcache/gossip_frame/acquired.rs +++ /dev/null @@ -1,141 +0,0 @@ -use std::{ - mem, - ops::{Deref, Range}, - ptr::NonNull, -}; - -use super::{ - AcquiredRange, AcquiredRead, GossipFrameSegment, GossipFrameView, RandomAccessConsumer, - SubReservationRef, TCacheRead, -}; - -pub enum AcquiredGossipSegment { - Framing(Range), - Data(AcquiredRange), -} - -#[derive(Debug)] -pub struct AcquiredGossipFrame { - view: GossipFrameView, - gossip: NonNull, - columns: Option>, - // Every non-framing descriptor in [next, acquired_end) owns one bucket - // count. The descriptor stays pinned until those counts are released. - next: usize, - acquired_end: usize, -} - -// As with AcquiredRead, consumers stay at stable addresses and outlive their -// reads. Acquisition, handoff, and drops remain on the consumer's tile. -unsafe impl Send for AcquiredGossipFrame {} - -impl AcquiredGossipFrame { - pub(super) fn new( - view: GossipFrameView, - gossip: &mut RandomAccessConsumer, - mut columns: Option<&mut RandomAccessConsumer>, - ) -> Option { - let mut frame = Self { - view, - gossip: NonNull::from(&mut *gossip), - columns: columns.as_deref_mut().map(NonNull::from), - next: 0, - acquired_end: 0, - }; - for segment in frame.view.segments() { - if segment.kind != 0 { - let range = segment.acquire(gossip, columns.as_deref_mut())?; - // No fallible work separates forgetting this owner and recording - // its count in the frame's acquired prefix. - mem::forget(range); - } - frame.acquired_end += 1; - } - Some(frame) - } - - pub fn take_next(&mut self) -> Option { - if self.next == self.acquired_end { - return None; - } - let segment = self.view.segment(self.next); - if let Some(range) = segment.framing_range() { - self.next += 1; - return Some(AcquiredGossipSegment::Framing(range)); - } - let mut range = self.take_range(&segment); - while self.next < self.acquired_end { - let next = self.view.segment(self.next); - if next.kind == 0 || - range.read.consumer != self.consumer(&next).as_ptr() || - range.read.seq() != next.seq || - range.offset + range.length != Self::offset(&next, range.read.read) - { - break; - } - let next = self.take_range(&next); - assert!(range.extend_contiguous(&next)); - } - Some(AcquiredGossipSegment::Data(range)) - } - - fn consumer(&self, segment: &GossipFrameSegment) -> NonNull { - if segment.kind == 1 { self.gossip } else { self.columns.expect("acquired column segment") } - } - - fn take_read(&mut self, segment: &GossipFrameSegment) -> AcquiredRead { - let consumer = self.consumer(segment); - // Each call transfers exactly one existing count. Nothing increments - // here, and frame cleanup excludes the transferred descriptor. - let read = AcquiredRead { - consumer: consumer.as_ptr(), - read: TCacheRead { tcache: unsafe { consumer.as_ref() }.cache, seq: segment.seq }, - acquired: self.view.read.acquired, - }; - self.next += 1; - read - } - - fn take_range(&mut self, segment: &GossipFrameSegment) -> AcquiredRange { - let read = self.take_read(segment); - let offset = Self::offset(segment, read.read); - AcquiredRange { read, offset, length: segment.length } - } - - fn offset(segment: &GossipFrameSegment, read: TCacheRead) -> usize { - if segment.kind != 3 { - return segment.offset; - } - let reference = SubReservationRef { read, header_bytes: (segment.metadata >> 32) as usize }; - // Admission validated this part and retains its count. Its immutable - // layout remains valid even after the reservation is closed. - let base = unsafe { - reference.acquired_offset( - ((segment.metadata as u32) >> 1) as usize, - segment.metadata & 1 != 0, - ) - }; - base + segment.offset - } -} - -impl Deref for AcquiredGossipFrame { - type Target = GossipFrameView; - - fn deref(&self) -> &Self::Target { - &self.view - } -} - -impl Drop for AcquiredGossipFrame { - fn drop(&mut self) { - while self.next < self.acquired_end { - let segment = self.view.segment(self.next); - if segment.kind == 0 { - self.next += 1; - } else { - drop(self.take_read(&segment)); - } - } - } -} diff --git a/crates/common/src/spine/tcache/gossip_frame/tests.rs b/crates/common/src/spine/tcache/gossip_frame/tests.rs deleted file mode 100644 index 9ed1cad5..00000000 --- a/crates/common/src/spine/tcache/gossip_frame/tests.rs +++ /dev/null @@ -1,445 +0,0 @@ -use std::{ - panic::{AssertUnwindSafe, catch_unwind}, - time::Duration, -}; - -use super::*; -use crate::{P2pSend, SubLayout, TCache}; - -fn write(producer: &mut Producer, bytes: &[u8]) -> TCacheRead { - let mut reservation = producer.reserve(bytes.len(), false).unwrap(); - reservation.write_all(bytes).unwrap(); - reservation.flush().unwrap(); - reservation.read() -} - -#[test] -fn copy_handle_round_trips_framing_and_source_ranges() { - fn is_copy() {} - is_copy::(); - is_copy::(); - let mut producer = TCache::producer("", 1 << 18); - let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); - let source = write(&mut producer, b"0123456789"); - let now = Instant::now(); - let frame = GossipFrameRef::write( - &mut producer, - now + Duration::from_secs(1), - b"ab--cd", - [ - GossipSegment::Framing { offset: 0, length: 2 }, - GossipSegment::Gossip { read: source, offset: 3, length: 4 }, - GossipSegment::Framing { offset: 4, length: 2 }, - ] - .into_iter(), - ) - .unwrap(); - let view = frame.acquire(&mut consumer, now).unwrap(); - assert_eq!(view.wire_len(), 8); - assert_eq!(view.segment_count(), 3); - let descriptor = view.descriptor_range(); - let mut wire = Vec::new(); - for segment in view.segments() { - if let Some(range) = segment.framing_range() { - wire.extend_from_slice(&descriptor.as_ref()[range]); - } else { - wire.extend_from_slice(segment.acquire(&mut consumer, None).unwrap().as_ref()); - } - } - assert_eq!(wire, b"ab3456cd"); - assert!(matches!( - frame.acquire(&mut consumer, now + Duration::from_secs(1)), - Err(GossipFrameError::Expired) - )); -} - -#[test] -fn shared_segments_expose_only_verified_subranges() { - let mut producer = TCache::producer("", 1 << 18); - let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); - let mut columns = TCache::producer("", 1 << 18); - let mut reader = Box::new(columns.cache_ref().retained_random_access("").unwrap()); - let reference = columns - .sub_reservation(SubLayout { parts: 2, first_len: 4, second_len: 2 }, b"", b"") - .unwrap(); - let pending = columns - .view_sub_reservation(reference) - .unwrap() - .claim(0) - .unwrap() - .write(b"cell", b"pf") - .unwrap(); - let now = Instant::now(); - let frame = GossipFrameRef::write( - &mut producer, - now + Duration::from_secs(1), - b"", - [ - GossipSegment::Shared { - reservation: reference, - part: 0, - second: false, - offset: 1, - length: 2, - }, - GossipSegment::Shared { - reservation: reference, - part: 0, - second: true, - offset: 0, - length: 2, - }, - ] - .into_iter(), - ) - .unwrap(); - let view = frame.acquire(&mut consumer, now).unwrap(); - assert!(view.segments().next().unwrap().acquire(&mut consumer, Some(&mut reader)).is_none()); - pending.acquire(&mut reader).unwrap().accept().unwrap(); - let ranges: Vec<_> = view - .segments() - .map(|segment| segment.acquire(&mut consumer, Some(&mut reader)).unwrap()) - .collect(); - assert_eq!(ranges[0].as_ref(), b"el"); - assert_eq!(ranges[1].as_ref(), b"pf"); - columns.view_sub_reservation(reference).unwrap().close(); - reader.advance_retention(columns.next_seq()); - assert!(view.segments().next().unwrap().acquire(&mut consumer, Some(&mut reader)).is_none()); - assert_eq!(ranges[0].as_ref(), b"el"); -} - -#[test] -fn descriptor_bounds_and_sources_are_checked() { - let mut producer = TCache::producer("", 1 << 18); - let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); - let mut other = TCache::producer("", 1 << 18); - let mut other_reader = Box::new(other.cache_ref().retained_random_access("").unwrap()); - let other_read = write(&mut other, b"data"); - let expires = Instant::now() + Duration::from_secs(1); - for segment in [ - GossipSegment::Framing { offset: usize::MAX, length: 1 }, - GossipSegment::Framing { offset: 0, length: 0 }, - GossipSegment::Framing { offset: 1, length: 4 }, - GossipSegment::Gossip { read: other_read, offset: 0, length: 4 }, - ] { - assert!( - GossipFrameRef::write(&mut producer, expires, b"data", [segment].into_iter()).is_err() - ); - } - assert!(GossipFrameRef::write(&mut producer, expires, b"", [].into_iter()).is_err()); - assert!( - GossipFrameRef::write( - &mut producer, - expires, - b"x", - std::iter::repeat_n( - GossipSegment::Framing { offset: 0, length: 1 }, - MAX_GOSSIP_SEGMENTS + 1, - ) - ) - .is_err() - ); - let frame = GossipFrameRef::write( - &mut producer, - expires, - b"", - [GossipSegment::DataColumns { read: other_read, offset: 2, length: 4 }].into_iter(), - ) - .unwrap(); - let view = frame.acquire(&mut consumer, Instant::now()).unwrap(); - let segment = view.segments().next().unwrap(); - assert!(segment.acquire(&mut consumer, None).is_none()); - assert!(segment.acquire(&mut consumer, Some(&mut other_reader)).is_none()); - - for bytes in [&b"not a descriptor"[..], &b"SGFRAME1\xff\xff\xff\xff\x01\x00\x00\x00"[..]] { - let malformed = GossipFrameRef { descriptor: write(&mut producer, bytes), expires }; - assert!(matches!( - malformed.acquire(&mut consumer, Instant::now()), - Err(GossipFrameError::InvalidDescriptor) - )); - } -} - -#[test] -fn owned_range_slices_check_bounds() { - let mut producer = TCache::producer("", 1 << 16); - let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); - let read = write(&mut producer, b"0123456789"); - let acquired = consumer.acquire_strict(read).unwrap(); - let range = acquired.with_range(2, 6).unwrap(); - assert_eq!(range.clone().slice(2, 3).unwrap().as_ref(), b"456"); - assert!(range.clone().slice(6, 0).unwrap().is_empty()); - assert!(range.clone().slice(6, 1).is_none()); - assert!(range.clone().slice(usize::MAX, 1).is_none()); - assert!(range.slice(1, usize::MAX).is_none()); - let mut first = acquired.with_range(0, 3).unwrap(); - let next = acquired.with_range(3, 4).unwrap(); - let gap = acquired.with_range(8, 1).unwrap(); - assert!(!first.extend_contiguous(&gap)); - assert!(first.extend_contiguous(&next)); - drop(next); - assert_eq!(first.as_ref(), b"0123456"); -} - -#[test] -fn failed_acquisition_releases_every_successful_prefix_without_touching_other_reads() { - let mut producer = TCache::producer("", 1 << 18); - let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); - let mut columns = TCache::producer("", 1 << 18); - let mut reader = Box::new(columns.cache_ref().retained_random_access("").unwrap()); - let gossip = write(&mut producer, b"gossip"); - let column = write(&mut columns, b"column"); - let gossip_guard = consumer.acquire_strict(gossip).unwrap(); - let column_guard = reader.acquire_strict(column).unwrap(); - let shared = columns - .sub_reservation(SubLayout { parts: 2, first_len: 4, second_len: 2 }, b"", b"") - .unwrap(); - columns - .view_sub_reservation(shared) - .unwrap() - .claim(0) - .unwrap() - .write(b"cell", b"pf") - .unwrap() - .acquire(&mut reader) - .unwrap() - .accept() - .unwrap(); - let now = Instant::now(); - let segments = [ - GossipSegment::Framing { offset: 0, length: 1 }, - GossipSegment::Gossip { read: gossip, offset: 1, length: 3 }, - GossipSegment::DataColumns { read: column, offset: 0, length: 4 }, - GossipSegment::Shared { reservation: shared, part: 0, second: false, offset: 0, length: 4 }, - GossipSegment::Shared { reservation: shared, part: 0, second: true, offset: 0, length: 2 }, - GossipSegment::Gossip { read: gossip, offset: 1, length: 3 }, - GossipSegment::Framing { offset: 0, length: 1 }, - ]; - for failed in 0..segments.len() { - let mut descriptors = segments; - descriptors[failed] = GossipSegment::DataColumns { read: column, offset: 100, length: 1 }; - let frame = GossipFrameRef::write( - &mut producer, - now + Duration::from_secs(1), - b"f", - descriptors.into_iter(), - ) - .unwrap(); - assert!( - frame - .acquire(&mut consumer, now) - .unwrap() - .acquire_segments(&mut consumer, Some(&mut reader)) - .is_none() - ); - assert_eq!(consumer.active_count(), 1, "failed descriptor {failed}"); - assert_eq!(reader.active_count(), 1, "failed descriptor {failed}"); - assert_eq!(gossip_guard.buffer().unwrap().0, b"gossip"); - assert_eq!(column_guard.buffer().unwrap().0, b"column"); - } - drop(gossip_guard); - drop(column_guard); - assert_eq!(consumer.active_count(), 0); - assert_eq!(reader.active_count(), 0); -} - -#[test] -fn handoff_transfers_counts_and_frame_drop_releases_only_the_remainder() { - let mut producer = TCache::producer("", 1 << 18); - let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); - let mut columns = TCache::producer("", 1 << 18); - let mut reader = Box::new(columns.cache_ref().retained_random_access("").unwrap()); - let gossip = write(&mut producer, b"gossip"); - let column = write(&mut columns, b"column"); - let now = Instant::now(); - let reference = GossipFrameRef::write( - &mut producer, - now + Duration::from_secs(1), - b"f", - [ - GossipSegment::Framing { offset: 0, length: 1 }, - GossipSegment::Gossip { read: gossip, offset: 0, length: 6 }, - GossipSegment::DataColumns { read: column, offset: 0, length: 6 }, - GossipSegment::DataColumns { read: column, offset: 0, length: 6 }, - GossipSegment::Framing { offset: 0, length: 1 }, - ] - .into_iter(), - ) - .unwrap(); - let mut frame = reference - .acquire(&mut consumer, now) - .unwrap() - .acquire_segments(&mut consumer, Some(&mut reader)) - .unwrap(); - assert_eq!(consumer.active_count(), 2); - assert_eq!(reader.active_count(), 2); - assert!(matches!(frame.take_next(), Some(AcquiredGossipSegment::Framing(_)))); - let Some(AcquiredGossipSegment::Data(gossip_range)) = frame.take_next() else { panic!() }; - let Some(AcquiredGossipSegment::Data(column_range)) = frame.take_next() else { panic!() }; - assert_eq!(consumer.active_count(), 2); - assert_eq!(reader.active_count(), 2); - reader.advance_retention(columns.next_seq()); - drop(frame); - assert_eq!(consumer.active_count(), 1); - assert_eq!(reader.active_count(), 1); - assert_eq!(gossip_range.as_ref(), b"gossip"); - assert_eq!(column_range.as_ref(), b"column"); - drop(gossip_range); - drop(column_range); - assert_eq!(consumer.active_count(), 0); - assert_eq!(reader.active_count(), 0); -} - -#[test] -fn shared_handoff_survives_closure_without_exposing_unverified_gaps() { - let mut producer = TCache::producer("", 1 << 18); - let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); - let mut columns = TCache::producer("", 1 << 18); - let mut reader = Box::new(columns.cache_ref().retained_random_access("").unwrap()); - let shared = columns - .sub_reservation(SubLayout { parts: 3, first_len: 4, second_len: 2 }, b"hdr", b"mid") - .unwrap(); - for part in [0, 2] { - columns - .view_sub_reservation(shared) - .unwrap() - .claim(part) - .unwrap() - .write(&[part as u8; 4], &[part as u8 + 10; 2]) - .unwrap() - .acquire(&mut reader) - .unwrap() - .accept() - .unwrap(); - } - let now = Instant::now(); - let reference = GossipFrameRef::write( - &mut producer, - now + Duration::from_secs(1), - b"", - [ - GossipSegment::Shared { - reservation: shared, - part: 0, - second: false, - offset: 1, - length: 3, - }, - GossipSegment::Shared { - reservation: shared, - part: 2, - second: false, - offset: 0, - length: 4, - }, - GossipSegment::Shared { - reservation: shared, - part: 0, - second: true, - offset: 0, - length: 2, - }, - GossipSegment::Shared { - reservation: shared, - part: 2, - second: true, - offset: 1, - length: 1, - }, - ] - .into_iter(), - ) - .unwrap(); - let mut frame = reference - .acquire(&mut consumer, now) - .unwrap() - .acquire_segments(&mut consumer, Some(&mut reader)) - .unwrap(); - assert_eq!(reader.active_count(), 4); - columns.view_sub_reservation(shared).unwrap().close(); - reader.advance_retention(columns.next_seq()); - assert!( - reference - .acquire(&mut consumer, now) - .unwrap() - .acquire_segments(&mut consumer, Some(&mut reader)) - .is_none() - ); - assert_eq!(reader.active_count(), 4); - for (index, expected) in [&[0; 3][..], &[2; 4], &[10; 2], &[12; 1]].into_iter().enumerate() { - let Some(AcquiredGossipSegment::Data(range)) = frame.take_next() else { panic!() }; - assert_eq!(reader.active_count(), 4 - index); - assert_eq!(range.as_ref(), expected); - drop(range); - assert_eq!(reader.active_count(), 3 - index); - } - assert!(frame.take_next().is_none()); - drop(frame); - assert_eq!(consumer.active_count(), 0); - assert_eq!(reader.active_count(), 0); -} - -#[test] -fn coalescing_transfers_one_pin_and_releases_redundant_pins() { - let mut producer = TCache::producer("", 1 << 18); - let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); - let source = write(&mut producer, b"0123456789"); - let now = Instant::now(); - let reference = GossipFrameRef::write( - &mut producer, - now + Duration::from_secs(1), - b"", - [ - GossipSegment::Gossip { read: source, offset: 1, length: 3 }, - GossipSegment::Gossip { read: source, offset: 4, length: 3 }, - GossipSegment::Gossip { read: source, offset: 7, length: 2 }, - ] - .into_iter(), - ) - .unwrap(); - let mut frame = reference - .acquire(&mut consumer, now) - .unwrap() - .acquire_segments(&mut consumer, None) - .unwrap(); - assert_eq!(consumer.active_count(), 4); - let Some(AcquiredGossipSegment::Data(range)) = frame.take_next() else { panic!() }; - assert_eq!(consumer.active_count(), 2); - assert!(frame.take_next().is_none()); - drop(frame); - assert_eq!(consumer.active_count(), 1); - assert_eq!(range.as_ref(), b"12345678"); - drop(range); - assert_eq!(consumer.active_count(), 0); -} - -#[test] -fn unwinding_releases_untransferred_pins_but_not_the_handed_off_read() { - let mut producer = TCache::producer("", 1 << 18); - let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); - let source = write(&mut producer, b"data"); - let now = Instant::now(); - let reference = GossipFrameRef::write( - &mut producer, - now + Duration::from_secs(1), - b"", - [GossipSegment::Gossip { read: source, offset: 0, length: 4 }; 3].into_iter(), - ) - .unwrap(); - let mut handed_off = None; - let result = catch_unwind(AssertUnwindSafe(|| { - let mut frame = reference - .acquire(&mut consumer, now) - .unwrap() - .acquire_segments(&mut consumer, None) - .unwrap(); - let Some(AcquiredGossipSegment::Data(range)) = frame.take_next() else { panic!() }; - handed_off = Some(range); - panic!("abort a partially handed-off frame"); - })); - assert!(result.is_err()); - assert_eq!(consumer.active_count(), 1); - assert_eq!(handed_off.as_ref().unwrap().as_ref(), b"data"); - drop(handed_off); - assert_eq!(consumer.active_count(), 0); -} diff --git a/crates/control/src/lib.rs b/crates/control/src/lib.rs index 0df2a7ae..7f336888 100644 --- a/crates/control/src/lib.rs +++ b/crates/control/src/lib.rs @@ -1,6 +1,5 @@ pub mod cell_ingress; pub mod cluster; -pub mod cell_ingress; mod counters; pub mod sync_engine; mod tile; diff --git a/crates/network/src/lib.rs b/crates/network/src/lib.rs index de8acc25..076fe04b 100644 --- a/crates/network/src/lib.rs +++ b/crates/network/src/lib.rs @@ -41,16 +41,16 @@ silver_common::declare_counters! { GossipStallDisconnect, // RPC codecs currently retained in the network-tile-wide free list. RpcCodecPoolIdle, - GossipSegmentedAdmitted, - GossipSegmentedRejected, - GossipSegmentedCapacity, - GossipSegmentedSegments, - GossipSegmentedOwnerAllocations, - GossipSegmentedFrames, + CacheSegmentedAdmitted, + CacheSegmentedRejected, + CacheSegmentedCapacity, + CacheSegmentedSegments, + CacheSegmentedOwnerAllocations, + CacheSegmentedFrames, // Reserved owner slots, including owners already held by Quinn. - GossipSegmentedOwners, + CacheSegmentedOwners, // Reserved ranges per recipient, including descriptors; not unique cache backing bytes. - GossipSegmentedRetainedBytes, + CacheSegmentedRetainedBytes, } } diff --git a/crates/network/src/p2p/mod.rs b/crates/network/src/p2p/mod.rs index ddc2a9aa..0780100e 100644 --- a/crates/network/src/p2p/mod.rs +++ b/crates/network/src/p2p/mod.rs @@ -18,8 +18,8 @@ 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::{ - CacheFrameRef, 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::{ @@ -381,7 +381,7 @@ impl P2p { pub fn enqueue_segmented_gossip( &mut self, peer_id: usize, - frame: GossipFrameRef, + frame: CacheFrameRef, context: &mut Context, ) -> SendResult { match self.peers.get_mut(&ConnectionHandle(peer_id)) { diff --git a/crates/network/src/p2p/quic/gossip_frame.rs b/crates/network/src/p2p/quic/gossip_frame.rs index db8a2887..8ee15335 100644 --- a/crates/network/src/p2p/quic/gossip_frame.rs +++ b/crates/network/src/p2p/quic/gossip_frame.rs @@ -2,7 +2,7 @@ use std::{cell::Cell, ptr::NonNull, time::Instant}; use bytes::Bytes; use silver_common::{ - AcquiredGossipFrame, AcquiredGossipSegment, AcquiredRange, GossipFrameView, TRead, + AcquiredCacheFrame, AcquiredCacheSegment, AcquiredRange, CacheFrameView, TRead, }; use super::{Leased, leased::OutboundLeaseWheel}; @@ -46,7 +46,7 @@ impl SegmentedGossipLimits { pub(crate) fn acquire( &self, - view: GossipFrameView, + view: CacheFrameView, context: &mut Context, wheel: &OutboundLeaseWheel, now: Instant, @@ -57,7 +57,7 @@ impl SegmentedGossipLimits { owners > self.max_owners.saturating_sub(self.owners.get()) || self.frames.get() >= self.max_frames { - NetworkCounters::GossipSegmentedCapacity.inc(); + NetworkCounters::CacheSegmentedCapacity.inc(); return None; } self.frames.set(self.frames.get() + 1); @@ -68,15 +68,15 @@ impl SegmentedGossipLimits { &mut context.gossip_consumer, context.data_columns_consumer.as_deref_mut(), )?; - NetworkCounters::GossipSegmentedAdmitted.inc(); - NetworkCounters::GossipSegmentedSegments.add(frame.segment_count() as u64); + 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::GossipSegmentedFrames.set(self.frames.get() as u64); - NetworkCounters::GossipSegmentedOwners.set(self.owners.get() as u64); - NetworkCounters::GossipSegmentedRetainedBytes.set(self.retained_bytes.get() as u64); + 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); } } @@ -90,7 +90,7 @@ impl Drop for SegmentedGossipLimits { #[derive(Debug)] pub(crate) struct SegmentedFrame { - segments: Leased, + segments: Leased, budget: FrameBudget, } @@ -118,13 +118,13 @@ impl SegmentedWriter { if self.current.is_empty() { let SegmentedFrame { segments, budget } = &mut self.frame; self.current = match segments.take_next()? { - AcquiredGossipSegment::Framing(range) => { + AcquiredCacheSegment::Framing(range) => { if self.descriptor.is_empty() { self.descriptor = budget.owner(segments.child(segments.descriptor_range())); } self.descriptor.slice(range) } - AcquiredGossipSegment::Data(range) => budget.owner(segments.child(range)), + AcquiredCacheSegment::Data(range) => budget.owner(segments.child(range)), }; } Some(&mut self.current) @@ -154,7 +154,7 @@ impl FrameBudget { self.owners -= 1; self.bytes -= data.len(); let owner = SegmentOwner { data, limits: self.limits }; - NetworkCounters::GossipSegmentedOwnerAllocations.inc(); + NetworkCounters::CacheSegmentedOwnerAllocations.inc(); Bytes::from_owner(owner) } } diff --git a/crates/network/src/p2p/quic/gossip_frame/tests.rs b/crates/network/src/p2p/quic/gossip_frame/tests.rs index a0541b4b..9610011f 100644 --- a/crates/network/src/p2p/quic/gossip_frame/tests.rs +++ b/crates/network/src/p2p/quic/gossip_frame/tests.rs @@ -8,7 +8,7 @@ use std::{ use quinn_proto::StreamId; use silver_common::{ - AcquiredWithOffset, GossipFrameRef, GossipSegment, P2pStreamId, StreamProtocol, SubLayout, + AcquiredWithOffset, CacheFrameRef, CacheSegment, P2pStreamId, StreamProtocol, SubLayout, SubReservationRef, TCache, TCacheProducer, TProducer, }; @@ -55,9 +55,13 @@ struct Harness { 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 { @@ -69,6 +73,12 @@ impl Harness { 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)), @@ -78,7 +88,7 @@ impl Harness { } } - fn assembly(&mut self, accept: bool) -> (GossipFrameRef, SubReservationRef) { + 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"") @@ -98,20 +108,20 @@ impl Harness { .accept() .unwrap(); } - let frame = GossipFrameRef::write( + let frame = CacheFrameRef::write( &mut self.gossip, self.now + Duration::from_secs(1), b"head", [ - GossipSegment::Framing { offset: 0, length: 4 }, - GossipSegment::Shared { + CacheSegment::Framing { offset: 0, length: 4 }, + CacheSegment::Shared { reservation: reference, part: 0, second: false, offset: 0, length: 64, }, - GossipSegment::Shared { + CacheSegment::Shared { reservation: reference, part: 0, second: true, @@ -125,7 +135,7 @@ impl Harness { (frame, reference) } - fn acquire(&mut self, frame: GossipFrameRef) -> Option { + 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) } @@ -160,6 +170,10 @@ impl MockIo { } 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]); @@ -371,11 +385,11 @@ fn adjacent_ranges_share_one_owner_without_gathering() { let mut reservation = h.gossip.reserve(64, true).unwrap(); reservation.write_all(&[0xab; 64]).unwrap(); let read = reservation.read(); - let reference = GossipFrameRef::write( + let reference = CacheFrameRef::write( &mut h.gossip, h.now + Duration::from_secs(1), b"", - [GossipSegment::Gossip { read, offset: 2, length: 4 }, GossipSegment::Gossip { + [CacheSegment::Gossip { read, offset: 2, length: 4 }, CacheSegment::Gossip { read, offset: 6, length: 8, @@ -399,19 +413,20 @@ fn adjacent_ranges_share_one_owner_without_gathering() { #[test] fn framing_fragments_share_one_lazy_owner_until_the_last_ack() { let mut h = Harness::new(); - let reference = GossipFrameRef::write( + let reference = CacheFrameRef::write( &mut h.gossip, h.now + Duration::from_secs(1), b"abcdef", [ - GossipSegment::Framing { offset: 0, length: 2 }, - GossipSegment::Framing { offset: 4, length: 2 }, - GossipSegment::Framing { offset: 2, length: 2 }, + 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(); @@ -457,11 +472,11 @@ fn dropping_before_the_prefix_allocates_no_owners_and_releases_the_lease() { #[test] fn partial_length_prefix_and_mixed_frames_preserve_boundaries() { let mut h = Harness::new(); - let reference = GossipFrameRef::write( + let reference = CacheFrameRef::write( &mut h.gossip, h.now + Duration::from_secs(1), &[0xab; 300], - [GossipSegment::Framing { offset: 0, length: 300 }].into_iter(), + [CacheSegment::Framing { offset: 0, length: 300 }].into_iter(), ) .unwrap(); let frame = h.acquire(reference).unwrap(); diff --git a/crates/network/src/p2p/quic/peer.rs b/crates/network/src/p2p/quic/peer.rs index 141667a7..0f363966 100644 --- a/crates/network/src/p2p/quic/peer.rs +++ b/crates/network/src/p2p/quic/peer.rs @@ -13,7 +13,7 @@ use quinn_proto::{ VarInt, }; use silver_common::{ - GossipFrameRef, P2pConnectionStats, P2pStreamId, PeerId, StreamProtocol, TRead, + CacheFrameRef, P2pConnectionStats, P2pStreamId, PeerId, StreamProtocol, TRead, rpc_rate_limit::RpcRateLimitSet, }; @@ -23,7 +23,7 @@ use crate::{ NetEvent, context::Context, quic::{ - OutboundGossip, SegmentedGossipLimits, SendResult, leased::OutboundLeaseWheel, + Leased, OutboundGossip, SegmentedGossipLimits, SendResult, leased::OutboundLeaseWheel, stream::StreamIoImpl, }, streams::{ @@ -177,7 +177,7 @@ impl Peer { pub(crate) fn send_segmented_gossip( &mut self, - frame: GossipFrameRef, + frame: CacheFrameRef, context: &mut Context, limits: &SegmentedGossipLimits, rpc_codec_pool: &mut RpcCodecPool, @@ -194,7 +194,7 @@ impl Peer { .ok() .and_then(|view| limits.acquire(view, context, &self.outbound_lease_wheel, now)); let Some(frame) = acquired else { - crate::NetworkCounters::GossipSegmentedRejected.inc(); + crate::NetworkCounters::CacheSegmentedRejected.inc(); return SendResult::MessageDropped; }; self.queue_gossip(OutboundGossip::Segmented(frame)) @@ -2182,14 +2182,14 @@ mod tests { write.flush().unwrap(); write.read() }; - let frame = GossipFrameRef::write( + let frame = CacheFrameRef::write( &mut client_h.gossip_out_producer, Instant::now() + Duration::from_secs(1), &payload[..2], [ - GossipSegment::Framing { offset: 0, length: 2 }, - GossipSegment::DataColumns { read, offset: 0, length: 2 }, - GossipSegment::DataColumns { read, offset: 2, length: 3 }, + CacheSegment::Framing { offset: 0, length: 2 }, + CacheSegment::DataColumns { read, offset: 0, length: 2 }, + CacheSegment::DataColumns { read, offset: 2, length: 3 }, ] .into_iter(), ) @@ -2235,11 +2235,11 @@ mod tests { 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 = GossipFrameRef::write( + let frame = CacheFrameRef::write( &mut h.gossip_out_producer, Instant::now() + Duration::from_secs(1), b"payload", - [GossipSegment::Framing { offset: 0, length: 7 }].into_iter(), + [CacheSegment::Framing { offset: 0, length: 7 }].into_iter(), ) .unwrap(); let peer = &mut pair.client_peer; diff --git a/crates/network/src/p2p/streams/gossip_in.rs b/crates/network/src/p2p/streams/gossip_in.rs index 089f9692..82f39dd8 100644 --- a/crates/network/src/p2p/streams/gossip_in.rs +++ b/crates/network/src/p2p/streams/gossip_in.rs @@ -194,7 +194,7 @@ mod tests { use std::net::SocketAddr; use quinn_proto::StreamId; - use silver_common::{StreamProtocol, TCache}; + use silver_common::{StreamProtocol, TCache, TRead}; use super::*; use crate::p2p::streams::AcquiredRpcOutbound; diff --git a/crates/network/src/p2p/streams/mod.rs b/crates/network/src/p2p/streams/mod.rs index b462a8d0..4f67ee38 100644 --- a/crates/network/src/p2p/streams/mod.rs +++ b/crates/network/src/p2p/streams/mod.rs @@ -3,7 +3,7 @@ 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}; +use silver_common::{AcquiredWithOffset, TCacheError, TRead}; use thiserror::Error; use crate::p2p::{ diff --git a/crates/network/src/p2p/streams/rpc/request_in.rs b/crates/network/src/p2p/streams/rpc/request_in.rs index 6cb31156..cb5eebf5 100644 --- a/crates/network/src/p2p/streams/rpc/request_in.rs +++ b/crates/network/src/p2p/streams/rpc/request_in.rs @@ -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 40c8dde3..420f5108 100644 --- a/crates/network/src/p2p/streams/rpc/response_in.rs +++ b/crates/network/src/p2p/streams/rpc/response_in.rs @@ -378,7 +378,7 @@ mod tests { Ok(data.as_ref().len()) } - fn cluster_next(&mut self) -> Option> { + fn cluster_next(&mut self) -> Option> { None } }