1use std::cell::Cell;
41use std::marker::PhantomData;
42use std::path::Path;
43use std::sync::{Arc, OnceLock};
44use std::sync::atomic::{AtomicBool, AtomicU8, AtomicU32, AtomicU64, AtomicUsize, Ordering};
45
46use arc_swap::ArcSwap;
47
48use crate::frame_ring::{FrameClass, LayoutHint};
49use crate::frame_region::FrameRegion;
50use crate::peer_directory::{
51 PeerDirectory, CONSUMER_SLOT_CEILING, OWNER_NONE,
52};
53
54use crate::ordering::{
55 default_stamp_kind, ordering_region_size, stamp_now, OrderingMode,
56 OrderingRegion, StampKind, STAMPED_PAYLOAD_BYTES, STAMP_BYTES,
57};
58use crate::qos_policy::Ordering as QosOrdering;
59use crate::shared_ring::{RingError, SharedRing, PAYLOAD_BYTES};
60use crate::spsc_ring::{SpscRingCore, SPSC_PAYLOAD_BYTES};
61
62pub const DRAINER_GRACE_EPOCHS: u64 = 3;
66
67#[repr(u8)]
70#[derive(Debug, Clone, Copy, PartialEq, Eq)]
71pub enum RingShape {
72 Spsc = 0,
74 Mpsc = 1,
76 Mpmc = 2,
78 Vyukov = 3,
80}
81
82impl RingShape {
83 fn from_u8(tag: u8) -> Self {
84 match tag {
85 0 => Self::Spsc,
86 1 => Self::Mpsc,
87 2 => Self::Mpmc,
88 3 => Self::Vyukov,
89 _ => panic!("AdaptiveRing shape_tag corrupted: {tag}"),
90 }
91 }
92}
93
94const STALE_NONE: u8 = u8::MAX;
108
109pub struct AdaptiveRing {
110 shape_tag: AtomicU8,
112
113 stale_shape_tag: AtomicU8,
122
123 pin_generation: AtomicU64,
126
127 frame_region: OnceLock<Arc<FrameRegion>>,
136
137 spsc: Arc<SpscRingCore>,
139
140 mpsc: Arc<MpscBacking>,
143
144 mpmc: Arc<MpmcBacking>,
147
148 vyukov: Arc<SharedRing>,
150
151 max_producers: usize,
156 max_consumers: usize,
157
158 capacity: usize,
161
162 directory: Arc<PeerDirectory>,
166
167 synced_epoch: AtomicU64,
170
171 grow_lock: parking_lot::Mutex<()>,
173
174 shape_auto: AtomicBool,
180
181 contract: Option<crate::ring_contract::RingContract>,
189
190 ordering: Option<Arc<OrderingState>>,
197
198 backing_id: BackingId,
201
202 header_sidecar: subetha_core::HandshakeHeader,
205 ring_sidecar: Box<subetha_core::ObservationRing>,
206}
207
208unsafe impl Send for AdaptiveRing {}
209unsafe impl Sync for AdaptiveRing {}
210
211impl subetha_sidecar::AdaptiveInstance for AdaptiveRing {
212 fn header(&self) -> &subetha_core::HandshakeHeader { &self.header_sidecar }
213 fn ring(&self) -> &subetha_core::ObservationRing { &self.ring_sidecar }
214 fn make_policy(&self) -> Box<dyn subetha_sidecar::Policy> {
215 Box::new(subetha_sidecar::NoMigrationPolicy)
216 }
217}
218
219enum BackingId {
222 Anon,
223 File { prefix: std::path::PathBuf, created: bool },
224 Shm { prefix: String },
225}
226
227struct OrderingState {
230 region: OrderingRegion,
231 seen: Vec<SeenLine>,
235}
236
237#[repr(align(64))]
238struct SeenLine {
239 stamp: AtomicU64,
240 mode_tag: AtomicU32,
241 lease_gen: AtomicU64,
247}
248
249impl SeenLine {
250 fn new() -> Self {
251 Self {
252 stamp: AtomicU64::new(0),
253 mode_tag: AtomicU32::new(OrderingMode::Unordered as u32),
254 lease_gen: AtomicU64::new(u64::MAX),
255 }
256 }
257}
258
259#[inline]
261fn drainer_token(consumer_id: usize) -> u64 {
262 ((std::process::id() as u64) << 32) | (consumer_id as u64 & 0xFFFF_FFFF)
263}
264
265struct MpscBacking {
266 rings: ArcSwap<Vec<Arc<SpscRingCore>>>,
271 next_drain: AtomicUsize,
272}
273
274struct MpmcBacking {
275 rings: ArcSwap<Vec<Arc<SpscRingCore>>>,
276 consumer_cursors: Vec<PaddedCursor>,
282}
283
284#[repr(align(64))]
289struct PaddedCursor(AtomicUsize, AtomicUsize);
290
291fn consumer_cursor_table() -> Vec<PaddedCursor> {
292 (0..CONSUMER_SLOT_CEILING)
293 .map(|_| PaddedCursor(AtomicUsize::new(0), AtomicUsize::new(0)))
294 .collect()
295}
296
297#[cfg(target_os = "linux")]
306fn hugepage_region(bytes: usize) -> std::io::Result<crate::hugepages::HugepageRegion> {
307 use crate::hugepages::{HugepageRegion, HugepageSize, HUGEPAGE_2MB};
308 let pages = bytes.div_ceil(HUGEPAGE_2MB).max(1);
309 HugepageRegion::allocate(pages, HugepageSize::Mb2)
310}
311
312#[cfg(windows)]
313fn hugepage_region(bytes: usize) -> std::io::Result<crate::large_pages::LargePageRegion> {
314 use crate::large_pages::{enable_lock_memory_privilege, LargePageRegion};
315 enable_lock_memory_privilege()?;
318 LargePageRegion::allocate(bytes)
319}
320
321#[cfg(any(target_os = "freebsd", target_os = "macos"))]
322fn hugepage_region(bytes: usize) -> std::io::Result<crate::super_pages::SuperPageRegion> {
323 crate::super_pages::SuperPageRegion::allocate(bytes)
330}
331
332impl AdaptiveRing {
333 pub fn create_anon(
339 max_producers: usize,
340 max_consumers: usize,
341 capacity: usize,
342 ) -> Result<Self, RingError> {
343 assert!(max_producers >= 1, "max_producers must be >= 1");
344 assert!(max_consumers >= 1, "max_consumers must be >= 1");
345
346 let spsc = Arc::new(SpscRingCore::create_anon(capacity)?);
347
348 let mpsc_rings: Vec<Arc<SpscRingCore>> = (0..max_producers)
349 .map(|_| SpscRingCore::create_anon(capacity).map(Arc::new))
350 .collect::<Result<Vec<_>, _>>()?;
351 let mpsc = Arc::new(MpscBacking {
352 rings: ArcSwap::from_pointee(mpsc_rings),
353 next_drain: AtomicUsize::new(0),
354 });
355
356 let mpmc_rings: Vec<Arc<SpscRingCore>> = (0..max_producers)
357 .map(|_| SpscRingCore::create_anon(capacity).map(Arc::new))
358 .collect::<Result<Vec<_>, _>>()?;
359 let mpmc = Arc::new(MpmcBacking {
360 rings: ArcSwap::from_pointee(mpmc_rings),
361 consumer_cursors: consumer_cursor_table(),
362 });
363
364 let vyukov = Arc::new(SharedRing::create_anon(capacity)?);
365
366 let directory = Arc::new(PeerDirectory::create_anon()?);
367 directory.publish_rings(max_producers);
368
369 Ok(Self {
370 shape_tag: AtomicU8::new(RingShape::Spsc as u8),
371 stale_shape_tag: AtomicU8::new(STALE_NONE),
372 pin_generation: AtomicU64::new(0),
373 frame_region: OnceLock::new(),
374 spsc,
375 mpsc,
376 mpmc,
377 vyukov,
378 max_producers,
379 max_consumers,
380 capacity,
381 directory,
382 synced_epoch: AtomicU64::new(u64::MAX),
383 grow_lock: parking_lot::Mutex::new(()),
384 contract: None,
385 shape_auto: AtomicBool::new(true),
386 ordering: None,
387 backing_id: BackingId::Anon,
388 header_sidecar: subetha_core::HandshakeHeader::new(),
389 ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
390 })
391 }
392
393 #[cfg(any(target_os = "linux", windows, target_os = "freebsd", target_os = "macos"))]
410 pub fn create_hugepage(
411 max_producers: usize,
412 max_consumers: usize,
413 capacity: usize,
414 ) -> Result<Self, RingError> {
415 assert!(max_producers >= 1, "max_producers must be >= 1");
416 assert!(max_consumers >= 1, "max_consumers must be >= 1");
417
418 let spsc_bytes = crate::spsc_ring::spsc_ring_file_size(capacity);
419 let vyukov_bytes = crate::shared_ring::ring_file_size(capacity);
420
421 let spsc = Arc::new(SpscRingCore::create_in_region(
422 hugepage_region(spsc_bytes)?, capacity)?);
423
424 let mut mpsc_rings = Vec::with_capacity(max_producers);
425 for _ in 0..max_producers {
426 mpsc_rings.push(Arc::new(SpscRingCore::create_in_region(
427 hugepage_region(spsc_bytes)?, capacity)?));
428 }
429 let mpsc = Arc::new(MpscBacking {
430 rings: ArcSwap::from_pointee(mpsc_rings),
431 next_drain: AtomicUsize::new(0),
432 });
433
434 let mut mpmc_rings = Vec::with_capacity(max_producers);
435 for _ in 0..max_producers {
436 mpmc_rings.push(Arc::new(SpscRingCore::create_in_region(
437 hugepage_region(spsc_bytes)?, capacity)?));
438 }
439 let mpmc = Arc::new(MpmcBacking {
440 rings: ArcSwap::from_pointee(mpmc_rings),
441 consumer_cursors: consumer_cursor_table(),
442 });
443
444 let vyukov = Arc::new(SharedRing::create_in_region(
445 hugepage_region(vyukov_bytes)?, capacity)?);
446
447 let directory = Arc::new(PeerDirectory::create_anon()?);
448 directory.publish_rings(max_producers);
449
450 Ok(Self {
451 shape_tag: AtomicU8::new(RingShape::Spsc as u8),
452 stale_shape_tag: AtomicU8::new(STALE_NONE),
453 pin_generation: AtomicU64::new(0),
454 frame_region: OnceLock::new(),
455 spsc,
456 mpsc,
457 mpmc,
458 vyukov,
459 max_producers,
460 max_consumers,
461 capacity,
462 directory,
463 synced_epoch: AtomicU64::new(u64::MAX),
464 grow_lock: parking_lot::Mutex::new(()),
465 contract: None,
466 shape_auto: AtomicBool::new(true),
467 ordering: None,
468 backing_id: BackingId::Anon,
473 header_sidecar: subetha_core::HandshakeHeader::new(),
474 ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
475 })
476 }
477
478 pub fn create(
483 path_prefix: impl AsRef<Path>,
484 max_producers: usize,
485 max_consumers: usize,
486 capacity: usize,
487 ) -> Result<Self, RingError> {
488 assert!(max_producers >= 1 && max_consumers >= 1);
489 let base = path_prefix.as_ref();
490
491 let spsc_path = with_suffix(base, ".spsc.bin");
492 let spsc = Arc::new(SpscRingCore::create(&spsc_path, capacity)?);
493
494 let mut mpsc_rings = Vec::with_capacity(max_producers);
495 for i in 0..max_producers {
496 let p = with_suffix(base, &format!(".mpsc.{i}.bin"));
497 mpsc_rings.push(Arc::new(SpscRingCore::create(&p, capacity)?));
498 }
499 let mpsc = Arc::new(MpscBacking {
500 rings: ArcSwap::from_pointee(mpsc_rings),
501 next_drain: AtomicUsize::new(0),
502 });
503
504 let mut mpmc_rings = Vec::with_capacity(max_producers);
505 for i in 0..max_producers {
506 let p = with_suffix(base, &format!(".mpmc.{i}.bin"));
507 mpmc_rings.push(Arc::new(SpscRingCore::create(&p, capacity)?));
508 }
509 let mpmc = Arc::new(MpmcBacking {
510 rings: ArcSwap::from_pointee(mpmc_rings),
511 consumer_cursors: consumer_cursor_table(),
512 });
513
514 let vyukov_path = with_suffix(base, ".vyukov.bin");
515 let vyukov = Arc::new(SharedRing::create(&vyukov_path, capacity)?);
516
517 let directory = Arc::new(
518 PeerDirectory::create(with_suffix(base, ".peers.bin"))?,
519 );
520 directory.publish_rings(max_producers);
521
522 Ok(Self {
523 shape_tag: AtomicU8::new(RingShape::Spsc as u8),
524 stale_shape_tag: AtomicU8::new(STALE_NONE),
525 pin_generation: AtomicU64::new(0),
526 frame_region: OnceLock::new(),
527 spsc,
528 mpsc,
529 mpmc,
530 vyukov,
531 max_producers,
532 max_consumers,
533 capacity,
534 directory,
535 synced_epoch: AtomicU64::new(u64::MAX),
536 grow_lock: parking_lot::Mutex::new(()),
537 contract: None,
538 shape_auto: AtomicBool::new(true),
539 ordering: None,
540 backing_id: BackingId::File {
541 prefix: base.to_path_buf(),
542 created: true,
543 },
544 header_sidecar: subetha_core::HandshakeHeader::new(),
545 ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
546 })
547 }
548
549 pub fn open(
561 path_prefix: impl AsRef<Path>,
562 max_producers: usize,
563 max_consumers: usize,
564 expected_capacity: usize,
565 ) -> Result<Self, RingError> {
566 assert!(max_producers >= 1 && max_consumers >= 1);
567 let base = path_prefix.as_ref();
568
569 let directory = Arc::new(
574 PeerDirectory::open(with_suffix(base, ".peers.bin"))?,
575 );
576 let n_rings = directory.published().max(1);
577
578 let spsc_path = with_suffix(base, ".spsc.bin");
579 let spsc = Arc::new(SpscRingCore::open(&spsc_path, expected_capacity)?);
580
581 let mut mpsc_rings = Vec::with_capacity(n_rings);
582 for i in 0..n_rings {
583 let p = with_suffix(base, &format!(".mpsc.{i}.bin"));
584 mpsc_rings.push(Arc::new(SpscRingCore::open(&p, expected_capacity)?));
585 }
586 let mpsc = Arc::new(MpscBacking {
587 rings: ArcSwap::from_pointee(mpsc_rings),
588 next_drain: AtomicUsize::new(0),
589 });
590
591 let mut mpmc_rings = Vec::with_capacity(n_rings);
592 for i in 0..n_rings {
593 let p = with_suffix(base, &format!(".mpmc.{i}.bin"));
594 mpmc_rings.push(Arc::new(SpscRingCore::open(&p, expected_capacity)?));
595 }
596 let mpmc = Arc::new(MpmcBacking {
597 rings: ArcSwap::from_pointee(mpmc_rings),
598 consumer_cursors: consumer_cursor_table(),
599 });
600
601 let vyukov_path = with_suffix(base, ".vyukov.bin");
602 let vyukov = Arc::new(SharedRing::open(&vyukov_path, expected_capacity)?);
603
604 Ok(Self {
605 shape_tag: AtomicU8::new(RingShape::Spsc as u8),
606 stale_shape_tag: AtomicU8::new(STALE_NONE),
607 pin_generation: AtomicU64::new(0),
608 frame_region: OnceLock::new(),
609 spsc,
610 mpsc,
611 mpmc,
612 vyukov,
613 max_producers,
614 max_consumers,
615 capacity: expected_capacity,
616 directory,
617 synced_epoch: AtomicU64::new(u64::MAX),
618 grow_lock: parking_lot::Mutex::new(()),
619 contract: None,
620 shape_auto: AtomicBool::new(true),
621 ordering: None,
622 backing_id: BackingId::File {
623 prefix: base.to_path_buf(),
624 created: false,
625 },
626 header_sidecar: subetha_core::HandshakeHeader::new(),
627 ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
628 })
629 }
630
631 pub fn create_shmfs(
640 name_prefix: &str,
641 max_producers: usize,
642 max_consumers: usize,
643 capacity: usize,
644 ) -> Result<Self, RingError> {
645 assert!(max_producers >= 1, "max_producers must be >= 1");
646 assert!(max_consumers >= 1, "max_consumers must be >= 1");
647
648 let spsc_size = crate::spsc_ring::spsc_ring_file_size(capacity);
649 let vyukov_size = crate::shared_ring::ring_file_size(capacity);
650
651 let spsc_shm = crate::shm_file::ShmFile::create_or_open_named(
653 &format!("{name_prefix}_spsc"), spsc_size,
654 ).map_err(|_| RingError::PayloadTooLarge)?;
655 let spsc = Arc::new(SpscRingCore::create_from_shm(spsc_shm, capacity)?);
656
657 let mut mpsc_rings = Vec::with_capacity(max_producers);
659 for i in 0..max_producers {
660 let shm = crate::shm_file::ShmFile::create_or_open_named(
661 &format!("{name_prefix}_mpsc_{i}"), spsc_size,
662 ).map_err(|_| RingError::PayloadTooLarge)?;
663 mpsc_rings.push(Arc::new(SpscRingCore::create_from_shm(shm, capacity)?));
664 }
665 let mpsc = Arc::new(MpscBacking {
666 rings: ArcSwap::from_pointee(mpsc_rings),
667 next_drain: AtomicUsize::new(0),
668 });
669
670 let mut mpmc_rings = Vec::with_capacity(max_producers);
672 for i in 0..max_producers {
673 let shm = crate::shm_file::ShmFile::create_or_open_named(
674 &format!("{name_prefix}_mpmc_{i}"), spsc_size,
675 ).map_err(|_| RingError::PayloadTooLarge)?;
676 mpmc_rings.push(Arc::new(SpscRingCore::create_from_shm(shm, capacity)?));
677 }
678 let mpmc = Arc::new(MpmcBacking {
679 rings: ArcSwap::from_pointee(mpmc_rings),
680 consumer_cursors: consumer_cursor_table(),
681 });
682
683 let vyukov_shm = crate::shm_file::ShmFile::create_or_open_named(
685 &format!("{name_prefix}_vyukov"), vyukov_size,
686 ).map_err(|_| RingError::PayloadTooLarge)?;
687 let vyukov = Arc::new(SharedRing::create_from_shm(vyukov_shm, capacity)?);
688
689 let directory = Arc::new(PeerDirectory::create_or_open_shm(
690 &format!("{name_prefix}_peers"),
691 )?);
692 directory.publish_rings(max_producers);
693
694 Ok(Self {
695 shape_tag: AtomicU8::new(RingShape::Spsc as u8),
696 stale_shape_tag: AtomicU8::new(STALE_NONE),
697 pin_generation: AtomicU64::new(0),
698 frame_region: OnceLock::new(),
699 spsc, mpsc, mpmc, vyukov,
700 max_producers, max_consumers,
701 capacity,
702 directory,
703 synced_epoch: AtomicU64::new(u64::MAX),
704 grow_lock: parking_lot::Mutex::new(()),
705 contract: None,
706 shape_auto: AtomicBool::new(true),
707 ordering: None,
708 backing_id: BackingId::Shm { prefix: name_prefix.to_owned() },
709 header_sidecar: subetha_core::HandshakeHeader::new(),
710 ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
711 })
712 }
713
714 pub fn open_shmfs(
734 name_prefix: &str,
735 max_producers: usize,
736 max_consumers: usize,
737 expected_capacity: usize,
738 ) -> Result<Self, RingError> {
739 assert!(max_producers >= 1, "max_producers must be >= 1");
740 assert!(max_consumers >= 1, "max_consumers must be >= 1");
741
742 let spsc_size = crate::spsc_ring::spsc_ring_file_size(expected_capacity);
743 let vyukov_size = crate::shared_ring::ring_file_size(expected_capacity);
744
745 let directory = Arc::new(PeerDirectory::create_or_open_shm(
750 &format!("{name_prefix}_peers"),
751 )?);
752 let n_rings = directory.published().max(1);
753
754 let spsc_shm = crate::shm_file::ShmFile::create_or_open_named(
756 &format!("{name_prefix}_spsc"), spsc_size,
757 ).map_err(|_| RingError::PayloadTooLarge)?;
758 let spsc = Arc::new(SpscRingCore::open_from_shm(spsc_shm, expected_capacity)?);
759
760 let mut mpsc_rings = Vec::with_capacity(n_rings);
762 for i in 0..n_rings {
763 let shm = crate::shm_file::ShmFile::create_or_open_named(
764 &format!("{name_prefix}_mpsc_{i}"), spsc_size,
765 ).map_err(|_| RingError::PayloadTooLarge)?;
766 mpsc_rings.push(Arc::new(SpscRingCore::open_from_shm(shm, expected_capacity)?));
767 }
768 let mpsc = Arc::new(MpscBacking {
769 rings: ArcSwap::from_pointee(mpsc_rings),
770 next_drain: AtomicUsize::new(0),
771 });
772
773 let mut mpmc_rings = Vec::with_capacity(n_rings);
775 for i in 0..n_rings {
776 let shm = crate::shm_file::ShmFile::create_or_open_named(
777 &format!("{name_prefix}_mpmc_{i}"), spsc_size,
778 ).map_err(|_| RingError::PayloadTooLarge)?;
779 mpmc_rings.push(Arc::new(SpscRingCore::open_from_shm(shm, expected_capacity)?));
780 }
781 let mpmc = Arc::new(MpmcBacking {
782 rings: ArcSwap::from_pointee(mpmc_rings),
783 consumer_cursors: consumer_cursor_table(),
784 });
785
786 let vyukov_shm = crate::shm_file::ShmFile::create_or_open_named(
788 &format!("{name_prefix}_vyukov"), vyukov_size,
789 ).map_err(|_| RingError::PayloadTooLarge)?;
790 let vyukov = Arc::new(SharedRing::open_from_shm(vyukov_shm, expected_capacity)?);
791
792 Ok(Self {
793 shape_tag: AtomicU8::new(RingShape::Spsc as u8),
794 stale_shape_tag: AtomicU8::new(STALE_NONE),
795 pin_generation: AtomicU64::new(0),
796 frame_region: OnceLock::new(),
797 spsc, mpsc, mpmc, vyukov,
798 max_producers, max_consumers,
799 capacity: expected_capacity,
800 directory,
801 synced_epoch: AtomicU64::new(u64::MAX),
802 grow_lock: parking_lot::Mutex::new(()),
803 contract: None,
804 shape_auto: AtomicBool::new(true),
805 ordering: None,
806 backing_id: BackingId::Shm { prefix: name_prefix.to_owned() },
807 header_sidecar: subetha_core::HandshakeHeader::new(),
808 ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
809 })
810 }
811
812 pub fn with_ordering_stamps(self) -> Result<Self, RingError> {
836 self.with_ordering_stamps_impl(None)
837 }
838
839 pub fn with_ordering_stamps_kind(self, kind: StampKind) -> Result<Self, RingError> {
846 self.with_ordering_stamps_impl(Some(kind))
847 }
848
849 fn with_ordering_stamps_impl(
850 mut self,
851 kind: Option<StampKind>,
852 ) -> Result<Self, RingError> {
853 if self.ordering.is_some() {
854 return Ok(self);
855 }
856 if self.current_shape() == RingShape::Vyukov {
857 return Err(RingError::LayoutMismatch);
858 }
859 let lines = crate::peer_directory::PRODUCER_SLOT_CEILING;
866 let region = match &self.backing_id {
867 BackingId::Anon => OrderingRegion::create_anon(
868 lines,
869 kind.unwrap_or_else(default_stamp_kind),
870 )?,
871 BackingId::File { prefix, created } => {
872 let path = with_suffix(prefix, ".ordering.bin");
873 if *created {
874 OrderingRegion::create(
875 &path,
876 lines,
877 kind.unwrap_or_else(default_stamp_kind),
878 )?
879 } else {
880 let region = OrderingRegion::open(&path, lines)?;
881 if let Some(k) = kind
882 && region.stamp_kind() != k
883 {
884 return Err(RingError::LayoutMismatch);
885 }
886 region
887 }
888 }
889 BackingId::Shm { prefix } => {
890 let size = ordering_region_size(lines);
891 let shm = crate::shm_file::ShmFile::create_or_open_named(
892 &format!("{prefix}_ordering"),
893 size,
894 ).map_err(|e| RingError::IoError(e.kind()))?;
895 OrderingRegion::create_shm(
896 shm,
897 lines,
898 kind.unwrap_or_else(default_stamp_kind),
899 )?
900 }
901 };
902 let seen = (0..CONSUMER_SLOT_CEILING).map(|_| SeenLine::new()).collect();
903 self.ordering = Some(Arc::new(OrderingState { region, seen }));
904 Ok(self)
905 }
906
907 pub fn current_shape(&self) -> RingShape {
909 RingShape::from_u8(self.shape_tag.load(Ordering::Acquire))
910 }
911
912 pub fn peek_spsc_slot(&self) -> Option<PeekedSpscSlot<'_>> {
923 if self.current_shape() != RingShape::Spsc {
924 return None;
925 }
926 self.spsc.peek_slot().map(|inner| PeekedSpscSlot { inner })
927 }
928
929 pub fn is_empty(&self) -> bool {
948 if let Some(stale) = self.stale_shape()
949 && !self.backing_is_empty(stale)
950 {
951 return false;
952 }
953 self.backing_is_empty(self.current_shape())
954 }
955
956 pub fn approx_len(&self) -> usize {
961 let stale_len = match self.stale_shape() {
962 Some(stale) if stale != self.current_shape() => {
963 self.backing_approx_len(stale)
964 }
965 _ => 0,
966 };
967 stale_len + self.backing_approx_len(self.current_shape())
968 }
969
970 fn backing_approx_len(&self, shape: RingShape) -> usize {
971 match shape {
972 RingShape::Spsc => self.spsc.approx_len(),
973 RingShape::Mpsc => self.mpsc.rings.load().iter().map(|r| r.approx_len()).sum(),
974 RingShape::Mpmc => self.mpmc.rings.load().iter().map(|r| r.approx_len()).sum(),
975 RingShape::Vyukov => self.vyukov.approx_len(),
976 }
977 }
978
979 pub fn sub_ring_capacity(&self) -> usize {
984 match self.current_shape() {
985 RingShape::Spsc => self.spsc.capacity(),
986 RingShape::Mpsc => self.mpsc.rings.load().first().map(|r| r.capacity()).unwrap_or(0),
987 RingShape::Mpmc => self.mpmc.rings.load().first().map(|r| r.capacity()).unwrap_or(0),
988 RingShape::Vyukov => self.vyukov.capacity(),
989 }
990 }
991
992 pub fn total_slot_capacity(&self) -> usize {
997 match self.current_shape() {
998 RingShape::Spsc => self.spsc.capacity(),
999 RingShape::Mpsc => self.mpsc.rings.load().iter().map(|r| r.capacity()).sum(),
1000 RingShape::Mpmc => self.mpmc.rings.load().iter().map(|r| r.capacity()).sum(),
1001 RingShape::Vyukov => self.vyukov.capacity(),
1002 }
1003 }
1004
1005 pub fn pin_generation(&self) -> u64 {
1008 self.pin_generation.load(Ordering::Acquire)
1009 }
1010
1011 pub fn max_producers(&self) -> usize { self.max_producers }
1015
1016 pub fn max_consumers(&self) -> usize { self.max_consumers }
1019
1020 pub fn published_producers(&self) -> usize {
1023 self.directory.published()
1024 }
1025
1026 pub fn contract(&self) -> crate::ring_contract::RingContract {
1031 self.contract.unwrap_or_else(crate::ring_contract::RingContract::unbounded)
1032 }
1033
1034 pub fn with_contract(mut self, contract: crate::ring_contract::RingContract) -> Self {
1040 self.contract = Some(contract);
1041 self
1042 }
1043
1044 pub fn contract_filtered_shape(&self, target: RingShape) -> RingShape {
1056 if self.contract().permits_shape(target) {
1057 return target;
1058 }
1059 if self.ordering.is_none() {
1060 RingShape::Vyukov
1061 } else {
1062 target
1063 }
1064 }
1065
1066 fn reshape_for_counts(&self) -> bool {
1078 if !self.shape_auto.load(Ordering::Relaxed)
1079 || self.current_shape() == RingShape::Vyukov
1080 {
1081 return true;
1082 }
1083 let p = self.directory.active_producers();
1084 let c = self.directory.active_consumers();
1085 if let Some(target) = DefaultRingShapePolicy::target_shape(p, c)
1086 && target != self.current_shape()
1087 {
1088 return self.morph_shape(self.contract_filtered_shape(target)).is_ok();
1089 }
1090 true
1091 }
1092
1093 pub fn pin_shape(&self) {
1098 self.shape_auto.store(false, Ordering::Relaxed);
1099 }
1100
1101 pub fn resume_auto_shape(&self) {
1107 self.shape_auto.store(true, Ordering::Relaxed);
1108 let p = self.directory.active_producers();
1109 let c = self.directory.active_consumers();
1110 if let Some(target) = DefaultRingShapePolicy::target_shape(p, c) {
1111 self.morph_shape(self.contract_filtered_shape(target)).ok();
1112 }
1113 }
1114
1115 pub fn shape_is_auto(&self) -> bool {
1119 self.shape_auto.load(Ordering::Relaxed)
1120 }
1121
1122 #[inline]
1127 fn ensure_synced(&self) {
1128 let e = self.directory.epoch();
1129 if e != self.synced_epoch.load(Ordering::Relaxed) {
1130 self.sync_topology(e);
1131 }
1132 }
1133
1134 #[cold]
1135 fn sync_topology(&self, epoch: u64) {
1136 self.directory.reap_dead_peers();
1137 let arrays_ok = self.refresh_local_arrays().is_ok();
1138 let shape_ok = self.reshape_for_counts();
1139 if arrays_ok && shape_ok {
1140 self.synced_epoch.store(epoch, Ordering::Relaxed);
1144 }
1145 }
1146
1147 fn refresh_local_arrays(&self) -> Result<(), RingError> {
1152 let published = self.directory.published();
1153 if self.mpsc.rings.load().len() >= published {
1154 return Ok(());
1155 }
1156 let _guard = self.grow_lock.lock();
1157 let cur_mpsc = self.mpsc.rings.load_full();
1158 let cur_mpmc = self.mpmc.rings.load_full();
1159 if cur_mpsc.len() >= published {
1160 return Ok(());
1161 }
1162 let mut mpsc_new = (*cur_mpsc).clone();
1163 let mut mpmc_new = (*cur_mpmc).clone();
1164 for i in cur_mpsc.len()..published {
1165 let (a, b) = self.open_ring_backing(i)?;
1166 mpsc_new.push(a);
1167 mpmc_new.push(b);
1168 }
1169 self.mpsc.rings.store(Arc::new(mpsc_new));
1170 self.mpmc.rings.store(Arc::new(mpmc_new));
1171 self.pin_generation.fetch_add(1, Ordering::AcqRel);
1172 Ok(())
1173 }
1174
1175 fn open_ring_backing(
1179 &self,
1180 i: usize,
1181 ) -> Result<(Arc<SpscRingCore>, Arc<SpscRingCore>), RingError> {
1182 match &self.backing_id {
1183 BackingId::File { prefix, .. } => {
1184 let a = SpscRingCore::open(
1185 with_suffix(prefix, &format!(".mpsc.{i}.bin")), self.capacity)?;
1186 let b = SpscRingCore::open(
1187 with_suffix(prefix, &format!(".mpmc.{i}.bin")), self.capacity)?;
1188 Ok((Arc::new(a), Arc::new(b)))
1189 }
1190 BackingId::Shm { prefix } => {
1191 let size = crate::spsc_ring::spsc_ring_file_size(self.capacity);
1192 let shm_a = crate::shm_file::ShmFile::create_or_open_named(
1193 &format!("{prefix}_mpsc_{i}"), size,
1194 ).map_err(|e| RingError::IoError(e.kind()))?;
1195 let shm_b = crate::shm_file::ShmFile::create_or_open_named(
1196 &format!("{prefix}_mpmc_{i}"), size,
1197 ).map_err(|e| RingError::IoError(e.kind()))?;
1198 let a = SpscRingCore::create_from_shm(shm_a, self.capacity)?;
1199 let b = SpscRingCore::create_from_shm(shm_b, self.capacity)?;
1200 Ok((Arc::new(a), Arc::new(b)))
1201 }
1202 BackingId::Anon => Err(RingError::LayoutMismatch),
1205 }
1206 }
1207
1208 fn create_ring_backing(
1212 &self,
1213 i: usize,
1214 ) -> Result<(Arc<SpscRingCore>, Arc<SpscRingCore>), RingError> {
1215 match &self.backing_id {
1216 BackingId::Anon => {
1217 let a = SpscRingCore::create_anon(self.capacity)?;
1218 let b = SpscRingCore::create_anon(self.capacity)?;
1219 Ok((Arc::new(a), Arc::new(b)))
1220 }
1221 BackingId::File { prefix, .. } => {
1222 let a = SpscRingCore::create(
1223 with_suffix(prefix, &format!(".mpsc.{i}.bin")), self.capacity)?;
1224 let b = SpscRingCore::create(
1225 with_suffix(prefix, &format!(".mpmc.{i}.bin")), self.capacity)?;
1226 Ok((Arc::new(a), Arc::new(b)))
1227 }
1228 BackingId::Shm { prefix } => {
1229 let size = crate::spsc_ring::spsc_ring_file_size(self.capacity);
1230 let shm_a = crate::shm_file::ShmFile::create_or_open_named(
1231 &format!("{prefix}_mpsc_{i}"), size,
1232 ).map_err(|e| RingError::IoError(e.kind()))?;
1233 let shm_b = crate::shm_file::ShmFile::create_or_open_named(
1234 &format!("{prefix}_mpmc_{i}"), size,
1235 ).map_err(|e| RingError::IoError(e.kind()))?;
1236 let a = SpscRingCore::create_from_shm(shm_a, self.capacity)?;
1237 let b = SpscRingCore::create_from_shm(shm_b, self.capacity)?;
1238 Ok((Arc::new(a), Arc::new(b)))
1239 }
1240 }
1241 }
1242
1243 fn grow_rings_to(&self, want: usize) -> Result<(), RingError> {
1248 let _guard = self.grow_lock.lock();
1249 let published = self.directory.published();
1250 let cur_mpsc = self.mpsc.rings.load_full();
1251 let cur_mpmc = self.mpmc.rings.load_full();
1252 let mut mpsc_new = (*cur_mpsc).clone();
1253 let mut mpmc_new = (*cur_mpmc).clone();
1254 for i in cur_mpsc.len()..published {
1257 let (a, b) = self.open_ring_backing(i)?;
1258 mpsc_new.push(a);
1259 mpmc_new.push(b);
1260 }
1261 for i in published..want {
1262 let (a, b) = self.create_ring_backing(i)?;
1263 mpsc_new.push(a);
1264 mpmc_new.push(b);
1265 }
1266 if mpsc_new.len() > cur_mpsc.len() {
1267 self.mpsc.rings.store(Arc::new(mpsc_new));
1268 self.mpmc.rings.store(Arc::new(mpmc_new));
1269 self.pin_generation.fetch_add(1, Ordering::AcqRel);
1270 }
1271 if want > published {
1272 self.directory.publish_rings(want);
1273 }
1274 Ok(())
1275 }
1276
1277 pub fn register_producer(&self) -> Result<usize, AdaptiveError> {
1288 let slot = self.directory.claim_producer_slot()
1289 .ok_or(AdaptiveError::TooManyProducers)?;
1290 if let Some(g) = self.contract
1291 && !g.permits_producer(self.directory.active_producers() - 1)
1292 {
1293 self.directory.release_producer_slot(slot);
1294 return Err(AdaptiveError::TooManyProducers);
1295 }
1296 if slot >= self.directory.published()
1297 && self.grow_rings_to(slot + 1).is_err()
1298 {
1299 self.directory.release_producer_slot(slot);
1300 return Err(AdaptiveError::GrowthFailed);
1301 }
1302 self.ensure_synced();
1303 self.reshape_for_counts();
1304 Ok(slot)
1305 }
1306
1307 pub fn unregister_producer(&self, producer_id: usize) {
1312 self.directory.release_producer_slot(producer_id);
1313 self.reshape_for_counts();
1314 }
1315
1316 pub fn register_consumer(&self) -> Result<usize, AdaptiveError> {
1323 let slot = self.directory.claim_consumer_slot()
1324 .ok_or(AdaptiveError::TooManyConsumers)?;
1325 if let Some(g) = self.contract
1326 && !g.permits_consumer(self.directory.active_consumers() - 1)
1327 {
1328 self.directory.release_consumer_slot(slot);
1329 return Err(AdaptiveError::TooManyConsumers);
1330 }
1331 self.rebalance_ownership();
1332 self.ensure_synced();
1333 self.reshape_for_counts();
1334 Ok(slot)
1335 }
1336
1337 pub fn unregister_consumer(&self, consumer_id: usize) {
1342 let me = consumer_id as u16;
1343 let remaining: Vec<u16> = self.directory.claimed_consumer_slots()
1344 .into_iter()
1345 .filter(|s| *s != me)
1346 .collect();
1347 let n = self.directory.published();
1348 for r in 0..n {
1349 let (owner, _) = self.directory.ring_owner(r);
1350 if owner == me {
1351 match remaining.get(r % remaining.len().max(1)) {
1352 Some(to) => self.directory.transfer_ring(r, me, *to),
1353 None => self.directory.transfer_ring(r, me, OWNER_NONE),
1354 }
1355 }
1356 }
1357 self.directory.release_consumer_slot(consumer_id);
1358 self.reshape_for_counts();
1359 }
1360
1361 fn rebalance_ownership(&self) {
1367 let slots = self.directory.claimed_consumer_slots();
1368 if slots.is_empty() {
1369 return;
1370 }
1371 let n = self.directory.published();
1372 for r in 0..n {
1373 let desired = slots[r % slots.len()];
1374 let (owner, pending) = self.directory.ring_owner(r);
1375 if owner == desired {
1376 continue;
1377 }
1378 if owner == OWNER_NONE {
1379 self.directory.try_claim_ring(r, desired);
1380 } else if pending != desired {
1381 self.directory.request_handoff(r, desired);
1382 }
1383 }
1384 }
1385
1386 pub fn active_producers(&self) -> usize {
1388 self.directory.active_producers()
1389 }
1390
1391 pub fn active_consumers(&self) -> usize {
1393 self.directory.active_consumers()
1394 }
1395
1396 pub fn is_stamped(&self) -> bool {
1398 self.ordering.is_some()
1399 }
1400
1401 pub fn stamp_kind(&self) -> Option<StampKind> {
1403 self.ordering.as_ref().map(|o| o.region.stamp_kind())
1404 }
1405
1406 pub fn ordering_mode(&self) -> Option<OrderingMode> {
1411 self.ordering.as_ref().map(|o| o.region.mode())
1412 }
1413
1414 pub fn set_ordering_mode(&self, mode: OrderingMode) -> Result<(), RingError> {
1420 let ord = self.ordering.as_ref().ok_or(RingError::NotStamped)?;
1421 ord.region.set_mode(mode);
1422 Ok(())
1423 }
1424
1425 pub fn inversions(&self) -> u64 {
1428 self.ordering.as_ref().map(|o| o.region.inversions()).unwrap_or(0)
1429 }
1430
1431 pub fn refresh_watermark(&self, producer_id: usize) -> Result<(), RingError> {
1434 let ord = self.ordering.as_ref().ok_or(RingError::NotStamped)?;
1435 if producer_id >= ord.region.max_producers() {
1436 return Err(RingError::PayloadTooLarge);
1437 }
1438 ord.region.refresh_watermark(producer_id);
1439 Ok(())
1440 }
1441
1442 pub fn retire_producer(&self, producer_id: usize) -> Result<(), RingError> {
1447 let ord = self.ordering.as_ref().ok_or(RingError::NotStamped)?;
1448 if producer_id >= ord.region.max_producers() {
1449 return Err(RingError::PayloadTooLarge);
1450 }
1451 ord.region.retire_producer(producer_id);
1452 Ok(())
1453 }
1454
1455 pub fn release_drainer(&self, consumer_id: usize) -> Result<bool, RingError> {
1459 let ord = self.ordering.as_ref().ok_or(RingError::NotStamped)?;
1460 Ok(ord.region.release_drainer(drainer_token(consumer_id)))
1461 }
1462
1463 pub fn tick_drainer_epoch(&self) -> Result<u64, RingError> {
1468 let ord = self.ordering.as_ref().ok_or(RingError::NotStamped)?;
1469 Ok(ord.region.tick_drainer_epoch())
1470 }
1471
1472 pub fn ordering_region(&self) -> Option<&OrderingRegion> {
1477 self.ordering.as_ref().map(|o| &o.region)
1478 }
1479
1480 pub fn try_send(&self, producer_id: usize, payload: &[u8]) -> Result<(), RingError> {
1492 self.ensure_synced();
1493 let shape = RingShape::from_u8(self.shape_tag.load(Ordering::Acquire));
1494 if let Some(ord) = &self.ordering {
1495 return self.stamped_send_inner(ord, shape, producer_id, payload);
1496 }
1497 match shape {
1498 RingShape::Spsc => self.spsc.try_push(payload),
1499 RingShape::Mpsc => {
1500 let rings = self.mpsc.rings.load();
1501 let ring = rings.get(producer_id)
1502 .ok_or(RingError::PayloadTooLarge)?; ring.try_push(payload)
1504 }
1505 RingShape::Mpmc => {
1506 let rings = self.mpmc.rings.load();
1507 let ring = rings.get(producer_id)
1508 .ok_or(RingError::PayloadTooLarge)?;
1509 ring.try_push(payload)
1510 }
1511 RingShape::Vyukov => self.vyukov.try_push(payload),
1512 }
1513 }
1514
1515 pub fn try_recv(&self, consumer_id: usize, out: &mut [u8]) -> Result<usize, RingError> {
1525 self.ensure_synced();
1526 let shape = RingShape::from_u8(self.shape_tag.load(Ordering::Acquire));
1527 if let Some(ord) = &self.ordering {
1528 return self.ordered_recv_inner(ord, shape, consumer_id, out)
1529 .map(|(n, _stamp)| n);
1530 }
1531 if let Some(stale) = self.stale_shape()
1535 && stale != shape
1536 && Self::may_walk_stale(stale, consumer_id)
1537 && let Ok(n) = self.shape_pop(stale, consumer_id, out)
1538 {
1539 return Ok(n);
1540 }
1541 self.shape_pop(shape, consumer_id, out)
1542 }
1543
1544 pub const FRAME_INLINE_BUDGET: usize = PAYLOAD_BYTES - 5;
1550
1551 pub const FRAME_DEFAULT_BLOCK_SIZE: usize = 8192;
1556
1557 pub fn with_frames(self, block_size: usize, block_count: usize) -> Self {
1563 self.frame_region.get_or_init(|| {
1564 Arc::new(
1565 FrameRegion::create_anon(block_size, block_count)
1566 .expect("frame region create"),
1567 )
1568 });
1569 self
1570 }
1571
1572 fn frame_region(&self) -> &FrameRegion {
1573 self.frame_region.get_or_init(|| {
1574 let blocks = self.spsc.capacity().max(16);
1575 Arc::new(
1576 FrameRegion::create_anon(Self::FRAME_DEFAULT_BLOCK_SIZE, blocks)
1577 .expect("frame region create"),
1578 )
1579 })
1580 }
1581
1582 pub fn send_frame(&self, producer_id: usize, payload: &[u8])
1596 -> Result<FrameClass, RingError>
1597 {
1598 self.send_frame_as(producer_id, payload, LayoutHint::Auto)
1599 }
1600
1601 pub fn send_frame_as(&self, producer_id: usize, payload: &[u8], hint: LayoutHint)
1604 -> Result<FrameClass, RingError>
1605 {
1606 if self.ordering.is_some() {
1607 return Err(RingError::LayoutMismatch);
1608 }
1609 let inline = match hint {
1610 LayoutHint::ForceInline => {
1611 if payload.len() > Self::FRAME_INLINE_BUDGET {
1612 return Err(RingError::PayloadTooLarge);
1613 }
1614 true
1615 }
1616 LayoutHint::ForceOffset => false,
1617 LayoutHint::Auto => payload.len() <= Self::FRAME_INLINE_BUDGET,
1618 };
1619 let len = payload.len() as u32;
1620 if inline {
1621 let mut buf = [0u8; PAYLOAD_BYTES];
1624 buf[0] = FrameClass::Inline as u8;
1625 buf[1..5].copy_from_slice(&len.to_le_bytes());
1626 buf[5..5 + payload.len()].copy_from_slice(payload);
1627 self.try_send(producer_id, &buf[..5 + payload.len()])?;
1628 Ok(FrameClass::Inline)
1629 } else {
1630 let region = self.frame_region();
1631 if payload.len() > region.block_size() {
1632 return Err(RingError::PayloadTooLarge);
1633 }
1634 let idx = region.alloc().ok_or(RingError::Full)?;
1635 region.write_block(idx, payload);
1636 let mut buf = [0u8; 9];
1638 buf[0] = FrameClass::Offset as u8;
1639 buf[1..5].copy_from_slice(&len.to_le_bytes());
1640 buf[5..9].copy_from_slice(&idx.to_le_bytes());
1641 match self.try_send(producer_id, &buf) {
1642 Ok(()) => Ok(FrameClass::Offset),
1643 Err(e) => {
1644 region.free(idx);
1647 Err(e)
1648 }
1649 }
1650 }
1651 }
1652
1653 pub fn recv_frame(&self, consumer_id: usize, out: &mut Vec<u8>)
1661 -> Result<FrameClass, RingError>
1662 {
1663 if self.ordering.is_some() {
1664 return Err(RingError::LayoutMismatch);
1665 }
1666 let mut slot = [0u8; SPSC_PAYLOAD_BYTES];
1668 self.try_recv(consumer_id, &mut slot)?;
1669 let len = u32::from_le_bytes([slot[1], slot[2], slot[3], slot[4]]) as usize;
1670 out.clear();
1671 if slot[0] == FrameClass::Inline as u8 {
1672 out.extend_from_slice(&slot[5..5 + len]);
1673 Ok(FrameClass::Inline)
1674 } else {
1675 let idx = u32::from_le_bytes([slot[5], slot[6], slot[7], slot[8]]);
1676 let region = self.frame_region();
1677 region.read_block_into(idx, len, out);
1678 region.free(idx);
1679 Ok(FrameClass::Offset)
1680 }
1681 }
1682
1683 pub fn try_recv_with_stamp(
1689 &self,
1690 consumer_id: usize,
1691 out: &mut [u8],
1692 ) -> Result<(usize, u64), RingError> {
1693 let ord = self.ordering.as_ref().ok_or(RingError::NotStamped)?;
1694 let shape = RingShape::from_u8(self.shape_tag.load(Ordering::Acquire));
1695 self.ordered_recv_inner(ord, shape, consumer_id, out)
1696 }
1697
1698 fn stamped_send_inner(
1705 &self,
1706 ord: &OrderingState,
1707 shape: RingShape,
1708 producer_id: usize,
1709 payload: &[u8],
1710 ) -> Result<(), RingError> {
1711 if payload.len() > STAMPED_PAYLOAD_BYTES {
1712 return Err(RingError::PayloadTooLarge);
1713 }
1714 if producer_id >= ord.region.max_producers() {
1715 return Err(RingError::PayloadTooLarge);
1716 }
1717 let mpsc_guard;
1718 let mpmc_guard;
1719 let ring: &SpscRingCore = match shape {
1720 RingShape::Spsc => &self.spsc,
1721 RingShape::Mpsc => {
1722 mpsc_guard = self.mpsc.rings.load();
1723 mpsc_guard.get(producer_id).ok_or(RingError::PayloadTooLarge)?
1724 }
1725 RingShape::Mpmc => {
1726 mpmc_guard = self.mpmc.rings.load();
1727 mpmc_guard.get(producer_id).ok_or(RingError::PayloadTooLarge)?
1728 }
1729 RingShape::Vyukov => return Err(RingError::LayoutMismatch),
1734 };
1735 let stamp = ord.region.next_stamp(producer_id);
1736 let mut buf = [0u8; SPSC_PAYLOAD_BYTES];
1737 buf[..STAMP_BYTES].copy_from_slice(&stamp.to_le_bytes());
1738 buf[STAMP_BYTES..STAMP_BYTES + payload.len()].copy_from_slice(payload);
1739 let result = ring.try_push(&buf[..STAMP_BYTES + payload.len()]);
1740 ord.region.publish_watermark(producer_id, stamp);
1741 result
1742 }
1743
1744 fn ordered_recv_inner(
1748 &self,
1749 ord: &OrderingState,
1750 shape: RingShape,
1751 consumer_id: usize,
1752 out: &mut [u8],
1753 ) -> Result<(usize, u64), RingError> {
1754 if consumer_id >= CONSUMER_SLOT_CEILING {
1755 return Err(RingError::PayloadTooLarge);
1756 }
1757 if out.len() < STAMPED_PAYLOAD_BYTES {
1758 return Err(RingError::PayloadTooLarge);
1759 }
1760 let mode = ord.region.mode();
1761 match mode {
1762 OrderingMode::Unordered => {
1763 let mut buf = [0u8; SPSC_PAYLOAD_BYTES];
1764 let popped = self.stale_shape()
1768 .filter(|stale| *stale != shape)
1769 .filter(|stale| Self::may_walk_stale(*stale, consumer_id))
1770 .and_then(|stale| {
1771 self.shape_pop(stale, consumer_id, &mut buf).ok()
1772 });
1773 if popped.is_none() {
1774 match shape {
1775 RingShape::Spsc => self.spsc.try_pop(&mut buf),
1776 RingShape::Mpsc => self.mpsc_pop(&mut buf),
1777 RingShape::Mpmc => self.mpmc_pop(consumer_id, &mut buf),
1778 RingShape::Vyukov => Err(RingError::LayoutMismatch),
1779 }?;
1780 }
1781 let stamp = u64::from_le_bytes(
1782 buf[..STAMP_BYTES].try_into().unwrap(),
1783 );
1784 self.note_stamp(ord, consumer_id, mode, stamp);
1785 out[..STAMPED_PAYLOAD_BYTES]
1786 .copy_from_slice(&buf[STAMP_BYTES..]);
1787 Ok((STAMPED_PAYLOAD_BYTES, stamp))
1788 }
1789 OrderingMode::MergeByStamp | OrderingMode::MergeStrict => {
1790 let seen = &ord.seen[consumer_id];
1801 let lease_gen_now = ord.region.lease_generation();
1802 if seen.lease_gen.load(Ordering::Relaxed) != lease_gen_now {
1803 if !ord.region.try_acquire_drainer(
1804 drainer_token(consumer_id),
1805 DRAINER_GRACE_EPOCHS,
1806 ) {
1807 return Err(RingError::NotDrainer);
1808 }
1809 seen.lease_gen.store(lease_gen_now, Ordering::Relaxed);
1810 }
1811 if let Some(stale) = self.stale_shape()
1817 && stale != shape
1818 {
1819 match self.merge_pop(ord, stale, consumer_id, mode, out) {
1820 Ok(result) => return Ok(result),
1821 Err(RingError::Empty) => {}
1822 Err(e) => return Err(e),
1823 }
1824 }
1825 self.merge_pop(ord, shape, consumer_id, mode, out)
1826 }
1827 }
1828 }
1829
1830
1831 fn merge_pop(
1862 &self,
1863 ord: &OrderingState,
1864 shape: RingShape,
1865 consumer_id: usize,
1866 mode: OrderingMode,
1867 out: &mut [u8],
1868 ) -> Result<(usize, u64), RingError> {
1869 let mpsc_guard;
1872 let mpmc_guard;
1873 let rings: &[Arc<SpscRingCore>] = match shape {
1874 RingShape::Spsc => std::slice::from_ref(&self.spsc),
1875 RingShape::Mpsc => {
1876 mpsc_guard = self.mpsc.rings.load();
1877 mpsc_guard.as_slice()
1878 }
1879 RingShape::Mpmc => {
1880 mpmc_guard = self.mpmc.rings.load();
1881 mpmc_guard.as_slice()
1882 }
1883 RingShape::Vyukov => &[],
1884 };
1885 if rings.is_empty() {
1886 return Err(RingError::LayoutMismatch);
1887 }
1888 let gate_lines = self.directory.published()
1893 .min(ord.region.max_producers());
1894 let kind = ord.region.stamp_kind();
1895 loop {
1896 let mut best: Option<(usize, u64)> = None;
1900 let mut any_empty = false;
1901 for (i, ring) in rings.iter().enumerate() {
1902 match ring.peek_slot() {
1903 Some(peek) => {
1904 let s = u64::from_le_bytes(
1905 peek[..STAMP_BYTES].try_into().unwrap(),
1906 );
1907 if best.is_none_or(|(_, bs)| s < bs) {
1908 best = Some((i, s));
1909 }
1910 }
1911 None => any_empty = true,
1912 }
1913 }
1914 let Some((idx, stamp)) = best else {
1915 return Err(RingError::Empty);
1916 };
1917
1918 for line in 0..gate_lines {
1919 if ord.region.in_flight_below(line, stamp) {
1920 return Err(RingError::Empty);
1921 }
1922 }
1923 if mode == OrderingMode::MergeStrict {
1924 for line in 0..gate_lines {
1925 if ord.region.issued(line) != 0
1930 && rings[line.min(rings.len() - 1)].approx_len() == 0
1931 && ord.region.watermark(line) < stamp
1932 {
1933 return Err(RingError::Empty);
1934 }
1935 }
1936 }
1937 if any_empty
1938 && let Some(guard) = kind.freshness_guard()
1939 && stamp_now(kind).wrapping_sub(stamp) < guard
1940 {
1941 std::hint::spin_loop();
1942 continue;
1943 }
1944
1945 let peek = rings[idx].peek_slot().expect(
1946 "single drainer holds the lease; a peeked head cannot vanish",
1947 );
1948 let confirmed_stamp = u64::from_le_bytes(
1949 peek[..STAMP_BYTES].try_into().unwrap(),
1950 );
1951 out[..STAMPED_PAYLOAD_BYTES].copy_from_slice(&peek[STAMP_BYTES..]);
1952 peek.confirm();
1953 self.note_stamp(ord, consumer_id, mode, confirmed_stamp);
1954 return Ok((STAMPED_PAYLOAD_BYTES, confirmed_stamp));
1955 }
1956 }
1957
1958 fn note_stamp(
1965 &self,
1966 ord: &OrderingState,
1967 consumer_id: usize,
1968 mode: OrderingMode,
1969 stamp: u64,
1970 ) {
1971 let line = &ord.seen[consumer_id];
1972 if line.mode_tag.swap(mode as u32, Ordering::Relaxed) != mode as u32 {
1973 line.stamp.store(0, Ordering::Relaxed);
1974 }
1975 let last = line.stamp.load(Ordering::Relaxed);
1976 if stamp < last {
1977 ord.region.record_inversion();
1978 self.ring_sidecar
1979 .push_op(crate::sidecar_ops::ordering::OP_ORDER_INVERSION, 0);
1980 }
1981 line.stamp.store(stamp, Ordering::Relaxed);
1982 }
1983
1984 fn mpsc_pop(&self, out: &mut [u8]) -> Result<usize, RingError> {
1985 let rings = self.mpsc.rings.load();
1986 self.mpsc_pop_in(&rings, out)
1987 }
1988
1989 fn mpsc_pop_in(
1990 &self,
1991 rings: &[Arc<SpscRingCore>],
1992 out: &mut [u8],
1993 ) -> Result<usize, RingError> {
1994 let n = rings.len();
1995 if n == 0 {
1996 return Err(RingError::Empty);
1997 }
1998 let start = self.mpsc.next_drain.load(Ordering::Relaxed);
1999 for i in 0..n {
2000 let idx = (start + i) % n;
2001 if let Ok(bytes) = rings[idx].try_pop(out) {
2002 self.mpsc.next_drain.store((idx + 1) % n, Ordering::Relaxed);
2003 return Ok(bytes);
2004 }
2005 }
2006 Err(RingError::Empty)
2007 }
2008
2009 fn mpmc_pop(&self, consumer_id: usize, out: &mut [u8]) -> Result<usize, RingError> {
2017 let rings = self.mpmc.rings.load();
2018 self.mpmc_pop_in(&rings, consumer_id, out)
2019 }
2020
2021 fn mpmc_pop_in(
2022 &self,
2023 rings: &[Arc<SpscRingCore>],
2024 consumer_id: usize,
2025 out: &mut [u8],
2026 ) -> Result<usize, RingError> {
2027 let cursor_line = self.mpmc.consumer_cursors.get(consumer_id)
2028 .ok_or(RingError::PayloadTooLarge)?;
2029 let me = consumer_id as u16;
2030 let n = rings.len();
2031 if n == 0 {
2032 return Err(RingError::Empty);
2033 }
2034 let start = cursor_line.0.load(Ordering::Relaxed) % n;
2035 let mut stuck: Option<(usize, u16)> = None;
2038 for i in 0..n {
2039 let idx = (start + i) % n;
2040 let (owner, pending) = self.directory.ring_owner(idx);
2041 if owner == me {
2042 if pending != OWNER_NONE
2043 && self.directory.apply_handoff(idx, me).is_some()
2044 {
2045 continue; }
2047 if let Ok(bytes) = rings[idx].try_pop(out) {
2048 cursor_line.0.store((idx + 1) % n, Ordering::Relaxed);
2049 return Ok(bytes);
2050 }
2051 } else if owner == OWNER_NONE {
2052 if self.directory.try_claim_ring(idx, me)
2053 && let Ok(bytes) = rings[idx].try_pop(out)
2054 {
2055 cursor_line.0.store((idx + 1) % n, Ordering::Relaxed);
2056 return Ok(bytes);
2057 }
2058 } else if stuck.is_none() && rings[idx].approx_len() > 0 {
2059 stuck = Some((idx, owner));
2060 }
2061 }
2062 if let Some((idx, owner)) = stuck {
2065 let probes = cursor_line.1.fetch_add(1, Ordering::Relaxed);
2066 if probes % 1024 == 1023
2067 && self.directory.try_takeover(idx, owner, me)
2068 && let Ok(bytes) = rings[idx].try_pop(out)
2069 {
2070 cursor_line.0.store((idx + 1) % n, Ordering::Relaxed);
2071 return Ok(bytes);
2072 }
2073 }
2074 Err(RingError::Empty)
2075 }
2076
2077 pub fn pin_current_shape(&self) -> PinnedRing<'_> {
2084 self.ensure_synced();
2085 let captured_gen = self.pin_generation.load(Ordering::Acquire);
2086 let shape = RingShape::from_u8(self.shape_tag.load(Ordering::Acquire));
2087 PinnedRing {
2088 parent: self,
2089 pinned_generation: captured_gen,
2090 shape,
2091 mpsc_rings: self.mpsc.rings.load_full(),
2092 mpmc_rings: self.mpmc.rings.load_full(),
2093 _not_sync: PhantomData,
2094 }
2095 }
2096
2097 pub fn morph_to(&self, new_shape: RingShape) -> Result<(), RingError> {
2124 self.shape_auto.store(false, Ordering::Relaxed);
2125 self.morph_shape(new_shape)
2126 }
2127
2128 pub(crate) fn morph_shape(&self, new_shape: RingShape) -> Result<(), RingError> {
2132 let old_shape = RingShape::from_u8(self.shape_tag.load(Ordering::Acquire));
2133 if old_shape == new_shape {
2134 return Ok(());
2135 }
2136
2137 if self.ordering.is_some() && new_shape == RingShape::Vyukov {
2143 return Err(RingError::LayoutMismatch);
2144 }
2145
2146 let prior_stale = self.stale_shape_tag.load(Ordering::Acquire);
2149 if prior_stale != STALE_NONE
2150 && !self.backing_is_empty(RingShape::from_u8(prior_stale))
2151 {
2152 return Err(RingError::StaleBacklog);
2153 }
2154
2155 self.pin_generation.fetch_add(1, Ordering::AcqRel);
2160 self.stale_shape_tag.store(old_shape as u8, Ordering::Release);
2161 self.shape_tag.store(new_shape as u8, Ordering::Release);
2162 Ok(())
2163 }
2164
2165 fn backing_is_empty(&self, shape: RingShape) -> bool {
2167 match shape {
2168 RingShape::Spsc => self.spsc.approx_len() == 0,
2169 RingShape::Mpsc => self.mpsc.rings.load().iter().all(|r| r.approx_len() == 0),
2170 RingShape::Mpmc => self.mpmc.rings.load().iter().all(|r| r.approx_len() == 0),
2171 RingShape::Vyukov => self.vyukov.approx_len() == 0,
2172 }
2173 }
2174
2175 fn stale_shape(&self) -> Option<RingShape> {
2177 let tag = self.stale_shape_tag.load(Ordering::Acquire);
2178 if tag == STALE_NONE {
2179 None
2180 } else {
2181 Some(RingShape::from_u8(tag))
2182 }
2183 }
2184
2185 fn may_walk_stale(shape: RingShape, consumer_id: usize) -> bool {
2190 match shape {
2191 RingShape::Spsc | RingShape::Mpsc => consumer_id == 0,
2192 RingShape::Mpmc | RingShape::Vyukov => true,
2193 }
2194 }
2195
2196 fn shape_pop(
2198 &self,
2199 shape: RingShape,
2200 consumer_id: usize,
2201 out: &mut [u8],
2202 ) -> Result<usize, RingError> {
2203 match shape {
2204 RingShape::Spsc => self.spsc.try_pop(out),
2205 RingShape::Mpsc => self.mpsc_pop(out),
2206 RingShape::Mpmc => self.mpmc_pop(consumer_id, out),
2207 RingShape::Vyukov => self.vyukov.try_pop(out),
2208 }
2209 }
2210}
2211
2212fn with_suffix(base: &std::path::Path, suffix: &str) -> std::path::PathBuf {
2213 let mut s = base.as_os_str().to_owned();
2214 s.push(suffix);
2215 std::path::PathBuf::from(s)
2216}
2217
2218#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2220pub enum AdaptiveError {
2221 TooManyProducers,
2226 TooManyConsumers,
2229 GrowthFailed,
2232}
2233
2234pub struct PinnedRing<'a> {
2241 parent: &'a AdaptiveRing,
2242 pinned_generation: u64,
2243 shape: RingShape,
2244 mpsc_rings: Arc<Vec<Arc<SpscRingCore>>>,
2248 mpmc_rings: Arc<Vec<Arc<SpscRingCore>>>,
2249 _not_sync: PhantomData<Cell<()>>,
2250}
2251
2252impl<'a> PinnedRing<'a> {
2253 pub fn shape(&self) -> RingShape { self.shape }
2255
2256 pub fn recv_signal(&self, shape: RingShape) -> &AtomicU64 {
2273 match shape {
2274 RingShape::Spsc => self.parent.spsc.head_signal(),
2275 RingShape::Mpsc => self.mpsc_rings[0].head_signal(),
2276 RingShape::Mpmc => self.mpmc_rings[0].head_signal(),
2277 RingShape::Vyukov => self.parent.vyukov.next_pop_signal(),
2278 }
2279 }
2280
2281 pub fn is_still_valid(&self) -> bool {
2285 self.parent.pin_generation.load(Ordering::Acquire) == self.pinned_generation
2286 }
2287
2288 pub fn spsc_try_push(&self, payload: &[u8]) -> Result<(), RingError> {
2291 self.parent.spsc.try_push(payload)
2292 }
2293
2294 pub fn spsc_try_pop(&self, out: &mut [u8]) -> Result<usize, RingError> {
2296 self.parent.spsc.try_pop(out)
2297 }
2298
2299 pub fn mpsc_try_push(&self, producer_id: usize, payload: &[u8]) -> Result<(), RingError> {
2301 let ring = self.mpsc_rings.get(producer_id)
2302 .ok_or(RingError::PayloadTooLarge)?;
2303 ring.try_push(payload)
2304 }
2305
2306 pub fn mpsc_try_pop(&self, out: &mut [u8]) -> Result<usize, RingError> {
2309 self.parent.mpsc_pop_in(&self.mpsc_rings, out)
2310 }
2311
2312 pub fn mpmc_try_push(&self, producer_id: usize, payload: &[u8]) -> Result<(), RingError> {
2314 let ring = self.mpmc_rings.get(producer_id)
2315 .ok_or(RingError::PayloadTooLarge)?;
2316 ring.try_push(payload)
2317 }
2318
2319 pub fn mpmc_try_pop(&self, consumer_id: usize, out: &mut [u8]) -> Result<usize, RingError> {
2323 self.parent.mpmc_pop_in(&self.mpmc_rings, consumer_id, out)
2324 }
2325
2326 pub fn vyukov_try_push(&self, payload: &[u8]) -> Result<(), RingError> {
2328 self.parent.vyukov.try_push(payload)
2329 }
2330
2331 pub fn vyukov_try_pop(&self, out: &mut [u8]) -> Result<usize, RingError> {
2333 self.parent.vyukov.try_pop(out)
2334 }
2335
2336 pub fn stamped_try_push(
2342 &self,
2343 producer_id: usize,
2344 payload: &[u8],
2345 ) -> Result<(), RingError> {
2346 let ord = self.parent.ordering.as_ref().ok_or(RingError::NotStamped)?;
2347 self.parent.stamped_send_inner(ord, self.shape, producer_id, payload)
2348 }
2349
2350 pub fn ordered_try_pop(
2358 &self,
2359 consumer_id: usize,
2360 out: &mut [u8],
2361 ) -> Result<usize, RingError> {
2362 let ord = self.parent.ordering.as_ref().ok_or(RingError::NotStamped)?;
2363 self.parent
2364 .ordered_recv_inner(ord, self.shape, consumer_id, out)
2365 .map(|(n, _stamp)| n)
2366 }
2367
2368 pub fn ordered_try_pop_with_stamp(
2372 &self,
2373 consumer_id: usize,
2374 out: &mut [u8],
2375 ) -> Result<(usize, u64), RingError> {
2376 let ord = self.parent.ordering.as_ref().ok_or(RingError::NotStamped)?;
2377 self.parent.ordered_recv_inner(ord, self.shape, consumer_id, out)
2378 }
2379}
2380
2381pub struct PeekedSpscSlot<'a> {
2386 inner: crate::spsc_ring::PeekedSlot<'a>,
2387}
2388
2389impl<'a> PeekedSpscSlot<'a> {
2390 pub fn as_slice(&self) -> &[u8] { self.inner.as_slice() }
2391 pub fn len(&self) -> usize { self.inner.len() }
2392 pub fn is_empty(&self) -> bool { self.inner.is_empty() }
2393 pub fn confirm(self) { self.inner.confirm() }
2394}
2395
2396impl<'a> std::ops::Deref for PeekedSpscSlot<'a> {
2397 type Target = [u8];
2398 fn deref(&self) -> &[u8] { &self.inner }
2399}
2400
2401pub const ADAPTIVE_SPSC_PAYLOAD_BYTES: usize = SPSC_PAYLOAD_BYTES;
2404
2405pub const ADAPTIVE_VYUKOV_PAYLOAD_BYTES: usize = PAYLOAD_BYTES;
2408
2409#[derive(Debug, Clone, Copy)]
2417pub struct PolicyObservation {
2418 pub active_producers: usize,
2419 pub active_consumers: usize,
2420 pub current_shape: RingShape,
2421 pub since_last_morph: std::time::Duration,
2422 pub stamped: bool,
2428}
2429
2430pub trait RingShapePolicy: Send + Sync + 'static {
2438 fn decide(&self, observation: &PolicyObservation) -> Option<RingShape>;
2439}
2440
2441pub struct DefaultRingShapePolicy {
2457 pub hysteresis: std::time::Duration,
2458}
2459
2460impl Default for DefaultRingShapePolicy {
2461 fn default() -> Self {
2462 Self { hysteresis: std::time::Duration::from_millis(100) }
2463 }
2464}
2465
2466impl DefaultRingShapePolicy {
2467 pub fn target_shape(producers: usize, consumers: usize) -> Option<RingShape> {
2468 match (producers, consumers) {
2469 (0, _) | (_, 0) => None,
2470 (1, 1) => Some(RingShape::Spsc),
2471 (_, 1) => Some(RingShape::Mpsc),
2472 (_, _) => Some(RingShape::Mpmc),
2473 }
2474 }
2475}
2476
2477impl RingShapePolicy for DefaultRingShapePolicy {
2478 fn decide(&self, obs: &PolicyObservation) -> Option<RingShape> {
2479 if obs.since_last_morph < self.hysteresis {
2480 return None;
2481 }
2482 let target = Self::target_shape(obs.active_producers, obs.active_consumers)?;
2483 if target == obs.current_shape {
2484 None
2485 } else {
2486 Some(target)
2487 }
2488 }
2489}
2490
2491pub struct QosRingShapePolicy {
2504 pub qos: Arc<crate::qos_policy::QosPolicy>,
2505 pub hysteresis: std::time::Duration,
2506}
2507
2508impl QosRingShapePolicy {
2509 pub fn new(qos: Arc<crate::qos_policy::QosPolicy>) -> Self {
2510 Self { qos, hysteresis: std::time::Duration::from_millis(100) }
2511 }
2512}
2513
2514impl RingShapePolicy for QosRingShapePolicy {
2515 fn decide(&self, obs: &PolicyObservation) -> Option<RingShape> {
2516 if obs.since_last_morph < self.hysteresis {
2517 return None;
2518 }
2519 let target = match self.qos.ordering() {
2520 QosOrdering::GlobalFifo if !obs.stamped => Some(RingShape::Vyukov),
2521 _ => DefaultRingShapePolicy::target_shape(
2522 obs.active_producers,
2523 obs.active_consumers,
2524 ),
2525 }?;
2526 if target == obs.current_shape {
2527 None
2528 } else {
2529 Some(target)
2530 }
2531 }
2532}
2533
2534#[derive(Debug, Clone, Copy)]
2537pub struct OrderingPolicyObservation {
2538 pub inversions_per_sec: f64,
2542 pub current_mode: OrderingMode,
2544 pub declared: QosOrdering,
2546 pub active_producers: usize,
2547 pub active_consumers: usize,
2548 pub since_last_change: std::time::Duration,
2550}
2551
2552pub trait OrderingPolicy: Send + Sync + 'static {
2557 fn decide(&self, observation: &OrderingPolicyObservation) -> Option<OrderingMode>;
2558}
2559
2560pub struct DefaultOrderingPolicy {
2575 pub hysteresis: std::time::Duration,
2576 pub auto_order_threshold: Option<f64>,
2577}
2578
2579impl Default for DefaultOrderingPolicy {
2580 fn default() -> Self {
2581 Self {
2582 hysteresis: std::time::Duration::from_millis(100),
2583 auto_order_threshold: None,
2584 }
2585 }
2586}
2587
2588impl OrderingPolicy for DefaultOrderingPolicy {
2589 fn decide(&self, obs: &OrderingPolicyObservation) -> Option<OrderingMode> {
2590 if obs.since_last_change < self.hysteresis {
2591 return None;
2592 }
2593 match obs.declared {
2594 QosOrdering::GlobalFifo => {
2595 if obs.current_mode == OrderingMode::Unordered {
2596 Some(OrderingMode::MergeByStamp)
2597 } else {
2598 None
2599 }
2600 }
2601 QosOrdering::PerProducer => {
2602 match self.auto_order_threshold {
2603 Some(threshold) => {
2604 if obs.current_mode == OrderingMode::Unordered
2605 && obs.inversions_per_sec > threshold
2606 {
2607 Some(OrderingMode::MergeByStamp)
2608 } else {
2609 None
2610 }
2611 }
2612 None => {
2613 if obs.current_mode != OrderingMode::Unordered {
2614 Some(OrderingMode::Unordered)
2615 } else {
2616 None
2617 }
2618 }
2619 }
2620 }
2621 }
2622 }
2623}
2624
2625pub struct AdaptiveRingSidecar {
2633 handle: Option<std::thread::JoinHandle<()>>,
2634 stop: Arc<std::sync::atomic::AtomicBool>,
2635 morphs_triggered: Arc<std::sync::atomic::AtomicU64>,
2636 ordering_flips: Arc<std::sync::atomic::AtomicU64>,
2637}
2638
2639impl AdaptiveRingSidecar {
2640 pub fn spawn<P: RingShapePolicy>(
2643 ring: Arc<AdaptiveRing>,
2644 policy: P,
2645 scan_interval: std::time::Duration,
2646 ) -> Self {
2647 Self::spawn_gated(
2648 ring,
2649 policy,
2650 scan_interval,
2651 crate::policy_gate::GateConfig::default(),
2652 )
2653 }
2654
2655 pub fn spawn_gated<P: RingShapePolicy>(
2659 ring: Arc<AdaptiveRing>,
2660 policy: P,
2661 scan_interval: std::time::Duration,
2662 gate_cfg: crate::policy_gate::GateConfig,
2663 ) -> Self {
2664 let stop = Arc::new(std::sync::atomic::AtomicBool::new(false));
2665 let morphs_triggered = Arc::new(std::sync::atomic::AtomicU64::new(0));
2666
2667 let stop_c = stop.clone();
2668 let morphs_c = morphs_triggered.clone();
2669 let handle = std::thread::spawn(move || {
2670 let mut last_morph = std::time::Instant::now();
2671 let mut gate = crate::policy_gate::ConfidenceGate::new(gate_cfg);
2672 while !stop_c.load(Ordering::Acquire) {
2673 let obs = PolicyObservation {
2674 active_producers: ring.active_producers(),
2675 active_consumers: ring.active_consumers(),
2676 current_shape: ring.current_shape(),
2677 since_last_morph: last_morph.elapsed(),
2678 stamped: ring.is_stamped(),
2679 };
2680 if let Some(new_shape) = gate
2684 .observe(policy.decide(&obs).map(|s| ring.contract_filtered_shape(s)))
2685 && ring.morph_shape(new_shape).is_ok()
2686 {
2687 last_morph = std::time::Instant::now();
2688 morphs_c.fetch_add(1, Ordering::Relaxed);
2689 }
2690 std::thread::sleep(scan_interval);
2691 }
2692 });
2693
2694 Self {
2695 handle: Some(handle),
2696 stop,
2697 morphs_triggered,
2698 ordering_flips: Arc::new(std::sync::atomic::AtomicU64::new(0)),
2699 }
2700 }
2701
2702 pub fn spawn_with_qos<P: RingShapePolicy, O: OrderingPolicy>(
2731 ring: Arc<AdaptiveRing>,
2732 shape_policy: P,
2733 ordering_policy: O,
2734 qos: Arc<crate::qos_policy::QosPolicy>,
2735 scan_interval: std::time::Duration,
2736 ) -> Self {
2737 Self::spawn_with_qos_core(
2738 ring,
2739 shape_policy,
2740 ordering_policy,
2741 qos,
2742 scan_interval,
2743 crate::policy_gate::GateConfig::default(),
2744 crate::policy_gate::GateConfig::enabled_with_arity(2),
2745 )
2746 }
2747
2748 pub fn spawn_with_qos_gated<P: RingShapePolicy, O: OrderingPolicy>(
2758 ring: Arc<AdaptiveRing>,
2759 shape_policy: P,
2760 ordering_policy: O,
2761 qos: Arc<crate::qos_policy::QosPolicy>,
2762 scan_interval: std::time::Duration,
2763 gate_cfg: crate::policy_gate::GateConfig,
2764 ) -> Self {
2765 Self::spawn_with_qos_core(
2766 ring,
2767 shape_policy,
2768 ordering_policy,
2769 qos,
2770 scan_interval,
2771 gate_cfg,
2772 gate_cfg,
2773 )
2774 }
2775
2776 fn spawn_with_qos_core<P: RingShapePolicy, O: OrderingPolicy>(
2777 ring: Arc<AdaptiveRing>,
2778 shape_policy: P,
2779 ordering_policy: O,
2780 qos: Arc<crate::qos_policy::QosPolicy>,
2781 scan_interval: std::time::Duration,
2782 shape_gate_cfg: crate::policy_gate::GateConfig,
2783 order_gate_cfg: crate::policy_gate::GateConfig,
2784 ) -> Self {
2785 let stop = Arc::new(std::sync::atomic::AtomicBool::new(false));
2786 let morphs_triggered = Arc::new(std::sync::atomic::AtomicU64::new(0));
2787 let ordering_flips = Arc::new(std::sync::atomic::AtomicU64::new(0));
2788
2789 let stop_c = stop.clone();
2790 let morphs_c = morphs_triggered.clone();
2791 let flips_c = ordering_flips.clone();
2792 let handle = std::thread::spawn(move || {
2793 let mut last_morph = std::time::Instant::now();
2794 let mut last_flip = std::time::Instant::now();
2795 let mut last_inversions = ring.inversions();
2796 let mut last_scan = std::time::Instant::now();
2797 let mut shape_gate = crate::policy_gate::ConfidenceGate::new(shape_gate_cfg);
2798 let mut order_gate = crate::policy_gate::ConfidenceGate::new(order_gate_cfg);
2799 let mut last_peers = (0usize, 0usize);
2800 let mut first_scan = true;
2801 while !stop_c.load(Ordering::Acquire) {
2802 let obs = PolicyObservation {
2803 active_producers: ring.active_producers(),
2804 active_consumers: ring.active_consumers(),
2805 current_shape: ring.current_shape(),
2806 since_last_morph: last_morph.elapsed(),
2807 stamped: ring.is_stamped(),
2808 };
2809 let peers = (obs.active_producers, obs.active_consumers);
2810 if !first_scan && peers != last_peers {
2811 shape_gate.shock();
2812 order_gate.shock();
2813 }
2814 last_peers = peers;
2815 first_scan = false;
2816
2817 if let Some(new_shape) = shape_gate
2818 .observe(shape_policy.decide(&obs).map(|s| ring.contract_filtered_shape(s)))
2819 && ring.morph_shape(new_shape).is_ok()
2820 {
2821 last_morph = std::time::Instant::now();
2822 morphs_c.fetch_add(1, Ordering::Relaxed);
2823 }
2824
2825 if let Some(current_mode) = ring.ordering_mode() {
2826 ring.tick_drainer_epoch().ok();
2827
2828 let now_inversions = ring.inversions();
2829 let elapsed = last_scan.elapsed().as_secs_f64().max(1e-9);
2830 let rate = now_inversions
2831 .saturating_sub(last_inversions) as f64 / elapsed;
2832 last_inversions = now_inversions;
2833 last_scan = std::time::Instant::now();
2834
2835 let declared = qos.ordering();
2836 let ord_obs = OrderingPolicyObservation {
2837 inversions_per_sec: rate,
2838 current_mode,
2839 declared,
2840 active_producers: obs.active_producers,
2841 active_consumers: obs.active_consumers,
2842 since_last_change: last_flip.elapsed(),
2843 };
2844 let decision = ordering_policy
2845 .decide(&ord_obs)
2846 .filter(|m| *m != current_mode);
2847
2848 let is_auto_arm = declared == QosOrdering::PerProducer
2857 && decision == Some(OrderingMode::MergeByStamp);
2858 let gated = if is_auto_arm {
2859 order_gate.observe(decision)
2860 } else {
2861 decision
2862 };
2863 if let Some(new_mode) = gated
2864 && ring.set_ordering_mode(new_mode).is_ok()
2865 {
2866 last_flip = std::time::Instant::now();
2867 flips_c.fetch_add(1, Ordering::Relaxed);
2868 }
2869 }
2870 std::thread::sleep(scan_interval);
2871 }
2872 });
2873
2874 Self {
2875 handle: Some(handle),
2876 stop,
2877 morphs_triggered,
2878 ordering_flips,
2879 }
2880 }
2881
2882 pub fn morphs_triggered(&self) -> u64 {
2885 self.morphs_triggered.load(Ordering::Acquire)
2886 }
2887
2888 pub fn ordering_flips(&self) -> u64 {
2891 self.ordering_flips.load(Ordering::Acquire)
2892 }
2893
2894 pub fn shutdown(mut self) {
2896 self.stop.store(true, Ordering::Release);
2897 if let Some(h) = self.handle.take() {
2898 h.join().ok();
2899 }
2900 }
2901}
2902
2903impl Drop for AdaptiveRingSidecar {
2904 fn drop(&mut self) {
2905 self.stop.store(true, Ordering::Release);
2906 if let Some(h) = self.handle.take() {
2907 h.join().ok();
2908 }
2909 }
2910}
2911
2912#[cfg(test)]
2913mod tests {
2914 use super::*;
2915
2916 #[test]
2917 fn create_starts_in_spsc_shape() {
2918 let ring = AdaptiveRing::create_anon(4, 4, 64).unwrap();
2919 assert_eq!(ring.current_shape(), RingShape::Spsc);
2920 assert_eq!(ring.pin_generation(), 0);
2921 }
2922
2923 #[test]
2924 fn adaptive_dispatch_round_trip_each_shape() {
2925 for shape in [RingShape::Spsc, RingShape::Mpsc, RingShape::Mpmc, RingShape::Vyukov] {
2926 let ring = AdaptiveRing::create_anon(4, 4, 64).unwrap();
2927 ring.shape_tag.store(shape as u8, Ordering::Release);
2928 let payload = [0xCDu8; 56];
2932 ring.try_send(0, &payload).unwrap();
2933 let mut out = [0u8; ADAPTIVE_SPSC_PAYLOAD_BYTES];
2934 let n = ring.try_recv(0, &mut out).unwrap();
2935 assert!(n > 0, "shape {:?} delivered zero bytes", shape);
2936 assert_eq!(&out[..payload.len()], &payload[..]);
2937 }
2938 }
2939
2940 #[test]
2941 fn pin_captures_shape_and_generation() {
2942 let ring = AdaptiveRing::create_anon(4, 4, 64).unwrap();
2943 let pinned = ring.pin_current_shape();
2944 assert_eq!(pinned.shape(), RingShape::Spsc);
2945 assert!(pinned.is_still_valid());
2946 let pinned = ring.pin_current_shape();
2948 assert!(pinned.is_still_valid());
2949 }
2950
2951 #[test]
2952 fn morph_invalidates_outstanding_pin() {
2953 let ring = AdaptiveRing::create_anon(4, 4, 64).unwrap();
2954 let pinned = ring.pin_current_shape();
2955 assert!(pinned.is_still_valid());
2956
2957 ring.morph_to(RingShape::Mpsc).unwrap();
2958 assert!(!pinned.is_still_valid(),
2959 "pin must invalidate after morph_to");
2960 assert_eq!(ring.current_shape(), RingShape::Mpsc);
2961 }
2962
2963 #[test]
2964 fn morph_to_same_shape_is_no_op() {
2965 let ring = AdaptiveRing::create_anon(4, 4, 64).unwrap();
2966 let gen_before = ring.pin_generation();
2967 ring.morph_to(RingShape::Spsc).unwrap();
2968 let gen_after = ring.pin_generation();
2969 assert_eq!(gen_before, gen_after,
2970 "morph_to(same shape) must not bump pin_generation");
2971 }
2972
2973 #[test]
2974 fn morph_preserves_in_flight_items_via_stale_walk() {
2975 let ring = AdaptiveRing::create_anon(4, 4, 64).unwrap();
2976
2977 for i in 0..3u32 {
2979 let mut buf = [0u8; ADAPTIVE_SPSC_PAYLOAD_BYTES];
2980 buf[..4].copy_from_slice(&i.to_le_bytes());
2981 ring.try_send(0, &buf).unwrap();
2982 }
2983
2984 ring.morph_to(RingShape::Mpsc).unwrap();
2987 assert_eq!(ring.current_shape(), RingShape::Mpsc);
2988 assert_eq!(ring.approx_len(), 3,
2989 "the stale backlog must stay visible through approx_len");
2990
2991 let mut buf = [0u8; ADAPTIVE_SPSC_PAYLOAD_BYTES];
2994 buf[..4].copy_from_slice(&99u32.to_le_bytes());
2995 ring.try_send(0, &buf).unwrap();
2996
2997 let mut seen = Vec::new();
2998 let mut out = [0u8; ADAPTIVE_SPSC_PAYLOAD_BYTES];
2999 while ring.try_recv(0, &mut out).is_ok() {
3000 seen.push(u32::from_le_bytes(out[..4].try_into().unwrap()));
3001 }
3002 assert_eq!(seen, vec![0u32, 1, 2, 99],
3003 "stale backlog must drain before post-morph items");
3004 assert!(ring.is_empty());
3005 }
3006
3007 #[test]
3008 fn frame_round_trip_all_shapes() {
3009 use crate::frame_ring::FrameClass;
3010 for shape in [RingShape::Spsc, RingShape::Mpsc, RingShape::Mpmc, RingShape::Vyukov] {
3011 let ring = AdaptiveRing::create_anon(4, 4, 64).unwrap();
3012 if shape != RingShape::Spsc {
3013 ring.morph_to(shape).unwrap();
3014 }
3015 let small = b"small inline payload".to_vec();
3016 let large = vec![0xABu8; 4000];
3017 assert_eq!(ring.send_frame(0, &small).unwrap(), FrameClass::Inline,
3018 "{shape:?} small should inline");
3019 assert_eq!(ring.send_frame(0, &large).unwrap(), FrameClass::Offset,
3020 "{shape:?} large should offset");
3021 let mut out = Vec::new();
3022 assert_eq!(ring.recv_frame(0, &mut out).unwrap(), FrameClass::Inline);
3023 assert_eq!(out, small, "{shape:?} small round-trip");
3024 assert_eq!(ring.recv_frame(0, &mut out).unwrap(), FrameClass::Offset);
3025 assert_eq!(out, large, "{shape:?} large round-trip");
3026 }
3027 }
3028
3029 #[test]
3030 fn frame_survives_morph() {
3031 let ring = AdaptiveRing::create_anon(4, 4, 64).unwrap();
3032 ring.send_frame(0, b"pre-morph small").unwrap();
3034 ring.send_frame(0, &vec![1u8; 3000]).unwrap();
3035 ring.morph_to(RingShape::Mpsc).unwrap();
3039 ring.send_frame(0, b"post-morph small").unwrap();
3040 ring.send_frame(0, &vec![2u8; 3000]).unwrap();
3041 let mut out = Vec::new();
3042 ring.recv_frame(0, &mut out).unwrap();
3043 assert_eq!(out, b"pre-morph small");
3044 ring.recv_frame(0, &mut out).unwrap();
3045 assert_eq!(out, vec![1u8; 3000]);
3046 ring.recv_frame(0, &mut out).unwrap();
3047 assert_eq!(out, b"post-morph small");
3048 ring.recv_frame(0, &mut out).unwrap();
3049 assert_eq!(out, vec![2u8; 3000]);
3050 }
3051
3052 #[test]
3053 fn frame_override_and_limits() {
3054 use crate::frame_ring::{FrameClass, LayoutHint};
3055 let ring = AdaptiveRing::create_anon(2, 2, 64).unwrap();
3056 let mut out = Vec::new();
3057 assert_eq!(ring.send_frame_as(0, b"tiny", LayoutHint::ForceOffset).unwrap(),
3059 FrameClass::Offset);
3060 assert_eq!(ring.recv_frame(0, &mut out).unwrap(), FrameClass::Offset);
3061 assert_eq!(out, b"tiny");
3062 let big = vec![0u8; AdaptiveRing::FRAME_INLINE_BUDGET + 1];
3064 assert_eq!(ring.send_frame_as(0, &big, LayoutHint::ForceInline).unwrap_err(),
3065 RingError::PayloadTooLarge);
3066 let at = vec![7u8; AdaptiveRing::FRAME_INLINE_BUDGET];
3068 assert_eq!(ring.send_frame(0, &at).unwrap(), FrameClass::Inline);
3069 ring.recv_frame(0, &mut out).unwrap();
3070 assert_eq!(out, at);
3071 }
3072
3073 #[test]
3074 fn frame_rejected_on_stamped_ring() {
3075 let ring = AdaptiveRing::create_anon(2, 2, 64)
3078 .unwrap()
3079 .with_ordering_stamps()
3080 .unwrap();
3081 assert_eq!(ring.send_frame(0, b"x").unwrap_err(), RingError::LayoutMismatch);
3082 let mut out = Vec::new();
3083 assert_eq!(ring.recv_frame(0, &mut out).unwrap_err(), RingError::LayoutMismatch);
3084 }
3085
3086 #[test]
3087 fn frame_vyukov_two_thread_mixed_size() {
3088 use std::sync::Arc;
3089 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering as AtOrd};
3090 use std::thread;
3091
3092 const PER: u32 = 5_000;
3093 const PRODUCERS: u32 = 2;
3094 let total = (PER * PRODUCERS) as usize;
3095
3096 let ring = Arc::new(AdaptiveRing::create_anon(2, 2, 256).unwrap());
3102 ring.morph_to(RingShape::Vyukov).unwrap();
3103 let seen: Arc<Vec<AtomicBool>> =
3106 Arc::new((0..total).map(|_| AtomicBool::new(false)).collect());
3107 let received = Arc::new(AtomicUsize::new(0));
3108
3109 let mut prods = Vec::new();
3110 for p in 0..PRODUCERS {
3111 let ring = ring.clone();
3112 prods.push(thread::spawn(move || {
3113 for i in 0..PER {
3114 let id = p * PER + i;
3115 let len = (id as usize % 200) + 4; let mut payload = vec![0u8; len];
3117 payload[0..4].copy_from_slice(&id.to_le_bytes());
3118 for k in 4..len {
3119 payload[k] = id.wrapping_add(k as u32) as u8;
3120 }
3121 while ring.send_frame(p as usize, &payload).is_err() {
3122 std::hint::spin_loop();
3123 }
3124 }
3125 }));
3126 }
3127
3128 let mut cons = Vec::new();
3129 for c in 0..2usize {
3130 let ring = ring.clone();
3131 let seen = seen.clone();
3132 let received = received.clone();
3133 cons.push(thread::spawn(move || {
3134 let mut out = Vec::new();
3135 while received.load(AtOrd::Acquire) < total {
3136 if ring.recv_frame(c, &mut out).is_ok() {
3137 let id = u32::from_le_bytes(out[0..4].try_into().unwrap());
3138 let len = (id as usize % 200) + 4;
3139 assert_eq!(out.len(), len, "id {id} length");
3140 for k in 4..len {
3141 assert_eq!(out[k], id.wrapping_add(k as u32) as u8,
3142 "id {id} byte {k}");
3143 }
3144 let already = seen[id as usize].swap(true, AtOrd::AcqRel);
3145 assert!(!already, "id {id} delivered twice");
3146 received.fetch_add(1, AtOrd::AcqRel);
3147 } else {
3148 std::hint::spin_loop();
3149 }
3150 }
3151 }));
3152 }
3153
3154 for p in prods { p.join().unwrap(); }
3155 for c in cons { c.join().unwrap(); }
3156 assert_eq!(received.load(AtOrd::Acquire), total);
3157 assert!(seen.iter().all(|b| b.load(AtOrd::Acquire)),
3158 "every id delivered exactly once");
3159 }
3160
3161 #[test]
3162 fn second_morph_blocked_until_stale_backlog_drains() {
3163 let ring = AdaptiveRing::create_anon(4, 4, 64).unwrap();
3164 ring.try_send(0, &[7u8; 8]).unwrap();
3165 ring.morph_to(RingShape::Mpsc).unwrap();
3166
3167 assert_eq!(ring.morph_to(RingShape::Mpmc).unwrap_err(),
3169 RingError::StaleBacklog);
3170
3171 let mut out = [0u8; ADAPTIVE_SPSC_PAYLOAD_BYTES];
3172 ring.try_recv(0, &mut out).unwrap();
3173 ring.morph_to(RingShape::Mpmc).unwrap();
3175 assert_eq!(ring.current_shape(), RingShape::Mpmc);
3176 }
3177
3178 #[test]
3179 fn register_producer_grows_past_hint_and_recycles_slots() {
3180 let ring = AdaptiveRing::create_anon(3, 1, 64).unwrap();
3181 let id0 = ring.register_producer().unwrap();
3182 let id1 = ring.register_producer().unwrap();
3183 let id2 = ring.register_producer().unwrap();
3184 assert_eq!((id0, id1, id2), (0, 1, 2));
3185
3186 let id3 = ring.register_producer().unwrap();
3189 assert_eq!(id3, 3);
3190 assert_eq!(ring.published_producers(), 4);
3191 ring.try_send(id3, &7u64.to_le_bytes()).unwrap();
3192 let mut out = [0u8; 64];
3193 let _c = ring.register_consumer().unwrap();
3197 let n = ring.try_recv(0, &mut out).unwrap();
3198 assert!(n >= 8, "popped record too short: {n}");
3199 assert_eq!(u64::from_le_bytes(out[..8].try_into().unwrap()), 7);
3200
3201 ring.unregister_producer(id1);
3204 assert_eq!(ring.register_producer().unwrap(), 1);
3205
3206 let pinned = AdaptiveRing::create_anon(2, 1, 64)
3208 .unwrap()
3209 .with_contract(crate::ring_contract::RingContract::from_counts(2, 1));
3210 pinned.register_producer().unwrap();
3211 pinned.register_producer().unwrap();
3212 assert_eq!(pinned.register_producer().unwrap_err(),
3213 AdaptiveError::TooManyProducers);
3214 }
3215
3216 #[test]
3217 fn default_policy_target_shape_per_peer_count() {
3218 assert_eq!(DefaultRingShapePolicy::target_shape(0, 1), None);
3220 assert_eq!(DefaultRingShapePolicy::target_shape(1, 0), None);
3221 assert_eq!(DefaultRingShapePolicy::target_shape(0, 0), None);
3222 assert_eq!(DefaultRingShapePolicy::target_shape(1, 1), Some(RingShape::Spsc));
3224 assert_eq!(DefaultRingShapePolicy::target_shape(2, 1), Some(RingShape::Mpsc));
3226 assert_eq!(DefaultRingShapePolicy::target_shape(8, 1), Some(RingShape::Mpsc));
3227 assert_eq!(DefaultRingShapePolicy::target_shape(1, 2), Some(RingShape::Mpmc));
3229 assert_eq!(DefaultRingShapePolicy::target_shape(4, 4), Some(RingShape::Mpmc));
3230 }
3231
3232 #[test]
3233 fn default_policy_returns_none_during_hysteresis() {
3234 let policy = DefaultRingShapePolicy {
3235 hysteresis: std::time::Duration::from_secs(1),
3236 };
3237 let obs = PolicyObservation {
3238 active_producers: 4,
3239 active_consumers: 4,
3240 current_shape: RingShape::Spsc,
3241 since_last_morph: std::time::Duration::from_millis(50),
3242 stamped: false,
3243 };
3244 assert_eq!(policy.decide(&obs), None);
3246 }
3247
3248 #[test]
3249 fn default_policy_returns_target_after_hysteresis() {
3250 let policy = DefaultRingShapePolicy {
3251 hysteresis: std::time::Duration::from_millis(10),
3252 };
3253 let obs = PolicyObservation {
3254 active_producers: 4,
3255 active_consumers: 4,
3256 current_shape: RingShape::Spsc,
3257 since_last_morph: std::time::Duration::from_secs(1),
3258 stamped: false,
3259 };
3260 assert_eq!(policy.decide(&obs), Some(RingShape::Mpmc));
3261 }
3262
3263 #[test]
3264 fn default_policy_returns_none_when_target_equals_current() {
3265 let policy = DefaultRingShapePolicy::default();
3266 let obs = PolicyObservation {
3267 active_producers: 1,
3268 active_consumers: 1,
3269 current_shape: RingShape::Spsc,
3270 since_last_morph: std::time::Duration::from_secs(1),
3271 stamped: false,
3272 };
3273 assert_eq!(policy.decide(&obs), None);
3274 }
3275
3276 #[test]
3277 fn shape_tracks_peer_counts_and_sidecar_stays_idle() {
3278 let ring = Arc::new(AdaptiveRing::create_anon(4, 4, 64).unwrap());
3279 let policy = DefaultRingShapePolicy {
3280 hysteresis: std::time::Duration::from_millis(5),
3281 };
3282 let sidecar = AdaptiveRingSidecar::spawn(
3283 ring.clone(),
3284 policy,
3285 std::time::Duration::from_millis(10),
3286 );
3287
3288 let _p0 = ring.register_producer().unwrap();
3290 let _c0 = ring.register_consumer().unwrap();
3291 assert_eq!(ring.current_shape(), RingShape::Spsc);
3292
3293 let _p1 = ring.register_producer().unwrap();
3296 assert_eq!(ring.current_shape(), RingShape::Mpsc,
3297 "2nd producer registration must morph to MPSC immediately");
3298
3299 let _c1 = ring.register_consumer().unwrap();
3300 assert_eq!(ring.current_shape(), RingShape::Mpmc,
3301 "2nd consumer registration must morph to MPMC immediately");
3302
3303 std::thread::sleep(std::time::Duration::from_millis(60));
3306 assert_eq!(sidecar.morphs_triggered(), 0,
3307 "register-path morphs left the sidecar nothing to do");
3308
3309 ring.unregister_consumer(1);
3312 ring.unregister_producer(1);
3313 assert_eq!(ring.current_shape(), RingShape::Spsc,
3314 "unregister must morph back down automatically");
3315
3316 sidecar.shutdown();
3317 }
3318
3319 #[test]
3320 fn pinned_native_paths_match_adaptive_paths() {
3321 let ring = AdaptiveRing::create_anon(4, 4, 64).unwrap();
3322 let pinned = ring.pin_current_shape();
3323 assert_eq!(pinned.shape(), RingShape::Spsc);
3324
3325 let payload = [0xAAu8; ADAPTIVE_SPSC_PAYLOAD_BYTES];
3326 pinned.spsc_try_push(&payload).unwrap();
3327 let mut out = [0u8; ADAPTIVE_SPSC_PAYLOAD_BYTES];
3328 let n = pinned.spsc_try_pop(&mut out).unwrap();
3329 assert_eq!(n, ADAPTIVE_SPSC_PAYLOAD_BYTES);
3330 assert_eq!(out, payload);
3331 assert!(pinned.is_still_valid());
3332 }
3333
3334 fn stamped_anon(
3339 max_producers: usize,
3340 max_consumers: usize,
3341 kind: StampKind,
3342 ) -> AdaptiveRing {
3343 AdaptiveRing::create_anon(max_producers, max_consumers, 64)
3344 .unwrap()
3345 .with_ordering_stamps_kind(kind)
3346 .unwrap()
3347 }
3348
3349 #[test]
3350 fn stamped_round_trip_strips_stamp_and_caps_payload() {
3351 let ring = stamped_anon(2, 1, StampKind::SharedCounter);
3352 assert!(ring.is_stamped());
3353 assert_eq!(ring.stamp_kind(), Some(StampKind::SharedCounter));
3354 assert_eq!(ring.ordering_mode(), Some(OrderingMode::Unordered));
3355
3356 let too_big = [0u8; STAMPED_PAYLOAD_BYTES + 1];
3358 assert_eq!(ring.try_send(0, &too_big).unwrap_err(),
3359 RingError::PayloadTooLarge);
3360
3361 let payload = [0xC3u8; STAMPED_PAYLOAD_BYTES];
3362 ring.try_send(0, &payload).unwrap();
3363 let mut out = [0u8; STAMPED_PAYLOAD_BYTES];
3364 let n = ring.try_recv(0, &mut out).unwrap();
3365 assert_eq!(n, STAMPED_PAYLOAD_BYTES,
3366 "stamped recv returns payload bytes only");
3367 assert_eq!(out, payload, "the stamp must be stripped, not leak into the payload");
3368 }
3369
3370 #[test]
3371 fn unstamped_ring_rejects_ordering_calls() {
3372 let ring = AdaptiveRing::create_anon(2, 1, 64).unwrap();
3373 assert!(!ring.is_stamped());
3374 assert_eq!(ring.inversions(), 0);
3375 assert_eq!(ring.set_ordering_mode(OrderingMode::MergeByStamp).unwrap_err(),
3376 RingError::NotStamped);
3377 assert_eq!(ring.refresh_watermark(0).unwrap_err(), RingError::NotStamped);
3378 let pinned = ring.pin_current_shape();
3379 let mut out = [0u8; STAMPED_PAYLOAD_BYTES];
3380 assert_eq!(pinned.ordered_try_pop(0, &mut out).unwrap_err(),
3381 RingError::NotStamped);
3382 assert_eq!(pinned.stamped_try_push(0, &[1u8; 8]).unwrap_err(),
3383 RingError::NotStamped);
3384 }
3385
3386 #[test]
3387 fn stamped_ring_rejects_vyukov_morph_and_vyukov_ring_rejects_stamps() {
3388 let ring = stamped_anon(2, 1, StampKind::Monotonic);
3389 assert_eq!(ring.morph_to(RingShape::Vyukov).unwrap_err(),
3390 RingError::LayoutMismatch);
3391
3392 let vyukov_first = AdaptiveRing::create_anon(2, 1, 64).unwrap();
3393 vyukov_first.morph_to(RingShape::Vyukov).unwrap();
3394 assert!(matches!(
3395 vyukov_first.with_ordering_stamps(),
3396 Err(RingError::LayoutMismatch)
3397 ));
3398 }
3399
3400 #[test]
3401 fn synthetic_interleave_fires_inversion_counter() {
3402 let ring = stamped_anon(2, 1, StampKind::SharedCounter);
3403 ring.morph_to(RingShape::Mpsc).unwrap();
3404
3405 ring.try_send(1, &1u64.to_le_bytes()).unwrap();
3410 ring.try_send(0, &2u64.to_le_bytes()).unwrap();
3411
3412 let mut out = [0u8; STAMPED_PAYLOAD_BYTES];
3413 ring.try_recv(0, &mut out).unwrap();
3414 assert_eq!(ring.inversions(), 0, "first pop has no predecessor");
3415 ring.try_recv(0, &mut out).unwrap();
3416 assert_eq!(ring.inversions(), 1,
3417 "older-after-newer must count as one inversion");
3418 }
3419
3420 #[test]
3421 fn merge_mode_delivers_global_stamp_order() {
3422 let ring = stamped_anon(4, 1, StampKind::SharedCounter);
3423 ring.morph_to(RingShape::Mpsc).unwrap();
3424 ring.set_ordering_mode(OrderingMode::MergeByStamp).unwrap();
3425
3426 for i in 0..32u64 {
3429 let producer = (i % 4) as usize;
3430 ring.try_send(producer, &i.to_le_bytes()).unwrap();
3431 }
3432
3433 let mut out = [0u8; STAMPED_PAYLOAD_BYTES];
3434 for expected in 0..32u64 {
3435 let n = ring.try_recv(0, &mut out).unwrap();
3436 assert_eq!(n, STAMPED_PAYLOAD_BYTES);
3437 let got = u64::from_le_bytes(out[..8].try_into().unwrap());
3438 assert_eq!(got, expected,
3439 "merge pop must deliver global push order");
3440 }
3441 assert_eq!(ring.try_recv(0, &mut out).unwrap_err(), RingError::Empty);
3442 assert_eq!(ring.inversions(), 0,
3443 "merged pops must observe zero inversions");
3444 }
3445
3446 #[test]
3447 fn flag_flip_orders_backlog_retroactively_without_loss() {
3448 let ring = stamped_anon(2, 1, StampKind::SharedCounter);
3449 ring.morph_to(RingShape::Mpsc).unwrap();
3450
3451 for i in 0..16u64 {
3454 let producer = ((i + 1) % 2) as usize;
3455 ring.try_send(producer, &i.to_le_bytes()).unwrap();
3456 }
3457
3458 let mut out = [0u8; STAMPED_PAYLOAD_BYTES];
3461 let mut popped = Vec::new();
3462 for _ in 0..2 {
3463 ring.try_recv(0, &mut out).unwrap();
3464 popped.push(u64::from_le_bytes(out[..8].try_into().unwrap()));
3465 }
3466 let inversions_before_flip = ring.inversions();
3467 assert!(inversions_before_flip > 0,
3468 "unordered interleave must show inversions before the flip");
3469
3470 ring.set_ordering_mode(OrderingMode::MergeByStamp).unwrap();
3473 let mut merged = Vec::new();
3474 while let Ok(_n) = ring.try_recv(0, &mut out) {
3475 merged.push(u64::from_le_bytes(out[..8].try_into().unwrap()));
3476 }
3477
3478 let mut all = popped.clone();
3480 all.extend(&merged);
3481 all.sort_unstable();
3482 assert_eq!(all, (0..16u64).collect::<Vec<_>>(),
3483 "no item may be lost across the mode flip");
3484 for pair in merged.windows(2) {
3487 assert!(pair[0] < pair[1],
3488 "post-flip pops must be globally ordered: {merged:?}");
3489 }
3490 assert_eq!(ring.inversions(), inversions_before_flip,
3491 "the flip itself and merged pops must add zero inversions");
3492 }
3493
3494 #[test]
3495 fn merge_strict_blocks_on_in_flight_stamp_then_releases() {
3496 let ring = stamped_anon(2, 1, StampKind::SharedCounter);
3497 ring.morph_to(RingShape::Mpsc).unwrap();
3498 ring.set_ordering_mode(OrderingMode::MergeStrict).unwrap();
3499 let region = ring.ordering_region().unwrap();
3500
3501 let stalled_stamp = region.next_stamp(1);
3505 ring.try_send(0, &42u64.to_le_bytes()).unwrap();
3507
3508 let mut out = [0u8; STAMPED_PAYLOAD_BYTES];
3511 assert_eq!(ring.try_recv(0, &mut out).unwrap_err(), RingError::Empty,
3512 "strict merge must hold the candidate while a smaller stamp is in flight");
3513
3514 region.publish_watermark(1, stalled_stamp);
3520 assert_eq!(ring.try_recv(0, &mut out).unwrap_err(), RingError::Empty,
3521 "strict watermark gate must hold until the silent producer vouches");
3522 ring.refresh_watermark(1).unwrap();
3525 let n = ring.try_recv(0, &mut out).unwrap();
3526 assert_eq!(n, STAMPED_PAYLOAD_BYTES);
3527 assert_eq!(u64::from_le_bytes(out[..8].try_into().unwrap()), 42);
3528 }
3529
3530 #[test]
3531 fn merge_by_stamp_in_flight_gate_blocks_descheduled_producer() {
3532 let ring = stamped_anon(2, 1, StampKind::SharedCounter);
3537 ring.morph_to(RingShape::Mpsc).unwrap();
3538 ring.set_ordering_mode(OrderingMode::MergeByStamp).unwrap();
3539 let region = ring.ordering_region().unwrap();
3540
3541 let stalled = region.next_stamp(1); ring.try_send(0, &9u64.to_le_bytes()).unwrap();
3543
3544 let mut out = [0u8; STAMPED_PAYLOAD_BYTES];
3545 assert_eq!(ring.try_recv(0, &mut out).unwrap_err(), RingError::Empty,
3546 "MergeByStamp must gate on in-flight stamps too");
3547 region.publish_watermark(1, stalled); let n = ring.try_recv(0, &mut out).unwrap();
3549 assert_eq!(n, STAMPED_PAYLOAD_BYTES);
3550 assert_eq!(u64::from_le_bytes(out[..8].try_into().unwrap()), 9);
3551 }
3552
3553 #[test]
3554 fn merge_strict_retired_producer_stops_gating() {
3555 let ring = stamped_anon(2, 1, StampKind::SharedCounter);
3556 ring.morph_to(RingShape::Mpsc).unwrap();
3557 ring.set_ordering_mode(OrderingMode::MergeStrict).unwrap();
3558
3559 ring.try_send(1, &1u64.to_le_bytes()).unwrap();
3563 let mut out = [0u8; STAMPED_PAYLOAD_BYTES];
3564 ring.try_recv(0, &mut out).unwrap();
3565
3566 ring.try_send(0, &2u64.to_le_bytes()).unwrap();
3568 assert_eq!(ring.try_recv(0, &mut out).unwrap_err(), RingError::Empty,
3569 "strict couples release to the slowest in-use producer");
3570
3571 ring.retire_producer(1).unwrap();
3574 let n = ring.try_recv(0, &mut out).unwrap();
3575 assert_eq!(n, STAMPED_PAYLOAD_BYTES);
3576 assert_eq!(u64::from_le_bytes(out[..8].try_into().unwrap()), 2);
3577 }
3578
3579 #[test]
3580 fn multi_consumer_merge_enforces_single_drainer() {
3581 let ring = stamped_anon(2, 2, StampKind::SharedCounter);
3582 ring.morph_to(RingShape::Mpmc).unwrap();
3583 ring.set_ordering_mode(OrderingMode::MergeByStamp).unwrap();
3584
3585 for i in 0..4u64 {
3586 ring.try_send((i % 2) as usize, &i.to_le_bytes()).unwrap();
3587 }
3588
3589 let mut out = [0u8; STAMPED_PAYLOAD_BYTES];
3591 ring.try_recv(0, &mut out).unwrap();
3592 assert_eq!(ring.try_recv(1, &mut out).unwrap_err(),
3594 RingError::NotDrainer);
3595 assert!(ring.release_drainer(0).unwrap());
3597 let n = ring.try_recv(1, &mut out).unwrap();
3598 assert_eq!(n, STAMPED_PAYLOAD_BYTES);
3599 assert_eq!(u64::from_le_bytes(out[..8].try_into().unwrap()), 1,
3600 "the new drainer continues in global stamp order");
3601 assert_eq!(ring.try_recv(0, &mut out).unwrap_err(),
3603 RingError::NotDrainer);
3604 }
3605
3606 #[test]
3607 fn mode_flip_does_not_invalidate_pins() {
3608 let ring = stamped_anon(2, 1, StampKind::SharedCounter);
3609 ring.morph_to(RingShape::Mpsc).unwrap();
3610 let pinned = ring.pin_current_shape();
3611 assert!(pinned.is_still_valid());
3612
3613 pinned.stamped_try_push(1, &1u64.to_le_bytes()).unwrap();
3615 pinned.stamped_try_push(0, &2u64.to_le_bytes()).unwrap();
3616
3617 ring.set_ordering_mode(OrderingMode::MergeByStamp).unwrap();
3621 assert!(pinned.is_still_valid(),
3622 "ordering-mode flips must not invalidate pins");
3623
3624 let mut out = [0u8; STAMPED_PAYLOAD_BYTES];
3625 pinned.ordered_try_pop(0, &mut out).unwrap();
3626 assert_eq!(u64::from_le_bytes(out[..8].try_into().unwrap()), 1,
3627 "pinned merge pop must deliver stamp order");
3628 pinned.ordered_try_pop(0, &mut out).unwrap();
3629 assert_eq!(u64::from_le_bytes(out[..8].try_into().unwrap()), 2);
3630 }
3631
3632 #[test]
3633 fn stamped_items_survive_shape_morphs() {
3634 let ring = stamped_anon(2, 1, StampKind::SharedCounter);
3635 for i in 0..3u64 {
3636 ring.try_send(0, &i.to_le_bytes()).unwrap();
3637 }
3638 ring.morph_to(RingShape::Mpsc).unwrap();
3639 let mut out = [0u8; STAMPED_PAYLOAD_BYTES];
3640 let mut got = Vec::new();
3641 while ring.try_recv(0, &mut out).is_ok() {
3642 got.push(u64::from_le_bytes(out[..8].try_into().unwrap()));
3643 }
3644 got.sort_unstable();
3645 assert_eq!(got, vec![0, 1, 2],
3646 "stamped slots must transfer intact across shape morphs");
3647 }
3648
3649 #[test]
3650 fn stamped_file_ring_open_adopts_creator_kind_and_shares_mode() {
3651 let mut prefix = std::env::temp_dir();
3652 prefix.push(format!(
3653 "subetha_stamped_open_{}_{}",
3654 std::process::id(),
3655 std::time::SystemTime::now()
3656 .duration_since(std::time::UNIX_EPOCH)
3657 .map(|d| d.as_nanos()).unwrap_or(0),
3658 ));
3659
3660 let creator = AdaptiveRing::create(&prefix, 2, 1, 64)
3661 .unwrap()
3662 .with_ordering_stamps_kind(StampKind::SharedCounter)
3663 .unwrap();
3664 creator.set_ordering_mode(OrderingMode::MergeByStamp).unwrap();
3665 creator.try_send(0, &7u64.to_le_bytes()).unwrap();
3666
3667 let opener = AdaptiveRing::open(&prefix, 2, 1, 64)
3668 .unwrap()
3669 .with_ordering_stamps()
3670 .unwrap();
3671 assert_eq!(opener.stamp_kind(), Some(StampKind::SharedCounter),
3672 "opener must adopt the creator's stamp kind");
3673 assert_eq!(opener.ordering_mode(), Some(OrderingMode::MergeByStamp),
3674 "the mode flag must be cross-process (region-resident)");
3675 assert!(matches!(
3677 AdaptiveRing::open(&prefix, 2, 1, 64)
3678 .unwrap()
3679 .with_ordering_stamps_kind(StampKind::Monotonic),
3680 Err(RingError::LayoutMismatch)
3681 ));
3682
3683 let mut out = [0u8; STAMPED_PAYLOAD_BYTES];
3684 let n = opener.try_recv(0, &mut out).unwrap();
3685 assert_eq!(n, STAMPED_PAYLOAD_BYTES);
3686 assert_eq!(u64::from_le_bytes(out[..8].try_into().unwrap()), 7);
3687
3688 drop(creator);
3689 drop(opener);
3690 for suffix in [".spsc.bin", ".mpsc.0.bin", ".mpsc.1.bin",
3691 ".mpmc.0.bin", ".mpmc.1.bin", ".vyukov.bin",
3692 ".ordering.bin"] {
3693 let mut p = prefix.as_os_str().to_owned();
3694 p.push(suffix);
3695 std::fs::remove_file(std::path::PathBuf::from(p)).ok();
3696 }
3697 }
3698
3699 #[test]
3700 fn qos_shape_policy_decision_matrix() {
3701 let qos = Arc::new(crate::qos_policy::QosPolicy::default());
3702 let policy = QosRingShapePolicy {
3703 qos: qos.clone(),
3704 hysteresis: std::time::Duration::from_millis(0),
3705 };
3706 let obs = |shape, stamped| PolicyObservation {
3707 active_producers: 2,
3708 active_consumers: 1,
3709 current_shape: shape,
3710 since_last_morph: std::time::Duration::from_secs(1),
3711 stamped,
3712 };
3713
3714 assert_eq!(policy.decide(&obs(RingShape::Spsc, false)),
3716 Some(RingShape::Mpsc));
3717 assert_eq!(policy.decide(&obs(RingShape::Mpsc, false)), None);
3718
3719 qos.set_ordering(crate::qos_policy::Ordering::GlobalFifo);
3721 assert_eq!(policy.decide(&obs(RingShape::Mpsc, false)),
3722 Some(RingShape::Vyukov));
3723 assert_eq!(policy.decide(&obs(RingShape::Vyukov, false)), None);
3724
3725 assert_eq!(policy.decide(&obs(RingShape::Spsc, true)),
3728 Some(RingShape::Mpsc));
3729 assert_eq!(policy.decide(&obs(RingShape::Mpsc, true)), None);
3730
3731 qos.set_ordering(crate::qos_policy::Ordering::PerProducer);
3733 assert_eq!(policy.decide(&obs(RingShape::Vyukov, false)),
3734 Some(RingShape::Mpsc));
3735
3736 let cold = QosRingShapePolicy {
3738 qos: qos.clone(),
3739 hysteresis: std::time::Duration::from_secs(10),
3740 };
3741 let mut o = obs(RingShape::Spsc, false);
3742 o.since_last_morph = std::time::Duration::from_millis(1);
3743 assert_eq!(cold.decide(&o), None);
3744 }
3745
3746 #[test]
3747 fn default_ordering_policy_decision_matrix() {
3748 let obs = |mode, declared, rate, since_ms| OrderingPolicyObservation {
3749 inversions_per_sec: rate,
3750 current_mode: mode,
3751 declared,
3752 active_producers: 2,
3753 active_consumers: 1,
3754 since_last_change: std::time::Duration::from_millis(since_ms),
3755 };
3756 let declarative = DefaultOrderingPolicy {
3757 hysteresis: std::time::Duration::from_millis(0),
3758 auto_order_threshold: None,
3759 };
3760 assert_eq!(
3762 declarative.decide(&obs(
3763 OrderingMode::Unordered, QosOrdering::GlobalFifo, 0.0, 500)),
3764 Some(OrderingMode::MergeByStamp),
3765 );
3766 assert_eq!(
3767 declarative.decide(&obs(
3768 OrderingMode::MergeByStamp, QosOrdering::GlobalFifo, 0.0, 500)),
3769 None,
3770 );
3771 assert_eq!(
3773 declarative.decide(&obs(
3774 OrderingMode::MergeByStamp, QosOrdering::PerProducer, 0.0, 500)),
3775 Some(OrderingMode::Unordered),
3776 );
3777
3778 let auto = DefaultOrderingPolicy {
3779 hysteresis: std::time::Duration::from_millis(0),
3780 auto_order_threshold: Some(100.0),
3781 };
3782 assert_eq!(
3784 auto.decide(&obs(
3785 OrderingMode::Unordered, QosOrdering::PerProducer, 50.0, 500)),
3786 None,
3787 );
3788 assert_eq!(
3790 auto.decide(&obs(
3791 OrderingMode::Unordered, QosOrdering::PerProducer, 250.0, 500)),
3792 Some(OrderingMode::MergeByStamp),
3793 );
3794 assert_eq!(
3796 auto.decide(&obs(
3797 OrderingMode::MergeByStamp, QosOrdering::PerProducer, 0.0, 500)),
3798 None,
3799 );
3800
3801 let cold = DefaultOrderingPolicy {
3803 hysteresis: std::time::Duration::from_secs(10),
3804 auto_order_threshold: Some(1.0),
3805 };
3806 assert_eq!(
3807 cold.decide(&obs(
3808 OrderingMode::Unordered, QosOrdering::GlobalFifo, 1e6, 1)),
3809 None,
3810 );
3811 }
3812
3813 #[test]
3814 fn sidecar_spawn_with_qos_flips_merge_flag_on_declaration() {
3815 let ring = Arc::new(stamped_anon(2, 1, StampKind::SharedCounter));
3816 ring.morph_to(RingShape::Mpsc).unwrap();
3817 let _p0 = ring.register_producer().unwrap();
3818 let _p1 = ring.register_producer().unwrap();
3819 let _c0 = ring.register_consumer().unwrap();
3820
3821 let qos = Arc::new(crate::qos_policy::QosPolicy::default());
3822 let sidecar = AdaptiveRingSidecar::spawn_with_qos(
3823 ring.clone(),
3824 QosRingShapePolicy {
3825 qos: qos.clone(),
3826 hysteresis: std::time::Duration::from_millis(5),
3827 },
3828 DefaultOrderingPolicy {
3829 hysteresis: std::time::Duration::from_millis(5),
3830 auto_order_threshold: None,
3831 },
3832 qos.clone(),
3833 std::time::Duration::from_millis(10),
3834 );
3835
3836 qos.set_ordering(crate::qos_policy::Ordering::GlobalFifo);
3839 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
3840 while std::time::Instant::now() < deadline
3841 && ring.ordering_mode() != Some(OrderingMode::MergeByStamp)
3842 {
3843 std::thread::sleep(std::time::Duration::from_millis(10));
3844 }
3845 assert_eq!(ring.ordering_mode(), Some(OrderingMode::MergeByStamp),
3846 "sidecar must arm the merge flag on the GlobalFifo declaration");
3847 assert_eq!(ring.current_shape(), RingShape::Mpsc,
3848 "stamped ring must stay composed (no Vyukov morph)");
3849 assert!(sidecar.ordering_flips() >= 1);
3850
3851 qos.set_ordering(crate::qos_policy::Ordering::PerProducer);
3853 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
3854 while std::time::Instant::now() < deadline
3855 && ring.ordering_mode() != Some(OrderingMode::Unordered)
3856 {
3857 std::thread::sleep(std::time::Duration::from_millis(10));
3858 }
3859 assert_eq!(ring.ordering_mode(), Some(OrderingMode::Unordered),
3860 "sidecar must disarm when the declaration is withdrawn");
3861 sidecar.shutdown();
3862 }
3863
3864 #[test]
3865 fn default_sidecar_gates_auto_arm_but_still_opens_on_sustained_inversions() {
3866 let ring = Arc::new(stamped_anon(2, 1, StampKind::SharedCounter));
3874 ring.morph_to(RingShape::Mpsc).unwrap();
3875 ring.register_producer().unwrap();
3876 ring.register_producer().unwrap();
3877 ring.register_consumer().unwrap();
3878 ring.set_ordering_mode(OrderingMode::Unordered).unwrap();
3879
3880 let qos = Arc::new(crate::qos_policy::QosPolicy::default());
3881 qos.set_ordering(crate::qos_policy::Ordering::PerProducer);
3882 let sidecar = AdaptiveRingSidecar::spawn_with_qos(
3883 ring.clone(),
3884 DefaultRingShapePolicy::default(),
3885 DefaultOrderingPolicy {
3886 hysteresis: std::time::Duration::from_millis(0),
3887 auto_order_threshold: Some(50.0),
3888 },
3889 qos.clone(),
3890 std::time::Duration::from_millis(5),
3891 );
3892
3893 let stop = Arc::new(std::sync::atomic::AtomicBool::new(false));
3894 let stop_c = stop.clone();
3895 let r = ring.clone();
3896 let consumer = std::thread::spawn(move || {
3897 let mut out = [0u8; 64];
3898 while !stop_c.load(Ordering::Acquire) {
3899 r.try_recv(0, &mut out).ok();
3900 }
3901 });
3902
3903 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(6);
3904 let mut seq = 0u64;
3905 while std::time::Instant::now() < deadline
3906 && ring.ordering_mode() != Some(OrderingMode::MergeByStamp)
3907 {
3908 ring.try_send(0, &seq.to_le_bytes()).ok();
3909 ring.try_send(1, &seq.to_le_bytes()).ok();
3910 seq += 1;
3911 }
3912 stop.store(true, Ordering::Release);
3913 consumer.join().unwrap();
3914
3915 assert_eq!(ring.ordering_mode(), Some(OrderingMode::MergeByStamp),
3916 "the gated auto-arm must still commit under sustained inversions");
3917 assert_eq!(sidecar.ordering_flips(), 1,
3918 "the one-way auto-arm fires exactly once");
3919 sidecar.shutdown();
3920 }
3921}