oxedyne/fe2o3/fe2o3_data/src/ring.rs
13.9 KiB, 41 runs
created by r1870400018:272, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | use oxedyne_fe2o3_core::{ |
| 2 | prelude::*, |
| 3 | mem::Extract, |
| 4 | }; |
| 5 | |
| 6 | use oxedyne_fe2o3_jdat::{ |
| 7 | prelude::*, |
| 8 | try_extract_tup2dat, |
| 9 | tup2dat, |
| 10 | }; |
| 11 | |
| 12 | use std::{ |
| 13 | fmt, |
| 14 | time::{ |
| 15 | Duration, |
| 16 | SystemTime, |
| 17 | }, |
| 18 | }; |
| 19 | |
| 20 | /// A generic ring buffer. |
| 21 | /// |
| 22 | /// |
| 23 | /// +---+---+---+---+---+---+---+---+---+---+---+---+ |
| 24 | /// | 0 | 1 | 2 | 3 | 4 | 5 | 6 | 7 | 8 | 9 |...| N | |
| 25 | /// +---+---+---+---+---+---+---+---+---+---+---+---+ |
| 26 | /// | | |
| 27 | /// curr | |
| 28 | /// next |
| 29 | /// +---+---+---+---+---+---+---+---+---+---+---+---+ |
| 30 | /// | 0 | 1 | 2 | 3 | 4 | 5 | 6 | 7 | 8 | 9 |...| N | |
| 31 | /// +---+---+---+---+---+---+---+---+---+---+---+---+ |
| 32 | /// | | |
| 33 | /// curr | |
| 34 | /// next |
| 35 | /// +---+---+---+---+---+---+---+---+---+---+---+---+ |
| 36 | /// | 0 | 1 | 2 | 3 | 4 | 5 | 6 | 7 | 8 | 9 |...| N | |
| 37 | /// +---+---+---+---+---+---+---+---+---+---+---+---+ |
| 38 | /// | | |
| 39 | /// | curr |
| 40 | /// next |
| 41 | /// |
| 42 | /// |
| 43 | #[derive(Clone, Copy, Debug, Eq, PartialEq)] |
| 44 | pub struct RingBuffer< |
| 45 | const N: usize, |
| 46 | D: Clone + fmt::Debug, |
| 47 | > { |
| 48 | pub curr: usize, |
| 49 | pub next: usize, |
| 50 | pub buf: [Option<D>; N], |
| 51 | } |
| 52 | |
| 53 | /// Fill buffer with the current time. |
| 54 | impl< |
| 55 | const N: usize, |
| 56 | D: Clone + fmt::Debug, |
| 57 | > |
| 58 | Default for RingBuffer<N, D> |
| 59 | { |
| 60 | fn default() -> Self { |
| 61 | let d0 = None::<D>; |
| 62 | let buf = std::array::from_fn(|_| d0.clone()); |
| 63 | let next = if N == 1 { 0 } else { 1 }; |
| 64 | Self { |
| 65 | curr: 0, |
| 66 | next, |
| 67 | buf, |
| 68 | } |
| 69 | } |
| 70 | } |
| 71 | |
| 72 | impl< |
| 73 | const N: usize, |
| 74 | D: Clone + fmt::Debug + ToDat, |
| 75 | > |
| 76 | ToDat for RingBuffer<N, D> |
| 77 | { |
| 78 | /// Converts `RingBuffer` to type: |
| 79 | /// |
| 80 | ///```ignore |
| 81 | /// Dat::Tup2(Box<[ |
| 82 | /// Dat::Vek(Vek<Vec<Dat::Opt(Box<Option<[ -+ |
| 83 | /// res!(D::to_dat()), +-- the buffer as a Vec<Option<D>> |
| 84 | /// ]>>)>>), -+ |
| 85 | /// Dat::U64(u64), -- the curr pointer (next can be reconstructed) |
| 86 | /// ]>) |
| 87 | ///``` |
| 88 | /// |
| 89 | fn to_dat(&self) -> Outcome<Dat> { |
| 90 | let mut v = Vec::with_capacity(N); |
| 91 | for opt in &self.buf { |
| 92 | let dat = match opt { |
| 93 | Some(d) => Dat::Opt(Box::new(Some(res!(d.to_dat())))), |
| 94 | None => Dat::Opt(Box::new(None)), |
| 95 | }; |
| 96 | v.push(dat) |
| 97 | } |
| 98 | Ok(tup2dat![ |
| 99 | Dat::Vek(Vek(v)), |
| 100 | Dat::U64(self.curr as u64), |
| 101 | ]) |
| 102 | } |
| 103 | } |
| 104 | |
| 105 | impl< |
| 106 | const N: usize, |
| 107 | D: Clone + fmt::Debug + FromDat, |
| 108 | > |
| 109 | FromDat for RingBuffer<N, D> |
| 110 | { |
| 111 | fn from_dat(dat: Dat) -> Outcome<Self> { |
| 112 | let mut v = try_extract_tup2dat!(dat); |
| 113 | let vek = try_extract_dat!(v[0].extract(), Vek); |
| 114 | let d0 = None::<D>; |
| 115 | let mut buf = std::array::from_fn(|_| d0.clone()); |
| 116 | for (i, dat1) in vek.into_iter().enumerate() { |
| 117 | match *try_extract_dat!(dat1, Opt) { |
| 118 | Some(dat2) => buf[i] = Some(res!(D::from_dat(dat2))), |
| 119 | None => buf[i] = d0.clone(), |
| 120 | } |
| 121 | } |
| 122 | let curr = try_extract_dat!(v[1].extract(), U64) as usize; |
| 123 | Ok(Self { |
| 124 | curr, |
| 125 | next: Self::next_index(curr), |
| 126 | buf, |
| 127 | }) |
| 128 | } |
| 129 | } |
| 130 | |
| 131 | impl< |
| 132 | const N: usize, |
| 133 | D: Clone + fmt::Debug, |
| 134 | > |
| 135 | RingBuffer<N, D> |
| 136 | { |
| 137 | /// Set data at current location. |
| 138 | pub fn set<I: Into<D>>(&mut self, new: I) { |
| 139 | let new = new.into(); |
| 140 | self.buf[self.curr] = Some(new); |
| 141 | } |
| 142 | |
| 143 | /// Set data at current location and advance pointer. |
| 144 | pub fn set_and_adv<I: Into<D>>(&mut self, new: I) { |
| 145 | self.set(new); |
| 146 | self.adv(); |
| 147 | } |
| 148 | |
| 149 | /// Advance pointer. |
| 150 | pub fn adv(&mut self) { |
| 151 | self.curr = self.next; |
| 152 | self.next = Self::next_index(self.curr); |
| 153 | } |
| 154 | |
| 155 | /// Return the current value of the pointer index. |
| 156 | pub fn curr(&self) -> usize { self.curr } |
| 157 | |
| 158 | /// Return the next value of the pointer index, given the current value. |
| 159 | pub fn next_index(curr: usize) -> usize { |
| 160 | if curr == N - 1 { |
| 161 | 0 |
| 162 | } else { |
| 163 | curr + 1 |
| 164 | } |
| 165 | } |
| 166 | |
| 167 | /// Return the previous value of the pointer index, given the current value. |
| 168 | pub fn prev_index(curr: usize) -> usize { |
| 169 | if curr == 0 { |
| 170 | N - 1 |
| 171 | } else { |
| 172 | curr - 1 |
| 173 | } |
| 174 | } |
| 175 | |
| 176 | /// Get data at current pointer location. |
| 177 | pub fn get(&self) -> Option<&D> { |
| 178 | self.buf[self.curr].as_ref() |
| 179 | } |
| 180 | |
| 181 | /// Get data at previous pointer location. |
| 182 | pub fn prev(&self) -> Option<&D> { |
| 183 | let index = if self.curr == 0 { |
| 184 | N - 1 |
| 185 | } else { |
| 186 | self.curr - 1 |
| 187 | }; |
| 188 | self.buf[index].as_ref() |
| 189 | } |
| 190 | |
| 191 | /// Get data at next pointer location. |
| 192 | pub fn next(&self) -> Option<&D> { |
| 193 | self.buf[self.next].as_ref() |
| 194 | } |
| 195 | |
| 196 | /// Copy the RingBuffer into a vector. |
| 197 | pub fn to_vec(&self) -> Vec<Option<D>> { |
| 198 | let mut result = Vec::with_capacity(N); |
| 199 | for opt in &self.buf { |
| 200 | result.push(opt.clone()) |
| 201 | } |
| 202 | result |
| 203 | } |
| 204 | |
| 205 | /// The occupied slots of a ring filled by [`Self::set_and_adv`], oldest first. |
| 206 | /// |
| 207 | /// Writing advances the pointer past the slot written, so the pointer always |
| 208 | /// rests on the oldest entry once the ring is full, and on the first empty slot |
| 209 | /// before then. Walking one lap from it therefore visits the entries in the |
| 210 | /// order they were written, whichever of the two states the ring is in. |
| 211 | /// [`Self::to_vec`] is storage order, which is chronological only until the |
| 212 | /// first wrap. |
| 213 | pub fn iter_chrono(&self) -> impl Iterator<Item = &D> + '_ { |
| 214 | (0..N).filter_map(move |i| self.buf[(self.curr + i) % N].as_ref()) |
| 215 | } |
| 216 | } |
| 217 | |
| 218 | /// A ring buffer consisting of timestamps. |
| 219 | /// |
| 220 | /// |
| 221 | /// +---+---+---+---+---+---+---+---+---+---+---+---+ |
| 222 | /// | 0 | 1 | 2 | 3 | 4 | 5 | 6 | 7 | 8 | 9 |...| N | |
| 223 | /// +---+---+---+---+---+---+---+---+---+---+---+---+ |
| 224 | /// i | |
| 225 | /// next |
| 226 | /// +---+---+---+---+---+---+---+---+---+---+---+---+ |
| 227 | /// | 0 | 1 | 2 | 3 | 4 | 5 | 6 | 7 | 8 | 9 |...| N | |
| 228 | /// +---+---+---+---+---+---+---+---+---+---+---+---+ |
| 229 | /// i | |
| 230 | /// next |
| 231 | /// +---+---+---+---+---+---+---+---+---+---+---+---+ |
| 232 | /// | 0 | 1 | 2 | 3 | 4 | 5 | 6 | 7 | 8 | 9 |...| N | |
| 233 | /// +---+---+---+---+---+---+---+---+---+---+---+---+ |
| 234 | /// | i |
| 235 | /// next |
| 236 | /// |
| 237 | /// |
| 238 | |
| 239 | #[derive(Clone, Debug, Default)] |
| 240 | pub struct RingTimer<const N: usize>(RingBuffer<N, SystemTime>); |
| 241 | |
| 242 | impl<const N: usize> std::ops::Deref for RingTimer<N> { |
| 243 | type Target = RingBuffer<N, SystemTime>; |
| 244 | fn deref(&self) -> &Self::Target { |
| 245 | &self.0 |
| 246 | } |
| 247 | } |
| 248 | |
| 249 | impl<const N: usize> RingTimer<N> { |
| 250 | |
| 251 | /// Add a timestamp to the ring buffer, returning the [Duration] since the previous entry. |
| 252 | pub fn update(&mut self) -> Duration { |
| 253 | // Note `self.0`, not a copy of it. `RingBuffer` is `Copy`, so binding |
| 254 | // it to a local writes the timestamp into a temporary that is then |
| 255 | // dropped, and the ring stays empty for ever. Every rate this type |
| 256 | // reports would be zero, which is not a rate limiter, it is an |
| 257 | // ornament. |
| 258 | self.0.set_and_adv(SystemTime::now()); |
| 259 | self.last_duration() |
| 260 | } |
| 261 | |
| 262 | fn extent(&self) -> (Option<SystemTime>, Option<SystemTime>, usize) { |
| 263 | let mut oldest: Option<SystemTime> = None; |
| 264 | let mut newest: Option<SystemTime> = None; |
| 265 | let mut count = 0; |
| 266 | for slot in self.0.buf.iter() { |
| 267 | let t = match slot { |
| 268 | Some(t) => *t, |
| 269 | None => continue, |
| 270 | }; |
| 271 | count += 1; |
| 272 | match oldest { |
| 273 | Some(o) if o <= t => (), |
| 274 | _ => oldest = Some(t), |
| 275 | } |
| 276 | match newest { |
| 277 | Some(n) if n >= t => (), |
| 278 | _ => newest = Some(t), |
| 279 | } |
| 280 | } |
| 281 | (oldest, newest, count) |
| 282 | } |
| 283 | |
| 284 | pub fn count(&self) -> usize { |
| 285 | let (_, _, count) = self.extent(); |
| 286 | count |
| 287 | } |
| 288 | |
| 289 | pub fn is_full(&self) -> bool { |
| 290 | self.count() == N |
| 291 | } |
| 292 | |
| 293 | pub fn last_duration(&self) -> Duration { |
| 294 | let newest_i = RingBuffer::<N, SystemTime>::prev_index(self.0.curr); |
| 295 | let prev_i = RingBuffer::<N, SystemTime>::prev_index(newest_i); |
| 296 | if newest_i == prev_i { |
| 297 | return Duration::ZERO; |
| 298 | } |
| 299 | match (self.0.buf[newest_i], self.0.buf[prev_i]) { |
| 300 | (Some(newest), Some(prev)) => match newest.duration_since(prev) { |
| 301 | Ok(duration) => duration, |
| 302 | Err(_) => Duration::ZERO, |
| 303 | }, |
| 304 | _ => Duration::ZERO, |
| 305 | } |
| 306 | } |
| 307 | |
| 308 | pub fn total_duration(&self) -> Duration { |
| 309 | let (oldest, newest, count) = self.extent(); |
| 310 | if count < 2 { |
| 311 | return Duration::ZERO; |
| 312 | } |
| 313 | match (oldest, newest) { |
| 314 | (Some(o), Some(n)) => match n.duration_since(o) { |
| 315 | Ok(duration) => duration, |
| 316 | Err(_) => Duration::ZERO, |
| 317 | }, |
| 318 | _ => Duration::ZERO, |
| 319 | } |
| 320 | } |
| 321 | |
| 322 | pub fn avg_rps(&self) -> u64 { |
| 323 | if !self.is_full() { |
| 324 | return 0; |
| 325 | } |
| 326 | let ms = self.total_duration().as_millis(); |
| 327 | if ms == 0 { |
| 328 | // A full ring inside a single millisecond is not a rate of zero, |
| 329 | // it is the fastest rate this clock can resolve. Report the |
| 330 | // maximum so that any finite limit is exceeded. |
| 331 | return u64::MAX; |
| 332 | } |
| 333 | // N timestamps span N-1 intervals. |
| 334 | let intervals = (N.saturating_sub(1)) as u128; |
| 335 | let rps = intervals.saturating_mul(1_000) / ms; |
| 336 | if rps > u64::MAX as u128 { |
| 337 | u64::MAX |
| 338 | } else { |
| 339 | rps as u64 |
| 340 | } |
| 341 | } |
| 342 | } |
| 343 | |
| 344 | |
| 345 | // ┌───────────────────────────────────────────────────────────────────────────┐ |
| 346 | // │ TESTS │ |
| 347 | // └───────────────────────────────────────────────────────────────────────────┘ |
| 348 | |
| 349 | #[cfg(test)] |
| 350 | mod tests { |
| 351 | use super::*; |
| 352 | |
| 353 | #[test] |
| 354 | fn test_update_records_a_timestamp_00() { |
| 355 | let mut timer = RingTimer::<4>::default(); |
| 356 | assert_eq!(timer.count(), 0); |
| 357 | timer.update(); |
| 358 | assert_eq!(timer.count(), 1, "update must write to the ring, not a copy"); |
| 359 | timer.update(); |
| 360 | assert_eq!(timer.count(), 2); |
| 361 | } |
| 362 | |
| 363 | #[test] |
| 364 | fn test_ring_fills_and_wraps_00() { |
| 365 | let mut timer = RingTimer::<4>::default(); |
| 366 | for _ in 0..4 { |
| 367 | timer.update(); |
| 368 | } |
| 369 | assert!(timer.is_full()); |
| 370 | assert_eq!(timer.count(), 4); |
| 371 | // Wrapping overwrites rather than growing. |
| 372 | timer.update(); |
| 373 | assert_eq!(timer.count(), 4); |
| 374 | } |
| 375 | |
| 376 | #[test] |
| 377 | fn test_avg_rps_is_zero_until_the_ring_fills_00() { |
| 378 | let mut timer = RingTimer::<8>::default(); |
| 379 | for _ in 0..7 { |
| 380 | timer.update(); |
| 381 | assert_eq!(timer.avg_rps(), 0, |
| 382 | "a partial ring must not report a rate"); |
| 383 | } |
| 384 | } |
| 385 | |
| 386 | #[test] |
| 387 | fn test_a_fast_burst_reports_a_high_rate_00() { |
| 388 | let mut timer = RingTimer::<8>::default(); |
| 389 | for _ in 0..8 { |
| 390 | timer.update(); |
| 391 | } |
| 392 | assert!(timer.is_full()); |
| 393 | // Eight timestamps taken as fast as the machine can loop span far less |
| 394 | // than a second. That is a very high rate, not a zero one. |
| 395 | let rps = timer.avg_rps(); |
| 396 | assert!(rps > 1_000, |
| 397 | "a sub-second burst must report a high rate, got {}", rps); |
| 398 | } |
| 399 | |
| 400 | #[test] |
| 401 | fn test_a_slow_caller_reports_a_low_rate_00() { |
| 402 | let mut timer = RingTimer::<4>::default(); |
| 403 | for _ in 0..4 { |
| 404 | timer.update(); |
| 405 | std::thread::sleep(Duration::from_millis(60)); |
| 406 | } |
| 407 | // Four timestamps spanning three 60 ms gaps is about 16 requests a |
| 408 | // second: comfortably measurable, and well under a 50/s limit. |
| 409 | let rps = timer.avg_rps(); |
| 410 | assert!(rps > 5 && rps < 40, |
| 411 | "expected roughly 16 rps from 60 ms spacing, got {}", rps); |
| 412 | } |
| 413 | |
| 414 | #[test] |
| 415 | fn test_iter_chrono_is_oldest_first_before_and_after_a_wrap_00() { |
| 416 | let mut ring = RingBuffer::<4, u32>::default(); |
| 417 | assert_eq!(ring.iter_chrono().count(), 0, "an empty ring yields nothing"); |
| 418 | for v in 1..=3u32 { |
| 419 | ring.set_and_adv(v); |
| 420 | } |
| 421 | let partial: Vec<u32> = ring.iter_chrono().copied().collect(); |
| 422 | assert_eq!(partial, vec![1, 2, 3], "a part-filled ring must read in write order"); |
| 423 | for v in 4..=6u32 { |
| 424 | ring.set_and_adv(v); |
| 425 | } |
| 426 | // Six writes into four slots: 1 and 2 are overwritten, and storage order |
| 427 | // is now [5, 6, 3, 4], which is not the order they were written in. |
| 428 | let wrapped: Vec<u32> = ring.iter_chrono().copied().collect(); |
| 429 | assert_eq!(wrapped, vec![3, 4, 5, 6], |
| 430 | "a wrapped ring must read oldest first, not in storage order"); |
| 431 | let stored: Vec<u32> = ring.to_vec().into_iter().flatten().collect(); |
| 432 | assert_ne!(stored, wrapped, "the test must exercise a ring whose storage order differs"); |
| 433 | } |
| 434 | |
| 435 | #[test] |
| 436 | fn test_last_duration_measures_the_most_recent_gap_00() { |
| 437 | let mut timer = RingTimer::<4>::default(); |
| 438 | // One timestamp: no interval to measure. |
| 439 | timer.update(); |
| 440 | assert_eq!(timer.last_duration(), Duration::ZERO); |
| 441 | std::thread::sleep(Duration::from_millis(40)); |
| 442 | timer.update(); |
| 443 | let gap = timer.last_duration(); |
| 444 | assert!(gap >= Duration::from_millis(35), |
| 445 | "expected a ~40 ms gap, got {:?}", gap); |
| 446 | assert!(gap < Duration::from_millis(500), |
| 447 | "expected a ~40 ms gap, got {:?}", gap); |
| 448 | } |
| 449 | } |