use aligned_vec::{AVec, ConstAlign};
use bytes::Bytes;
use core::ops::Range;
use thiserror::Error;
use crate::value::{EventType, Metadata, Payload, SchemaVersion};
pub const PAYLOAD_ALIGN: usize = 16;
pub(crate) const HEADER_FIXED_SIZE: usize = 11;
pub(crate) const VERSION_OFFSET: usize = 0;
pub(crate) const SCHEMA_VERSION_OFFSET: usize = 1;
pub(crate) const EVENT_TYPE_LEN_OFFSET: usize = 5;
pub(crate) const META_LEN_OFFSET: usize = 7;
const HEADER_FIXED_SIZE_V1: usize = 19;
const SCHEMA_VERSION_OFFSET_V1: usize = 9;
const EVENT_TYPE_LEN_OFFSET_V1: usize = 13;
const META_LEN_OFFSET_V1: usize = 15;
pub(crate) const META_LEN_ABSENT: u32 = u32::MAX;
#[inline]
const fn align_padding(offset: usize, align: usize) -> usize {
(align - (offset % align)) % align
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum FrameFormatVersion {
V1,
V2,
}
impl FrameFormatVersion {
pub(crate) const CURRENT: Self = Self::V2;
#[inline]
const fn to_u8(self) -> u8 {
match self {
Self::V1 => 1,
Self::V2 => 2,
}
}
#[inline]
const fn from_u8(byte: u8) -> Option<Self> {
match byte {
1 => Some(Self::V1),
2 => Some(Self::V2),
_ => None,
}
}
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct FrameHeader {
pub(crate) format_version: FrameFormatVersion,
pub(crate) schema_version: u32,
event_type_len: u16,
metadata_len: Option<u32>,
}
impl FrameHeader {
pub(crate) const SIZE: usize = HEADER_FIXED_SIZE;
fn from_validated_lengths(
format_version: FrameFormatVersion,
schema_version: u32,
event_type_len: usize,
metadata_len: Option<usize>,
) -> Self {
#[allow(
clippy::expect_used,
reason = "validated by EventType::from_bytes invariant: length ≤ u16::MAX"
)]
let event_type_len_u16 = u16::try_from(event_type_len)
.expect("event_type length validated by EventType invariant");
let metadata_len_u32 = metadata_len.map(|n| {
#[allow(
clippy::expect_used,
reason = "validated by Metadata::from_bytes invariant: length ≤ MAX_METADATA_LEN"
)]
let v = u32::try_from(n).expect("metadata length validated by Metadata invariant");
v
});
Self {
format_version,
schema_version,
event_type_len: event_type_len_u16,
metadata_len: metadata_len_u32,
}
}
fn write_into(&self, buf: &mut AVec<u8, ConstAlign<PAYLOAD_ALIGN>>) {
let meta_field = self.metadata_len.unwrap_or(META_LEN_ABSENT);
buf.extend_from_slice(&[self.format_version.to_u8()]);
buf.extend_from_slice(&self.schema_version.to_le_bytes());
buf.extend_from_slice(&self.event_type_len.to_le_bytes());
buf.extend_from_slice(&meta_field.to_le_bytes());
}
pub(crate) fn read_from(value: &[u8]) -> Result<Self, DecodeError> {
if value.len() < Self::SIZE {
return Err(DecodeError::ValueTooShort {
min: Self::SIZE,
actual: value.len(),
});
}
let version_byte = value[VERSION_OFFSET];
let format_version = FrameFormatVersion::from_u8(version_byte).ok_or(
DecodeError::UnsupportedFrameVersion {
version: version_byte,
},
)?;
let (schema_off, et_len_off, meta_off) = match format_version {
FrameFormatVersion::V1 => {
if value.len() < HEADER_FIXED_SIZE_V1 {
return Err(DecodeError::ValueTooShort {
min: HEADER_FIXED_SIZE_V1,
actual: value.len(),
});
}
(
SCHEMA_VERSION_OFFSET_V1,
EVENT_TYPE_LEN_OFFSET_V1,
META_LEN_OFFSET_V1,
)
}
FrameFormatVersion::V2 => (
SCHEMA_VERSION_OFFSET,
EVENT_TYPE_LEN_OFFSET,
META_LEN_OFFSET,
),
};
let schema_version = u32::from_le_bytes([
value[schema_off],
value[schema_off + 1],
value[schema_off + 2],
value[schema_off + 3],
]);
let event_type_len = u16::from_le_bytes([value[et_len_off], value[et_len_off + 1]]);
let meta_field = u32::from_le_bytes([
value[meta_off],
value[meta_off + 1],
value[meta_off + 2],
value[meta_off + 3],
]);
let metadata_len = if meta_field == META_LEN_ABSENT {
None
} else {
Some(meta_field)
};
Ok(Self {
format_version,
schema_version,
event_type_len,
metadata_len,
})
}
}
#[derive(Debug, Clone)]
struct FrameLayout {
event_type: Range<u32>,
metadata: Option<Range<u32>>,
payload: Range<u32>,
padding: usize,
total: usize,
}
#[inline]
const fn length_overflow(header: usize, padding: usize, payload: usize) -> WireError {
WireError::FrameLengthOverflow {
header,
padding,
payload,
}
}
impl FrameLayout {
fn compute_from_validated_lengths(
event_type_len: usize,
metadata_len: Option<usize>,
payload_len: usize,
) -> Result<Self, WireError> {
let meta_len_usize = metadata_len.unwrap_or(0);
let pre_payload_len = HEADER_FIXED_SIZE
.checked_add(event_type_len)
.and_then(|n| n.checked_add(meta_len_usize))
.ok_or_else(|| length_overflow(HEADER_FIXED_SIZE, 0, payload_len))?;
let padding = align_padding(pre_payload_len, PAYLOAD_ALIGN);
let total = pre_payload_len
.checked_add(padding)
.and_then(|n| n.checked_add(payload_len))
.ok_or_else(|| length_overflow(pre_payload_len, padding, payload_len))?;
let overflow = || length_overflow(pre_payload_len, padding, payload_len);
let event_type_start = u32::try_from(HEADER_FIXED_SIZE).map_err(|_| overflow())?;
let event_type_len_u32 = u32::try_from(event_type_len).map_err(|_| overflow())?;
let event_type_end = event_type_start
.checked_add(event_type_len_u32)
.ok_or_else(overflow)?;
let metadata_range = metadata_len
.map(|n| -> Result<Range<u32>, WireError> {
let n_u32 = u32::try_from(n).map_err(|_| overflow())?;
let end = event_type_end.checked_add(n_u32).ok_or_else(overflow)?;
Ok(event_type_end..end)
})
.transpose()?;
let payload_start_usize = pre_payload_len.checked_add(padding).ok_or_else(overflow)?;
let payload_start = u32::try_from(payload_start_usize).map_err(|_| overflow())?;
let payload_len_u32 = u32::try_from(payload_len).map_err(|_| overflow())?;
let payload_end = payload_start
.checked_add(payload_len_u32)
.ok_or_else(overflow)?;
Ok(Self {
event_type: event_type_start..event_type_end,
metadata: metadata_range,
payload: payload_start..payload_end,
padding,
total,
})
}
}
#[derive(Debug)]
pub struct EncodedFrame {
pub value: Bytes,
pub offsets: FrameOffsets,
}
#[derive(Debug, Clone)]
pub struct FrameOffsets {
pub event_type: Range<u32>,
pub metadata: Option<Range<u32>>,
pub payload: Range<u32>,
}
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum WireError {
#[error(
"frame length overflow combining header={header}, padding={padding}, payload={payload}"
)]
FrameLengthOverflow {
header: usize,
padding: usize,
payload: usize,
},
}
#[derive(Debug)]
pub struct DecodedFrame {
pub schema_version: SchemaVersion,
pub offsets: FrameOffsets,
}
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum DecodeError {
#[error("value too short: need at least {min} bytes, got {actual}")]
ValueTooShort { min: usize, actual: usize },
#[error(
"unsupported frame format version on wire: got {version}, this build supports up to {}",
FrameFormatVersion::CURRENT.to_u8()
)]
UnsupportedFrameVersion { version: u8 },
#[error("event type length {et_len} extends past value (len={value_len})")]
EventTypeTruncated { et_len: usize, value_len: usize },
#[error("metadata length {meta_len} extends past value (len={value_len})")]
MetadataTruncated { meta_len: u32, value_len: usize },
#[error("computed offset overflows u32 (value len={value_len})")]
OffsetOverflow { value_len: usize },
#[error("corrupt schema_version on wire: got 0, must be > 0")]
CorruptSchemaVersion,
}
#[derive(Debug)]
struct FramePlan<'a> {
header: FrameHeader,
event_type_bytes: &'a [u8],
metadata: Option<&'a [u8]>,
payload: &'a [u8],
layout: FrameLayout,
}
fn plan<'a>(
schema_version: SchemaVersion,
event_type: &'a EventType,
payload: &'a Payload,
metadata: Option<&'a Metadata>,
) -> Result<FramePlan<'a>, WireError> {
let event_type_bytes = event_type.as_bytes();
let metadata_bytes = metadata.map(Metadata::as_slice);
let payload_bytes = payload.as_slice();
let layout = FrameLayout::compute_from_validated_lengths(
event_type_bytes.len(),
metadata_bytes.map(<[u8]>::len),
payload_bytes.len(),
)?;
let header = FrameHeader::from_validated_lengths(
FrameFormatVersion::CURRENT,
schema_version.get(),
event_type_bytes.len(),
metadata_bytes.map(<[u8]>::len),
);
Ok(FramePlan {
header,
event_type_bytes,
metadata: metadata_bytes,
payload: payload_bytes,
layout,
})
}
fn execute(plan: FramePlan<'_>) -> EncodedFrame {
let mut buf: AVec<u8, ConstAlign<PAYLOAD_ALIGN>> =
AVec::with_capacity(PAYLOAD_ALIGN, plan.layout.total);
plan.header.write_into(&mut buf);
buf.extend_from_slice(plan.event_type_bytes);
if let Some(m) = plan.metadata {
buf.extend_from_slice(m);
}
buf.resize(buf.len() + plan.layout.padding, 0u8);
buf.extend_from_slice(plan.payload);
EncodedFrame {
value: Bytes::from_owner(buf),
offsets: FrameOffsets {
event_type: plan.layout.event_type,
metadata: plan.layout.metadata,
payload: plan.layout.payload,
},
}
}
pub fn encode_frame(
schema_version: SchemaVersion,
event_type: &EventType,
payload: &Payload,
metadata: Option<&Metadata>,
) -> Result<EncodedFrame, WireError> {
plan(schema_version, event_type, payload, metadata).map(execute)
}
pub fn decode_frame(value: &[u8]) -> Result<DecodedFrame, DecodeError> {
let header = FrameHeader::read_from(value)?;
let header_size = match header.format_version {
FrameFormatVersion::V1 => HEADER_FIXED_SIZE_V1,
FrameFormatVersion::V2 => HEADER_FIXED_SIZE,
};
decode_frame_body(value, header, header_size)
}
fn decode_frame_body(
value: &[u8],
header: FrameHeader,
header_size: usize,
) -> Result<DecodedFrame, DecodeError> {
let schema_version = SchemaVersion::from_u32(header.schema_version)
.map_err(|_| DecodeError::CorruptSchemaVersion)?;
let et_len = usize::from(header.event_type_len);
let et_start = header_size;
let et_end = et_start
.checked_add(et_len)
.ok_or(DecodeError::OffsetOverflow {
value_len: value.len(),
})?;
if value.len() < et_end {
return Err(DecodeError::EventTypeTruncated {
et_len,
value_len: value.len(),
});
}
let (metadata_range, post_meta) = match header.metadata_len {
None => (None, et_end),
Some(meta_len) => {
let meta_len_usize =
usize::try_from(meta_len).map_err(|_| DecodeError::OffsetOverflow {
value_len: value.len(),
})?;
let meta_end =
et_end
.checked_add(meta_len_usize)
.ok_or(DecodeError::OffsetOverflow {
value_len: value.len(),
})?;
if value.len() < meta_end {
return Err(DecodeError::MetadataTruncated {
meta_len,
value_len: value.len(),
});
}
let m_start_u32 = u32::try_from(et_end).map_err(|_| DecodeError::OffsetOverflow {
value_len: value.len(),
})?;
let m_end_u32 = u32::try_from(meta_end).map_err(|_| DecodeError::OffsetOverflow {
value_len: value.len(),
})?;
(Some(m_start_u32..m_end_u32), meta_end)
}
};
let padding = align_padding(post_meta, PAYLOAD_ALIGN);
let payload_start = post_meta
.checked_add(padding)
.ok_or(DecodeError::OffsetOverflow {
value_len: value.len(),
})?;
let payload_end = value.len();
if payload_start > payload_end {
return Err(DecodeError::OffsetOverflow {
value_len: value.len(),
});
}
let et_start_u32 = u32::try_from(et_start).map_err(|_| DecodeError::OffsetOverflow {
value_len: value.len(),
})?;
let et_end_u32 = u32::try_from(et_end).map_err(|_| DecodeError::OffsetOverflow {
value_len: value.len(),
})?;
let payload_start_u32 =
u32::try_from(payload_start).map_err(|_| DecodeError::OffsetOverflow {
value_len: value.len(),
})?;
let payload_end_u32 = u32::try_from(payload_end).map_err(|_| DecodeError::OffsetOverflow {
value_len: value.len(),
})?;
Ok(DecodedFrame {
schema_version,
offsets: FrameOffsets {
event_type: et_start_u32..et_end_u32,
metadata: metadata_range,
payload: payload_start_u32..payload_end_u32,
},
})
}
#[cfg(test)]
#[allow(
clippy::as_conversions,
clippy::cast_possible_truncation,
clippy::panic,
clippy::redundant_clone,
clippy::single_match_else,
reason = "test code: index arithmetic, prop_assert_eq macro expansions, \
and `panic!(\"expected X, got {other:?}\")` arms surface failing test diagnostics"
)]
mod tests {
use super::*;
use crate::value::{MAX_EVENT_TYPE_LEN, MAX_METADATA_LEN, MAX_PAYLOAD_LEN};
use proptest::prelude::*;
fn payload_ptr_aligned(frame: &EncodedFrame) -> bool {
let start = usize::try_from(frame.offsets.payload.start).expect("u32 fits usize");
let end = usize::try_from(frame.offsets.payload.end).expect("u32 fits usize");
let payload_slice = &frame.value[start..end];
payload_slice.as_ptr().addr().is_multiple_of(PAYLOAD_ALIGN)
}
fn et(s: &str) -> EventType {
EventType::from_bytes(Bytes::copy_from_slice(s.as_bytes())).expect("test event_type valid")
}
fn pl(b: &[u8]) -> Payload {
Payload::from_bytes(Bytes::copy_from_slice(b)).expect("test payload valid")
}
fn md(b: &[u8]) -> Metadata {
Metadata::from_bytes(Bytes::copy_from_slice(b)).expect("test metadata non-empty + valid")
}
fn sv1() -> SchemaVersion {
SchemaVersion::INITIAL
}
fn align_and_offset() -> impl Strategy<Value = (usize, usize)> {
(0u32..16).prop_flat_map(|align_pow| {
let align = 1usize << align_pow;
let offset = prop_oneof![
1 => Just(0usize),
1 => Just(1usize),
1 => Just(align.saturating_sub(1)),
1 => Just(align),
1 => Just(align + 1),
10 => 0usize..1_000_000,
];
(Just(align), offset)
})
}
fn u32_strategy() -> impl Strategy<Value = u32> {
prop_oneof![
1 => Just(0u32),
1 => Just(1u32),
1 => Just(u32::MAX - 1),
1 => Just(u32::MAX),
10 => any::<u32>(),
]
}
fn schema_version_strategy() -> impl Strategy<Value = SchemaVersion> {
prop_oneof![
1 => Just(1u32),
1 => Just(2u32),
1 => Just(u32::MAX - 1),
1 => Just(u32::MAX),
10 => 1u32..=u32::MAX,
]
.prop_map(|v| SchemaVersion::from_u32(v).expect("nonzero strategy"))
}
fn u16_strategy() -> impl Strategy<Value = u16> {
prop_oneof![
1 => Just(0u16),
1 => Just(1u16),
1 => Just(u16::MAX - 1),
1 => Just(u16::MAX),
10 => any::<u16>(),
]
}
fn frame_body_length() -> impl Strategy<Value = usize> {
prop_oneof![
1 => Just(0usize),
1 => Just(1usize),
1 => Just(PAYLOAD_ALIGN - 1),
1 => Just(PAYLOAD_ALIGN),
1 => Just(PAYLOAD_ALIGN + 1),
10 => 0usize..=4096,
]
}
fn event_type_str_strategy() -> impl Strategy<Value = String> {
prop_oneof![
1 => Just(String::new()),
1 => Just("a".to_owned()),
10 => prop::collection::vec(any::<char>(), 0..=256)
.prop_map(|chars| chars.into_iter().collect::<String>()),
]
}
fn metadata_bytes_strategy() -> impl Strategy<Value = Option<Vec<u8>>> {
prop_oneof![
1 => Just(None),
1 => Just(Some(vec![0u8])),
10 => prop::option::of(prop::collection::vec(any::<u8>(), 1..512)),
]
}
fn payload_bytes_strategy() -> impl Strategy<Value = Vec<u8>> {
prop_oneof![
1 => Just(Vec::<u8>::new()),
1 => Just(vec![0u8]),
10 => prop::collection::vec(any::<u8>(), 0..2048),
]
}
prop_compose! {
fn valid_frame_inputs()(
schema_version in schema_version_strategy(),
event_type in event_type_str_strategy(),
metadata in metadata_bytes_strategy(),
payload in payload_bytes_strategy(),
) -> (SchemaVersion, String, Option<Vec<u8>>, Vec<u8>) {
(schema_version, event_type, metadata, payload)
}
}
proptest! {
#[test]
fn payload_pointer_is_16_aligned(
(schema_version, event_type, metadata, payload) in valid_frame_inputs(),
) {
let et_v = et(&event_type);
let pl_v = pl(&payload);
let md_v = metadata.as_deref().map(md);
let frame = encode_frame(schema_version, &et_v, &pl_v, md_v.as_ref())
.expect("encode_frame succeeds on bounded inputs");
prop_assert!(payload_ptr_aligned(&frame));
}
#[test]
fn ranges_recover_each_field(
(schema_version, event_type, metadata, payload) in valid_frame_inputs(),
) {
let et_v = et(&event_type);
let pl_v = pl(&payload);
let md_v = metadata.as_deref().map(md);
let frame = encode_frame(schema_version, &et_v, &pl_v, md_v.as_ref())
.expect("encode_frame succeeds on bounded inputs");
let v = &frame.value;
prop_assert_eq!(
&v[frame.offsets.event_type.start as usize..frame.offsets.event_type.end as usize],
event_type.as_bytes()
);
prop_assert_eq!(
&v[frame.offsets.payload.start as usize..frame.offsets.payload.end as usize],
payload.as_slice()
);
if let (Some(meta), Some(range)) = (metadata.as_deref(), frame.offsets.metadata) {
prop_assert_eq!(
&v[range.start as usize..range.end as usize],
meta
);
}
}
#[test]
fn header_fields_are_recoverable(
(schema_version, event_type, metadata, payload) in valid_frame_inputs(),
) {
let et_v = et(&event_type);
let pl_v = pl(&payload);
let md_v = metadata.as_deref().map(md);
let frame = encode_frame(schema_version, &et_v, &pl_v, md_v.as_ref())
.expect("encode_frame succeeds on bounded inputs");
let v = &frame.value;
let mut sv_buf = [0u8; 4];
sv_buf.copy_from_slice(&v[SCHEMA_VERSION_OFFSET..EVENT_TYPE_LEN_OFFSET]);
prop_assert_eq!(u32::from_le_bytes(sv_buf), schema_version.get());
let mut et_len_buf = [0u8; 2];
et_len_buf.copy_from_slice(&v[EVENT_TYPE_LEN_OFFSET..EVENT_TYPE_LEN_OFFSET + 2]);
prop_assert_eq!(usize::from(u16::from_le_bytes(et_len_buf)), event_type.len());
let mut ml_buf = [0u8; 4];
ml_buf.copy_from_slice(&v[META_LEN_OFFSET..META_LEN_OFFSET + 4]);
let ml = u32::from_le_bytes(ml_buf);
match metadata.as_deref() {
Some(m) => prop_assert_eq!(usize::try_from(ml).unwrap(), m.len()),
None => prop_assert_eq!(ml, META_LEN_ABSENT),
}
}
#[test]
fn encoded_frame_carries_v2_version_byte(
(schema_version, event_type, metadata, payload) in valid_frame_inputs(),
) {
let et_v = et(&event_type);
let pl_v = pl(&payload);
let md_v = metadata.as_deref().map(md);
let frame = encode_frame(schema_version, &et_v, &pl_v, md_v.as_ref())
.expect("encode_frame succeeds on bounded inputs");
prop_assert_eq!(frame.value[VERSION_OFFSET], 2);
let header = FrameHeader::read_from(&frame.value)
.expect("header reads back from a freshly built frame");
prop_assert_eq!(header.format_version, FrameFormatVersion::V2);
}
}
#[test]
fn empty_payload_still_aligned() {
let frame = encode_frame(sv1(), &et("X"), &pl(b""), None).expect("trivial frame builds");
assert!(payload_ptr_aligned(&frame));
assert_eq!(frame.offsets.payload.start, frame.offsets.payload.end);
}
#[test]
fn empty_event_type_permitted() {
let frame = encode_frame(sv1(), &et(""), &pl(b"data"), None)
.expect("empty event_type accepted at wire layer");
assert!(payload_ptr_aligned(&frame));
}
#[test]
fn max_event_type_accepted() {
let huge = "a".repeat(MAX_EVENT_TYPE_LEN);
encode_frame(sv1(), &et(&huge), &pl(b"d"), None).expect("max-length event_type accepted");
}
#[test]
fn meta_len_u32_max_is_absent_sentinel() {
let frame =
encode_frame(sv1(), &et("X"), &pl(b"d"), None).expect("none-metadata frame builds");
let mut ml_buf = [0u8; 4];
ml_buf.copy_from_slice(&frame.value[META_LEN_OFFSET..META_LEN_OFFSET + 4]);
assert_eq!(u32::from_le_bytes(ml_buf), META_LEN_ABSENT);
assert!(frame.offsets.metadata.is_none());
}
proptest! {
#[test]
fn build_then_decode_round_trips(
(schema_version, event_type, metadata, payload) in valid_frame_inputs(),
) {
let et_v = et(&event_type);
let pl_v = pl(&payload);
let md_v = metadata.as_deref().map(md);
let frame = encode_frame(schema_version, &et_v, &pl_v, md_v.as_ref())
.expect("encode_frame succeeds on bounded inputs");
let decoded = decode_frame(&frame.value).expect("decode_frame succeeds on a built frame");
prop_assert_eq!(decoded.schema_version, schema_version);
prop_assert_eq!(decoded.offsets.event_type.clone(), frame.offsets.event_type.clone());
prop_assert_eq!(decoded.offsets.metadata.clone(), frame.offsets.metadata.clone());
prop_assert_eq!(decoded.offsets.payload.clone(), frame.offsets.payload.clone());
}
}
#[test]
fn decode_rejects_truncated_value() {
let too_short = vec![0u8; HEADER_FIXED_SIZE - 1];
assert!(matches!(
decode_frame(&too_short),
Err(DecodeError::ValueTooShort { .. })
));
}
#[test]
fn decode_rejects_truncated_event_type() {
let mut buf = vec![0u8; HEADER_FIXED_SIZE];
buf[VERSION_OFFSET] = 2;
buf[SCHEMA_VERSION_OFFSET..EVENT_TYPE_LEN_OFFSET].copy_from_slice(&1u32.to_le_bytes());
buf[EVENT_TYPE_LEN_OFFSET..EVENT_TYPE_LEN_OFFSET + 2]
.copy_from_slice(&100u16.to_le_bytes());
buf[META_LEN_OFFSET..META_LEN_OFFSET + 4].copy_from_slice(&META_LEN_ABSENT.to_le_bytes());
assert!(matches!(
decode_frame(&buf),
Err(DecodeError::EventTypeTruncated { .. })
));
}
#[test]
fn decode_rejects_truncated_metadata() {
let mut buf = vec![0u8; HEADER_FIXED_SIZE];
buf[VERSION_OFFSET] = 2;
buf[SCHEMA_VERSION_OFFSET..EVENT_TYPE_LEN_OFFSET].copy_from_slice(&1u32.to_le_bytes());
buf[EVENT_TYPE_LEN_OFFSET..EVENT_TYPE_LEN_OFFSET + 2].copy_from_slice(&0u16.to_le_bytes());
buf[META_LEN_OFFSET..META_LEN_OFFSET + 4].copy_from_slice(&100u32.to_le_bytes());
assert!(matches!(
decode_frame(&buf),
Err(DecodeError::MetadataTruncated { .. })
));
}
#[test]
fn decode_rejects_corrupt_schema_version_zero() {
let frame = encode_frame(sv1(), &et("X"), &pl(b"p"), None).expect("encode");
let mut bytes_vec = frame.value.to_vec();
bytes_vec[SCHEMA_VERSION_OFFSET..EVENT_TYPE_LEN_OFFSET].fill(0);
let tampered = Bytes::from(bytes_vec);
assert!(matches!(
decode_frame(&tampered),
Err(DecodeError::CorruptSchemaVersion)
));
}
fn adversarial_decode_bytes() -> impl Strategy<Value = Vec<u8>> {
let header_shaped = (
prop_oneof![10 => Just(2u8), 1 => any::<u8>()],
any::<u32>(),
0u16..=64,
prop_oneof![Just(META_LEN_ABSENT), 0u32..=64],
prop::collection::vec(any::<u8>(), 0..=512),
)
.prop_map(|(version, sv, et_len, meta_len, body)| {
let mut buf = Vec::with_capacity(HEADER_FIXED_SIZE + body.len());
buf.extend_from_slice(&[version]);
buf.extend_from_slice(&sv.to_le_bytes());
buf.extend_from_slice(&et_len.to_le_bytes());
buf.extend_from_slice(&meta_len.to_le_bytes());
buf.extend_from_slice(&body);
buf
});
prop_oneof![
1 => Just(Vec::<u8>::new()),
1 => Just(vec![0u8]),
1 => prop::collection::vec(any::<u8>(), HEADER_FIXED_SIZE - 1..=HEADER_FIXED_SIZE - 1),
1 => prop::collection::vec(any::<u8>(), HEADER_FIXED_SIZE..=HEADER_FIXED_SIZE),
1 => prop::collection::vec(any::<u8>(), HEADER_FIXED_SIZE + 1..=HEADER_FIXED_SIZE + 1),
5 => prop::collection::vec(any::<u8>(), 0..=4096),
5 => header_shaped,
]
}
proptest! {
#[test]
fn decode_never_panics(bytes in adversarial_decode_bytes()) {
let _ = decode_frame(&bytes);
}
#[test]
fn decode_offsets_in_bounds_on_success(bytes in adversarial_decode_bytes()) {
if let Ok(decoded) = decode_frame(&bytes) {
let len_u32 = u32::try_from(bytes.len()).unwrap_or(u32::MAX);
prop_assert!(decoded.offsets.event_type.start <= decoded.offsets.event_type.end);
prop_assert!(decoded.offsets.event_type.end <= len_u32);
if let Some(meta) = decoded.offsets.metadata {
prop_assert!(meta.start <= meta.end);
prop_assert!(meta.end <= len_u32);
}
prop_assert!(decoded.offsets.payload.start <= decoded.offsets.payload.end);
prop_assert!(decoded.offsets.payload.end <= len_u32);
}
}
}
#[test]
fn align_padding_zero_offset_yields_zero() {
assert_eq!(align_padding(0, PAYLOAD_ALIGN), 0);
}
#[test]
fn align_padding_one_below_boundary_yields_one() {
assert_eq!(align_padding(15, PAYLOAD_ALIGN), 1);
}
#[test]
fn align_padding_on_boundary_yields_zero() {
assert_eq!(align_padding(PAYLOAD_ALIGN, PAYLOAD_ALIGN), 0);
}
#[test]
fn align_padding_one_above_boundary_yields_fifteen() {
assert_eq!(align_padding(PAYLOAD_ALIGN + 1, PAYLOAD_ALIGN), 15);
}
proptest! {
#[test]
fn align_padding_invariants(
(align, offset) in align_and_offset(),
) {
let pad = align_padding(offset, align);
prop_assert!(
(offset + pad).is_multiple_of(align),
"offset={offset} align={align} pad={pad} not multiple",
);
prop_assert!(pad < align, "pad {pad} >= align {align}");
prop_assert_eq!(pad == 0, offset.is_multiple_of(align));
}
}
fn fresh_buf() -> AVec<u8, ConstAlign<PAYLOAD_ALIGN>> {
AVec::with_capacity(PAYLOAD_ALIGN, 64)
}
#[test]
fn frame_header_write_into_writes_all_fields_at_correct_offsets() {
let header = FrameHeader {
format_version: FrameFormatVersion::V2,
schema_version: 0x090A_0B0C,
event_type_len: 0x0D0E,
metadata_len: Some(0x0F10_1112),
};
let mut buf = fresh_buf();
header.write_into(&mut buf);
assert_eq!(buf.len(), FrameHeader::SIZE);
assert_eq!(buf[VERSION_OFFSET], 2);
assert_eq!(
&buf[SCHEMA_VERSION_OFFSET..EVENT_TYPE_LEN_OFFSET],
&0x090A_0B0Cu32.to_le_bytes(),
);
assert_eq!(
&buf[EVENT_TYPE_LEN_OFFSET..EVENT_TYPE_LEN_OFFSET + 2],
&0x0D0Eu16.to_le_bytes(),
);
assert_eq!(
&buf[META_LEN_OFFSET..META_LEN_OFFSET + 4],
&0x0F10_1112u32.to_le_bytes(),
);
}
#[test]
fn frame_header_none_metadata_encodes_sentinel() {
let header = FrameHeader {
format_version: FrameFormatVersion::V2,
schema_version: 1,
event_type_len: 0,
metadata_len: None,
};
let mut buf = fresh_buf();
header.write_into(&mut buf);
let mut ml = [0u8; 4];
ml.copy_from_slice(&buf[META_LEN_OFFSET..META_LEN_OFFSET + 4]);
assert_eq!(u32::from_le_bytes(ml), META_LEN_ABSENT);
let read = FrameHeader::read_from(&buf).expect("read back");
assert!(read.metadata_len.is_none());
}
#[test]
fn frame_header_some_zero_metadata_distinct_from_none() {
let with_empty = FrameHeader {
format_version: FrameFormatVersion::V2,
schema_version: 1,
event_type_len: 0,
metadata_len: Some(0),
};
let mut buf = fresh_buf();
with_empty.write_into(&mut buf);
let mut ml = [0u8; 4];
ml.copy_from_slice(&buf[META_LEN_OFFSET..META_LEN_OFFSET + 4]);
assert_eq!(u32::from_le_bytes(ml), 0);
assert_ne!(u32::from_le_bytes(ml), META_LEN_ABSENT);
let read = FrameHeader::read_from(&buf).expect("read back");
assert_eq!(read.metadata_len, Some(0));
}
#[test]
fn frame_header_read_from_rejects_buffer_below_size() {
for too_short_len in 0..FrameHeader::SIZE {
let buf = vec![0u8; too_short_len];
match FrameHeader::read_from(&buf) {
Err(DecodeError::ValueTooShort { min, actual }) => {
assert_eq!(min, FrameHeader::SIZE);
assert_eq!(actual, too_short_len);
}
other => panic!("expected ValueTooShort for len={too_short_len}, got {other:?}"),
}
}
}
#[test]
fn frame_header_read_from_accepts_exactly_size() {
let mut buf = vec![0u8; FrameHeader::SIZE];
buf[VERSION_OFFSET] = 2;
let header = FrameHeader::read_from(&buf).expect("accepts at SIZE");
assert_eq!(header.format_version, FrameFormatVersion::V2);
assert_eq!(header.schema_version, 0);
assert_eq!(header.event_type_len, 0);
assert_eq!(header.metadata_len, Some(0));
}
proptest! {
#[test]
fn frame_header_round_trip(
schema_version in u32_strategy(),
et_raw in u16_strategy(),
meta_choice in 0u32..4,
) {
let metadata_len = match meta_choice {
0 => None,
1 => Some(0u32),
2 => Some(u32::MAX - 2),
_ => Some((u32::MAX - 1) / 2),
};
let original = FrameHeader {
format_version: FrameFormatVersion::V2,
schema_version,
event_type_len: et_raw,
metadata_len,
};
let mut buf = fresh_buf();
original.write_into(&mut buf);
prop_assert_eq!(buf.len(), FrameHeader::SIZE);
let read = FrameHeader::read_from(&buf).expect("round-trip read");
prop_assert_eq!(read.format_version, original.format_version);
prop_assert_eq!(read.schema_version, original.schema_version);
prop_assert_eq!(read.event_type_len, original.event_type_len);
prop_assert_eq!(read.metadata_len, original.metadata_len);
}
}
#[test]
fn layout_concrete_no_metadata_example() {
let layout = FrameLayout::compute_from_validated_lengths(2, None, 1).expect("ok");
assert_eq!(layout.padding, 3);
assert_eq!(layout.event_type, 11..13);
assert_eq!(layout.metadata, None);
assert_eq!(layout.payload, 16..17);
assert_eq!(layout.total, 17);
}
#[test]
fn layout_concrete_with_metadata_example() {
let layout = FrameLayout::compute_from_validated_lengths(2, Some(3), 4).expect("ok");
assert_eq!(layout.event_type, 11..13);
assert_eq!(layout.metadata, Some(13..16));
assert_eq!(layout.padding, 0);
assert_eq!(layout.payload, 16..20);
assert_eq!(layout.total, 20);
}
proptest! {
#[test]
fn layout_structural_invariants(
et_len_raw in frame_body_length(),
meta in prop::option::of(frame_body_length()),
payload_len in frame_body_length(),
) {
let et_len = et_len_raw.min(MAX_EVENT_TYPE_LEN);
let meta_capped = meta.map(|n| n.min(MAX_METADATA_LEN));
let payload_len_capped = payload_len.min(MAX_PAYLOAD_LEN);
let layout = FrameLayout::compute_from_validated_lengths(
et_len,
meta_capped,
payload_len_capped,
).expect("bounded inputs compute");
prop_assert_eq!(
usize::try_from(layout.event_type.start).unwrap(),
HEADER_FIXED_SIZE,
);
prop_assert_eq!(
(layout.event_type.end - layout.event_type.start) as usize,
et_len,
);
match (meta_capped, layout.metadata.clone()) {
(None, None) => {},
(Some(meta_len), Some(range)) => {
prop_assert_eq!((range.end - range.start) as usize, meta_len);
}
_ => prop_assert!(false, "metadata Option mismatch between input and layout"),
}
prop_assert_eq!(
(layout.payload.end - layout.payload.start) as usize,
payload_len_capped,
);
if let Some(m) = layout.metadata.clone() {
prop_assert!(layout.event_type.end <= m.start);
prop_assert!(m.end <= layout.payload.start);
} else {
prop_assert!(layout.event_type.end <= layout.payload.start);
}
let payload_start = usize::try_from(layout.payload.start).unwrap();
prop_assert!(payload_start.is_multiple_of(PAYLOAD_ALIGN));
prop_assert!(layout.padding < PAYLOAD_ALIGN);
prop_assert_eq!(layout.total, usize::try_from(layout.payload.end).unwrap());
let body_total = et_len
+ meta_capped.unwrap_or(0)
+ layout.padding
+ payload_len_capped;
prop_assert_eq!(layout.total, HEADER_FIXED_SIZE + body_total);
}
}
#[test]
fn plan_then_execute_matches_encode_frame_concrete() {
let sv = SchemaVersion::from_u32(2).expect("nonzero");
let et_v = et("Evt");
let pl_v = pl(b"payload");
let md_v = md(b"meta");
let one_shot = encode_frame(sv, &et_v, &pl_v, Some(&md_v)).expect("ok");
let staged = execute(plan(sv, &et_v, &pl_v, Some(&md_v)).expect("plan ok"));
assert_eq!(one_shot.value.as_ref(), staged.value.as_ref());
assert_eq!(one_shot.offsets.event_type, staged.offsets.event_type);
assert_eq!(one_shot.offsets.metadata, staged.offsets.metadata);
assert_eq!(one_shot.offsets.payload, staged.offsets.payload);
}
#[test]
fn execute_buffer_length_equals_layout_total() {
let cases: Vec<(EventType, Option<Metadata>, Payload)> = vec![
(et(""), None, pl(b"")),
(et("X"), None, pl(b"")),
(et("Evt"), Some(md(b"meta")), pl(b"payload")),
(et("LongerType"), Some(md(b"x")), pl(b"x")),
];
for (et_v, md_v, pl_v) in cases {
let p = plan(sv1(), &et_v, &pl_v, md_v.as_ref()).expect("plan ok");
let total = p.layout.total;
let frame = execute(p);
assert_eq!(frame.value.len(), total);
}
}
#[test]
fn execute_padding_bytes_are_zero() {
let frame = encode_frame(sv1(), &et("x"), &pl(b"payload"), None).expect("ok");
let pad_start = usize::try_from(frame.offsets.event_type.end).unwrap();
let pad_end = usize::try_from(frame.offsets.payload.start).unwrap();
assert!(pad_end > pad_start, "expected at least one padding byte");
for (i, byte) in frame.value[pad_start..pad_end].iter().enumerate() {
assert_eq!(
*byte,
0,
"padding byte at offset {} is {:#x}",
pad_start + i,
byte
);
}
}
proptest! {
#[test]
fn plan_execute_equals_encode_frame(
(schema_version, event_type, metadata, payload) in valid_frame_inputs(),
) {
let et_v = et(&event_type);
let pl_v = pl(&payload);
let md_v = metadata.as_deref().map(md);
let one_shot = encode_frame(
schema_version, &et_v, &pl_v, md_v.as_ref(),
).expect("valid inputs encode");
let staged = execute(
plan(schema_version, &et_v, &pl_v, md_v.as_ref())
.expect("valid inputs plan"),
);
prop_assert_eq!(one_shot.value.as_ref(), staged.value.as_ref());
prop_assert_eq!(one_shot.offsets.event_type, staged.offsets.event_type);
prop_assert_eq!(one_shot.offsets.metadata, staged.offsets.metadata);
prop_assert_eq!(one_shot.offsets.payload, staged.offsets.payload);
}
#[test]
fn execute_invariants(
(schema_version, event_type, metadata, payload) in valid_frame_inputs(),
) {
let et_v = et(&event_type);
let pl_v = pl(&payload);
let md_v = metadata.as_deref().map(md);
let p = plan(schema_version, &et_v, &pl_v, md_v.as_ref())
.expect("valid inputs plan");
let layout_total = p.layout.total;
let event_type_range = p.layout.event_type.clone();
let metadata_range = p.layout.metadata.clone();
let payload_range = p.layout.payload.clone();
let frame = execute(p);
prop_assert_eq!(frame.value.len(), layout_total);
let payload_slice_start = usize::try_from(payload_range.start).unwrap();
let ptr = frame.value[payload_slice_start..].as_ptr().addr();
prop_assert!(ptr.is_multiple_of(PAYLOAD_ALIGN));
let et_start = usize::try_from(event_type_range.start).unwrap();
let et_end = usize::try_from(event_type_range.end).unwrap();
prop_assert_eq!(&frame.value[et_start..et_end], event_type.as_bytes());
if let (Some(range), Some(meta)) = (metadata_range.clone(), metadata.as_deref()) {
let s = usize::try_from(range.start).unwrap();
let e = usize::try_from(range.end).unwrap();
prop_assert_eq!(&frame.value[s..e], meta);
}
let p_start = usize::try_from(payload_range.start).unwrap();
let p_end = usize::try_from(payload_range.end).unwrap();
prop_assert_eq!(&frame.value[p_start..p_end], payload.as_slice());
let pad_start = metadata_range
.as_ref()
.map_or(et_end, |r| usize::try_from(r.end).unwrap());
for byte in &frame.value[pad_start..p_start] {
prop_assert_eq!(*byte, 0u8);
}
}
}
#[test]
fn encode_frame_accepts_value_newtypes() {
let et_v = EventType::from_static_str("UserCreated");
let payload = Payload::from_bytes(Bytes::from_static(b"hello")).expect("valid");
let metadata = Metadata::from_bytes(Bytes::from_static(b"m")).expect("valid");
let sv = SchemaVersion::INITIAL;
let frame = encode_frame(sv, &et_v, &payload, Some(&metadata)).expect("valid frame");
let decoded = decode_frame(&frame.value).expect("decodes");
assert_eq!(decoded.schema_version, sv);
}
#[test]
fn decode_frame_rejects_corrupt_schema_version_zero() {
let et_v = EventType::from_static_str("X");
let payload = Payload::from_bytes(Bytes::from_static(b"p")).expect("valid");
let sv_one = SchemaVersion::INITIAL;
let frame = encode_frame(sv_one, &et_v, &payload, None).expect("valid frame for tamper");
let mut bytes_vec = frame.value.to_vec();
bytes_vec[SCHEMA_VERSION_OFFSET..EVENT_TYPE_LEN_OFFSET].fill(0);
let tampered = Bytes::from(bytes_vec);
let err = decode_frame(&tampered).expect_err("schema_version=0 on wire rejected");
assert!(matches!(err, DecodeError::CorruptSchemaVersion));
}
#[test]
fn decode_rejects_every_unknown_version_byte() {
let frame = encode_frame(sv1(), &et("Evt"), &pl(b"payload"), Some(&md(b"m")))
.expect("valid frame for tamper base");
for bad in (0u8..=u8::MAX).filter(|b| *b != 1 && *b != 2) {
let mut bytes_vec = frame.value.to_vec();
bytes_vec[VERSION_OFFSET] = bad;
let tampered = Bytes::from(bytes_vec);
match decode_frame(&tampered) {
Err(DecodeError::UnsupportedFrameVersion { version }) => {
assert_eq!(version, bad);
}
other => panic!("version byte {bad} should be rejected, got {other:?}"),
}
}
}
#[test]
fn decode_reads_a_v1_frame_dropping_its_global_seq() {
let event_type = b"Created";
let payload = b"data-bytes";
let mut buf = Vec::new();
buf.push(1u8); buf.extend_from_slice(&999u64.to_le_bytes()); buf.extend_from_slice(&7u32.to_le_bytes()); buf.extend_from_slice(&u16::try_from(event_type.len()).unwrap().to_le_bytes()); buf.extend_from_slice(&META_LEN_ABSENT.to_le_bytes()); buf.extend_from_slice(event_type);
let post_et = buf.len();
buf.resize(post_et + align_padding(post_et, PAYLOAD_ALIGN), 0u8);
buf.extend_from_slice(payload);
let decoded = decode_frame(&buf).expect("a well-formed V1 frame must still decode");
assert_eq!(decoded.schema_version.get(), 7);
let et = &buf
[decoded.offsets.event_type.start as usize..decoded.offsets.event_type.end as usize];
assert_eq!(et, event_type);
let pl = &buf[decoded.offsets.payload.start as usize..decoded.offsets.payload.end as usize];
assert_eq!(pl, payload);
}
#[test]
fn decode_empty_buffer_is_too_short_not_version_error() {
match decode_frame(&[]) {
Err(DecodeError::ValueTooShort { min, actual }) => {
assert_eq!(min, HEADER_FIXED_SIZE);
assert_eq!(actual, 0);
}
other => panic!("empty buffer should be ValueTooShort, got {other:?}"),
}
}
#[test]
fn corrupt_version_byte_surfaces_unsupported_not_panic() {
let frame = encode_frame(sv1(), &et("X"), &pl(b"p"), None).expect("encode");
let mut bytes_vec = frame.value.to_vec();
bytes_vec[VERSION_OFFSET] = 0xFF;
let tampered = Bytes::from(bytes_vec);
let err = decode_frame(&tampered).expect_err("corrupt version rejected");
assert!(matches!(
err,
DecodeError::UnsupportedFrameVersion { version: 0xFF }
));
}
}