templar_proxy_oracle_kernel/proxy/circuit_breaker/
ring_buffer.rs

1use alloc::collections::VecDeque;
2
3#[cfg(any(feature = "borsh", feature = "schemars"))]
4use alloc::{format, string::ToString};
5
6serialize! {
7    #[derive(Debug, Clone, PartialEq, Eq)]
8    pub struct UncheckedRingBuffer<T> {
9        pub capacity: u32,
10        pub entries: VecDeque<T>,
11    }
12}
13
14#[cfg_attr(
15    feature = "serde",
16    derive(::serde::Deserialize, ::serde::Serialize),
17    serde(
18        try_from = "UncheckedRingBuffer<T>",
19        into = "UncheckedRingBuffer<T>",
20        bound(
21            serialize = "T: Clone + ::serde::Serialize",
22            deserialize = "T: ::serde::Deserialize<'de>"
23        )
24    )
25)]
26#[cfg_attr(
27    feature = "schemars",
28    derive(::schemars::JsonSchema),
29    schemars(transparent)
30)]
31#[cfg_attr(
32    feature = "borsh",
33    derive(::borsh::BorshSerialize, ::borsh::BorshSchema)
34)]
35#[derive(Debug, Clone, PartialEq, Eq)]
36pub struct RingBuffer<T>(UncheckedRingBuffer<T>);
37
38#[derive(Debug, Clone, Copy, PartialEq, Eq)]
39pub enum RingBufferParseError {
40    EntriesExceedCapacity,
41}
42
43impl core::fmt::Display for RingBufferParseError {
44    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
45        match self {
46            Self::EntriesExceedCapacity => write!(f, "entries exceed ring buffer capacity"),
47        }
48    }
49}
50
51impl<T> TryFrom<UncheckedRingBuffer<T>> for RingBuffer<T> {
52    type Error = RingBufferParseError;
53
54    fn try_from(value: UncheckedRingBuffer<T>) -> Result<Self, Self::Error> {
55        if value.entries.len() > value.capacity as usize {
56            return Err(RingBufferParseError::EntriesExceedCapacity);
57        }
58        Ok(Self(value))
59    }
60}
61
62impl<T> From<RingBuffer<T>> for UncheckedRingBuffer<T> {
63    fn from(value: RingBuffer<T>) -> Self {
64        value.0
65    }
66}
67
68impl<T> RingBuffer<T> {
69    #[must_use]
70    pub fn new(capacity: u32) -> Self {
71        Self(UncheckedRingBuffer {
72            capacity,
73            entries: VecDeque::new(),
74        })
75    }
76
77    pub fn push(&mut self, item: T) {
78        let capacity = self.0.capacity as usize;
79        if capacity == 0 {
80            return;
81        }
82
83        if self.0.entries.len() == capacity {
84            self.0.entries.pop_front();
85        }
86
87        self.0.entries.push_back(item);
88    }
89
90    pub fn set_capacity(&mut self, capacity: u32) {
91        self.0.capacity = capacity;
92        let capacity = capacity as usize;
93        let excess = self.0.entries.len().saturating_sub(capacity);
94        for _ in 0..excess {
95            self.0.entries.pop_front();
96        }
97    }
98
99    pub fn clear(&mut self) {
100        self.0.entries.clear();
101    }
102
103    #[must_use]
104    pub fn capacity(&self) -> u32 {
105        self.0.capacity
106    }
107
108    #[must_use]
109    pub fn last(&self) -> Option<&T> {
110        self.0.entries.back()
111    }
112
113    #[must_use]
114    pub fn len(&self) -> usize {
115        self.0.entries.len()
116    }
117
118    #[must_use]
119    pub fn is_empty(&self) -> bool {
120        self.0.entries.is_empty()
121    }
122
123    #[must_use]
124    pub fn get(&self, index: usize) -> Option<&T> {
125        self.0.entries.get(index)
126    }
127
128    pub fn iter(&self) -> alloc::collections::vec_deque::Iter<'_, T> {
129        self.0.entries.iter()
130    }
131}
132
133impl<'a, T> IntoIterator for &'a RingBuffer<T> {
134    type Item = &'a T;
135    type IntoIter = alloc::collections::vec_deque::Iter<'a, T>;
136
137    fn into_iter(self) -> Self::IntoIter {
138        self.iter()
139    }
140}
141
142#[cfg(feature = "borsh")]
143impl<T: ::borsh::BorshDeserialize> ::borsh::BorshDeserialize for RingBuffer<T> {
144    fn deserialize_reader<Reader: ::borsh::io::Read>(
145        reader: &mut Reader,
146    ) -> ::borsh::io::Result<Self> {
147        let unchecked =
148            <UncheckedRingBuffer<T> as ::borsh::BorshDeserialize>::deserialize_reader(reader)?;
149        unchecked.try_into().map_err(|_| {
150            ::borsh::io::Error::new(
151                ::borsh::io::ErrorKind::InvalidData,
152                "could not parse ring buffer",
153            )
154        })
155    }
156}
157
158#[cfg(test)]
159mod tests {
160    use alloc::{vec, vec::Vec};
161
162    use super::*;
163
164    #[test]
165    fn push_noops_when_capacity_is_zero() {
166        let mut buffer = RingBuffer::new(0);
167
168        buffer.push(1);
169        buffer.push(2);
170
171        assert!(buffer.is_empty());
172        assert_eq!(buffer.last(), None);
173    }
174
175    #[test]
176    fn push_preserves_insertion_order_before_capacity() {
177        let mut buffer = RingBuffer::new(3);
178
179        buffer.push(1);
180        buffer.push(2);
181
182        assert_eq!(buffer.len(), 2);
183        assert_eq!(buffer.iter().collect::<Vec<_>>(), vec![&1, &2]);
184        assert_eq!(buffer.last(), Some(&2));
185    }
186
187    #[test]
188    fn push_drops_oldest_entry_at_capacity() {
189        let mut buffer = RingBuffer::new(3);
190
191        buffer.push(1);
192        buffer.push(2);
193        buffer.push(3);
194        buffer.push(4);
195        buffer.push(5);
196
197        assert_eq!(buffer.len(), 3);
198        assert_eq!(buffer.iter().collect::<Vec<_>>(), vec![&3, &4, &5]);
199        assert_eq!(buffer.last(), Some(&5));
200    }
201
202    #[test]
203    fn parse_rejects_entries_exceeding_capacity() {
204        assert_eq!(
205            RingBuffer::try_from(UncheckedRingBuffer {
206                capacity: 1,
207                entries: vec![1, 2].into(),
208            }),
209            Err(RingBufferParseError::EntriesExceedCapacity)
210        );
211    }
212
213    #[cfg(feature = "borsh")]
214    #[test]
215    fn borsh_rejects_entries_exceeding_capacity() {
216        let unchecked = UncheckedRingBuffer {
217            capacity: 1,
218            entries: vec![1_u32, 2].into(),
219        };
220        let bytes = borsh::to_vec(&unchecked).unwrap();
221
222        assert!(borsh::from_slice::<RingBuffer<u32>>(&bytes).is_err());
223    }
224
225    #[cfg(feature = "serde")]
226    #[test]
227    fn serde_rejects_entries_exceeding_capacity() {
228        let unchecked = UncheckedRingBuffer {
229            capacity: 1,
230            entries: vec![1_u32, 2].into(),
231        };
232        let bytes = serde_json::to_vec(&unchecked).unwrap();
233
234        assert!(serde_json::from_slice::<RingBuffer<u32>>(&bytes).is_err());
235    }
236
237    #[cfg(feature = "serde")]
238    #[test]
239    fn serde_serializes_like_unchecked_representation() {
240        let mut buffer = RingBuffer::new(2);
241        buffer.push(1_u32);
242        buffer.push(2);
243        let unchecked = UncheckedRingBuffer::from(buffer.clone());
244
245        assert_eq!(
246            serde_json::to_value(&buffer).unwrap(),
247            serde_json::to_value(&unchecked).unwrap()
248        );
249    }
250
251    #[test]
252    fn set_capacity_grow_preserves_existing_entries() {
253        let mut buffer = RingBuffer::new(2);
254        buffer.push(1);
255        buffer.push(2);
256
257        buffer.set_capacity(4);
258        buffer.push(3);
259        buffer.push(4);
260
261        assert_eq!(buffer.iter().collect::<Vec<_>>(), vec![&1, &2, &3, &4]);
262    }
263
264    #[test]
265    fn set_capacity_shrink_keeps_newest_entries() {
266        let mut buffer = RingBuffer::new(4);
267        for item in 1..=4 {
268            buffer.push(item);
269        }
270
271        buffer.set_capacity(2);
272
273        assert_eq!(buffer.len(), 2);
274        assert_eq!(buffer.iter().collect::<Vec<_>>(), vec![&3, &4]);
275        assert_eq!(buffer.last(), Some(&4));
276    }
277
278    #[test]
279    fn set_capacity_zero_clears_entries_and_blocks_future_pushes() {
280        let mut buffer = RingBuffer::new(2);
281        buffer.push(1);
282        buffer.push(2);
283
284        buffer.set_capacity(0);
285        buffer.push(3);
286
287        assert!(buffer.is_empty());
288    }
289
290    #[test]
291    fn set_capacity_after_zero_allows_future_entries() {
292        let mut buffer = RingBuffer::new(0);
293        buffer.push(1);
294
295        buffer.set_capacity(2);
296        buffer.push(2);
297        buffer.push(3);
298
299        assert_eq!(buffer.iter().collect::<Vec<_>>(), vec![&2, &3]);
300    }
301
302    #[test]
303    fn set_capacity_to_same_value_is_stable() {
304        let mut buffer = RingBuffer::new(3);
305        for item in 1..=3 {
306            buffer.push(item);
307        }
308        buffer.set_capacity(3);
309
310        assert_eq!(buffer.iter().collect::<Vec<_>>(), vec![&1, &2, &3]);
311    }
312}