use core::{
ptr,
sync::atomic::{AtomicU32, Ordering},
};
use crate::{consts::BUF_SIZE, MODE_BLOCK_IF_FULL, MODE_MASK};
#[repr(C)]
pub(crate) struct Channel {
pub name: *const u8,
pub buffer: *mut u8,
pub size: u32,
pub write: AtomicU32,
pub read: AtomicU32,
pub flags: AtomicU32,
}
#[cfg(not(feature = "drop-on-contention"))]
impl Channel {
pub fn write_all(&self, mut bytes: &[u8]) {
let write = match self.host_is_connected() {
_ if cfg!(feature = "disable-blocking-mode") => Self::nonblocking_write,
true => Self::blocking_write,
false => Self::nonblocking_write,
};
while !bytes.is_empty() {
let consumed = write(self, bytes);
if consumed != 0 {
bytes = &bytes[consumed..];
}
}
}
fn blocking_write(&self, bytes: &[u8]) -> usize {
if bytes.is_empty() {
return 0;
}
let read = self.read.load(Ordering::Relaxed) as usize;
let write = self.write.load(Ordering::Acquire) as usize;
let available = available_buffer_size(read, write);
if available == 0 {
return 0;
}
self.write_impl(bytes, write, available)
}
fn nonblocking_write(&self, bytes: &[u8]) -> usize {
let write = self.write.load(Ordering::Acquire) as usize;
self.write_impl(bytes, write, BUF_SIZE)
}
fn write_impl(&self, bytes: &[u8], cursor: usize, available: usize) -> usize {
let len = bytes.len().min(available);
unsafe {
self.copy_wrapping(&bytes[..len], cursor);
}
self.write.store(
(cursor.wrapping_add(len) % BUF_SIZE) as u32,
Ordering::Release,
);
len
}
}
#[cfg(feature = "drop-on-contention")]
impl Channel {
pub fn stage_bytes(&self, cursor: &mut usize, bytes: &[u8]) -> bool {
if bytes.is_empty() || bytes.len() >= BUF_SIZE {
return bytes.is_empty();
}
if cfg!(feature = "disable-blocking-mode") || !self.host_is_connected() {
let read = self.read.load(Ordering::Relaxed) as usize;
if available_buffer_size(read, *cursor) < bytes.len() {
return false;
}
} else {
while {
let read = self.read.load(Ordering::Relaxed) as usize;
available_buffer_size(read, *cursor) < bytes.len()
} {}
}
unsafe {
self.copy_wrapping(bytes, *cursor);
}
*cursor = cursor.wrapping_add(bytes.len()) % BUF_SIZE;
true
}
pub fn commit(&self, cursor: usize) {
self.write.store(cursor as u32, Ordering::Release);
}
}
impl Channel {
unsafe fn copy_wrapping(&self, bytes: &[u8], cursor: usize) {
if cursor + bytes.len() > BUF_SIZE {
let pivot = BUF_SIZE - cursor;
ptr::copy_nonoverlapping(bytes.as_ptr(), self.buffer.add(cursor), pivot);
ptr::copy_nonoverlapping(bytes.as_ptr().add(pivot), self.buffer, bytes.len() - pivot);
} else {
ptr::copy_nonoverlapping(bytes.as_ptr(), self.buffer.add(cursor), bytes.len());
}
}
pub fn flush(&self) {
if !self.host_is_connected() {
return;
}
let read = || self.read.load(Ordering::Relaxed);
let write = || self.write.load(Ordering::Relaxed);
while read() != write() {}
}
fn host_is_connected(&self) -> bool {
self.flags.load(Ordering::Relaxed) & MODE_MASK == MODE_BLOCK_IF_FULL
}
}
pub(crate) fn available_buffer_size(read_cursor: usize, write_cursor: usize) -> usize {
if read_cursor > write_cursor {
read_cursor - write_cursor - 1
} else {
BUF_SIZE - write_cursor - 1 + read_cursor
}
}
#[cfg(test)]
mod tests {
use super::available_buffer_size;
use crate::consts::BUF_SIZE;
#[test]
fn test_rtt_available_buffer_size() {
let avail = |read: usize, write: usize| available_buffer_size(read, write);
assert_eq!(avail(0, 0), BUF_SIZE - 1);
assert_eq!(avail(10, 10), BUF_SIZE - 1);
assert_eq!(avail(BUF_SIZE - 1, BUF_SIZE - 1), BUF_SIZE - 1);
assert_eq!(avail(0, BUF_SIZE - 1), 0); assert_eq!(avail(5, 4), 0); assert_eq!(avail(1, 0), 0);
assert_eq!(avail(10, 5), 10 - 5 - 1);
assert_eq!(avail(BUF_SIZE - 1, 0), (BUF_SIZE - 1) - 0 - 1);
assert_eq!(avail(5, 10), BUF_SIZE - 10 - 1 + 5);
assert_eq!(avail(0, 1), BUF_SIZE - 1 - 1 + 0); assert_eq!(avail(1, BUF_SIZE - 1), BUF_SIZE - (BUF_SIZE - 1) - 1 + 1);
assert_eq!(avail(1, BUF_SIZE - 1), 1); assert_eq!(avail(2, BUF_SIZE - 1), 2);
let data_in_buffer = |read: usize, write: usize| (write + BUF_SIZE - read) % BUF_SIZE;
let free_should_be = |read: usize, write: usize| BUF_SIZE - 1 - data_in_buffer(read, write);
for read in 0..BUF_SIZE.min(64) {
for write in 0..BUF_SIZE.min(64) {
let expected = free_should_be(read, write);
let actual = avail(read, write);
assert_eq!(actual, expected, "Mismatch at read={read}, write={write}");
}
}
}
}
#[cfg(all(test, feature = "drop-on-contention"))]
mod test_drop_on_contention {
use super::available_buffer_size;
use crate::consts::BUF_SIZE;
use super::Channel;
use crate::MODE_NON_BLOCKING_TRIM;
use core::{
ptr,
sync::atomic::{AtomicU32, Ordering},
};
#[test]
fn staged_bytes_stay_hidden_until_commit() {
let mut buffer = [0u8; BUF_SIZE];
let channel = Channel {
name: ptr::null(),
buffer: buffer.as_mut_ptr(),
size: BUF_SIZE as u32,
write: AtomicU32::new(0),
read: AtomicU32::new(0),
flags: AtomicU32::new(MODE_NON_BLOCKING_TRIM),
};
let mut cursor = 0;
assert!(channel.stage_bytes(&mut cursor, b"abc"));
assert_eq!(cursor, 3);
assert_eq!(channel.write.load(Ordering::Relaxed), 0);
assert_eq!(&buffer[..3], b"abc");
channel.commit(cursor);
assert_eq!(channel.write.load(Ordering::Relaxed), 3);
}
#[test]
fn staged_bytes_wrap_without_publishing() {
let mut buffer = [0u8; BUF_SIZE];
let channel = Channel {
name: ptr::null(),
buffer: buffer.as_mut_ptr(),
size: BUF_SIZE as u32,
write: AtomicU32::new(0),
read: AtomicU32::new(8),
flags: AtomicU32::new(MODE_NON_BLOCKING_TRIM),
};
let mut cursor = BUF_SIZE - 2;
assert!(channel.stage_bytes(&mut cursor, b"wxyz"));
assert_eq!(cursor, 2);
assert_eq!(channel.write.load(Ordering::Relaxed), 0);
assert_eq!(&buffer[BUF_SIZE - 2..], b"wx");
assert_eq!(&buffer[..2], b"yz");
}
#[test]
fn staged_bytes_rejects_oversized_frame_without_side_effects() {
let mut buffer = [0xAA; BUF_SIZE];
let channel = Channel {
name: ptr::null(),
buffer: buffer.as_mut_ptr(),
size: BUF_SIZE as u32,
write: AtomicU32::new(7),
read: AtomicU32::new(7),
flags: AtomicU32::new(MODE_NON_BLOCKING_TRIM),
};
let mut cursor = 7;
let bytes = [0x55; BUF_SIZE];
assert!(!channel.stage_bytes(&mut cursor, &bytes));
assert_eq!(cursor, 7);
assert_eq!(channel.write.load(Ordering::Relaxed), 7);
assert!(buffer.iter().all(|&b| b == 0xAA));
}
#[test]
fn staged_bytes_rejects_when_space_is_insufficient_without_partial_copy() {
let mut buffer = [0xAA; BUF_SIZE];
let channel = Channel {
name: ptr::null(),
buffer: buffer.as_mut_ptr(),
size: BUF_SIZE as u32,
write: AtomicU32::new(0),
read: AtomicU32::new(4),
flags: AtomicU32::new(MODE_NON_BLOCKING_TRIM),
};
let mut cursor = 0;
assert_eq!(available_buffer_size(4, cursor), 3);
assert!(!channel.stage_bytes(&mut cursor, b"wxyz"));
assert_eq!(cursor, 0);
assert_eq!(channel.write.load(Ordering::Relaxed), 0);
assert!(buffer.iter().all(|&b| b == 0xAA));
}
}