Skip to main content

s2_api/v1/
config.rs

1use std::{str::FromStr, time::Duration};
2
3use compact_str::CompactString;
4use http::{HeaderName, HeaderValue};
5use s2_common::{http::ParseableHeader, maybe::Maybe};
6use serde::{Deserialize, Serialize};
7
8#[rustfmt::skip]
9#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
10#[cfg_attr(feature = "utoipa", derive(utoipa::ToSchema))]
11#[serde(rename_all = "kebab-case")]
12pub enum RetentionPolicy {
13    /// Age in seconds for automatic trimming of records older than this threshold.
14    /// This must be set to a value greater than 0 seconds.
15    Age(u64),
16    /// Retain records unless explicitly trimmed.
17    Infinite(InfiniteRetention)
18}
19
20#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
21#[cfg_attr(feature = "utoipa", derive(utoipa::ToSchema))]
22#[serde(rename_all = "kebab-case")]
23pub struct InfiniteRetention {}
24
25impl TryFrom<RetentionPolicy> for s2_common::config::RetentionPolicy {
26    type Error = s2_common::ValidationError;
27
28    fn try_from(value: RetentionPolicy) -> Result<Self, Self::Error> {
29        let policy = match value {
30            RetentionPolicy::Age(age) => Self::Age(Duration::from_secs(age)),
31            RetentionPolicy::Infinite(_) => Self::Infinite(),
32        };
33        policy.validate()
34    }
35}
36
37impl From<s2_common::config::RetentionPolicy> for RetentionPolicy {
38    fn from(value: s2_common::config::RetentionPolicy) -> Self {
39        match value {
40            s2_common::config::RetentionPolicy::Age(age) => Self::Age(age.as_secs()),
41            s2_common::config::RetentionPolicy::Infinite() => Self::Infinite(InfiniteRetention {}),
42        }
43    }
44}
45
46#[rustfmt::skip]
47#[derive(Debug, Default, PartialEq, Eq, Clone, Serialize, Deserialize)]
48#[cfg_attr(feature = "utoipa", derive(utoipa::ToSchema))]
49#[serde(rename_all = "kebab-case")]
50pub enum TimestampingMode {
51    /// Prefer client-specified timestamp if present otherwise use arrival time.
52    #[default]
53    ClientPrefer,
54    /// Require a client-specified timestamp and reject the append if it is missing.
55    ClientRequire,
56    /// Use the arrival time and ignore any client-specified timestamp.
57    Arrival,
58}
59
60impl From<TimestampingMode> for s2_common::config::TimestampingMode {
61    fn from(value: TimestampingMode) -> Self {
62        match value {
63            TimestampingMode::ClientPrefer => Self::ClientPrefer,
64            TimestampingMode::ClientRequire => Self::ClientRequire,
65            TimestampingMode::Arrival => Self::Arrival,
66        }
67    }
68}
69
70impl From<s2_common::config::TimestampingMode> for TimestampingMode {
71    fn from(value: s2_common::config::TimestampingMode) -> Self {
72        match value {
73            s2_common::config::TimestampingMode::ClientPrefer => Self::ClientPrefer,
74            s2_common::config::TimestampingMode::ClientRequire => Self::ClientRequire,
75            s2_common::config::TimestampingMode::Arrival => Self::Arrival,
76        }
77    }
78}
79
80#[rustfmt::skip]
81#[derive(Debug, Default, PartialEq, Eq, Clone, Serialize, Deserialize)]
82#[cfg_attr(feature = "utoipa", derive(utoipa::ToSchema))]
83pub struct TimestampingConfig {
84    /// Timestamping mode for appends that influences how timestamps are handled.
85    pub mode: Option<TimestampingMode>,
86    /// Allow client-specified timestamps to exceed the arrival time.
87    /// If this is `false` or not set, client timestamps will be capped at the arrival time.
88    pub uncapped: Option<bool>,
89}
90
91impl TimestampingConfig {
92    pub fn to_opt(config: s2_common::config::OptionalTimestampingConfig) -> Option<Self> {
93        let config = TimestampingConfig {
94            mode: config.mode.map(Into::into),
95            uncapped: config.uncapped,
96        };
97        if config == Self::default() {
98            None
99        } else {
100            Some(config)
101        }
102    }
103}
104
105impl From<s2_common::config::TimestampingConfig> for TimestampingConfig {
106    fn from(value: s2_common::config::TimestampingConfig) -> Self {
107        Self {
108            mode: Some(value.mode.into()),
109            uncapped: Some(value.uncapped),
110        }
111    }
112}
113
114impl From<s2_common::config::OptionalTimestampingConfig> for TimestampingConfig {
115    fn from(value: s2_common::config::OptionalTimestampingConfig) -> Self {
116        Self {
117            mode: value.mode.map(Into::into),
118            uncapped: value.uncapped,
119        }
120    }
121}
122
123impl From<TimestampingConfig> for s2_common::config::OptionalTimestampingConfig {
124    fn from(value: TimestampingConfig) -> Self {
125        Self {
126            mode: value.mode.map(Into::into),
127            uncapped: value.uncapped,
128        }
129    }
130}
131
132#[rustfmt::skip]
133#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
134#[cfg_attr(feature = "utoipa", derive(utoipa::ToSchema))]
135pub struct TimestampingReconfiguration {
136    /// Timestamping mode for appends that influences how timestamps are handled.
137    #[serde(default, skip_serializing_if = "Maybe::is_unspecified")]
138    #[cfg_attr(feature = "utoipa", schema(value_type = Option<TimestampingMode>))]
139    pub mode: Maybe<Option<TimestampingMode>>,
140    /// Allow client-specified timestamps to exceed the arrival time.
141    #[serde(default, skip_serializing_if = "Maybe::is_unspecified")]
142    #[cfg_attr(feature = "utoipa", schema(value_type = Option<bool>))]
143    pub uncapped: Maybe<Option<bool>>,
144}
145
146impl From<TimestampingReconfiguration> for s2_common::config::TimestampingReconfiguration {
147    fn from(value: TimestampingReconfiguration) -> Self {
148        Self {
149            mode: value.mode.map_opt(Into::into),
150            uncapped: value.uncapped,
151        }
152    }
153}
154
155impl From<s2_common::config::TimestampingReconfiguration> for TimestampingReconfiguration {
156    fn from(value: s2_common::config::TimestampingReconfiguration) -> Self {
157        Self {
158            mode: value.mode.map_opt(Into::into),
159            uncapped: value.uncapped,
160        }
161    }
162}
163
164#[rustfmt::skip]
165#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
166#[cfg_attr(feature = "utoipa", derive(utoipa::ToSchema))]
167pub struct DeleteOnEmptyConfig {
168    /// Minimum age in seconds before an empty stream can be deleted.
169    /// Set to 0 (default) to disable delete-on-empty (don't delete automatically).
170    #[serde(default)]
171    pub min_age_secs: u64,
172}
173
174impl DeleteOnEmptyConfig {
175    pub fn to_opt(config: s2_common::config::OptionalDeleteOnEmptyConfig) -> Option<Self> {
176        config.min_age.map(|min_age| DeleteOnEmptyConfig {
177            min_age_secs: min_age.as_secs(),
178        })
179    }
180}
181
182impl From<s2_common::config::DeleteOnEmptyConfig> for DeleteOnEmptyConfig {
183    fn from(value: s2_common::config::DeleteOnEmptyConfig) -> Self {
184        Self {
185            min_age_secs: value.min_age.as_secs(),
186        }
187    }
188}
189
190impl From<s2_common::config::OptionalDeleteOnEmptyConfig> for DeleteOnEmptyConfig {
191    fn from(value: s2_common::config::OptionalDeleteOnEmptyConfig) -> Self {
192        Self {
193            min_age_secs: value.min_age.unwrap_or_default().as_secs(),
194        }
195    }
196}
197
198impl From<DeleteOnEmptyConfig> for s2_common::config::DeleteOnEmptyConfig {
199    fn from(value: DeleteOnEmptyConfig) -> Self {
200        Self {
201            min_age: Duration::from_secs(value.min_age_secs),
202        }
203    }
204}
205
206impl From<DeleteOnEmptyConfig> for s2_common::config::OptionalDeleteOnEmptyConfig {
207    fn from(value: DeleteOnEmptyConfig) -> Self {
208        Self {
209            min_age: Some(Duration::from_secs(value.min_age_secs)),
210        }
211    }
212}
213
214#[rustfmt::skip]
215#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
216#[cfg_attr(feature = "utoipa", derive(utoipa::ToSchema))]
217pub struct DeleteOnEmptyReconfiguration {
218    /// Minimum age in seconds before an empty stream can be deleted.
219    /// Set to 0 to disable delete-on-empty (don't delete automatically).
220    #[serde(default, skip_serializing_if = "Maybe::is_unspecified")]
221    #[cfg_attr(feature = "utoipa", schema(value_type = Option<u64>))]
222    pub min_age_secs: Maybe<Option<u64>>,
223}
224
225impl From<DeleteOnEmptyReconfiguration> for s2_common::config::DeleteOnEmptyReconfiguration {
226    fn from(value: DeleteOnEmptyReconfiguration) -> Self {
227        Self {
228            min_age: value.min_age_secs.map_opt(Duration::from_secs),
229        }
230    }
231}
232
233impl From<s2_common::config::DeleteOnEmptyReconfiguration> for DeleteOnEmptyReconfiguration {
234    fn from(value: s2_common::config::DeleteOnEmptyReconfiguration) -> Self {
235        Self {
236            min_age_secs: value.min_age.map_opt(|d| d.as_secs()),
237        }
238    }
239}
240
241#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
242#[cfg_attr(feature = "utoipa", derive(utoipa::ToSchema))]
243pub enum EncryptionAlgorithm {
244    /// AEGIS-256 authenticated encryption.
245    #[serde(rename = "aegis-256")]
246    Aegis256,
247    /// AES-256-GCM authenticated encryption.
248    #[serde(rename = "aes-256-gcm")]
249    Aes256Gcm,
250}
251
252impl From<EncryptionAlgorithm> for s2_common::encryption::EncryptionAlgorithm {
253    fn from(value: EncryptionAlgorithm) -> Self {
254        match value {
255            EncryptionAlgorithm::Aegis256 => Self::Aegis256,
256            EncryptionAlgorithm::Aes256Gcm => Self::Aes256Gcm,
257        }
258    }
259}
260
261impl From<s2_common::encryption::EncryptionAlgorithm> for EncryptionAlgorithm {
262    fn from(value: s2_common::encryption::EncryptionAlgorithm) -> Self {
263        match value {
264            s2_common::encryption::EncryptionAlgorithm::Aegis256 => Self::Aegis256,
265            s2_common::encryption::EncryptionAlgorithm::Aes256Gcm => Self::Aes256Gcm,
266        }
267    }
268}
269
270#[rustfmt::skip]
271#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
272#[cfg_attr(feature = "utoipa", derive(utoipa::ToSchema))]
273pub struct StreamConfig {
274    /// [Storage class](https://s2.dev/docs/storage-classes) for recent writes.
275    #[cfg_attr(feature = "utoipa", schema(value_type = Option<String>))]
276    pub storage_class: Option<CompactString>,
277    /// Retention policy for the stream.
278    /// If unspecified, the default is to retain records for 7 days.
279    pub retention_policy: Option<RetentionPolicy>,
280    /// Timestamping behavior.
281    pub timestamping: Option<TimestampingConfig>,
282    /// Delete-on-empty configuration.
283    #[serde(default)]
284    pub delete_on_empty: Option<DeleteOnEmptyConfig>,
285}
286
287impl StreamConfig {
288    pub fn to_opt(config: s2_common::config::OptionalStreamConfig) -> Option<Self> {
289        let s2_common::config::OptionalStreamConfig {
290            storage_class,
291            retention_policy,
292            timestamping,
293            delete_on_empty,
294        } = config;
295
296        let config = StreamConfig {
297            storage_class,
298            retention_policy: retention_policy.map(Into::into),
299            timestamping: TimestampingConfig::to_opt(timestamping),
300            delete_on_empty: DeleteOnEmptyConfig::to_opt(delete_on_empty),
301        };
302        if config == Self::default() {
303            None
304        } else {
305            Some(config)
306        }
307    }
308
309    /// Encode as compact JSON for the `s2-stream-config` header.
310    pub fn to_header_value(&self) -> HeaderValue {
311        let json = serde_json::to_string(self).expect("StreamConfig serializes to JSON");
312        HeaderValue::from_str(&json).expect("compact JSON of StreamConfig is a valid header value")
313    }
314}
315
316impl From<s2_common::config::StreamConfig> for StreamConfig {
317    fn from(value: s2_common::config::StreamConfig) -> Self {
318        let s2_common::config::StreamConfig {
319            storage_class,
320            retention_policy,
321            timestamping,
322            delete_on_empty,
323        } = value;
324
325        Self {
326            storage_class,
327            retention_policy: Some(retention_policy.into()),
328            timestamping: Some(timestamping.into()),
329            delete_on_empty: Some(delete_on_empty.into()),
330        }
331    }
332}
333
334pub static STREAM_CONFIG_HEADER: HeaderName = HeaderName::from_static("s2-stream-config");
335
336/// Value of the `s2-stream-config` header: a JSON-encoded [`StreamConfig`] to apply over the
337/// basin's default stream config if the stream is created on append or read. Ignored if the
338/// stream already exists.
339#[derive(Debug, Clone, PartialEq, Eq)]
340pub struct StreamConfigHeader(pub s2_common::config::OptionalStreamConfig);
341
342impl FromStr for StreamConfigHeader {
343    type Err = s2_common::ValidationError;
344
345    fn from_str(s: &str) -> Result<Self, Self::Err> {
346        let config: StreamConfig =
347            serde_json::from_str(s).map_err(|e| format!("invalid JSON: {e}"))?;
348        Ok(Self(config.try_into()?))
349    }
350}
351
352impl ParseableHeader for StreamConfigHeader {
353    fn name() -> &'static HeaderName {
354        &STREAM_CONFIG_HEADER
355    }
356}
357
358impl TryFrom<StreamConfig> for s2_common::config::OptionalStreamConfig {
359    type Error = s2_common::ValidationError;
360
361    fn try_from(value: StreamConfig) -> Result<Self, Self::Error> {
362        let StreamConfig {
363            storage_class,
364            retention_policy,
365            timestamping,
366            delete_on_empty,
367        } = value;
368
369        let retention_policy = match retention_policy {
370            None => None,
371            Some(policy) => Some(policy.try_into()?),
372        };
373
374        let config = Self {
375            storage_class,
376            retention_policy,
377            timestamping: timestamping.map(Into::into).unwrap_or_default(),
378            delete_on_empty: delete_on_empty.map(Into::into).unwrap_or_default(),
379        };
380        config.validate()?;
381        Ok(config)
382    }
383}
384
385#[rustfmt::skip]
386#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
387#[cfg_attr(feature = "utoipa", derive(utoipa::ToSchema))]
388pub struct StreamReconfiguration {
389    /// [Storage class](https://s2.dev/docs/storage-classes) for recent writes.
390    #[serde(default, skip_serializing_if = "Maybe::is_unspecified")]
391    #[cfg_attr(feature = "utoipa", schema(value_type = Option<String>))]
392    pub storage_class: Maybe<Option<CompactString>>,
393    /// Retention policy for the stream.
394    /// If unspecified, the default is to retain records for 7 days.
395    #[serde(default, skip_serializing_if = "Maybe::is_unspecified")]
396    #[cfg_attr(feature = "utoipa", schema(value_type = Option<RetentionPolicy>))]
397    pub retention_policy: Maybe<Option<RetentionPolicy>>,
398    /// Timestamping behavior.
399    #[serde(default, skip_serializing_if = "Maybe::is_unspecified")]
400    #[cfg_attr(feature = "utoipa", schema(value_type = Option<TimestampingReconfiguration>))]
401    pub timestamping: Maybe<Option<TimestampingReconfiguration>>,
402    /// Delete-on-empty configuration.
403    #[serde(default, skip_serializing_if = "Maybe::is_unspecified")]
404    #[cfg_attr(feature = "utoipa", schema(value_type = Option<DeleteOnEmptyReconfiguration>))]
405    pub delete_on_empty: Maybe<Option<DeleteOnEmptyReconfiguration>>,
406}
407
408impl TryFrom<StreamReconfiguration> for s2_common::config::StreamReconfiguration {
409    type Error = s2_common::ValidationError;
410
411    fn try_from(value: StreamReconfiguration) -> Result<Self, Self::Error> {
412        let StreamReconfiguration {
413            storage_class,
414            retention_policy,
415            timestamping,
416            delete_on_empty,
417        } = value;
418
419        Ok(Self {
420            storage_class,
421            retention_policy: retention_policy.try_map_opt(TryInto::try_into)?,
422            timestamping: timestamping.map_opt(Into::into),
423            delete_on_empty: delete_on_empty.map_opt(Into::into),
424        })
425    }
426}
427
428impl From<s2_common::config::StreamReconfiguration> for StreamReconfiguration {
429    fn from(value: s2_common::config::StreamReconfiguration) -> Self {
430        let s2_common::config::StreamReconfiguration {
431            storage_class,
432            retention_policy,
433            timestamping,
434            delete_on_empty,
435        } = value;
436
437        Self {
438            storage_class,
439            retention_policy: retention_policy.map_opt(Into::into),
440            timestamping: timestamping.map_opt(Into::into),
441            delete_on_empty: delete_on_empty.map_opt(Into::into),
442        }
443    }
444}
445
446#[rustfmt::skip]
447#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
448#[cfg_attr(feature = "utoipa", derive(utoipa::ToSchema))]
449pub struct BasinConfig {
450    /// Default stream configuration.
451    pub default_stream_config: Option<StreamConfig>,
452    /// Encryption algorithm to apply to newly created streams in the basin.
453    pub stream_cipher: Option<EncryptionAlgorithm>,
454    /// Create stream on append if it doesn't exist, using the default stream configuration.
455    #[serde(default)]
456    #[cfg_attr(feature = "utoipa", schema(default = false))]
457    pub create_stream_on_append: bool,
458    /// Create stream on read if it doesn't exist, using the default stream configuration.
459    #[serde(default)]
460    #[cfg_attr(feature = "utoipa", schema(default = false))]
461    pub create_stream_on_read: bool,
462}
463
464impl TryFrom<BasinConfig> for s2_common::config::BasinConfig {
465    type Error = s2_common::ValidationError;
466
467    fn try_from(value: BasinConfig) -> Result<Self, Self::Error> {
468        let BasinConfig {
469            default_stream_config,
470            stream_cipher,
471            create_stream_on_append,
472            create_stream_on_read,
473        } = value;
474
475        let config = Self {
476            default_stream_config: match default_stream_config {
477                Some(config) => config.try_into()?,
478                None => Default::default(),
479            },
480            stream_cipher: stream_cipher.map(Into::into),
481            create_stream_on_append,
482            create_stream_on_read,
483        };
484        config.validate()?;
485        Ok(config)
486    }
487}
488
489impl From<s2_common::config::BasinConfig> for BasinConfig {
490    fn from(value: s2_common::config::BasinConfig) -> Self {
491        let s2_common::config::BasinConfig {
492            default_stream_config,
493            stream_cipher,
494            create_stream_on_append,
495            create_stream_on_read,
496        } = value;
497
498        Self {
499            default_stream_config: StreamConfig::to_opt(default_stream_config),
500            stream_cipher: stream_cipher.map(Into::into),
501            create_stream_on_append,
502            create_stream_on_read,
503        }
504    }
505}
506
507#[rustfmt::skip]
508#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
509#[cfg_attr(feature = "utoipa", derive(utoipa::ToSchema))]
510pub struct BasinReconfiguration {
511    /// Basin configuration.
512    #[serde(default, skip_serializing_if = "Maybe::is_unspecified")]
513    #[cfg_attr(feature = "utoipa", schema(value_type = Option<StreamReconfiguration>))]
514    pub default_stream_config: Maybe<Option<StreamReconfiguration>>,
515    /// Encryption algorithm to apply to newly created streams in the basin.
516    #[serde(default, skip_serializing_if = "Maybe::is_unspecified")]
517    #[cfg_attr(feature = "utoipa", schema(value_type = Option<EncryptionAlgorithm>))]
518    pub stream_cipher: Maybe<Option<EncryptionAlgorithm>>,
519    /// Create a stream on append.
520    #[serde(default, skip_serializing_if = "Maybe::is_unspecified")]
521    #[cfg_attr(feature = "utoipa", schema(value_type = Option<bool>))]
522    pub create_stream_on_append: Maybe<bool>,
523    /// Create a stream on read.
524    #[serde(default, skip_serializing_if = "Maybe::is_unspecified")]
525    #[cfg_attr(feature = "utoipa", schema(value_type = Option<bool>))]
526    pub create_stream_on_read: Maybe<bool>,
527}
528
529impl TryFrom<BasinReconfiguration> for s2_common::config::BasinReconfiguration {
530    type Error = s2_common::ValidationError;
531
532    fn try_from(value: BasinReconfiguration) -> Result<Self, Self::Error> {
533        let BasinReconfiguration {
534            default_stream_config,
535            stream_cipher,
536            create_stream_on_append,
537            create_stream_on_read,
538        } = value;
539
540        Ok(Self {
541            default_stream_config: default_stream_config.try_map_opt(TryInto::try_into)?,
542            stream_cipher: stream_cipher.map_opt(Into::into),
543            create_stream_on_append: create_stream_on_append.map(Into::into),
544            create_stream_on_read: create_stream_on_read.map(Into::into),
545        })
546    }
547}
548
549impl From<s2_common::config::BasinReconfiguration> for BasinReconfiguration {
550    fn from(value: s2_common::config::BasinReconfiguration) -> Self {
551        let s2_common::config::BasinReconfiguration {
552            default_stream_config,
553            stream_cipher,
554            create_stream_on_append,
555            create_stream_on_read,
556        } = value;
557
558        Self {
559            default_stream_config: default_stream_config.map_opt(Into::into),
560            stream_cipher: stream_cipher.map_opt(Into::into),
561            create_stream_on_append: create_stream_on_append.map(Into::into),
562            create_stream_on_read: create_stream_on_read.map(Into::into),
563        }
564    }
565}
566
567#[cfg(test)]
568mod tests {
569    use proptest::prelude::*;
570
571    use super::*;
572
573    fn gen_storage_class() -> impl Strategy<Value = CompactString> {
574        "[a-z][a-z0-9-]{0,30}".prop_map(CompactString::from)
575    }
576
577    fn gen_timestamping_mode() -> impl Strategy<Value = TimestampingMode> {
578        prop_oneof![
579            Just(TimestampingMode::ClientPrefer),
580            Just(TimestampingMode::ClientRequire),
581            Just(TimestampingMode::Arrival),
582        ]
583    }
584
585    fn gen_retention_policy() -> impl Strategy<Value = RetentionPolicy> {
586        prop_oneof![
587            any::<u64>().prop_map(RetentionPolicy::Age),
588            Just(RetentionPolicy::Infinite(InfiniteRetention {})),
589        ]
590    }
591
592    fn gen_timestamping_config() -> impl Strategy<Value = TimestampingConfig> {
593        (
594            proptest::option::of(gen_timestamping_mode()),
595            proptest::option::of(any::<bool>()),
596        )
597            .prop_map(|(mode, uncapped)| TimestampingConfig { mode, uncapped })
598    }
599
600    fn gen_delete_on_empty_config() -> impl Strategy<Value = DeleteOnEmptyConfig> {
601        any::<u64>().prop_map(|min_age_secs| DeleteOnEmptyConfig { min_age_secs })
602    }
603
604    fn gen_encryption_algorithm() -> impl Strategy<Value = EncryptionAlgorithm> {
605        prop_oneof![
606            Just(EncryptionAlgorithm::Aegis256),
607            Just(EncryptionAlgorithm::Aes256Gcm),
608        ]
609    }
610
611    fn gen_stream_config() -> impl Strategy<Value = StreamConfig> {
612        (
613            proptest::option::of(gen_storage_class()),
614            proptest::option::of(gen_retention_policy()),
615            proptest::option::of(gen_timestamping_config()),
616            proptest::option::of(gen_delete_on_empty_config()),
617        )
618            .prop_map(
619                |(storage_class, retention_policy, timestamping, delete_on_empty)| StreamConfig {
620                    storage_class,
621                    retention_policy,
622                    timestamping,
623                    delete_on_empty,
624                },
625            )
626    }
627
628    fn gen_basin_config() -> impl Strategy<Value = BasinConfig> {
629        (
630            proptest::option::of(gen_stream_config()),
631            proptest::option::of(gen_encryption_algorithm()),
632            any::<bool>(),
633            any::<bool>(),
634        )
635            .prop_map(
636                |(
637                    default_stream_config,
638                    stream_cipher,
639                    create_stream_on_append,
640                    create_stream_on_read,
641                )| {
642                    BasinConfig {
643                        default_stream_config,
644                        stream_cipher,
645                        create_stream_on_append,
646                        create_stream_on_read,
647                    }
648                },
649            )
650    }
651
652    fn gen_maybe<T: std::fmt::Debug + Clone + 'static>(
653        inner: impl Strategy<Value = T>,
654    ) -> impl Strategy<Value = Maybe<Option<T>>> {
655        prop_oneof![
656            Just(Maybe::Unspecified),
657            Just(Maybe::Specified(None)),
658            inner.prop_map(|v| Maybe::Specified(Some(v))),
659        ]
660    }
661
662    fn gen_stream_reconfiguration() -> impl Strategy<Value = StreamReconfiguration> {
663        (
664            gen_maybe(gen_storage_class()),
665            gen_maybe(gen_retention_policy()),
666            gen_maybe(gen_timestamping_reconfiguration()),
667            gen_maybe(gen_delete_on_empty_reconfiguration()),
668        )
669            .prop_map(
670                |(storage_class, retention_policy, timestamping, delete_on_empty)| {
671                    StreamReconfiguration {
672                        storage_class,
673                        retention_policy,
674                        timestamping,
675                        delete_on_empty,
676                    }
677                },
678            )
679    }
680
681    fn gen_timestamping_reconfiguration() -> impl Strategy<Value = TimestampingReconfiguration> {
682        (gen_maybe(gen_timestamping_mode()), gen_maybe(any::<bool>()))
683            .prop_map(|(mode, uncapped)| TimestampingReconfiguration { mode, uncapped })
684    }
685
686    fn gen_delete_on_empty_reconfiguration() -> impl Strategy<Value = DeleteOnEmptyReconfiguration>
687    {
688        gen_maybe(any::<u64>())
689            .prop_map(|min_age_secs| DeleteOnEmptyReconfiguration { min_age_secs })
690    }
691
692    fn gen_basin_reconfiguration() -> impl Strategy<Value = BasinReconfiguration> {
693        (
694            gen_maybe(gen_stream_reconfiguration()),
695            gen_maybe(gen_encryption_algorithm()),
696            prop_oneof![
697                Just(Maybe::Unspecified),
698                any::<bool>().prop_map(Maybe::Specified),
699            ],
700            prop_oneof![
701                Just(Maybe::Unspecified),
702                any::<bool>().prop_map(Maybe::Specified),
703            ],
704        )
705            .prop_map(
706                |(
707                    default_stream_config,
708                    stream_cipher,
709                    create_stream_on_append,
710                    create_stream_on_read,
711                )| BasinReconfiguration {
712                    default_stream_config,
713                    stream_cipher,
714                    create_stream_on_append,
715                    create_stream_on_read,
716                },
717            )
718    }
719
720    fn gen_internal_optional_stream_config()
721    -> impl Strategy<Value = s2_common::config::OptionalStreamConfig> {
722        (
723            proptest::option::of(gen_storage_class()),
724            proptest::option::of(gen_retention_policy()),
725            proptest::option::of(gen_timestamping_mode()),
726            proptest::option::of(any::<bool>()),
727            proptest::option::of(any::<u64>()),
728        )
729            .prop_map(|(sc, rp, ts_mode, ts_uncapped, doe)| {
730                s2_common::config::OptionalStreamConfig {
731                    storage_class: sc,
732                    retention_policy: rp.map(|rp| match rp {
733                        RetentionPolicy::Age(secs) => {
734                            s2_common::config::RetentionPolicy::Age(Duration::from_secs(secs))
735                        }
736                        RetentionPolicy::Infinite(_) => {
737                            s2_common::config::RetentionPolicy::Infinite()
738                        }
739                    }),
740                    timestamping: s2_common::config::OptionalTimestampingConfig {
741                        mode: ts_mode.map(Into::into),
742                        uncapped: ts_uncapped,
743                    },
744                    delete_on_empty: s2_common::config::OptionalDeleteOnEmptyConfig {
745                        min_age: doe.map(Duration::from_secs),
746                    },
747                }
748            })
749    }
750
751    proptest! {
752        #[test]
753        fn stream_config_conversion_validates(config in gen_stream_config()) {
754            let has_zero_age = matches!(config.retention_policy, Some(RetentionPolicy::Age(0)));
755            let result: Result<s2_common::config::OptionalStreamConfig, _> = config.try_into();
756
757            if has_zero_age {
758                prop_assert!(result.is_err());
759            } else {
760                prop_assert!(result.is_ok());
761            }
762        }
763
764        #[test]
765        fn basin_config_conversion_validates(config in gen_basin_config()) {
766            let has_invalid_config = config.default_stream_config.as_ref().is_some_and(|sc| {
767                matches!(sc.retention_policy, Some(RetentionPolicy::Age(0)))
768            });
769
770            let result: Result<s2_common::config::BasinConfig, _> = config.try_into();
771
772            if has_invalid_config {
773                prop_assert!(result.is_err());
774            } else {
775                prop_assert!(result.is_ok());
776            }
777        }
778
779        #[test]
780        fn stream_reconfiguration_conversion_validates(reconfig in gen_stream_reconfiguration()) {
781            let has_zero_age = matches!(
782                reconfig.retention_policy,
783                Maybe::Specified(Some(RetentionPolicy::Age(0)))
784            );
785            let result: Result<s2_common::config::StreamReconfiguration, _> = reconfig.try_into();
786
787            if has_zero_age {
788                prop_assert!(result.is_err());
789            } else {
790                prop_assert!(result.is_ok());
791            }
792        }
793
794        #[test]
795        fn merge_stream_or_basin_or_default(
796            stream in gen_internal_optional_stream_config(),
797            basin in gen_internal_optional_stream_config(),
798        ) {
799            let merged = stream.clone().merge(basin.clone());
800
801            prop_assert_eq!(
802                merged.storage_class,
803                stream.storage_class.or(basin.storage_class)
804            );
805            prop_assert_eq!(
806                merged.retention_policy,
807                stream.retention_policy.or(basin.retention_policy).unwrap_or_default()
808            );
809            prop_assert_eq!(
810                merged.timestamping.mode,
811                stream.timestamping.mode.or(basin.timestamping.mode).unwrap_or_default()
812            );
813            prop_assert_eq!(
814                merged.timestamping.uncapped,
815                stream.timestamping.uncapped.or(basin.timestamping.uncapped).unwrap_or_default()
816            );
817            prop_assert_eq!(
818                merged.delete_on_empty.min_age,
819                stream.delete_on_empty.min_age.or(basin.delete_on_empty.min_age).unwrap_or_default()
820            );
821        }
822
823        #[test]
824        fn reconfigure_unspecified_preserves_base(base in gen_internal_optional_stream_config()) {
825            let reconfig = s2_common::config::StreamReconfiguration::default();
826            let result = base.clone().reconfigure(reconfig);
827
828            prop_assert_eq!(result.storage_class, base.storage_class);
829            prop_assert_eq!(result.retention_policy, base.retention_policy);
830            prop_assert_eq!(result.timestamping.mode, base.timestamping.mode);
831            prop_assert_eq!(result.timestamping.uncapped, base.timestamping.uncapped);
832            prop_assert_eq!(result.delete_on_empty.min_age, base.delete_on_empty.min_age);
833        }
834
835        #[test]
836        fn reconfigure_specified_none_clears(base in gen_internal_optional_stream_config()) {
837            let reconfig = s2_common::config::StreamReconfiguration {
838                storage_class: Maybe::Specified(None),
839                retention_policy: Maybe::Specified(None),
840                timestamping: Maybe::Specified(None),
841                delete_on_empty: Maybe::Specified(None),
842            };
843            let result = base.reconfigure(reconfig);
844
845            prop_assert!(result.storage_class.is_none());
846            prop_assert!(result.retention_policy.is_none());
847            prop_assert!(result.timestamping.mode.is_none());
848            prop_assert!(result.timestamping.uncapped.is_none());
849            prop_assert!(result.delete_on_empty.min_age.is_none());
850        }
851
852        #[test]
853        fn reconfigure_specified_some_sets_value(
854            base in gen_internal_optional_stream_config(),
855            new_sc in gen_storage_class(),
856            new_rp_secs in 1u64..u64::MAX,
857        ) {
858            let reconfig = s2_common::config::StreamReconfiguration {
859                storage_class: Maybe::Specified(Some(new_sc.clone())),
860                retention_policy: Maybe::Specified(Some(
861                    s2_common::config::RetentionPolicy::Age(Duration::from_secs(new_rp_secs))
862                )),
863                ..Default::default()
864            };
865            let result = base.reconfigure(reconfig);
866
867            prop_assert_eq!(result.storage_class, Some(new_sc));
868            prop_assert_eq!(
869                result.retention_policy,
870                Some(s2_common::config::RetentionPolicy::Age(Duration::from_secs(new_rp_secs)))
871            );
872        }
873
874        #[test]
875        fn to_opt_returns_some_for_non_defaults(
876            sc in gen_storage_class(),
877            doe_secs in 1u64..u64::MAX,
878            ts_mode in gen_timestamping_mode(),
879        ) {
880            // non-default storage class -> Some
881            let internal = s2_common::config::OptionalStreamConfig {
882                storage_class: Some(sc),
883                ..Default::default()
884            };
885            prop_assert!(StreamConfig::to_opt(internal).is_some());
886
887            // non-zero delete_on_empty -> Some
888            let internal = s2_common::config::OptionalDeleteOnEmptyConfig {
889                min_age: Some(Duration::from_secs(doe_secs)),
890            };
891            let api = DeleteOnEmptyConfig::to_opt(internal);
892            prop_assert!(api.is_some());
893            prop_assert_eq!(api.unwrap().min_age_secs, doe_secs);
894
895            // non-default timestamping -> Some
896            let internal = s2_common::config::OptionalTimestampingConfig {
897                mode: Some(ts_mode.into()),
898                uncapped: None,
899            };
900            prop_assert!(TimestampingConfig::to_opt(internal).is_some());
901        }
902
903        #[test]
904        fn basin_reconfiguration_conversion_validates(reconfig in gen_basin_reconfiguration()) {
905            let has_zero_age = matches!(
906                &reconfig.default_stream_config,
907                Maybe::Specified(Some(sr)) if matches!(
908                    sr.retention_policy,
909                    Maybe::Specified(Some(RetentionPolicy::Age(0)))
910                )
911            );
912            let result: Result<s2_common::config::BasinReconfiguration, _> = reconfig.try_into();
913
914            if has_zero_age {
915                prop_assert!(result.is_err());
916            } else {
917                prop_assert!(result.is_ok());
918            }
919        }
920
921        #[test]
922        fn reconfigure_basin_unspecified_preserves(
923            base_sc in proptest::option::of(gen_storage_class()),
924            base_algorithm in proptest::option::of(gen_encryption_algorithm()),
925            base_on_append in any::<bool>(),
926            base_on_read in any::<bool>(),
927        ) {
928            let base = s2_common::config::BasinConfig {
929                default_stream_config: s2_common::config::OptionalStreamConfig {
930                    storage_class: base_sc,
931                    ..Default::default()
932                },
933                stream_cipher: base_algorithm.map(Into::into),
934                create_stream_on_append: base_on_append,
935                create_stream_on_read: base_on_read,
936            };
937
938            let reconfig = s2_common::config::BasinReconfiguration::default();
939            let result = base.clone().reconfigure(reconfig);
940
941            prop_assert_eq!(result.default_stream_config.storage_class, base.default_stream_config.storage_class);
942            prop_assert_eq!(result.stream_cipher, base.stream_cipher);
943            prop_assert_eq!(result.create_stream_on_append, base.create_stream_on_append);
944            prop_assert_eq!(result.create_stream_on_read, base.create_stream_on_read);
945        }
946
947        #[test]
948        fn reconfigure_basin_specified_updates(
949            base_on_append in any::<bool>(),
950            new_on_append in any::<bool>(),
951            new_sc in gen_storage_class(),
952            new_algorithm in gen_encryption_algorithm(),
953        ) {
954            let base = s2_common::config::BasinConfig {
955                create_stream_on_append: base_on_append,
956                ..Default::default()
957            };
958
959            let reconfig = s2_common::config::BasinReconfiguration {
960                default_stream_config: Maybe::Specified(Some(s2_common::config::StreamReconfiguration {
961                    storage_class: Maybe::Specified(Some(new_sc.clone())),
962                    ..Default::default()
963                })),
964                stream_cipher: Maybe::Specified(Some(new_algorithm.into())),
965                create_stream_on_append: Maybe::Specified(new_on_append),
966                ..Default::default()
967            };
968            let result = base.reconfigure(reconfig);
969
970            prop_assert_eq!(result.default_stream_config.storage_class, Some(new_sc));
971            prop_assert_eq!(result.stream_cipher, Some(new_algorithm.into()));
972            prop_assert_eq!(result.create_stream_on_append, new_on_append);
973        }
974
975        #[test]
976        fn reconfigure_nested_partial_update(
977            base_mode in gen_timestamping_mode(),
978            base_uncapped in any::<bool>(),
979            new_mode in gen_timestamping_mode(),
980        ) {
981            let base = s2_common::config::OptionalStreamConfig {
982                timestamping: s2_common::config::OptionalTimestampingConfig {
983                    mode: Some(base_mode.into()),
984                    uncapped: Some(base_uncapped),
985                },
986                ..Default::default()
987            };
988
989            let expected_mode: s2_common::config::TimestampingMode = new_mode.into();
990
991            let reconfig = s2_common::config::StreamReconfiguration {
992                timestamping: Maybe::Specified(Some(s2_common::config::TimestampingReconfiguration {
993                    mode: Maybe::Specified(Some(expected_mode)),
994                    uncapped: Maybe::Unspecified,
995                })),
996                ..Default::default()
997            };
998            let result = base.reconfigure(reconfig);
999
1000            prop_assert_eq!(result.timestamping.mode, Some(expected_mode));
1001            prop_assert_eq!(result.timestamping.uncapped, Some(base_uncapped));
1002        }
1003    }
1004
1005    #[test]
1006    fn to_opt_returns_none_for_defaults() {
1007        // default stream config -> None
1008        assert!(StreamConfig::to_opt(s2_common::config::OptionalStreamConfig::default()).is_none());
1009
1010        // delete_on_empty: None -> None
1011        let doe_none = s2_common::config::OptionalDeleteOnEmptyConfig { min_age: None };
1012        assert!(DeleteOnEmptyConfig::to_opt(doe_none).is_none());
1013
1014        // default timestamping -> None
1015        assert!(
1016            TimestampingConfig::to_opt(s2_common::config::OptionalTimestampingConfig::default())
1017                .is_none()
1018        );
1019    }
1020
1021    #[test]
1022    fn optional_stream_config_to_opt_preserves_explicit_zero_delete_on_empty() {
1023        let api = StreamConfig::to_opt(s2_common::config::OptionalStreamConfig {
1024            delete_on_empty: s2_common::config::OptionalDeleteOnEmptyConfig {
1025                min_age: Some(Duration::ZERO),
1026            },
1027            ..Default::default()
1028        })
1029        .unwrap();
1030
1031        assert_eq!(
1032            api.delete_on_empty,
1033            Some(DeleteOnEmptyConfig { min_age_secs: 0 })
1034        );
1035    }
1036
1037    #[test]
1038    fn empty_json_converts_to_all_none() {
1039        let json = serde_json::json!({});
1040        let parsed: StreamConfig = serde_json::from_value(json).unwrap();
1041        let internal: s2_common::config::OptionalStreamConfig = parsed.try_into().unwrap();
1042
1043        assert!(
1044            internal.storage_class.is_none(),
1045            "storage_class should be None"
1046        );
1047        assert!(
1048            internal.retention_policy.is_none(),
1049            "retention_policy should be None"
1050        );
1051        assert!(
1052            internal.timestamping.mode.is_none(),
1053            "timestamping.mode should be None"
1054        );
1055        assert!(
1056            internal.timestamping.uncapped.is_none(),
1057            "timestamping.uncapped should be None"
1058        );
1059        assert!(
1060            internal.delete_on_empty.min_age.is_none(),
1061            "delete_on_empty.min_age should be None"
1062        );
1063    }
1064
1065    #[test]
1066    fn stream_config_header_parses_and_validates() {
1067        let header: StreamConfigHeader =
1068            r#"{"retention_policy":{"age":3600},"delete_on_empty":{"min_age_secs":300}}"#
1069                .parse()
1070                .unwrap();
1071        assert_eq!(
1072            header.0,
1073            s2_common::config::OptionalStreamConfig {
1074                retention_policy: Some(s2_common::config::RetentionPolicy::Age(
1075                    Duration::from_secs(3600)
1076                )),
1077                delete_on_empty: s2_common::config::OptionalDeleteOnEmptyConfig {
1078                    min_age: Some(Duration::from_secs(300)),
1079                },
1080                ..Default::default()
1081            }
1082        );
1083
1084        for spaced in [
1085            r#"{ "retention_policy": { "age": 3600 }, "delete_on_empty": { "min_age_secs": 300 } }"#,
1086            "{\t\"delete_on_empty\":\t{\"min_age_secs\":\t300},\t\"retention_policy\":\t{\"age\":\t3600}\t}",
1087            "  {\"retention_policy\":{\"age\":3600},\"delete_on_empty\":{\"min_age_secs\":300}}  ",
1088        ] {
1089            let parsed: StreamConfigHeader = spaced.parse().unwrap();
1090            assert_eq!(parsed, header, "{spaced:?}");
1091        }
1092
1093        let empty: StreamConfigHeader = "{}".parse().unwrap();
1094        assert_eq!(empty.0, Default::default());
1095
1096        let invalid_json = "not json".parse::<StreamConfigHeader>().unwrap_err();
1097        assert!(invalid_json.to_string().contains("invalid JSON"));
1098
1099        let invalid_age =
1100            r#"{"retention_policy":{"age":0}}"#.parse::<StreamConfigHeader>().unwrap_err();
1101        assert!(
1102            invalid_age
1103                .to_string()
1104                .contains("age must be greater than 0 seconds"),
1105            "{invalid_age}"
1106        );
1107    }
1108
1109    #[test]
1110    fn stream_config_header_value_roundtrips() {
1111        let config = StreamConfig {
1112            storage_class: Some("express".into()),
1113            retention_policy: Some(RetentionPolicy::Infinite(InfiniteRetention {})),
1114            timestamping: Some(TimestampingConfig {
1115                mode: Some(TimestampingMode::ClientRequire),
1116                uncapped: Some(true),
1117            }),
1118            delete_on_empty: Some(DeleteOnEmptyConfig { min_age_secs: 60 }),
1119        };
1120        let value = config.to_header_value();
1121        let parsed: StreamConfigHeader = value.to_str().unwrap().parse().unwrap();
1122        assert_eq!(
1123            parsed.0,
1124            s2_common::config::OptionalStreamConfig::try_from(config).unwrap()
1125        );
1126    }
1127}