Skip to main content

sim_lib_stream_host/
cassette.rs

1//! Host callback cassette and replay support.
2
3use sim_kernel::Result;
4use sim_lib_stream_core::{
5    PushResult, StreamCassette, StreamItem, StreamMetadata, StreamPacket, StreamStats,
6    TransportProfile,
7};
8
9use crate::HostCallbackQueue;
10
11/// Deterministic recording of host callback items.
12#[derive(Clone, Debug, Default, PartialEq, Eq)]
13pub struct HostCallbackCassette {
14    items: Vec<StreamItem>,
15}
16
17/// Outcome counts from deterministic host callback cassette replay.
18#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
19pub struct HostCallbackReplayReport {
20    /// Callback items accepted by the target queue.
21    pub accepted: u64,
22    /// Callback items dropped by the target queue's drop-newest policy.
23    pub dropped_newest: u64,
24    /// Callback items evicted from the target queue's drop-oldest policy.
25    pub dropped_oldest: u64,
26    /// Callback items rejected by the target queue's overflow policy.
27    pub rejected: u64,
28    /// Callback items refused because the target queue was closed.
29    pub closed: u64,
30}
31
32impl HostCallbackCassette {
33    /// Creates an empty cassette.
34    pub fn new() -> Self {
35        Self::default()
36    }
37
38    /// Creates a cassette from callback items.
39    pub fn from_items(items: Vec<StreamItem>) -> Self {
40        Self { items }
41    }
42
43    /// Appends one callback item.
44    pub fn record_item(&mut self, item: StreamItem) {
45        self.items.push(item);
46    }
47
48    /// Appends one packet with no explicit ticks.
49    pub fn record_packet(&mut self, packet: StreamPacket) {
50        self.record_item(StreamItem::new(packet));
51    }
52
53    /// Returns recorded callback items.
54    pub fn items(&self) -> &[StreamItem] {
55        &self.items
56    }
57
58    /// Converts a shared stream cassette into a callback cassette.
59    pub fn from_stream_cassette(cassette: &StreamCassette) -> Result<Self> {
60        Ok(Self {
61            items: cassette.items()?,
62        })
63    }
64
65    /// Converts recorded callbacks into the shared stream cassette format.
66    pub fn to_stream_cassette(
67        &self,
68        metadata: StreamMetadata,
69        profile: TransportProfile,
70    ) -> Result<StreamCassette> {
71        StreamCassette::from_items(
72            metadata,
73            self.items.clone(),
74            profile,
75            StreamStats {
76                pushed: self.items.len() as u64,
77                accepted: self.items.len() as u64,
78                ..StreamStats::default()
79            },
80        )
81    }
82
83    /// Replays every recorded item into a host callback queue.
84    pub fn replay(&self, queue: &HostCallbackQueue) -> Result<HostCallbackReplayReport> {
85        let mut report = HostCallbackReplayReport::default();
86        for item in &self.items {
87            report.record(queue.callback_item(item.clone())?);
88        }
89        Ok(report)
90    }
91}
92
93impl HostCallbackReplayReport {
94    fn record(&mut self, result: PushResult) {
95        match result {
96            PushResult::Accepted => self.accepted = self.accepted.saturating_add(1),
97            PushResult::DroppedNewest(_) => {
98                self.dropped_newest = self.dropped_newest.saturating_add(1);
99            }
100            PushResult::DroppedOldest(_) => {
101                self.dropped_oldest = self.dropped_oldest.saturating_add(1);
102            }
103            PushResult::Rejected(_) => self.rejected = self.rejected.saturating_add(1),
104            PushResult::Closed(_) => self.closed = self.closed.saturating_add(1),
105        }
106    }
107}