1use std::cmp::Reverse;
5
6use reifydb_codec::{
7 key::encoded::{EncodedKey, EncodedKeyRange},
8 row::pod::EncodedPodRow,
9};
10use reifydb_value::{
11 Result,
12 byte_size::ByteSize,
13 value::{datetime::DateTime, row_number::RowNumber},
14};
15
16use crate::key::operator::{
17 keyspace::{RootSibling, root_sibling_of},
18 state::{GroupId, GroupStateKey, group_data_inner_range, group_inner_range, keyspace_inner_range_split},
19};
20
21#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
22#[repr(u8)]
23pub enum TimerKind {
24 Seal = 0,
25 Grace = 1,
26 RowTtl = 2,
27 Maintenance = 3,
28}
29
30impl TimerKind {
31 pub fn is_maintenance(&self) -> bool {
32 matches!(self, Self::Maintenance)
33 }
34
35 pub fn from_u8(value: u8) -> Option<Self> {
36 match value {
37 0 => Some(Self::Seal),
38 1 => Some(Self::Grace),
39 2 => Some(Self::RowTtl),
40 3 => Some(Self::Maintenance),
41 _ => None,
42 }
43 }
44}
45
46pub struct GroupSweep {
47 pub rows: Vec<(GroupStateKey, EncodedPodRow)>,
48 pub complete: bool,
49}
50
51impl GroupSweep {
52 pub fn of(mut rows: Vec<(GroupStateKey, EncodedPodRow)>, limit: usize) -> Self {
53 let complete = rows.len() <= limit;
54 rows.truncate(limit);
55 Self {
56 rows,
57 complete,
58 }
59 }
60}
61
62pub fn sweep_order(groups: &[GroupId]) -> Vec<GroupId> {
63 let mut ordered = groups.to_vec();
64 ordered.sort_by_key(|group| Reverse(*group.as_bytes()));
65 ordered.dedup();
66 ordered
67}
68
69pub trait StateStore {
70 fn state_get(&mut self, key: &GroupStateKey) -> Result<Option<EncodedPodRow>>;
71
72 fn state_get_many_visit(
73 &mut self,
74 keys: &[GroupStateKey],
75 visit: &mut dyn FnMut(GroupStateKey, EncodedPodRow) -> Result<()>,
76 ) -> Result<()>;
77
78 fn state_classify(&mut self, _key: &GroupStateKey, _pre: Option<ByteSize>) {}
79
80 fn state_set(&mut self, key: &GroupStateKey, payload: EncodedPodRow) -> Result<()>;
81
82 fn state_remove(&mut self, key: &GroupStateKey) -> Result<()>;
83
84 fn state_page(
85 &mut self,
86 range: EncodedKeyRange,
87 limit: Option<usize>,
88 ) -> Result<Vec<(GroupStateKey, EncodedPodRow)>> {
89 debug_assert!(
90 keyspace_inner_range_split(&range).is_some(),
91 "a state page must stay inside one group and one keyspace; {range:?} spans more than one"
92 );
93 self.state_page_inner(range, limit)
94 }
95
96 fn state_page_inner(
97 &mut self,
98 range: EncodedKeyRange,
99 limit: Option<usize>,
100 ) -> Result<Vec<(GroupStateKey, EncodedPodRow)>>;
101
102 fn group_sweep(
103 &mut self,
104 group: GroupId,
105 data_only: bool,
106 limit: Option<usize>,
107 ) -> Result<Vec<(GroupStateKey, EncodedPodRow)>> {
108 let range = match data_only {
109 true => group_data_inner_range(group),
110 false => group_inner_range(group),
111 };
112 self.state_page_inner(range, limit)
113 }
114
115 fn group_sweep_many(&mut self, groups: &[GroupId], limit: usize) -> Result<GroupSweep> {
116 let mut rows = Vec::new();
117 for group in sweep_order(groups) {
118 if rows.len() > limit {
119 break;
120 }
121 let remaining = limit.saturating_add(1).saturating_sub(rows.len());
122 rows.extend(self.group_sweep(group, false, Some(remaining))?);
123 }
124 Ok(GroupSweep::of(rows, limit))
125 }
126
127 fn remove_root_siblings(&mut self, swept: &[(GroupStateKey, EncodedPodRow)]) -> Result<()> {
128 for (key, row) in swept {
129 let Some(sibling) = root_sibling_of(key, row) else {
130 continue;
131 };
132 if let RootSibling::Derived(sibling) = sibling {
133 self.state_remove(&sibling)?;
134 }
135 }
136 Ok(())
137 }
138
139 fn state_last(&mut self, range: EncodedKeyRange) -> Result<Option<(GroupStateKey, EncodedPodRow)>> {
140 Ok(self.state_page(range, None)?.pop())
141 }
142
143 fn get_or_create_row_numbers(&mut self, group: GroupId, keys: &[EncodedKey]) -> Result<Vec<(RowNumber, bool)>>;
144
145 fn get_or_create_row_numbers_for_groups(&mut self, groups: &[GroupId]) -> Result<Vec<(RowNumber, bool)>>;
146
147 fn remove_row_number(&mut self, group: GroupId, key: &EncodedKey) -> Result<()>;
148
149 fn remove_row_number_for_group(&mut self, group: GroupId) -> Result<()>;
150
151 fn remove_row_numbers(&mut self, group: GroupId, keys: &[EncodedKey]) -> Result<()> {
152 for key in keys {
153 self.remove_row_number(group, key)?;
154 }
155 Ok(())
156 }
157
158 fn written_at(&self) -> DateTime;
159}
160
161pub trait TimerStore {
162 fn arm_timer(&mut self, due: DateTime, kind: TimerKind, key: &EncodedKey) -> Result<()>;
163
164 fn disarm_timer(&mut self, due: DateTime, kind: TimerKind, key: &EncodedKey) -> Result<()>;
165
166 fn flow_watermark(&mut self) -> Result<Option<DateTime>>;
167}