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