Skip to content
File

Blob: firmware/vendor/str0m/src/streams/rtx_cache_buf.rs

rust387 lines
1#![allow(missing_docs)]
2 
3use std::mem;
4use std::time::{Duration, Instant};
5 
6use crate::util::already_happened;
7 
8/// Fixed size buffer that evicts the oldest entries based on time.
9///
10/// The buffer is ringbuffer-esque in that all elements are inserted with a position that is modulo to an
11/// insert index into the buffer. That makes lookups very fast since they are fixed offsets.
12///
13/// The buffer evicts values based on time. If the size of the buffer is too small and would
14/// overwrite an entry that has not been evicted due to age, the buffer grows.
15#[derive(Debug)]
16pub struct EvictingBuffer<T> {
17 buf: Vec<Option<Entry<T>>>,
18 /// How long to keep entries for.
19 max_age: Duration,
20 /// The maximum size allowed to grow to. Once this sized is reached, new pushed
21 /// entries will overwrite older entries even if they haven't reached max_age.
22 max_size: usize,
23 /// Last inserted position
24 last_position: Option<u64>,
25 /// Next element to evict.
26 next_evict: Option<u64>,
27 // Last timeout when we evicted elements.
28 last_timeout: Instant,
29}
30 
31/// Entry in the buffer to keep track of position and time separately.
32#[derive(Debug)]
33struct Entry<T> {
34 /// Position is offset modulus the buffer size.
35 position: u64,
36 /// Entry time. Used for eviction.
37 timestamp: Instant,
38 /// The value held at this entry.
39 value: T,
40}
41 
42// We don't want to require that T is Clone/Copy, which means we
43// must do this instead of using the vec![] macro.
44fn prepare_buf<T>(len: usize) -> Vec<Option<T>> {
45 let mut buf = Vec::with_capacity(len);
46 for _ in 0..len {
47 buf.push(None);
48 }
49 buf
50}
51 
52impl<T> EvictingBuffer<T> {
53 /// Creates a new buffer with an initial size.
54 pub fn new(initial_size: usize, max_age: Duration, max_size: usize) -> Self {
55 Self {
56 buf: prepare_buf(initial_size),
57 max_age,
58 max_size,
59 last_position: None,
60 next_evict: None,
61 last_timeout: already_happened(),
62 }
63 }
64 
65 fn index_for_position(&self, position: u64) -> usize {
66 (position % self.buf.len() as u64) as usize
67 }
68 
69 #[inline(always)]
70 fn is_inert(&self) -> bool {
71 self.buf.is_empty() || self.max_age.is_zero()
72 }
73 
74 /// Push a new entry.
75 ///
76 /// Position is an increasing sequence number. The sequence can be out of order.
77 pub fn push(&mut self, position: u64, timestamp: Instant, value: T) {
78 if self.is_inert() {
79 return;
80 }
81 
82 if timestamp < self.last_timeout {
83 // Value is already considered evicted.
84 return;
85 }
86 
87 let next_evict = if let Some(v) = self.next_evict {
88 v
89 } else {
90 // First ever cached value sets the initial evict position.
91 self.next_evict = Some(position);
92 position
93 };
94 
95 if position < next_evict {
96 // Do not cache values preceding evict position.
97 return;
98 }
99 
100 let mut index = self.index_for_position(position);
101 
102 if let Some(entry) = &self.buf[index] {
103 // If the position is exactly the same, we allow it since it's
104 // replacing the current T value. If position differs, we've
105 // wrapped around.
106 if entry.position != position {
107 // Make space to continue.
108 self.grow();
109 index = self.index_for_position(position);
110 }
111 }
112 self.last_position = Some(position);
113 
114 self.buf[index] = Some(Entry {
115 position,
116 timestamp,
117 value,
118 });
119 }
120 
121 /// Get the entry for the previously inserted position.
122 #[allow(unused)]
123 pub fn get(&self, position: u64) -> Option<&T> {
124 if self.is_inert() {
125 return None;
126 }
127 
128 let index = (position % self.buf.len() as u64) as usize;
129 if let Some(entry) = &self.buf[index] {
130 if entry.position == position {
131 return Some(&entry.value);
132 }
133 }
134 None
135 }
136 
137 pub fn get_mut(&mut self, position: u64) -> Option<&mut T> {
138 if self.is_inert() {
139 return None;
140 }
141 
142 let index = (position % self.buf.len() as u64) as usize;
143 if let Some(entry) = &mut self.buf[index] {
144 if entry.position == position {
145 return Some(&mut entry.value);
146 }
147 }
148 None
149 }
150 
151 pub fn maybe_evict(&mut self, now: Instant) {
152 if self.is_inert() {
153 return;
154 }
155 
156 if now < self.last_timeout {
157 // Time cannot go backwards.
158 return;
159 }
160 self.last_timeout = now;
161 
162 self.evict(now);
163 }
164 
165 fn evict(&mut self, now: Instant) {
166 let Some(start_position) = self.next_evict else {
167 // Before first element been pushed.
168 return;
169 };
170 
171 let mut position = start_position;
172 let start_index = self.index_for_position(position);
173 
174 loop {
175 let index = self.index_for_position(position);
176 
177 if index == start_index && position > start_position {
178 // looped around without finding the end. Means there are no elements.
179 break;
180 }
181 
182 let Some(entry) = &self.buf[index] else {
183 // No entry means we might have a gap. We can't break on gaps because
184 // there might be entries to evict later. Skip the position and check
185 // until we find a timestamp.
186 position += 1;
187 continue;
188 };
189 
190 let age = now.saturating_duration_since(entry.timestamp);
191 
192 if age > self.max_age {
193 // Evict.
194 self.buf[index] = None;
195 } else {
196 // We assume entries are roughly in time order (some jumble is allowed).
197 // Once we reach an element we should not evict, stop.
198 break;
199 }
200 
201 position += 1;
202 }
203 
204 self.next_evict = Some(position);
205 }
206 
207 fn grow(&mut self) {
208 if self.buf.len() >= self.max_size {
209 // No growing.
210 return;
211 }
212 
213 // This is the new sized buffer. We can make other strategies for growing.
214 let new_size = self
215 .max_size
216 .min((self.buf.len() * 133) / 100)
217 .max(self.buf.len() + 1);
218 
219 let old_buffer = mem::replace(&mut self.buf, prepare_buf(new_size));
220 
221 // Move all entries over to the new buffer. Changing the buffer size might alter the
222 // index position. However max_position and next_evict are already containing positions
223 // not index, which means they stay correct.
224 for e in old_buffer.into_iter().flatten() {
225 let index = self.index_for_position(e.position);
226 self.buf[index] = Some(e);
227 }
228 }
229 
230 pub fn contains(&self, position: u64) -> bool {
231 if self.is_inert() {
232 return false;
233 }
234 
235 let index = self.index_for_position(position);
236 self.buf[index].is_some()
237 }
238 
239 pub fn last_position(&self) -> Option<u64> {
240 if self.is_inert() {
241 return None;
242 }
243 
244 self.last_position
245 }
246 
247 pub fn last(&self) -> Option<&T> {
248 if self.is_inert() {
249 return None;
250 }
251 
252 let last = self.last_position?;
253 self.get(last)
254 }
255 
256 pub fn clear(&mut self) {
257 for i in 0..self.buf.len() {
258 self.buf[i] = None;
259 }
260 self.last_position = None;
261 self.next_evict = None;
262 self.last_timeout = already_happened();
263 }
264}
265 
266#[cfg(test)]
267mod test {
268 use super::*;
269 
270 #[test]
271 fn push_and_get() {
272 let mut buf = EvictingBuffer::new(1, Duration::from_secs(10), 10);
273 let now = Instant::now();
274 
275 buf.push(5, now, 'A');
276 
277 assert_eq!(buf.index_for_position(5), buf.index_for_position(3));
278 assert_eq!(buf.get(5), Some(&'A'));
279 assert_eq!(buf.get(3), None); // modulo to same index
280 }
281 
282 #[test]
283 fn push_over_capacity() {
284 let mut buf = EvictingBuffer::new(2, Duration::from_secs(10), 10);
285 let now = Instant::now();
286 
287 buf.push(5, now + Duration::from_secs(0), 'A');
288 buf.push(6, now + Duration::from_secs(1), 'B');
289 buf.push(7, now + Duration::from_secs(2), 'C');
290 
291 // The size is 2, we should have grown to accomodate.
292 assert_eq!(buf.get(5), Some(&'A'));
293 assert_eq!(buf.get(6), Some(&'B'));
294 assert_eq!(buf.get(7), Some(&'C'));
295 }
296 
297 #[test]
298 fn push_before_next_evict() {
299 let mut buf = EvictingBuffer::new(2, Duration::from_secs(10), 10);
300 let now = Instant::now();
301 
302 buf.push(6, now + Duration::from_secs(0), 'B');
303 assert_eq!(buf.next_evict, Some(6));
304 
305 // Before the initial evict position, thus ignored.
306 buf.push(5, now + Duration::from_secs(1), 'A');
307 
308 assert_eq!(buf.get(5), None);
309 }
310 
311 #[test]
312 fn evict_oldest() {
313 let mut buf = EvictingBuffer::new(2, Duration::from_secs(10), 10);
314 let now = Instant::now();
315 
316 buf.push(5, now + Duration::from_secs(0), 'A');
317 buf.push(6, now + Duration::from_secs(1), 'B');
318 
319 // Nothing should go
320 buf.maybe_evict(now + Duration::from_secs(1));
321 assert_eq!(buf.get(5), Some(&'A'));
322 assert_eq!(buf.get(6), Some(&'B'));
323 
324 // One entry gone.
325 buf.maybe_evict(now + Duration::from_secs(11));
326 assert_eq!(buf.get(5), None);
327 assert_eq!(buf.get(6), Some(&'B'));
328 }
329 
330 #[test]
331 fn evict_with_gap() {
332 let mut buf = EvictingBuffer::new(4, Duration::from_secs(10), 10);
333 let now = Instant::now();
334 
335 buf.push(5, now + Duration::from_secs(0), 'A');
336 // GAP
337 buf.push(7, now + Duration::from_secs(2), 'C');
338 buf.push(8, now + Duration::from_secs(3), 'D');
339 
340 // Should evict A and C
341 buf.maybe_evict(now + Duration::from_secs(13));
342 
343 assert_eq!(buf.get(5), None);
344 assert_eq!(buf.get(7), None);
345 assert_eq!(buf.get(8), Some(&'D'));
346 }
347 
348 #[test]
349 fn evict_all() {
350 let mut buf = EvictingBuffer::new(4, Duration::from_secs(10), 10);
351 let now = Instant::now();
352 
353 buf.push(5, now + Duration::from_secs(0), 'A');
354 buf.push(6, now + Duration::from_secs(1), 'B');
355 
356 // Should evict A and B
357 buf.maybe_evict(now + Duration::from_secs(12));
358 
359 assert_eq!(buf.get(5), None);
360 assert_eq!(buf.get(6), None);
361 }
362 
363 fn buffer_cmp(b: &EvictingBuffer<char>) -> Vec<Option<char>> {
364 b.buf.iter().map(|e| e.as_ref().map(|v| v.value)).collect()
365 }
366 
367 #[test]
368 fn grow_and_reindex() {
369 let mut buf = EvictingBuffer::new(2, Duration::from_secs(10), 10);
370 let now = Instant::now();
371 
372 buf.push(2, now + Duration::from_secs(0), 'A');
373 buf.push(3, now + Duration::from_secs(1), 'B');
374 
375 assert_eq!(buffer_cmp(&buf), &[Some('A'), Some('B')]);
376 
377 // overwrites 2, thus grows
378 buf.push(4, now + Duration::from_secs(2), 'C');
379 
380 assert_eq!(buffer_cmp(&buf), &[Some('B'), Some('C'), Some('A')]);
381 
382 assert_eq!(buf.get(2), Some(&'A'));
383 assert_eq!(buf.get(3), Some(&'B'));
384 assert_eq!(buf.get(4), Some(&'C'));
385 }
386}