use bytes::Bytes;
use serde::{Deserialize, Serialize};
use crate::error::{MessageReject, Result, RiftError};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Event {
pub event: String,
pub message_id: String,
pub schema: String,
pub payload: serde_json::Value,
pub dedupe_key: Option<String>,
pub ordering_key: Option<String>,
pub ttl_ms: Option<u32>,
}
impl Event {
pub fn new(
event: impl Into<String>,
message_id: impl Into<String>,
schema: impl Into<String>,
payload: serde_json::Value,
) -> Self {
Self {
event: event.into(),
message_id: message_id.into(),
schema: schema.into(),
payload,
dedupe_key: None,
ordering_key: None,
ttl_ms: None,
}
}
pub fn size_hint(&self) -> usize {
self.event.len()
+ self.message_id.len()
+ self.schema.len()
+ serde_json::to_vec(&self.payload)
.map(|v| v.len())
.unwrap_or(0)
}
}
pub fn encode_event_body(e: &Event) -> Result<Bytes> {
Ok(Bytes::from(serde_json::to_vec(e)?))
}
pub fn decode_event_body(bytes: &[u8]) -> Result<Event> {
if bytes.is_empty() {
return Err(RiftError::Message(MessageReject::Rejected(
"empty event payload".into(),
)));
}
Ok(serde_json::from_slice(bytes)?)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn round_trip() {
let e = Event::new(
"chat.message.created",
"01HZZZZZZZZZZZZZZZZZZZZZZ",
"chat.message.created@1.0",
serde_json::json!({"text": "hi"}),
);
let bytes = encode_event_body(&e).unwrap();
let back = decode_event_body(&bytes).unwrap();
assert_eq!(back.event, e.event);
assert_eq!(back.message_id, e.message_id);
}
#[test]
fn decode_empty_fails() {
let r = decode_event_body(&[]);
assert!(r.is_err());
}
}