File
Blob: firmware/vendor/str0m/src/streams/rtx_cache_buf.rs
| 1 | #![allow(missing_docs)] |
| 2 | |
| 3 | use std::mem; |
| 4 | use std::time::{Duration, Instant}; |
| 5 | |
| 6 | use 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)] |
| 16 | pub 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)] |
| 33 | struct 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. |
| 44 | fn 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 | |
| 52 | impl<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)] |
| 267 | mod 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 | } |