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