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(u64),
16 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 #[default]
53 ClientPrefer,
54 ClientRequire,
56 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 pub mode: Option<TimestampingMode>,
86 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 #[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 #[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 #[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 #[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 #[serde(rename = "aegis-256")]
246 Aegis256,
247 #[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 #[cfg_attr(feature = "utoipa", schema(value_type = Option<String>))]
276 pub storage_class: Option<CompactString>,
277 pub retention_policy: Option<RetentionPolicy>,
280 pub timestamping: Option<TimestampingConfig>,
282 #[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 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#[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 #[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 #[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 #[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 #[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 pub default_stream_config: Option<StreamConfig>,
452 pub stream_cipher: Option<EncryptionAlgorithm>,
454 #[serde(default)]
456 #[cfg_attr(feature = "utoipa", schema(default = false))]
457 pub create_stream_on_append: bool,
458 #[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 #[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 #[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 #[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 #[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 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 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 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 assert!(StreamConfig::to_opt(s2_common::config::OptionalStreamConfig::default()).is_none());
1009
1010 let doe_none = s2_common::config::OptionalDeleteOnEmptyConfig { min_age: None };
1012 assert!(DeleteOnEmptyConfig::to_opt(doe_none).is_none());
1013
1014 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}