use std::fs::{File, OpenOptions};
use std::path::Path;
use std::sync::atomic::{AtomicU32, AtomicU64, Ordering as AtomOrd};
use memmap2::{MmapMut, MmapOptions};
use crate::shared_ring::RingError;
pub const ORDERING_MAGIC: u64 = 0x4F52_4452_0000_0001;
pub const STAMPED_PAYLOAD_BYTES: usize = 56;
pub const STAMP_BYTES: usize = 8;
pub const TSC_FRESHNESS_GUARD_CYCLES: u64 = 8192;
pub const MONOTONIC_FRESHNESS_GUARD_NANOS: u64 = 2_000;
#[repr(u32)]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OrderingMode {
Unordered = 0,
MergeByStamp = 1,
MergeStrict = 2,
}
impl OrderingMode {
fn from_u32(tag: u32) -> Self {
match tag {
0 => Self::Unordered,
1 => Self::MergeByStamp,
2 => Self::MergeStrict,
_ => panic!("OrderingHeader.mode corrupted: {tag}"),
}
}
}
#[repr(u32)]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StampKind {
Tsc = 0,
SharedCounter = 1,
Monotonic = 2,
}
impl StampKind {
fn from_u32(tag: u32) -> Option<Self> {
match tag {
0 => Some(Self::Tsc),
1 => Some(Self::SharedCounter),
2 => Some(Self::Monotonic),
_ => None,
}
}
pub fn has_time_semantics(self) -> bool {
matches!(self, Self::Tsc | Self::Monotonic)
}
pub fn freshness_guard(self) -> Option<u64> {
match self {
Self::Tsc => Some(TSC_FRESHNESS_GUARD_CYCLES),
Self::Monotonic => Some(MONOTONIC_FRESHNESS_GUARD_NANOS),
Self::SharedCounter => None,
}
}
}
pub fn default_stamp_kind() -> StampKind {
#[cfg(any(target_arch = "x86", target_arch = "x86_64"))]
{
if has_invariant_tsc() {
StampKind::Tsc
} else {
StampKind::SharedCounter
}
}
#[cfg(not(any(target_arch = "x86", target_arch = "x86_64")))]
{
StampKind::Monotonic
}
}
#[cfg(any(target_arch = "x86", target_arch = "x86_64"))]
pub fn has_invariant_tsc() -> bool {
#[cfg(target_arch = "x86_64")]
use core::arch::x86_64::__cpuid;
#[cfg(target_arch = "x86")]
use core::arch::x86::__cpuid;
let max_extended = __cpuid(0x8000_0000).eax;
if max_extended < 0x8000_0007 {
return false;
}
let power = __cpuid(0x8000_0007);
(power.edx & (1 << 8)) != 0
}
#[cfg(target_arch = "aarch64")]
pub fn has_invariant_tsc() -> bool {
true
}
#[cfg(not(any(
target_arch = "x86",
target_arch = "x86_64",
target_arch = "aarch64"
)))]
pub fn has_invariant_tsc() -> bool {
false
}
#[cfg(any(target_arch = "x86", target_arch = "x86_64"))]
#[inline]
pub fn read_tsc() -> u64 {
#[cfg(target_arch = "x86_64")]
unsafe {
core::arch::x86_64::_rdtsc()
}
#[cfg(target_arch = "x86")]
unsafe {
core::arch::x86::_rdtsc()
}
}
#[cfg(target_arch = "aarch64")]
#[inline]
pub fn read_tsc() -> u64 {
let v: u64;
unsafe {
core::arch::asm!(
"mrs {v}, cntvct_el0",
v = out(reg) v,
options(nomem, nostack, preserves_flags),
);
}
v
}
#[cfg(target_arch = "aarch64")]
#[inline]
pub fn counter_frequency_hz() -> u64 {
let v: u64;
unsafe {
core::arch::asm!(
"mrs {v}, cntfrq_el0",
v = out(reg) v,
options(nomem, nostack, preserves_flags),
);
}
v
}
#[cfg(not(any(
target_arch = "x86",
target_arch = "x86_64",
target_arch = "aarch64"
)))]
#[inline]
pub fn read_tsc() -> u64 {
monotonic_nanos()
}
#[cfg(unix)]
#[inline]
pub fn monotonic_nanos() -> u64 {
let mut ts = libc::timespec { tv_sec: 0, tv_nsec: 0 };
let rc = unsafe { libc::clock_gettime(libc::CLOCK_MONOTONIC, &mut ts) };
assert_eq!(rc, 0, "clock_gettime(CLOCK_MONOTONIC) failed");
(ts.tv_sec as u64).wrapping_mul(1_000_000_000).wrapping_add(ts.tv_nsec as u64)
}
#[cfg(windows)]
#[inline]
pub fn monotonic_nanos() -> u64 {
use windows_sys::Win32::System::Performance::{
QueryPerformanceCounter, QueryPerformanceFrequency,
};
use std::sync::OnceLock;
static FREQ: OnceLock<i64> = OnceLock::new();
let freq = *FREQ.get_or_init(|| {
let mut f: i64 = 0;
let ok = unsafe { QueryPerformanceFrequency(&mut f) };
assert!(ok != 0 && f > 0, "QueryPerformanceFrequency failed");
f
});
let mut count: i64 = 0;
let ok = unsafe { QueryPerformanceCounter(&mut count) };
assert!(ok != 0, "QueryPerformanceCounter failed");
let secs = (count as u64) / (freq as u64);
let rem = (count as u64) % (freq as u64);
secs.wrapping_mul(1_000_000_000)
.wrapping_add(rem.wrapping_mul(1_000_000_000) / (freq as u64))
}
#[inline]
pub(crate) fn stamp_now(kind: StampKind) -> u64 {
match kind {
StampKind::Tsc => read_tsc(),
StampKind::Monotonic => monotonic_nanos(),
StampKind::SharedCounter => 0,
}
}
#[repr(C, align(64))]
pub struct OrderingHeader {
pub magic: u64,
pub mode: AtomicU32,
pub stamp_kind: u32,
pub inversions: AtomicU64,
pub shared_stamp: AtomicU64,
pub drainer_token: AtomicU64,
pub drainer_heartbeat: AtomicU64,
pub drainer_epoch: AtomicU64,
_pad: [u8; 8],
}
#[repr(C, align(64))]
pub struct LeaseGenLine {
pub lease_generation: AtomicU64,
_pad: [u8; 56],
}
#[repr(C, align(64))]
pub struct ProducerLine {
pub issued: AtomicU64,
pub watermark: AtomicU64,
_pad: [u8; 48],
}
pub const fn ordering_region_size(max_producers: usize) -> usize {
std::mem::size_of::<OrderingHeader>()
+ std::mem::size_of::<LeaseGenLine>()
+ max_producers * std::mem::size_of::<ProducerLine>()
}
#[allow(dead_code)]
enum OrderingBacking {
Anon(MmapMut),
File(File, MmapMut),
Shm(crate::shm_file::ShmFile),
}
pub struct OrderingRegion {
_backing: OrderingBacking,
raw_ptr: *mut u8,
max_producers: usize,
kind: StampKind,
}
unsafe impl Send for OrderingRegion {}
unsafe impl Sync for OrderingRegion {}
unsafe fn init_ordering_layout(ptr: *mut u8, max_producers: usize, kind: StampKind) {
let header_ptr = ptr as *mut OrderingHeader;
unsafe {
std::ptr::write(header_ptr, OrderingHeader {
magic: ORDERING_MAGIC,
mode: AtomicU32::new(OrderingMode::Unordered as u32),
stamp_kind: kind as u32,
inversions: AtomicU64::new(0),
shared_stamp: AtomicU64::new(0),
drainer_token: AtomicU64::new(0),
drainer_heartbeat: AtomicU64::new(0),
drainer_epoch: AtomicU64::new(0),
_pad: [0; 8],
});
}
let gen_ptr = unsafe {
ptr.add(std::mem::size_of::<OrderingHeader>()) as *mut LeaseGenLine
};
unsafe {
std::ptr::write(gen_ptr, LeaseGenLine {
lease_generation: AtomicU64::new(0),
_pad: [0; 56],
});
}
let lines_base = unsafe {
ptr.add(std::mem::size_of::<OrderingHeader>()
+ std::mem::size_of::<LeaseGenLine>())
};
for i in 0..max_producers {
let line_ptr = unsafe {
lines_base.add(i * std::mem::size_of::<ProducerLine>()) as *mut ProducerLine
};
unsafe {
std::ptr::write(line_ptr, ProducerLine {
issued: AtomicU64::new(0),
watermark: AtomicU64::new(0),
_pad: [0; 48],
});
}
}
}
impl OrderingRegion {
pub fn create_anon(max_producers: usize, kind: StampKind) -> Result<Self, RingError> {
let total = ordering_region_size(max_producers);
let mut mmap = MmapOptions::new().len(total).map_anon()?;
unsafe { init_ordering_layout(mmap.as_mut_ptr(), max_producers, kind) };
let raw_ptr = mmap.as_mut_ptr();
Ok(Self {
_backing: OrderingBacking::Anon(mmap),
raw_ptr,
max_producers,
kind,
})
}
pub fn create(
path: impl AsRef<Path>,
max_producers: usize,
kind: StampKind,
) -> Result<Self, RingError> {
let total = ordering_region_size(max_producers);
let file = OpenOptions::new()
.read(true).write(true).create(true).truncate(true)
.open(path.as_ref())?;
file.set_len(total as u64)?;
let mut mmap = unsafe { MmapOptions::new().len(total).map_mut(&file)? };
unsafe { init_ordering_layout(mmap.as_mut_ptr(), max_producers, kind) };
let raw_ptr = mmap.as_mut_ptr();
Ok(Self {
_backing: OrderingBacking::File(file, mmap),
raw_ptr,
max_producers,
kind,
})
}
pub fn open(
path: impl AsRef<Path>,
max_producers: usize,
) -> Result<Self, RingError> {
let total = ordering_region_size(max_producers);
let file = OpenOptions::new().read(true).write(true).open(path.as_ref())?;
if (file.metadata()?.len() as usize) < total {
return Err(RingError::LayoutMismatch);
}
let mmap = unsafe { MmapOptions::new().len(total).map_mut(&file)? };
let header = unsafe { &*(mmap.as_ptr() as *const OrderingHeader) };
if header.magic != ORDERING_MAGIC {
return Err(RingError::LayoutMismatch);
}
let kind = StampKind::from_u32(header.stamp_kind)
.ok_or(RingError::LayoutMismatch)?;
let raw_ptr = mmap.as_ptr() as *mut u8;
Ok(Self {
_backing: OrderingBacking::File(file, mmap),
raw_ptr,
max_producers,
kind,
})
}
pub fn create_shm(
shm: crate::shm_file::ShmFile,
max_producers: usize,
kind: StampKind,
) -> Result<Self, RingError> {
let total = ordering_region_size(max_producers);
let mut shm = shm;
if shm.len() < total {
return Err(RingError::LayoutMismatch);
}
let raw_ptr = shm.as_mut_slice().as_mut_ptr();
unsafe { init_ordering_layout(raw_ptr, max_producers, kind) };
Ok(Self {
_backing: OrderingBacking::Shm(shm),
raw_ptr,
max_producers,
kind,
})
}
pub fn stamp_kind(&self) -> StampKind { self.kind }
pub fn max_producers(&self) -> usize { self.max_producers }
pub(crate) fn header(&self) -> &OrderingHeader {
unsafe { &*(self.raw_ptr as *const OrderingHeader) }
}
pub(crate) fn line(&self, producer_id: usize) -> &ProducerLine {
assert!(producer_id < self.max_producers,
"producer_id {producer_id} out of range (max {})", self.max_producers);
let lines_base = unsafe {
self.raw_ptr.add(std::mem::size_of::<OrderingHeader>()
+ std::mem::size_of::<LeaseGenLine>())
};
unsafe {
&*(lines_base.add(producer_id * std::mem::size_of::<ProducerLine>())
as *const ProducerLine)
}
}
fn lease_gen_line(&self) -> &LeaseGenLine {
unsafe {
&*(self.raw_ptr.add(std::mem::size_of::<OrderingHeader>())
as *const LeaseGenLine)
}
}
#[inline]
pub fn lease_generation(&self) -> u64 {
self.lease_gen_line().lease_generation.load(AtomOrd::Acquire)
}
fn bump_lease_generation(&self) {
self.lease_gen_line().lease_generation.fetch_add(1, AtomOrd::AcqRel);
}
pub fn mode(&self) -> OrderingMode {
OrderingMode::from_u32(self.header().mode.load(AtomOrd::Acquire))
}
pub fn set_mode(&self, mode: OrderingMode) {
self.header().mode.store(mode as u32, AtomOrd::Release);
}
pub fn inversions(&self) -> u64 {
self.header().inversions.load(AtomOrd::Relaxed)
}
pub(crate) fn record_inversion(&self) {
self.header().inversions.fetch_add(1, AtomOrd::Relaxed);
}
#[inline]
pub(crate) fn next_stamp(&self, producer_id: usize) -> u64 {
let line = self.line(producer_id);
let floor = line.issued.load(AtomOrd::Relaxed) + 1;
line.issued.store(floor, AtomOrd::Release);
let raw = match self.kind {
StampKind::Tsc => read_tsc(),
StampKind::Monotonic => monotonic_nanos(),
StampKind::SharedCounter => {
self.header().shared_stamp.fetch_add(1, AtomOrd::Relaxed) + 1
}
};
let stamp = raw.max(floor);
line.issued.store(stamp, AtomOrd::Release);
stamp
}
#[inline]
pub(crate) fn publish_watermark(&self, producer_id: usize, stamp: u64) {
self.line(producer_id).watermark.store(stamp, AtomOrd::Release);
}
#[inline]
pub(crate) fn in_flight_below(&self, producer_id: usize, candidate: u64) -> bool {
let line = self.line(producer_id);
let issued = line.issued.load(AtomOrd::Acquire);
if issued >= candidate {
return false;
}
issued != line.watermark.load(AtomOrd::Acquire)
}
pub fn refresh_watermark(&self, producer_id: usize) {
let line = self.line(producer_id);
let raw = match self.kind {
StampKind::Tsc => read_tsc(),
StampKind::Monotonic => monotonic_nanos(),
StampKind::SharedCounter => {
self.header().shared_stamp.fetch_add(1, AtomOrd::Relaxed) + 1
}
};
let floor = line.issued.load(AtomOrd::Relaxed) + 1;
let stamp = raw.max(floor);
line.issued.store(stamp, AtomOrd::Release);
line.watermark.store(stamp, AtomOrd::Release);
}
pub fn watermark(&self, producer_id: usize) -> u64 {
self.line(producer_id).watermark.load(AtomOrd::Acquire)
}
pub fn issued(&self, producer_id: usize) -> u64 {
self.line(producer_id).issued.load(AtomOrd::Acquire)
}
pub fn retire_producer(&self, producer_id: usize) {
let line = self.line(producer_id);
line.issued.store(u64::MAX, AtomOrd::Release);
line.watermark.store(u64::MAX, AtomOrd::Release);
}
pub fn seed_from(&self, other: &OrderingRegion) {
self.header().shared_stamp.store(
other.header().shared_stamp.load(AtomOrd::Acquire),
AtomOrd::Release,
);
self.header().inversions.store(
other.header().inversions.load(AtomOrd::Relaxed),
AtomOrd::Relaxed,
);
let n = self.max_producers.min(other.max_producers);
for i in 0..n {
self.line(i).issued.store(
other.line(i).issued.load(AtomOrd::Relaxed),
AtomOrd::Relaxed,
);
self.line(i).watermark.store(
other.line(i).watermark.load(AtomOrd::Acquire),
AtomOrd::Release,
);
}
self.set_mode(other.mode());
}
pub fn try_acquire_drainer(&self, token: u64, grace_epochs: u64) -> bool {
assert!(token != 0, "drainer token 0 is reserved for unleased");
let header = self.header();
loop {
let cur = header.drainer_token.load(AtomOrd::Acquire);
if cur == token {
let global = header.drainer_epoch.load(AtomOrd::Acquire);
if header.drainer_heartbeat.load(AtomOrd::Relaxed) != global {
header.drainer_heartbeat.store(global, AtomOrd::Release);
}
return true;
}
let can_claim = if cur == 0 {
true
} else {
let beat = header.drainer_heartbeat.load(AtomOrd::Acquire);
let global = header.drainer_epoch.load(AtomOrd::Acquire);
global.saturating_sub(beat) > grace_epochs
};
if !can_claim {
return false;
}
if header.drainer_token.compare_exchange(
cur, token, AtomOrd::AcqRel, AtomOrd::Acquire,
).is_ok() {
let global = header.drainer_epoch.load(AtomOrd::Acquire);
header.drainer_heartbeat.store(global, AtomOrd::Release);
self.bump_lease_generation();
return true;
}
std::hint::spin_loop();
}
}
pub fn release_drainer(&self, token: u64) -> bool {
let released = self.header().drainer_token
.compare_exchange(token, 0, AtomOrd::AcqRel, AtomOrd::Acquire)
.is_ok();
if released {
self.bump_lease_generation();
}
released
}
pub fn current_drainer(&self) -> u64 {
self.header().drainer_token.load(AtomOrd::Acquire)
}
pub fn drainer_beat(&self, token: u64) -> bool {
let header = self.header();
if header.drainer_token.load(AtomOrd::Acquire) != token {
return false;
}
let global = header.drainer_epoch.load(AtomOrd::Acquire);
header.drainer_heartbeat.store(global, AtomOrd::Release);
true
}
pub fn tick_drainer_epoch(&self) -> u64 {
let epoch = self.header().drainer_epoch.fetch_add(1, AtomOrd::AcqRel) + 1;
self.bump_lease_generation();
epoch
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn header_is_one_cache_line_and_lines_are_padded() {
assert_eq!(std::mem::size_of::<OrderingHeader>(), 64);
assert_eq!(std::mem::size_of::<LeaseGenLine>(), 64);
assert_eq!(std::mem::size_of::<ProducerLine>(), 64);
assert_eq!(ordering_region_size(4), 64 + 64 + 4 * 64);
}
#[test]
fn lease_generation_bumps_on_claim_release_and_tick_only() {
let region = OrderingRegion::create_anon(2, StampKind::SharedCounter).unwrap();
let g0 = region.lease_generation();
assert!(region.try_acquire_drainer(7, 3));
let g1 = region.lease_generation();
assert!(g1 > g0, "claim must bump the generation");
assert!(region.try_acquire_drainer(7, 3));
assert_eq!(region.lease_generation(), g1);
region.tick_drainer_epoch();
let g2 = region.lease_generation();
assert!(g2 > g1, "epoch tick must bump so the holder re-beats");
assert!(region.release_drainer(7));
assert!(region.lease_generation() > g2, "release must bump");
}
#[test]
fn probe_runs_without_fault_and_tsc_reads_advance() {
let invariant = has_invariant_tsc();
if invariant {
let a = read_tsc();
let mut spin = 0u64;
for i in 0..10_000u64 { spin = spin.wrapping_add(i); }
std::hint::black_box(spin);
let b = read_tsc();
assert!(b > a, "TSC must advance across real work: {a} -> {b}");
}
}
#[test]
fn monotonic_clock_advances() {
let a = monotonic_nanos();
std::thread::sleep(std::time::Duration::from_millis(2));
let b = monotonic_nanos();
assert!(b > a, "monotonic clock must advance: {a} -> {b}");
}
#[test]
fn default_kind_matches_probe_chain() {
let kind = default_stamp_kind();
#[cfg(any(target_arch = "x86", target_arch = "x86_64"))]
{
if has_invariant_tsc() {
assert_eq!(kind, StampKind::Tsc);
} else {
assert_eq!(kind, StampKind::SharedCounter);
}
}
#[cfg(not(any(target_arch = "x86", target_arch = "x86_64")))]
assert_eq!(kind, StampKind::Monotonic);
}
#[test]
fn stamps_strictly_increase_per_producer_every_kind() {
for kind in [StampKind::Tsc, StampKind::SharedCounter, StampKind::Monotonic] {
let region = OrderingRegion::create_anon(2, kind).unwrap();
let mut last = 0u64;
for _ in 0..10_000 {
let s = region.next_stamp(0);
assert!(s > last, "{kind:?} stamp must strictly increase: {last} -> {s}");
last = s;
}
}
}
#[test]
fn counter_stamps_are_globally_unique_across_producers() {
let region = OrderingRegion::create_anon(4, StampKind::SharedCounter).unwrap();
let mut seen = std::collections::HashSet::new();
for p in 0..4 {
for _ in 0..100 {
assert!(seen.insert(region.next_stamp(p)),
"counter stamps must never repeat");
}
}
}
#[test]
fn watermark_publishes_and_refreshes() {
let region = OrderingRegion::create_anon(2, StampKind::Monotonic).unwrap();
assert_eq!(region.watermark(0), 0);
let s = region.next_stamp(0);
region.publish_watermark(0, s);
assert_eq!(region.watermark(0), s);
region.refresh_watermark(0);
assert!(region.watermark(0) > s, "refresh must advance the watermark");
}
#[test]
fn mode_flips_round_trip() {
let region = OrderingRegion::create_anon(1, StampKind::Monotonic).unwrap();
assert_eq!(region.mode(), OrderingMode::Unordered);
region.set_mode(OrderingMode::MergeByStamp);
assert_eq!(region.mode(), OrderingMode::MergeByStamp);
region.set_mode(OrderingMode::MergeStrict);
assert_eq!(region.mode(), OrderingMode::MergeStrict);
region.set_mode(OrderingMode::Unordered);
assert_eq!(region.mode(), OrderingMode::Unordered);
}
#[test]
fn file_region_open_validates_and_adopts_kind() {
let p = std::env::temp_dir().join(format!(
"subetha-ordering-open-{}-{}.bin",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos()).unwrap_or(0),
));
{
let creator =
OrderingRegion::create(&p, 3, StampKind::SharedCounter).unwrap();
creator.set_mode(OrderingMode::MergeByStamp);
let s = creator.next_stamp(1);
creator.publish_watermark(1, s);
}
let opened = OrderingRegion::open(&p, 3).unwrap();
assert_eq!(opened.stamp_kind(), StampKind::SharedCounter,
"opener must adopt the creator's stamp kind");
assert_eq!(opened.mode(), OrderingMode::MergeByStamp,
"open must not re-initialise the live mode flag");
assert!(opened.watermark(1) > 0,
"open must not wipe watermarks");
std::fs::remove_file(&p).ok();
}
#[test]
fn open_rejects_garbage() {
let p = std::env::temp_dir().join(format!(
"subetha-ordering-garbage-{}.bin", std::process::id(),
));
std::fs::write(&p, vec![0u8; ordering_region_size(2)]).unwrap();
assert!(matches!(
OrderingRegion::open(&p, 2),
Err(RingError::LayoutMismatch)
));
std::fs::remove_file(&p).ok();
}
#[test]
fn drainer_lease_claim_refresh_release_takeover() {
let region = OrderingRegion::create_anon(1, StampKind::Monotonic).unwrap();
let a = 100u64 << 32;
let b = 200u64 << 32;
assert!(region.try_acquire_drainer(a, 3));
assert_eq!(region.current_drainer(), a);
assert!(region.try_acquire_drainer(a, 3));
assert!(!region.try_acquire_drainer(b, 3));
assert!(region.release_drainer(a));
assert!(region.try_acquire_drainer(b, 3));
for _ in 0..5 { region.tick_drainer_epoch(); }
assert!(region.try_acquire_drainer(a, 3),
"stale drainer must be preemptible after grace epochs");
assert_eq!(region.current_drainer(), a);
assert!(!region.drainer_beat(b));
assert!(region.drainer_beat(a));
}
#[test]
fn seed_from_carries_counter_and_watermarks() {
let old = OrderingRegion::create_anon(2, StampKind::SharedCounter).unwrap();
for _ in 0..50 { old.next_stamp(0); }
let w = old.next_stamp(1);
old.publish_watermark(1, w);
old.set_mode(OrderingMode::MergeStrict);
let fresh = OrderingRegion::create_anon(2, StampKind::SharedCounter).unwrap();
fresh.seed_from(&old);
assert_eq!(fresh.mode(), OrderingMode::MergeStrict);
assert_eq!(fresh.watermark(1), w);
let next = fresh.next_stamp(0);
assert!(next > w, "seeded counter must continue past the old region's stamps");
}
}