#![expect(clippy::cast_ptr_alignment, reason = "TODO COME BACK TO THIS")]
use errno::{Errno, errno};
use std::ffi::{CStr, c_void};
use std::mem::size_of;
use std::ptr;
use std::sync::atomic::{self, AtomicU16};
use crate::shm::shm_header::{CLOCKBOUND_SHM_LATEST_VERSION, ShmHeader};
use crate::{
shm::{
ClockErrorBound, ClockErrorBoundGeneric, ClockErrorBoundLayoutVersion, ClockErrorBoundV2,
ClockErrorBoundV3, ShmError,
},
syserror,
};
struct FdGuard(i32);
impl FdGuard {
fn new(path: &CStr) -> Result<Self, ShmError> {
let fd = unsafe { libc::open(path.as_ptr(), libc::O_RDONLY) };
if fd < 0 {
return syserror!(format!("Faild to open file at {:?}", path));
}
Ok(FdGuard(fd))
}
}
impl Drop for FdGuard {
fn drop(&mut self) {
unsafe {
let ret = libc::close(self.0);
assert!(ret == 0 || errno() == Errno(libc::EINTR));
}
}
}
#[derive(Debug)]
struct MmapGuard {
segment: *mut c_void,
segsize: usize,
}
impl MmapGuard {
fn new(fdguard: &FdGuard) -> Result<Self, ShmError> {
let header = ShmHeader::read(fdguard.0)?;
let segsize = header.segsize.into_inner() as usize;
let segment: *mut c_void = unsafe {
libc::mmap(
ptr::null_mut(),
segsize,
libc::PROT_READ,
libc::MAP_SHARED,
fdguard.0,
0,
)
};
if segment == libc::MAP_FAILED {
return syserror!(String::from("Failed to mmap the SHM segment"));
}
Ok(MmapGuard { segment, segsize })
}
}
impl Drop for MmapGuard {
fn drop(&mut self) {
unsafe {
let ret = libc::munmap(self.segment, self.segsize);
assert_eq!(0, ret);
}
}
}
#[derive(Debug)]
pub struct ShmReader {
_marker: std::marker::PhantomData<*const ()>,
_guard: MmapGuard,
version: *const atomic::AtomicU16,
generation: *const atomic::AtomicU16,
ceb_shm: *const ClockErrorBound,
snapshot_ceb: ClockErrorBound,
snapshot_gen: u16,
}
impl ShmReader {
#[expect(clippy::missing_errors_doc, reason = "todo")]
pub fn new(path: &CStr) -> Result<ShmReader, ShmError> {
let (mmap_guard, version, generation, ceb_shm) = ShmReader::map_segment(path)?;
let shm_version = unsafe { &*version };
let shm_version = shm_version.load(atomic::Ordering::Acquire);
let min_version = shm_version >> 8;
let max_version = shm_version & 0x00ff;
if CLOCKBOUND_SHM_LATEST_VERSION < min_version
|| CLOCKBOUND_SHM_LATEST_VERSION > max_version
{
let msg = format!(
"Clockbound shared memory segment supports versions {min_version} to {max_version} which does not include this reader version {CLOCKBOUND_SHM_LATEST_VERSION}",
);
return Err(ShmError::SegmentVersionNotSupported(msg));
}
let shm_version = ClockErrorBoundLayoutVersion::try_from(CLOCKBOUND_SHM_LATEST_VERSION)?;
Ok(ShmReader {
_marker: std::marker::PhantomData,
_guard: mmap_guard,
version,
generation,
ceb_shm,
snapshot_ceb: ClockErrorBoundGeneric::builder().build(shm_version),
snapshot_gen: 0,
})
}
#[expect(clippy::missing_errors_doc, reason = "todo")]
pub fn new_with_max_version_unchecked(path: &CStr) -> Result<(ShmReader, u16), ShmError> {
let (mmap_guard, version, generation, ceb_shm) = ShmReader::map_segment(path)?;
let shm_version = unsafe { &*version };
let shm_version = shm_version.load(atomic::Ordering::Acquire);
let max_version = shm_version & 0x00ff;
let current_version =
ClockErrorBoundLayoutVersion::try_from(CLOCKBOUND_SHM_LATEST_VERSION)?;
Ok((
ShmReader {
_marker: std::marker::PhantomData,
_guard: mmap_guard,
version,
generation,
ceb_shm,
snapshot_ceb: ClockErrorBoundGeneric::builder().build(current_version),
snapshot_gen: 0,
},
max_version,
))
}
fn map_segment(
path: &CStr,
) -> Result<
(
MmapGuard,
*const AtomicU16,
*const AtomicU16,
*const ClockErrorBound,
),
ShmError,
> {
let fdguard = FdGuard::new(path)?;
let mmap_guard = MmapGuard::new(&fdguard)?;
let mut cursor: *const u8 = mmap_guard.segment.cast();
let version = unsafe { ptr::addr_of!((*cursor.cast::<ShmHeader>()).version) };
let generation = unsafe { ptr::addr_of!((*cursor.cast::<ShmHeader>()).generation) };
let shm_version = unsafe { &*version };
let shm_version = shm_version.load(atomic::Ordering::Acquire);
let shm_version = ClockErrorBoundLayoutVersion::try_from(shm_version & 0x00ff)?;
let layout_size = match shm_version {
ClockErrorBoundLayoutVersion::V2 => size_of::<ClockErrorBoundV2>(),
ClockErrorBoundLayoutVersion::V3 => size_of::<ClockErrorBoundV3>(),
};
if mmap_guard.segsize < size_of::<ShmHeader>() + layout_size {
let msg = format!(
"Clockbound segment size is smaller than expected [{} < {}].",
mmap_guard.segsize,
size_of::<ShmHeader>() + layout_size
);
return Err(ShmError::SegmentMalformed(msg));
}
cursor = unsafe { cursor.add(size_of::<ShmHeader>()) };
let ceb_shm = ptr::addr_of!(*cursor.cast::<ClockErrorBound>());
Ok((mmap_guard, version, generation, ceb_shm))
}
#[expect(clippy::missing_errors_doc, reason = "todo")]
pub fn snapshot(&mut self) -> Result<&ClockErrorBound, ShmError> {
let version = unsafe { &*self.version };
let version = version.load(atomic::Ordering::Acquire);
if version == 0 {
return Ok(&self.snapshot_ceb);
}
let min_version = version >> 8;
let max_version = version & 0x00ff;
if CLOCKBOUND_SHM_LATEST_VERSION < min_version
|| CLOCKBOUND_SHM_LATEST_VERSION < max_version
{
let msg = format!(
"Clockbound shared memory segment supports versions {min_version} to {max_version} which does not include this reader version {CLOCKBOUND_SHM_LATEST_VERSION}",
);
return Err(ShmError::SegmentVersionNotSupported(msg));
}
let generation = unsafe { &*self.generation };
let mut first_gen = generation.load(atomic::Ordering::Acquire);
if first_gen == 0 {
return Ok(&self.snapshot_ceb);
}
if first_gen == self.snapshot_gen {
return Ok(&self.snapshot_ceb);
}
if first_gen & 0x0001 == 1 {
return Ok(&self.snapshot_ceb);
}
let mut retries = 1_000_000;
while retries > 0 {
let snapshot = unsafe { self.ceb_shm.read_volatile() };
let second_gen = generation.load(atomic::Ordering::Acquire);
if first_gen == second_gen {
self.snapshot_gen = first_gen;
self.snapshot_ceb = snapshot;
return Ok(&self.snapshot_ceb);
}
if second_gen & 0x0001 == 0 {
first_gen = second_gen;
}
retries -= 1;
}
Err(ShmError::SegmentNotInitialized(String::from(
"Failed to read the SHM segment after all attempts",
)))
}
}
#[cfg(test)]
mod t_reader {
use super::*;
use crate::shm::ClockStatus;
use byteorder::{NativeEndian, WriteBytesExt};
use nix::sys::time::TimeSpec;
use std::ffi::CString;
use std::fs::OpenOptions;
use std::io::Seek;
use std::io::Write;
use tempfile::NamedTempFile;
macro_rules! write_memory_segment {
($file:ident,
$magic_0:literal,
$magic_1:literal,
$segsize:literal,
$version:literal,
$generation:literal,
($as_of_sec:literal, $as_of_nsec:literal),
($void_after_sec:literal, $void_after_nsec:literal),
$bound_nsec:literal,
$max_drift: literal) => {
let ceb = ClockErrorBoundGeneric::builder()
.as_of(TimeSpec::new($as_of_sec, $as_of_nsec))
.void_after(TimeSpec::new($void_after_sec, $void_after_nsec))
.bound_nsec($bound_nsec)
.disruption_marker(0)
.max_drift_ppb($max_drift)
.clock_status(ClockStatus::Synchronized)
.clock_disruption_support_enabled(true)
.build(ClockErrorBoundLayoutVersion::V2);
let slice = unsafe {
::core::slice::from_raw_parts(
(&ceb as *const ClockErrorBound) as *const u8,
::core::mem::size_of::<ClockErrorBound>(),
)
};
$file
.write_u32::<NativeEndian>($magic_0)
.expect("Write failed magic_0");
$file
.write_u32::<NativeEndian>($magic_1)
.expect("Write failed magic_1");
$file
.write_u32::<NativeEndian>($segsize)
.expect("Write failed segsize");
$file
.write_u16::<NativeEndian>($version)
.expect("Write failed version");
$file
.write_u16::<NativeEndian>($generation)
.expect("Write failed generation");
$file
.write_all(slice)
.expect("Write failed ClockErrorBound");
$file.sync_all().expect("Sync to disk failed");
};
}
#[test]
fn test_reader_new_shm_v2() {
let clockbound_shm_tempfile = NamedTempFile::new().expect("create clockbound file failed");
let clockbound_shm_temppath = clockbound_shm_tempfile.into_temp_path();
let clockbound_shm_path = clockbound_shm_temppath.to_str().unwrap();
let mut clockbound_shm_file = OpenOptions::new()
.write(true)
.open(clockbound_shm_path)
.expect("open clockbound file failed");
write_memory_segment!(
clockbound_shm_file,
0x414D5A4E,
0x43420200,
400,
0x0002,
10,
(0, 0),
(0, 0),
123,
0
);
let path = CString::new(clockbound_shm_path).expect("CString failed");
let res = ShmReader::new(&path);
assert!(res.is_err());
}
#[test]
fn test_reader_new_shm_v3() {
let clockbound_shm_tempfile = NamedTempFile::new().expect("create clockbound file failed");
let clockbound_shm_temppath = clockbound_shm_tempfile.into_temp_path();
let clockbound_shm_path = clockbound_shm_temppath.to_str().unwrap();
let mut clockbound_shm_file = OpenOptions::new()
.write(true)
.open(clockbound_shm_path)
.expect("open clockbound file failed");
write_memory_segment!(
clockbound_shm_file,
0x414D5A4E,
0x43420200,
400,
0x0303,
10,
(0, 0),
(0, 0),
123,
0
);
let path = CString::new(clockbound_shm_path).expect("CString failed");
let reader = ShmReader::new(&path).expect("Failed to create ShmReader");
let version = unsafe { &*reader.version };
let generation = unsafe { &*reader.generation };
let ceb = unsafe { *reader.ceb_shm };
assert_eq!(version.load(atomic::Ordering::Relaxed), 0x0303);
assert_eq!(generation.load(atomic::Ordering::Relaxed), 10);
assert_eq!(ceb.bound_nsec(), 123);
}
#[test]
fn test_reader_new_of_unsupported_shm_version() {
let clockbound_shm_tempfile = NamedTempFile::new().expect("create clockbound file failed");
let clockbound_shm_temppath = clockbound_shm_tempfile.into_temp_path();
let clockbound_shm_path = clockbound_shm_temppath.to_str().unwrap();
let mut clockbound_shm_file = OpenOptions::new()
.write(true)
.open(clockbound_shm_path)
.expect("open clockbound file failed");
write_memory_segment!(
clockbound_shm_file,
0x414D5A4E,
0x43420200,
400,
9999,
10,
(0, 0),
(0, 0),
123,
0
);
let path = CString::new(clockbound_shm_path).expect("CString failed");
let result = ShmReader::new(&path);
assert!(matches!(
result.unwrap_err(),
ShmError::SegmentVersionNotSupported(_)
));
}
#[test]
fn test_reader_snapshot_of_unsupported_shm_version() {
let clockbound_shm_tempfile = NamedTempFile::new().expect("create clockbound file failed");
let clockbound_shm_temppath = clockbound_shm_tempfile.into_temp_path();
let clockbound_shm_path = clockbound_shm_temppath.to_str().unwrap();
let mut clockbound_shm_file = OpenOptions::new()
.write(true)
.open(clockbound_shm_path)
.expect("open clockbound file failed");
write_memory_segment!(
clockbound_shm_file,
0x414D5A4E,
0x43420200,
400,
0x0303,
10,
(0, 0),
(0, 0),
123,
0
);
let path = CString::new(clockbound_shm_path).expect("CString failed");
let mut reader = ShmReader::new(&path).expect("Failed to create ShmReader");
let version = unsafe { &*reader.version };
assert_eq!(version.load(atomic::Ordering::Relaxed), 0x0303);
let result = reader.snapshot();
assert!(result.is_ok());
let _ = clockbound_shm_file.rewind();
write_memory_segment!(
clockbound_shm_file,
0x414D5A4E,
0x43420200,
400,
9999,
10,
(0, 0),
(0, 0),
123,
0
);
let version = unsafe { &*reader.version };
assert_eq!(version.load(atomic::Ordering::Relaxed), 9999);
let result = reader.snapshot();
assert!(
matches!(result.unwrap_err(), ShmError::SegmentVersionNotSupported(msg) if msg.starts_with("Clockbound shared memory segment supports versions"))
);
}
}