#![warn(clippy::undocumented_unsafe_blocks)]
use libc::{c_void, munmap};
use std::marker::PhantomData;
use std::ops::{Deref, DerefMut};
use std::sync::atomic::AtomicU8;
mod alloc;
thread_local! {
pub(crate) static SIMULATE_FAIL: std::cell::Cell<bool> = const { std::cell::Cell::new(false) };
}
pub const HUGEPAGE_STATE_STANDARD: u8 = 0;
pub const HUGEPAGE_STATE_THP: u8 = 1;
pub const HUGEPAGE_STATE_HUGETLB: u8 = 2;
static MIRROR_BUF_HUGEPAGE_STATE: AtomicU8 = AtomicU8::new(HUGEPAGE_STATE_STANDARD);
pub fn sync_huge_page_flag(rt_status: &crate::common::spsc::RtStatusFlags) {
let state = MIRROR_BUF_HUGEPAGE_STATE.load(std::sync::atomic::Ordering::Relaxed);
match state {
HUGEPAGE_STATE_HUGETLB => {
rt_status.set_flag(crate::common::spsc::RT_STATUS_HUGEPAGE_OK);
}
HUGEPAGE_STATE_THP => {
rt_status.set_flag(crate::common::spsc::RT_STATUS_THP_ACTIVE);
}
_ => {}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MirrorHugePageStatus {
Explicit2MB,
Transparent,
Standard,
}
pub fn is_huge_page_active() -> bool {
let state = MIRROR_BUF_HUGEPAGE_STATE.load(std::sync::atomic::Ordering::Relaxed);
state == HUGEPAGE_STATE_THP || state == HUGEPAGE_STATE_HUGETLB
}
pub fn huge_page_status() -> MirrorHugePageStatus {
match MIRROR_BUF_HUGEPAGE_STATE.load(std::sync::atomic::Ordering::Relaxed) {
HUGEPAGE_STATE_HUGETLB => MirrorHugePageStatus::Explicit2MB,
HUGEPAGE_STATE_THP => MirrorHugePageStatus::Transparent,
_ => MirrorHugePageStatus::Standard,
}
}
pub struct MirroredBuffer<T> {
ptr: *mut T,
size_elements: usize,
_marker: PhantomData<T>,
}
pub fn set_simulate_fail(fail: bool) {
SIMULATE_FAIL.with(|f| f.set(fail));
}
impl<T> MirroredBuffer<T> {
pub fn size(&self) -> usize {
self.size_elements
}
}
impl<T> std::fmt::Debug for MirroredBuffer<T> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("MirroredBuffer")
.field("ptr", &self.ptr)
.field("size_elements", &self.size_elements)
.field("capacity_virtual", &(self.size_elements * 2))
.finish()
}
}
impl<T> Deref for MirroredBuffer<T> {
type Target = [T];
#[inline(always)]
fn deref(&self) -> &Self::Target {
unsafe { std::slice::from_raw_parts(self.ptr, self.size_elements * 2) }
}
}
impl<T> DerefMut for MirroredBuffer<T> {
#[inline(always)]
fn deref_mut(&mut self) -> &mut Self::Target {
unsafe { std::slice::from_raw_parts_mut(self.ptr, self.size_elements * 2) }
}
}
impl<T> Drop for MirroredBuffer<T> {
fn drop(&mut self) {
let element_size = std::mem::size_of::<T>();
let size_bytes = self.size_elements * element_size;
unsafe {
munmap(self.ptr as *mut c_void, size_bytes * 2);
}
}
}
impl<T: Clone> MirroredBuffer<T> {
pub fn try_clone(&self) -> std::io::Result<Self> {
let mut new_buf = Self::new(self.size_elements)?;
new_buf[..self.size_elements].clone_from_slice(&self[..self.size_elements]);
Ok(new_buf)
}
}
impl<T: Clone> Clone for MirroredBuffer<T> {
#[cold]
fn clone(&self) -> Self {
self.try_clone()
.expect("MirroredBuffer::clone: allocation failed (use try_clone for fallible path)")
}
}
unsafe impl<T: Send> Send for MirroredBuffer<T> {}
unsafe impl<T: Sync> Sync for MirroredBuffer<T> {}
#[cfg(all(test, target_os = "linux"))]
#[path = "mirror_buf_test.rs"]
mod mirror_buf_test;