1use std::time::Duration;
27
28use crate::{ValidationError, encryption::EncryptionAlgorithm, maybe::Maybe};
29
30#[derive(
31 Debug,
32 Default,
33 Clone,
34 Copy,
35 strum::Display,
36 strum::IntoStaticStr,
37 strum::EnumIter,
38 strum::FromRepr,
39 strum::EnumString,
40 PartialEq,
41 Eq,
42 Hash,
43)]
44#[cfg_attr(feature = "clap", derive(clap::ValueEnum))]
45#[repr(u8)]
46pub enum StorageClass {
47 #[strum(serialize = "standard")]
48 Standard = 1,
49 #[default]
50 #[strum(serialize = "express")]
51 Express = 2,
52}
53
54#[derive(Debug, Clone, Copy, PartialEq, Eq)]
55pub enum RetentionPolicy {
56 Age(Duration),
57 Infinite(),
58}
59
60impl RetentionPolicy {
61 pub fn age(&self) -> Option<Duration> {
62 match self {
63 Self::Age(duration) => Some(*duration),
64 Self::Infinite() => None,
65 }
66 }
67
68 pub fn validate(self) -> Result<Self, ValidationError> {
69 match self {
70 Self::Age(duration) if duration.is_zero() => Err(ValidationError(
71 "age must be greater than 0 seconds".to_string(),
72 )),
73 policy => Ok(policy),
74 }
75 }
76}
77
78impl Default for RetentionPolicy {
79 fn default() -> Self {
80 const ONE_WEEK: Duration = Duration::from_secs(7 * 24 * 60 * 60);
81
82 Self::Age(ONE_WEEK)
83 }
84}
85
86#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
87pub enum TimestampingMode {
88 #[default]
89 ClientPrefer,
90 ClientRequire,
91 Arrival,
92}
93
94#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
95pub struct TimestampingConfig {
96 pub mode: TimestampingMode,
97 pub uncapped: bool,
98}
99
100#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
101pub struct DeleteOnEmptyConfig {
102 pub min_age: Duration,
103}
104
105impl DeleteOnEmptyConfig {
106 pub fn min_age(&self) -> Option<Duration> {
107 Some(self.min_age).filter(|age| !age.is_zero())
108 }
109}
110
111#[derive(Debug, Clone, Default, PartialEq, Eq)]
112pub struct StreamConfig {
113 pub storage_class: StorageClass,
114 pub retention_policy: RetentionPolicy,
115 pub timestamping: TimestampingConfig,
116 pub delete_on_empty: DeleteOnEmptyConfig,
117}
118
119#[derive(Debug, Clone, Default)]
120pub struct TimestampingReconfiguration {
121 pub mode: Maybe<Option<TimestampingMode>>,
122 pub uncapped: Maybe<Option<bool>>,
123}
124
125#[derive(Debug, Clone, Default)]
126pub struct DeleteOnEmptyReconfiguration {
127 pub min_age: Maybe<Option<Duration>>,
128}
129
130#[derive(Debug, Clone, Default)]
131pub struct StreamReconfiguration {
132 pub storage_class: Maybe<Option<StorageClass>>,
133 pub retention_policy: Maybe<Option<RetentionPolicy>>,
134 pub timestamping: Maybe<Option<TimestampingReconfiguration>>,
135 pub delete_on_empty: Maybe<Option<DeleteOnEmptyReconfiguration>>,
136}
137
138#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
139pub struct OptionalTimestampingConfig {
140 pub mode: Option<TimestampingMode>,
141 pub uncapped: Option<bool>,
142}
143
144impl OptionalTimestampingConfig {
145 pub fn reconfigure(mut self, reconfiguration: TimestampingReconfiguration) -> Self {
146 if let Maybe::Specified(mode) = reconfiguration.mode {
147 self.mode = mode;
148 }
149 if let Maybe::Specified(uncapped) = reconfiguration.uncapped {
150 self.uncapped = uncapped;
151 }
152 self
153 }
154
155 pub fn merge(self, basin_defaults: Self) -> TimestampingConfig {
156 let mode = self.mode.or(basin_defaults.mode).unwrap_or_default();
157 let uncapped = self
158 .uncapped
159 .or(basin_defaults.uncapped)
160 .unwrap_or_default();
161 TimestampingConfig { mode, uncapped }
162 }
163}
164
165impl From<OptionalTimestampingConfig> for TimestampingConfig {
166 fn from(value: OptionalTimestampingConfig) -> Self {
167 Self {
168 mode: value.mode.unwrap_or_default(),
169 uncapped: value.uncapped.unwrap_or_default(),
170 }
171 }
172}
173
174impl From<TimestampingConfig> for OptionalTimestampingConfig {
175 fn from(value: TimestampingConfig) -> Self {
176 Self {
177 mode: Some(value.mode),
178 uncapped: Some(value.uncapped),
179 }
180 }
181}
182
183impl From<OptionalTimestampingConfig> for TimestampingReconfiguration {
184 fn from(value: OptionalTimestampingConfig) -> Self {
185 Self {
186 mode: value.mode.into(),
187 uncapped: value.uncapped.into(),
188 }
189 }
190}
191
192#[derive(Debug, Clone, Default, PartialEq, Eq)]
193pub struct OptionalDeleteOnEmptyConfig {
194 pub min_age: Option<Duration>,
195}
196
197impl OptionalDeleteOnEmptyConfig {
198 pub fn reconfigure(mut self, reconfiguration: DeleteOnEmptyReconfiguration) -> Self {
199 if let Maybe::Specified(min_age) = reconfiguration.min_age {
200 self.min_age = min_age;
201 }
202 self
203 }
204
205 pub fn merge(self, basin_defaults: Self) -> DeleteOnEmptyConfig {
206 let min_age = self.min_age.or(basin_defaults.min_age).unwrap_or_default();
207 DeleteOnEmptyConfig { min_age }
208 }
209}
210
211impl From<OptionalDeleteOnEmptyConfig> for DeleteOnEmptyConfig {
212 fn from(value: OptionalDeleteOnEmptyConfig) -> Self {
213 Self {
214 min_age: value.min_age.unwrap_or_default(),
215 }
216 }
217}
218
219impl From<DeleteOnEmptyConfig> for OptionalDeleteOnEmptyConfig {
220 fn from(value: DeleteOnEmptyConfig) -> Self {
221 Self {
222 min_age: Some(value.min_age),
223 }
224 }
225}
226
227impl From<OptionalDeleteOnEmptyConfig> for DeleteOnEmptyReconfiguration {
228 fn from(value: OptionalDeleteOnEmptyConfig) -> Self {
229 Self {
230 min_age: value.min_age.into(),
231 }
232 }
233}
234
235#[derive(Debug, Clone, Default, PartialEq, Eq)]
236pub struct OptionalStreamConfig {
237 pub storage_class: Option<StorageClass>,
238 pub retention_policy: Option<RetentionPolicy>,
239 pub timestamping: OptionalTimestampingConfig,
240 pub delete_on_empty: OptionalDeleteOnEmptyConfig,
241}
242
243impl OptionalStreamConfig {
244 pub fn validate(&self) -> Result<(), ValidationError> {
245 if let Some(retention_policy) = self.retention_policy {
246 retention_policy.validate()?;
247 }
248 Ok(())
249 }
250
251 pub fn reconfigure(mut self, reconfiguration: StreamReconfiguration) -> Self {
252 let StreamReconfiguration {
253 storage_class,
254 retention_policy,
255 timestamping,
256 delete_on_empty,
257 } = reconfiguration;
258 if let Maybe::Specified(storage_class) = storage_class {
259 self.storage_class = storage_class;
260 }
261 if let Maybe::Specified(retention_policy) = retention_policy {
262 self.retention_policy = retention_policy;
263 }
264 if let Maybe::Specified(timestamping) = timestamping {
265 self.timestamping = timestamping
266 .map(|ts| self.timestamping.reconfigure(ts))
267 .unwrap_or_default();
268 }
269 if let Maybe::Specified(delete_on_empty_reconfig) = delete_on_empty {
270 self.delete_on_empty = delete_on_empty_reconfig
271 .map(|reconfig| self.delete_on_empty.reconfigure(reconfig))
272 .unwrap_or_default();
273 }
274 self
275 }
276
277 pub fn merge(self, basin_defaults: Self) -> StreamConfig {
278 let storage_class = self
279 .storage_class
280 .or(basin_defaults.storage_class)
281 .unwrap_or_default();
282
283 let retention_policy = self
284 .retention_policy
285 .or(basin_defaults.retention_policy)
286 .unwrap_or_default();
287
288 let timestamping = self.timestamping.merge(basin_defaults.timestamping);
289
290 let delete_on_empty = self.delete_on_empty.merge(basin_defaults.delete_on_empty);
291
292 StreamConfig {
293 storage_class,
294 retention_policy,
295 timestamping,
296 delete_on_empty,
297 }
298 }
299}
300
301impl From<OptionalStreamConfig> for StreamReconfiguration {
302 fn from(value: OptionalStreamConfig) -> Self {
303 let OptionalStreamConfig {
304 storage_class,
305 retention_policy,
306 timestamping,
307 delete_on_empty,
308 } = value;
309
310 Self {
311 storage_class: storage_class.into(),
312 retention_policy: retention_policy.into(),
313 timestamping: Some(timestamping.into()).into(),
314 delete_on_empty: Some(delete_on_empty.into()).into(),
315 }
316 }
317}
318
319impl From<OptionalStreamConfig> for StreamConfig {
320 fn from(value: OptionalStreamConfig) -> Self {
321 let OptionalStreamConfig {
322 storage_class,
323 retention_policy,
324 timestamping,
325 delete_on_empty,
326 } = value;
327
328 Self {
329 storage_class: storage_class.unwrap_or_default(),
330 retention_policy: retention_policy.unwrap_or_default(),
331 timestamping: timestamping.into(),
332 delete_on_empty: delete_on_empty.into(),
333 }
334 }
335}
336
337impl From<StreamConfig> for OptionalStreamConfig {
338 fn from(value: StreamConfig) -> Self {
339 let StreamConfig {
340 storage_class,
341 retention_policy,
342 timestamping,
343 delete_on_empty,
344 } = value;
345
346 Self {
347 storage_class: Some(storage_class),
348 retention_policy: Some(retention_policy),
349 timestamping: timestamping.into(),
350 delete_on_empty: delete_on_empty.into(),
351 }
352 }
353}
354
355#[derive(Debug, Clone, Default, PartialEq, Eq)]
356pub struct BasinConfig {
357 pub default_stream_config: OptionalStreamConfig,
358 pub stream_cipher: Option<EncryptionAlgorithm>,
359 pub create_stream_on_append: bool,
360 pub create_stream_on_read: bool,
361}
362
363impl BasinConfig {
364 pub fn validate(&self) -> Result<(), ValidationError> {
365 self.default_stream_config.validate()
366 }
367
368 pub fn reconfigure(mut self, reconfiguration: BasinReconfiguration) -> Self {
369 let BasinReconfiguration {
370 default_stream_config,
371 stream_cipher,
372 create_stream_on_append,
373 create_stream_on_read,
374 } = reconfiguration;
375
376 if let Maybe::Specified(default_stream_config) = default_stream_config {
377 self.default_stream_config = default_stream_config
378 .map(|reconfig| self.default_stream_config.reconfigure(reconfig))
379 .unwrap_or_default();
380 }
381
382 if let Maybe::Specified(stream_cipher) = stream_cipher {
383 self.stream_cipher = stream_cipher;
384 }
385
386 if let Maybe::Specified(create_stream_on_append) = create_stream_on_append {
387 self.create_stream_on_append = create_stream_on_append;
388 }
389
390 if let Maybe::Specified(create_stream_on_read) = create_stream_on_read {
391 self.create_stream_on_read = create_stream_on_read;
392 }
393
394 self
395 }
396}
397
398impl From<BasinConfig> for BasinReconfiguration {
399 fn from(value: BasinConfig) -> Self {
400 let BasinConfig {
401 default_stream_config,
402 stream_cipher,
403 create_stream_on_append,
404 create_stream_on_read,
405 } = value;
406
407 Self {
408 default_stream_config: Some(default_stream_config.into()).into(),
409 stream_cipher: stream_cipher.into(),
410 create_stream_on_append: create_stream_on_append.into(),
411 create_stream_on_read: create_stream_on_read.into(),
412 }
413 }
414}
415
416#[derive(Debug, Clone, Default)]
417pub struct BasinReconfiguration {
418 pub default_stream_config: Maybe<Option<StreamReconfiguration>>,
419 pub stream_cipher: Maybe<Option<EncryptionAlgorithm>>,
420 pub create_stream_on_append: Maybe<bool>,
421 pub create_stream_on_read: Maybe<bool>,
422}