use core::iter::{Chain, Once, once};
use core::num::NonZeroUsize;
use core::ops::Range;
use core::slice::Iter;
use bytes::Bytes;
use mnesis::{DomainEvent, Version};
use thiserror::Error;
use crate::value::{
EventType, MAX_EVENT_TYPE_LEN, MAX_METADATA_LEN, Metadata, Payload, SchemaVersion, ValueError,
};
#[allow(
clippy::as_conversions,
reason = "u32→usize is lossless on all Mnesis target platforms (32-bit+)"
)]
#[inline]
const fn idx(n: u32) -> usize {
n as usize
}
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum EnvelopeError {
#[error("range {start}..{end} exceeds buffer length {len}")]
RangeOutOfBounds { start: u32, end: u32, len: usize },
#[error("invalid UTF-8 in event_type at bytes {start}..{end}")]
InvalidUtf8 {
start: u32,
end: u32,
#[source]
source: core::str::Utf8Error,
},
#[error("event_type range length {actual} exceeds maximum {max}")]
EventTypeRangeTooLong { actual: u32, max: usize },
#[error("metadata range length {actual} exceeds maximum {max}")]
MetadataRangeTooLong { actual: u32, max: usize },
#[error("metadata range is empty; use None to represent absent metadata")]
MetadataRangeEmpty,
#[error(transparent)]
Value(#[from] ValueError),
}
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum ForDecodeError {
#[error(transparent)]
Value(#[from] ValueError),
#[error(transparent)]
Wire(#[from] crate::wire::WireError),
#[error(transparent)]
Envelope(#[from] EnvelopeError),
}
#[derive(Debug, Clone)]
pub struct PendingEnvelope {
version: Version,
event_type: EventType,
schema_version: SchemaVersion,
payload: Payload,
metadata: Option<Metadata>,
}
impl PendingEnvelope {
#[must_use]
pub const fn version(&self) -> Version {
self.version
}
#[must_use]
pub fn event_type(&self) -> &str {
self.event_type.as_str()
}
#[must_use]
pub fn event_type_value(&self) -> EventType {
self.event_type.clone()
}
#[must_use]
pub fn payload(&self) -> &[u8] {
self.payload.as_slice()
}
#[must_use]
pub fn payload_bytes(&self) -> Bytes {
self.payload.clone().into_bytes()
}
#[must_use]
pub fn payload_value(&self) -> Payload {
self.payload.clone()
}
#[must_use]
pub fn metadata(&self) -> Option<&[u8]> {
self.metadata.as_ref().map(Metadata::as_slice)
}
#[must_use]
pub fn metadata_bytes(&self) -> Option<Bytes> {
self.metadata.as_ref().map(|m| m.clone().into_bytes())
}
#[must_use]
pub fn metadata_value(&self) -> Option<Metadata> {
self.metadata.clone()
}
#[must_use]
pub const fn schema_version(&self) -> u32 {
self.schema_version.get()
}
#[must_use]
pub const fn schema_version_value(&self) -> SchemaVersion {
self.schema_version
}
#[cfg(feature = "import")]
#[must_use]
pub(crate) fn from_persisted(persisted: &PersistedEnvelope) -> Self {
Self {
version: persisted.version(),
event_type: persisted.event_type_value(),
schema_version: persisted.schema_version_value(),
payload: persisted.payload_value(),
metadata: persisted.metadata_value(),
}
}
}
pub type PendingBatchIter<'a> = Chain<Once<&'a PendingEnvelope>, Iter<'a, PendingEnvelope>>;
#[derive(Debug, Clone, Copy)]
pub struct PendingBatch<'a> {
first: &'a PendingEnvelope,
rest: &'a [PendingEnvelope],
}
impl<'a> PendingBatch<'a> {
#[must_use]
pub const fn of(only: &'a PendingEnvelope) -> Self {
Self {
first: only,
rest: &[],
}
}
#[must_use]
pub const fn from_parts(first: &'a PendingEnvelope, rest: &'a [PendingEnvelope]) -> Self {
Self { first, rest }
}
#[must_use]
pub fn new(envelopes: &'a [PendingEnvelope]) -> Option<Self> {
envelopes
.split_first()
.map(|(first, rest)| Self { first, rest })
}
#[must_use]
pub const fn first(&self) -> &'a PendingEnvelope {
self.first
}
#[must_use]
pub fn last(&self) -> &'a PendingEnvelope {
self.rest.last().unwrap_or(self.first)
}
#[must_use]
pub const fn len(&self) -> NonZeroUsize {
NonZeroUsize::MIN.saturating_add(self.rest.len())
}
pub fn iter(&self) -> PendingBatchIter<'a> {
once(self.first).chain(self.rest)
}
}
impl<'a> IntoIterator for PendingBatch<'a> {
type Item = &'a PendingEnvelope;
type IntoIter = PendingBatchIter<'a>;
fn into_iter(self) -> Self::IntoIter {
self.iter()
}
}
impl<'a> IntoIterator for &PendingBatch<'a> {
type Item = &'a PendingEnvelope;
type IntoIter = PendingBatchIter<'a>;
fn into_iter(self) -> Self::IntoIter {
self.iter()
}
}
#[derive(Debug)]
pub struct WithVersion {
version: Version,
}
#[derive(Debug)]
pub struct WithEventType {
version: Version,
event_type: EventType,
}
#[derive(Debug)]
pub struct WithPayload {
version: Version,
event_type: EventType,
payload: Bytes,
schema_version: SchemaVersion,
metadata: Option<Bytes>,
}
impl WithVersion {
#[must_use]
pub fn event_type(self, event_type: &'static str) -> WithEventType {
WithEventType {
version: self.version,
event_type: EventType::from_static_str(event_type),
}
}
#[must_use]
pub fn event<E: DomainEvent + ?Sized>(self, event: &E) -> WithEventType {
self.event_type(event.name())
}
pub fn event_type_bytes(self, bytes: Bytes) -> Result<WithEventType, EnvelopeError> {
let event_type = EventType::from_bytes(bytes)?;
Ok(WithEventType {
version: self.version,
event_type,
})
}
}
impl WithEventType {
#[must_use]
pub fn payload(self, payload: impl Into<Bytes>) -> WithPayload {
WithPayload {
version: self.version,
event_type: self.event_type,
payload: payload.into(),
schema_version: SchemaVersion::INITIAL,
metadata: None,
}
}
}
impl WithPayload {
#[must_use]
pub const fn schema_version(mut self, schema_version: SchemaVersion) -> Self {
self.schema_version = schema_version;
self
}
#[must_use]
pub fn metadata(mut self, metadata: impl Into<Bytes>) -> Self {
self.metadata = Some(metadata.into());
self
}
pub fn build(self) -> Result<PendingEnvelope, EnvelopeError> {
let payload = Payload::from_bytes(self.payload)?;
let metadata = self.metadata.map(Metadata::from_bytes).transpose()?;
Ok(PendingEnvelope {
version: self.version,
event_type: self.event_type,
schema_version: self.schema_version,
payload,
metadata,
})
}
}
#[must_use]
pub const fn pending_envelope(version: Version) -> WithVersion {
WithVersion { version }
}
#[derive(Debug, Clone)]
pub struct PersistedEnvelope {
version: Version,
schema_version: SchemaVersion,
value: Bytes,
event_type_range: Range<u32>,
payload_range: Range<u32>,
metadata_range: Option<Range<u32>>,
}
impl PersistedEnvelope {
#[allow(
clippy::too_many_arguments,
reason = "all 6 fields are required to construct a validated PersistedEnvelope; \
a builder would add indirection with no type-safety benefit here"
)]
pub fn try_new(
version: Version,
value: Bytes,
schema_version: SchemaVersion,
event_type_range: Range<u32>,
payload_range: Range<u32>,
metadata_range: Option<Range<u32>>,
) -> Result<Self, EnvelopeError> {
let len = value.len();
check_range(&event_type_range, len)?;
check_range(&payload_range, len)?;
if let Some(ref m) = metadata_range {
check_range(m, len)?;
}
let et_range_len = event_type_range.end - event_type_range.start;
if idx(et_range_len) > MAX_EVENT_TYPE_LEN {
return Err(EnvelopeError::EventTypeRangeTooLong {
actual: et_range_len,
max: MAX_EVENT_TYPE_LEN,
});
}
if let Some(ref m) = metadata_range {
let meta_range_len = m.end - m.start;
if meta_range_len == 0 {
return Err(EnvelopeError::MetadataRangeEmpty);
}
if idx(meta_range_len) > MAX_METADATA_LEN {
return Err(EnvelopeError::MetadataRangeTooLong {
actual: meta_range_len,
max: MAX_METADATA_LEN,
});
}
}
let et_start = idx(event_type_range.start);
let et_end = idx(event_type_range.end);
core::str::from_utf8(&value[et_start..et_end]).map_err(|e| EnvelopeError::InvalidUtf8 {
start: event_type_range.start,
end: event_type_range.end,
source: e,
})?;
Ok(Self {
version,
schema_version,
value,
event_type_range,
payload_range,
metadata_range,
})
}
#[must_use]
pub const fn version(&self) -> Version {
self.version
}
#[must_use]
pub const fn schema_version(&self) -> u32 {
self.schema_version.get()
}
#[must_use]
pub const fn schema_version_value(&self) -> SchemaVersion {
self.schema_version
}
#[must_use]
pub fn event_type(&self) -> &str {
let start = idx(self.event_type_range.start);
let end = idx(self.event_type_range.end);
#[allow(
unsafe_code,
reason = "UTF-8 invariant established at construction; ranges validated"
)]
unsafe {
core::str::from_utf8_unchecked(&self.value[start..end])
}
}
#[must_use]
pub fn event_type_bytes(&self) -> Bytes {
self.slice_range(&self.event_type_range)
}
#[must_use]
pub fn payload(&self) -> &[u8] {
let start = idx(self.payload_range.start);
let end = idx(self.payload_range.end);
&self.value[start..end]
}
#[must_use]
pub fn payload_bytes(&self) -> Bytes {
self.slice_range(&self.payload_range)
}
#[must_use]
pub fn metadata(&self) -> Option<&[u8]> {
self.metadata_range.as_ref().map(|r| {
let start = idx(r.start);
let end = idx(r.end);
&self.value[start..end]
})
}
#[must_use]
pub fn metadata_bytes(&self) -> Option<Bytes> {
self.metadata_range.as_ref().map(|r| self.slice_range(r))
}
#[must_use]
pub fn event_type_value(&self) -> EventType {
#[allow(
unsafe_code,
reason = "UTF-8 and length cap both established by try_new"
)]
unsafe {
EventType::from_validated_bytes(self.event_type_bytes())
}
}
#[must_use]
pub fn payload_value(&self) -> Payload {
Payload::from_validated_bytes(self.payload_bytes())
}
#[must_use]
pub fn metadata_value(&self) -> Option<Metadata> {
self.metadata_bytes().map(|b| {
Metadata::from_validated_bytes(b)
})
}
#[must_use]
pub fn schema_version_as_version(&self) -> Version {
Version::from(self.schema_version)
}
pub fn for_decode(event_type: &str, payload: &[u8]) -> Result<Self, ForDecodeError> {
let et = EventType::from_bytes(Bytes::copy_from_slice(event_type.as_bytes()))?;
let pl = Payload::from_bytes(Bytes::copy_from_slice(payload))?;
let sv = SchemaVersion::INITIAL;
let frame = crate::wire::encode_frame(sv, &et, &pl, None)?;
Ok(Self::try_new(
Version::INITIAL,
frame.value,
sv,
frame.offsets.event_type,
frame.offsets.payload,
None,
)?)
}
fn slice_range(&self, range: &Range<u32>) -> Bytes {
self.value.slice(idx(range.start)..idx(range.end))
}
}
const fn check_range(range: &Range<u32>, len: usize) -> Result<(), EnvelopeError> {
if idx(range.end) > len || range.start > range.end {
return Err(EnvelopeError::RangeOutOfBounds {
start: range.start,
end: range.end,
len,
});
}
Ok(())
}
#[cfg(test)]
#[allow(
clippy::unwrap_used,
clippy::expect_used,
clippy::panic,
reason = "test code asserts exact values"
)]
mod tests {
use super::*;
use crate::value::MAX_EVENT_TYPE_LEN;
use bytes::Bytes;
use mnesis::Version;
#[derive(Debug)]
struct TestEvent;
impl mnesis::Message for TestEvent {}
impl DomainEvent for TestEvent {
fn name(&self) -> &'static str {
"TestEvent"
}
}
#[test]
fn pending_envelope_builds_with_metadata() {
let env = pending_envelope(Version::INITIAL)
.event_type("UserCreated")
.payload(Bytes::from_static(b"payload-bytes"))
.metadata(Bytes::from_static(b"meta-bytes"))
.build()
.expect("valid envelope");
assert_eq!(env.event_type(), "UserCreated");
assert_eq!(env.payload(), b"payload-bytes");
assert_eq!(env.metadata(), Some(b"meta-bytes".as_slice()));
assert_eq!(env.schema_version(), 1);
}
#[test]
fn pending_envelope_builds_without_metadata() {
let env = pending_envelope(Version::INITIAL)
.event_type("X")
.payload(Bytes::from_static(b"p"))
.build()
.expect("valid envelope");
assert_eq!(env.metadata(), None);
}
#[test]
fn pending_envelope_typed_accessors_roundtrip() {
let env = pending_envelope(Version::INITIAL)
.event_type("UserCreated")
.payload(Bytes::from_static(b"payload-bytes"))
.metadata(Bytes::from_static(b"meta-bytes"))
.build()
.expect("valid envelope");
assert_eq!(env.event_type_value().as_str(), "UserCreated");
assert_eq!(env.payload_value().as_slice(), b"payload-bytes");
assert_eq!(
env.metadata_value().map(|m| m.as_slice().to_vec()),
Some(b"meta-bytes".to_vec())
);
assert_eq!(env.schema_version_value().get(), 1);
}
#[test]
fn pending_envelope_rejects_oversize_event_type() {
let oversized = "x".repeat(MAX_EVENT_TYPE_LEN + 1);
let err = pending_envelope(Version::INITIAL)
.event_type_bytes(Bytes::from(oversized))
.expect_err("oversized must be rejected");
assert!(matches!(err, EnvelopeError::Value(_)));
}
#[test]
fn event_derives_same_event_type_as_event_type_name() {
let via_event = pending_envelope(Version::INITIAL)
.event(&TestEvent)
.payload(Bytes::from_static(b"p"))
.build()
.expect("valid envelope");
let via_name = pending_envelope(Version::INITIAL)
.event_type(TestEvent.name())
.payload(Bytes::from_static(b"p"))
.build()
.expect("valid envelope");
assert_eq!(via_event.event_type(), via_name.event_type());
assert_eq!(via_event.event_type(), "TestEvent");
}
#[test]
fn build_rejects_oversize_payload() {
let oversized = vec![0u8; crate::value::MAX_PAYLOAD_LEN + 1];
let err = pending_envelope(Version::INITIAL)
.event_type("X")
.payload(Bytes::from(oversized))
.build()
.expect_err("oversized payload must be rejected");
assert!(matches!(err, EnvelopeError::Value(_)));
}
#[test]
fn build_rejects_empty_metadata() {
let err = pending_envelope(Version::INITIAL)
.event_type("X")
.payload(Bytes::from_static(b"p"))
.metadata(Bytes::new())
.build()
.expect_err("empty metadata must be rejected");
assert!(matches!(err, EnvelopeError::Value(_)));
}
#[test]
fn persisted_envelope_accessors_return_views_into_value() {
let value = Bytes::from_static(b"TYPEpayloadmeta");
let env = PersistedEnvelope::try_new(
Version::INITIAL,
value,
SchemaVersion::INITIAL,
0..4,
4..11,
Some(11..15),
)
.expect("valid construction");
assert_eq!(env.event_type(), "TYPE");
assert_eq!(env.payload(), b"payload");
assert_eq!(env.metadata(), Some(b"meta".as_slice()));
}
#[test]
fn persisted_envelope_payload_bytes_shares_arc() {
let env = PersistedEnvelope::try_new(
Version::INITIAL,
Bytes::from_static(b"TYPEpayload"),
SchemaVersion::INITIAL,
0..4,
4..11,
None,
)
.unwrap();
let payload = env.payload_bytes();
assert_eq!(payload.as_ref(), b"payload");
}
#[test]
fn persisted_envelope_rejects_range_past_buffer() {
let value = Bytes::from_static(b"short");
let err = PersistedEnvelope::try_new(
Version::INITIAL,
value,
SchemaVersion::INITIAL,
0..4,
4..100,
None,
)
.expect_err("must reject out-of-bounds range");
assert!(matches!(err, EnvelopeError::RangeOutOfBounds { .. }));
}
#[test]
fn try_new_rejects_event_type_range_too_long() {
let len = MAX_EVENT_TYPE_LEN + 1;
let mut buf = vec![0u8; len];
buf[0] = b'A';
let value = Bytes::from(buf);
let too_long_end = u32::try_from(len).expect("len fits in u32 by construction");
let err = PersistedEnvelope::try_new(
Version::INITIAL,
value,
SchemaVersion::INITIAL,
0..too_long_end,
too_long_end..too_long_end,
None,
)
.expect_err("event_type range > MAX_EVENT_TYPE_LEN must be rejected");
let expected_actual = u32::try_from(len).expect("len fits in u32 by construction (above)");
assert!(matches!(
err,
EnvelopeError::EventTypeRangeTooLong { actual, max }
if actual == expected_actual && max == MAX_EVENT_TYPE_LEN
));
}
#[test]
fn persisted_envelope_rejects_empty_metadata_range() {
let env = PersistedEnvelope::try_new(
Version::INITIAL,
Bytes::from_static(b"TYPEpayload"),
SchemaVersion::INITIAL,
0..4,
4..11,
Some(4..4),
);
assert!(matches!(env, Err(EnvelopeError::MetadataRangeEmpty)));
}
#[test]
fn persisted_envelope_rejects_invalid_utf8_in_event_type() {
let value = Bytes::from_static(&[0xFFu8, 0xFF, b'p', b'a', b'y']);
let err = PersistedEnvelope::try_new(
Version::INITIAL,
value,
SchemaVersion::INITIAL,
0..2,
2..5,
None,
)
.expect_err("must reject non-UTF-8 event_type");
assert!(matches!(err, EnvelopeError::InvalidUtf8 { .. }));
}
#[test]
fn persisted_envelope_event_type_value_returns_validated_value_newtype() {
let env = PersistedEnvelope::try_new(
Version::INITIAL,
Bytes::from_static(b"TYPEpayload"),
SchemaVersion::INITIAL,
0..4,
4..11,
None,
)
.expect("valid");
let et = env.event_type_value();
assert_eq!(et.as_str(), "TYPE");
}
#[test]
fn persisted_envelope_payload_value_returns_validated_value_newtype() {
let env = PersistedEnvelope::try_new(
Version::INITIAL,
Bytes::from_static(b"TYPEpayload"),
SchemaVersion::INITIAL,
0..4,
4..11,
None,
)
.expect("valid");
let p = env.payload_value();
assert_eq!(p.as_slice(), b"payload");
}
#[test]
fn persisted_envelope_metadata_value_returns_some_when_present() {
let env = PersistedEnvelope::try_new(
Version::INITIAL,
Bytes::from_static(b"TYPEpayloadMETA"),
SchemaVersion::INITIAL,
0..4,
4..11,
Some(11..15),
)
.expect("valid");
let m = env.metadata_value().expect("present");
assert_eq!(m.as_slice(), b"META");
}
#[test]
fn persisted_envelope_metadata_value_returns_none_when_absent() {
let env = PersistedEnvelope::try_new(
Version::INITIAL,
Bytes::from_static(b"TYPEpayload"),
SchemaVersion::INITIAL,
0..4,
4..11,
None,
)
.expect("valid");
assert!(env.metadata_value().is_none());
}
#[cfg(feature = "import")]
#[test]
fn from_persisted_preserves_fields_drops_global_seq_zero_copy() {
let value = Bytes::from_static(b"TYPEpayloadmeta");
let persisted = PersistedEnvelope::try_new(
Version::new(7).expect("nonzero"),
value,
crate::value::SchemaVersion::from_u32(3).expect("nonzero"),
0..4,
4..11,
Some(11..15),
)
.expect("valid");
let pending = PendingEnvelope::from_persisted(&persisted);
assert_eq!(pending.version(), persisted.version());
assert_eq!(pending.event_type(), "TYPE");
assert_eq!(pending.payload(), b"payload");
assert_eq!(pending.metadata(), Some(b"meta".as_slice()));
assert_eq!(pending.schema_version(), 3);
assert!(
std::ptr::eq(
pending.payload_bytes().as_ptr(),
persisted.payload().as_ptr()
),
"from_persisted must reuse the Arc-shared payload, not deep-copy",
);
}
}
#[cfg(test)]
#[allow(clippy::expect_used, reason = "test code")]
mod persisted_exhaustive_tests {
use super::{EnvelopeError, PersistedEnvelope};
use crate::value::{MAX_EVENT_TYPE_LEN, Metadata, SchemaVersion};
use bytes::Bytes;
use mnesis::Version;
use proptest::prelude::*;
use static_assertions::assert_impl_all;
use std::ops::Range;
fn v(n: u64) -> Version {
Version::new(n).expect("test version must be nonzero")
}
fn sv(n: u32) -> SchemaVersion {
SchemaVersion::from_u32(n).expect("test schema_version must be nonzero")
}
fn assemble(
event_type: &[u8],
payload: &[u8],
metadata: Option<&[u8]>,
) -> (Bytes, Range<u32>, Range<u32>, Option<Range<u32>>) {
let mut buf = Vec::new();
buf.extend_from_slice(event_type);
buf.extend_from_slice(payload);
if let Some(m) = metadata {
buf.extend_from_slice(m);
}
let et_end = u32::try_from(event_type.len()).expect("event_type len fits u32");
let pl_end = et_end + u32::try_from(payload.len()).expect("payload len fits u32");
let meta_range = metadata.map(|m| {
let end = pl_end + u32::try_from(m.len()).expect("metadata len fits u32");
pl_end..end
});
(Bytes::from(buf), 0..et_end, et_end..pl_end, meta_range)
}
fn build(
version: Version,
schema: SchemaVersion,
event_type: &[u8],
payload: &[u8],
metadata: Option<&[u8]>,
) -> PersistedEnvelope {
let (value, et, pl, meta) = assemble(event_type, payload, metadata);
PersistedEnvelope::try_new(version, value, schema, et, pl, meta)
.expect("components form a valid PersistedEnvelope")
}
assert_impl_all!(PersistedEnvelope: Send, Sync, Clone, std::fmt::Debug);
#[test]
fn persisted_envelope_clones_across_thread_boundary() {
let event = build(
v(3),
sv(2),
b"MoneyDeposited",
b"payload-bytes",
Some(b"meta"),
);
let moved = event.clone();
let (ver, schema, etype, payload, meta) = std::thread::scope(|scope| {
scope
.spawn(move || {
(
moved.version(),
moved.schema_version(),
moved.event_type().to_owned(),
moved.payload().to_vec(),
moved.metadata().map(<[u8]>::to_vec),
)
})
.join()
.expect("worker thread must not panic")
});
assert_eq!(ver, v(3));
assert_eq!(schema, 2);
assert_eq!(etype, "MoneyDeposited");
assert_eq!(payload, b"payload-bytes");
assert_eq!(meta, Some(b"meta".to_vec()));
assert_eq!(event.event_type(), "MoneyDeposited");
}
#[test]
fn repeated_accessor_calls_are_consistent() {
let event = build(v(9), sv(4), b"TYPE", b"payload", Some(b"meta"));
assert_eq!(event.event_type(), event.event_type());
assert_eq!(event.payload(), event.payload());
assert_eq!(event.metadata(), event.metadata());
assert_eq!(
event.event_type().as_bytes(),
event.event_type_bytes().as_ref()
);
assert_eq!(event.payload(), event.payload_bytes().as_ref());
assert_eq!(event.metadata(), event.metadata_bytes().as_deref());
assert_eq!(event.event_type(), event.event_type_value().as_str());
assert_eq!(event.payload(), event.payload_value().as_slice());
assert_eq!(
event.metadata(),
event.metadata_value().as_ref().map(Metadata::as_slice),
);
assert_eq!(event.version(), v(9));
assert_eq!(event.version(), event.version());
assert_eq!(event.schema_version(), 4);
assert_eq!(event.schema_version_value(), sv(4));
}
#[test]
fn clone_is_field_for_field_equal_view() {
let original = build(
v(42),
sv(7),
b"OrderPlaced",
b"the-payload",
Some(b"the-meta"),
);
let cloned = original.clone();
assert_eq!(cloned.version(), original.version());
assert_eq!(cloned.schema_version(), original.schema_version());
assert_eq!(
cloned.schema_version_value(),
original.schema_version_value()
);
assert_eq!(cloned.event_type(), original.event_type());
assert_eq!(cloned.payload(), original.payload());
assert_eq!(cloned.metadata(), original.metadata());
assert_eq!(
cloned.metadata_value().map(|m| m.as_slice().to_vec()),
original.metadata_value().map(|m| m.as_slice().to_vec()),
);
}
#[test]
fn clone_shares_the_same_backing_buffer() {
let original = build(v(1), sv(1), b"TYPE", b"payload", None);
let cloned = original.clone();
let from_original = original.payload_bytes();
let from_clone = cloned.payload_bytes();
assert!(
std::ptr::eq(from_original.as_ptr(), from_clone.as_ptr()),
"clone must share the parent buffer, not deep-copy it",
);
assert_eq!(from_original.as_ref(), from_clone.as_ref());
}
#[test]
fn rejects_inverted_range_start_after_end() {
let value = Bytes::from_static(b"TYPEpayload");
let (bad_start, bad_end) = (7u32, 2u32);
let err = PersistedEnvelope::try_new(
Version::INITIAL,
value,
SchemaVersion::INITIAL,
0..4,
bad_start..bad_end,
None,
)
.expect_err("start > end must be rejected");
assert!(matches!(
err,
EnvelopeError::RangeOutOfBounds {
start: 7,
end: 2,
..
}
));
}
#[test]
fn event_type_cap_uses_range_length_not_endpoint_sum() {
let start: u32 = 40_000;
let end: u32 = 70_000;
let buf_len = usize::try_from(end).expect("70_000 fits usize");
let value = Bytes::from(vec![b'A'; buf_len]);
let event = PersistedEnvelope::try_new(
Version::INITIAL,
value,
SchemaVersion::INITIAL,
start..end,
0..0,
None,
)
.expect("event_type length 30_000 is within MAX_EVENT_TYPE_LEN (65_535)");
assert_eq!(event.event_type().len(), 30_000);
}
#[test]
fn rejects_event_type_range_past_buffer_end() {
let value = Bytes::from_static(b"short");
let err = PersistedEnvelope::try_new(
Version::INITIAL,
value,
SchemaVersion::INITIAL,
0..100,
0..0,
None,
)
.expect_err("event_type range past buffer must be rejected");
assert!(matches!(
err,
EnvelopeError::RangeOutOfBounds {
start: 0,
end: 100,
len: 5
}
));
}
#[test]
fn rejects_metadata_range_past_buffer_end() {
let value = Bytes::from_static(b"TYPEpayload");
let err = PersistedEnvelope::try_new(
Version::INITIAL,
value,
SchemaVersion::INITIAL,
0..4,
4..11,
Some(11..100),
)
.expect_err("metadata range past buffer must be rejected");
assert!(matches!(
err,
EnvelopeError::RangeOutOfBounds {
start: 11,
end: 100,
len: 11
}
));
}
#[test]
fn rejects_inverted_metadata_range() {
let value = Bytes::from_static(b"TYPEpayload");
let (bad_start, bad_end) = (9u32, 5u32);
let err = PersistedEnvelope::try_new(
Version::INITIAL,
value,
SchemaVersion::INITIAL,
0..4,
4..11,
Some(bad_start..bad_end),
)
.expect_err("metadata start > end must be rejected");
assert!(matches!(
err,
EnvelopeError::RangeOutOfBounds {
start: 9,
end: 5,
..
}
));
}
#[test]
fn ranges_are_independent_and_may_overlap() {
let value = Bytes::from_static(b"TYPE");
let event = PersistedEnvelope::try_new(
Version::INITIAL,
value,
SchemaVersion::INITIAL,
0..4,
0..4,
Some(0..4),
)
.expect("overlapping ranges are structurally permitted");
assert_eq!(event.event_type(), "TYPE");
assert_eq!(event.payload(), b"TYPE");
assert_eq!(event.metadata(), Some(b"TYPE".as_slice()));
}
#[test]
fn empty_payload_range_is_accepted_unlike_empty_metadata() {
let value = Bytes::from_static(b"TYPE");
let event = PersistedEnvelope::try_new(
Version::INITIAL,
value,
SchemaVersion::INITIAL,
0..4,
4..4,
None,
)
.expect("empty payload range is accepted");
assert_eq!(event.payload(), b"");
assert_eq!(event.payload_value().as_slice(), b"");
}
#[test]
fn metadata_absent_is_none_everywhere() {
let event = build(v(1), sv(1), b"TYPE", b"payload", None);
assert!(event.metadata().is_none());
assert!(event.metadata_bytes().is_none());
assert!(event.metadata_value().is_none());
}
#[test]
fn metadata_present_threads_through_every_accessor() {
let event = build(v(1), sv(1), b"TYPE", b"payload", Some(b"META"));
assert_eq!(event.metadata(), Some(b"META".as_slice()));
assert_eq!(event.metadata_bytes().expect("some").as_ref(), b"META");
assert_eq!(event.metadata_value().expect("some").as_slice(), b"META");
}
#[test]
fn empty_event_type_is_accepted() {
let event = build(v(1), sv(1), b"", b"payload", None);
assert_eq!(event.event_type(), "");
assert_eq!(event.event_type_value().as_str(), "");
assert_eq!(event.payload(), b"payload");
}
#[test]
fn boundary_version_and_schema_version_carried_verbatim() {
let event = build(v(u64::MAX), sv(u32::MAX), b"TYPE", b"payload", None);
assert_eq!(event.version().as_u64(), u64::MAX);
assert_eq!(event.schema_version(), u32::MAX);
assert_eq!(event.schema_version_value().get(), u32::MAX);
}
#[test]
fn owned_views_alias_the_single_backing_buffer() {
let (value, et, pl, meta) = assemble(b"TYPE", b"payload", Some(b"meta"));
let base = value.clone(); let event = PersistedEnvelope::try_new(
Version::INITIAL,
value,
SchemaVersion::INITIAL,
et,
pl,
meta,
)
.expect("valid");
let et_bytes = event.event_type_bytes();
let pl_bytes = event.payload_bytes();
let meta_bytes = event.metadata_bytes().expect("metadata present");
assert!(
std::ptr::eq(et_bytes.as_ptr(), base.as_ptr().wrapping_add(0)),
"event_type_bytes must alias buffer offset 0",
);
assert!(
std::ptr::eq(pl_bytes.as_ptr(), base.as_ptr().wrapping_add(4)),
"payload_bytes must alias buffer offset 4",
);
assert!(
std::ptr::eq(meta_bytes.as_ptr(), base.as_ptr().wrapping_add(11)),
"metadata_bytes must alias buffer offset 11",
);
assert!(std::ptr::eq(
event.event_type_value().as_bytes().as_ptr(),
base.as_ptr().wrapping_add(0),
));
assert!(std::ptr::eq(
event.payload_value().as_slice().as_ptr(),
base.as_ptr().wrapping_add(4),
));
}
fn boundary_u64() -> impl Strategy<Value = u64> {
prop_oneof![
Just(1u64),
Just(2u64),
Just(u64::MAX - 1),
Just(u64::MAX),
1u64..=u64::MAX,
]
}
fn boundary_u32_nonzero() -> impl Strategy<Value = u32> {
prop_oneof![
Just(1u32),
Just(2u32),
Just(u32::MAX - 1),
Just(u32::MAX),
1u32..=u32::MAX,
]
}
fn event_type_bytes_strategy() -> impl Strategy<Value = Vec<u8>> {
prop_oneof![
Just(0usize),
Just(1usize),
Just(255usize),
Just(MAX_EVENT_TYPE_LEN),
]
.prop_flat_map(|n| proptest::collection::vec(b'A'..=b'Z', n))
}
fn payload_bytes_strategy() -> impl Strategy<Value = Vec<u8>> {
prop_oneof![Just(0usize), Just(1usize), Just(1024usize)]
.prop_flat_map(|n| proptest::collection::vec(any::<u8>(), n))
}
fn metadata_bytes_strategy() -> impl Strategy<Value = Option<Vec<u8>>> {
prop_oneof![
Just(None),
Just(Some(vec![0u8])),
(1usize..=2048)
.prop_flat_map(|n| proptest::collection::vec(any::<u8>(), n))
.prop_map(Some),
]
}
proptest! {
#[test]
fn persisted_envelope_roundtrips_every_accessor(
event_type in event_type_bytes_strategy(),
payload in payload_bytes_strategy(),
metadata in metadata_bytes_strategy(),
version_raw in boundary_u64(),
schema_raw in boundary_u32_nonzero(),
) {
let version = v(version_raw);
let schema = sv(schema_raw);
let event = build(
version,
schema,
&event_type,
&payload,
metadata.as_deref(),
);
prop_assert_eq!(event.version(), version);
prop_assert_eq!(event.version().as_u64(), version_raw);
prop_assert_eq!(event.schema_version(), schema_raw);
prop_assert_eq!(event.schema_version_value(), schema);
let et_bytes = event.event_type_bytes();
let et_value = event.event_type_value();
prop_assert_eq!(event.event_type().as_bytes(), event_type.as_slice());
prop_assert_eq!(et_bytes.as_ref(), event_type.as_slice());
prop_assert_eq!(et_value.as_bytes(), event_type.as_slice());
let pl_bytes = event.payload_bytes();
let pl_value = event.payload_value();
prop_assert_eq!(event.payload(), payload.as_slice());
prop_assert_eq!(pl_bytes.as_ref(), payload.as_slice());
prop_assert_eq!(pl_value.as_slice(), payload.as_slice());
if let Some(ref m) = metadata {
prop_assert_eq!(event.metadata(), Some(m.as_slice()));
let meta_bytes = event.metadata_bytes().expect("some");
prop_assert_eq!(meta_bytes.as_ref(), m.as_slice());
let meta_value = event.metadata_value().expect("some");
prop_assert_eq!(meta_value.as_slice(), m.as_slice());
} else {
prop_assert!(event.metadata().is_none());
prop_assert!(event.metadata_bytes().is_none());
prop_assert!(event.metadata_value().is_none());
}
}
}
fn arbitrary_range() -> impl Strategy<Value = Range<u32>> {
(0u32..=260, 0u32..=260).prop_map(|(start, end)| start..end)
}
fn arbitrary_optional_range() -> impl Strategy<Value = Option<Range<u32>>> {
prop_oneof![Just(None), arbitrary_range().prop_map(Some)]
}
proptest! {
#[test]
fn try_new_never_panics_and_accepted_events_have_sound_accessors(
value_bytes in proptest::collection::vec(any::<u8>(), 0..256),
event_type_range in arbitrary_range(),
payload_range in arbitrary_range(),
metadata_range in arbitrary_optional_range(),
version_raw in boundary_u64(),
schema_raw in boundary_u32_nonzero(),
) {
let result = PersistedEnvelope::try_new(
v(version_raw),
Bytes::from(value_bytes),
sv(schema_raw),
event_type_range,
payload_range,
metadata_range,
);
match result {
Ok(event) => {
let et = event.event_type();
prop_assert!(std::str::from_utf8(et.as_bytes()).is_ok());
let _ = event.payload();
let _ = event.metadata();
let _ = event.event_type_value();
let _ = event.payload_value();
let _ = event.metadata_value();
let _ = event.event_type_bytes();
let _ = event.payload_bytes();
let _ = event.metadata_bytes();
}
Err(err) => {
match err {
EnvelopeError::RangeOutOfBounds { .. }
| EnvelopeError::InvalidUtf8 { .. }
| EnvelopeError::EventTypeRangeTooLong { .. }
| EnvelopeError::MetadataRangeTooLong { .. }
| EnvelopeError::MetadataRangeEmpty
| EnvelopeError::Value(_) => {}
}
}
}
}
}
}
#[cfg(test)]
#[allow(
clippy::unwrap_used,
clippy::expect_used,
clippy::panic,
reason = "test code asserts exact values"
)]
mod pending_batch_tests {
use super::{PendingBatch, PendingEnvelope, pending_envelope};
use bytes::Bytes;
use mnesis::Version;
fn env(version: u64) -> PendingEnvelope {
pending_envelope(Version::new(version).expect("version is non-zero"))
.event_type("E")
.payload(Bytes::from_static(b"p"))
.build()
.expect("valid envelope")
}
#[test]
fn empty_slice_yields_no_batch() {
assert!(PendingBatch::new(&[]).is_none());
}
#[test]
fn batch_keeps_every_envelope_in_order() {
let envs = [env(1), env(2), env(3)];
let batch = PendingBatch::new(&envs).expect("three envelopes are non-empty");
assert_eq!(batch.len().get(), 3);
let mut versions = batch.iter().map(|e| e.version().as_u64());
assert_eq!(versions.next(), Some(1));
assert_eq!(versions.next(), Some(2));
assert_eq!(versions.next(), Some(3));
assert_eq!(versions.next(), None);
}
#[test]
fn first_and_last_are_the_boundary_envelopes() {
let envs = [env(7), env(8), env(9)];
let batch = PendingBatch::new(&envs).expect("non-empty");
assert_eq!(batch.first().version().as_u64(), 7);
assert_eq!(batch.last().version().as_u64(), 9);
}
#[test]
fn single_envelope_batch_is_infallible() {
let only = env(4);
let batch = PendingBatch::of(&only);
assert_eq!(batch.len().get(), 1);
assert_eq!(batch.first().version().as_u64(), 4);
assert_eq!(batch.last().version().as_u64(), 4);
}
#[test]
fn from_parts_agrees_with_new_over_the_same_envelopes() {
let envs = [env(1), env(2)];
let (first, rest) = envs.split_first().expect("non-empty");
let parts = PendingBatch::from_parts(first, rest);
let whole = PendingBatch::new(&envs).expect("non-empty");
assert_eq!(parts.len(), whole.len());
let mut a = parts.iter().map(|e| e.version().as_u64());
let mut b = whole.iter().map(|e| e.version().as_u64());
assert_eq!(a.next(), b.next());
assert_eq!(a.next(), b.next());
assert_eq!(a.next(), None);
assert_eq!(b.next(), None);
}
}