use std::fmt;
use std::str;
use thiserror::Error;
pub use tephra_types::{EventType, MAX_NAME_LEN, NameError, Tag, Tags, TagsError};
const FIXED_HEADER: usize = 4;
fn header_len(tag_count: usize) -> usize {
FIXED_HEADER + 2 * tag_count
}
#[derive(Debug, Error, PartialEq, Eq)]
pub enum EncodeError {
#[error("event has {count} tags, exceeding the maximum of {max}")]
TooManyTags { count: usize, max: usize },
#[error("encoded event data region of {size} bytes exceeds the {max}-byte maximum")]
TooLarge { size: u64, max: u64 },
}
#[derive(Clone, Debug, Error, PartialEq, Eq)]
pub enum DecodeError {
#[error("buffer is shorter than the encoded event claims")]
Truncated,
#[error("encoded event data region exceeds the addressable maximum")]
TooLarge,
#[error("event type is empty")]
EmptyType,
#[error("a tag is empty")]
EmptyTag,
#[error("type or tag bytes are not valid UTF-8")]
InvalidUtf8,
#[error("tags are not in strictly ascending order")]
TagsNotSorted,
}
#[derive(Clone, PartialEq, Eq)]
pub struct Event {
buf: Box<[u8]>,
data_offset: u32,
}
impl Event {
pub fn new(event_type: &EventType, tags: &Tags, payload: &[u8]) -> Result<Self, EncodeError> {
let ty = event_type.as_str();
let tag_slice = tags.as_slice();
let tag_count = tag_slice.len();
if tag_count > u16::MAX as usize {
return Err(EncodeError::TooManyTags {
count: tag_count,
max: u16::MAX as usize,
});
}
let type_len = ty.len();
let tags_total: u64 = tag_slice.iter().map(|t| t.as_str().len() as u64).sum();
let data_start = header_len(tag_count) as u64 + type_len as u64 + tags_total;
if data_start > u32::MAX as u64 {
return Err(EncodeError::TooLarge {
size: data_start,
max: u32::MAX as u64,
});
}
let data_offset = data_start as u32;
let mut buf = Vec::with_capacity(data_start as usize + payload.len());
buf.extend_from_slice(&(type_len as u16).to_le_bytes());
buf.extend_from_slice(&(tag_count as u16).to_le_bytes());
for t in tag_slice {
buf.extend_from_slice(&(t.as_str().len() as u16).to_le_bytes());
}
buf.extend_from_slice(ty.as_bytes());
for t in tag_slice {
buf.extend_from_slice(t.as_str().as_bytes());
}
buf.extend_from_slice(payload);
Ok(Event {
buf: buf.into_boxed_slice(),
data_offset,
})
}
pub fn as_ref(&self) -> EventRef<'_> {
EventRef {
buf: &self.buf,
data_offset: self.data_offset,
}
}
pub fn event_type(&self) -> &str {
decode_type(&self.buf)
}
pub fn tags(&self) -> TagsRef<'_> {
decode_tags(&self.buf)
}
pub fn data(&self) -> &[u8] {
decode_data(&self.buf, self.data_offset)
}
pub fn as_bytes(&self) -> &[u8] {
&self.buf
}
}
impl fmt::Debug for Event {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
self.as_ref().fmt(f)
}
}
#[derive(Clone, Copy, PartialEq, Eq)]
pub struct EventRef<'a> {
buf: &'a [u8],
data_offset: u32,
}
impl<'a> EventRef<'a> {
pub fn from_bytes(buf: &'a [u8]) -> Result<Self, DecodeError> {
if buf.len() < FIXED_HEADER {
return Err(DecodeError::Truncated);
}
let type_len = read_u16(buf, 0) as usize;
let tag_count = read_u16(buf, 2) as usize;
let hlen = header_len(tag_count);
if buf.len() < hlen {
return Err(DecodeError::Truncated);
}
if type_len == 0 {
return Err(DecodeError::EmptyType);
}
let mut tags_total: usize = 0;
for i in 0..tag_count {
let len = read_u16(buf, FIXED_HEADER + 2 * i) as usize;
if len == 0 {
return Err(DecodeError::EmptyTag);
}
tags_total = tags_total.checked_add(len).ok_or(DecodeError::TooLarge)?;
}
let data_start = hlen
.checked_add(type_len)
.and_then(|x| x.checked_add(tags_total))
.ok_or(DecodeError::TooLarge)?;
if buf.len() < data_start {
return Err(DecodeError::Truncated);
}
let data_offset = u32::try_from(data_start).map_err(|_| DecodeError::TooLarge)?;
let type_end = hlen + type_len;
str::from_utf8(&buf[hlen..type_end]).map_err(|_| DecodeError::InvalidUtf8)?;
let mut pos = type_end;
let mut prev: Option<&str> = None;
for i in 0..tag_count {
let len = read_u16(buf, FIXED_HEADER + 2 * i) as usize;
let end = pos + len;
let tag = str::from_utf8(&buf[pos..end]).map_err(|_| DecodeError::InvalidUtf8)?;
if let Some(p) = prev
&& tag <= p
{
return Err(DecodeError::TagsNotSorted);
}
prev = Some(tag);
pos = end;
}
Ok(EventRef { buf, data_offset })
}
pub fn event_type(&self) -> &'a str {
decode_type(self.buf)
}
pub fn tags(&self) -> TagsRef<'a> {
decode_tags(self.buf)
}
pub fn data(&self) -> &'a [u8] {
decode_data(self.buf, self.data_offset)
}
pub fn as_bytes(&self) -> &'a [u8] {
self.buf
}
pub fn to_owned(&self) -> Event {
Event {
buf: Box::from(self.buf),
data_offset: self.data_offset,
}
}
}
impl fmt::Debug for EventRef<'_> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("EventRef")
.field("event_type", &self.event_type())
.field("tags", &self.tags().collect::<Vec<_>>())
.field("data_len", &self.data().len())
.finish()
}
}
fn read_u16(buf: &[u8], at: usize) -> u16 {
u16::from_le_bytes([buf[at], buf[at + 1]])
}
fn decode_type(buf: &[u8]) -> &str {
let type_len = read_u16(buf, 0) as usize;
let tag_count = read_u16(buf, 2) as usize;
let start = header_len(tag_count);
unsafe { str::from_utf8_unchecked(&buf[start..start + type_len]) }
}
fn decode_tags(buf: &[u8]) -> TagsRef<'_> {
let type_len = read_u16(buf, 0) as usize;
let tag_count = read_u16(buf, 2) as usize;
let lens = &buf[FIXED_HEADER..FIXED_HEADER + 2 * tag_count];
let start = header_len(tag_count) + type_len;
TagsRef {
data: &buf[start..],
lens,
idx: 0,
}
}
fn decode_data(buf: &[u8], data_offset: u32) -> &[u8] {
&buf[data_offset as usize..]
}
#[derive(Clone, Copy)]
pub struct TagsRef<'a> {
data: &'a [u8],
lens: &'a [u8],
idx: usize,
}
impl<'a> Iterator for TagsRef<'a> {
type Item = &'a str;
fn next(&mut self) -> Option<&'a str> {
let off = self.idx * 2;
if off + 2 > self.lens.len() {
return None;
}
let len = u16::from_le_bytes([self.lens[off], self.lens[off + 1]]) as usize;
self.idx += 1;
let (head, tail) = self.data.split_at(len);
self.data = tail;
Some(unsafe { str::from_utf8_unchecked(head) })
}
fn size_hint(&self) -> (usize, Option<usize>) {
let rem = self.lens.len() / 2 - self.idx;
(rem, Some(rem))
}
}
impl ExactSizeIterator for TagsRef<'_> {}
#[cfg(test)]
mod tests {
use super::*;
use smallvec::SmallVec;
fn ty(s: &str) -> EventType {
EventType::new(s).unwrap()
}
fn tag(s: &str) -> Tag {
Tag::new(s).unwrap()
}
fn tags(items: &[&str]) -> Tags {
Tags::new(items.iter().map(|s| tag(s)).collect::<SmallVec<[Tag; 4]>>()).unwrap()
}
#[test]
fn round_trip_full_event() {
let event = Event::new(
&ty("Registered"),
&tags(&["course:c1", "student:s1"]),
b"payload",
)
.unwrap();
let decoded = EventRef::from_bytes(event.as_bytes()).unwrap();
assert_eq!(decoded.event_type(), "Registered");
assert_eq!(
decoded.tags().collect::<Vec<_>>(),
vec!["course:c1", "student:s1"]
);
assert_eq!(decoded.data(), b"payload");
assert_eq!(decoded.tags().len(), 2);
}
#[test]
fn round_trip_no_tags_empty_payload() {
let event = Event::new(&ty("Ping"), &Tags::empty(), b"").unwrap();
let decoded = EventRef::from_bytes(event.as_bytes()).unwrap();
assert_eq!(decoded.event_type(), "Ping");
assert_eq!(decoded.tags().count(), 0);
assert_eq!(decoded.data(), b"");
}
#[test]
fn owned_and_borrowed_agree() {
let event = Event::new(&ty("T"), &tags(&["a", "bb", "ccc"]), b"data").unwrap();
let borrowed = EventRef::from_bytes(event.as_bytes()).unwrap();
let owned = borrowed.to_owned();
assert_eq!(owned.event_type(), borrowed.event_type());
assert_eq!(
owned.tags().collect::<Vec<_>>(),
borrowed.tags().collect::<Vec<_>>()
);
assert_eq!(owned.data(), borrowed.data());
assert_eq!(owned.as_bytes(), borrowed.as_bytes());
assert_eq!(owned, event);
}
#[test]
fn encoding_is_canonical() {
let a = Event::new(&ty("T"), &tags(&["z", "a", "m"]), b"p").unwrap();
let b = Event::new(&ty("T"), &tags(&["a", "m", "z"]), b"p").unwrap();
assert_eq!(a.as_bytes(), b.as_bytes());
}
#[test]
fn tags_never_touch_payload() {
let event = Event::new(&ty("T"), &tags(&["aa"]), b"\xff\xff\xff\xff").unwrap();
let decoded = EventRef::from_bytes(event.as_bytes()).unwrap();
assert_eq!(decoded.tags().collect::<Vec<_>>(), vec!["aa"]);
assert_eq!(decoded.data(), b"\xff\xff\xff\xff");
}
#[test]
fn decode_rejects_short_header() {
assert_eq!(EventRef::from_bytes(&[]), Err(DecodeError::Truncated));
assert_eq!(
EventRef::from_bytes(&[1, 0, 0]),
Err(DecodeError::Truncated)
);
}
#[test]
fn decode_rejects_truncated_tag_lens() {
let buf = [1u8, 0, 2, 0, 5, 0];
assert_eq!(EventRef::from_bytes(&buf), Err(DecodeError::Truncated));
}
#[test]
fn decode_rejects_lying_type_len() {
let buf = [10u8, 0, 0, 0];
assert_eq!(EventRef::from_bytes(&buf), Err(DecodeError::Truncated));
}
#[test]
fn decode_rejects_lying_tag_len() {
let mut buf = vec![1u8, 0, 1, 0, 9, 0];
buf.push(b'T');
buf.extend_from_slice(b"short");
assert_eq!(EventRef::from_bytes(&buf), Err(DecodeError::Truncated));
}
#[test]
fn decode_rejects_empty_type() {
let buf = [0u8, 0, 0, 0];
assert_eq!(EventRef::from_bytes(&buf), Err(DecodeError::EmptyType));
}
#[test]
fn decode_rejects_empty_tag() {
let buf = [1u8, 0, 1, 0, 0, 0, b'T'];
assert_eq!(EventRef::from_bytes(&buf), Err(DecodeError::EmptyTag));
}
#[test]
fn decode_rejects_invalid_utf8() {
let buf = [1u8, 0, 0, 0, 0xFF];
assert_eq!(EventRef::from_bytes(&buf), Err(DecodeError::InvalidUtf8));
}
#[test]
fn decode_rejects_unsorted_tags() {
let buf = [1u8, 0, 2, 0, 1, 0, 1, 0, b'T', b'b', b'a'];
assert_eq!(EventRef::from_bytes(&buf), Err(DecodeError::TagsNotSorted));
}
#[test]
fn decode_rejects_duplicate_tags() {
let buf = [1u8, 0, 2, 0, 1, 0, 1, 0, b'T', b'a', b'a'];
assert_eq!(EventRef::from_bytes(&buf), Err(DecodeError::TagsNotSorted));
}
#[test]
fn decode_accepts_manually_built_bytes() {
let buf = [1u8, 0, 2, 0, 1, 0, 1, 0, b'T', b'a', b'b', b'x', b'y'];
let decoded = EventRef::from_bytes(&buf).unwrap();
assert_eq!(decoded.event_type(), "T");
assert_eq!(decoded.tags().collect::<Vec<_>>(), vec!["a", "b"]);
assert_eq!(decoded.data(), b"xy");
}
}