Skip to main content

reifydb_flow/window/
mint.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use reifydb_core::{key::operator_state::GroupId, state::store::StateStore};
5use reifydb_value::{Result, value::row_number::RowNumber};
6
7use crate::window::{
8	coord::{EventCoord, OrdinalCoord, TimeStamped},
9	meta::WindowMeta,
10};
11
12pub struct Mint<'a> {
13	meta: &'a mut WindowMeta,
14}
15
16impl<'a> Mint<'a> {
17	pub fn new(meta: &'a mut WindowMeta) -> Self {
18		Self {
19			meta,
20		}
21	}
22
23	pub fn event(row: &impl TimeStamped) -> EventCoord {
24		EventCoord::of(row)
25	}
26
27	pub fn ordinal(&mut self, store: &mut dyn StateStore, group: GroupId) -> Result<OrdinalCoord> {
28		Ok(OrdinalCoord::from_arrival_counter(self.meta.get_and_increment_count(store, group)?))
29	}
30
31	pub fn membership(
32		&mut self,
33		store: &mut dyn StateStore,
34		group: GroupId,
35		row_number: RowNumber,
36	) -> Result<Vec<u64>> {
37		self.meta.lookup_row_index(store, group, row_number)
38	}
39
40	pub fn record_membership(
41		&mut self,
42		store: &mut dyn StateStore,
43		group: GroupId,
44		row_number: RowNumber,
45		window_id: u64,
46	) -> Result<()> {
47		self.meta.store_row_index(store, group, row_number, window_id)
48	}
49
50	pub fn drop_membership(
51		&mut self,
52		store: &mut dyn StateStore,
53		group: GroupId,
54		row_number: RowNumber,
55	) -> Result<()> {
56		self.meta.drop_row_index(store, group, row_number)
57	}
58}
59
60#[cfg(test)]
61mod tests {
62	use reifydb_value::value::datetime::DateTime;
63
64	use super::*;
65	use crate::operator::state::mock::MockStore;
66
67	const GROUP: GroupId = GroupId(42);
68	const OTHER: GroupId = GroupId(43);
69
70	fn meta() -> WindowMeta {
71		WindowMeta::new()
72	}
73
74	#[test]
75	fn the_arrival_counter_starts_at_zero_and_never_repeats_within_a_group() {
76		// The ordinal IS the coordinate for a count window, so a repeat aliases two rows onto one
77		// slot and a skip leaves a hole the sweep never reaches. The counter is read-then-increment,
78		// so starting at 1 would shift every window boundary by one row.
79		let mut meta = meta();
80		let mut mint = Mint::new(&mut meta);
81		let mut store = MockStore::default();
82
83		let minted: Vec<u64> =
84			(0..4).map(|_| mint.ordinal(&mut store, GROUP).unwrap().value()).collect::<Vec<_>>();
85
86		assert_eq!(minted, vec![0, 1, 2, 3]);
87	}
88
89	#[test]
90	fn each_group_counts_independently() {
91		// A count window holds the last N rows per group. One shared counter would let a busy group
92		// shove a quiet group's rows across a boundary they never crossed.
93		let mut meta = meta();
94		let mut mint = Mint::new(&mut meta);
95		let mut store = MockStore::default();
96
97		mint.ordinal(&mut store, GROUP).unwrap();
98		mint.ordinal(&mut store, GROUP).unwrap();
99
100		assert_eq!(mint.ordinal(&mut store, OTHER).unwrap().value(), 0);
101		assert_eq!(mint.ordinal(&mut store, GROUP).unwrap().value(), 2);
102	}
103
104	#[test]
105	fn an_event_coordinate_is_minted_from_the_row_and_never_from_the_counter() {
106		// A time window's coordinate comes from the row's own instant, never the arrival counter one
107		// call away. The signatures enforce it: `event` cannot see the store, `ordinal` cannot see a row.
108		let row = DateTime::from_millis(5_000);
109
110		assert_eq!(Mint::event(&row).at(), DateTime::from_millis(5_000));
111	}
112
113	#[test]
114	fn a_row_records_every_window_it_joined_and_never_the_same_one_twice() {
115		// Sliding windows overlap, so one row joins several and retraction must find all of them. A
116		// duplicated id subtracts the row's contribution twice from one window.
117		let mut meta = meta();
118		let mut mint = Mint::new(&mut meta);
119		let mut store = MockStore::default();
120
121		mint.record_membership(&mut store, GROUP, RowNumber(7), 100).unwrap();
122		mint.record_membership(&mut store, GROUP, RowNumber(7), 200).unwrap();
123		mint.record_membership(&mut store, GROUP, RowNumber(7), 100).unwrap();
124
125		assert_eq!(mint.membership(&mut store, GROUP, RowNumber(7)).unwrap(), vec![100, 200]);
126	}
127
128	#[test]
129	fn a_row_that_joined_no_window_reports_an_empty_membership() {
130		// Retraction runs for every removed row, including rows the gate refused. Defaulting to some
131		// window would retract a contribution the row never made.
132		let mut meta = meta();
133		let mut mint = Mint::new(&mut meta);
134		let mut store = MockStore::default();
135
136		assert!(mint.membership(&mut store, GROUP, RowNumber(7)).unwrap().is_empty());
137	}
138
139	#[test]
140	fn membership_is_scoped_to_its_group() {
141		// Row numbers are unique per source, not per group, so two groups routinely see the same
142		// RowNumber. One shared list would retract a row from a group it never entered.
143		let mut meta = meta();
144		let mut mint = Mint::new(&mut meta);
145		let mut store = MockStore::default();
146
147		mint.record_membership(&mut store, GROUP, RowNumber(7), 100).unwrap();
148
149		assert!(mint.membership(&mut store, OTHER, RowNumber(7)).unwrap().is_empty());
150	}
151}