1use std::{
5 fmt,
6 fmt::{Display, Formatter},
7 num::ParseIntError,
8 str::FromStr,
9};
10
11use reifydb_value::value::duration::Duration;
12use serde::{Deserialize, Deserializer, Serialize, Serializer, de::Visitor};
13
14#[repr(transparent)]
15#[derive(Debug, Copy, Clone, PartialOrd, PartialEq, Ord, Eq, Hash)]
16pub struct CommitVersion(pub u64);
17
18impl FromStr for CommitVersion {
19 type Err = ParseIntError;
20
21 fn from_str(s: &str) -> Result<Self, Self::Err> {
22 Ok(CommitVersion(u64::from_str(s)?))
23 }
24}
25
26impl Display for CommitVersion {
27 fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
28 Display::fmt(&self.0, f)
29 }
30}
31
32impl PartialEq<i32> for CommitVersion {
33 fn eq(&self, other: &i32) -> bool {
34 self.0 == *other as u64
35 }
36}
37
38impl PartialEq<CommitVersion> for i32 {
39 fn eq(&self, other: &CommitVersion) -> bool {
40 *self as u64 == other.0
41 }
42}
43
44impl PartialEq<u64> for CommitVersion {
45 fn eq(&self, other: &u64) -> bool {
46 self.0.eq(other)
47 }
48}
49
50impl From<CommitVersion> for u64 {
51 fn from(value: CommitVersion) -> Self {
52 value.0
53 }
54}
55
56impl From<i32> for CommitVersion {
57 fn from(value: i32) -> Self {
58 Self(value as u64)
59 }
60}
61
62impl From<u64> for CommitVersion {
63 fn from(value: u64) -> Self {
64 Self(value)
65 }
66}
67
68impl Serialize for CommitVersion {
69 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
70 where
71 S: Serializer,
72 {
73 serializer.serialize_u64(self.0)
74 }
75}
76
77impl<'de> Deserialize<'de> for CommitVersion {
78 fn deserialize<D>(deserializer: D) -> Result<CommitVersion, D::Error>
79 where
80 D: Deserializer<'de>,
81 {
82 struct U64Visitor;
83
84 impl Visitor<'_> for U64Visitor {
85 type Value = CommitVersion;
86
87 fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result {
88 formatter.write_str("an unsigned 64-bit number")
89 }
90
91 fn visit_u64<E>(self, value: u64) -> Result<Self::Value, E> {
92 Ok(CommitVersion(value))
93 }
94 }
95
96 deserializer.deserialize_u64(U64Visitor)
97 }
98}
99
100#[repr(transparent)]
101#[derive(Debug, Copy, Clone, PartialOrd, PartialEq, Ord, Eq, Hash, Serialize, Deserialize)]
102#[serde(transparent)]
103pub struct SourceVersion(pub u64);
104
105impl From<CommitVersion> for SourceVersion {
106 fn from(version: CommitVersion) -> Self {
107 Self(version.0)
108 }
109}
110
111#[derive(Debug, Copy, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
112pub struct ChangeVersion {
113 pub commit: CommitVersion,
114 pub source: SourceVersion,
115}
116
117impl From<CommitVersion> for ChangeVersion {
118 fn from(commit: CommitVersion) -> Self {
119 Self {
120 commit,
121 source: SourceVersion::from(commit),
122 }
123 }
124}
125
126#[derive(Debug, Clone, Copy, PartialEq, Serialize, Deserialize, Default)]
127pub enum JoinType {
128 Inner,
129 #[default]
130 Left,
131}
132
133#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Default)]
134pub enum IndexType {
135 #[default]
136 Index,
137 Unique,
138 Primary,
139}
140
141#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
142pub enum WindowSize {
143 Duration(Duration),
144 Count(u64),
145}
146
147impl WindowSize {
148 pub fn is_count(&self) -> bool {
149 matches!(self, WindowSize::Count(_))
150 }
151
152 pub fn as_duration(&self) -> Option<Duration> {
153 match self {
154 WindowSize::Duration(d) => Some(*d),
155 _ => None,
156 }
157 }
158
159 pub fn as_count(&self) -> Option<u64> {
160 match self {
161 WindowSize::Count(c) => Some(*c),
162 _ => None,
163 }
164 }
165}
166
167#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
168pub enum TimeDomain {
169 None,
170 Event,
171 Processing,
172}
173
174impl TimeDomain {
175 pub fn to_u8(self) -> u8 {
176 match self {
177 TimeDomain::None => 0,
178 TimeDomain::Event => 1,
179 TimeDomain::Processing => 2,
180 }
181 }
182
183 pub fn from_u8(value: u8) -> Self {
184 match value {
185 1 => TimeDomain::Event,
186 2 => TimeDomain::Processing,
187 _ => TimeDomain::None,
188 }
189 }
190
191 pub fn as_str(&self) -> &'static str {
192 match self {
193 TimeDomain::None => "none",
194 TimeDomain::Event => "event",
195 TimeDomain::Processing => "processing",
196 }
197 }
198}
199
200#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
201pub enum TimeSource {
202 None,
203 Event {
204 ts: String,
205 },
206 Processing,
207}
208
209impl TimeSource {
210 pub fn domain(&self) -> TimeDomain {
211 match self {
212 TimeSource::None => TimeDomain::None,
213 TimeSource::Event {
214 ..
215 } => TimeDomain::Event,
216 TimeSource::Processing => TimeDomain::Processing,
217 }
218 }
219
220 pub fn ts(&self) -> Option<&str> {
221 match self {
222 TimeSource::Event {
223 ts,
224 } => Some(ts.as_str()),
225 TimeSource::None | TimeSource::Processing => None,
226 }
227 }
228}
229
230#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
231pub enum WindowKind {
232 Tumbling {
233 size: WindowSize,
234 },
235
236 Sliding {
237 size: WindowSize,
238 slide: WindowSize,
239 },
240
241 Rolling {
242 size: WindowSize,
243 #[serde(default)]
244 lag: Option<Duration>,
245 },
246
247 Session {
248 gap: Duration,
249 },
250}
251
252impl WindowKind {
253 pub fn size(&self) -> Option<&WindowSize> {
254 match self {
255 WindowKind::Tumbling {
256 size,
257 ..
258 } => Some(size),
259 WindowKind::Sliding {
260 size,
261 ..
262 } => Some(size),
263 WindowKind::Rolling {
264 size,
265 ..
266 } => Some(size),
267 WindowKind::Session {
268 ..
269 } => None,
270 }
271 }
272}