Skip to content

Instantly share code, notes, and snippets.

@mohashari
Created July 23, 2026 01:04
Show Gist options
  • Select an option

  • Save mohashari/ae3e8f73a393f8196918536f7dbbc2f5 to your computer and use it in GitHub Desktop.

Select an option

Save mohashari/ae3e8f73a393f8196918536f7dbbc2f5 to your computer and use it in GitHub Desktop.
Implementing a Zero-Copy Serializer in Rust Using FlatBuffers and Shared Memory for IPC — code snippets
namespace telemetry;
table SensorReading {
sensor_id: uint64;
timestamp: int64;
temperature: double;
humidity: double;
firmware_version: string;
raw_payload: [ubyte];
}
root_type SensorReading;
use std::fs::File;
use std::os::unix::io::FromRawFd;
use nix::sys::mman::{shm_open, shm_unlink, OMode, ShmOFlags};
use nix::sys::stat::Mode;
use memmap2::MmapMut;
pub struct SharedMemory {
name: String,
mmap: MmapMut,
}
impl SharedMemory {
pub fn new(name: &str, size: usize, create: bool) -> Result<Self, Box<dyn std::error::Error>> {
let flags = if create {
ShmOFlags::O_CREAT | ShmOFlags::O_RDWR | ShmOFlags::O_TRUNC
} else {
ShmOFlags::O_RDWR
};
let mode = Mode::S_IRUSR | Mode::S_IWUSR;
// Open the POSIX shared memory file descriptor
let fd = shm_open(name, flags, mode)?;
// Set the size of the shared memory object
if create {
unsafe {
let res = libc::ftruncate(fd, size as libc::off_t);
if res != 0 {
return Err(std::io::Error::last_os_error().into());
}
}
}
// Convert raw file descriptor to standard File for memmap2
let file = unsafe { File::from_raw_fd(fd) };
let mmap = unsafe { MmapMut::map_mut(&file)? };
Ok(Self {
name: name.to_string(),
mmap,
})
}
pub fn as_mut_slice(&mut self) -> &mut [u8] {
&mut self.mmap
}
pub fn as_slice(&self) -> &[u8] {
&self.mmap
}
}
impl Drop for SharedMemory {
fn drop(&mut self) {
// Attempt to unlink the SHM file. In production, only the owner
// process should call shm_unlink, but here we attempt cleanup gracefully.
let _ = shm_unlink(&self.name);
}
}
use std::sync::atomic::{AtomicU64, Ordering};
#[repr(C)]
#[repr(align(64))]
pub struct RingBufferHeader {
// Padded to align to CPU cache line size (64 bytes) to prevent false sharing
pub write_offset: AtomicU64,
_pad1: [u8; 56],
pub read_offset: AtomicU64,
_pad2: [u8; 56],
pub capacity: u64,
}
pub struct RingBufferWriter<'a> {
header: &'a mut RingBufferHeader,
data: &'a mut [u8],
}
impl<'a> RingBufferWriter<'a> {
pub fn new(shm_slice: &'a mut [u8], capacity: u64) -> Self {
let header_size = std::mem::size_of::<RingBufferHeader>();
assert!(shm_slice.len() >= header_size + capacity as usize);
let (header_bytes, data_bytes) = shm_slice.split_at_mut(header_size);
// Safely cast raw slice to our aligned header representation
let header = unsafe { &mut *(header_bytes.as_mut_ptr() as *mut RingBufferHeader) };
header.capacity = capacity;
Self {
header,
data: &mut data_bytes[..capacity as usize],
}
}
}
use flatbuffers::FlatBufferBuilder;
use std::sync::atomic::Ordering;
impl<'a> RingBufferWriter<'a> {
pub fn write_message(&mut self, serialized: &[u8]) -> Result<(), &'static str> {
let len = serialized.len();
// 4 bytes to store the length prefix of the message
let required_space = len + 4;
let write_ptr = self.header.write_offset.load(Ordering::Relaxed);
let read_ptr = self.header.read_offset.load(Ordering::Acquire);
// Calculate used and free space in the circular buffer
let used_space = write_ptr.checked_sub(read_ptr).unwrap_or(0);
let free_space = self.header.capacity - used_space;
if free_space < required_space as u64 {
return Err("Ring buffer full (backpressure)");
}
let mut offset = (write_ptr % self.header.capacity) as usize;
let remaining_tail = self.header.capacity as usize - offset;
if remaining_tail < required_space {
// Early wrap: write sentinel (0 length) to tail and write message at buffer start
let zero_sentinel = 0u32.to_le_bytes();
if remaining_tail >= 4 {
self.data[offset..offset+4].copy_from_slice(&zero_sentinel);
}
// Advance write pointer past the padded tail space
let pad_bytes = remaining_tail as u64;
let new_write_ptr = write_ptr + pad_bytes;
// Check if there is enough space at the start of the buffer
let current_read_ptr = self.header.read_offset.load(Ordering::Acquire);
let updated_used = new_write_ptr - current_read_ptr;
if (self.header.capacity - updated_used) < required_space as u64 {
return Err("Ring buffer full after padding");
}
// Write length and payload at the beginning
let len_bytes = (len as u32).to_le_bytes();
self.data[0..4].copy_from_slice(&len_bytes);
self.data[4..4+len].copy_from_slice(serialized);
// Commit write pointer
self.header.write_offset.store(new_write_ptr + required_space as u64, Ordering::Release);
} else {
// Contiguous write fits at current tail
let len_bytes = (len as u32).to_le_bytes();
self.data[offset..offset+4].copy_from_slice(&len_bytes);
self.data[offset+4..offset+required_space].copy_from_slice(serialized);
self.header.write_offset.store(write_ptr + required_space as u64, Ordering::Release);
}
Ok(())
}
}
use std::sync::atomic::Ordering;
use flatbuffers::VerifierOptions;
pub struct RingBufferReader<'a> {
header: &'a RingBufferHeader,
data: &'a [u8],
}
impl<'a> RingBufferReader<'a> {
pub fn new(shm_slice: &'a [u8], capacity: u64) -> Self {
let header_size = std::mem::size_of::<RingBufferHeader>();
let (header_bytes, data_bytes) = shm_slice.split_at(header_size);
let header = unsafe { &*(header_bytes.as_ptr() as *const RingBufferHeader) };
Self {
header,
data: &data_bytes[..capacity as usize],
}
}
pub fn get_next_reading(&self) -> Result<Option<(telemetry::SensorReading<'a>, usize)>, &'static str> {
let write_ptr = self.header.write_offset.load(Ordering::Acquire);
let read_ptr = self.header.read_offset.load(Ordering::Relaxed);
if read_ptr == write_ptr {
return Ok(None); // No new data
}
let mut offset = (read_ptr % self.header.capacity) as usize;
// Read 4-byte length prefix
let msg_len = u32::from_le_bytes(
self.data[offset..offset+4].try_into().map_err(|_| "Failed to parse length")?
) as usize;
let (payload, advanced_bytes) = if msg_len == 0 {
// Wrap sentinel detected; real message is at the beginning of the buffer
let remaining_tail = self.header.capacity as usize - offset;
let real_len = u32::from_le_bytes(
self.data[0..4].try_into().map_err(|_| "Failed to parse length")?
) as usize;
(&self.data[4..4+real_len], remaining_tail + 4 + real_len)
} else {
(&self.data[offset+4..offset+4+msg_len], 4 + msg_len)
};
// Verify FlatBuffer structure safety before casting
let verifier_opts = VerifierOptions::default();
let mut verifier = flatbuffers::Verifier::new(&verifier_opts, payload);
if let Err(_) = flatbuffers::root_with_opts::<telemetry::SensorReading>(&mut verifier) {
return Err("FlatBuffer layout verification failed");
}
let reading = flatbuffers::root::<telemetry::SensorReading>(payload)
.map_err(|_| "Failed to cast FlatBuffer")?;
Ok(Some((reading, advanced_bytes)))
}
pub fn commit_read(&self, bytes_consumed: usize) {
let read_ptr = self.header.read_offset.load(Ordering::Relaxed);
self.header.read_offset.store(read_ptr + bytes_consumed as u64, Ordering::Release);
}
}
use nix::sys::eventfd::{eventfd, EfdFlags};
use std::fs::File;
use std::os::unix::io::{FromRawFd, RawFd};
use std::io::{Read, Write};
pub struct EventfdChannel {
fd: RawFd,
file: File,
}
impl EventfdChannel {
pub fn new() -> Result<Self, std::io::Error> {
let fd = eventfd(0, EfdFlags::EFD_SEMAPHORE | EfdFlags::EFD_CLOEXEC)
.map_err(|e| std::io::Error::from_raw_os_error(e as i32))?;
let file = unsafe { File::from_raw_fd(fd) };
Ok(Self { fd, file })
}
pub fn notify(&mut self) -> Result<(), std::io::Error> {
let buffer = 1u64.to_ne_bytes();
(&self.file).write_all(&buffer)
}
pub fn wait(&mut self) -> Result<(), std::io::Error> {
let mut buffer = [0u8; 8];
match (&self.file).read_exact(&mut buffer) {
Ok(_) => Ok(()),
Err(e) => Err(e),
}
}
pub fn raw_fd(&self) -> RawFd {
self.fd
}
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment