use std::fmt;
use bytes::Bytes;
use serde::{Deserialize, Serialize};
use crate::frame::{Codec, FrameFlags, FrameType, Priority};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Frame {
pub version: u16,
pub frame_id: u64,
pub frame_type: FrameType,
pub flags: FrameFlags,
pub codec: Codec,
pub session_id: Option<String>,
pub stream_id: Option<String>,
pub topic: Option<String>,
pub event: Option<String>,
pub message_id: Option<String>,
pub correlation_id: Option<String>,
pub trace_id: Option<String>,
pub timestamp: i64,
pub ttl_ms: Option<u32>,
pub priority: Option<Priority>,
pub payload: Option<Bytes>,
}
impl Default for Frame {
fn default() -> Self {
Self {
version: 0,
frame_id: 0,
frame_type: FrameType::Control,
flags: FrameFlags::empty(),
codec: Codec::Json,
session_id: None,
stream_id: None,
topic: None,
event: None,
message_id: None,
correlation_id: None,
trace_id: None,
timestamp: 0,
ttl_ms: None,
priority: None,
payload: None,
}
}
}
impl Frame {
pub fn control() -> Self {
Self {
frame_type: FrameType::Control,
..Self::default()
}
}
pub fn data() -> Self {
Self {
frame_type: FrameType::Data,
..Self::default()
}
}
pub fn ack() -> Self {
Self {
frame_type: FrameType::Ack,
..Self::default()
}
}
pub fn flow() -> Self {
Self {
frame_type: FrameType::Flow,
..Self::default()
}
}
pub fn error() -> Self {
Self {
frame_type: FrameType::Error,
..Self::default()
}
}
pub fn requires_ack(&self) -> bool {
self.flags.contains(FrameFlags::REQUIRES_ACK)
}
pub fn is_replay(&self) -> bool {
self.flags.contains(FrameFlags::REPLAYED)
}
pub fn mark_replay(&mut self) {
self.flags.set(FrameFlags::REPLAYED);
}
pub fn is_snapshot(&self) -> bool {
self.flags.contains(FrameFlags::SNAPSHOT)
}
pub fn mark_snapshot(&mut self) {
self.flags.set(FrameFlags::SNAPSHOT);
}
}
impl fmt::Display for Frame {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(
f,
"Frame[id={} type={} codec={} topic={} event={} msg_id={} corr_id={} flags={} payload={}B]",
self.frame_id,
self.frame_type,
self.codec,
self.topic.as_deref().unwrap_or("-"),
self.event.as_deref().unwrap_or("-"),
self.message_id.as_deref().unwrap_or("-"),
self.correlation_id.as_deref().unwrap_or("-"),
self.flags,
self.payload.as_ref().map(|p| p.len()).unwrap_or(0),
)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn flag_round_trip() {
let mut f = FrameFlags::empty();
f.set(FrameFlags::COMPRESSED);
f.set(FrameFlags::REQUIRES_ACK);
assert!(f.contains(FrameFlags::COMPRESSED));
assert!(f.contains(FrameFlags::REQUIRES_ACK));
assert!(!f.contains(FrameFlags::ENCRYPTED));
assert_eq!(f.bits(), FrameFlags::COMPRESSED | FrameFlags::REQUIRES_ACK);
}
#[test]
fn priority_default() {
assert_eq!(Priority::default(), Priority::Normal);
}
#[test]
fn frame_type_tag() {
for t in [
FrameType::Control,
FrameType::Data,
FrameType::Ack,
FrameType::Flow,
FrameType::Error,
] {
assert_eq!(FrameType::from_tag(t.tag()), Some(t));
}
assert_eq!(FrameType::from_tag(b'X'), None);
}
#[test]
fn codec_tag() {
assert_eq!(Codec::from_tag(Codec::Json.tag()), Some(Codec::Json));
assert_eq!(Codec::from_tag(Codec::Cbor.tag()), Some(Codec::Cbor));
assert_eq!(Codec::from_tag(b'?'), None);
}
#[test]
fn frame_constructors() {
assert_eq!(Frame::control().frame_type, FrameType::Control);
assert_eq!(Frame::data().frame_type, FrameType::Data);
assert_eq!(Frame::ack().frame_type, FrameType::Ack);
assert_eq!(Frame::flow().frame_type, FrameType::Flow);
assert_eq!(Frame::error().frame_type, FrameType::Error);
}
}