Oregami
Repositories/oxedyne/fe2o3

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

1use oxedyne_fe2o3_core::{
2 prelude::*,
3 mem::Extract,
4};
5
6use oxedyne_fe2o3_jdat::{
7 prelude::*,
8 try_extract_tup2dat,
9 tup2dat,
10};
11
12use 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)]
44pub 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.
54impl<
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
72impl<
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
105impl<
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
131impl<
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)]
240pub struct RingTimer<const N: usize>(RingBuffer<N, SystemTime>);
241
242impl<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
249impl<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)]
350mod 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}