Skip to main content

reifydb_core/state/
timer.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}