reifydb_flow/window/
mint.rs1use 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 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 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 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 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 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 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}