use alloc::vec::Vec;
use core::convert::Infallible;
use core::ops::Range;
use bytes::Bytes;
use futures::Stream;
use futures::StreamExt as _;
use minicbor::Decoder;
use minicbor::Encoder;
use minicbor::data::Type;
use thiserror::Error;
use crate::envelope::PersistedEnvelope;
use crate::import::{ImportBlock, StreamSection};
use crate::value::SchemaVersion;
use mnesis::Version;
const MAGIC: &[u8] = b"nxch";
const FORMAT_VERSION: u32 = 1;
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum ChunkError {
#[error("malformed chunk: {0}")]
Malformed(&'static str),
}
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum WriteError<E> {
#[error("chunk write failed: {0}")]
Encode(#[from] minicbor::encode::Error<E>),
}
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum SectionError<E, R> {
#[error(transparent)]
Write(#[from] WriteError<E>),
#[error("event stream read failed: {0}")]
Read(#[source] R),
}
fn widen_encode_err<E>(e: minicbor::encode::Error<Infallible>) -> WriteError<E> {
WriteError::Encode(minicbor::encode::Error::message(e))
}
#[derive(Debug, Clone)]
pub struct ChunkHeader {
pub format_version: u32,
pub origin: Option<Bytes>,
}
#[derive(minicbor::Encode, minicbor::Decode)]
#[cbor(map)]
struct HeaderRepr<'a> {
#[n(0)]
#[cbor(with = "minicbor::bytes")]
magic: &'a [u8],
#[n(1)]
format_version: u32,
#[n(2)]
#[cbor(with = "minicbor::bytes")]
origin: Option<&'a [u8]>,
}
fn validate_header(magic: &[u8], format_version: u32) -> Result<(), ChunkError> {
if magic != MAGIC {
return Err(ChunkError::Malformed("bad magic"));
}
if format_version != FORMAT_VERSION {
return Err(ChunkError::Malformed("unknown format version"));
}
Ok(())
}
pub fn decode_header(bytes: &[u8]) -> Result<ChunkHeader, ChunkError> {
let repr: HeaderRepr =
minicbor::decode(bytes).map_err(|_| ChunkError::Malformed("unreadable header"))?;
validate_header(repr.magic, repr.format_version)?;
Ok(ChunkHeader {
format_version: repr.format_version,
origin: repr.origin.map(Bytes::copy_from_slice),
})
}
#[derive(Debug)]
pub struct ChunkWriter<W> {
enc: Encoder<W>,
}
impl<W: minicbor::encode::Write> ChunkWriter<W> {
pub fn new(sink: W, origin: Option<&[u8]>) -> Result<Self, WriteError<W::Error>> {
let mut enc = Encoder::new(sink);
enc.encode(HeaderRepr {
magic: MAGIC,
format_version: FORMAT_VERSION,
origin,
})?;
Ok(Self { enc })
}
pub fn section(
&mut self,
stream_id: &[u8],
) -> Result<SectionWriter<'_, W>, WriteError<W::Error>> {
self.enc.encode(HeadingRepr { stream_id })?;
Ok(SectionWriter { enc: &mut self.enc })
}
pub fn into_sink(self) -> W {
self.enc.into_writer()
}
}
#[derive(Debug)]
pub struct SectionWriter<'a, W> {
enc: &'a mut Encoder<W>,
}
impl<W: minicbor::encode::Write> SectionWriter<'_, W> {
pub async fn try_extend<S, R>(
&mut self,
events: S,
) -> Result<&mut Self, SectionError<W::Error, R>>
where
S: Stream<Item = Result<PersistedEnvelope, R>>,
{
futures::pin_mut!(events);
while let Some(item) = events.next().await {
let env = item.map_err(SectionError::Read)?;
self.block(&env).map_err(SectionError::Write)?;
}
Ok(self)
}
pub fn block(&mut self, event: &PersistedEnvelope) -> Result<&mut Self, WriteError<W::Error>> {
let body = BodyRepr {
version: event.version().as_u64(),
schema_version: event.schema_version(),
event_type: event.event_type(),
metadata: event.metadata(),
payload: event.payload(),
};
let body_bytes = minicbor::to_vec(&body).map_err(widen_encode_err)?;
self.enc.encode(BlockRepr {
crc: crc32c::crc32c(&body_bytes),
body: &body_bytes,
})?;
Ok(self)
}
}
#[derive(minicbor::Encode, minicbor::Decode)]
#[cbor(map)]
struct HeadingRepr<'a> {
#[n(0)]
#[cbor(with = "minicbor::bytes")]
stream_id: &'a [u8],
}
#[derive(minicbor::Encode, minicbor::Decode)]
#[cbor(map)]
struct BodyRepr<'a> {
#[n(0)]
version: u64,
#[n(1)]
schema_version: u32,
#[n(2)]
event_type: &'a str,
#[n(3)]
#[cbor(with = "minicbor::bytes")]
metadata: Option<&'a [u8]>,
#[n(4)]
#[cbor(with = "minicbor::bytes")]
payload: &'a [u8],
}
#[derive(minicbor::Encode)]
#[cbor(array)]
struct BlockRepr<'a> {
#[n(0)]
crc: u32,
#[n(1)]
#[cbor(with = "minicbor::bytes")]
body: &'a [u8],
}
fn decode_block(d: &mut Decoder<'_>) -> Result<Option<ImportBlock>, &'static str> {
match d.array() {
Ok(Some(2)) => {}
Ok(Some(_)) => return Err("block array must have exactly 2 elements"),
Ok(None) => return Err("indefinite-length block array"),
Err(e) if e.is_end_of_input() => return Ok(None),
Err(_) => return Err("malformed block array"),
}
let crc = match d.u32() {
Ok(c) => c,
Err(e) if e.is_end_of_input() => return Ok(None),
Err(_) => return Err("malformed block crc"),
};
let body: &[u8] = match d.bytes() {
Ok(b) => b,
Err(e) if e.is_end_of_input() => return Ok(None),
Err(_) => return Err("malformed block body bytes"),
};
if crc32c::crc32c(body) != crc {
return Ok(Some(ImportBlock::Corrupt));
}
let parsed: BodyRepr = minicbor::decode(body).map_err(|_| "crc-valid body failed to decode")?;
let envelope = reconstruct(&parsed).ok_or("crc-valid body has invalid fields")?;
Ok(Some(ImportBlock::Event(envelope)))
}
fn reconstruct(body: &BodyRepr<'_>) -> Option<PersistedEnvelope> {
let version = Version::new(body.version)?;
let schema = SchemaVersion::from_u32(body.schema_version).ok()?;
let event_type = body.event_type.as_bytes();
let et_end = u32::try_from(event_type.len()).ok()?;
let mut buf = Vec::with_capacity(
event_type.len() + body.metadata.map_or(0, <[u8]>::len) + body.payload.len(),
);
buf.extend_from_slice(event_type);
let et_range: Range<u32> = 0..et_end;
let meta_range = match body.metadata {
Some(m) => {
let start = u32::try_from(buf.len()).ok()?;
buf.extend_from_slice(m);
let end = u32::try_from(buf.len()).ok()?;
Some(start..end)
}
None => None,
};
let pl_start = u32::try_from(buf.len()).ok()?;
buf.extend_from_slice(body.payload);
let pl_end = u32::try_from(buf.len()).ok()?;
PersistedEnvelope::try_new(
version,
Bytes::from(buf),
schema,
et_range,
pl_start..pl_end,
meta_range,
)
.ok()
}
pub fn decode_chunk(bytes: &[u8]) -> Result<Vec<StreamSection>, ChunkError> {
let mut d = Decoder::new(bytes);
let header: HeaderRepr = d
.decode()
.map_err(|_| ChunkError::Malformed("unreadable header"))?;
validate_header(header.magic, header.format_version)?;
let mut sections: Vec<StreamSection> = Vec::new();
while d.position() < bytes.len() {
match d.datatype() {
Ok(Type::Map) => match d.decode::<HeadingRepr>() {
Ok(heading) => sections.push(StreamSection {
origin: Bytes::copy_from_slice(heading.stream_id),
blocks: Vec::new(),
}),
Err(e) if e.is_end_of_input() => break,
Err(_) => return Err(ChunkError::Malformed("malformed section heading")),
},
Ok(Type::Array) => match decode_block(&mut d) {
Ok(Some(block)) => match sections.last_mut() {
Some(section) => section.blocks.push(block),
None => return Err(ChunkError::Malformed("block before section heading")),
},
Ok(None) => break,
Err(hint) => return Err(ChunkError::Malformed(hint)),
},
Ok(_) => return Err(ChunkError::Malformed("unexpected item type")),
Err(e) if e.is_end_of_input() => break,
Err(_) => return Err(ChunkError::Malformed("decode error")),
}
}
Ok(sections)
}
#[cfg(test)]
#[allow(
clippy::unwrap_used,
clippy::expect_used,
clippy::panic,
reason = "test code asserts exact values"
)]
mod tests {
use super::*;
#[test]
fn writer_multi_section_round_trips() {
let mut w = ChunkWriter::new(Vec::new(), Some(b"dev")).expect("new");
{
let mut s = w.section(b"task-1").expect("section");
s.block(&persisted(1, 1, "E", None, b"a1")).expect("block");
s.block(&persisted(2, 1, "E", Some(b"m"), b"a2"))
.expect("block");
}
{
let mut s = w.section(b"task-2").expect("section");
s.block(&persisted(1, 2, "E", None, b"b1")).expect("block");
}
let chunk = Bytes::from(w.into_sink());
let sections = decode_chunk(&chunk).expect("decode");
assert_eq!(sections.len(), 2);
assert_eq!(sections[0].origin.as_ref(), b"task-1");
assert_eq!(sections[0].blocks.len(), 2);
assert_eq!(sections[1].origin.as_ref(), b"task-2");
match (§ions[0].blocks[1], §ions[1].blocks[0]) {
(ImportBlock::Event(a2), ImportBlock::Event(b1)) => {
assert_eq!(a2.metadata(), Some(b"m".as_slice()));
assert_eq!(a2.payload(), b"a2");
assert_eq!(b1.schema_version(), 2);
assert_eq!(b1.payload(), b"b1");
}
_ => panic!("expected Event blocks"),
}
}
#[test]
fn writer_empty_section_then_stream() {
let mut w = ChunkWriter::new(Vec::new(), None).expect("new");
{
let _s = w.section(b"empty").expect("section");
}
{
let mut s = w.section(b"real").expect("section");
s.block(&persisted(1, 1, "E", None, b"x")).expect("block");
}
let sections = decode_chunk(&Bytes::from(w.into_sink())).expect("decode");
assert_eq!(sections.len(), 2);
assert_eq!(sections[0].origin.as_ref(), b"empty");
assert!(sections[0].blocks.is_empty());
assert_eq!(sections[1].blocks.len(), 1);
}
#[test]
fn writer_header_decodes() {
let w = ChunkWriter::new(Vec::new(), Some(b"dev-1")).expect("new");
let bytes = Bytes::from(w.into_sink());
let header = decode_header(&bytes).expect("decode header");
assert_eq!(header.format_version, 1);
assert_eq!(header.origin.as_deref(), Some(b"dev-1".as_slice()));
}
#[test]
fn write_error_displays_and_is_error() {
fn assert_error<E: std::error::Error>() {}
assert_error::<WriteError<core::convert::Infallible>>();
assert_error::<SectionError<core::convert::Infallible, std::io::Error>>();
}
fn persisted(
version: u64,
schema: u32,
event_type: &str,
metadata: Option<&[u8]>,
payload: &[u8],
) -> PersistedEnvelope {
let mut buf = Vec::new();
buf.extend_from_slice(event_type.as_bytes());
let et_end = u32::try_from(buf.len()).expect("fits");
let meta_range = metadata.map(|m| {
let start = u32::try_from(buf.len()).expect("fits");
buf.extend_from_slice(m);
start..u32::try_from(buf.len()).expect("fits")
});
let pl_start = u32::try_from(buf.len()).expect("fits");
buf.extend_from_slice(payload);
let pl_end = u32::try_from(buf.len()).expect("fits");
PersistedEnvelope::try_new(
Version::new(version).expect("nonzero"),
Bytes::from(buf),
SchemaVersion::from_u32(schema).expect("nonzero"),
0..et_end,
pl_start..pl_end,
meta_range,
)
.expect("valid persisted")
}
fn decode_one_block(bytes: &[u8]) -> ImportBlock {
let mut d = Decoder::new(bytes);
decode_block(&mut d)
.expect("not malformed")
.expect("not torn")
}
fn header_bytes(origin: Option<&[u8]>) -> Vec<u8> {
ChunkWriter::new(Vec::new(), origin)
.expect("header writer")
.into_sink()
}
fn heading_bytes(stream_id: &[u8]) -> Vec<u8> {
let mut w = ChunkWriter::new(Vec::new(), None).expect("writer");
w.section(stream_id).expect("section");
let full = w.into_sink();
full[header_bytes(None).len()..].to_vec()
}
fn block_bytes(event: &PersistedEnvelope) -> Vec<u8> {
let mut w = ChunkWriter::new(Vec::new(), None).expect("writer");
{
let mut s = w.section(b"k").expect("section");
s.block(event).expect("block");
}
let full = w.into_sink();
let prefix = {
let mut w2 = ChunkWriter::new(Vec::new(), None).expect("writer");
w2.section(b"k").expect("section");
w2.into_sink().len()
};
full[prefix..].to_vec()
}
#[test]
fn block_round_trips_all_fields() {
let event = persisted(7, 3, "AccountOpened", Some(b"hlc=42"), b"balance:100");
let bytes = block_bytes(&event);
match decode_one_block(&bytes) {
ImportBlock::Event(got) => {
assert_eq!(got.version().as_u64(), 7);
assert_eq!(got.schema_version(), 3);
assert_eq!(got.event_type(), "AccountOpened");
assert_eq!(got.metadata(), Some(b"hlc=42".as_slice()));
assert_eq!(got.payload(), b"balance:100");
}
ImportBlock::Corrupt => panic!("expected Event, got Corrupt"),
}
}
#[test]
fn block_round_trips_without_metadata_and_empty_payload() {
let event = persisted(1, 1, "E", None, b"");
let bytes = block_bytes(&event);
match decode_one_block(&bytes) {
ImportBlock::Event(got) => {
assert_eq!(got.metadata(), None);
assert_eq!(got.payload(), b"");
assert_eq!(got.version().as_u64(), 1);
}
ImportBlock::Corrupt => panic!("expected Event"),
}
}
#[test]
fn block_with_flipped_body_byte_is_corrupt() {
let event = persisted(2, 1, "E", None, b"hello");
let mut v = block_bytes(&event);
let last = v.len() - 1;
v[last] ^= 0xFF;
assert!(matches!(decode_one_block(&v), ImportBlock::Corrupt));
}
#[test]
fn header_round_trips_without_origin() {
let bytes = header_bytes(None);
let header = decode_header(&bytes).expect("decode");
assert_eq!(header.format_version, 1);
assert_eq!(header.origin, None);
}
#[test]
fn header_round_trips_with_origin() {
let bytes = header_bytes(Some(b"phone-7"));
let header = decode_header(&bytes).expect("decode");
assert_eq!(header.format_version, 1);
assert_eq!(header.origin.as_deref(), Some(b"phone-7".as_slice()));
}
#[test]
fn header_rejects_bad_magic() {
let mut v = header_bytes(None);
let pos = v.iter().position(|&b| b == b'n').expect("magic present");
v[pos] = b'X';
let err = decode_header(&v).expect_err("bad magic rejected");
assert!(matches!(err, ChunkError::Malformed("bad magic")));
}
#[test]
fn header_rejects_truncated() {
let bytes = header_bytes(Some(b"x"));
let err = decode_header(&bytes[..bytes.len() / 2]).expect_err("truncated rejected");
assert!(matches!(err, ChunkError::Malformed(_)));
}
fn encode_chunk(origin: Option<&[u8]>, streams: &[(&[u8], Vec<PersistedEnvelope>)]) -> Bytes {
let mut w = ChunkWriter::new(Vec::new(), origin).expect("writer");
for (stream_id, events) in streams {
let mut s = w.section(stream_id).expect("section");
for e in events {
s.block(e).expect("block");
}
}
Bytes::from(w.into_sink())
}
#[test]
fn decode_chunk_round_trips_multi_stream() {
let a = vec![
persisted(1, 1, "E", None, b"a1"),
persisted(2, 1, "E", Some(b"m"), b"a2"),
];
let b = vec![persisted(1, 2, "E", None, b"b1")];
let chunk = encode_chunk(
Some(b"dev-1"),
&[(b"task-1".as_slice(), a), (b"task-2".as_slice(), b)],
);
let sections = decode_chunk(&chunk).expect("decode");
assert_eq!(sections.len(), 2);
assert_eq!(sections[0].origin.as_ref(), b"task-1");
assert_eq!(sections[0].blocks.len(), 2);
match (§ions[0].blocks[0], §ions[0].blocks[1]) {
(ImportBlock::Event(e1), ImportBlock::Event(e2)) => {
assert_eq!(e1.version().as_u64(), 1);
assert_eq!(e1.payload(), b"a1");
assert_eq!(e2.metadata(), Some(b"m".as_slice()));
assert_eq!(e2.payload(), b"a2");
}
_ => panic!("expected two Event blocks"),
}
assert_eq!(sections[1].origin.as_ref(), b"task-2");
match §ions[1].blocks[0] {
ImportBlock::Event(e) => {
assert_eq!(e.schema_version(), 2);
assert_eq!(e.payload(), b"b1");
}
ImportBlock::Corrupt => panic!("expected Event"),
}
}
#[test]
fn decode_header_only_chunk_is_empty_vec() {
let chunk = header_bytes(None);
let sections = decode_chunk(&chunk).expect("decode");
assert!(sections.is_empty());
}
#[test]
fn decode_block_before_heading_is_malformed() {
let mut chunk = header_bytes(None);
let event = persisted(1, 1, "E", None, b"x");
chunk.extend_from_slice(&block_bytes(&event));
let err = decode_chunk(&chunk).expect_err("block before heading");
assert!(matches!(
err,
ChunkError::Malformed("block before section heading")
));
}
#[test]
fn decode_empty_section_then_stream() {
let chunk = encode_chunk(
None,
&[
(b"empty".as_slice(), vec![]),
(b"real".as_slice(), vec![persisted(1, 1, "E", None, b"r")]),
],
);
let sections = decode_chunk(&chunk).expect("decode");
assert_eq!(sections.len(), 2);
assert!(sections[0].blocks.is_empty());
assert_eq!(sections[0].origin.as_ref(), b"empty");
assert_eq!(sections[1].blocks.len(), 1);
}
#[test]
fn flipped_body_byte_in_chunk_decodes_to_corrupt_block() {
let chunk = encode_chunk(
None,
&[(b"s".as_slice(), vec![persisted(1, 1, "E", None, b"hello")])],
);
let mut v = chunk.to_vec();
let last = v.len() - 1;
v[last] ^= 0xFF;
let sections = decode_chunk(&v).expect("framing intact");
assert_eq!(sections.len(), 1);
assert!(matches!(sections[0].blocks[0], ImportBlock::Corrupt));
}
#[test]
fn bad_magic_chunk_is_malformed() {
let mut v = encode_chunk(
None,
&[(b"s".as_slice(), vec![persisted(1, 1, "E", None, b"x")])],
)
.to_vec();
let pos = v.iter().position(|&b| b == b'n').expect("magic");
v[pos] = b'Z';
assert!(matches!(
decode_chunk(&v),
Err(ChunkError::Malformed("bad magic"))
));
}
#[test]
fn unknown_format_version_is_malformed() {
let repr = HeaderRepr {
magic: MAGIC,
format_version: 2,
origin: None,
};
let bytes = minicbor::to_vec(&repr).expect("encode");
assert!(matches!(
decode_chunk(&bytes),
Err(ChunkError::Malformed("unknown format version"))
));
}
#[test]
fn unexpected_top_level_item_is_malformed() {
let mut v = header_bytes(None);
v.push(0x01); assert!(matches!(
decode_chunk(&v),
Err(ChunkError::Malformed("unexpected item type"))
));
}
#[test]
fn crc_valid_but_body_invalid_is_malformed() {
let mut body = Vec::new();
{
let mut e = minicbor::Encoder::new(&mut body);
e.map(4)
.expect("map")
.u32(0)
.expect("k0")
.u64(0)
.expect("v0")
.u32(1)
.expect("k1")
.u32(1)
.expect("v1")
.u32(2)
.expect("k2")
.str("E")
.expect("v2")
.u32(4)
.expect("k4")
.bytes(b"")
.expect("v4");
}
let block = BlockRepr {
crc: crc32c::crc32c(&body),
body: &body,
};
let mut chunk = header_bytes(None);
chunk.extend_from_slice(&heading_bytes(b"s"));
chunk.extend_from_slice(&minicbor::to_vec(&block).expect("block"));
assert!(matches!(
decode_chunk(&chunk),
Err(ChunkError::Malformed(_))
));
}
fn prefix_len_after_n_blocks(
stream_id: &[u8],
events: &[PersistedEnvelope],
n: usize,
) -> usize {
let mut w = ChunkWriter::new(Vec::new(), None).expect("writer");
{
let mut s = w.section(stream_id).expect("section");
for e in events.iter().take(n) {
s.block(e).expect("block");
}
}
w.into_sink().len()
}
#[test]
fn every_block_boundary_prefix_is_valid() {
let events: Vec<_> = (1..=4).map(|v| persisted(v, 1, "E", None, b"p")).collect();
let chunk = encode_chunk(None, &[(b"s".as_slice(), events.clone())]);
for n in 0..=events.len() {
let cut = prefix_len_after_n_blocks(b"s", &events, n);
let sections = decode_chunk(&chunk[..cut]).expect("prefix valid");
let got = sections.first().map_or(0, |s| s.blocks.len());
assert_eq!(got, n, "prefix after {n} blocks must decode to {n} blocks");
}
}
#[test]
fn torn_final_block_is_dropped_earlier_survive() {
let events: Vec<_> = (1..=3)
.map(|v| persisted(v, 1, "E", None, b"payload"))
.collect();
let chunk = encode_chunk(None, &[(b"s".as_slice(), events)]);
let sections = decode_chunk(&chunk[..chunk.len() - 1]).expect("valid prefix");
assert_eq!(sections[0].blocks.len(), 2, "torn 3rd block dropped");
assert!(matches!(sections[0].blocks[0], ImportBlock::Event(_)));
}
#[test]
fn empty_input_is_malformed_header() {
assert!(matches!(decode_chunk(&[]), Err(ChunkError::Malformed(_))));
}
fn hex_of(bytes: &[u8]) -> String {
bytes
.iter()
.map(|b| format!("{b:02x}"))
.collect::<Vec<_>>()
.join(" ")
}
#[test]
fn golden_header_bytes() {
let bytes = header_bytes(Some(b"dev"));
insta::assert_snapshot!("header_with_origin_hex", hex_of(&bytes));
}
#[test]
fn golden_block_bytes() {
let event = persisted(1, 1, "E", None, b"hi");
let bytes = block_bytes(&event);
insta::assert_snapshot!("block_v1_hex", hex_of(&bytes));
}
use proptest::prelude::*;
fn event_strategy() -> impl Strategy<Value = (u64, u32, String, Option<Vec<u8>>, Vec<u8>)> {
(
prop_oneof![
Just(1u64),
Just(2),
Just(u64::MAX - 1),
Just(u64::MAX),
1u64..1000
],
prop_oneof![Just(1u32), Just(7), Just(u32::MAX), 1u32..100],
prop_oneof![
Just(String::new()),
Just("E".to_owned()),
"[A-Za-z]{0,40}",
Just("z".repeat(300)),
],
prop_oneof![
Just(None),
proptest::collection::vec(any::<u8>(), 1..32).prop_map(Some),
proptest::collection::vec(any::<u8>(), 256..400).prop_map(Some),
],
prop_oneof![
proptest::collection::vec(any::<u8>(), 0..64),
proptest::collection::vec(any::<u8>(), 256..512),
],
)
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(128))]
#[test]
fn vopr_chunk_round_trips(
origin in prop_oneof![
Just(None),
proptest::collection::vec(any::<u8>(), 1..16).prop_map(Some)
],
streams in proptest::collection::vec(
(proptest::collection::vec(any::<u8>(), 0..12),
proptest::collection::vec(event_strategy(), 0..5)),
0..4,
),
) {
let built: Vec<(Vec<u8>, Vec<PersistedEnvelope>)> = streams
.iter()
.map(|(sid, evs)| {
let events = evs.iter().map(|(v, sc, et, md, pl)| {
persisted(*v, *sc, et, md.as_deref(), pl)
}).collect();
(sid.clone(), events)
})
.collect();
let refs: Vec<(&[u8], Vec<PersistedEnvelope>)> =
built.iter().map(|(s, e)| (s.as_slice(), e.clone())).collect();
let chunk = encode_chunk(origin.as_deref(), &refs);
let sections = decode_chunk(&chunk).expect("decode");
prop_assert_eq!(sections.len(), built.len());
for (section, (sid, events)) in sections.iter().zip(built.iter()) {
prop_assert_eq!(section.origin.as_ref(), sid.as_slice());
prop_assert_eq!(section.blocks.len(), events.len());
for (block, original) in section.blocks.iter().zip(events.iter()) {
match block {
ImportBlock::Event(got) => {
prop_assert_eq!(got.version(), original.version());
prop_assert_eq!(got.schema_version(), original.schema_version());
prop_assert_eq!(got.event_type(), original.event_type());
prop_assert_eq!(got.metadata(), original.metadata());
prop_assert_eq!(got.payload(), original.payload());
}
ImportBlock::Corrupt => prop_assert!(false, "unexpected corrupt"),
}
}
}
}
#[test]
fn vopr_single_byte_body_mutation_is_corrupt_never_silently_wrong(
(ev_v, ev_sc, ev_et, ev_md, ev_pl) in event_strategy(),
flip_pick in any::<prop::sample::Index>(),
) {
let event = persisted(ev_v, ev_sc, &ev_et, ev_md.as_deref(), &ev_pl);
let mut v = block_bytes(&event);
let start = v.len() / 2;
let idx = start + flip_pick.index(v.len() - start);
v[idx] ^= 0xFF;
let mut d = Decoder::new(&v);
match decode_block(&mut d) {
Ok(Some(ImportBlock::Event(got))) => {
prop_assert_eq!(got.payload(), event.payload());
prop_assert_eq!(got.event_type(), event.event_type());
}
Ok(Some(ImportBlock::Corrupt) | None) | Err(_) => {}
}
}
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(512))]
#[test]
fn decode_chunk_never_panics_on_arbitrary_bytes(
bytes in proptest::collection::vec(any::<u8>(), 0..256),
) {
let _ = decode_chunk(&bytes);
}
#[test]
fn decode_chunk_never_panics_on_valid_header_plus_garbage(
garbage in proptest::collection::vec(any::<u8>(), 0..128),
) {
let mut v = header_bytes(None);
v.extend_from_slice(&garbage);
let _ = decode_chunk(&v);
}
}
#[test]
fn body_with_unknown_extra_key_still_decodes_to_event() {
let mut body = Vec::new();
{
let mut e = minicbor::Encoder::new(&mut body);
e.map(5)
.expect("map")
.u32(0)
.expect("k0")
.u64(3)
.expect("v0")
.u32(1)
.expect("k1")
.u32(1)
.expect("v1")
.u32(2)
.expect("k2")
.str("E")
.expect("v2")
.u32(4)
.expect("k4")
.bytes(b"data")
.expect("v4")
.u32(5)
.expect("k5")
.u32(999)
.expect("v5"); }
let block = BlockRepr {
crc: crc32c::crc32c(&body),
body: &body,
};
let mut chunk = header_bytes(None);
chunk.extend_from_slice(&heading_bytes(b"s"));
chunk.extend_from_slice(&minicbor::to_vec(&block).expect("block"));
let sections = decode_chunk(&chunk).expect("decode");
match §ions[0].blocks[0] {
ImportBlock::Event(e) => {
assert_eq!(e.version().as_u64(), 3);
assert_eq!(e.payload(), b"data");
}
ImportBlock::Corrupt => panic!("unknown key must be skipped, not corrupt"),
}
}
#[test]
fn box_round_trips_every_payload_length_band() {
for size in [0usize, 23, 24, 255, 256, 1024, 65536] {
let payload = vec![0xABu8; size];
let event = persisted(1, 1, "E", None, &payload);
let bytes = block_bytes(&event);
match decode_one_block(&bytes) {
ImportBlock::Event(got) => {
assert_eq!(got.payload().len(), size, "payload size {size} length");
assert_eq!(got.payload(), payload.as_slice(), "payload {size} bytes");
}
ImportBlock::Corrupt => panic!("payload size {size} must round-trip"),
}
}
}
#[test]
fn box_round_trips_event_type_at_2byte_band_and_max() {
for et_len in [256usize, crate::value::MAX_EVENT_TYPE_LEN] {
let et = "a".repeat(et_len);
let event = persisted(1, 1, &et, Some(b"m"), b"p");
let bytes = block_bytes(&event);
match decode_one_block(&bytes) {
ImportBlock::Event(got) => {
assert_eq!(got.event_type().len(), et_len, "event_type len {et_len}");
assert_eq!(got.event_type(), et);
assert_eq!(got.metadata(), Some(b"m".as_slice()));
}
ImportBlock::Corrupt => panic!("event_type len {et_len} must round-trip"),
}
}
}
#[test]
fn box_round_trips_metadata_at_2byte_band() {
let meta = vec![0x07u8; 300];
let event = persisted(1, 1, "E", Some(&meta), b"p");
match decode_one_block(&block_bytes(&event)) {
ImportBlock::Event(got) => assert_eq!(got.metadata(), Some(meta.as_slice())),
ImportBlock::Corrupt => panic!("metadata 2-byte band must round-trip"),
}
}
#[test]
fn decode_block_rejects_indefinite_array() {
let mut d = Decoder::new(&[0x9fu8]);
assert!(matches!(
decode_block(&mut d),
Err("indefinite-length block array")
));
}
#[test]
fn decode_block_rejects_wrong_arity() {
let mut buf = Vec::new();
{
let mut e = minicbor::Encoder::new(&mut buf);
e.array(3).expect("array header");
}
let mut d = Decoder::new(&buf);
assert!(matches!(
decode_block(&mut d),
Err("block array must have exactly 2 elements")
));
}
#[test]
fn indefinite_array_at_top_level_is_malformed() {
let mut chunk = header_bytes(None);
chunk.extend_from_slice(&heading_bytes(b"s"));
chunk.push(0x9f); chunk.push(0xff); assert!(matches!(
decode_chunk(&chunk),
Err(ChunkError::Malformed("unexpected item type"))
));
}
#[test]
fn torn_mid_heading_stops_at_valid_prefix() {
let chunk = encode_chunk(
None,
&[(b"done".as_slice(), vec![persisted(1, 1, "E", None, b"x")])],
);
let mut v = chunk.to_vec();
let heading = heading_bytes(b"truncated-stream-id");
v.extend_from_slice(&heading[..heading.len() - 4]); let sections = decode_chunk(&v).expect("valid prefix");
assert_eq!(sections.len(), 1, "torn heading produces no section");
assert_eq!(sections[0].origin.as_ref(), b"done");
assert_eq!(sections[0].blocks.len(), 1);
}
#[test]
fn header_missing_magic_key_is_malformed() {
let mut bytes = Vec::new();
{
let mut e = minicbor::Encoder::new(&mut bytes);
e.map(1)
.expect("map")
.u32(1)
.expect("k1")
.u32(1)
.expect("v1");
}
assert!(matches!(
decode_header(&bytes),
Err(ChunkError::Malformed(_))
));
assert!(matches!(
decode_chunk(&bytes),
Err(ChunkError::Malformed(_))
));
}
#[tokio::test]
async fn try_extend_drains_stream_into_section() {
let events = vec![
Ok::<_, std::io::Error>(persisted(1, 1, "E", None, b"x1")),
Ok(persisted(2, 1, "E", Some(b"m"), b"x2")),
];
let mut w = ChunkWriter::new(Vec::new(), None).expect("new");
w.section(b"s")
.expect("section")
.try_extend(futures::stream::iter(events))
.await
.expect("extend");
let chunk = Bytes::from(w.into_sink());
let sections = decode_chunk(&chunk).expect("decode");
assert_eq!(sections[0].blocks.len(), 2);
match §ions[0].blocks[1] {
ImportBlock::Event(e) => assert_eq!(e.payload(), b"x2"),
ImportBlock::Corrupt => panic!("expected Event"),
}
}
#[tokio::test]
async fn try_extend_surfaces_read_error_distinctly() {
let boom = std::io::Error::other("boom");
let events = vec![Ok(persisted(1, 1, "E", None, b"x1")), Err(boom)];
let mut w = ChunkWriter::new(Vec::new(), None).expect("new");
let err = w
.section(b"s")
.expect("section")
.try_extend(futures::stream::iter(events))
.await
.expect_err("read error propagates");
match err {
SectionError::Read(io) => assert_eq!(io.to_string(), "boom"),
SectionError::Write(_) => panic!("a read failure must not be a write failure"),
}
}
#[test]
fn crc_valid_body_with_empty_metadata_is_malformed() {
let mut body = Vec::new();
{
let mut e = minicbor::Encoder::new(&mut body);
e.map(5)
.expect("map")
.u32(0)
.expect("k0")
.u64(1)
.expect("v0")
.u32(1)
.expect("k1")
.u32(1)
.expect("v1")
.u32(2)
.expect("k2")
.str("E")
.expect("v2")
.u32(3)
.expect("k3")
.bytes(b"")
.expect("v3") .u32(4)
.expect("k4")
.bytes(b"p")
.expect("v4");
}
let block = BlockRepr {
crc: crc32c::crc32c(&body),
body: &body,
};
let mut chunk = header_bytes(None);
chunk.extend_from_slice(&heading_bytes(b"s"));
chunk.extend_from_slice(&minicbor::to_vec(&block).expect("block"));
assert!(matches!(
decode_chunk(&chunk),
Err(ChunkError::Malformed(_))
));
}
}