Skip to main content

reifydb_core/window/engine/
rolling.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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		// Before(0) evicts nothing at apply (all coords > 0), so the buffer keeps 10,20,30.
633		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		// A tick with no new events for this group evicts coords <= 20; coord 30 survives.
646		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		// The next tick evicts the last coord: the group empties and is removed.
666		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		// A further tick finds nothing due.
682		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		// Cutoff 5 is due only for group 2 (oldest coord 5); group 1 (oldest 100) is untouched.
706		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		// Evict accumulators stamped <= 20; the stamp-30 entry survives.
735		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}