1use std::{
5 collections::{BTreeMap, BTreeSet, HashMap},
6 fmt::Debug,
7 hash::Hash,
8 marker::PhantomData,
9};
10
11use reifydb_value::{Result, reifydb_assertions, value::row_number::RowNumber};
12use serde::{Deserialize, Serialize, de::DeserializeOwned};
13
14use crate::{
15 encoded::key::{EncodedKey, IntoEncodedKey},
16 window::{
17 accumulator::WindowAccumulator,
18 engine::{
19 AccumulatorEvent, EmitKind, GroupMeta, LatePolicy, MetaKey, expiry_due_range, expiry_key,
20 meta_key_for,
21 },
22 span::Slot,
23 state::StateCache,
24 store::WindowStore,
25 },
26};
27
28pub type RollingBuffer<C, Accumulator> = BTreeMap<C, Accumulator>;
29
30pub type RollingBuckets<G, C, Contribution> = BTreeMap<(G, C), Vec<AccumulatorEvent<Contribution>>>;
31
32pub struct RollingResult<G, Output> {
33 pub row_number: RowNumber,
34 pub group: G,
35 pub value: Output,
36 pub prior: Option<Output>,
37 pub kind: EmitKind,
38}
39
40pub enum RollingEviction<C: Slot> {
41 Capacity(usize),
42 Before(C),
43 BeforeStamp(u64),
44}
45
46pub enum RollingExpiry<G, Output> {
47 Update {
48 row_number: RowNumber,
49 group: G,
50 value: Output,
51 },
52 Remove {
53 row_number: RowNumber,
54 group: G,
55 },
56}
57
58#[derive(Clone, Copy)]
59enum IndexMode {
60 Coord,
61 Stamp,
62}
63
64#[derive(Serialize, Deserialize)]
65#[serde(bound(serialize = "G: Serialize", deserialize = "G: DeserializeOwned"))]
66struct RollingIndexEntry<G> {
67 group: G,
68 row_number: u64,
69}
70
71fn coord_min_key<C: Slot, A>(buffer: &RollingBuffer<C, A>) -> Option<u64> {
72 buffer.keys().next().map(|c| c.order_key())
73}
74
75fn stamp_min_key<C, A: WindowAccumulator>(buffer: &RollingBuffer<C, A>) -> Option<u64> {
76 buffer.values().filter_map(|a| a.stamp()).min()
77}
78
79type MetaLoaded<G, C> = HashMap<G, GroupMeta<C>>;
80type BufferRows<G> = HashMap<G, (RowNumber, bool)>;
81
82struct GroupSlot<C, Accumulator> {
83 row_number: RowNumber,
84 is_new: bool,
85 buffer: RollingBuffer<C, Accumulator>,
86 was_empty_before: bool,
87 buffer_changed: bool,
88 prior_index_key: Option<u64>,
89}
90
91pub struct RollingEngine<G, C, Accumulator> {
92 buffers: StateCache<RowNumber, RollingBuffer<C, Accumulator>>,
93 meta: StateCache<MetaKey, GroupMeta<C>>,
94 late_policy: LatePolicy,
95 _pd: PhantomData<G>,
96}
97
98impl<G, C, Accumulator> Default for RollingEngine<G, C, Accumulator>
99where
100 G: Clone + Eq + Ord + Hash + Debug + Serialize + DeserializeOwned,
101 C: Slot + Hash + Serialize + DeserializeOwned,
102 Accumulator: WindowAccumulator,
103 for<'a> &'a G: IntoEncodedKey,
104{
105 fn default() -> Self {
106 Self::new()
107 }
108}
109
110impl<G, C, Accumulator> RollingEngine<G, C, Accumulator>
111where
112 G: Clone + Eq + Ord + Hash + Debug + Serialize + DeserializeOwned,
113 C: Slot + Hash + Serialize + DeserializeOwned,
114 Accumulator: WindowAccumulator,
115 for<'a> &'a G: IntoEncodedKey,
116{
117 pub fn new() -> Self {
118 Self::with_late_policy(LatePolicy::Drop)
119 }
120
121 pub fn with_late_policy(late_policy: LatePolicy) -> Self {
122 Self {
123 buffers: StateCache::<RowNumber, RollingBuffer<C, Accumulator>>::new(8),
124 meta: StateCache::<MetaKey, GroupMeta<C>>::new_internal(64),
125 late_policy,
126 _pd: PhantomData,
127 }
128 }
129
130 pub fn apply<S, K, CB, Output>(
131 &mut self,
132 store: &mut S,
133 buckets: RollingBuckets<G, C, Accumulator::Contribution>,
134 capacity: usize,
135 row_key: K,
136 combine: CB,
137 ) -> Result<Vec<RollingResult<G, Output>>>
138 where
139 S: WindowStore,
140 K: Fn(&G) -> EncodedKey,
141 CB: Fn(&G, &RollingBuffer<C, Accumulator>) -> Option<Output>,
142 {
143 self.apply_evicting(
144 store,
145 buckets,
146 RollingEviction::Capacity(capacity),
147 row_key,
148 Accumulator::default,
149 combine,
150 )
151 }
152
153 pub fn apply_evicting<S, K, NA, CB, Output>(
154 &mut self,
155 store: &mut S,
156 buckets: RollingBuckets<G, C, Accumulator::Contribution>,
157 eviction: RollingEviction<C>,
158 row_key: K,
159 new_accumulator: NA,
160 combine: CB,
161 ) -> Result<Vec<RollingResult<G, Output>>>
162 where
163 S: WindowStore,
164 K: Fn(&G) -> EncodedKey,
165 NA: Fn() -> Accumulator,
166 CB: Fn(&G, &RollingBuffer<C, Accumulator>) -> Option<Output>,
167 {
168 if buckets.is_empty() {
169 return Ok(Vec::new());
170 }
171 let index_mode = match eviction {
172 RollingEviction::Capacity(_) => None,
173 RollingEviction::Before(_) => Some(IndexMode::Coord),
174 RollingEviction::BeforeStamp(_) => Some(IndexMode::Stamp),
175 };
176 let mut meta_loaded = self.warm_and_load_meta(store, &buckets)?;
177 let buffer_rows = self.resolve_buffer_rows(store, &buckets, &meta_loaded, &row_key)?;
178 let group_slots = self.apply_events_into_buffers(
179 store,
180 buckets,
181 &mut meta_loaded,
182 &buffer_rows,
183 &row_key,
184 &eviction,
185 &new_accumulator,
186 index_mode,
187 )?;
188 let results = self.combine_and_collect(store, group_slots, &combine, index_mode)?;
189 self.persist_meta(store, meta_loaded)?;
190 Ok(results)
191 }
192
193 pub fn flush<S: WindowStore>(&mut self, store: &mut S) -> Result<()> {
194 self.buffers.flush(store)?;
195 self.meta.flush(store)?;
196 Ok(())
197 }
198
199 fn warm_and_load_meta<S: WindowStore>(
200 &mut self,
201 store: &mut S,
202 buckets: &RollingBuckets<G, C, Accumulator::Contribution>,
203 ) -> Result<MetaLoaded<G, C>> {
204 let meta_keys: Vec<MetaKey> = buckets
205 .keys()
206 .map(|(group, _)| group)
207 .collect::<BTreeSet<_>>()
208 .into_iter()
209 .map(meta_key_for)
210 .collect();
211 self.meta.warm(store, &meta_keys)?;
212
213 let mut meta_loaded: MetaLoaded<G, C> = HashMap::new();
214 for (group, _) in buckets.keys() {
215 if !meta_loaded.contains_key(group) {
216 let m = self.meta.get(store, &meta_key_for(group))?.unwrap_or_default();
217 meta_loaded.insert(group.clone(), m);
218 }
219 }
220 Ok(meta_loaded)
221 }
222
223 fn resolve_buffer_rows<S, K>(
224 &mut self,
225 store: &mut S,
226 buckets: &RollingBuckets<G, C, Accumulator::Contribution>,
227 meta_loaded: &MetaLoaded<G, C>,
228 row_key: &K,
229 ) -> Result<BufferRows<G>>
230 where
231 S: WindowStore,
232 K: Fn(&G) -> EncodedKey,
233 {
234 let mut buffer_rows: BufferRows<G> = HashMap::new();
235 let mut resolve_order: Vec<G> = Vec::new();
236 let mut group_keys: Vec<EncodedKey> = Vec::new();
237 let mut seen: BTreeSet<G> = BTreeSet::new();
238 for (group, coord) in buckets.keys() {
239 let initial_high_water = meta_loaded.get(group).and_then(|m| m.high_water);
240 if initial_high_water.is_none_or(|hw| *coord >= hw) && seen.insert(group.clone()) {
241 resolve_order.push(group.clone());
242 group_keys.push(row_key(group));
243 }
244 }
245 let resolved_rows = store.get_or_create_row_numbers(&group_keys)?;
246 reifydb_assertions! {
247 let requested = group_keys.len();
248 let resolved = resolved_rows.len();
249 assert!(
250 requested == resolved,
251 "get_or_create_row_numbers returned a different count than requested, so the resolve_order \
252 zip would silently truncate buffer_rows and survivor groups would be re-resolved one at a \
253 time in apply_events_into_buffers, changing the per-batch row-number lookup cost \
254 (requested={requested}, resolved={resolved})"
255 );
256 }
257 let buffer_keys: Vec<RowNumber> = resolved_rows.iter().map(|(rn, _)| *rn).collect();
258 for (group, resolved) in resolve_order.into_iter().zip(resolved_rows) {
259 buffer_rows.insert(group, resolved);
260 }
261 self.buffers.warm(store, &buffer_keys)?;
262 Ok(buffer_rows)
263 }
264
265 #[allow(clippy::too_many_arguments)]
266 fn apply_events_into_buffers<S, K, NA>(
267 &mut self,
268 store: &mut S,
269 buckets: RollingBuckets<G, C, Accumulator::Contribution>,
270 meta_loaded: &mut MetaLoaded<G, C>,
271 buffer_rows: &BufferRows<G>,
272 row_key: &K,
273 eviction: &RollingEviction<C>,
274 new_accumulator: &NA,
275 index_mode: Option<IndexMode>,
276 ) -> Result<BTreeMap<G, GroupSlot<C, Accumulator>>>
277 where
278 S: WindowStore,
279 K: Fn(&G) -> EncodedKey,
280 NA: Fn() -> Accumulator,
281 {
282 let mut group_slots: BTreeMap<G, GroupSlot<C, Accumulator>> = BTreeMap::new();
283
284 for ((group, coord), events) in buckets {
285 let meta = meta_loaded.entry(group.clone()).or_default();
286
287 let slot = match group_slots.get_mut(&group) {
288 Some(s) => s,
289 None => {
290 let (row_number, is_new) = match buffer_rows.get(&group) {
291 Some(&resolved) => resolved,
292 None => {
293 let key = row_key(&group);
294 store.get_or_create_row_number(&key)?
295 }
296 };
297 let buffer: RollingBuffer<C, Accumulator> =
298 self.buffers.get(store, &row_number)?.unwrap_or_default();
299 let was_empty_before = buffer.is_empty();
300 let prior_index_key = match index_mode {
301 Some(IndexMode::Coord) => coord_min_key(&buffer),
302 Some(IndexMode::Stamp) => stamp_min_key(&buffer),
303 None => None,
304 };
305 group_slots.insert(
306 group.clone(),
307 GroupSlot {
308 row_number,
309 is_new,
310 buffer,
311 was_empty_before,
312 buffer_changed: false,
313 prior_index_key,
314 },
315 );
316 group_slots.get_mut(&group).expect("just inserted")
317 }
318 };
319
320 let late = matches!(meta.high_water, Some(hw) if coord < hw)
321 && matches!(self.late_policy, LatePolicy::Drop)
322 && !slot.buffer.contains_key(&coord);
323
324 let mut accumulator = slot.buffer.remove(&coord).unwrap_or_else(new_accumulator);
325 let mut touched = false;
326 for event in events {
327 match event {
328 AccumulatorEvent::Add(c) => {
329 if late {
330 continue;
331 }
332 accumulator.add(&c);
333 touched = true;
334 }
335 AccumulatorEvent::Remove(c) => {
336 if accumulator.is_empty() {
337 continue;
338 }
339 accumulator.remove(&c);
340 touched = true;
341 }
342 }
343 }
344 if !accumulator.is_empty() {
345 slot.buffer.insert(coord, accumulator);
346 }
347 if !touched {
348 continue;
349 }
350 match eviction {
351 RollingEviction::Capacity(cap) => {
352 while slot.buffer.len() > *cap {
353 slot.buffer.pop_first();
354 }
355 }
356 RollingEviction::Before(cutoff) => {
357 while let Some((&oldest, _)) = slot.buffer.iter().next() {
358 if oldest <= *cutoff {
359 slot.buffer.pop_first();
360 } else {
361 break;
362 }
363 }
364 }
365 RollingEviction::BeforeStamp(cutoff) => {
366 let stale: Vec<C> = slot
367 .buffer
368 .iter()
369 .filter(|(_, accumulator)| {
370 accumulator.stamp().is_some_and(|s| s <= *cutoff)
371 })
372 .map(|(coord, _)| *coord)
373 .collect();
374 for coord in stale {
375 slot.buffer.remove(&coord);
376 }
377 }
378 }
379 slot.buffer_changed = true;
380
381 meta.high_water = Some(match meta.high_water {
382 Some(hw) if hw > coord => hw,
383 _ => coord,
384 });
385 }
386 Ok(group_slots)
387 }
388
389 fn combine_and_collect<S, CB, Output>(
390 &mut self,
391 store: &mut S,
392 group_slots: BTreeMap<G, GroupSlot<C, Accumulator>>,
393 combine: &CB,
394 index_mode: Option<IndexMode>,
395 ) -> Result<Vec<RollingResult<G, Output>>>
396 where
397 S: WindowStore,
398 CB: Fn(&G, &RollingBuffer<C, Accumulator>) -> Option<Output>,
399 {
400 let mut results: Vec<RollingResult<G, Output>> = Vec::new();
401 for (group, slot) in group_slots {
402 if !slot.buffer_changed {
403 continue;
404 }
405 if let Some(mode) = index_mode {
406 let new_index_key = match mode {
407 IndexMode::Coord => coord_min_key(&slot.buffer),
408 IndexMode::Stamp => stamp_min_key(&slot.buffer),
409 };
410 if new_index_key != slot.prior_index_key {
411 if let Some(old) = slot.prior_index_key {
412 store.internal_drop(&expiry_key(old, &group, &[]))?;
413 }
414 if let Some(new) = new_index_key {
415 store.internal_set(
416 &expiry_key(new, &group, &[]),
417 &RollingIndexEntry {
418 group: group.clone(),
419 row_number: slot.row_number.0,
420 },
421 )?;
422 }
423 }
424 }
425 let output = combine(&group, &slot.buffer);
426 self.buffers.put(store, &slot.row_number, slot.buffer)?;
427
428 if let Some(out) = output {
429 let kind = if slot.is_new || slot.was_empty_before {
430 EmitKind::Insert
431 } else {
432 EmitKind::Update
433 };
434 results.push(RollingResult {
435 row_number: slot.row_number,
436 group,
437 value: out,
438 prior: None,
439 kind,
440 });
441 }
442 }
443 Ok(results)
444 }
445
446 pub fn expire_before<S, CB, Output>(
447 &mut self,
448 store: &mut S,
449 cutoff: C,
450 combine: CB,
451 ) -> Result<Vec<RollingExpiry<G, Output>>>
452 where
453 S: WindowStore,
454 CB: Fn(&G, &RollingBuffer<C, Accumulator>) -> Option<Output>,
455 {
456 let mut due: Vec<(EncodedKey, RollingIndexEntry<G>)> = Vec::new();
457 store.internal_range_visit::<RollingIndexEntry<G>>(
458 expiry_due_range(cutoff.order_key()),
459 &mut |key, entry| {
460 due.push((key, entry));
461 Ok(())
462 },
463 )?;
464
465 let mut out: Vec<RollingExpiry<G, Output>> = Vec::new();
466 for (index_key, entry) in due {
467 let row_number = RowNumber(entry.row_number);
468 store.internal_drop(&index_key)?;
469 let Some(mut buffer) = self.buffers.get(store, &row_number)? else {
470 continue;
471 };
472 let before = buffer.len();
473 buffer.retain(|&coord, _| coord > cutoff);
474 if buffer.len() == before {
475 if let Some(new) = coord_min_key(&buffer) {
476 store.internal_set(
477 &expiry_key(new, &entry.group, &[]),
478 &RollingIndexEntry {
479 group: entry.group.clone(),
480 row_number: entry.row_number,
481 },
482 )?;
483 }
484 continue;
485 }
486 match combine(&entry.group, &buffer) {
487 Some(value) if !buffer.is_empty() => {
488 if let Some(new) = coord_min_key(&buffer) {
489 store.internal_set(
490 &expiry_key(new, &entry.group, &[]),
491 &RollingIndexEntry {
492 group: entry.group.clone(),
493 row_number: entry.row_number,
494 },
495 )?;
496 }
497 self.buffers.put(store, &row_number, buffer)?;
498 out.push(RollingExpiry::Update {
499 row_number,
500 group: entry.group,
501 value,
502 });
503 }
504 _ => {
505 self.buffers.remove(store, &row_number)?;
506 out.push(RollingExpiry::Remove {
507 row_number,
508 group: entry.group,
509 });
510 }
511 }
512 }
513 Ok(out)
514 }
515
516 pub fn expire_before_stamp<S, CB, Output>(
517 &mut self,
518 store: &mut S,
519 cutoff: u64,
520 combine: CB,
521 ) -> Result<Vec<RollingExpiry<G, Output>>>
522 where
523 S: WindowStore,
524 CB: Fn(&G, &RollingBuffer<C, Accumulator>) -> Option<Output>,
525 {
526 let mut due: Vec<(EncodedKey, RollingIndexEntry<G>)> = Vec::new();
527 store.internal_range_visit::<RollingIndexEntry<G>>(expiry_due_range(cutoff), &mut |key, entry| {
528 due.push((key, entry));
529 Ok(())
530 })?;
531
532 let mut out: Vec<RollingExpiry<G, Output>> = Vec::new();
533 for (index_key, entry) in due {
534 let row_number = RowNumber(entry.row_number);
535 store.internal_drop(&index_key)?;
536 let Some(mut buffer) = self.buffers.get(store, &row_number)? else {
537 continue;
538 };
539 let before = buffer.len();
540 buffer.retain(|_, accumulator| accumulator.stamp().is_none_or(|s| s > cutoff));
541 if buffer.len() == before {
542 if let Some(new) = stamp_min_key(&buffer) {
543 store.internal_set(
544 &expiry_key(new, &entry.group, &[]),
545 &RollingIndexEntry {
546 group: entry.group.clone(),
547 row_number: entry.row_number,
548 },
549 )?;
550 }
551 continue;
552 }
553 match combine(&entry.group, &buffer) {
554 Some(value) if !buffer.is_empty() => {
555 if let Some(new) = stamp_min_key(&buffer) {
556 store.internal_set(
557 &expiry_key(new, &entry.group, &[]),
558 &RollingIndexEntry {
559 group: entry.group.clone(),
560 row_number: entry.row_number,
561 },
562 )?;
563 }
564 self.buffers.put(store, &row_number, buffer)?;
565 out.push(RollingExpiry::Update {
566 row_number,
567 group: entry.group,
568 value,
569 });
570 }
571 _ => {
572 self.buffers.remove(store, &row_number)?;
573 out.push(RollingExpiry::Remove {
574 row_number,
575 group: entry.group,
576 });
577 }
578 }
579 }
580 Ok(out)
581 }
582
583 fn persist_meta<S: WindowStore>(&mut self, store: &mut S, meta_loaded: MetaLoaded<G, C>) -> Result<()> {
584 for (group, meta) in meta_loaded {
585 self.meta.set(store, &meta_key_for(&group), &meta)?;
586 }
587 Ok(())
588 }
589}
590
591#[cfg(test)]
592mod tests {
593 use std::collections::BTreeMap;
594
595 use crate::{
596 encoded::key::EncodedKey,
597 window::engine::{
598 AccumulatorEvent,
599 rolling::{RollingBuckets, RollingBuffer, RollingEngine, RollingEviction, RollingExpiry},
600 test_support::{MockStore, StampedSum, SumAccumulator},
601 },
602 };
603
604 fn row_key(group: &u32) -> EncodedKey {
605 EncodedKey::builder().u32(*group).build()
606 }
607
608 fn sum_combine(_group: &u32, buffer: &RollingBuffer<u64, SumAccumulator>) -> Option<i64> {
609 if buffer.is_empty() {
610 None
611 } else {
612 Some(buffer.values().map(|a| a.sum).sum())
613 }
614 }
615
616 fn stamped_combine(_group: &u32, buffer: &RollingBuffer<u64, StampedSum>) -> Option<i64> {
617 if buffer.is_empty() {
618 None
619 } else {
620 Some(buffer.values().map(|a| a.sum).sum())
621 }
622 }
623
624 #[test]
625 fn expire_before_evicts_a_quiet_group_then_rekeys_then_removes() {
626 let mut store = MockStore::default();
627 let mut engine = RollingEngine::<u32, u64, SumAccumulator>::new();
628 let mut buckets: RollingBuckets<u32, u64, i64> = BTreeMap::new();
629 buckets.insert((1u32, 10u64), vec![AccumulatorEvent::Add(1)]);
630 buckets.insert((1u32, 20u64), vec![AccumulatorEvent::Add(2)]);
631 buckets.insert((1u32, 30u64), vec![AccumulatorEvent::Add(3)]);
632 engine.apply_evicting(
634 &mut store,
635 buckets,
636 RollingEviction::Before(0),
637 row_key,
638 SumAccumulator::default,
639 sum_combine,
640 )
641 .unwrap();
642 engine.flush(&mut store).unwrap();
643 assert_eq!(store.index_entry_count(), 1, "the group is indexed by its oldest coord");
644
645 let mut engine = RollingEngine::<u32, u64, SumAccumulator>::new();
647 let out = engine.expire_before(&mut store, 20, sum_combine).unwrap();
648 engine.flush(&mut store).unwrap();
649 assert_eq!(out.len(), 1);
650 match &out[0] {
651 RollingExpiry::Update {
652 group,
653 value,
654 ..
655 } => {
656 assert_eq!(*group, 1);
657 assert_eq!(*value, 3, "only the surviving coord 30 contributes");
658 }
659 RollingExpiry::Remove {
660 ..
661 } => panic!("group still has a live coord"),
662 }
663 assert_eq!(store.index_entry_count(), 1, "still one entry, re-keyed to coord 30");
664
665 let mut engine = RollingEngine::<u32, u64, SumAccumulator>::new();
667 let out = engine.expire_before(&mut store, 30, sum_combine).unwrap();
668 engine.flush(&mut store).unwrap();
669 assert_eq!(out.len(), 1);
670 match &out[0] {
671 RollingExpiry::Remove {
672 group,
673 ..
674 } => assert_eq!(*group, 1),
675 RollingExpiry::Update {
676 ..
677 } => panic!("the group is empty and must be removed"),
678 }
679 assert_eq!(store.index_entry_count(), 0, "the emptied group leaves no index entry");
680
681 let mut engine = RollingEngine::<u32, u64, SumAccumulator>::new();
683 assert!(engine.expire_before(&mut store, 1000, sum_combine).unwrap().is_empty());
684 }
685
686 #[test]
687 fn expire_before_leaves_groups_whose_oldest_coord_is_not_due() {
688 let mut store = MockStore::default();
689 let mut engine = RollingEngine::<u32, u64, SumAccumulator>::new();
690 let mut buckets: RollingBuckets<u32, u64, i64> = BTreeMap::new();
691 buckets.insert((1u32, 100u64), vec![AccumulatorEvent::Add(1)]);
692 buckets.insert((2u32, 5u64), vec![AccumulatorEvent::Add(9)]);
693 engine.apply_evicting(
694 &mut store,
695 buckets,
696 RollingEviction::Before(0),
697 row_key,
698 SumAccumulator::default,
699 sum_combine,
700 )
701 .unwrap();
702 engine.flush(&mut store).unwrap();
703 assert_eq!(store.index_entry_count(), 2);
704
705 let mut engine = RollingEngine::<u32, u64, SumAccumulator>::new();
707 let out = engine.expire_before(&mut store, 5, sum_combine).unwrap();
708 engine.flush(&mut store).unwrap();
709 assert_eq!(out.len(), 1, "only the group with a due coord is processed");
710 assert!(matches!(&out[0], RollingExpiry::Remove { group, .. } if *group == 2));
711 assert_eq!(store.index_entry_count(), 1, "group 1 keeps its index entry");
712 }
713
714 #[test]
715 fn expire_before_stamp_evicts_by_accumulator_stamp() {
716 let mut store = MockStore::default();
717 let mut engine = RollingEngine::<u32, u64, StampedSum>::new();
718 let mut buckets: RollingBuckets<u32, u64, (i64, u64)> = BTreeMap::new();
719 buckets.insert((1u32, 1u64), vec![AccumulatorEvent::Add((1, 10))]);
720 buckets.insert((1u32, 2u64), vec![AccumulatorEvent::Add((2, 20))]);
721 buckets.insert((1u32, 3u64), vec![AccumulatorEvent::Add((3, 30))]);
722 engine.apply_evicting(
723 &mut store,
724 buckets,
725 RollingEviction::BeforeStamp(0),
726 row_key,
727 StampedSum::default,
728 stamped_combine,
729 )
730 .unwrap();
731 engine.flush(&mut store).unwrap();
732 assert_eq!(store.index_entry_count(), 1, "indexed by the minimum stamp");
733
734 let mut engine = RollingEngine::<u32, u64, StampedSum>::new();
736 let out = engine.expire_before_stamp(&mut store, 20, stamped_combine).unwrap();
737 engine.flush(&mut store).unwrap();
738 assert_eq!(out.len(), 1);
739 match &out[0] {
740 RollingExpiry::Update {
741 value,
742 ..
743 } => assert_eq!(*value, 3),
744 RollingExpiry::Remove {
745 ..
746 } => panic!("a live entry remains"),
747 }
748 assert_eq!(store.index_entry_count(), 1, "re-keyed to the surviving stamp");
749 }
750}