#[derive(Debug, Clone, Copy, Default)]
pub struct PacketMeta {
pub timestamp_ns: u64,
pub queue_id: u16,
pub original_len: u32,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PushResult {
Ok,
Truncated,
Full,
}
pub struct RawBatch {
bufs: Vec<Vec<u8>>,
lens: Vec<usize>,
meta: Vec<PacketMeta>,
count: usize,
max_frame_size: usize,
received: u64,
dropped: u64,
truncated: u64,
}
impl RawBatch {
pub fn new(capacity: usize, max_frame_size: usize) -> Self {
assert!(capacity > 0, "RawBatch capacity must be > 0");
assert!(max_frame_size > 0, "RawBatch max_frame_size must be > 0");
let bufs = (0..capacity).map(|_| vec![0u8; max_frame_size]).collect();
let lens = vec![0usize; capacity];
let meta = vec![PacketMeta::default(); capacity];
Self {
bufs,
lens,
meta,
count: 0,
max_frame_size,
received: 0,
dropped: 0,
truncated: 0,
}
}
pub fn reset(&mut self, _max_frame_size: usize) {
self.count = 0;
}
pub fn capacity(&self) -> usize {
self.bufs.len()
}
pub fn max_frame_size(&self) -> usize {
self.max_frame_size
}
pub fn len(&self) -> usize {
self.count
}
pub fn is_empty(&self) -> bool {
self.count == 0
}
pub fn received(&self) -> u64 {
self.received
}
pub fn dropped(&self) -> u64 {
self.dropped
}
pub fn truncated(&self) -> u64 {
self.truncated
}
pub fn record_drop(&mut self) {
self.dropped = self.dropped.saturating_add(1);
}
pub fn packets(&self) -> impl Iterator<Item = (&[u8], &PacketMeta)> {
(0..self.count).map(move |i| (&self.bufs[i][..self.lens[i]], &self.meta[i]))
}
pub fn packet_mut(&mut self, index: usize) -> Option<(&mut [u8], &mut PacketMeta)> {
if index >= self.count {
return None;
}
let len = self.lens[index];
Some((&mut self.bufs[index][..len], &mut self.meta[index]))
}
pub fn truncate(&mut self, new_len: usize) {
if new_len < self.count {
self.count = new_len;
}
}
pub fn push(&mut self, data: &[u8], mut meta: PacketMeta) -> PushResult {
if self.count >= self.bufs.len() {
return PushResult::Full;
}
let slot = &mut self.bufs[self.count];
let true_len = data.len();
if meta.original_len == 0 || (meta.original_len as usize) < true_len {
meta.original_len = true_len.min(u32::MAX as usize) as u32;
}
let copy_len = true_len.min(slot.len());
slot[..copy_len].copy_from_slice(&data[..copy_len]);
self.lens[self.count] = copy_len;
self.meta[self.count] = meta;
self.count += 1;
self.received += 1;
if copy_len < true_len {
self.truncated += 1;
PushResult::Truncated
} else {
PushResult::Ok
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn push_and_iterate() {
let mut batch = RawBatch::new(4, 64);
let meta = PacketMeta {
timestamp_ns: 1000,
queue_id: 0,
original_len: 10,
};
assert_eq!(batch.push(b"hello", meta), PushResult::Ok);
assert_eq!(batch.len(), 1);
assert_eq!(batch.received(), 1);
let packets: Vec<_> = batch.packets().collect();
assert_eq!(packets.len(), 1);
assert_eq!(packets[0].0, b"hello");
assert_eq!(packets[0].1.timestamp_ns, 1000);
}
#[test]
fn full_batch_returns_full() {
let mut batch = RawBatch::new(2, 64);
let meta = PacketMeta::default();
assert_eq!(batch.push(b"a", meta), PushResult::Ok);
assert_eq!(batch.push(b"b", meta), PushResult::Ok);
assert_eq!(batch.push(b"c", meta), PushResult::Full);
assert_eq!(batch.len(), 2);
}
#[test]
fn reset_reuses_allocation() {
let mut batch = RawBatch::new(4, 64);
let meta = PacketMeta::default();
batch.push(b"packet1", meta);
batch.push(b"packet2", meta);
assert_eq!(batch.len(), 2);
batch.reset(64);
assert_eq!(batch.len(), 0);
assert_eq!(batch.push(b"packet3", meta), PushResult::Ok);
let packets: Vec<_> = batch.packets().collect();
assert_eq!(packets[0].0, b"packet3");
}
#[test]
fn truncates_and_sets_original_len() {
let mut batch = RawBatch::new(1, 4);
let meta = PacketMeta::default();
assert_eq!(batch.push(b"hello world", meta), PushResult::Truncated);
let packets: Vec<_> = batch.packets().collect();
assert_eq!(packets[0].0, b"hell");
assert_eq!(packets[0].1.original_len, 11);
assert_eq!(batch.truncated(), 1);
}
#[test]
fn received_counter_increments() {
let mut batch = RawBatch::new(4, 64);
let meta = PacketMeta::default();
for _ in 0..3 {
batch.push(b"x", meta);
}
assert_eq!(batch.received(), 3);
batch.reset(64);
assert_eq!(batch.received(), 3);
}
}