use core::mem::MaybeUninit;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[must_use = "a completed block is the publishable window; dropping it discards N samples"]
pub struct Block<T: Copy, const N: usize> {
first_sequence: u32,
last_sequence: u32,
samples: [T; N],
}
impl<T: Copy, const N: usize> Block<T, N> {
#[must_use]
pub const fn first_sequence(&self) -> u32 {
self.first_sequence
}
#[must_use]
pub const fn last_sequence(&self) -> u32 {
self.last_sequence
}
#[must_use]
pub const fn samples(&self) -> &[T; N] {
&self.samples
}
#[must_use]
pub fn into_samples(self) -> [T; N] {
self.samples
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[must_use = "the rejected sample rides in this error; dropping it unseen loses the sample"]
pub enum FillError<T: Copy> {
ReservedSequence {
sample: T,
},
Discontinuous {
expected: u32,
received: u32,
sample: T,
},
}
pub struct BlockBuilder<T: Copy, const N: usize> {
samples: [MaybeUninit<T>; N],
len: usize,
first_sequence: u32,
last_sequence: u32,
}
impl<T: Copy, const N: usize> BlockBuilder<T, N> {
#[must_use]
pub const fn new() -> Self {
const { assert!(N > 0, "BlockBuilder capacity must be greater than zero") };
Self {
samples: [const { MaybeUninit::uninit() }; N],
len: 0,
first_sequence: 0,
last_sequence: 0,
}
}
pub fn push(&mut self, sequence: u32, sample: T) -> Result<Option<Block<T, N>>, FillError<T>> {
if sequence == 0 {
return Err(FillError::ReservedSequence { sample });
}
if self.len != 0 {
let expected = next_sequence(self.last_sequence);
if sequence != expected {
return Err(FillError::Discontinuous {
expected,
received: sequence,
sample,
});
}
} else {
self.first_sequence = sequence;
}
self.samples[self.len].write(sample);
self.len += 1;
self.last_sequence = sequence;
if self.len != N {
return Ok(None);
}
let samples = unsafe { self.samples.as_ptr().cast::<[T; N]>().read() };
let block = Block {
first_sequence: self.first_sequence,
last_sequence: self.last_sequence,
samples,
};
self.clear();
Ok(Some(block))
}
pub fn clear(&mut self) {
self.len = 0;
self.first_sequence = 0;
self.last_sequence = 0;
}
#[must_use]
pub const fn len(&self) -> usize {
self.len
}
#[must_use]
pub const fn is_empty(&self) -> bool {
self.len == 0
}
#[must_use]
pub const fn capacity(&self) -> usize {
N
}
#[must_use]
pub const fn expected_sequence(&self) -> Option<u32> {
if self.len == 0 {
None
} else {
Some(next_sequence(self.last_sequence))
}
}
}
impl<T: Copy, const N: usize> Default for BlockBuilder<T, N> {
fn default() -> Self {
Self::new()
}
}
impl<T: Copy, const N: usize> core::fmt::Debug for BlockBuilder<T, N> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("BlockBuilder")
.field("len", &self.len)
.field("capacity", &N)
.field("first_sequence", &self.first_sequence)
.field("last_sequence", &self.last_sequence)
.finish_non_exhaustive()
}
}
const fn next_sequence(sequence: u32) -> u32 {
if sequence == u32::MAX {
1
} else {
sequence + 1
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::EventBuf;
#[test]
fn completes_only_after_n_contiguous_samples() {
let mut fill = BlockBuilder::<u16, 3>::new();
assert_eq!(fill.push(41, 4), Ok(None));
assert_eq!(fill.push(42, 5), Ok(None));
let block = fill.push(43, 6).unwrap().unwrap();
assert_eq!(block.first_sequence(), 41);
assert_eq!(block.last_sequence(), 43);
assert_eq!(block.samples(), &[4, 5, 6]);
assert!(fill.is_empty());
}
#[test]
fn discontinuity_check_is_modular_over_the_span() {
let mut fill = BlockBuilder::<u16, 2>::new();
assert_eq!(fill.push(u32::MAX, 1), Ok(None));
assert!(matches!(
fill.push(2, 9),
Err(FillError::Discontinuous {
expected: 1,
received: 2,
..
})
));
let block = fill.push(1, 2).unwrap().unwrap();
assert_eq!(block.samples(), &[1, 2]);
}
#[test]
fn completion_resets_for_the_next_block() {
let mut fill = BlockBuilder::<u8, 1>::new();
assert_eq!(fill.push(7, 9).unwrap().unwrap().into_samples(), [9]);
assert_eq!(fill.push(8, 10).unwrap().unwrap().into_samples(), [10]);
}
#[test]
fn rejects_reserved_zero_without_changing_partial_block() {
let mut fill = BlockBuilder::<u8, 2>::new();
assert_eq!(fill.push(9, 1), Ok(None));
assert_eq!(
fill.push(0, 2),
Err(FillError::ReservedSequence { sample: 2 })
);
assert_eq!(fill.len(), 1);
assert_eq!(fill.expected_sequence(), Some(10));
}
#[test]
fn rejects_gap_without_hiding_loss_policy() {
let mut fill = BlockBuilder::<u8, 2>::new();
assert_eq!(fill.push(9, 1), Ok(None));
assert_eq!(
fill.push(11, 2),
Err(FillError::Discontinuous {
expected: 10,
received: 11,
sample: 2
})
);
assert_eq!(fill.len(), 1);
}
#[test]
fn clear_discards_a_partial_block() {
let mut fill = BlockBuilder::<u8, 4>::new();
assert_eq!(fill.push(1, 1), Ok(None));
fill.clear();
assert!(fill.is_empty());
assert_eq!(fill.expected_sequence(), None);
assert_eq!(fill.push(20, 2), Ok(None));
}
#[test]
fn sequence_wrap_skips_zero() {
let mut fill = BlockBuilder::<u8, 2>::new();
assert_eq!(fill.push(u32::MAX, 1), Ok(None));
let block = fill.push(1, 2).unwrap().unwrap();
assert_eq!(block.first_sequence(), u32::MAX);
assert_eq!(block.last_sequence(), 1);
}
#[test]
fn works_without_default_bound() {
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
struct NoDefault(u8);
let mut fill = BlockBuilder::<NoDefault, 1>::new();
let block = fill.push(1, NoDefault(3)).unwrap().unwrap();
assert_eq!(block.samples(), &[NoDefault(3)]);
}
#[test]
fn default_and_capacity_match_new() {
let fill = BlockBuilder::<u8, 8>::default();
assert_eq!(fill.capacity(), 8);
assert_eq!(fill.len(), 0);
}
#[test]
fn event_buf_composition_queues_and_returns_a_rejected_block() {
let mut fill = BlockBuilder::<u8, 2>::new();
assert_eq!(fill.push(1, 10), Ok(None));
let first = fill.push(2, 11).unwrap().unwrap();
assert_eq!(fill.push(3, 12), Ok(None));
let second = fill.push(4, 13).unwrap().unwrap();
let queue = EventBuf::<Block<u8, 2>, 1>::new();
let producer = queue.try_producer().unwrap();
let consumer = queue.try_consumer().unwrap();
assert_eq!(producer.push(first), Ok(()));
assert_eq!(producer.push(second), Err(second));
assert_eq!(consumer.pop(), Some(first));
}
}