Skip to main content

subtr_actor/stats/calculators/
event_stream.rs

1use std::ops::{Deref, DerefMut};
2
3use serde::{Serialize, Serializer};
4
5/// Append-only buffer of emitted events exposing all events and newly added ones.
6#[derive(Debug, Clone, PartialEq)]
7pub struct EventStream<E> {
8    events: Vec<E>,
9    update_start: usize,
10}
11
12impl<E> Default for EventStream<E> {
13    fn default() -> Self {
14        Self {
15            events: Vec::new(),
16            update_start: 0,
17        }
18    }
19}
20
21impl<E> EventStream<E> {
22    pub fn new() -> Self {
23        Self::default()
24    }
25
26    pub fn from_vec(events: Vec<E>) -> Self {
27        Self {
28            update_start: 0,
29            events,
30        }
31    }
32
33    pub fn begin_update(&mut self) {
34        self.update_start = self.events.len();
35    }
36
37    pub fn push(&mut self, event: E) {
38        self.events.push(event);
39    }
40
41    pub fn extend(&mut self, events: impl IntoIterator<Item = E>) {
42        self.events.extend(events);
43    }
44
45    pub fn replace_all_assuming_append_only(&mut self, events: Vec<E>) {
46        self.update_start = self.events.len().min(events.len());
47        self.events = events;
48    }
49
50    pub fn all(&self) -> &[E] {
51        &self.events
52    }
53
54    pub fn new_events(&self) -> &[E] {
55        &self.events[self.update_start..]
56    }
57
58    pub fn into_vec(self) -> Vec<E> {
59        self.events
60    }
61}
62
63impl<E> From<Vec<E>> for EventStream<E> {
64    fn from(events: Vec<E>) -> Self {
65        Self::from_vec(events)
66    }
67}
68
69impl<E: Serialize> Serialize for EventStream<E> {
70    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
71    where
72        S: Serializer,
73    {
74        self.events.serialize(serializer)
75    }
76}
77
78impl<'a, E> IntoIterator for &'a EventStream<E> {
79    type Item = &'a E;
80    type IntoIter = std::slice::Iter<'a, E>;
81
82    fn into_iter(self) -> Self::IntoIter {
83        self.events.iter()
84    }
85}
86
87impl<E> Deref for EventStream<E> {
88    type Target = [E];
89
90    fn deref(&self) -> &Self::Target {
91        self.all()
92    }
93}
94
95impl<E> DerefMut for EventStream<E> {
96    fn deref_mut(&mut self) -> &mut Self::Target {
97        &mut self.events
98    }
99}
100
101#[cfg(test)]
102mod tests {
103    use super::*;
104
105    #[test]
106    fn new_events_tracks_events_after_update_start() {
107        let mut stream = EventStream::new();
108        stream.push(1);
109        stream.push(2);
110        assert_eq!(stream.all(), &[1, 2]);
111        assert_eq!(stream.new_events(), &[1, 2]);
112
113        stream.begin_update();
114        assert!(stream.new_events().is_empty());
115        stream.push(3);
116        stream.extend([4, 5]);
117        assert_eq!(stream.all(), &[1, 2, 3, 4, 5]);
118        assert_eq!(stream.new_events(), &[3, 4, 5]);
119    }
120
121    #[test]
122    fn append_only_replacement_exposes_suffix_after_previous_len() {
123        let mut stream = EventStream::from_vec(vec![1, 2]);
124        stream.replace_all_assuming_append_only(vec![1, 2, 3, 4]);
125        assert_eq!(stream.all(), &[1, 2, 3, 4]);
126        assert_eq!(stream.new_events(), &[3, 4]);
127    }
128}