Skip to main content

s2_common/
config.rs

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