use core::cell::{Cell, UnsafeCell};
use core::fmt;
use core::marker::PhantomData;
use core::sync::atomic::{AtomicU32, AtomicU64, Ordering, fence};
use crate::layout::{AcknowledgementRoute, RoleId};
#[cfg(not(target_has_atomic = "64"))]
compile_error!("native-ipc-core requires lock-free 64-bit atomic support");
pub const SLOT_HEADER_SIZE: u64 = 64;
pub const ACKNOWLEDGEMENT_CELL_SIZE: u64 = 64;
#[repr(C, align(64))]
#[derive(Debug)]
pub struct SlotMetadata {
generation: AtomicU64,
payload_len: AtomicU32,
reserved_word: UnsafeCell<u32>,
published_sequence: AtomicU64,
reserved: UnsafeCell<[u8; 40]>,
}
impl SlotMetadata {
pub const fn new(generation: u64) -> Self {
Self {
generation: AtomicU64::new(generation),
payload_len: AtomicU32::new(0),
reserved_word: UnsafeCell::new(0),
published_sequence: AtomicU64::new(0),
reserved: UnsafeCell::new([0; 40]),
}
}
pub unsafe fn initialize(&mut self, generation: u64) -> Result<(), SlotError> {
if generation == 0 {
return Err(SlotError::ZeroGeneration);
}
*self.generation.get_mut() = generation;
*self.payload_len.get_mut() = 0;
unsafe { *self.reserved_word.get() = 0 };
*self.published_sequence.get_mut() = 0;
unsafe { *self.reserved.get() = [0; 40] };
Ok(())
}
}
unsafe impl Sync for SlotMetadata {}
#[repr(C, align(64))]
pub struct AcknowledgementCell {
sequence: AtomicU64,
reserved: UnsafeCell<[u8; 56]>,
}
impl AcknowledgementCell {
pub const fn new() -> Self {
Self {
sequence: AtomicU64::new(0),
reserved: UnsafeCell::new([0; 56]),
}
}
pub unsafe fn initialize(&mut self) {
*self.sequence.get_mut() = 0;
unsafe { *self.reserved.get() = [0; 56] };
}
}
unsafe impl Sync for AcknowledgementCell {}
impl Default for AcknowledgementCell {
fn default() -> Self {
Self::new()
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct WriterSlotBinding(SlotBinding);
impl WriterSlotBinding {
pub(crate) const fn validated(
role: RoleId,
generation: u64,
payload_capacity: u32,
slot_index: u32,
slot_count: u32,
acknowledgement_owner: RoleId,
acknowledgement_cell_index: u32,
) -> Self {
Self(SlotBinding {
role,
generation,
payload_capacity,
slot_index,
slot_count,
acknowledgement_owner: Some(acknowledgement_owner),
acknowledgement_cell_index: Some(acknowledgement_cell_index),
})
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct ReaderSlotBinding(SlotBinding);
impl ReaderSlotBinding {
pub(crate) const fn validated(
role: RoleId,
generation: u64,
payload_capacity: u32,
slot_index: u32,
slot_count: u32,
) -> Self {
Self(SlotBinding {
role,
generation,
payload_capacity,
slot_index,
slot_count,
acknowledgement_owner: None,
acknowledgement_cell_index: None,
})
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
struct SlotBinding {
role: RoleId,
generation: u64,
payload_capacity: u32,
slot_index: u32,
slot_count: u32,
acknowledgement_owner: Option<RoleId>,
acknowledgement_cell_index: Option<u32>,
}
pub struct WriterSlot<'a> {
header: &'a SlotMetadata,
binding: SlotBinding,
_not_sync: PhantomData<Cell<()>>,
}
impl<'a> WriterSlot<'a> {
pub unsafe fn bind(
header: &'a SlotMetadata,
binding: WriterSlotBinding,
) -> Result<Self, SlotError> {
validate_bound_generation(header, binding.0.generation)?;
Ok(Self {
header,
binding: binding.0,
_not_sync: PhantomData,
})
}
pub fn prepare_publish(
&mut self,
sequence: u64,
acknowledgement: Option<AcknowledgementObservation>,
) -> Result<PublishReservation<'_>, SlotError> {
validate_bound_generation(self.header, self.binding.generation)?;
validate_sequence_slot(self.binding, sequence)?;
let current = self.header.published_sequence.load(Ordering::Relaxed);
if current == 0 {
let expected = u64::from(self.binding.slot_index) + 1;
if sequence != expected {
return Err(SlotError::UnexpectedFirstSequence {
expected,
actual: sequence,
});
}
} else {
let expected = current
.checked_add(u64::from(self.binding.slot_count))
.ok_or(SlotError::SequenceWrap)?;
if sequence != expected {
return Err(SlotError::UnexpectedNextSequence {
expected,
actual: sequence,
});
}
let acknowledgement =
acknowledgement.ok_or(SlotError::MissingAcknowledgement { sequence: current })?;
if acknowledgement.target != self.binding.role {
return Err(SlotError::WrongAcknowledgementTarget);
}
if acknowledgement.owner != self.binding.acknowledgement_owner.unwrap() {
return Err(SlotError::WrongAcknowledgementOwner);
}
if acknowledgement.slot_index != self.binding.slot_index {
return Err(SlotError::WrongAcknowledgementSlot);
}
if acknowledgement.cell_index != self.binding.acknowledgement_cell_index.unwrap() {
return Err(SlotError::WrongAcknowledgementCell);
}
if acknowledgement.generation != self.binding.generation {
return Err(SlotError::StaleAcknowledgementGeneration);
}
if acknowledgement.sequence < current {
return Err(SlotError::LaggingAcknowledgement {
expected: current,
actual: acknowledgement.sequence,
});
}
if acknowledgement.sequence > current {
return Err(SlotError::FutureAcknowledgement {
expected: current,
actual: acknowledgement.sequence,
});
}
}
Ok(PublishReservation {
header: self.header,
sequence,
capacity: self.binding.payload_capacity,
_exclusive: PhantomData,
})
}
}
pub struct ReaderSlot<'a> {
header: &'a SlotMetadata,
binding: SlotBinding,
}
impl<'a> ReaderSlot<'a> {
pub unsafe fn bind(
header: &'a SlotMetadata,
binding: ReaderSlotBinding,
) -> Result<Self, SlotError> {
validate_bound_generation(header, binding.0.generation)?;
Ok(Self {
header,
binding: binding.0,
})
}
pub fn observe(&self, expected_sequence: u64) -> Result<SlotObservation, SlotError> {
validate_sequence_slot(self.binding, expected_sequence)?;
let sequence = self.header.published_sequence.load(Ordering::Acquire);
if sequence != expected_sequence {
return Err(SlotError::StaleSequence {
expected: expected_sequence,
actual: sequence,
});
}
validate_bound_generation(self.header, self.binding.generation)?;
let payload_len = self.header.payload_len.load(Ordering::Relaxed);
if payload_len > self.binding.payload_capacity {
return Err(SlotError::PayloadTooLarge {
length: payload_len,
capacity: self.binding.payload_capacity,
});
}
Ok(SlotObservation {
role: self.binding.role,
slot_index: self.binding.slot_index,
generation: self.binding.generation,
sequence,
payload_len,
})
}
pub fn recheck(&self, observation: SlotObservation) -> Result<(), SlotError> {
if observation.role != self.binding.role {
return Err(SlotError::WrongObservationRole);
}
if observation.generation != self.binding.generation {
return Err(SlotError::StaleGeneration {
expected: self.binding.generation,
actual: observation.generation,
});
}
fence(Ordering::SeqCst);
let sequence = self.header.published_sequence.load(Ordering::Acquire);
if sequence != observation.sequence {
return Err(SlotError::StaleSequence {
expected: observation.sequence,
actual: sequence,
});
}
validate_bound_generation(self.header, observation.generation)?;
let payload_len = self.header.payload_len.load(Ordering::Relaxed);
if payload_len != observation.payload_len {
return Err(SlotError::ChangedPayloadLength {
expected: observation.payload_len,
actual: payload_len,
});
}
Ok(())
}
}
#[must_use = "payload bytes remain unpublished until publish is called"]
#[derive(Debug)]
pub struct PublishReservation<'a> {
header: &'a SlotMetadata,
sequence: u64,
capacity: u32,
_exclusive: PhantomData<&'a mut ()>,
}
impl PublishReservation<'_> {
pub fn publish(self, payload_len: u32) -> Result<(), SlotError> {
if payload_len > self.capacity {
return Err(SlotError::PayloadTooLarge {
length: payload_len,
capacity: self.capacity,
});
}
self.header
.payload_len
.store(payload_len, Ordering::Relaxed);
self.header
.published_sequence
.store(self.sequence, Ordering::Release);
Ok(())
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct SlotObservation {
role: RoleId,
slot_index: u32,
generation: u64,
sequence: u64,
payload_len: u32,
}
impl SlotObservation {
pub const fn role(self) -> RoleId {
self.role
}
pub const fn slot_index(self) -> u32 {
self.slot_index
}
pub const fn generation(self) -> u64 {
self.generation
}
pub const fn sequence(self) -> u64 {
self.sequence
}
pub const fn payload_len(self) -> u32 {
self.payload_len
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct AcknowledgementWriterBinding(AcknowledgementBinding);
impl AcknowledgementWriterBinding {
pub(crate) const fn validated(route: AcknowledgementRoute, generation: u64) -> Self {
Self(AcknowledgementBinding {
owner: route.owner(),
target: route.target(),
slot_index: route.slot_index(),
cell_index: route.cell_index(),
generation,
})
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct AcknowledgementReaderBinding(AcknowledgementBinding);
impl AcknowledgementReaderBinding {
pub(crate) const fn validated(route: AcknowledgementRoute, generation: u64) -> Self {
Self(AcknowledgementBinding {
owner: route.owner(),
target: route.target(),
slot_index: route.slot_index(),
cell_index: route.cell_index(),
generation,
})
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
struct AcknowledgementBinding {
owner: RoleId,
target: RoleId,
slot_index: u32,
cell_index: u32,
generation: u64,
}
pub struct AcknowledgementWriter<'a> {
cell: &'a AcknowledgementCell,
binding: AcknowledgementBinding,
_not_sync: PhantomData<Cell<()>>,
}
impl<'a> AcknowledgementWriter<'a> {
pub unsafe fn bind(
cell: &'a AcknowledgementCell,
binding: AcknowledgementWriterBinding,
) -> Self {
Self {
cell,
binding: binding.0,
_not_sync: PhantomData,
}
}
pub fn acknowledge(
&mut self,
observation: SlotObservation,
) -> Result<(), AcknowledgementError> {
if observation.role != self.binding.target {
return Err(AcknowledgementError::WrongTarget);
}
if observation.generation != self.binding.generation {
return Err(AcknowledgementError::StaleGeneration);
}
if observation.slot_index != self.binding.slot_index {
return Err(AcknowledgementError::WrongSlot);
}
if observation.sequence == 0 {
return Err(AcknowledgementError::UnpublishedSequence);
}
let current = self.cell.sequence.load(Ordering::Relaxed);
if observation.sequence < current {
return Err(AcknowledgementError::NonMonotonic {
current,
next: observation.sequence,
});
}
if observation.sequence == current {
return Ok(());
}
self.cell
.sequence
.store(observation.sequence, Ordering::Release);
Ok(())
}
}
pub struct AcknowledgementReader<'a> {
cell: &'a AcknowledgementCell,
binding: AcknowledgementBinding,
}
impl<'a> AcknowledgementReader<'a> {
pub unsafe fn bind(
cell: &'a AcknowledgementCell,
binding: AcknowledgementReaderBinding,
) -> Self {
Self {
cell,
binding: binding.0,
}
}
pub fn observe(&self) -> AcknowledgementObservation {
AcknowledgementObservation {
owner: self.binding.owner,
target: self.binding.target,
generation: self.binding.generation,
slot_index: self.binding.slot_index,
cell_index: self.binding.cell_index,
sequence: self.cell.sequence.load(Ordering::Acquire),
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct AcknowledgementObservation {
owner: RoleId,
target: RoleId,
generation: u64,
slot_index: u32,
cell_index: u32,
sequence: u64,
}
impl AcknowledgementObservation {
pub const fn owner(self) -> RoleId {
self.owner
}
pub const fn target(self) -> RoleId {
self.target
}
pub const fn generation(self) -> u64 {
self.generation
}
pub const fn slot_index(self) -> u32 {
self.slot_index
}
pub const fn cell_index(self) -> u32 {
self.cell_index
}
pub const fn sequence(self) -> u64 {
self.sequence
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum SlotError {
ZeroGeneration,
UnpublishedSequence,
StaleGeneration {
expected: u64,
actual: u64,
},
SequenceWrap,
WrongSlot {
expected: u32,
actual: u32,
},
UnexpectedFirstSequence {
expected: u64,
actual: u64,
},
UnexpectedNextSequence {
expected: u64,
actual: u64,
},
MissingAcknowledgement {
sequence: u64,
},
WrongAcknowledgementTarget,
WrongAcknowledgementOwner,
WrongAcknowledgementSlot,
WrongAcknowledgementCell,
StaleAcknowledgementGeneration,
LaggingAcknowledgement {
expected: u64,
actual: u64,
},
FutureAcknowledgement {
expected: u64,
actual: u64,
},
StaleSequence {
expected: u64,
actual: u64,
},
PayloadTooLarge {
length: u32,
capacity: u32,
},
ChangedPayloadLength {
expected: u32,
actual: u32,
},
WrongObservationRole,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum AcknowledgementError {
WrongTarget,
StaleGeneration,
WrongSlot,
UnpublishedSequence,
NonMonotonic {
current: u64,
next: u64,
},
}
impl fmt::Display for SlotError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(formatter, "slot operation failed: {self:?}")
}
}
impl fmt::Display for AcknowledgementError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(formatter, "acknowledgement failed: {self:?}")
}
}
#[cfg(feature = "std")]
impl std::error::Error for SlotError {}
#[cfg(feature = "std")]
impl std::error::Error for AcknowledgementError {}
fn validate_bound_generation(header: &SlotMetadata, expected: u64) -> Result<(), SlotError> {
if expected == 0 {
return Err(SlotError::ZeroGeneration);
}
let actual = header.generation.load(Ordering::Relaxed);
if actual == expected {
Ok(())
} else {
Err(SlotError::StaleGeneration { expected, actual })
}
}
fn validate_sequence_slot(binding: SlotBinding, sequence: u64) -> Result<(), SlotError> {
if sequence == 0 {
return Err(SlotError::UnpublishedSequence);
}
let expected = ((sequence - 1) % u64::from(binding.slot_count)) as u32;
if binding.slot_index == expected {
Ok(())
} else {
Err(SlotError::WrongSlot {
expected,
actual: binding.slot_index,
})
}
}
const _: () = assert!(core::mem::size_of::<SlotMetadata>() == SLOT_HEADER_SIZE as usize);
const _: () = assert!(core::mem::align_of::<SlotMetadata>() == 64);
const _: () = assert!(core::mem::offset_of!(SlotMetadata, generation) == 0);
const _: () = assert!(core::mem::offset_of!(SlotMetadata, payload_len) == 8);
const _: () = assert!(core::mem::offset_of!(SlotMetadata, reserved_word) == 12);
const _: () = assert!(core::mem::offset_of!(SlotMetadata, published_sequence) == 16);
const _: () = assert!(core::mem::offset_of!(SlotMetadata, reserved) == 24);
const _: () = assert!(core::mem::size_of::<AcknowledgementCell>() == 64);
const _: () = assert!(core::mem::align_of::<AcknowledgementCell>() == 64);
const _: () = assert!(core::mem::offset_of!(AcknowledgementCell, sequence) == 0);
const _: () = assert!(core::mem::offset_of!(AcknowledgementCell, reserved) == 8);
#[cfg(test)]
#[path = "slot_test.rs"]
mod tests;