From 56008a904e46594106f7285b244bb0191f9905a1 Mon Sep 17 00:00:00 2001 From: vladimir-ea Date: Fri, 11 Sep 2026 14:48:32 +0100 Subject: [PATCH 1/5] 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 3702aa9f5517d4cf20e619b845e9605869780c03 Mon Sep 17 00:00:00 2001 From: vladimir-ea Date: Fri, 11 Sep 2026 14:49:28 +0100 Subject: [PATCH 2/5] 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 ea227c1aae7e2e8e61de5acab7248119997c38d3 Mon Sep 17 00:00:00 2001 From: vladimir-ea Date: Fri, 11 Sep 2026 14:50:30 +0100 Subject: [PATCH 3/5] 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/network/src/p2p/streams/gossip_in.rs | 2 +- crates/network/src/p2p/streams/mod.rs | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/crates/network/src/p2p/streams/gossip_in.rs b/crates/network/src/p2p/streams/gossip_in.rs index 82f39dd8..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; diff --git a/crates/network/src/p2p/streams/mod.rs b/crates/network/src/p2p/streams/mod.rs index 4f67ee38..b462a8d0 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, TRead}; +use silver_common::{AcquiredWithOffset, TCacheError}; use thiserror::Error; use crate::p2p::{ From 28e75fa0fb2cbe40818b4719ae17aba152aee955 Mon Sep 17 00:00:00 2001 From: vladimir-ea Date: Fri, 11 Sep 2026 14:50:39 +0100 Subject: [PATCH 4/5] partial payloads: wire cell store into tiles Main creates the data_columns tcache sized from the blob schedule, with retained consumers for the columns and network tiles. Control owns the store through CellIngress; the columns and network tiles advance their retention boundaries from RetentionEvent. Co-Authored-By: Claude Fable 5 --- crates/bin/src/main.rs | 32 ++++++++++++++++++++++++-------- crates/columns/src/tile.rs | 14 ++++++++++++++ crates/control/src/tile.rs | 29 ++++++++++++++++++++++++++++- 3 files changed, 66 insertions(+), 9 deletions(-) diff --git a/crates/bin/src/main.rs b/crates/bin/src/main.rs index 8468c94d..da835165 100644 --- a/crates/bin/src/main.rs +++ b/crates/bin/src/main.rs @@ -17,12 +17,15 @@ use rand::RngCore; use silver_application_boundary::ApplicationBoundaryTile; use silver_beacon_state::{BeaconStateTile, SlotTicker}; use silver_beacon_state_data::{BeaconState, SLOTS_PER_EPOCH}; -use silver_columns::tile::{ColumnConsumers, DataColumnsTile}; +use silver_columns::{ + cell_store::CellStoreConfig, + tile::{ColumnConsumers, DataColumnsTile}, +}; #[cfg(feature = "alloc-profile")] use silver_common::metrics::CountingAllocator; use silver_common::{ - APP_NAME, Enr, ProtoIdentify, SilverSpine, TCache, TCacheProducer, profiler::enable_profiler, - tracing::initialise_tracing_log, + APP_NAME, Enr, ProtoIdentify, SilverSpine, TCache, TCacheProducer, + cells::GOSSIP_DELIVERY_RETENTION, profiler::enable_profiler, tracing::initialise_tracing_log, }; use silver_config::Config; use silver_control::{Controller, cluster::AttestationClusterConfig, sync_engine::SyncEngine}; @@ -150,6 +153,7 @@ fn main() -> Result<(), Box> { tracing::info!(enr = local_enr.to_base64(), "local ENR on startup"); let chain_config = config.chain_config(); + let spec = Arc::new(chain_config.spec.clone()); sleep_until_genesis(chain_config.genesis_unix_secs); let ticker = SlotTicker::new( chain_config.genesis_unix_secs, @@ -216,8 +220,7 @@ fn main() -> Result<(), Box> { trusted_ips, ); let identify = config.identify()?; - let p2p_context = Context { - data_columns_consumer: None, + let mut p2p_context = Context { gossip_producer: incoming_gossip_producer, gossip_consumer: outgoing_gossip_producer .cache_ref() @@ -228,6 +231,7 @@ fn main() -> Result<(), Box> { cluster_nodes: cluster_nodes.map(ClusterNodes::new), cluster_inbound_producer, cluster_outbound_consumer, + data_columns_consumer: None, }; let now = Instant::now(); @@ -249,6 +253,16 @@ fn main() -> Result<(), Box> { } } + let cell_config = + CellStoreConfig::new(spec.clone(), das_custody_groups, GOSSIP_DELIVERY_RETENTION) + .map_err(|error| format!("cell store configuration: {error:?}"))?; + let data_columns_producer = TCache::producer("data_columns", cell_config.cache_capacity()); + let columns_consumer = + data_columns_producer.cache_ref().retained_random_access("columns_cells")?; + p2p_context.data_columns_consumer = + Some(Box::new(data_columns_producer.cache_ref().retained_random_access("network_cells")?)); + let (cell_slot, cell_slot_start) = ticker.current_slot_start(); + let network_tile = NetworkTile::new(discv5_addr, discv5, p2p_addr, p2p_endpoint, p2p_context)?; let (checkpoint, checkpoint_pubkeys) = load_checkpoint(&config)?; @@ -256,8 +270,6 @@ fn main() -> Result<(), Box> { tracing::info!("booting from local checkpoint: {booting_from_local_checkpoint}"); - let spec = Arc::new(chain_config.spec.clone()); - let gossip_handler = GossipHandler::new( incoming_gossip_consumer, ssz_gossip_producer, @@ -289,6 +301,9 @@ fn main() -> Result<(), Box> { spec.clone(), ), )?; + control_tile = control_tile + .with_data_columns_cache(cell_config, data_columns_producer, cell_slot, cell_slot_start) + .map_err(|error| format!("cell store construction: {error:?}"))?; control_tile.set_pending_subnet_topics( silver_common::attnet_subnets(subnets) .map(silver_common::GossipTopic::BeaconAttestation) @@ -349,7 +364,8 @@ fn main() -> Result<(), Box> { chain_config.slot_duration(), chain_config.playload_lookahead(), ), - ); + ) + .with_data_columns_consumer(columns_consumer); let beacon_api_binds = config.beacon_api_bind().iter().map(String::as_str).map(Bind::parse).collect::>(); diff --git a/crates/columns/src/tile.rs b/crates/columns/src/tile.rs index 77e4d313..42d1e241 100644 --- a/crates/columns/src/tile.rs +++ b/crates/columns/src/tile.rs @@ -15,6 +15,7 @@ use silver_common::{ EngineResp, GossipTopic, IngestionTime, NewGossipMsg, Origin, P2pStreamId, PeerEvent, RequestId, RpcInbound, RpcSeverity, SilverSpine, SilverSpineProducers, StreamProtocol, SyncNeed, SyncUpdate, TCacheRead, TProducer, TRandomAccess, TRead, Wheel, + cells::RetentionEvent, column_util::{self as util, KzgScratch}, ssz_view::{NUMBER_OF_COLUMNS, SignedBeaconBlockView, StatusView}, ticker::SlotTicker, @@ -88,6 +89,7 @@ pub struct DataColumnsTile { // Declared last so it drops last: a parked column's read releases through // the consumer it was acquired from. consumers: ColumnConsumers, + data_columns_consumer: Option>, } impl DataColumnsTile { @@ -114,9 +116,16 @@ impl DataColumnsTile { el_fetcher: ElBlobFetcher::new(engine_resp_consumer), el_column_producer, kzg_scratch: KzgScratch::default(), + data_columns_consumer: None, } } + pub fn with_data_columns_consumer(mut self, consumer: TRandomAccess) -> Self { + assert!(consumer.is_retained()); + self.data_columns_consumer = Some(Box::new(consumer)); + self + } + #[timed] fn beacon_block( &mut self, @@ -676,6 +685,11 @@ impl Tile for DataColumnsTile { fn loop_body(&mut self, adapter: &mut SpineAdapter) { self.consumers.free(); + if let Some(consumer) = &mut self.data_columns_consumer { + adapter.consume(|event: RetentionEvent, _| { + consumer.advance_retention(event.retain_from); + }); + } adapter.consume(|gossip: NewGossipMsg, producers| match gossip.topic { silver_common::GossipTopic::BeaconBlock if self.sync_state.is_synced() => { diff --git a/crates/control/src/tile.rs b/crates/control/src/tile.rs index 269f6f21..012b10f5 100644 --- a/crates/control/src/tile.rs +++ b/crates/control/src/tile.rs @@ -1,6 +1,7 @@ use std::time::{Duration, Instant}; use flux::{spine::SpineAdapter, tile::Tile}; +use silver_columns::cell_store::{CellStoreConfig, StoreError}; use silver_common::{ BeaconApiRequest, BeaconStateEvent, GossipTopic, LOCAL_GOSSIP_STREAM_ID, Nanos, P2pSend, PeerControl, PeerEvent, PeerStats, RpcInbound, RpcOutbound, RpcRequest, RpcRequestOutbound, @@ -18,6 +19,10 @@ use crate::{ }; mod attestation_cluster; +use crate::{ + cell_ingress::CellIngress, + sync_engine::{SyncAction, SyncEngine}, +}; const PEER_PERSIST_INTERVAL: Duration = Duration::from_secs(300); @@ -52,6 +57,7 @@ pub struct Controller { /// meshes would earn P3 deficit at peers since nothing validates or /// forwards until then. Drained into the PM on the first transition. pending_subnet_topics: Vec, + cell_ingress: Option, } impl Controller { @@ -90,7 +96,19 @@ impl Controller { last_peer_persist: now, auto_ping: true, pending_subnet_topics: Vec::new(), - }) + cell_ingress: None, + } + } + + pub fn with_data_columns_cache( + mut self, + config: CellStoreConfig, + producer: TProducer, + slot: u64, + slot_start: Instant, + ) -> Result { + self.cell_ingress = Some(CellIngress::new(config, producer, slot, slot_start)?); + Ok(self) } pub fn set_pending_subnet_topics(&mut self, topics: Vec) { @@ -116,6 +134,9 @@ impl Controller { fn handle_latest_status(&mut self, latest_status_event: Option<([u8; 92], u64, u64)>) -> bool { if let Some((ssz, latest_block_slot, wall_slot)) = latest_status_event { + if let Some(ingress) = &mut self.cell_ingress { + ingress.set_min_slot(StatusView::finalized_epoch(&ssz) * SLOTS_PER_EPOCH); + } tracing::debug!(wall_slot, latest_block_slot, "new status set"); // PM still tracks our Status (peer-Status validation) + applied head // (custody-peer eligibility); the wall slot is the engine's only. @@ -158,6 +179,12 @@ impl Tile for Controller { let now = Instant::now(); self.rpc_ssz_consumer.free(); self.attestation_cluster.free(); + if let Some(ingress) = &mut self.cell_ingress { + ingress.spin(now, &adapter.producers); + adapter.consume(|event: CellStoreEvent, producers| { + ingress.handle(event, now, producers); + }); + } // Local status must land before the sync drive below: issuance is // capped against the imported head, and a one-loop-stale watermark From b3a248d8f7df8265bd061add218e66bc1b9ab530 Mon Sep 17 00:00:00 2001 From: vladimir-ea Date: Tue, 15 Sep 2026 10:46:35 +0100 Subject: [PATCH 5/5] rebase snafu --- .../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/control/src/tile.rs | 12 +- .../src/p2p/quic/gossip_frame/tests.rs | 1 - crates/network/src/p2p/streams/gossip_in.rs | 2 +- crates/network/src/p2p/streams/mod.rs | 2 +- 8 files changed, 7 insertions(+), 896 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/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/control/src/tile.rs b/crates/control/src/tile.rs index 012b10f5..990fdb58 100644 --- a/crates/control/src/tile.rs +++ b/crates/control/src/tile.rs @@ -5,8 +5,9 @@ use silver_columns::cell_store::{CellStoreConfig, StoreError}; use silver_common::{ BeaconApiRequest, BeaconStateEvent, GossipTopic, LOCAL_GOSSIP_STREAM_ID, Nanos, P2pSend, PeerControl, PeerEvent, PeerStats, RpcInbound, RpcOutbound, RpcRequest, RpcRequestOutbound, - RpcResponse, RpcResponseInbound, SilverSpine, SilverSpineProducers, SyncNeed, SyncUpdate, - TMultiProducer, TProducer, TRandomAccess, + RpcResponse, RpcResponseInbound, SLOTS_PER_EPOCH, SilverSpine, SilverSpineProducers, SyncNeed, + SyncUpdate, TMultiProducer, TProducer, TRandomAccess, + cells::CellStoreEvent, ssz_view::{METADATA_SIZE, STATUS_V2_SIZE, StatusView}, }; use silver_gossip::{GossipHandler, GossipHandlerEvent}; @@ -14,15 +15,12 @@ use silver_peer::PeerManager; use self::attestation_cluster::AttestationClusterHandler; use crate::{ + cell_ingress::CellIngress, cluster::{AttestationClusterConfig, ClusterError}, sync_engine::{SyncAction, SyncEngine}, }; mod attestation_cluster; -use crate::{ - cell_ingress::CellIngress, - sync_engine::{SyncAction, SyncEngine}, -}; const PEER_PERSIST_INTERVAL: Duration = Duration::from_secs(300); @@ -97,7 +95,7 @@ impl Controller { auto_ping: true, pending_subnet_topics: Vec::new(), cell_ingress: None, - } + }) } pub fn with_data_columns_cache( diff --git a/crates/network/src/p2p/quic/gossip_frame/tests.rs b/crates/network/src/p2p/quic/gossip_frame/tests.rs index 9610011f..0672f2e3 100644 --- a/crates/network/src/p2p/quic/gossip_frame/tests.rs +++ b/crates/network/src/p2p/quic/gossip_frame/tests.rs @@ -426,7 +426,6 @@ fn framing_fragments_share_one_lazy_owner_until_the_last_ack() { ) .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(); 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::{