Created
July 23, 2026 01:04
-
-
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
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| namespace telemetry; | |
| table SensorReading { | |
| sensor_id: uint64; | |
| timestamp: int64; | |
| temperature: double; | |
| humidity: double; | |
| firmware_version: string; | |
| raw_payload: [ubyte]; | |
| } | |
| root_type SensorReading; |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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); | |
| } | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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], | |
| } | |
| } | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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(()) | |
| } | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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); | |
| } | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 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