templar_proxy_oracle_kernel/proxy/circuit_breaker/
ring_buffer.rs1use 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}