Skip to main content

s2_common/
config.rs

1//! Stream and basin configuration types.
2//!
3//! Stream configuration uses three representations:
4//!
5//! - Resolved (`StreamConfig`, `TimestampingConfig`, `DeleteOnEmptyConfig`): concrete values,
6//!   produced by merging optional configs with defaults using `merge()`.
7//!
8//! - Optional (`OptionalStreamConfig`, `OptionalTimestampingConfig`,
9//!   `OptionalDeleteOnEmptyConfig`): partial configuration layers, where `None` means "not set at
10//!   this layer; fall back to defaults."
11//!
12//! - Reconfiguration (`StreamReconfiguration`, `TimestampingReconfiguration`,
13//!   `DeleteOnEmptyReconfiguration`): PATCH-style updates applied with `reconfigure()`.
14//!
15//! Reconfiguration of nested fields (e.g. `timestamping`, `delete_on_empty`,
16//! `default_stream_config`) is applied recursively: `Specified(Some(inner_reconfig))`
17//! applies the inner reconfiguration to the existing value, while `Specified(None)`
18//! clears it to the default.
19//!
20//! `merge()` resolves optional configs into resolved configs with precedence:
21//! stream-level → basin-level → system default (via `Option::or` chaining).
22//!
23//! Basin config also carries basin-level knobs like `stream_cipher`,
24//! `create_stream_on_append`, and `create_stream_on_read`.
25
26use 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}