use super::message::RYNK_HEADER_SIZE;
pub struct Deframer {
consumed: usize,
filled: usize,
discarding: bool,
}
impl Deframer {
pub const fn new() -> Self {
Self {
consumed: 0,
filled: 0,
discarding: false,
}
}
pub fn tail<'b>(&mut self, buf: &'b mut [u8]) -> &'b mut [u8] {
if self.consumed > 0 {
buf.copy_within(self.consumed..self.filled, 0);
self.filled -= self.consumed;
self.consumed = 0;
}
&mut buf[self.filled..]
}
pub fn commit(&mut self, n: usize) {
self.filled += n;
}
pub fn has_pending(&self) -> bool {
!self.discarding && self.consumed < self.filled
}
pub fn park_pending(&mut self, buf: &mut [u8]) -> usize {
let len = self.filled - self.consumed;
buf.copy_within(self.consumed..self.filled, buf.len() - len);
self.consumed = 0;
self.filled = 0;
len
}
pub fn unpark_pending(&mut self, buf: &mut [u8], parked: usize) {
buf.copy_within(buf.len() - parked.., 0);
self.consumed = 0;
self.filled = parked;
}
pub fn next(&mut self, buf: &mut [u8]) -> Option<usize> {
loop {
let Some(i) = buf[self.consumed..self.filled].iter().position(|&b| b == 0) else {
if self.discarding || (self.consumed == 0 && self.filled == buf.len()) {
self.discarding = true;
self.consumed = self.filled;
}
return None;
};
let (start, delim) = (self.consumed, self.consumed + i);
self.consumed = delim + 1;
if self.discarding {
self.discarding = false;
continue;
}
match cobs::decode_in_place(&mut buf[start..delim]) {
Ok(len) if len >= RYNK_HEADER_SIZE => {
if start > 0 {
buf.copy_within(start..start + len, 0);
}
return Some(len);
}
_ => {}
}
}
}
}
impl Default for Deframer {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
extern crate alloc;
use super::*;
use crate::protocol::rynk::{Cmd, RynkHeader, encode_frame};
const CMD: Cmd = Cmd::from_raw(0x0102);
fn encode(buf: &mut [u8], cmd: Cmd, seq: u8, payload: &[u8]) -> usize {
encode_frame(buf, RynkHeader { cmd, seq }, &payload).unwrap()
}
fn feed(df: &mut Deframer, buf: &mut [u8], cmd: Cmd, seq: u8, payload: &[u8]) {
let tail = df.tail(buf);
let n = encode(tail, cmd, seq, payload);
df.commit(n);
}
fn assert_frame(frame: &[u8], cmd: Cmd, seq: u8, payload: &[u8]) {
assert_eq!(Cmd::from_le_bytes([frame[0], frame[1]]), cmd);
assert_eq!(frame[2], seq);
let decoded: heapless::Vec<u8, 32> = postcard::from_bytes(&frame[RYNK_HEADER_SIZE..]).unwrap();
assert_eq!(&decoded[..], payload);
}
#[test]
fn round_trip() {
let mut buf = [0u8; 64];
let mut df = Deframer::new();
feed(&mut df, &mut buf, CMD, 0x42, &[1, 2, 3]);
let len = df.next(&mut buf).expect("one frame");
assert_frame(&buf[..len], CMD, 0x42, &[1, 2, 3]);
assert!(df.next(&mut buf).is_none(), "no second frame");
}
#[test]
fn reassembles_byte_by_byte() {
let mut src = [0u8; 64];
let n = encode(&mut src, CMD, 7, &[9, 8, 7, 6]);
let mut buf = [0u8; 64];
let mut df = Deframer::new();
for i in 0..n {
df.tail(&mut buf)[0] = src[i];
df.commit(1);
if i + 1 < n {
assert!(df.next(&mut buf).is_none(), "no frame before the delimiter");
}
}
let len = df.next(&mut buf).expect("frame after the last byte");
assert_frame(&buf[..len], CMD, 7, &[9, 8, 7, 6]);
}
#[test]
fn two_pipelined_frames_in_one_buffer() {
let mut buf = [0u8; 128];
let mut df = Deframer::new();
feed(&mut df, &mut buf, CMD, 1, &[1]);
feed(&mut df, &mut buf, CMD, 2, &[2, 2]);
let len = df.next(&mut buf).expect("first frame");
assert_frame(&buf[..len], CMD, 1, &[1]);
let len = df.next(&mut buf).expect("second frame");
assert_frame(&buf[..len], CMD, 2, &[2, 2]);
assert!(df.next(&mut buf).is_none());
}
#[test]
fn frames_shorter_than_a_header_are_skipped() {
let mut buf = [0u8; 64];
let mut df = Deframer::new();
let short = [
0x00, 0x02, 0xAA, 0x00, 0x03, 0xAA, 0xBB, 0x00, ];
df.tail(&mut buf)[..short.len()].copy_from_slice(&short);
df.commit(short.len());
assert!(df.next(&mut buf).is_none(), "sub-header frames must be skipped");
let header_only = [0x04, 0x02, 0x01, 0x42, 0x00];
df.tail(&mut buf)[..header_only.len()].copy_from_slice(&header_only);
df.commit(header_only.len());
let len = df.next(&mut buf).expect("header-only frame is valid");
assert_eq!(len, RYNK_HEADER_SIZE);
assert_eq!(Cmd::from_le_bytes([buf[0], buf[1]]), CMD);
assert_eq!(buf[2], 0x42);
}
#[test]
fn park_and_unpark_preserve_a_pipelined_frame() {
let mut buf = [0u8; 128];
let mut df = Deframer::new();
feed(&mut df, &mut buf, CMD, 1, &[1]);
feed(&mut df, &mut buf, CMD, 2, &[2, 2]);
let len = df.next(&mut buf).expect("first frame");
assert_frame(&buf[..len], CMD, 1, &[1]);
let parked = df.park_pending(&mut buf);
assert!(parked > 0);
let window = buf.len() - parked;
buf[..window].iter_mut().for_each(|b| *b = 0xEE);
df.unpark_pending(&mut buf, parked);
let len = df.next(&mut buf).expect("second frame survives the reply");
assert_frame(&buf[..len], CMD, 2, &[2, 2]);
assert!(df.next(&mut buf).is_none());
}
#[test]
fn park_preserves_a_partial_frame() {
let mut src = [0u8; 64];
let n = encode(&mut src, CMD, 9, &[7, 7, 7]);
let mut buf = [0u8; 64];
let mut df = Deframer::new();
feed(&mut df, &mut buf, CMD, 8, &[5]);
let split = n / 2;
df.tail(&mut buf)[..split].copy_from_slice(&src[..split]);
df.commit(split);
let len = df.next(&mut buf).expect("complete frame");
assert_frame(&buf[..len], CMD, 8, &[5]);
let parked = df.park_pending(&mut buf);
assert!(parked > 0);
let window = buf.len() - parked;
buf[..window].iter_mut().for_each(|b| *b = 0xEE);
df.unpark_pending(&mut buf, parked);
df.tail(&mut buf)[..n - split].copy_from_slice(&src[split..n]);
df.commit(n - split);
let len = df.next(&mut buf).expect("split frame completes after the park");
assert_frame(&buf[..len], CMD, 9, &[7, 7, 7]);
}
#[test]
fn park_preserves_the_overflow_drain() {
let mut buf = [0u8; 32];
let mut df = Deframer::new();
df.tail(&mut buf).iter_mut().for_each(|b| *b = 0xFF);
df.commit(32);
assert!(df.next(&mut buf).is_none(), "overflow: draining");
let parked = df.park_pending(&mut buf);
assert_eq!(parked, 0, "drained bytes are dead, nothing to park");
df.unpark_pending(&mut buf, parked);
df.tail(&mut buf)[..2].copy_from_slice(&[0xFF, 0x00]);
df.commit(2);
feed(&mut df, &mut buf, CMD, 3, &[9]);
let len = df.next(&mut buf).expect("frame after the drain clears");
assert_frame(&buf[..len], CMD, 3, &[9]);
}
#[test]
fn resyncs_past_injected_garbage() {
let mut buf = [0u8; 128];
let mut df = Deframer::new();
feed(&mut df, &mut buf, CMD, 10, &[0xAA, 0xBB]);
let tail = df.tail(&mut buf);
tail[..4].copy_from_slice(&[0xDE, 0xAD, 0xBE, 0x00]);
df.commit(4);
feed(&mut df, &mut buf, CMD, 11, &[0xCC]);
let len = df.next(&mut buf).expect("frame A");
assert_frame(&buf[..len], CMD, 10, &[0xAA, 0xBB]);
let len = df.next(&mut buf).expect("frame B survives the garbage");
assert_frame(&buf[..len], CMD, 11, &[0xCC]);
}
#[test]
fn resyncs_after_buffer_overflow() {
let mut buf = [0u8; 32];
let mut df = Deframer::new();
df.tail(&mut buf).iter_mut().for_each(|b| *b = 0xFF);
df.commit(32);
assert!(df.next(&mut buf).is_none(), "overflow yields no frame");
df.tail(&mut buf)[0] = 0x00;
df.commit(1);
feed(&mut df, &mut buf, CMD, 12, &[5, 6]);
let len = df.next(&mut buf).expect("frame after overflow resync");
assert_frame(&buf[..len], CMD, 12, &[5, 6]);
}
#[test]
fn has_pending_tracks_partial_frames_not_garbage() {
let mut buf = [0u8; 32];
let mut df = Deframer::new();
assert!(!df.has_pending(), "empty buffer: nothing pending");
let mut src = [0u8; 32];
let n = encode(&mut src, CMD, 1, &[1, 2, 3]);
df.tail(&mut buf)[..n - 1].copy_from_slice(&src[..n - 1]);
df.commit(n - 1);
assert!(df.next(&mut buf).is_none());
assert!(df.has_pending(), "half-received frame is pending");
df.tail(&mut buf)[0] = 0x00;
df.commit(1);
let len = df.next(&mut buf).expect("frame completes");
assert_frame(&buf[..len], CMD, 1, &[1, 2, 3]);
assert!(df.next(&mut buf).is_none());
assert!(!df.has_pending(), "consumed frame is not pending");
df.tail(&mut buf).iter_mut().for_each(|b| *b = 0xFF);
df.commit(32);
assert!(df.next(&mut buf).is_none());
assert!(!df.has_pending(), "overflow drain is not pending");
}
#[test]
fn discard_drain_preserves_bytes_committed_after_the_scan() {
let mut buf = [0u8; 32];
let mut df = Deframer::new();
df.tail(&mut buf).iter_mut().for_each(|b| *b = 0xFF);
df.commit(32);
assert!(df.next(&mut buf).is_none(), "overflow: enter discard mode");
let mut src = [0u8; 32];
let n = encode(&mut src, CMD, 3, &[7]);
{
let tail = df.tail(&mut buf);
assert_eq!(tail.len(), 32, "drained garbage is reclaimed by tail()");
tail[0] = 0x00;
tail[1..=n].copy_from_slice(&src[..n]);
}
df.commit(1 + n);
let len = df.next(&mut buf).expect("frame after the drain delimiter");
assert_frame(&buf[..len], CMD, 3, &[7]);
assert!(df.next(&mut buf).is_none());
}
#[test]
fn compacts_instead_of_dropping_a_buffer_filling_frame() {
let mut src = [0u8; 64];
let n = encode(&mut src, CMD, 9, &[1, 2, 3, 4, 5]);
let cap = 1 + (n - 1);
let mut store = [0u8; 64];
let mut df = Deframer::new();
{
let tail = df.tail(&mut store[..cap]);
tail[0] = 0x00; tail[1..].copy_from_slice(&src[..n - 1]); df.commit(cap);
}
assert!(
df.next(&mut store[..cap]).is_none(),
"buffer full with a consumed prefix must compact, not overflow"
);
{
let tail = df.tail(&mut store[..cap]);
assert!(!tail.is_empty(), "compaction must free the skipped byte");
tail[0] = 0x00; df.commit(1);
}
let len = df.next(&mut store[..cap]).expect("frame completes after compaction");
assert_frame(&store[..len], CMD, 9, &[1, 2, 3, 4, 5]);
}
struct Rng(u64);
impl Rng {
fn next_u64(&mut self) -> u64 {
self.0 ^= self.0 >> 12;
self.0 ^= self.0 << 25;
self.0 ^= self.0 >> 27;
self.0.wrapping_mul(0x2545_F491_4F6C_DD1D)
}
fn below(&mut self, n: usize) -> usize {
(self.next_u64() % n as u64) as usize
}
}
#[test]
fn fuzz_recovers_every_wellformed_frame() {
use alloc::vec;
use alloc::vec::Vec;
for seed in 1..=64u64 {
for &cap in &[16usize, 24, 48, 96] {
let mut rng = Rng(seed.wrapping_mul(0x9E37_79B9_7F4A_7C15) | 1);
let mut stream: Vec<u8> = Vec::new();
let mut expected: Vec<Vec<u8>> = Vec::new();
for seq in 0..30u8 {
match rng.below(4) {
0 => stream.push(0x00),
1 => {
let glen = 1 + rng.below(cap - 2);
stream.push(0xFF);
for _ in 0..glen {
stream.push(1 + rng.below(255) as u8);
}
stream.push(0x00);
}
2 => {
for _ in 0..cap + 1 + rng.below(cap) {
stream.push(0xFF);
}
stream.push(0x00);
}
_ => {
let mut payload = [0u8; 96];
let plen = rng.below(cap);
payload[..plen].iter_mut().for_each(|b| *b = rng.next_u64() as u8);
let mut tmp = [0u8; 256];
let n = encode(&mut tmp, CMD, seq, &payload[..plen]);
if n <= cap {
stream.extend_from_slice(&tmp[..n]);
let mut logical = tmp[..n - 1].to_vec();
let llen = cobs::decode_in_place(&mut logical).unwrap();
logical.truncate(llen);
expected.push(logical);
} else {
stream.push(0x00);
}
}
}
}
let mut buf = vec![0u8; cap];
let mut df = Deframer::new();
let mut got: Vec<Vec<u8>> = Vec::new();
let mut pos = 0;
while pos < stream.len() {
let tail = df.tail(&mut buf);
assert!(!tail.is_empty(), "reader must never be offered an empty tail");
let n = tail.len().min(1 + rng.below(7)).min(stream.len() - pos);
tail[..n].copy_from_slice(&stream[pos..pos + n]);
pos += n;
df.commit(n);
while let Some(len) = df.next(&mut buf) {
got.push(buf[..len].to_vec());
}
}
assert_eq!(got, expected, "seed {seed} cap {cap}");
}
}
}
}