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::timer::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::{util::hash::Hash128, value::datetime::DateTime};
63
64	use super::*;
65	use crate::operator::state::mock::MockStore;
66
67	fn group() -> GroupId {
68		GroupId::hashed(Hash128(42))
69	}
70
71	fn other() -> GroupId {
72		GroupId::hashed(Hash128(43))
73	}
74
75	fn meta() -> WindowMeta {
76		WindowMeta::new()
77	}
78
79	#[test]
80	fn the_arrival_counter_starts_at_zero_and_never_repeats_within_a_group() {
81		// The ordinal IS the coordinate for a count window, so a repeat aliases two rows onto one
82		// slot and a skip leaves a hole the sweep never reaches. The counter is read-then-increment,
83		// so starting at 1 would shift every window boundary by one row.
84		let mut meta = meta();
85		let mut mint = Mint::new(&mut meta);
86		let mut store = MockStore::default();
87
88		let minted: Vec<u64> =
89			(0..4).map(|_| mint.ordinal(&mut store, group()).unwrap().value()).collect::<Vec<_>>();
90
91		assert_eq!(minted, vec![0, 1, 2, 3]);
92	}
93
94	#[test]
95	fn each_group_counts_independently() {
96		// A count window holds the last N rows per group. One shared counter would let a busy group
97		// shove a quiet group's rows across a boundary they never crossed.
98		let mut meta = meta();
99		let mut mint = Mint::new(&mut meta);
100		let mut store = MockStore::default();
101
102		mint.ordinal(&mut store, group()).unwrap();
103		mint.ordinal(&mut store, group()).unwrap();
104
105		assert_eq!(mint.ordinal(&mut store, other()).unwrap().value(), 0);
106		assert_eq!(mint.ordinal(&mut store, group()).unwrap().value(), 2);
107	}
108
109	#[test]
110	fn an_event_coordinate_is_minted_from_the_row_and_never_from_the_counter() {
111		// A time window's coordinate comes from the row's own instant, never the arrival counter one
112		// call away. The signatures enforce it: `event` cannot see the store, `ordinal` cannot see a row.
113		let row = DateTime::from_millis(5_000);
114
115		assert_eq!(Mint::event(&row).at(), DateTime::from_millis(5_000));
116	}
117
118	#[test]
119	fn a_row_records_every_window_it_joined_and_never_the_same_one_twice() {
120		// Sliding windows overlap, so one row joins several and retraction must find all of them. A
121		// duplicated id subtracts the row's contribution twice from one window.
122		let mut meta = meta();
123		let mut mint = Mint::new(&mut meta);
124		let mut store = MockStore::default();
125
126		mint.record_membership(&mut store, group(), RowNumber(7), 100).unwrap();
127		mint.record_membership(&mut store, group(), RowNumber(7), 200).unwrap();
128		mint.record_membership(&mut store, group(), RowNumber(7), 100).unwrap();
129
130		assert_eq!(mint.membership(&mut store, group(), RowNumber(7)).unwrap(), vec![100, 200]);
131	}
132
133	#[test]
134	fn a_row_that_joined_no_window_reports_an_empty_membership() {
135		// Retraction runs for every removed row, including rows the gate refused. Defaulting to some
136		// window would retract a contribution the row never made.
137		let mut meta = meta();
138		let mut mint = Mint::new(&mut meta);
139		let mut store = MockStore::default();
140
141		assert!(mint.membership(&mut store, group(), RowNumber(7)).unwrap().is_empty());
142	}
143
144	#[test]
145	fn membership_is_scoped_to_its_group() {
146		// Row numbers are unique per source, not per group, so two groups routinely see the same
147		// RowNumber. One shared list would retract a row from a group it never entered.
148		let mut meta = meta();
149		let mut mint = Mint::new(&mut meta);
150		let mut store = MockStore::default();
151
152		mint.record_membership(&mut store, group(), RowNumber(7), 100).unwrap();
153
154		assert!(mint.membership(&mut store, other(), RowNumber(7)).unwrap().is_empty());
155	}
156}