use sim_kernel::Result;
use sim_lib_stream_core::{
PushResult, StreamCassette, StreamItem, StreamMetadata, StreamPacket, StreamStats,
TransportProfile,
};
use crate::HostCallbackQueue;
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct HostCallbackCassette {
items: Vec<StreamItem>,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct HostCallbackReplayReport {
pub accepted: u64,
pub dropped_newest: u64,
pub dropped_oldest: u64,
pub rejected: u64,
pub closed: u64,
}
impl HostCallbackCassette {
pub fn new() -> Self {
Self::default()
}
pub fn from_items(items: Vec<StreamItem>) -> Self {
Self { items }
}
pub fn record_item(&mut self, item: StreamItem) {
self.items.push(item);
}
pub fn record_packet(&mut self, packet: StreamPacket) {
self.record_item(StreamItem::new(packet));
}
pub fn items(&self) -> &[StreamItem] {
&self.items
}
pub fn from_stream_cassette(cassette: &StreamCassette) -> Result<Self> {
Ok(Self {
items: cassette.items()?,
})
}
pub fn to_stream_cassette(
&self,
metadata: StreamMetadata,
profile: TransportProfile,
) -> Result<StreamCassette> {
StreamCassette::from_items(
metadata,
self.items.clone(),
profile,
StreamStats {
pushed: self.items.len() as u64,
accepted: self.items.len() as u64,
..StreamStats::default()
},
)
}
pub fn replay(&self, queue: &HostCallbackQueue) -> Result<HostCallbackReplayReport> {
let mut report = HostCallbackReplayReport::default();
for item in &self.items {
report.record(queue.callback_item(item.clone())?);
}
Ok(report)
}
}
impl HostCallbackReplayReport {
fn record(&mut self, result: PushResult) {
match result {
PushResult::Accepted => self.accepted = self.accepted.saturating_add(1),
PushResult::DroppedNewest(_) => {
self.dropped_newest = self.dropped_newest.saturating_add(1);
}
PushResult::DroppedOldest(_) => {
self.dropped_oldest = self.dropped_oldest.saturating_add(1);
}
PushResult::Rejected(_) => self.rejected = self.rejected.saturating_add(1),
PushResult::Closed(_) => self.closed = self.closed.saturating_add(1),
}
}
}