diff --git a/crates/common/src/lib.rs b/crates/common/src/lib.rs index 19250510..c1f8d929 100644 --- a/crates/common/src/lib.rs +++ b/crates/common/src/lib.rs @@ -1,5 +1,10 @@ extern crate self as silver_common; +pub use spine::{ + AcquiredCacheFrame, AcquiredCacheSegment, CacheFrameError, CacheFrameRef, CacheFrameSegment, + CacheFrameView, CacheSegment, MAX_CACHE_SEGMENTS, +}; + pub use crate::{ error::Error, gossip::{ @@ -16,30 +21,31 @@ pub use crate::{ }, request::{DataKind, Origin, RequestId, Scope, SyncRequest}, spine::{ - ALL_PROTOCOLS, AcquiredRead as TRead, AcquiredWithOffset, AgentString, BeaconApiRequest, - BeaconApiResponse, BeaconStateEvent, BlockSource, BlockStage, ClusterIn, ClusterMsgIn, - ClusterMsgOut, ColumnSource, Consumer as TConsumer, DataColumnsEvent, ELSyncStatus, - EngineFcuReq, EngineFcuResp, EngineGetBlobsReq, EngineGetBlobsResp, - EngineGetPayloadBodiesByHashReq, EngineGetPayloadBodiesByRangeReq, - EngineGetPayloadBodiesResp, EngineGetPayloadReq, EngineGetPayloadResp, EngineHealthEvent, - EngineNewPayloadEnvelopeReq, EngineNewPayloadReq, EngineNewPayloadResp, - EnginePreparePayloadReq, EngineReq, EngineResp, Error as TCacheError, GossipMsgIn, - GossipMsgOut, IpBytes, LOCAL_GOSSIP_STREAM_ID, LocalAttestationFailure, + ALL_PROTOCOLS, AcquiredRange, AcquiredRead as TRead, AcquiredSubReservation, + AcquiredWithOffset, AgentString, BeaconApiRequest, BeaconApiResponse, BeaconStateEvent, + BlockSource, BlockStage, ClusterIn, ClusterMsgIn, ClusterMsgOut, ColumnSource, + Consumer as TConsumer, DataColumnsEvent, ELSyncStatus, EngineFcuReq, EngineFcuResp, + EngineGetBlobsReq, EngineGetBlobsResp, EngineGetPayloadBodiesByHashReq, + EngineGetPayloadBodiesByRangeReq, EngineGetPayloadBodiesResp, EngineGetPayloadReq, + EngineGetPayloadResp, EngineHealthEvent, EngineNewPayloadEnvelopeReq, EngineNewPayloadReq, + EngineNewPayloadResp, EnginePreparePayloadReq, EngineReq, EngineResp, Error as TCacheError, + GossipMsgIn, GossipMsgOut, IpBytes, LOCAL_GOSSIP_STREAM_ID, LocalAttestationFailure, LocalAttestationResult, MAX_BLOBS_PER_BLOCK, MAX_PAYLOAD_BODIES_PER_REQ, MULTISTREAM_V1, MultiProducer as TMultiProducer, NewGossipMsg, P2pConnectionStats, P2pSend, P2pStreamId, PREFILL_SLOTS, PayloadValidationStatus, PeerControl, PeerEvent, PeerScores, PeerStats, - PeerStatus, PeerTopicScores, Prefill, Producer as TProducer, REJECT_RESPONSE, - RPC_PROTOCOLS, RandomAccessConsumer as TRandomAccess, ReplayBlock, + PeerStatus, PeerTopicScores, PendingSubReservation, Prefill, Producer as TProducer, + REJECT_RESPONSE, RPC_PROTOCOLS, RandomAccessConsumer as TRandomAccess, ReplayBlock, Reservation as TReservation, RpcInbound, RpcOutbound, RpcRequest, RpcRequestInbound, RpcRequestOutbound, RpcResponse, RpcResponseInbound, RpcResponseOutbound, RpcSeverity, - SilverSpine, SilverSpineProducers, StreamProtocol, SyncNeed, SyncUpdate, SyncingStrategy, - TCache, TCacheProducer, TCacheRead, TCacheRef, WithdrawalInline, + SilverSpine, SilverSpineProducers, StreamProtocol, SubLayout, SubReservation, + SubReservationError, SubReservationRef, SubReservationView, SubValidation, SubWrite, + SyncNeed, SyncUpdate, SyncingStrategy, TCache, TCacheProducer, TCacheRead, TCacheRef, + WithdrawalInline, }, util::{create_self_signed_certificate, decode_varint, encode_varint, hex32}, wheel::Wheel, wither::{CountingWitherFilter, WitherFilter}, }; - pub mod column_util; mod enr; mod error; diff --git a/crates/common/src/spine.rs b/crates/common/src/spine.rs index 5af583fb..96637914 100644 --- a/crates/common/src/spine.rs +++ b/crates/common/src/spine.rs @@ -20,8 +20,12 @@ pub use stream_protocol::{ ALL_PROTOCOLS, MULTISTREAM_V1, REJECT_RESPONSE, RPC_PROTOCOLS, StreamProtocol, }; pub use tcache::{ - AcquiredRead, AcquiredWithOffset, Consumer, Error, MultiProducer, Producer, - RandomAccessConsumer, Reservation, TCache, TCacheProducer, TCacheRead, TCacheRef, + AcquiredCacheFrame, AcquiredCacheSegment, AcquiredRange, AcquiredRead, AcquiredSubReservation, + AcquiredWithOffset, CacheFrameError, CacheFrameRef, CacheFrameSegment, CacheFrameView, + CacheSegment, Consumer, Error, MAX_CACHE_SEGMENTS, MultiProducer, PendingSubReservation, + Producer, RandomAccessConsumer, Reservation, SubLayout, SubReservation, SubReservationError, + SubReservationRef, SubReservationView, SubValidation, SubWrite, TCache, TCacheProducer, + TCacheRead, TCacheRef, }; mod messages; diff --git a/crates/common/src/spine/tcache.rs b/crates/common/src/spine/tcache.rs index 1f1819a2..ff148ab3 100644 --- a/crates/common/src/spine/tcache.rs +++ b/crates/common/src/spine/tcache.rs @@ -5,12 +5,22 @@ use std::{ ops::Deref, ptr::addr_of, slice, - sync::atomic::{AtomicU64, Ordering}, + sync::atomic::{AtomicU8, AtomicU64, Ordering}, }; -pub use consumer::{AcquiredRead, AcquiredWithOffset, Consumer, RandomAccessConsumer, TCacheRead}; +pub use cache_frame::{ + AcquiredCacheFrame, AcquiredCacheSegment, CacheFrameError, CacheFrameRef, CacheFrameSegment, + CacheFrameView, CacheSegment, MAX_CACHE_SEGMENTS, +}; +pub use consumer::{ + AcquiredRange, AcquiredRead, AcquiredWithOffset, Consumer, RandomAccessConsumer, TCacheRead, +}; use flux::{Timer, timing::Nanos, tracing}; pub use producer::{MultiProducer, Producer, Reservation, TCacheProducer}; +pub use sub_reservation::{ + AcquiredSubReservation, PendingSubReservation, SubLayout, SubReservation, SubReservationError, + SubReservationRef, SubReservationView, SubValidation, SubWrite, +}; use thiserror::Error; use crate::spine::tcache::consumer::Buckets; @@ -34,9 +44,11 @@ const fn lag_threshold(len: u32) -> u64 { (len as u64 / 10) * 9 } +mod cache_frame; mod consumer; mod metrics; mod producer; +mod sub_reservation; use metrics::TCacheMetrics; @@ -126,6 +138,12 @@ pub enum Error { UnexpectedCacheRef, #[error("stale seq: {seq} < {tail}")] StaleSeq { name: &'static str, seq: u64, tail: u64 }, + #[error("reservation is incomplete")] + Incomplete, + #[error("reservation is already committed")] + Committed, + #[error("range exceeds reservation")] + InvalidRange, } impl TCache { @@ -133,35 +151,37 @@ impl TCache { /// (e.g. `"gossip_in"`); the metrics layer uses it to produce /// `counters-tcache-{name}`. pub fn producer(name: &'static str, n: usize) -> Producer { - let tcache = Self::alloc_heap(name, n); - let space = tcache.len; - Producer { cache: Box::into_raw(tcache), seq: 0, published_seq: 0, space } + Producer::new(Self::alloc_heap(name, n)) } /// Create a multi-producer t-cache. pub fn multi_producer(name: &'static str, n: usize) -> MultiProducer { - let tcache = Self::alloc_heap(name, n); - let len = tcache.len; - MultiProducer::new(Box::into_raw(tcache), len) + MultiProducer::new(Self::producer(name, n)) } pub fn name(&self) -> &'static str { self.name } + #[inline] + pub fn capacity(&self) -> usize { + self.len as usize + } + /// Attach to a named shmem segment as a producer, creating it if needed. /// Either side (producer or consumer) may start first. `n` must be /// identical on both sides. + /// Writers must be trusted cooperating processes. Pointer-bearing payloads + /// are process-local, not portable between processes. #[cfg(unix)] pub fn shm_producer(name: &'static str, n: usize) -> Producer { - let tcache = Self::attach_shmem(name, n); - let space = tcache.len; - Producer { cache: Box::into_raw(tcache), seq: 0, published_seq: 0, space } + Producer::new(Self::attach_shmem(name, n)) } /// Attach to a named shmem segment as a random-access consumer, creating /// it if needed. Either side (producer or consumer) may start first. `n` /// must be identical on both sides. + /// The trust and payload restrictions of [`Self::shm_producer`] also apply. #[cfg(unix)] pub fn shm_consumer(name: &str, n: usize) -> Result { let tcache = Box::into_raw(Self::attach_shmem(name, n)); @@ -218,6 +238,15 @@ impl TCache { self.ra_consumer(name, auto_free, true) } + pub fn retained_random_access( + &self, + name: &'static str, + ) -> Result { + let mut consumer = self.ra_consumer(name, true, true)?; + consumer.retain(); + Ok(consumer) + } + fn ra_consumer( &self, name: &'static str, @@ -292,8 +321,8 @@ impl TCache { (seq & (self.len - 1) as u64) as usize } - fn space(&self, head_seq: u64) -> u32 { - let min_tail = self.min_tail(head_seq); + fn space(&self, head_seq: u64, min_allocation: u64) -> u32 { + let min_tail = self.min_tail(min_allocation); debug_assert!( head_seq - min_tail <= self.len as u64, "{head_seq} - {min_tail} > {}", @@ -321,7 +350,11 @@ impl TCache { if slot_seq != seq { return Err(Error::WrongSeq { expected: seq, slot: slot_seq }); } - if slot.skip != 0 { + let skip = slot.skip.load(Ordering::Acquire); + if skip != 0 { + if skip == sub_reservation::INCOMPLETE { + return Err(Error::Incomplete); + } return Ok((&[], slot.reservation_len as u64, slot.reserve_ns)); } @@ -337,6 +370,27 @@ impl TCache { )) } + // Callers only expose ranges whose writers have permanently relinquished + // ownership. + #[inline] + fn read_range(&self, seq: u64, offset: usize, length: usize) -> Result<&[u8], Error> { + let slot = self.slot_at(self.index(seq)); + if slot.magic != MAGIC { + return Err(Error::NoMagic); + } + let actual = slot.seq.load(Ordering::Acquire); + if actual != seq { + return Err(Error::WrongSeq { expected: seq, slot: actual }); + } + let end = offset.checked_add(length).ok_or(Error::InvalidRange)?; + if end > (slot.data_end - slot.data_start) as usize { + return Err(Error::InvalidRange); + } + Ok(unsafe { + slice::from_raw_parts(self.data_ptr().add(slot.data_start as usize + offset), length) + }) + } + fn slot_ts(&self, seq: u64) -> Result { let idx = self.index(seq); let slot: &Slot = self.slot_at(idx); @@ -405,7 +459,7 @@ impl TCache { slot.reserve_ns = Nanos::now(); slot.data_start = start as u32; slot.data_end = end as u32; - slot.skip = 1; + slot.skip = AtomicU8::new(1); slot.magic = MAGIC; (reserve_seq, reserve_len) @@ -448,16 +502,39 @@ impl TCache { // Update the slot ts - used in the consumer to measure queue latency. slot.reserve_ns = now; - slot.skip = 0; + slot.skip.store(0, Ordering::Relaxed); } let new_head = seq + slot.reservation_len as u64; - slot.seq = AtomicU64::new(seq); + slot.seq.store(seq, Ordering::Release); // Track the producer's actual progress for the metrics layer. // `head.seq` (visible to joining consumers) is only updated by // `publish_head` on out-of-space, but surfer wants the live // production cursor. self.record_head(new_head); + self.notify_readers(); + } + + fn complete_sub_reservation(&self, seq: u64, success: bool) -> Result<(), u8> { + let slot = self.slot_at(self.index(seq)); + slot.skip.compare_exchange( + sub_reservation::INCOMPLETE, + u8::from(!success), + Ordering::AcqRel, + Ordering::Acquire, + )?; + if success { + if let Some(timer) = &self.timer { + // Strict readers can already observe this timestamp; do not rewrite it. + timer.emit_latency_from_nanos_without_first(slot.reserve_ns, Nanos::now()); + } + } + self.notify_readers(); + Ok(()) + } + + #[inline] + fn notify_readers(&self) { #[cfg(feature = "thread_park")] flux::park::SIGNAL.signal(); } @@ -478,10 +555,7 @@ impl TCache { #[inline] fn slot_at(&self, idx: usize) -> &Slot { - unsafe { - let ptr = self.data_ptr().add(idx); - &*(slice::from_raw_parts(ptr, size_of::()).as_ptr() as *const Slot) - } + unsafe { &*self.data_ptr().add(idx).cast::() } } // --- allocators --- @@ -688,7 +762,7 @@ struct Slot { data_start: u32, data_end: u32, reservation_len: u32, - skip: u8, + skip: AtomicU8, magic: [u8; 3], } @@ -700,7 +774,7 @@ impl Default for Slot { data_start: 0, data_end: 0, reservation_len: 0, - skip: 0, + skip: AtomicU8::new(0), magic: MAGIC, } } @@ -714,7 +788,7 @@ impl Clone for Slot { data_end: self.data_end, reserve_ns: self.reserve_ns, reservation_len: self.reservation_len, - skip: self.skip, + skip: AtomicU8::new(self.skip.load(Ordering::Relaxed)), magic: MAGIC, } } @@ -1007,8 +1081,7 @@ mod tests { .collect(); // Spawn producer threads, each with its own clone of the - // MultiProducer (Arc>-backed; multiple writers - // coordinate seq+space). + // MultiProducer; allocation is shared, but payload writes are independent. let producer_threads: Vec<_> = (0..PRODUCERS) .map(|p| { let mut mp_clone = mp.clone(); diff --git a/crates/common/src/spine/tcache/cache_frame.rs b/crates/common/src/spine/tcache/cache_frame.rs new file mode 100644 index 00000000..cec65744 --- /dev/null +++ b/crates/common/src/spine/tcache/cache_frame.rs @@ -0,0 +1,300 @@ +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::{AcquiredCacheFrame, AcquiredCacheSegment}; + +/// N.B. Sized for partial data columns gossip +pub const MAX_CACHE_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 CacheFrameError { + InvalidDescriptor, + TooLarge, + CacheFull, + Expired, + Stale, +} + +#[derive(Clone, Copy, Debug)] +pub enum CacheSegment { + 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 CacheFrameRef { + descriptor: TCacheRead, + expires: Instant, +} + +impl CacheFrameRef { + pub fn write( + producer: &mut Producer, + expires: Instant, + framing: &[u8], + segments: impl ExactSizeIterator, + ) -> Result { + let count = segments.len(); + if count == 0 || count > MAX_CACHE_SEGMENTS || framing.len() > MAX_GOSSIP_FRAME_SIZE { + return Err(CacheFrameError::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(CacheFrameError::CacheFull)?; + let buffer = reservation.buffer().map_err(|_| CacheFrameError::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(CacheFrameError::InvalidDescriptor); + } + let (kind, read, metadata, offset, length) = match segment { + CacheSegment::Framing { offset, length } => { + if offset.checked_add(length).is_none_or(|end| end > framing.len()) { + return Err(CacheFrameError::InvalidDescriptor); + } + (0u64, None, 0, offset, length) + } + CacheSegment::Gossip { read, offset, length } => { + if read.tcache.cache != cache.cache { + return Err(CacheFrameError::InvalidDescriptor); + } + (1, Some(read), 0, offset, length) + } + CacheSegment::DataColumns { read, offset, length } => { + (2, Some(read), 0, offset, length) + } + CacheSegment::Shared { reservation, part, second, offset, length } => { + if part >= 128 { + return Err(CacheFrameError::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(CacheFrameError::InvalidDescriptor); + } + total = total.checked_add(length).ok_or(CacheFrameError::TooLarge)?; + if total > MAX_GOSSIP_FRAME_SIZE { + return Err(CacheFrameError::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(CacheFrameError::InvalidDescriptor); + } + buffer[12..16].copy_from_slice(&(total as u32).to_le_bytes()); + reservation.flush().map_err(|_| CacheFrameError::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(CacheFrameError::Expired); + } + if !consumer.is_strict() || consumer.cache.cache != self.descriptor.tcache.cache { + return Err(CacheFrameError::InvalidDescriptor); + } + let read = consumer.acquire_strict(self.descriptor).ok_or(CacheFrameError::Stale)?; + let buffer = read.buffer().map_err(|_| CacheFrameError::Stale)?.0; + if buffer.len() < HEADER_BYTES || buffer[..8] != MAGIC { + return Err(CacheFrameError::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_CACHE_SEGMENTS { + return Err(CacheFrameError::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(CacheFrameError::InvalidDescriptor); + } + let descriptor_len = buffer.len(); + let view = CacheFrameView { 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(CacheFrameError::InvalidDescriptor); + } + total = total.checked_add(segment.length).ok_or(CacheFrameError::TooLarge)?; + } + if total != wire_len { + return Err(CacheFrameError::InvalidDescriptor); + } + Ok(view) + } +} + +#[derive(Debug)] +pub struct CacheFrameView { + read: AcquiredRead, + count: usize, + wire_len: usize, + framing_start: usize, + descriptor_len: usize, +} + +impl CacheFrameView { + 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 { + AcquiredCacheFrame::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| CacheFrameSegment::decode(entry, self.framing_start)) + } + + fn segment(&self, index: usize) -> CacheFrameSegment { + assert!(index < self.count); + let start = HEADER_BYTES + index * SEGMENT_BYTES; + let buffer = self.read.buffer().expect("acquired descriptor").0; + CacheFrameSegment::decode(&buffer[start..start + SEGMENT_BYTES], self.framing_start) + } +} + +pub struct CacheFrameSegment { + kind: u64, + cache: u64, + seq: u64, + metadata: u64, + offset: usize, + length: usize, + framing_start: usize, +} + +impl CacheFrameSegment { + 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/cache_frame/acquired.rs b/crates/common/src/spine/tcache/cache_frame/acquired.rs new file mode 100644 index 00000000..8146d571 --- /dev/null +++ b/crates/common/src/spine/tcache/cache_frame/acquired.rs @@ -0,0 +1,141 @@ +use std::{ + mem, + ops::{Deref, Range}, + ptr::NonNull, +}; + +use super::{ + AcquiredRange, AcquiredRead, CacheFrameSegment, CacheFrameView, RandomAccessConsumer, + SubReservationRef, TCacheRead, +}; + +pub enum AcquiredCacheSegment { + Framing(Range), + Data(AcquiredRange), +} + +#[derive(Debug)] +pub struct AcquiredCacheFrame { + view: CacheFrameView, + 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 AcquiredCacheFrame {} + +impl AcquiredCacheFrame { + pub(super) fn new( + view: CacheFrameView, + 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(AcquiredCacheSegment::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(AcquiredCacheSegment::Data(range)) + } + + fn consumer(&self, segment: &CacheFrameSegment) -> NonNull { + if segment.kind == 1 { self.gossip } else { self.columns.expect("acquired column segment") } + } + + fn take_read(&mut self, segment: &CacheFrameSegment) -> 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: &CacheFrameSegment) -> AcquiredRange { + let read = self.take_read(segment); + let offset = Self::offset(segment, read.read); + AcquiredRange { read, offset, length: segment.length } + } + + fn offset(segment: &CacheFrameSegment, 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 AcquiredCacheFrame { + type Target = CacheFrameView; + + fn deref(&self) -> &Self::Target { + &self.view + } +} + +impl Drop for AcquiredCacheFrame { + 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/cache_frame/tests.rs b/crates/common/src/spine/tcache/cache_frame/tests.rs new file mode 100644 index 00000000..928e58c6 --- /dev/null +++ b/crates/common/src/spine/tcache/cache_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 = CacheFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"ab--cd", + [ + CacheSegment::Framing { offset: 0, length: 2 }, + CacheSegment::Gossip { read: source, offset: 3, length: 4 }, + CacheSegment::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(CacheFrameError::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 = CacheFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"", + [ + CacheSegment::Shared { + reservation: reference, + part: 0, + second: false, + offset: 1, + length: 2, + }, + CacheSegment::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 [ + CacheSegment::Framing { offset: usize::MAX, length: 1 }, + CacheSegment::Framing { offset: 0, length: 0 }, + CacheSegment::Framing { offset: 1, length: 4 }, + CacheSegment::Gossip { read: other_read, offset: 0, length: 4 }, + ] { + assert!( + CacheFrameRef::write(&mut producer, expires, b"data", [segment].into_iter()).is_err() + ); + } + assert!(CacheFrameRef::write(&mut producer, expires, b"", [].into_iter()).is_err()); + assert!( + CacheFrameRef::write( + &mut producer, + expires, + b"x", + std::iter::repeat_n( + CacheSegment::Framing { offset: 0, length: 1 }, + MAX_CACHE_SEGMENTS + 1, + ) + ) + .is_err() + ); + let frame = CacheFrameRef::write( + &mut producer, + expires, + b"", + [CacheSegment::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 = CacheFrameRef { descriptor: write(&mut producer, bytes), expires }; + assert!(matches!( + malformed.acquire(&mut consumer, Instant::now()), + Err(CacheFrameError::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 = [ + CacheSegment::Framing { offset: 0, length: 1 }, + CacheSegment::Gossip { read: gossip, offset: 1, length: 3 }, + CacheSegment::DataColumns { read: column, offset: 0, length: 4 }, + CacheSegment::Shared { reservation: shared, part: 0, second: false, offset: 0, length: 4 }, + CacheSegment::Shared { reservation: shared, part: 0, second: true, offset: 0, length: 2 }, + CacheSegment::Gossip { read: gossip, offset: 1, length: 3 }, + CacheSegment::Framing { offset: 0, length: 1 }, + ]; + for failed in 0..segments.len() { + let mut descriptors = segments; + descriptors[failed] = CacheSegment::DataColumns { read: column, offset: 100, length: 1 }; + let frame = CacheFrameRef::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 = CacheFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"f", + [ + CacheSegment::Framing { offset: 0, length: 1 }, + CacheSegment::Gossip { read: gossip, offset: 0, length: 6 }, + CacheSegment::DataColumns { read: column, offset: 0, length: 6 }, + CacheSegment::DataColumns { read: column, offset: 0, length: 6 }, + CacheSegment::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(AcquiredCacheSegment::Framing(_)))); + let Some(AcquiredCacheSegment::Data(gossip_range)) = frame.take_next() else { panic!() }; + let Some(AcquiredCacheSegment::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 = CacheFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"", + [ + CacheSegment::Shared { + reservation: shared, + part: 0, + second: false, + offset: 1, + length: 3, + }, + CacheSegment::Shared { + reservation: shared, + part: 2, + second: false, + offset: 0, + length: 4, + }, + CacheSegment::Shared { + reservation: shared, + part: 0, + second: true, + offset: 0, + length: 2, + }, + CacheSegment::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(AcquiredCacheSegment::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 = CacheFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"", + [ + CacheSegment::Gossip { read: source, offset: 1, length: 3 }, + CacheSegment::Gossip { read: source, offset: 4, length: 3 }, + CacheSegment::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(AcquiredCacheSegment::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 = CacheFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"", + [CacheSegment::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(AcquiredCacheSegment::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/common/src/spine/tcache/consumer.rs b/crates/common/src/spine/tcache/consumer.rs index fac47e5c..a9bf11c1 100644 --- a/crates/common/src/spine/tcache/consumer.rs +++ b/crates/common/src/spine/tcache/consumer.rs @@ -136,6 +136,37 @@ pub struct RandomAccessConsumer { } impl RandomAccessConsumer { + pub fn cache_ref(&self) -> TCacheRef { + self.cache + } + + pub fn is_strict(&self) -> bool { + self.strict + } + + pub fn is_retained(&self) -> bool { + matches!(self.active.guard, TailGuard::Fixed(_)) + } + + pub(super) fn retain(&mut self) { + assert!(self.strict); + self.active.guard = TailGuard::Fixed(self.active.tail_seq); + } + + /// The producer captures `seq` before reserving the next retained region. + /// Live reads still stop rollup before this boundary. + pub fn advance_retention(&mut self, seq: u64) { + let TailGuard::Fixed(boundary) = &mut self.active.guard else { + panic!("consumer has no fixed retention boundary"); + }; + if seq > *boundary { + *boundary = seq; + self.active.head_seq = self.active.head_seq.max(seq); + self.active.rollup(self.active.head_seq); + self.free(); + } + } + pub fn acquire(&mut self, read: TCacheRead) -> AcquiredRead { let now = Nanos::now(); self.last_read = now; @@ -200,6 +231,19 @@ impl RandomAccessConsumer { fn release(&mut self, seq: u64) { self.active.release(seq, self.name); } + + #[cfg(test)] + pub(super) fn active_count(&self) -> usize { + self.active.buckets.iter().map(|count| usize::from(*count)).sum() + } + + #[inline] + fn warn_below_tail(&self, seq: u64) { + if seq < self.active.tail_seq { + let e = TCacheError::StaleSeq { name: self.name, seq, tail: self.active.tail_seq }; + tracing::warn!("reading below current tail: {:?}", e); + } + } } impl std::fmt::Debug for RandomAccessConsumer { @@ -229,28 +273,38 @@ impl Drop for RandomAccessConsumer { /// of reads before consumer. #[derive(Debug)] pub struct AcquiredRead { - consumer: *const RandomAccessConsumer, + pub(super) consumer: *const RandomAccessConsumer, pub read: TCacheRead, pub acquired: Nanos, } impl AcquiredRead { + pub fn is_strict(&self) -> bool { + unsafe { &*self.consumer }.strict + } + pub fn buffer(&self) -> Result<(&[u8], Nanos), TCacheError> { let consumer = unsafe { &*self.consumer }; - if self.read.seq < consumer.active.tail_seq { - let e = TCacheError::StaleSeq { - name: consumer.name, - seq: self.read.seq, - tail: consumer.active.tail_seq, - }; - tracing::warn!("reading below current tail: {:?}", e); - } + consumer.warn_below_tail(self.read.seq); consumer.cache.read(self.read.seq).map(|(data, _, ts)| (data, ts)) } + #[inline] pub fn with_offset(&self, offset: usize) -> Option { let consumer = unsafe { &mut *(self.consumer as *mut RandomAccessConsumer) }; - consumer.acquire_strict(self.read).map(|read| AcquiredWithOffset { read, offset }) + let read = consumer.acquire_strict(self.read)?; + let length = read.buffer().ok()?.0.len().checked_sub(offset)?; + Some(AcquiredRange { read, offset, length }) + } + + #[inline] + pub fn with_range(&self, offset: usize, length: usize) -> Option { + let mut range = self.with_offset(offset)?; + if length > range.length { + return None; + } + range.length = length; + Some(range) } } @@ -280,6 +334,7 @@ impl Drop for AcquiredRead { } impl Clone for AcquiredRead { + #[inline] fn clone(&self) -> Self { unsafe { let consumer = &mut *(self.consumer as *mut RandomAccessConsumer); @@ -289,18 +344,60 @@ impl Clone for AcquiredRead { } } -pub struct AcquiredWithOffset { - read: AcquiredRead, - offset: usize, +pub type AcquiredWithOffset = AcquiredRange; + +#[derive(Clone, Debug)] +pub struct AcquiredRange { + pub(super) read: AcquiredRead, + pub(super) offset: usize, + pub(super) length: usize, } -impl AsRef<[u8]> for AcquiredWithOffset { - fn as_ref(&self) -> &[u8] { - match self.read.buffer() { - Ok((buffer, _)) => &buffer[self.offset..], - Err(_) => &[], +impl AcquiredRange { + #[inline] + pub fn extend_contiguous(&mut self, next: &Self) -> bool { + if self.read.consumer != next.read.consumer || + self.read.seq() != next.read.seq() || + self.offset + self.length != next.offset + { + return false; } + self.length += next.length; + true } + + pub fn slice(mut self, offset: usize, length: usize) -> Option { + if offset.checked_add(length)? > self.length { + return None; + } + self.offset += offset; + self.length = length; + Some(self) + } + + #[inline] + pub fn len(&self) -> usize { + self.length + } + + #[inline] + pub fn is_empty(&self) -> bool { + self.length == 0 + } +} + +impl AsRef<[u8]> for AcquiredRange { + #[inline] + fn as_ref(&self) -> &[u8] { + let consumer = unsafe { &*self.read.consumer }; + consumer.warn_below_tail(self.read.seq()); + consumer.cache.read_range(self.read.seq(), self.offset, self.length).unwrap_or(&[]) + } +} + +enum TailGuard { + Sliding(u64), + Fixed(u64), } pub(super) struct Buckets { @@ -314,11 +411,9 @@ pub(super) struct Buckets { // for 'strict' consumers this is set to cache length so that it never // triggers lag_threshold: u64, - // Out-of-order acquire lookback: the tail never advances within this - // many seqs of the newest acquire's bucket, so late acquires up to - // this far behind still land at or above the tail. 20% of capacity, - // rounded up to a bucket. - guard: u64, + // Sliding consumers keep 20% lookback. Retained consumers keep everything + // from a fixed boundary, independent of acquire order. + guard: TailGuard, } impl Buckets { @@ -348,7 +443,9 @@ impl Buckets { } else { lag_threshold(cache_capacity as u32) }, - guard: (cache_capacity / 5).next_multiple_of(bucket_size).max(bucket_size), + guard: TailGuard::Sliding( + (cache_capacity / 5).next_multiple_of(bucket_size).max(bucket_size), + ), } } @@ -369,7 +466,7 @@ impl Buckets { self.head_seq = self.head_seq.max(seq); if self.tail_seq == u64::MAX { - self.tail_seq = self.bucket_start_seq(seq).saturating_sub(self.guard); + self.tail_seq = self.rollup_limit(seq); } self.rollup(seq); true @@ -386,12 +483,17 @@ impl Buckets { self.rollup(self.head_seq); } + #[inline] + fn rollup_limit(&self, seq: u64) -> u64 { + match self.guard { + TailGuard::Sliding(distance) => self.bucket_start_seq(seq).saturating_sub(distance), + TailGuard::Fixed(boundary) => self.bucket_start_seq(boundary), + } + } + fn rollup(&mut self, seq: u64) { - // Rollup tail for completed buckets — keep `guard` seqs of slack - // below the newest acquire's bucket so bounded out-of-order - // acquires never land behind the tail. - let head_bucket_seq = self.bucket_start_seq(seq); - while head_bucket_seq > self.tail_seq.saturating_add(self.guard) { + let limit = self.rollup_limit(seq); + while self.tail_seq < limit { let tail_bucket = self.bucket_index(self.tail_seq); if self.head_seq - self.tail_seq > self.lag_threshold { tracing::warn!( @@ -438,6 +540,8 @@ mod tests { time::{Duration, Instant}, }; + use bytes::Bytes; + use super::*; use crate::spine::tcache::{Producer, TCache, producer::TCacheProducer}; @@ -604,6 +708,288 @@ mod tests { consumer.free(); } + #[test] + fn retention_advance_skips_unacquired_records_without_releasing_live_reads() { + let mut producer = TCache::producer("", 1 << 18); + let mut consumer = producer.cache_ref().retained_random_access("").unwrap(); + let read = write_marker(&mut producer, 32, 0xab); + let pinned = consumer.acquire_strict(read).unwrap(); + for _ in 0..20 { + write_marker(&mut producer, 8192, 0xcd); + } + let boundary = producer.next_seq(); + consumer.advance_retention(boundary); + assert_eq!(consumer.active.tail_seq, 0); + assert_eq!(pinned.buffer().unwrap().0, &[0xab; 32]); + drop(pinned); + assert_eq!(consumer.active.tail_seq, consumer.active.bucket_start_seq(boundary)); + assert!(consumer.acquire_strict(read).is_none()); + } + + #[test] + fn fixed_retention_replaces_the_sliding_guard_on_acquire_and_drop() { + let mut producer = TCache::producer("", 1 << 18); + let mut consumer = producer.cache_ref().retained_random_access("").unwrap(); + let old = write_marker(&mut producer, 32, 0xab); + for _ in 0..20 { + let newer = write_marker(&mut producer, 8192, 0xcd); + drop(consumer.acquire_strict(newer).unwrap()); + consumer.free(); + } + assert_eq!(consumer.active.tail_seq, 0); + assert_eq!(consumer.cache.head().tails[consumer.index].load(Ordering::Acquire), 0); + assert_eq!(consumer.acquire_strict(old).unwrap().buffer().unwrap().0, &[0xab; 32]); + + let boundary = producer.next_seq(); + consumer.advance_retention(boundary); + assert_eq!(consumer.active.tail_seq, consumer.active.bucket_start_seq(boundary)); + assert!(consumer.acquire_strict(old).is_none()); + } + + #[test] + fn delayed_boundary_cannot_discard_next_region_or_move_backwards() { + let mut producer = TCache::producer("", 1 << 18); + let mut consumer = producer.cache_ref().retained_random_access("").unwrap(); + for _ in 0..10 { + write_marker(&mut producer, 8192, 0xab); + } + let boundary = producer.next_seq(); + assert_ne!(boundary % consumer.active.bucket_size, 0); + let next = write_marker(&mut producer, 32, 0xcd); + for _ in 0..10 { + let newer = write_marker(&mut producer, 8192, 0xef); + drop(consumer.acquire_strict(newer).unwrap()); + } + consumer.advance_retention(boundary); + consumer.advance_retention(boundary - 8192); + assert_eq!(consumer.active.tail_seq, consumer.active.bucket_start_seq(boundary)); + assert_eq!(consumer.acquire_strict(next).unwrap().buffer().unwrap().0, &[0xcd; 32]); + } + + #[test] + fn acquired_ranges_share_cell_and_proof_bytes() { + let mut producer = TCache::producer("", 1 << 16); + let mut consumer = producer.cache_ref().strict_random_access("", true).unwrap(); + let mut reservation = producer.reserve(2098, true).unwrap(); + let buffer = reservation.buffer().unwrap(); + buffer.fill(0xff); + buffer[1..2049].fill(0x11); + buffer[2049..2097].fill(0x22); + reservation.increment_offset(2098); + + let acquired = consumer.acquire_strict(reservation.read()).unwrap(); + let cell = acquired.with_range(1, 2048).unwrap(); + let proof = acquired.with_range(2049, 48).unwrap(); + let buffer = acquired.buffer().unwrap().0; + + assert_eq!(cell.as_ref(), &[0x11; 2048]); + assert_eq!(proof.as_ref(), &[0x22; 48]); + assert_eq!(cell.as_ref().as_ptr(), buffer[1..].as_ptr()); + assert_eq!(proof.as_ref().as_ptr(), buffer[2049..].as_ptr()); + assert_eq!(cell.len(), 2048); + assert_eq!(proof.len(), 48); + assert!(!cell.is_empty()); + assert!(!proof.is_empty()); + } + + #[test] + fn acquired_range_checks_bounds_without_leaking_pins() { + let mut producer = TCache::producer("", 1 << 16); + let mut consumer = producer.cache_ref().strict_random_access("", true).unwrap(); + let read = write_marker(&mut producer, 32, 0xab); + let acquired = consumer.acquire_strict(read).unwrap(); + let bucket = consumer.active.bucket_index(read.seq()); + + for (offset, length) in [(0, 0), (0, 32), (5, 7), (31, 1), (32, 0)] { + let range = acquired.with_range(offset, length).unwrap(); + assert_eq!(range.as_ref(), &acquired.buffer().unwrap().0[offset..offset + length]); + assert_eq!(range.len(), length); + assert_eq!(range.is_empty(), length == 0); + assert_eq!(consumer.active.buckets[bucket], 2); + drop(range); + assert_eq!(consumer.active.buckets[bucket], 1); + } + + for (offset, length) in [ + (33, 0), + (32, 1), + (31, 2), + (0, 33), + (0, usize::MAX), + (usize::MAX, 0), + (usize::MAX, 1), + (1, usize::MAX), + ] { + assert!(acquired.with_range(offset, length).is_none(), "{offset}, {length}"); + assert_eq!(consumer.active.buckets[bucket], 1); + } + drop(acquired); + assert_eq!(consumer.active.buckets[bucket], 0); + } + + #[test] + fn acquired_offset_preserves_suffix_access() { + let mut producer = TCache::producer("", 1 << 16); + let mut consumer = producer.cache_ref().strict_random_access("", true).unwrap(); + let read = write_marker(&mut producer, 32, 0xab); + let acquired = consumer.acquire_strict(read).unwrap(); + let bucket = consumer.active.bucket_index(read.seq()); + + for offset in [0, 1, 31, 32] { + let range: AcquiredWithOffset = acquired.with_offset(offset).unwrap(); + assert_eq!(range.as_ref(), &acquired.buffer().unwrap().0[offset..]); + assert_eq!(range.len(), 32 - offset); + assert_eq!(range.is_empty(), offset == 32); + } + for offset in [33, usize::MAX] { + assert!(acquired.with_offset(offset).is_none()); + assert_eq!(consumer.active.buckets[bucket], 1); + } + } + + #[test] + fn acquired_range_clones_release_exactly_once() { + let mut producer = TCache::producer("", 1 << 18); + let mut consumer = producer.cache_ref().strict_random_access("", false).unwrap(); + let read = write_marker(&mut producer, 2096, 0xab); + let acquired = consumer.acquire_strict(read).unwrap(); + let cell = acquired.with_range(0, 2048).unwrap(); + let proof = acquired.with_range(2048, 48).unwrap(); + let cell_clone = cell.clone(); + let bucket = consumer.active.bucket_index(read.seq()); + assert_eq!(consumer.active.buckets[bucket], 4); + + write_marker(&mut producer, 3 * 32 * 1024, 0xcd); + let newer = write_marker(&mut producer, 32, 0xef); + drop(consumer.acquire_strict(newer).unwrap()); + + drop(acquired); + assert_eq!(consumer.active.buckets[bucket], 3); + drop(cell); + assert_eq!(consumer.active.buckets[bucket], 2); + drop(proof); + assert_eq!(consumer.active.buckets[bucket], 1); + assert_eq!(cell_clone.as_ref(), &[0xab; 2048]); + consumer.free(); + assert_eq!(consumer.cache.head().tails[consumer.index].load(Ordering::Acquire), 0); + + drop(cell_clone); + assert_eq!(consumer.active.buckets[bucket], 0); + consumer.free(); + assert!(consumer.cache.head().tails[consumer.index].load(Ordering::Acquire) > read.seq()); + } + + #[test] + fn acquired_range_bytes_slices_share_one_pin() { + let mut producer = TCache::producer("", 1 << 16); + let mut consumer = producer.cache_ref().strict_random_access("", true).unwrap(); + let read = write_marker(&mut producer, 2096, 0xab); + let acquired = consumer.acquire_strict(read).unwrap(); + let bytes = Bytes::from_owner(acquired.with_range(0, 2096).unwrap()); + let cell = bytes.slice(..2048); + let proof = bytes.slice(2048..); + let clone = cell.clone(); + let bucket = consumer.active.bucket_index(read.seq()); + assert_eq!(consumer.active.buckets[bucket], 2); + + drop(acquired); + drop(bytes); + drop(cell); + drop(proof); + assert_eq!(consumer.active.buckets[bucket], 1); + assert_eq!(clone.as_ref(), &[0xab; 2048]); + + drop(clone); + assert_eq!(consumer.active.buckets[bucket], 0); + } + + #[test] + fn strict_acquired_ranges_block_overwrite_until_last_drop() { + const CAPACITY: usize = 1 << 18; + const MESSAGE_LEN: usize = 8 * 1024; + + let mut producer = TCache::producer("", CAPACITY); + let mut consumer = producer.cache_ref().strict_random_access("", true).unwrap(); + let read = write_marker(&mut producer, 2096, 0xab); + let acquired = consumer.acquire_strict(read).unwrap(); + let cell = acquired.with_range(0, 2048).unwrap(); + let proof = acquired.with_range(2048, 48).unwrap(); + drop(acquired); + + let mut produced = 0; + while let Some(mut reservation) = producer.reserve(MESSAGE_LEN, true) { + reservation.buffer().unwrap().fill(0xcd); + reservation.increment_offset(MESSAGE_LEN); + drop(consumer.acquire_strict(reservation.read()).unwrap()); + produced += 1; + assert!(produced <= CAPACITY / MESSAGE_LEN, "overwrote a pinned record"); + } + + assert!(produced > 0); + assert_eq!(cell.as_ref(), &[0xab; 2048]); + drop(cell); + assert!(producer.reserve(MESSAGE_LEN, true).is_none()); + assert_eq!(proof.as_ref(), &[0xab; 48]); + + drop(proof); + assert!(producer.reserve(MESSAGE_LEN, true).is_some()); + assert!(consumer.acquire_strict(read).is_none()); + } + + #[test] + fn acquired_range_rejects_stale_reads_and_hides_overwritten_data() { + const CAPACITY: usize = 1 << 18; + + let mut producer = TCache::producer("", CAPACITY); + let mut consumer = producer.cache_ref().random_access("", true).unwrap(); + let read = write_marker(&mut producer, 32, 0xab); + let acquired = consumer.acquire_strict(read).unwrap(); + let range = acquired.with_range(16, 16).unwrap(); + + write_marker(&mut producer, CAPACITY - 16 * 1024, 0xcd); + let newer = write_marker(&mut producer, 32, 0xef); + let newer = consumer.acquire_strict(newer).unwrap(); + consumer.free(); + assert!(consumer.active.tail_seq > read.seq()); + assert!(consumer.cache.check_seq(read.seq())); + assert!(acquired.with_range(0, 1).is_none()); + assert!(acquired.with_range(0, 0).is_none()); + assert!(acquired.with_offset(0).is_none()); + + write_marker(&mut producer, 32 * 1024, 0x11); + assert!(!consumer.cache.check_seq(read.seq())); + assert!(range.as_ref().is_empty()); + let clone = range.clone(); + assert!(clone.as_ref().is_empty()); + drop(acquired); + drop(range); + drop(clone); + assert_eq!(consumer.active.buckets.iter().sum::(), 1); + assert_eq!(newer.buffer().unwrap().0, &[0xef; 32]); + drop(newer); + assert_eq!(consumer.active.buckets.iter().sum::(), 0); + } + + #[test] + fn acquired_range_rejects_uncommitted_reads_without_leaking_pins() { + let mut producer = TCache::producer("", 1 << 16); + let mut consumer = producer.cache_ref().strict_random_access("", true).unwrap(); + let mut reservation = producer.reserve(32, true).unwrap(); + let acquired = consumer.acquire(reservation.read()); + let bucket = consumer.active.bucket_index(reservation.seq()); + + assert!(acquired.with_range(0, 1).is_none()); + assert!(acquired.with_offset(0).is_none()); + assert_eq!(consumer.active.buckets[bucket], 1); + + reservation.buffer().unwrap().fill(0xab); + reservation.increment_offset(32); + let range = acquired.with_range(0, 1).unwrap(); + assert_eq!(range.as_ref(), &[0xab]); + drop(range); + assert_eq!(consumer.active.buckets[bucket], 1); + } + /// A consumer that keeps acquiring without ever releasing must not /// stall the producer — the Buckets lag-threshold force-evicts the /// tail so the producer can reclaim space. diff --git a/crates/common/src/spine/tcache/producer.rs b/crates/common/src/spine/tcache/producer.rs index 73e9a398..933eae4e 100644 --- a/crates/common/src/spine/tcache/producer.rs +++ b/crates/common/src/spine/tcache/producer.rs @@ -1,32 +1,34 @@ -use std::{io::Write, sync::Arc}; - -use flux::communication::Seqlock; +use std::io::Write; use super::*; +mod multi; +pub use multi::MultiProducer; + +#[cfg(test)] +mod tests; + #[allow(private_bounds)] pub trait TCacheProducer: SealedProducer { + /// May publish skipped wrap padding even when no payload space is + /// available. fn reserve(&mut self, len: usize, auto_commit: bool) -> Option; /// Publish the head sequence for joining consumers. - fn publish_head(&self) { - let tcache = unsafe { &*self.tcache() }; - let seq = self.seq(); - tcache.head().seq.store(seq, Ordering::Release); - tcache.record_head(seq); - } + fn publish_head(&self); fn cache_ref(&self) -> TCacheRef { TCacheRef { cache: self.tcache() as *const c_void } } - #[allow(clippy::mut_from_ref)] - fn reservation_buffer(&self, reservation: &mut Reservation) -> Result<&mut [u8], Error> { + fn reservation_buffer<'a>( + &self, + reservation: &'a mut Reservation, + ) -> Result<&'a mut [u8], Error> { if reservation.cache.cache != (self.tcache() as *const c_void) { return Err(Error::UnexpectedCacheRef); } - let tcache = unsafe { &*self.tcache() }; - let buffer = tcache.write(reservation.seq)?; + let buffer = reservation.writable_buffer()?; Ok(&mut buffer[reservation.offset..]) } } @@ -34,127 +36,155 @@ pub trait TCacheProducer: SealedProducer { /// Private trait. trait SealedProducer { fn tcache(&self) -> *const TCache; - fn seq(&self) -> u64; } #[derive(Debug)] pub struct Producer { pub(super) cache: *const TCache, - pub(super) seq: u64, - pub(super) published_seq: u64, - pub(super) space: u32, + state: AllocationState, } unsafe impl Send for Producer {} unsafe impl Sync for Producer {} -impl SealedProducer for Producer { - fn tcache(&self) -> *const TCache { - self.cache +impl Producer { + pub(super) fn new(cache: Box) -> Self { + let state = + AllocationState { seq: 0, min_allocation: 0, published_seq: 0, space: cache.len }; + Self { cache: Box::into_raw(cache), state } } - fn seq(&self) -> u64 { - self.seq + pub fn next_seq(&self) -> u64 { + self.state.seq } -} -impl TCacheProducer for Producer { - /// Return requested buffer space, if available. - /// If None is returned, caller should retry. - /// if `auto_commit` the reservation will be commited as soon as it is - /// filled. otherwise it must ber manually committed by calling `flush`. - fn reserve(&mut self, len: usize, auto_commit: bool) -> Option { - let tcache = unsafe { &*self.cache }; - if tcache.reserve_len(self.seq, len) > self.space as usize || - self.seq - self.published_seq > (tcache.len >> 4) as u64 - // for 32MB buffer, publish head for every 2MB reserved - { - // try reclaim space. - self.publish_head(); - self.published_seq = self.seq; - self.space = tcache.space(self.seq); + #[inline] + pub fn read_buffer(&self, read: TCacheRead) -> Result<&[u8], Error> { + if read.tcache.cache != self.cache.cast() { + return Err(Error::UnexpectedCacheRef); } - tcache.reserve(self.seq, self.space, len as u32).map(|(seq, reservation_len)| { - self.seq += reservation_len as u64; - self.space -= reservation_len as u32; - Reservation { cache: self.cache_ref(), seq, offset: 0, committed: false, auto_commit } - }) + let cache = unsafe { &*self.cache }; + cache.read(read.seq).map(|(bytes, _, _)| bytes) } -} - -#[derive(Clone, Debug)] -pub struct MultiProducer { - cache: *const TCache, - state: Arc>, -} -#[derive(Clone, Copy, Debug)] -struct MultiProducerState { - seq: u64, - space: u32, -} - -unsafe impl Send for MultiProducer {} -unsafe impl Sync for MultiProducer {} + /// A producer-side claim cannot survive an allocation, even after the view + /// has been consumed. + /// + /// ```compile_fail + /// use silver_common::{SubLayout, TCache, TCacheProducer}; + /// let mut producer = TCache::producer("", 1 << 16); + /// let reference = producer.sub_reservation( + /// SubLayout { parts: 1, first_len: 4, second_len: 2 }, b"", b"" + /// ).unwrap(); + /// let claim = producer.view_sub_reservation(reference).unwrap().claim(0).unwrap(); + /// let next = producer.reserve(32, false); + /// claim.write(b"cell", b"pf").unwrap(); + /// ``` + #[inline] + pub fn view_sub_reservation( + &self, + reference: SubReservationRef, + ) -> Result, SubReservationError> { + SubReservationView::from_producer(self, reference) + } -impl MultiProducer { - pub(super) fn new(cache: *const TCache, len: u32) -> Self { - let state = Arc::new(Seqlock::new(MultiProducerState { seq: 0, space: len })); - Self { cache, state } + /// The descriptor does not pin storage. Retention boundaries or acquired + /// owners must protect it across subsequent allocations. + pub fn sub_reservation( + &mut self, + layout: SubLayout, + prefix: &[u8], + middle: &[u8], + ) -> Result { + let length = layout + .reservation_bytes(prefix.len(), middle.len()) + .ok_or(SubReservationError::InvalidLayout)?; + let reservation = self.reserve(length, false).ok_or(SubReservationError::CacheFull)?; + Ok(SubReservationRef::new(reservation, layout, prefix, middle)) } } -impl SealedProducer for MultiProducer { +impl SealedProducer for Producer { fn tcache(&self) -> *const TCache { self.cache } +} - fn seq(&self) -> u64 { - self.state.read_copy().unwrap().0.seq +impl TCacheProducer for Producer { + fn publish_head(&self) { + self.state.publish_head(&self.cache_ref()); } -} -impl TCacheProducer for MultiProducer { /// Return requested buffer space, if available. /// If None is returned, caller should retry. /// if `auto_commit` the reservation will be commited as soon as it is /// filled. otherwise it must ber manually committed by calling `flush`. + #[inline] fn reserve(&mut self, len: usize, auto_commit: bool) -> Option { - let tcache = unsafe { &*self.cache }; - - loop { - let (mut state, version) = self.state.read_copy().ok()?; - let reservation_len = tcache.reserve_len(state.seq, len); - if reservation_len as u32 > state.space { - self.publish_head(); - state.space = tcache.space(state.seq); - self.state.write_at_version(&state, version); - } - if reservation_len as u32 > state.space { - // failed to reclaim enough space - return None; + let cache = self.cache_ref(); + self.state.reserve(cache, len, auto_commit) + } +} + +#[derive(Clone, Copy, Debug)] +struct AllocationState { + seq: u64, + // Stops reuse at the first reservation not yet observed committed or aborted. + min_allocation: u64, + published_seq: u64, + space: u32, +} + +impl AllocationState { + fn publish_head(&self, cache: &TCache) { + cache.head().seq.store(self.seq, Ordering::Release); + cache.record_head(self.seq); + } + + #[inline] + fn reserve(&mut self, cache: TCacheRef, len: usize, auto_commit: bool) -> Option { + if len > cache.capacity() - size_of::() { + return None; + } + let reservation_len = cache.reserve_len(self.seq, len); + if reservation_len > self.space as usize || + self.seq - self.published_seq > (cache.len >> 4) as u64 + // for 32MB buffer, publish head for every 2MB reserved + { + self.reclaim(&cache); + if reservation_len > cache.capacity() { + // The payload fits, but cannot share a record with wrap padding. + let padding = cache.capacity() - cache.index(self.seq); + let (seq, reserved) = + cache.reserve(self.seq, self.space, (padding - size_of::()) as u32)?; + self.seq += reserved as u64; + self.space -= reserved as u32; + cache.commit(seq, false); + self.reclaim(&cache); } + } + cache.reserve(self.seq, self.space, len as u32).map(|(seq, reservation_len)| { + self.seq += reservation_len as u64; + self.space -= reservation_len as u32; + Reservation { cache, seq, offset: 0, committed: false, auto_commit } + }) + } - let alloc_seq = state.seq; - state.seq += reservation_len as u64; - state.space -= reservation_len as u32; - - if self.state.write_at_version(&state, version) { - // We've claimed [alloc_seq, alloc_seq + reservation_len). - return tcache.reserve(alloc_seq, reservation_len as u32, len as u32).map( - |(seq, res)| { - assert_eq!(reservation_len, res); - Reservation { - cache: self.cache_ref(), - seq, - offset: 0, - committed: false, - auto_commit, - } - }, - ); + fn reclaim(&mut self, cache: &TCache) { + self.publish_head(cache); + self.published_seq = self.seq; + while self.min_allocation < self.seq { + let slot = cache.slot_at(cache.index(self.min_allocation)); + if slot.seq.load(Ordering::Acquire) != self.min_allocation { + break; } + debug_assert!( + slot.reservation_len > 0 && + slot.reservation_len as u64 <= self.seq - self.min_allocation + ); + self.min_allocation += slot.reservation_len as u64; } + self.space = cache.space(self.seq, self.min_allocation); } } @@ -163,7 +193,7 @@ pub struct Reservation { pub(super) cache: TCacheRef, pub(super) seq: u64, offset: usize, - committed: bool, + pub(super) committed: bool, auto_commit: bool, } @@ -176,30 +206,39 @@ impl Reservation { } pub fn remaining(&self) -> Result { - let buffer = self.cache.write(self.seq).map_err(std::io::Error::other)?; + let buffer = self.buffer()?; Ok(buffer.len() - self.offset) } pub fn increment_offset(&mut self, len: usize) { + let Ok(buffer_len) = self.buffer().map(|buffer| buffer.len()) else { + return; + }; self.offset += len; - if let Ok(len) = self.buffer().map(|b| b.len()) { - if self.auto_commit && self.offset == len { - tracing::trace!(seq = self.seq, len, "recv committed"); - self.cache.commit(self.seq, true); - self.committed = true; - } + if self.auto_commit && self.offset == buffer_len { + tracing::trace!(seq = self.seq, len = buffer_len, "recv committed"); + self.cache.commit(self.seq, true); + self.committed = true; + } + } + + #[inline] + fn writable_buffer(&self) -> Result<&mut [u8], Error> { + if self.committed { + return Err(Error::Committed); } + self.cache.write(self.seq) } pub fn buffer(&self) -> Result<&mut [u8], std::io::Error> { - self.cache.write(self.seq).map_err(std::io::Error::other) + self.writable_buffer().map_err(std::io::Error::other) } /// Buffer slice from the current write offset to the end of the /// reservation. Use this when successive writes must not overwrite /// earlier bytes (e.g. a framed header followed by body chunks). pub fn remaining_buffer(&self) -> Result<&mut [u8], std::io::Error> { - let buf = self.cache.write(self.seq).map_err(std::io::Error::other)?; + let buf = self.buffer()?; Ok(&mut buf[self.offset..]) } @@ -215,8 +254,9 @@ impl Reservation { impl Write for Reservation { fn write(&mut self, buf: &[u8]) -> std::io::Result { - let buffer = self.cache.write(self.seq).map_err(std::io::Error::other)?; - if buf.len() + self.offset > buffer.len() { + let buffer = self.buffer()?; + let buffer_len = buffer.len(); + if buf.len() + self.offset > buffer_len { tracing::error!( reservation_len = buffer.len(), offset = self.offset, @@ -228,7 +268,7 @@ impl Write for Reservation { buffer[self.offset..self.offset + buf.len()].copy_from_slice(buf); self.offset += buf.len(); - if self.auto_commit && self.offset == buffer.len() { + if self.auto_commit && self.offset == buffer_len { self.cache.commit(self.seq, true); self.committed = true; } @@ -237,8 +277,10 @@ impl Write for Reservation { } fn flush(&mut self) -> std::io::Result<()> { - self.cache.commit(self.seq, true); - self.committed = true; + if !self.committed { + self.cache.commit(self.seq, true); + self.committed = true; + } Ok(()) } } diff --git a/crates/common/src/spine/tcache/producer/multi.rs b/crates/common/src/spine/tcache/producer/multi.rs new file mode 100644 index 00000000..90f9e3b5 --- /dev/null +++ b/crates/common/src/spine/tcache/producer/multi.rs @@ -0,0 +1,103 @@ +use std::{ + fmt, + hint::spin_loop, + sync::{Arc, atomic::AtomicU64}, +}; + +use flux::communication::Seqlock; + +use super::*; + +#[cfg(test)] +mod tests; + +#[derive(Clone)] +pub struct MultiProducer { + cache: *const TCache, + state: Arc>, +} + +unsafe impl Send for MultiProducer {} +unsafe impl Sync for MultiProducer {} + +impl MultiProducer { + pub(in super::super) fn new(producer: Producer) -> Self { + Self { cache: producer.cache, state: Arc::new(Seqlock::new(producer.state)) } + } + + #[inline] + fn claim(&self) -> AllocationGuard<'_> { + loop { + if let Some(guard) = self.try_claim() { + return guard; + } + spin_loop(); + } + } + + #[inline] + fn try_claim(&self) -> Option> { + let version = self.state.version.load(Ordering::Relaxed); + if version & 1 != 0 { + return None; + } + self.state + .version + .compare_exchange(version, version.wrapping_add(1), Ordering::AcqRel, Ordering::Relaxed) + .ok()?; + + Some(AllocationGuard { + version: &self.state.version, + next_version: version.wrapping_add(2), + // SAFETY: Every state access holds this exclusive writer claim, + // including Debug. No optimistic copies race with these mutations. + state: unsafe { &mut *self.state.data.get() }, + }) + } +} + +impl SealedProducer for MultiProducer { + fn tcache(&self) -> *const TCache { + self.cache + } +} + +impl TCacheProducer for MultiProducer { + fn publish_head(&self) { + self.claim().state.publish_head(&self.cache_ref()); + } + + /// Contention retries internally; None means the allocation cannot fit. + #[inline] + fn reserve(&mut self, len: usize, auto_commit: bool) -> Option { + self.claim().state.reserve(self.cache_ref(), len, auto_commit) + } +} + +impl fmt::Debug for MultiProducer { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + let mut debug = f.debug_struct("MultiProducer"); + debug.field("cache", &self.cache); + if let Some(guard) = self.try_claim() { + debug.field("state", &guard.state); + } else { + debug.field("state", &""); + } + debug.finish() + } +} + +// The claim covers reclamation and header initialization, never payload writes. +// Releasing it publishes the advanced allocation head only after headers exist. +struct AllocationGuard<'a> { + version: &'a AtomicU64, + next_version: u64, + state: &'a mut AllocationState, +} + +impl Drop for AllocationGuard<'_> { + #[inline] + fn drop(&mut self) { + self.version.store(self.next_version, Ordering::Release); + } +} diff --git a/crates/common/src/spine/tcache/producer/multi/tests.rs b/crates/common/src/spine/tcache/producer/multi/tests.rs new file mode 100644 index 00000000..709250c4 --- /dev/null +++ b/crates/common/src/spine/tcache/producer/multi/tests.rs @@ -0,0 +1,184 @@ +use std::{ + panic::{AssertUnwindSafe, catch_unwind}, + sync::{Barrier, mpsc}, + thread, + time::{Duration, Instant}, +}; + +use super::*; + +#[test] +fn multi_producer_clones_share_the_oldest_unfinished_reservation() { + let mut producer = TCache::multi_producer("", 256); + let mut other = producer.clone(); + let mut first = producer.reserve(32, false).unwrap(); + first.write_all(&[0xaa; 16]).unwrap(); + other.reserve(32, true).unwrap().write_all(&[0xbb; 32]).unwrap(); + let third = producer.reserve(32, false).unwrap(); + third.buffer().unwrap().fill(0xcc); + other.reserve(32, true).unwrap().write_all(&[0xdd; 32]).unwrap(); + + assert!(producer.reserve(32, false).is_none()); + assert!(other.reserve(32, false).is_none()); + assert_eq!(producer.claim().state.min_allocation, first.seq()); + first.write_all(&[0xab; 16]).unwrap(); + first.flush().unwrap(); + + let next = other.reserve(32, false).unwrap(); + assert_eq!(next.seq(), 256); + assert_eq!(producer.claim().state.min_allocation, third.seq()); + assert_eq!(third.buffer().unwrap(), &[0xcc; 32]); + drop(third); + assert!(producer.reserve(32, false).is_some()); + assert_eq!(other.claim().state.min_allocation, next.seq()); +} + +#[test] +fn multi_producer_retries_contention_before_reserving_or_publishing() { + for allocate in [false, true] { + let producer = TCache::multi_producer("", 256); + let cache = producer.cache_ref(); + let allocator = producer.claim(); + let reservation = allocator.state.reserve(cache, 32, false).unwrap(); + assert!(producer.try_claim().is_none()); + assert_eq!(cache.head().seq.load(Ordering::Acquire), 0); + let mut other = producer.clone(); + let (started, ready) = mpsc::channel(); + let (send, receive) = mpsc::channel(); + let worker = thread::spawn(move || { + started.send(()).unwrap(); + let next = if allocate { other.reserve(32, false) } else { None }; + other.publish_head(); + send.send(next).unwrap(); + }); + + ready.recv_timeout(Duration::from_secs(2)).unwrap(); + let early = receive.recv_timeout(Duration::from_millis(50)); + let early_head = cache.head().seq.load(Ordering::Acquire); + drop(allocator); + assert!(matches!(early, Err(mpsc::RecvTimeoutError::Timeout))); + assert_eq!(early_head, 0); + let next = receive.recv_timeout(Duration::from_secs(2)).unwrap(); + worker.join().unwrap(); + + if allocate { + assert_eq!(next.as_ref().unwrap().seq(), 64); + } else { + assert!(next.is_none()); + } + let slot = cache.slot_at(cache.index(reservation.seq())); + assert_eq!(slot.magic, MAGIC); + assert_eq!(slot.seq.load(Ordering::Acquire), u64::MAX); + assert_eq!(slot.reservation_len, 64); + assert_eq!(cache.head().seq.load(Ordering::Acquire), if allocate { 128 } else { 64 }); + } +} + +#[test] +fn multi_producer_releases_claim_when_reservation_cannot_fit() { + let mut producer = TCache::multi_producer("", 256); + for len in [225, 256, usize::MAX] { + assert!(producer.reserve(len, false).is_none()); + assert_eq!(producer.try_claim().unwrap().state.seq, 0); + } + let full = producer.reserve(224, false).unwrap(); + assert!(producer.reserve(0, false).is_none()); + assert_eq!(producer.try_claim().unwrap().state.seq, 256); + drop(full); + assert_eq!(producer.reserve(224, false).unwrap().seq(), 256); +} + +#[test] +fn multi_producer_releases_claim_after_unwind() { + let mut producer = TCache::multi_producer("", 256); + let result = catch_unwind(AssertUnwindSafe(|| { + let allocator = producer.claim(); + let _reservation = allocator.state.reserve(producer.cache_ref(), 224, false).unwrap(); + panic!("unwind with an allocation claim"); + })); + assert!(result.is_err()); + assert_eq!(producer.try_claim().unwrap().state.seq, 256); + assert_eq!(producer.reserve(224, false).unwrap().seq(), 256); +} + +#[test] +fn multi_producer_debug_does_not_read_claimed_state() { + let producer = TCache::multi_producer("", 256); + let allocator = producer.claim(); + assert!(format!("{producer:?}").contains("")); + drop(allocator); + assert!(format!("{producer:?}").contains("min_allocation: 0")); +} + +#[test] +fn multi_producer_walks_and_header_initialization_survive_concurrent_wraps() { + const CAPACITY: usize = 4096; + const WORKERS: usize = 4; + let producer = TCache::multi_producer("", CAPACITY); + let start = Barrier::new(WORKERS); + + thread::scope(|scope| { + for id in 0..WORKERS { + let mut producer = producer.clone(); + let start = &start; + scope.spawn(move || { + let pattern = [id as u8 + 1; 544]; + let cache = producer.cache_ref(); + start.wait(); + let deadline = Instant::now() + Duration::from_secs(20); + for index in 0..1024 { + let len = 32 + (index * 131 + id * 17) % 513; + let mut reservation = loop { + if let Some(reservation) = producer.reserve(len, index % 3 == 2) { + break reservation; + } + assert!(Instant::now() < deadline, "allocator stopped making progress"); + thread::yield_now(); + }; + + let allocator = producer.claim(); + let mut seq = allocator.state.min_allocation; + while seq < allocator.state.seq { + let slot = cache.slot_at(cache.index(seq)); + let actual = slot.seq.load(Ordering::Acquire); + assert!(actual == seq || actual == u64::MAX); + assert_eq!(slot.magic, MAGIC); + assert!(slot.reservation_len > 0); + assert!(slot.data_start <= slot.data_end); + assert!(slot.data_end <= CAPACITY as u32); + seq += slot.reservation_len as u64; + assert!(seq <= allocator.state.seq); + } + drop(allocator); + + let middle = len / 2; + reservation.write_all(&pattern[..middle]).unwrap(); + thread::yield_now(); + assert_eq!(&reservation.buffer().unwrap()[..middle], &pattern[..middle]); + reservation.write_all(&pattern[middle..len]).unwrap(); + if index % 3 == 1 { + reservation.flush().unwrap(); + } + } + }); + } + }); + + producer.publish_head(); + assert!(producer.cache_ref().head().seq.load(Ordering::Acquire) > 100 * CAPACITY as u64); +} + +#[test] +fn multi_producer_consumer_tail_still_limits_reuse() { + let mut producer = TCache::multi_producer("", 256); + let mut consumer = producer.cache_ref().consumer("").unwrap(); + for _ in 0..4 { + producer.reserve(32, true).unwrap().write_all(&[0xab; 32]).unwrap(); + } + assert!(producer.reserve(32, false).is_none()); + assert_eq!(producer.claim().state.min_allocation, 256); + assert_eq!(consumer.read().unwrap().0, &[0xab; 32]); + consumer.free(); + assert_eq!(producer.clone().reserve(32, false).unwrap().seq(), 256); + assert!(producer.reserve(32, false).is_none()); +} diff --git a/crates/common/src/spine/tcache/producer/tests.rs b/crates/common/src/spine/tcache/producer/tests.rs new file mode 100644 index 00000000..0f5c17dc --- /dev/null +++ b/crates/common/src/spine/tcache/producer/tests.rs @@ -0,0 +1,221 @@ +use super::*; + +#[test] +fn unfinished_reservations_stop_reuse_despite_out_of_order_commits() { + let mut producer = TCache::producer("", 256); + let mut first = producer.reserve(32, false).unwrap(); + first.write_all(&[0xaa; 16]).unwrap(); + let mut second = producer.reserve(32, true).unwrap(); + second.write_all(&[0xbb; 32]).unwrap(); + let mut third = producer.reserve(32, false).unwrap(); + third.buffer().unwrap().fill(0xcc); + let mut fourth = producer.reserve(32, true).unwrap(); + fourth.write_all(&[0xdd; 32]).unwrap(); + + assert!(producer.reserve(32, false).is_none()); + assert_eq!(producer.state.min_allocation, first.seq()); + assert_eq!(&first.buffer().unwrap()[..16], &[0xaa; 16]); + + first.write_all(&[0xab; 16]).unwrap(); + first.flush().unwrap(); + let next = producer.reserve(32, false).unwrap(); + assert_eq!(next.seq(), 256); + assert_eq!(producer.state.min_allocation, third.seq()); + assert_eq!(third.buffer().unwrap(), &[0xcc; 32]); + + third.flush().unwrap(); + assert!(producer.reserve(32, false).is_some()); + assert_eq!(producer.state.min_allocation, next.seq()); +} + +#[test] +fn allocation_min_passes_manual_auto_and_aborted_commits() { + let mut producer = TCache::producer("", 256); + let mut manual = producer.reserve(32, false).unwrap(); + let mut automatic_write = producer.reserve(32, true).unwrap(); + let mut automatic_offset = producer.reserve(32, true).unwrap(); + let aborted = producer.reserve(32, false).unwrap(); + + automatic_write.write_all(&[1; 32]).unwrap(); + automatic_offset.buffer().unwrap().fill(2); + automatic_offset.increment_offset(32); + drop(aborted); + assert!(producer.reserve(32, false).is_none()); + + manual.write_all(&[3; 32]).unwrap(); + manual.flush().unwrap(); + let head = producer.next_seq(); + let full = producer.reserve(256 - size_of::(), false).unwrap(); + assert_eq!(full.seq(), head); + assert_eq!(producer.state.min_allocation, head); + assert_eq!(producer.state.space, 0); +} + +#[test] +fn allocation_min_walks_wrap_padding_on_commit_or_abort() { + for abort in [false, true] { + let mut producer = TCache::producer("", 256); + producer.reserve(160, false).unwrap().flush().unwrap(); + let mut wrapped = producer.reserve(96, false).unwrap(); + assert_eq!(wrapped.seq(), 192); + assert_eq!(producer.next_seq(), 384); + wrapped.buffer().unwrap().fill(0xab); + + producer.reserve(32, false).unwrap().flush().unwrap(); + assert!(producer.reserve(32, false).is_none()); + assert_eq!(producer.state.min_allocation, wrapped.seq()); + assert_eq!(wrapped.buffer().unwrap(), &[0xab; 96]); + + if !abort { + wrapped.flush().unwrap(); + } + drop(wrapped); + let head = producer.next_seq(); + assert_eq!(head, 448); + assert_eq!(producer.reserve(32, false).unwrap().seq(), head); + assert_eq!(producer.state.min_allocation, head); + } +} + +#[test] +fn consumer_tail_still_limits_reuse_after_all_writes_finish() { + let mut producer = TCache::producer("", 256); + let mut consumer = producer.cache_ref().consumer("").unwrap(); + for _ in 0..4 { + producer.reserve(32, true).unwrap().write_all(&[0xab; 32]).unwrap(); + } + assert!(producer.reserve(32, false).is_none()); + assert_eq!(producer.state.min_allocation, producer.next_seq()); + + assert_eq!(consumer.read().unwrap().0, &[0xab; 32]); + consumer.free(); + assert_eq!(producer.reserve(32, false).unwrap().seq(), 256); + assert!(producer.reserve(32, false).is_none()); +} + +#[test] +fn committed_reservation_cannot_modify_a_reused_slot() { + let mut producer = TCache::producer("", 256); + let mut old = producer.reserve(224, false).unwrap(); + old.write_all(&[0xaa; 32]).unwrap(); + old.flush().unwrap(); + let mut replacement = producer.reserve(224, false).unwrap(); + replacement.buffer().unwrap().fill(0xbb); + assert_eq!(replacement.seq(), 256); + + assert!(old.buffer().is_err()); + assert!(old.remaining_buffer().is_err()); + assert!(old.remaining().is_err()); + assert!(old.write(b"x").is_err()); + assert!(matches!(producer.reservation_buffer(&mut old), Err(Error::Committed))); + old.increment_offset(1); + assert_eq!(old.offset, 32); + old.flush().unwrap(); + drop(old); + + assert_eq!(replacement.buffer().unwrap(), &[0xbb; 224]); + replacement.flush().unwrap(); + assert_eq!(producer.read_buffer(replacement.read()).unwrap(), &[0xbb; 224]); +} + +#[test] +fn oversized_reservations_do_not_advance_the_allocator() { + let mut producer = TCache::producer("", 256); + for length in [225, 256, usize::MAX] { + assert!(producer.reserve(length, false).is_none()); + assert_eq!(producer.next_seq(), 0); + assert_eq!(producer.state.min_allocation, 0); + assert_eq!(producer.state.space, 256); + } + assert!(producer.reserve(224, false).is_some()); +} + +#[test] +fn every_valid_payload_size_fits_at_every_empty_cache_offset() { + fn check(mut producer: impl TCacheProducer, offset: usize, len: usize) { + if offset != 0 { + producer.reserve(offset - size_of::(), false).unwrap().flush().unwrap(); + } + let mut reservation = producer.reserve(len, false).expect("empty cache must fit payload"); + reservation.write_all(&[0xab; 224][..len]).unwrap(); + reservation.flush().unwrap(); + assert_eq!(reservation.read().len().unwrap(), len); + } + + for offset in (0..256).step_by(ALIGN) { + for len in 0..=224 { + check(TCache::producer("", 256), offset, len); + check(TCache::multi_producer("", 256), offset, len); + } + } +} + +#[test] +fn padding_waits_for_pending_writes_and_linear_consumer_release() { + fn check(mut producer: impl TCacheProducer) { + let cache = producer.cache_ref(); + let mut consumer = cache.consumer("").unwrap(); + let mut pending = producer.reserve(96, false).unwrap(); + pending.buffer().unwrap().fill(0xab); + assert!(producer.reserve(160, false).is_none()); + assert_eq!(cache.head().seq.load(Ordering::Acquire), 256); + let padding = cache.slot_at(128); + assert_eq!(padding.seq.load(Ordering::Acquire), 128); + assert_eq!(padding.reservation_len, 128); + assert_eq!(padding.skip.load(Ordering::Acquire), 1); + assert_eq!(pending.buffer().unwrap(), &[0xab; 96]); + + pending.flush().unwrap(); + assert!(producer.reserve(160, false).is_none()); + assert_eq!(consumer.read().unwrap().0, &[0xab; 96]); + consumer.free(); + assert!(consumer.read().is_err()); + consumer.free(); + let mut large = producer.reserve(160, true).unwrap(); + assert_eq!(large.seq(), 256); + large.write_all(&[0xcd; 160]).unwrap(); + assert_eq!(consumer.read().unwrap().0, &[0xcd; 160]); + } + + check(TCache::producer("", 256)); + check(TCache::multi_producer("", 256)); +} + +#[test] +fn padding_itself_must_fit_without_overwriting_consumer_data() { + fn check(mut producer: impl TCacheProducer) { + let mut consumer = producer.cache_ref().consumer("").unwrap(); + producer.reserve(96, true).unwrap().write_all(&[1; 96]).unwrap(); + assert_eq!(consumer.read().unwrap().0, &[1; 96]); + consumer.free(); + for value in [2, 3] { + producer.reserve(96, true).unwrap().write_all(&[value; 96]).unwrap(); + } + assert!(producer.reserve(160, false).is_none()); + assert_eq!(producer.cache_ref().head().seq.load(Ordering::Acquire), 384); + assert_eq!(consumer.read().unwrap().0, &[2; 96]); + consumer.free(); + assert!(producer.reserve(160, false).is_none()); + assert_eq!(producer.cache_ref().head().seq.load(Ordering::Acquire), 512); + assert_eq!(consumer.read().unwrap().0, &[3; 96]); + consumer.free(); + assert!(consumer.read().is_err()); + consumer.free(); + assert_eq!(producer.reserve(160, false).unwrap().seq(), 512); + } + + check(TCache::producer("", 256)); + check(TCache::multi_producer("", 256)); +} + +#[test] +fn sub_reservations_use_separate_padding_when_needed() { + let mut producer = TCache::producer("", 512); + producer.reserve(224, false).unwrap().flush().unwrap(); + let reference = producer + .sub_reservation(SubLayout { parts: 0, first_len: 4, second_len: 2 }, &[0xab; 300], b"") + .unwrap(); + assert_eq!(reference.read().seq(), 512); + let read = producer.view_sub_reservation(reference).unwrap().finish().unwrap(); + assert_eq!(producer.read_buffer(read).unwrap(), &[0xab; 300]); +} diff --git a/crates/common/src/spine/tcache/sub_reservation.rs b/crates/common/src/spine/tcache/sub_reservation.rs new file mode 100644 index 00000000..ed7c649f --- /dev/null +++ b/crates/common/src/spine/tcache/sub_reservation.rs @@ -0,0 +1,530 @@ +use std::{ + marker::PhantomData, + ptr, slice, + sync::atomic::{AtomicBool, AtomicU64, Ordering}, +}; + +use super::{ + AcquiredRange, AcquiredRead, Producer, RandomAccessConsumer, Reservation, Slot, TCacheRead, +}; + +pub(super) const INCOMPLETE: u8 = 2; +const STATE_BITS: u32 = 3; +const STATE_MASK: u64 = (1 << STATE_BITS) - 1; +const WRITING: u64 = 1; +const PENDING: u64 = 2; +const VALIDATING: u64 = 3; +const VERIFIED: u64 = 4; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum SubReservationError { + InvalidLayout, + WrongConsumer, + WrongProducer, + CacheFull, + Stale, + Closed, + Claimed, + Published, + Incomplete, +} + +/// Each part owns one element in each array: prefix, first array, middle, +/// second array. +#[derive(Clone, Copy, Debug)] +pub struct SubLayout { + pub parts: usize, + pub first_len: usize, + pub second_len: usize, +} + +impl SubLayout { + fn header_bytes(self) -> usize { + (size_of::
() + self.parts * size_of::()).next_multiple_of(super::ALIGN) + } + + pub fn reservation_bytes(self, prefix_len: usize, middle_len: usize) -> Option { + if self.parts > 128 || + self.first_len == 0 || + self.second_len == 0 || + self.first_len > u32::MAX as usize || + self.second_len > u32::MAX as usize + { + return None; + } + self.first_len + .checked_add(self.second_len)? + .checked_mul(self.parts)? + .checked_add(prefix_len)? + .checked_add(middle_len)? + .checked_add(self.header_bytes()) + .filter(|len| *len <= u32::MAX as usize) + } +} + +#[repr(C, align(32))] +struct Header { + ready: [AtomicU64; 2], + closed: AtomicBool, + parts: u32, + first_offset: u32, + first_len: u32, + second_offset: u32, + second_len: u32, +} + +impl Header { + #[inline] + fn ready(&self) -> u128 { + self.ready[0].load(Ordering::Acquire) as u128 | + ((self.ready[1].load(Ordering::Acquire) as u128) << 64) + } + + #[inline] + fn complete(&self) -> bool { + self.ready() == u128::MAX.checked_shr(128 - self.parts).unwrap_or(0) + } + + #[inline] + fn ranges(&self, part: usize) -> [(usize, usize); 2] { + [ + (self.first_offset as usize + part * self.first_len as usize, self.first_len as usize), + ( + self.second_offset as usize + part * self.second_len as usize, + self.second_len as usize, + ), + ] + } +} + +#[derive(Clone, Copy, Debug)] +pub struct SubReservationRef { + pub(super) read: TCacheRead, + pub(super) header_bytes: usize, +} + +impl SubReservationRef { + pub(super) fn new( + mut reservation: Reservation, + layout: SubLayout, + prefix: &[u8], + middle: &[u8], + ) -> Self { + let read = reservation.read(); + let cache = read.tcache; + let slot_ptr = unsafe { cache.data_ptr().add(cache.index(read.seq)).cast::() }; + // The sole allocator owns this unpublished record during initialization. + unsafe { + let slot = &mut *slot_ptr; + let payload = cache.data_ptr().add(slot.data_start as usize); + let first_end = prefix.len() + layout.parts * layout.first_len; + ptr::write(payload.cast::
(), Header { + ready: [AtomicU64::new(0), AtomicU64::new(0)], + closed: AtomicBool::new(false), + parts: layout.parts as u32, + first_offset: prefix.len() as u32, + first_len: layout.first_len as u32, + second_offset: (first_end + middle.len()) as u32, + second_len: layout.second_len as u32, + }); + let states = payload.add(size_of::
()).cast::(); + for part in 0..layout.parts { + ptr::write(states.add(part), AtomicU64::new(0)); + } + let data = payload.add(layout.header_bytes()); + ptr::copy_nonoverlapping(prefix.as_ptr(), data, prefix.len()); + ptr::copy_nonoverlapping(middle.as_ptr(), data.add(first_end), middle.len()); + slot.data_start += layout.header_bytes() as u32; + slot.skip.store(INCOMPLETE, Ordering::Relaxed); + slot.seq.store(read.seq, Ordering::Release); + cache.record_head(read.seq + slot.reservation_len as u64); + } + reservation.committed = true; + Self { read, header_bytes: layout.header_bytes() } + } + + pub fn read(self) -> TCacheRead { + self.read + } + + // The caller must hold a pin and have validated this part. Closure does + // not change the layout or revoke previously acquired, verified bytes. + pub(super) unsafe fn acquired_offset(self, part: usize, second: bool) -> usize { + let view = SubReservationView { reference: self, _scope: PhantomData }; + view.header().ranges(part)[usize::from(second)].0 + } + + pub fn acquire( + self, + consumer: &mut RandomAccessConsumer, + ) -> Result { + if !consumer.strict || consumer.cache.cache != self.read.tcache.cache { + return Err(SubReservationError::WrongConsumer); + } + let pin = consumer.acquire_strict(self.read).ok_or(SubReservationError::Stale)?; + let acquired = AcquiredSubReservation { pin, reference: self, _local: PhantomData }; + if acquired.view().header().closed.load(Ordering::Acquire) { + return Err(SubReservationError::Closed); + } + Ok(acquired) + } +} + +/// Owns closure. Remote tiles acquire the descriptor through their own strict +/// consumers. +pub struct SubReservation { + acquired: AcquiredSubReservation, +} + +impl SubReservation { + #[inline] + pub fn new(acquired: AcquiredSubReservation) -> Self { + Self { acquired } + } + + #[inline] + pub fn reference(&self) -> SubReservationRef { + self.acquired.reference + } + + #[inline] + pub fn acquired(&self) -> &AcquiredSubReservation { + &self.acquired + } + + pub fn finish(&self) -> Result { + self.acquired.view().finish() + } + + pub fn close(&self) { + self.acquired.view().close(); + } +} + +impl Drop for SubReservation { + fn drop(&mut self) { + self.close(); + } +} + +pub struct AcquiredSubReservation { + pin: AcquiredRead, + reference: SubReservationRef, + // Consumer accounting, including writer and validator drops, remains on the acquiring thread. + _local: PhantomData<*const ()>, +} + +impl AcquiredSubReservation { + #[inline] + fn view(&self) -> SubReservationView<'_> { + SubReservationView { reference: self.reference, _scope: PhantomData } + } + + #[inline] + pub fn ready(&self) -> u128 { + self.view().ready() + } + + #[inline] + pub fn len(&self) -> usize { + self.view().len() + } + + #[inline] + pub fn is_empty(&self) -> bool { + self.len() == 0 + } + + pub fn claim(&self, part: usize) -> Result, SubReservationError> { + self.view().claim(part) + } + + pub fn ranges(&self, part: usize) -> Option<[AcquiredRange; 2]> { + let header = self.view().header(); + if part >= header.parts as usize || header.ready() & (1u128 << part) == 0 { + return None; + } + Some(header.ranges(part).map(|(offset, length)| AcquiredRange { + read: self.pin.clone(), + offset, + length, + })) + } +} + +/// Borrows either the sole allocator or an acquired owner. Neither can release +/// this record while the view or one of its write claims remains in use. +pub struct SubReservationView<'a> { + reference: SubReservationRef, + _scope: PhantomData<&'a *const ()>, +} + +impl<'a> SubReservationView<'a> { + #[inline] + pub(super) fn from_producer( + producer: &'a Producer, + reference: SubReservationRef, + ) -> Result { + if reference.read.tcache.cache != producer.cache.cast() { + return Err(SubReservationError::WrongProducer); + } + if !reference.read.tcache.check_seq(reference.read.seq) { + return Err(SubReservationError::Stale); + } + Ok(Self { reference, _scope: PhantomData }) + } + + #[inline] + fn data(&self) -> *mut u8 { + let read = self.reference.read; + let slot = read.tcache.slot_at(read.tcache.index(read.seq)); + unsafe { read.tcache.data_ptr().add(slot.data_start as usize) } + } + + #[inline] + fn header(&self) -> &'a Header { + unsafe { &*self.data().sub(self.reference.header_bytes).cast::
() } + } + + #[inline] + fn state(&self, part: usize) -> &'a AtomicU64 { + assert!(part < self.header().parts as usize); + // Derive this pointer from the allocation, not a reference limited to the fixed + // header. + unsafe { + &*self + .data() + .sub(self.reference.header_bytes) + .add(size_of::
()) + .cast::() + .add(part) + } + } + + #[inline] + pub fn ready(&self) -> u128 { + self.header().ready() + } + + #[inline] + pub fn len(&self) -> usize { + let header = self.header(); + header.second_offset as usize + header.parts as usize * header.second_len as usize + } + + #[inline] + pub fn is_empty(&self) -> bool { + self.len() == 0 + } + + pub fn claim(self, part: usize) -> Result, SubReservationError> { + let header = self.header(); + if part >= header.parts as usize { + return Err(SubReservationError::InvalidLayout); + } + if header.closed.load(Ordering::Acquire) { + return Err(SubReservationError::Closed); + } + let state = self.state(part); + let previous = state.load(Ordering::Relaxed); + if previous & STATE_MASK != 0 { + return Err(if previous & STATE_MASK == VERIFIED { + SubReservationError::Published + } else { + SubReservationError::Claimed + }); + } + // Missing retains its attempt so delayed completions cannot match a retry. + let attempt = previous.checked_add(1 << STATE_BITS).ok_or(SubReservationError::Closed)?; + state + .compare_exchange(previous, attempt | WRITING, Ordering::Acquire, Ordering::Relaxed) + .map_err(|_| SubReservationError::Claimed)?; + let write = SubWrite { view: self, part, attempt, staged: false }; + if header.closed.load(Ordering::Acquire) { + return Err(SubReservationError::Closed); + } + Ok(write) + } + + pub fn finish(&self) -> Result { + let header = self.header(); + if header.closed.load(Ordering::Acquire) { + return Err(SubReservationError::Closed); + } + if !header.complete() { + return Err(SubReservationError::Incomplete); + } + let read = self.reference.read; + match read.tcache.complete_sub_reservation(read.seq, true) { + Ok(()) | Err(0) => {} + Err(_) => return Err(SubReservationError::Closed), + } + Ok(read) + } + + pub fn close(&self) { + self.header().closed.store(true, Ordering::Release); + let read = self.reference.read; + let _ = read.tcache.complete_sub_reservation(read.seq, false); + } + + pub fn cancel(&self, pending: PendingSubReservation) -> Result { + if pending.reservation.read.tcache.cache != self.reference.read.tcache.cache || + pending.reservation.read.seq != self.reference.read.seq || + pending.part >= self.header().parts as usize + { + return Err(SubReservationError::Stale); + } + if self.header().closed.load(Ordering::Acquire) { + return Err(SubReservationError::Closed); + } + Ok(self + .state(pending.part) + .compare_exchange( + pending.attempt | PENDING, + pending.attempt, + Ordering::AcqRel, + Ordering::Relaxed, + ) + .is_ok()) + } +} + +pub struct SubWrite<'a> { + view: SubReservationView<'a>, + part: usize, + attempt: u64, + staged: bool, +} + +impl SubWrite<'_> { + pub fn write( + mut self, + first: &[u8], + second: &[u8], + ) -> Result { + let header = self.view.header(); + if first.len() != header.first_len as usize || second.len() != header.second_len as usize { + return Err(SubReservationError::InvalidLayout); + } + let [first_range, second_range] = header.ranges(self.part); + // Only the claimed ranges are mutable. Shared references never span another + // writer's ranges. + unsafe { + ptr::copy_nonoverlapping( + first.as_ptr(), + self.view.data().add(first_range.0), + first_range.1, + ); + ptr::copy_nonoverlapping( + second.as_ptr(), + self.view.data().add(second_range.0), + second_range.1, + ); + } + if header.closed.load(Ordering::Acquire) { + return Err(SubReservationError::Closed); + } + self.view.state(self.part).store(self.attempt | PENDING, Ordering::Release); + self.staged = true; + Ok(PendingSubReservation { + reservation: self.view.reference, + part: self.part, + attempt: self.attempt, + }) + } +} + +impl Drop for SubWrite<'_> { + fn drop(&mut self) { + if !self.staged { + self.view.state(self.part).store(self.attempt, Ordering::Release); + } + } +} + +/// Copyable queue descriptor, not a pin. Retention boundaries or acquired +/// owners must protect it until handoff or expiry. +#[derive(Clone, Copy, Debug)] +pub struct PendingSubReservation { + reservation: SubReservationRef, + part: usize, + attempt: u64, +} + +impl PendingSubReservation { + pub fn reservation(self) -> SubReservationRef { + self.reservation + } + + pub fn part(self) -> usize { + self.part + } + + pub fn acquire( + self, + consumer: &mut RandomAccessConsumer, + ) -> Result { + let acquired = self.reservation.acquire(consumer)?; + acquired + .view() + .state(self.part) + .compare_exchange( + self.attempt | PENDING, + self.attempt | VALIDATING, + Ordering::Acquire, + Ordering::Relaxed, + ) + .map_err(|_| SubReservationError::Stale)?; + Ok(SubValidation { acquired, pending: self, validated: false }) + } + + /// Cancels queued work only; an active validator owns its bytes until it + /// finishes. + pub fn cancel(self, consumer: &mut RandomAccessConsumer) -> Result { + let acquired = self.reservation.acquire(consumer)?; + acquired.view().cancel(self) + } +} + +pub struct SubValidation { + acquired: AcquiredSubReservation, + pending: PendingSubReservation, + validated: bool, +} + +impl SubValidation { + pub fn buffers(&self) -> [&[u8]; 2] { + let view = self.acquired.view(); + view.header().ranges(self.pending.part).map(|(offset, length)| unsafe { + slice::from_raw_parts(view.data().add(offset), length) + }) + } + + pub fn accept(mut self) -> Result<(), SubReservationError> { + let view = self.acquired.view(); + let header = view.header(); + if header.closed.load(Ordering::Acquire) { + return Err(SubReservationError::Closed); + } + view.state(self.pending.part).store(self.pending.attempt | VERIFIED, Ordering::Release); + header.ready[self.pending.part / 64] + .fetch_or(1u64 << (self.pending.part % 64), Ordering::Release); + self.validated = true; + Ok(()) + } +} + +impl Drop for SubValidation { + fn drop(&mut self) { + if !self.validated { + self.acquired + .view() + .state(self.pending.part) + .store(self.pending.attempt, Ordering::Release); + } + } +} + +#[cfg(test)] +mod tests; diff --git a/crates/common/src/spine/tcache/sub_reservation/tests.rs b/crates/common/src/spine/tcache/sub_reservation/tests.rs new file mode 100644 index 00000000..fc09c317 --- /dev/null +++ b/crates/common/src/spine/tcache/sub_reservation/tests.rs @@ -0,0 +1,405 @@ +use std::{ + io::Write, + sync::{Arc, Barrier}, +}; + +use super::{ + super::{Producer, TCacheProducer}, + *, +}; +use crate::TCache; + +struct Harness { + owner: Option, + consumer: Box, + producer: Producer, +} + +impl Harness { + fn new(parts: usize) -> Self { + let mut producer = TCache::producer("", 1 << 17); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + let reference = producer + .sub_reservation(SubLayout { parts, first_len: 4, second_len: 2 }, b"prefix", b"middle") + .unwrap(); + let owner = SubReservation::new(reference.acquire(&mut consumer).unwrap()); + Self { owner: Some(owner), consumer, producer } + } + + fn owner(&self) -> &SubReservation { + self.owner.as_ref().unwrap() + } + + fn stage(&self, part: usize) -> PendingSubReservation { + self.owner() + .acquired() + .claim(part) + .unwrap() + .write(&[part as u8; 4], &[part as u8; 2]) + .unwrap() + } + + fn accept(&mut self, pending: PendingSubReservation) { + pending.acquire(&mut self.consumer).unwrap().accept().unwrap(); + } +} + +#[test] +fn producer_returns_an_unpinned_descriptor() { + let mut producer = TCache::producer("", 1 << 17); + let reference = producer + .sub_reservation(SubLayout { parts: 1, first_len: 4, second_len: 2 }, b"prefix", b"middle") + .unwrap(); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + drop(reference.acquire(&mut consumer).unwrap()); + drop(reference.acquire(&mut consumer).unwrap()); + + for _ in 0..64 { + let mut reservation = producer.reserve(4096, true).unwrap(); + reservation.buffer().unwrap().fill(0xcc); + reservation.increment_offset(4096); + drop(consumer.acquire_strict(reservation.read()).unwrap()); + } + assert!(matches!(reference.acquire(&mut consumer), Err(SubReservationError::Stale))); +} + +#[test] +fn producer_views_stage_cancel_and_finish_without_local_pins() { + let mut producer = TCache::producer("", 1 << 16); + let mut consumer = Box::new(producer.cache_ref().retained_random_access("").unwrap()); + let reference = producer + .sub_reservation(SubLayout { parts: 1, first_len: 4, second_len: 2 }, b"prefix", b"middle") + .unwrap(); + let old = producer + .view_sub_reservation(reference) + .unwrap() + .claim(0) + .unwrap() + .write(b"bad!", b"pf") + .unwrap(); + let view = producer.view_sub_reservation(reference).unwrap(); + assert_eq!(view.ready(), 0); + assert_eq!(view.len(), 18); + assert!(view.cancel(old).unwrap()); + let retry = view.claim(0).unwrap().write(b"cell", b"pf").unwrap(); + assert!(!producer.view_sub_reservation(reference).unwrap().cancel(old).unwrap()); + let validation = retry.acquire(&mut consumer).unwrap(); + assert!(!producer.view_sub_reservation(reference).unwrap().cancel(retry).unwrap()); + assert_eq!(validation.buffers(), [&b"cell"[..], &b"pf"[..]]); + validation.accept().unwrap(); + let view = producer.view_sub_reservation(reference).unwrap(); + assert_eq!(view.ready(), 1); + let read = view.finish().unwrap(); + assert_eq!(producer.read_buffer(read).unwrap(), b"prefixcellmiddlepf"); + view.close(); + assert!(matches!(reference.acquire(&mut consumer), Err(SubReservationError::Closed))); + + let other = TCache::producer("", 1 << 16); + assert!(matches!( + other.view_sub_reservation(reference), + Err(SubReservationError::WrongProducer) + )); + assert!(matches!(other.read_buffer(read), Err(super::super::Error::UnexpectedCacheRef))); + consumer.advance_retention(producer.next_seq()); + for _ in 0..32 { + let mut write = producer.reserve(4096, false).unwrap(); + write.buffer().unwrap().fill(0xcc); + write.flush().unwrap(); + drop(write); + consumer.advance_retention(producer.next_seq()); + } + assert!(matches!(producer.view_sub_reservation(reference), Err(SubReservationError::Stale))); + assert!(producer.read_buffer(read).is_err()); +} + +#[test] +fn ordinary_reservations_commit_or_abort() { + let mut producer = TCache::producer("", 1 << 16); + assert!(producer.reserve(usize::MAX, false).is_none()); + let read = { + let mut write = producer.reserve(4, false).unwrap(); + write.write_all(b"cell").unwrap(); + write.flush().unwrap(); + write.read() + }; + assert_eq!(producer.read_buffer(read).unwrap(), b"cell"); + let aborted = { + let write = producer.reserve(4, false).unwrap(); + write.buffer().unwrap().fill(0xcc); + write.read() + }; + assert!(producer.read_buffer(aborted).unwrap().is_empty()); + assert!(producer.reserve(4, false).is_some()); +} + +#[test] +fn writing_does_not_publish_and_completion_excludes_the_sub_header() { + let mut h = Harness::new(2); + let read = h.owner().reference().read(); + assert!(matches!(read.len(), Err(super::super::Error::Incomplete))); + let pending = h.stage(1); + assert!(h.owner().acquired().ranges(1).is_none()); + assert_eq!(h.owner().acquired().ready(), 0); + h.accept(pending); + let [cell, proof] = h.owner().acquired().ranges(1).unwrap(); + let bytes = cell.as_ref(); + let pending = h.stage(0); + h.accept(pending); + assert_eq!(bytes, &[1; 4]); + assert_eq!(proof.as_ref(), &[1; 2]); + let full = h.owner().finish().unwrap(); + let full = h.consumer.acquire_strict(full).unwrap(); + assert_eq!(full.buffer().unwrap().0, b"prefix\0\0\0\0\x01\x01\x01\x01middle\0\0\x01\x01"); +} + +#[test] +fn failed_validation_allows_retry_but_old_notifications_cannot_touch_it() { + let mut h = Harness::new(1); + let old = h.stage(0); + assert!(matches!(h.owner().acquired().claim(0), Err(SubReservationError::Claimed))); + let validating = old.acquire(&mut h.consumer).unwrap(); + assert_eq!(validating.buffers(), [&[0; 4][..], &[0; 2][..]]); + assert!(!old.cancel(&mut h.consumer).unwrap()); + assert!(matches!(old.acquire(&mut h.consumer), Err(SubReservationError::Stale))); + drop(validating); + let retry = h.stage(0); + assert!(!old.cancel(&mut h.consumer).unwrap()); + assert!(matches!(old.acquire(&mut h.consumer), Err(SubReservationError::Stale))); + h.accept(retry); + assert!(matches!(h.owner().acquired().claim(0), Err(SubReservationError::Published))); + assert_eq!(h.owner().acquired().ready(), 1); +} + +#[test] +fn aborted_writes_and_cancelled_queue_entries_release_their_claims() { + let mut h = Harness::new(1); + drop(h.owner().acquired().claim(0).unwrap()); + assert_eq!( + h.owner().acquired().claim(0).unwrap().write(&[], &[]).unwrap_err(), + SubReservationError::InvalidLayout + ); + let pending = h.stage(0); + assert!(pending.cancel(&mut h.consumer).unwrap()); + assert!(!pending.cancel(&mut h.consumer).unwrap()); + let retry = h.stage(0); + h.accept(retry); +} + +#[test] +fn closure_preserves_existing_ranges_and_prevents_new_work() { + let mut h = Harness::new(2); + let reference = h.owner().reference(); + let pending = h.stage(0); + h.accept(pending); + let ranges = h.owner().acquired().ranges(0).unwrap(); + let pending = h.stage(1); + let validation = pending.acquire(&mut h.consumer).unwrap(); + h.owner.take(); + assert_eq!(validation.buffers()[0], &[1; 4]); + assert_eq!(validation.accept(), Err(SubReservationError::Closed)); + assert!(matches!(reference.acquire(&mut h.consumer), Err(SubReservationError::Closed))); + assert_eq!(ranges[0].as_ref(), &[0; 4]); +} + +#[test] +fn closing_an_incomplete_record_unblocks_linear_consumers() { + for drop_owner in [false, true] { + let mut producer = TCache::producer("", 4096); + let mut linear = producer.cache_ref().consumer("").unwrap(); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + let reference = producer + .sub_reservation(SubLayout { parts: 1, first_len: 4, second_len: 2 }, b"", b"") + .unwrap(); + let owner = SubReservation::new(reference.acquire(&mut consumer).unwrap()); + producer.reserve(5, true).unwrap().write_all(b"after").unwrap(); + assert!(matches!(linear.read(), Err(super::super::Error::Incomplete))); + if !drop_owner { + owner.close(); + owner.close(); + } + drop(owner); + assert_eq!(linear.read().unwrap().0, b"after"); + linear.free(); + assert_eq!(reference.read().len().unwrap(), 0); + assert!(matches!( + producer.view_sub_reservation(reference).unwrap().finish(), + Err(SubReservationError::Closed) + )); + } +} + +#[test] +fn close_preserves_finished_records_and_their_timestamp() { + let mut producer = TCache::producer("", 4096); + let mut linear = producer.cache_ref().consumer("").unwrap(); + let reference = producer + .sub_reservation(SubLayout { parts: 0, first_len: 4, second_len: 2 }, b"full", b"") + .unwrap(); + let timestamp = reference.read().cache_ts().unwrap(); + let view = producer.view_sub_reservation(reference).unwrap(); + view.finish().unwrap(); + view.finish().unwrap(); + view.close(); + view.close(); + assert_eq!(reference.read().cache_ts().unwrap(), timestamp); + assert_eq!(linear.read().unwrap().0, b"full"); +} + +#[test] +fn competing_finish_and_close_do_not_reopen_skipped_records() { + let mut producer = TCache::producer("", 1 << 17); + for _ in 0..128 { + let reference = producer + .sub_reservation(SubLayout { parts: 0, first_len: 4, second_len: 2 }, b"full", b"") + .unwrap(); + let start = Barrier::new(2); + let finished = std::thread::scope(|scope| { + let producer = &producer; + let start = &start; + let worker = scope.spawn(move || { + start.wait(); + producer.view_sub_reservation(reference).unwrap().finish().is_ok() + }); + start.wait(); + producer.view_sub_reservation(reference).unwrap().close(); + worker.join().unwrap() + }); + assert_eq!(reference.read().len().unwrap(), if finished { 4 } else { 0 }); + let view = producer.view_sub_reservation(reference).unwrap(); + assert!(matches!(view.finish(), Err(SubReservationError::Closed))); + view.close(); + assert_eq!(reference.read().len().unwrap(), if finished { 4 } else { 0 }); + } +} + +#[test] +fn cross_thread_writers_claim_once_and_validator_reads_published_input() { + let mut h = Harness::new(64); + let reference = h.owner().reference(); + let cache = h.producer.cache_ref(); + let barrier = Arc::new(Barrier::new(4)); + let results = std::thread::scope(|scope| { + let mut threads = Vec::new(); + for _ in 0..3 { + let barrier = barrier.clone(); + threads.push(scope.spawn(move || { + let mut consumer = Box::new(cache.strict_random_access("", true).unwrap()); + let acquired = reference.acquire(&mut consumer).unwrap(); + barrier.wait(); + let mut pending = Vec::new(); + for part in 0..64 { + if let Ok(write) = acquired.claim(part) { + pending.push(write.write(&[part as u8; 4], &[part as u8; 2]).unwrap()); + } + } + pending + })); + } + barrier.wait(); + let mut pending = Vec::new(); + for part in 0..64 { + if let Ok(write) = h.producer.view_sub_reservation(reference).unwrap().claim(part) { + pending.push(write.write(&[part as u8; 4], &[part as u8; 2]).unwrap()); + } + } + for thread in threads { + pending.extend(thread.join().unwrap()); + } + pending + }); + assert_eq!(results.len(), 64); + assert_eq!(h.owner().acquired().ready(), 0); + for pending in results { + let validation = pending.acquire(&mut h.consumer).unwrap(); + assert_eq!(validation.buffers()[0], &[pending.part() as u8; 4]); + validation.accept().unwrap(); + } + assert_eq!(h.owner().acquired().ready(), u64::MAX as u128); + h.owner().finish().unwrap(); +} + +#[test] +fn zero_and_128_parts_finish_without_shifting_overflow() { + let h = Harness::new(0); + assert_eq!(h.owner().finish().unwrap().len().unwrap(), 12); + let mut h = Harness::new(128); + for part in 0..128 { + let pending = h.stage(part); + h.accept(pending); + } + assert_eq!(h.owner().acquired().ready(), u128::MAX); + h.owner().finish().unwrap(); + assert!(matches!(h.owner().acquired().claim(128), Err(SubReservationError::InvalidLayout))); +} + +#[test] +fn wrong_cache_and_non_strict_consumers_are_rejected() { + let h = Harness::new(1); + let reference = h.owner().reference(); + let mut non_strict = h.producer.cache_ref().random_access("", true).unwrap(); + assert!(matches!(reference.acquire(&mut non_strict), Err(SubReservationError::WrongConsumer))); + let other = TCache::producer("", 1 << 16); + let mut foreign = other.cache_ref().strict_random_access("", true).unwrap(); + assert!(matches!(reference.acquire(&mut foreign), Err(SubReservationError::WrongConsumer))); + assert_eq!(size_of::(), 32); +} + +#[test] +fn invalid_layouts_do_not_allocate() { + let mut producer = TCache::producer("", 1 << 17); + let seq = producer.next_seq(); + for layout in [ + SubLayout { parts: 129, first_len: 1, second_len: 1 }, + SubLayout { parts: 1, first_len: usize::MAX, second_len: 1 }, + SubLayout { parts: 1, first_len: 0, second_len: 1 }, + ] { + assert!(layout.reservation_bytes(1, 1).is_none()); + assert!(matches!( + producer.sub_reservation(layout, b"a", b"b"), + Err(SubReservationError::InvalidLayout) + )); + assert_eq!(producer.next_seq(), seq); + } +} + +#[test] +fn writers_validators_and_ranges_pin_expired_storage_until_their_last_drop() { + let mut h = Harness::new(3); + let reference = h.owner().reference(); + let pending = h.stage(0); + h.accept(pending); + let ranges = h.owner().acquired().ranges(0).unwrap(); + let pending = h.stage(1); + let validation = pending.acquire(&mut h.consumer).unwrap(); + let writer = reference.acquire(&mut h.consumer).unwrap(); + let writing = writer.claim(2).unwrap(); + h.owner.take(); + let mut blocked = false; + for _ in 0..64 { + let Some(mut reservation) = h.producer.reserve(4096, true) else { + blocked = true; + break; + }; + reservation.buffer().unwrap().fill(0xcc); + reservation.increment_offset(4096); + drop(h.consumer.acquire_strict(reservation.read()).unwrap()); + } + assert!(blocked); + assert_eq!(validation.buffers()[0], &[1; 4]); + assert_eq!(ranges[0].as_ref(), &[0; 4]); + drop(writing); + drop(writer); + assert!(h.producer.reserve(4096, false).is_none()); + drop(validation); + assert!(h.producer.reserve(4096, false).is_none()); + drop(ranges); + h.consumer.free(); + for _ in 0..64 { + let mut reservation = h.producer.reserve(4096, true).unwrap(); + reservation.buffer().unwrap().fill(0xcc); + reservation.increment_offset(4096); + drop(h.consumer.acquire_strict(reservation.read()).unwrap()); + } + assert!(matches!(reference.acquire(&mut h.consumer), Err(SubReservationError::Stale))); + assert!(matches!(pending.cancel(&mut h.consumer), Err(SubReservationError::Stale))); +} diff --git a/crates/common/tests/tcache_publication.rs b/crates/common/tests/tcache_publication.rs new file mode 100644 index 00000000..3eec9953 --- /dev/null +++ b/crates/common/tests/tcache_publication.rs @@ -0,0 +1,35 @@ +#![cfg(feature = "thread_park")] + +use flux::park::SIGNAL; +use silver_common::{SubLayout, SubReservationError, TCache}; + +#[test] +fn sub_reservations_notify_once_on_terminal_publication() { + let mut producer = TCache::producer("", 4096); + let initial = SIGNAL.read_counter(); + let finished = producer + .sub_reservation(SubLayout { parts: 0, first_len: 4, second_len: 2 }, b"full", b"") + .unwrap(); + let abandoned = producer + .sub_reservation(SubLayout { parts: 1, first_len: 4, second_len: 2 }, b"", b"") + .unwrap(); + assert_eq!(SIGNAL.read_counter(), initial); + let view = producer.view_sub_reservation(abandoned).unwrap(); + assert!(matches!(view.finish(), Err(SubReservationError::Incomplete))); + assert_eq!(SIGNAL.read_counter(), initial); + + let view = producer.view_sub_reservation(finished).unwrap(); + view.finish().unwrap(); + assert_eq!(SIGNAL.read_counter(), initial.wrapping_add(1)); + view.finish().unwrap(); + view.close(); + view.close(); + assert_eq!(SIGNAL.read_counter(), initial.wrapping_add(1)); + + let view = producer.view_sub_reservation(abandoned).unwrap(); + view.close(); + assert_eq!(SIGNAL.read_counter(), initial.wrapping_add(2)); + view.close(); + assert!(matches!(view.finish(), Err(SubReservationError::Closed))); + assert_eq!(SIGNAL.read_counter(), initial.wrapping_add(2)); +} diff --git a/crates/common/tests/tcache_ranges_alloc.rs b/crates/common/tests/tcache_ranges_alloc.rs new file mode 100644 index 00000000..46de35df --- /dev/null +++ b/crates/common/tests/tcache_ranges_alloc.rs @@ -0,0 +1,236 @@ +use std::{ + alloc::{GlobalAlloc, Layout, System}, + array, + cell::Cell, + hint::black_box, + io::Write, + time::{Duration, Instant}, +}; + +use silver_common::{AcquiredCacheSegment, CacheFrameRef, CacheSegment, TCache, TCacheProducer}; + +thread_local! { + static ALLOCATION_EVENTS: Cell = const { Cell::new(0) }; +} + +struct CountingAllocator; + +unsafe impl GlobalAlloc for CountingAllocator { + unsafe fn alloc(&self, layout: Layout) -> *mut u8 { + ALLOCATION_EVENTS.with(|count| count.set(count.get() + 1)); + unsafe { System.alloc(layout) } + } + + unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) { + unsafe { System.dealloc(ptr, layout) } + } + + unsafe fn alloc_zeroed(&self, layout: Layout) -> *mut u8 { + ALLOCATION_EVENTS.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 { + ALLOCATION_EVENTS.with(|count| count.set(count.get() + 1)); + unsafe { System.realloc(ptr, layout, new_size) } + } +} + +#[global_allocator] +static ALLOCATOR: CountingAllocator = CountingAllocator; + +#[test] +fn allocation_min_walk_and_out_of_order_completion_allocate_nothing() { + let mut producer = TCache::producer("", 256); + let before = ALLOCATION_EVENTS.with(Cell::get); + for _ in 0..128 { + let first = producer.reserve(32, false).unwrap(); + for _ in 0..3 { + producer.reserve(32, true).unwrap().write_all(&[0xab; 32]).unwrap(); + } + assert!(producer.reserve(32, false).is_none()); + drop(first); + } + assert_eq!(ALLOCATION_EVENTS.with(Cell::get) - before, 0); +} + +#[test] +fn multi_producer_allocation_and_reclamation_allocate_nothing() { + let producer = TCache::multi_producer("", 256); + let mut writers: [_; 4] = array::from_fn(|_| producer.clone()); + let before = ALLOCATION_EVENTS.with(Cell::get); + for _ in 0..128 { + let first = writers[0].reserve(32, false).unwrap(); + for writer in &mut writers[1..] { + writer.reserve(32, true).unwrap().write_all(&[0xab; 32]).unwrap(); + } + assert!(writers[0].reserve(32, false).is_none()); + drop(first); + } + producer.publish_head(); + assert_eq!(ALLOCATION_EVENTS.with(Cell::get) - before, 0); +} + +#[test] +fn separate_wrap_padding_allocates_nothing() { + fn check(mut producer: impl TCacheProducer) { + let before = ALLOCATION_EVENTS.with(Cell::get); + for _ in 0..128 { + producer.reserve(96, true).unwrap().write_all(&[1; 96]).unwrap(); + producer.reserve(160, true).unwrap().write_all(&[2; 160]).unwrap(); + producer.reserve(32, true).unwrap().write_all(&[3; 32]).unwrap(); + } + assert_eq!(ALLOCATION_EVENTS.with(Cell::get) - before, 0); + } + + check(TCache::producer("", 256)); + check(TCache::multi_producer("", 256)); +} + +#[test] +fn descriptor_construction_and_acquisition_allocate_nothing() { + let mut producer = TCache::producer("", 1 << 18); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + let now = Instant::now(); + let before = ALLOCATION_EVENTS.with(Cell::get); + for _ in 0..128 { + let frame = CacheFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"framing", + [CacheSegment::Framing { offset: 0, length: 3 }, CacheSegment::Framing { + offset: 3, + length: 4, + }] + .into_iter(), + ) + .unwrap(); + let view = frame.acquire(&mut consumer, now).unwrap(); + black_box(view.descriptor_range()); + black_box(view.segments().count()); + } + assert_eq!(ALLOCATION_EVENTS.with(Cell::get) - before, 0); +} + +#[test] +fn frame_acquisition_handoff_and_rollback_allocate_nothing() { + let mut producer = TCache::producer("", 1 << 18); + let mut consumer = Box::new(producer.cache_ref().strict_random_access("", true).unwrap()); + let source = { + let mut reservation = producer.reserve(32, false).unwrap(); + reservation.write_all(&[0xab; 32]).unwrap(); + reservation.flush().unwrap(); + reservation.read() + }; + let _source_pin = consumer.acquire_strict(source).unwrap(); + let now = Instant::now(); + let before = ALLOCATION_EVENTS.with(Cell::get); + for _ in 0..32 { + let reference = CacheFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"framing", + [ + CacheSegment::Gossip { read: source, offset: 0, length: 4 }, + CacheSegment::Gossip { read: source, offset: 4, length: 4 }, + CacheSegment::Framing { offset: 0, length: 7 }, + CacheSegment::Gossip { read: source, offset: 16, length: 4 }, + ] + .into_iter(), + ) + .unwrap(); + let mut frame = reference + .acquire(&mut consumer, now) + .unwrap() + .acquire_segments(&mut consumer, None) + .unwrap(); + while let Some(segment) = frame.take_next() { + match segment { + AcquiredCacheSegment::Framing(range) => { + black_box(range); + } + AcquiredCacheSegment::Data(range) => { + black_box(range.as_ref()); + } + } + } + drop(frame); + let invalid = CacheFrameRef::write( + &mut producer, + now + Duration::from_secs(1), + b"", + [CacheSegment::Gossip { read: source, offset: 0, length: 4 }, CacheSegment::Gossip { + read: source, + offset: 64, + length: 1, + }] + .into_iter(), + ) + .unwrap(); + assert!( + invalid + .acquire(&mut consumer, now) + .unwrap() + .acquire_segments(&mut consumer, None) + .is_none() + ); + } + assert_eq!(ALLOCATION_EVENTS.with(Cell::get) - before, 0); +} + +#[test] +fn range_creation_cloning_and_dropping_allocates_nothing() { + let mut producer = TCache::producer("", 1 << 16); + let mut consumer = producer.cache_ref().strict_random_access("", true).unwrap(); + let mut reservation = producer.reserve(2096, true).unwrap(); + reservation.buffer().unwrap().fill(0xab); + reservation.increment_offset(2096); + + let before = ALLOCATION_EVENTS.with(Cell::get); + assert!(before > 0); + + let acquired = consumer.acquire_strict(reservation.read()).unwrap(); + let cell = acquired.with_range(0, 2048).unwrap(); + let proof = acquired.with_range(2048, 48).unwrap(); + let clone = cell.clone(); + let suffix = acquired.with_offset(2048).unwrap(); + assert!(acquired.with_range(1, usize::MAX).is_none()); + + black_box(cell.as_ref()); + black_box(proof.as_ref()); + black_box(suffix.as_ref()); + drop(acquired); + drop(cell); + drop(proof); + drop(suffix); + black_box(clone.as_ref()); + drop(clone); + + assert_eq!(ALLOCATION_EVENTS.with(Cell::get) - before, 0); +} + +#[test] +fn slot_retention_acquisition_and_expiry_allocate_nothing_after_construction() { + let mut producer = TCache::producer("", 1 << 18); + let mut readers = array::from_fn::<_, 2, _>(|_| { + Box::new(producer.cache_ref().retained_random_access("").unwrap()) + }); + let before = ALLOCATION_EVENTS.with(Cell::get); + + for _ in 0..512 { + let boundary = producer.next_seq(); + for reader in &mut readers { + reader.advance_retention(boundary); + } + let mut reservation = producer.reserve(8192, false).unwrap(); + reservation.buffer().unwrap().fill(0xab); + reservation.flush().unwrap(); + for reader in &mut readers { + let acquired = reader.acquire_strict(reservation.read()).unwrap(); + let range = acquired.with_range(7, 31).unwrap(); + let clone = range.clone(); + black_box(clone.as_ref()); + } + } + assert_eq!(ALLOCATION_EVENTS.with(Cell::get) - before, 0); +} diff --git a/crates/gossip/src/message.rs b/crates/gossip/src/message.rs index a660704f..6f377b8c 100644 --- a/crates/gossip/src/message.rs +++ b/crates/gossip/src/message.rs @@ -149,9 +149,8 @@ fn decompress_to_reservation( let output_buffer = producer.reservation_buffer(reservation)?; let decompressed_len = snap_decoder.decompress(data, output_buffer)?; - reservation.increment_offset(decompressed_len); - let msg_id = msg_id_valid_snappy(topic, &output_buffer[..decompressed_len]); + reservation.increment_offset(decompressed_len); Ok(msg_id) }