use crate::streaming::frame::StreamFrame;
#[derive(Debug)]
pub enum MpscFrame<T> {
Item(T),
SenderError(String),
Detached,
Dropped(Option<String>),
}
impl<T> MpscFrame<T> {
pub(crate) fn from_stream_frame(frame: StreamFrame<T>) -> Option<Self> {
match frame {
StreamFrame::Heartbeat => None,
StreamFrame::Item(t) => Some(Self::Item(t)),
StreamFrame::SenderError(s) => Some(Self::SenderError(s)),
StreamFrame::Detached | StreamFrame::Finalized => Some(Self::Detached),
StreamFrame::Dropped => Some(Self::Dropped(None)),
StreamFrame::TransportError(s) => Some(Self::Dropped(Some(s))),
}
}
pub(crate) fn is_sender_exit(&self) -> bool {
matches!(self, Self::Detached | Self::Dropped(_))
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn heartbeat_filtered() {
let frame: StreamFrame<u32> = StreamFrame::Heartbeat;
assert!(MpscFrame::from_stream_frame(frame).is_none());
}
#[test]
fn finalized_maps_to_detached() {
let frame: StreamFrame<u32> = StreamFrame::Finalized;
let mapped = MpscFrame::from_stream_frame(frame).expect("mapped");
assert!(matches!(mapped, MpscFrame::Detached));
}
#[test]
fn transport_error_folds_into_dropped() {
let frame: StreamFrame<u32> = StreamFrame::TransportError("boom".into());
let mapped = MpscFrame::from_stream_frame(frame).expect("mapped");
match mapped {
MpscFrame::Dropped(Some(reason)) => assert_eq!(reason, "boom"),
other => panic!("expected Dropped(Some), got {:?}", other),
}
}
#[test]
fn dropped_no_reason() {
let frame: StreamFrame<u32> = StreamFrame::Dropped;
let mapped = MpscFrame::from_stream_frame(frame).expect("mapped");
assert!(matches!(mapped, MpscFrame::Dropped(None)));
}
#[test]
fn item_round_trip() {
let frame: StreamFrame<u32> = StreamFrame::Item(7);
let mapped = MpscFrame::from_stream_frame(frame).expect("mapped");
match mapped {
MpscFrame::Item(v) => assert_eq!(v, 7),
other => panic!("expected Item, got {:?}", other),
}
}
#[test]
fn sender_error_maps_to_sender_error() {
let frame: StreamFrame<u32> = StreamFrame::SenderError("soft".into());
let mapped = MpscFrame::from_stream_frame(frame).expect("mapped");
match mapped {
MpscFrame::SenderError(msg) => assert_eq!(msg, "soft"),
other => panic!("expected SenderError, got {:?}", other),
}
}
#[test]
fn detach_maps_to_detached() {
let frame: StreamFrame<u32> = StreamFrame::Detached;
let mapped = MpscFrame::from_stream_frame(frame).expect("mapped");
assert!(matches!(mapped, MpscFrame::Detached));
}
#[test]
fn is_sender_exit_predicate() {
let item: MpscFrame<u32> = MpscFrame::Item(1);
assert!(!item.is_sender_exit());
let err: MpscFrame<u32> = MpscFrame::SenderError("x".into());
assert!(!err.is_sender_exit());
let detached: MpscFrame<u32> = MpscFrame::Detached;
assert!(detached.is_sender_exit());
let dropped_none: MpscFrame<u32> = MpscFrame::Dropped(None);
assert!(dropped_none.is_sender_exit());
let dropped_some: MpscFrame<u32> = MpscFrame::Dropped(Some("why".into()));
assert!(dropped_some.is_sender_exit());
}
}