1use std::{
4 collections::HashSet,
5 env::VarError,
6 fmt,
7 num::NonZeroU32,
8 ops::{Deref, RangeTo},
9 pin::Pin,
10 str::FromStr,
11 sync::Arc,
12 time::Duration,
13};
14
15use bytes::Bytes;
16use http::{
17 header::HeaderValue,
18 uri::{Authority, Scheme},
19};
20use rand::RngExt;
21use s2_api::{v1 as api, v1::stream::s2s::CompressionAlgorithm};
22pub use s2_common::ValidationError;
24pub use s2_common::access::AccessTokenId;
28pub use s2_common::access::AccessTokenIdPrefix;
30pub use s2_common::access::AccessTokenIdStartAfter;
32pub use s2_common::basin::BasinName;
37pub use s2_common::basin::BasinNamePrefix;
39pub use s2_common::basin::BasinNameStartAfter;
41pub use s2_common::location::LocationName;
46pub use s2_common::stream::StreamName;
50pub use s2_common::stream::StreamNamePrefix;
52pub use s2_common::stream::StreamNameStartAfter;
54pub use s2_common::{
55 caps::RECORD_BATCH_MAX,
56 encryption::{EncryptionAlgorithm, EncryptionKey},
57};
58
59pub(crate) const ONE_MIB: u32 = 1024 * 1024;
60
61use s2_common::{
62 maybe::Maybe,
63 record::{MAX_FENCING_TOKEN_LENGTH, Metered, MeteredSize},
64 resources::ProvisionResult,
65};
66use secrecy::SecretString;
67
68use crate::error::RequestError;
69
70#[derive(Debug, Clone, Copy, PartialEq, Eq)]
76pub struct S2DateTime(time::OffsetDateTime);
77
78impl TryFrom<time::OffsetDateTime> for S2DateTime {
79 type Error = ValidationError;
80
81 fn try_from(dt: time::OffsetDateTime) -> Result<Self, Self::Error> {
82 dt.format(&time::format_description::well_known::Rfc3339)
83 .map_err(|e| ValidationError(format!("not a valid RFC 3339 datetime: {e}")))?;
84 Ok(Self(dt))
85 }
86}
87
88impl From<S2DateTime> for time::OffsetDateTime {
89 fn from(dt: S2DateTime) -> Self {
90 dt.0
91 }
92}
93
94impl FromStr for S2DateTime {
95 type Err = ValidationError;
96
97 fn from_str(s: &str) -> Result<Self, Self::Err> {
98 time::OffsetDateTime::parse(s, &time::format_description::well_known::Rfc3339)
99 .map(Self)
100 .map_err(|e| ValidationError(format!("not a valid RFC 3339 datetime: {e}")))
101 }
102}
103
104impl fmt::Display for S2DateTime {
105 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
106 write!(
107 f,
108 "{}",
109 self.0
110 .format(&time::format_description::well_known::Rfc3339)
111 .expect("RFC3339 formatting should not fail for S2DateTime")
112 )
113 }
114}
115
116#[derive(Debug, Clone, PartialEq)]
118pub(crate) enum BasinAuthority {
119 ParentZone(Authority),
121 Direct(Authority),
123}
124
125#[derive(Debug, Clone)]
127pub struct AccountEndpoint {
128 scheme: Scheme,
129 authority: Authority,
130}
131
132impl AccountEndpoint {
133 pub fn new(endpoint: &str) -> Result<Self, ValidationError> {
135 endpoint.parse()
136 }
137}
138
139impl FromStr for AccountEndpoint {
140 type Err = ValidationError;
141
142 fn from_str(s: &str) -> Result<Self, Self::Err> {
143 let (scheme, authority) = match s.find("://") {
144 Some(idx) => {
145 let scheme: Scheme = s[..idx]
146 .parse()
147 .map_err(|_| "invalid account endpoint scheme".to_string())?;
148 (scheme, &s[idx + 3..])
149 }
150 None => (Scheme::HTTPS, s),
151 };
152 Ok(Self {
153 scheme,
154 authority: authority
155 .parse()
156 .map_err(|e| format!("invalid account endpoint authority: {e}"))?,
157 })
158 }
159}
160
161#[derive(Debug, Clone)]
163pub struct BasinEndpoint {
164 scheme: Scheme,
165 authority: BasinAuthority,
166}
167
168impl BasinEndpoint {
169 pub fn new(endpoint: &str) -> Result<Self, ValidationError> {
171 endpoint.parse()
172 }
173}
174
175impl FromStr for BasinEndpoint {
176 type Err = ValidationError;
177
178 fn from_str(s: &str) -> Result<Self, Self::Err> {
179 let (scheme, authority) = match s.find("://") {
180 Some(idx) => {
181 let scheme: Scheme = s[..idx]
182 .parse()
183 .map_err(|_| "invalid basin endpoint scheme".to_string())?;
184 (scheme, &s[idx + 3..])
185 }
186 None => (Scheme::HTTPS, s),
187 };
188 let authority = if let Some(authority) = authority.strip_prefix("{basin}.") {
189 BasinAuthority::ParentZone(
190 authority
191 .parse()
192 .map_err(|e| format!("invalid basin endpoint authority: {e}"))?,
193 )
194 } else {
195 BasinAuthority::Direct(
196 authority
197 .parse()
198 .map_err(|e| format!("invalid basin endpoint authority: {e}"))?,
199 )
200 };
201 Ok(Self { scheme, authority })
202 }
203}
204
205#[derive(Debug, Clone)]
206#[non_exhaustive]
207pub struct S2Endpoints {
209 pub(crate) scheme: Scheme,
210 pub(crate) account_authority: Authority,
211 pub(crate) basin_authority: BasinAuthority,
212}
213
214impl S2Endpoints {
215 pub fn new(
217 account_endpoint: AccountEndpoint,
218 basin_endpoint: BasinEndpoint,
219 ) -> Result<Self, ValidationError> {
220 if account_endpoint.scheme != basin_endpoint.scheme {
221 return Err("account and basin endpoints must have the same scheme".into());
222 }
223 Ok(Self {
224 scheme: account_endpoint.scheme,
225 account_authority: account_endpoint.authority,
226 basin_authority: basin_endpoint.authority,
227 })
228 }
229
230 pub fn for_endpoint(endpoint: &str) -> Result<Self, ValidationError> {
234 Self::new(
235 AccountEndpoint::new(endpoint)?,
236 BasinEndpoint::new(endpoint)?,
237 )
238 }
239
240 pub fn from_env() -> Result<Self, ValidationError> {
246 let account_endpoint: AccountEndpoint = match std::env::var("S2_ACCOUNT_ENDPOINT") {
247 Ok(endpoint) => endpoint.parse()?,
248 Err(VarError::NotPresent) => return Err("S2_ACCOUNT_ENDPOINT env var not set".into()),
249 Err(VarError::NotUnicode(_)) => {
250 return Err("S2_ACCOUNT_ENDPOINT is not valid unicode".into());
251 }
252 };
253
254 let basin_endpoint: BasinEndpoint = match std::env::var("S2_BASIN_ENDPOINT") {
255 Ok(endpoint) => endpoint.parse()?,
256 Err(VarError::NotPresent) => return Err("S2_BASIN_ENDPOINT env var not set".into()),
257 Err(VarError::NotUnicode(_)) => {
258 return Err("S2_BASIN_ENDPOINT is not valid unicode".into());
259 }
260 };
261
262 if account_endpoint.scheme != basin_endpoint.scheme {
263 return Err(
264 "S2_ACCOUNT_ENDPOINT and S2_BASIN_ENDPOINT must have the same scheme".into(),
265 );
266 }
267
268 Ok(Self {
269 scheme: account_endpoint.scheme,
270 account_authority: account_endpoint.authority,
271 basin_authority: basin_endpoint.authority,
272 })
273 }
274
275 pub fn for_cloud() -> Self {
277 Self {
278 scheme: Scheme::HTTPS,
279 account_authority: "a.s2.dev".try_into().expect("valid authority"),
280 basin_authority: BasinAuthority::ParentZone(
281 "b.s2.dev".try_into().expect("valid authority"),
282 ),
283 }
284 }
285}
286
287#[derive(Debug, Clone, Copy)]
288pub enum Compression {
290 None,
292 Gzip,
294 Zstd,
296}
297
298impl From<Compression> for CompressionAlgorithm {
299 fn from(value: Compression) -> Self {
300 match value {
301 Compression::None => CompressionAlgorithm::None,
302 Compression::Gzip => CompressionAlgorithm::Gzip,
303 Compression::Zstd => CompressionAlgorithm::Zstd,
304 }
305 }
306}
307
308#[derive(Debug, Clone, Copy, PartialEq)]
309#[non_exhaustive]
310pub enum AppendRetryPolicy {
313 All,
315 NoSideEffects,
325}
326
327#[derive(Debug, Clone)]
328#[non_exhaustive]
329pub struct RetryConfig {
338 pub max_attempts: NonZeroU32,
342 pub min_base_delay: Duration,
346 pub max_base_delay: Duration,
350 pub append_retry_policy: AppendRetryPolicy,
355}
356
357impl Default for RetryConfig {
358 fn default() -> Self {
359 Self {
360 max_attempts: NonZeroU32::new(3).expect("valid non-zero u32"),
361 min_base_delay: Duration::from_millis(100),
362 max_base_delay: Duration::from_secs(1),
363 append_retry_policy: AppendRetryPolicy::All,
364 }
365 }
366}
367
368impl RetryConfig {
369 pub fn new() -> Self {
371 Self::default()
372 }
373
374 pub(crate) fn max_retries(&self) -> u32 {
375 self.max_attempts.get() - 1
376 }
377
378 pub fn with_max_attempts(self, max_attempts: NonZeroU32) -> Self {
380 Self {
381 max_attempts,
382 ..self
383 }
384 }
385
386 pub fn with_min_base_delay(self, min_base_delay: Duration) -> Self {
388 Self {
389 min_base_delay,
390 ..self
391 }
392 }
393
394 pub fn with_max_base_delay(self, max_base_delay: Duration) -> Self {
396 Self {
397 max_base_delay,
398 ..self
399 }
400 }
401
402 pub fn with_append_retry_policy(self, append_retry_policy: AppendRetryPolicy) -> Self {
405 Self {
406 append_retry_policy,
407 ..self
408 }
409 }
410}
411
412#[derive(Debug, Clone)]
413#[non_exhaustive]
414pub struct S2Config {
416 pub(crate) access_token: SecretString,
417 pub(crate) endpoints: S2Endpoints,
418 pub(crate) connection_timeout: Duration,
419 pub(crate) request_timeout: Duration,
420 pub(crate) retry: RetryConfig,
421 pub(crate) compression: Compression,
422 pub(crate) user_agent: HeaderValue,
423 pub(crate) insecure_skip_cert_verification: bool,
424 pub(crate) rustls_crypto_provider: Option<Arc<rustls::crypto::CryptoProvider>>,
425}
426
427impl S2Config {
428 pub fn new(access_token: impl Into<String>) -> Self {
430 Self {
431 access_token: access_token.into().into(),
432 endpoints: S2Endpoints::for_cloud(),
433 connection_timeout: Duration::from_secs(3),
434 request_timeout: Duration::from_secs(5),
435 retry: RetryConfig::new(),
436 compression: Compression::None,
437 user_agent: concat!("s2-sdk-rust/", env!("CARGO_PKG_VERSION"))
438 .parse()
439 .expect("valid user agent"),
440 insecure_skip_cert_verification: false,
441 rustls_crypto_provider: default_rustls_crypto_provider(),
442 }
443 }
444
445 pub fn with_endpoints(self, endpoints: S2Endpoints) -> Self {
447 Self { endpoints, ..self }
448 }
449
450 pub fn with_connection_timeout(self, connection_timeout: Duration) -> Self {
454 Self {
455 connection_timeout,
456 ..self
457 }
458 }
459
460 pub fn with_request_timeout(self, request_timeout: Duration) -> Self {
464 Self {
465 request_timeout,
466 ..self
467 }
468 }
469
470 pub fn with_retry(self, retry: RetryConfig) -> Self {
474 Self { retry, ..self }
475 }
476
477 pub fn with_compression(self, compression: Compression) -> Self {
481 Self {
482 compression,
483 ..self
484 }
485 }
486
487 pub fn with_insecure_skip_cert_verification(self, skip: bool) -> Self {
499 Self {
500 insecure_skip_cert_verification: skip,
501 ..self
502 }
503 }
504
505 pub fn with_rustls_crypto_provider(
515 self,
516 provider: impl Into<Arc<rustls::crypto::CryptoProvider>>,
517 ) -> Self {
518 Self {
519 rustls_crypto_provider: Some(provider.into()),
520 ..self
521 }
522 }
523
524 #[cfg(feature = "rustls-aws-lc-rs")]
528 pub fn with_rustls_aws_lc_rs_crypto_provider(self) -> Self {
529 self.with_rustls_crypto_provider(rustls::crypto::aws_lc_rs::default_provider())
530 }
531
532 #[cfg(feature = "rustls-ring")]
536 pub fn with_rustls_ring_crypto_provider(self) -> Self {
537 self.with_rustls_crypto_provider(rustls::crypto::ring::default_provider())
538 }
539
540 #[doc(hidden)]
541 #[cfg(feature = "_hidden")]
542 pub fn with_user_agent(self, user_agent: impl Into<String>) -> Result<Self, ValidationError> {
543 let user_agent = user_agent
544 .into()
545 .parse()
546 .map_err(|e| ValidationError(format!("invalid user agent: {e}")))?;
547 Ok(Self { user_agent, ..self })
548 }
549}
550
551#[cfg(feature = "rustls-aws-lc-rs")]
552fn default_rustls_crypto_provider() -> Option<Arc<rustls::crypto::CryptoProvider>> {
553 Some(Arc::new(rustls::crypto::aws_lc_rs::default_provider()))
554}
555
556#[cfg(all(not(feature = "rustls-aws-lc-rs"), feature = "rustls-ring"))]
557fn default_rustls_crypto_provider() -> Option<Arc<rustls::crypto::CryptoProvider>> {
558 Some(Arc::new(rustls::crypto::ring::default_provider()))
559}
560
561#[cfg(all(not(feature = "rustls-aws-lc-rs"), not(feature = "rustls-ring")))]
562fn default_rustls_crypto_provider() -> Option<Arc<rustls::crypto::CryptoProvider>> {
563 None
564}
565
566#[derive(Debug, Default, Clone, PartialEq, Eq)]
567#[non_exhaustive]
568pub struct Page<T> {
570 pub values: Vec<T>,
572 pub has_more: bool,
574}
575
576impl<T> Page<T> {
577 pub(crate) fn new(values: impl Into<Vec<T>>, has_more: bool) -> Self {
578 Self {
579 values: values.into(),
580 has_more,
581 }
582 }
583}
584
585#[derive(Debug, Clone, Copy, PartialEq, Eq)]
586pub enum StorageClass {
588 Standard,
590 Express,
592}
593
594impl From<api::config::StorageClass> for StorageClass {
595 fn from(value: api::config::StorageClass) -> Self {
596 match value {
597 api::config::StorageClass::Standard => StorageClass::Standard,
598 api::config::StorageClass::Express => StorageClass::Express,
599 }
600 }
601}
602
603impl From<StorageClass> for api::config::StorageClass {
604 fn from(value: StorageClass) -> Self {
605 match value {
606 StorageClass::Standard => api::config::StorageClass::Standard,
607 StorageClass::Express => api::config::StorageClass::Express,
608 }
609 }
610}
611
612#[derive(Debug, Clone, Copy, PartialEq, Eq)]
613pub enum RetentionPolicy {
615 Age(u64),
617 Infinite,
619}
620
621impl From<api::config::RetentionPolicy> for RetentionPolicy {
622 fn from(value: api::config::RetentionPolicy) -> Self {
623 match value {
624 api::config::RetentionPolicy::Age(secs) => RetentionPolicy::Age(secs),
625 api::config::RetentionPolicy::Infinite(_) => RetentionPolicy::Infinite,
626 }
627 }
628}
629
630impl From<RetentionPolicy> for api::config::RetentionPolicy {
631 fn from(value: RetentionPolicy) -> Self {
632 match value {
633 RetentionPolicy::Age(secs) => api::config::RetentionPolicy::Age(secs),
634 RetentionPolicy::Infinite => {
635 api::config::RetentionPolicy::Infinite(api::config::InfiniteRetention {})
636 }
637 }
638 }
639}
640
641#[derive(Debug, Clone, Copy, PartialEq, Eq)]
642pub enum TimestampingMode {
644 ClientPrefer,
646 ClientRequire,
648 Arrival,
650}
651
652impl From<api::config::TimestampingMode> for TimestampingMode {
653 fn from(value: api::config::TimestampingMode) -> Self {
654 match value {
655 api::config::TimestampingMode::ClientPrefer => TimestampingMode::ClientPrefer,
656 api::config::TimestampingMode::ClientRequire => TimestampingMode::ClientRequire,
657 api::config::TimestampingMode::Arrival => TimestampingMode::Arrival,
658 }
659 }
660}
661
662impl From<TimestampingMode> for api::config::TimestampingMode {
663 fn from(value: TimestampingMode) -> Self {
664 match value {
665 TimestampingMode::ClientPrefer => api::config::TimestampingMode::ClientPrefer,
666 TimestampingMode::ClientRequire => api::config::TimestampingMode::ClientRequire,
667 TimestampingMode::Arrival => api::config::TimestampingMode::Arrival,
668 }
669 }
670}
671
672#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
673#[non_exhaustive]
674pub struct TimestampingConfig {
676 pub mode: Option<TimestampingMode>,
680 pub uncapped: Option<bool>,
684}
685
686impl TimestampingConfig {
687 pub fn new() -> Self {
689 Self::default()
690 }
691
692 pub fn with_mode(self, mode: TimestampingMode) -> Self {
694 Self {
695 mode: Some(mode),
696 ..self
697 }
698 }
699
700 pub fn with_uncapped(self, uncapped: bool) -> Self {
702 Self {
703 uncapped: Some(uncapped),
704 ..self
705 }
706 }
707}
708
709impl From<api::config::TimestampingConfig> for TimestampingConfig {
710 fn from(value: api::config::TimestampingConfig) -> Self {
711 Self {
712 mode: value.mode.map(Into::into),
713 uncapped: value.uncapped,
714 }
715 }
716}
717
718impl From<TimestampingConfig> for api::config::TimestampingConfig {
719 fn from(value: TimestampingConfig) -> Self {
720 Self {
721 mode: value.mode.map(Into::into),
722 uncapped: value.uncapped,
723 }
724 }
725}
726
727#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
728#[non_exhaustive]
729pub struct DeleteOnEmptyConfig {
731 pub min_age_secs: u64,
735}
736
737impl DeleteOnEmptyConfig {
738 pub fn new() -> Self {
740 Self::default()
741 }
742
743 pub fn with_min_age(self, min_age: Duration) -> Self {
745 Self {
746 min_age_secs: min_age.as_secs(),
747 }
748 }
749}
750
751impl From<api::config::DeleteOnEmptyConfig> for DeleteOnEmptyConfig {
752 fn from(value: api::config::DeleteOnEmptyConfig) -> Self {
753 Self {
754 min_age_secs: value.min_age_secs,
755 }
756 }
757}
758
759impl From<DeleteOnEmptyConfig> for api::config::DeleteOnEmptyConfig {
760 fn from(value: DeleteOnEmptyConfig) -> Self {
761 Self {
762 min_age_secs: value.min_age_secs,
763 }
764 }
765}
766
767#[derive(Debug, Clone, Default, PartialEq, Eq)]
768#[non_exhaustive]
769pub struct StreamConfig {
771 pub storage_class: Option<StorageClass>,
775 pub retention_policy: Option<RetentionPolicy>,
779 pub timestamping: Option<TimestampingConfig>,
783 pub delete_on_empty: Option<DeleteOnEmptyConfig>,
787}
788
789impl StreamConfig {
790 pub fn new() -> Self {
792 Self::default()
793 }
794
795 pub fn with_storage_class(self, storage_class: StorageClass) -> Self {
797 Self {
798 storage_class: Some(storage_class),
799 ..self
800 }
801 }
802
803 pub fn with_retention_policy(self, retention_policy: RetentionPolicy) -> Self {
805 Self {
806 retention_policy: Some(retention_policy),
807 ..self
808 }
809 }
810
811 pub fn with_timestamping(self, timestamping: TimestampingConfig) -> Self {
813 Self {
814 timestamping: Some(timestamping),
815 ..self
816 }
817 }
818
819 pub fn with_delete_on_empty(self, delete_on_empty: DeleteOnEmptyConfig) -> Self {
821 Self {
822 delete_on_empty: Some(delete_on_empty),
823 ..self
824 }
825 }
826}
827
828impl From<api::config::StreamConfig> for StreamConfig {
829 fn from(value: api::config::StreamConfig) -> Self {
830 Self {
831 storage_class: value.storage_class.map(Into::into),
832 retention_policy: value.retention_policy.map(Into::into),
833 timestamping: value.timestamping.map(Into::into),
834 delete_on_empty: value.delete_on_empty.map(Into::into),
835 }
836 }
837}
838
839impl From<StreamConfig> for api::config::StreamConfig {
840 fn from(value: StreamConfig) -> Self {
841 Self {
842 storage_class: value.storage_class.map(Into::into),
843 retention_policy: value.retention_policy.map(Into::into),
844 timestamping: value.timestamping.map(Into::into),
845 delete_on_empty: value.delete_on_empty.map(Into::into),
846 }
847 }
848}
849
850#[derive(Debug, Clone, Default, PartialEq, Eq)]
851#[non_exhaustive]
852pub struct BasinConfig {
854 pub default_stream_config: Option<StreamConfig>,
858 pub stream_cipher: Option<EncryptionAlgorithm>,
860 pub create_stream_on_append: bool,
864 pub create_stream_on_read: bool,
868}
869
870impl BasinConfig {
871 pub fn new() -> Self {
873 Self::default()
874 }
875
876 pub fn with_default_stream_config(self, config: StreamConfig) -> Self {
878 Self {
879 default_stream_config: Some(config),
880 ..self
881 }
882 }
883
884 pub fn with_stream_cipher(self, stream_cipher: EncryptionAlgorithm) -> Self {
886 Self {
887 stream_cipher: Some(stream_cipher),
888 ..self
889 }
890 }
891
892 pub fn with_create_stream_on_append(self, create_stream_on_append: bool) -> Self {
895 Self {
896 create_stream_on_append,
897 ..self
898 }
899 }
900
901 pub fn with_create_stream_on_read(self, create_stream_on_read: bool) -> Self {
903 Self {
904 create_stream_on_read,
905 ..self
906 }
907 }
908}
909
910impl From<api::config::BasinConfig> for BasinConfig {
911 fn from(value: api::config::BasinConfig) -> Self {
912 Self {
913 default_stream_config: value.default_stream_config.map(Into::into),
914 stream_cipher: value.stream_cipher.map(Into::into),
915 create_stream_on_append: value.create_stream_on_append,
916 create_stream_on_read: value.create_stream_on_read,
917 }
918 }
919}
920
921impl From<BasinConfig> for api::config::BasinConfig {
922 fn from(value: BasinConfig) -> Self {
923 Self {
924 default_stream_config: value.default_stream_config.map(Into::into),
925 stream_cipher: value.stream_cipher.map(Into::into),
926 create_stream_on_append: value.create_stream_on_append,
927 create_stream_on_read: value.create_stream_on_read,
928 }
929 }
930}
931
932#[derive(Debug, Clone)]
933#[non_exhaustive]
934pub struct CreateBasinInput {
936 pub name: BasinName,
938 pub config: Option<BasinConfig>,
942 pub location: Option<LocationName>,
946 idempotency_token: String,
947}
948
949impl CreateBasinInput {
950 pub fn new(name: BasinName) -> Self {
952 Self {
953 name,
954 config: None,
955 location: None,
956 idempotency_token: idempotency_token(),
957 }
958 }
959
960 pub fn with_config(self, config: BasinConfig) -> Self {
962 Self {
963 config: Some(config),
964 ..self
965 }
966 }
967
968 pub fn with_location<S>(self, location: S) -> Result<Self, ValidationError>
970 where
971 S: TryInto<LocationName>,
972 S::Error: fmt::Display,
973 {
974 let location = location
975 .try_into()
976 .map_err(|e| ValidationError(e.to_string()))?;
977 Ok(Self {
978 location: Some(location),
979 ..self
980 })
981 }
982}
983
984impl From<CreateBasinInput> for (api::basin::CreateBasinRequest, String) {
985 fn from(value: CreateBasinInput) -> Self {
986 (
987 api::basin::CreateBasinRequest {
988 basin: value.name,
989 config: value.config.map(Into::into),
990 location: value.location,
991 },
992 value.idempotency_token,
993 )
994 }
995}
996
997#[derive(Debug, Clone)]
998#[non_exhaustive]
999pub struct EnsureBasinInput {
1001 pub name: BasinName,
1003 pub config: Option<BasinConfig>,
1007 pub location: Option<LocationName>,
1012}
1013
1014impl EnsureBasinInput {
1015 pub fn new(name: BasinName) -> Self {
1017 Self {
1018 name,
1019 config: None,
1020 location: None,
1021 }
1022 }
1023
1024 pub fn with_config(self, config: BasinConfig) -> Self {
1026 Self {
1027 config: Some(config),
1028 ..self
1029 }
1030 }
1031
1032 pub fn with_location<S>(self, location: S) -> Result<Self, ValidationError>
1034 where
1035 S: TryInto<LocationName>,
1036 S::Error: fmt::Display,
1037 {
1038 let location = location
1039 .try_into()
1040 .map_err(|e| ValidationError(e.to_string()))?;
1041 Ok(Self {
1042 location: Some(location),
1043 ..self
1044 })
1045 }
1046}
1047
1048impl From<EnsureBasinInput> for (BasinName, Option<api::basin::EnsureBasinRequest>) {
1049 fn from(value: EnsureBasinInput) -> Self {
1050 let config = value.config;
1051 let request = if config.is_some() || value.location.is_some() {
1052 Some(api::basin::EnsureBasinRequest {
1053 config: config.map(Into::into),
1054 location: value.location,
1055 })
1056 } else {
1057 None
1058 };
1059 (value.name, request)
1060 }
1061}
1062
1063#[derive(Debug, Clone)]
1064pub enum EnsureOutput<T> {
1067 Created(T),
1069 ConfigUpdated(T),
1071 ConfigUnchanged(T),
1073}
1074
1075impl<T> From<ProvisionResult<T>> for EnsureOutput<T> {
1076 fn from(result: ProvisionResult<T>) -> Self {
1077 match result {
1078 ProvisionResult::Created(info) => EnsureOutput::Created(info),
1079 ProvisionResult::Updated(info) => EnsureOutput::ConfigUpdated(info),
1080 ProvisionResult::Noop(info) => EnsureOutput::ConfigUnchanged(info),
1081 }
1082 }
1083}
1084
1085#[derive(Debug, Clone, Default)]
1086#[non_exhaustive]
1087pub struct ListBasinsInput {
1089 pub prefix: BasinNamePrefix,
1093 pub start_after: BasinNameStartAfter,
1097 pub limit: Option<usize>,
1101}
1102
1103impl ListBasinsInput {
1104 pub fn new() -> Self {
1106 Self::default()
1107 }
1108
1109 pub fn with_prefix(self, prefix: BasinNamePrefix) -> Self {
1111 Self { prefix, ..self }
1112 }
1113
1114 pub fn with_start_after(self, start_after: BasinNameStartAfter) -> Self {
1117 Self {
1118 start_after,
1119 ..self
1120 }
1121 }
1122
1123 pub fn with_limit(self, limit: usize) -> Self {
1125 Self {
1126 limit: Some(limit),
1127 ..self
1128 }
1129 }
1130}
1131
1132impl From<ListBasinsInput> for api::basin::ListBasinsRequest {
1133 fn from(value: ListBasinsInput) -> Self {
1134 Self {
1135 prefix: Some(value.prefix),
1136 start_after: Some(value.start_after),
1137 limit: value.limit,
1138 }
1139 }
1140}
1141
1142#[derive(Debug, Clone, Default)]
1143pub struct ListAllBasinsInput {
1145 pub prefix: BasinNamePrefix,
1149 pub start_after: BasinNameStartAfter,
1153 pub include_deleted: bool,
1157}
1158
1159impl ListAllBasinsInput {
1160 pub fn new() -> Self {
1162 Self::default()
1163 }
1164
1165 pub fn with_prefix(self, prefix: BasinNamePrefix) -> Self {
1167 Self { prefix, ..self }
1168 }
1169
1170 pub fn with_start_after(self, start_after: BasinNameStartAfter) -> Self {
1173 Self {
1174 start_after,
1175 ..self
1176 }
1177 }
1178
1179 pub fn with_include_deleted(self, include_deleted: bool) -> Self {
1181 Self {
1182 include_deleted,
1183 ..self
1184 }
1185 }
1186}
1187
1188#[derive(Debug, Clone, PartialEq, Eq)]
1189#[non_exhaustive]
1190pub struct BasinInfo {
1192 pub name: BasinName,
1194 pub location: Option<LocationName>,
1196 pub created_at: S2DateTime,
1198 pub deleted_at: Option<S2DateTime>,
1200}
1201
1202impl TryFrom<api::basin::BasinInfo> for BasinInfo {
1203 type Error = ValidationError;
1204
1205 fn try_from(value: api::basin::BasinInfo) -> Result<Self, Self::Error> {
1206 Ok(Self {
1207 name: value.name,
1208 location: value.location,
1209 created_at: value.created_at.try_into()?,
1210 deleted_at: value.deleted_at.map(S2DateTime::try_from).transpose()?,
1211 })
1212 }
1213}
1214
1215#[derive(Debug, Clone)]
1216#[non_exhaustive]
1217pub struct DeleteBasinInput {
1219 pub name: BasinName,
1221 pub ignore_not_found: bool,
1223}
1224
1225impl DeleteBasinInput {
1226 pub fn new(name: BasinName) -> Self {
1228 Self {
1229 name,
1230 ignore_not_found: false,
1231 }
1232 }
1233
1234 pub fn with_ignore_not_found(self, ignore_not_found: bool) -> Self {
1236 Self {
1237 ignore_not_found,
1238 ..self
1239 }
1240 }
1241}
1242
1243#[derive(Debug, Clone, Default)]
1244#[non_exhaustive]
1245pub struct TimestampingReconfiguration {
1247 pub mode: Maybe<Option<TimestampingMode>>,
1249 pub uncapped: Maybe<Option<bool>>,
1251}
1252
1253impl TimestampingReconfiguration {
1254 pub fn new() -> Self {
1256 Self::default()
1257 }
1258
1259 pub fn with_mode(self, mode: TimestampingMode) -> Self {
1261 Self {
1262 mode: Maybe::Specified(Some(mode)),
1263 ..self
1264 }
1265 }
1266
1267 pub fn with_uncapped(self, uncapped: bool) -> Self {
1269 Self {
1270 uncapped: Maybe::Specified(Some(uncapped)),
1271 ..self
1272 }
1273 }
1274}
1275
1276impl From<TimestampingReconfiguration> for api::config::TimestampingReconfiguration {
1277 fn from(value: TimestampingReconfiguration) -> Self {
1278 Self {
1279 mode: value.mode.map(|m| m.map(Into::into)),
1280 uncapped: value.uncapped,
1281 }
1282 }
1283}
1284
1285#[derive(Debug, Clone, Default)]
1286#[non_exhaustive]
1287pub struct DeleteOnEmptyReconfiguration {
1289 pub min_age_secs: Maybe<Option<u64>>,
1291}
1292
1293impl DeleteOnEmptyReconfiguration {
1294 pub fn new() -> Self {
1296 Self::default()
1297 }
1298
1299 pub fn with_min_age(self, min_age: Duration) -> Self {
1301 Self {
1302 min_age_secs: Maybe::Specified(Some(min_age.as_secs())),
1303 }
1304 }
1305}
1306
1307impl From<DeleteOnEmptyReconfiguration> for api::config::DeleteOnEmptyReconfiguration {
1308 fn from(value: DeleteOnEmptyReconfiguration) -> Self {
1309 Self {
1310 min_age_secs: value.min_age_secs,
1311 }
1312 }
1313}
1314
1315#[derive(Debug, Clone, Default)]
1316#[non_exhaustive]
1317pub struct StreamReconfiguration {
1319 pub storage_class: Maybe<Option<StorageClass>>,
1321 pub retention_policy: Maybe<Option<RetentionPolicy>>,
1323 pub timestamping: Maybe<Option<TimestampingReconfiguration>>,
1325 pub delete_on_empty: Maybe<Option<DeleteOnEmptyReconfiguration>>,
1327}
1328
1329impl StreamReconfiguration {
1330 pub fn new() -> Self {
1332 Self::default()
1333 }
1334
1335 pub fn with_storage_class(self, storage_class: StorageClass) -> Self {
1337 Self {
1338 storage_class: Maybe::Specified(Some(storage_class)),
1339 ..self
1340 }
1341 }
1342
1343 pub fn with_retention_policy(self, retention_policy: RetentionPolicy) -> Self {
1345 Self {
1346 retention_policy: Maybe::Specified(Some(retention_policy)),
1347 ..self
1348 }
1349 }
1350
1351 pub fn with_timestamping(self, timestamping: TimestampingReconfiguration) -> Self {
1353 Self {
1354 timestamping: Maybe::Specified(Some(timestamping)),
1355 ..self
1356 }
1357 }
1358
1359 pub fn with_delete_on_empty(self, delete_on_empty: DeleteOnEmptyReconfiguration) -> Self {
1361 Self {
1362 delete_on_empty: Maybe::Specified(Some(delete_on_empty)),
1363 ..self
1364 }
1365 }
1366}
1367
1368impl From<StreamReconfiguration> for api::config::StreamReconfiguration {
1369 fn from(value: StreamReconfiguration) -> Self {
1370 Self {
1371 storage_class: value.storage_class.map(|m| m.map(Into::into)),
1372 retention_policy: value.retention_policy.map(|m| m.map(Into::into)),
1373 timestamping: value.timestamping.map(|m| m.map(Into::into)),
1374 delete_on_empty: value.delete_on_empty.map(|m| m.map(Into::into)),
1375 }
1376 }
1377}
1378
1379#[derive(Debug, Clone, Default)]
1380#[non_exhaustive]
1381pub struct BasinReconfiguration {
1383 pub default_stream_config: Maybe<Option<StreamReconfiguration>>,
1385 pub stream_cipher: Maybe<Option<EncryptionAlgorithm>>,
1387 pub create_stream_on_append: Maybe<bool>,
1390 pub create_stream_on_read: Maybe<bool>,
1392}
1393
1394impl BasinReconfiguration {
1395 pub fn new() -> Self {
1397 Self::default()
1398 }
1399
1400 pub fn with_default_stream_config(self, config: StreamReconfiguration) -> Self {
1403 Self {
1404 default_stream_config: Maybe::Specified(Some(config)),
1405 ..self
1406 }
1407 }
1408
1409 pub fn with_stream_cipher(self, stream_cipher: EncryptionAlgorithm) -> Self {
1411 Self {
1412 stream_cipher: Maybe::Specified(Some(stream_cipher)),
1413 ..self
1414 }
1415 }
1416
1417 pub fn with_create_stream_on_append(self, create_stream_on_append: bool) -> Self {
1420 Self {
1421 create_stream_on_append: Maybe::Specified(create_stream_on_append),
1422 ..self
1423 }
1424 }
1425
1426 pub fn with_create_stream_on_read(self, create_stream_on_read: bool) -> Self {
1429 Self {
1430 create_stream_on_read: Maybe::Specified(create_stream_on_read),
1431 ..self
1432 }
1433 }
1434}
1435
1436impl From<BasinReconfiguration> for api::config::BasinReconfiguration {
1437 fn from(value: BasinReconfiguration) -> Self {
1438 Self {
1439 default_stream_config: value.default_stream_config.map(|m| m.map(Into::into)),
1440 stream_cipher: value.stream_cipher.map(|m| m.map(Into::into)),
1441 create_stream_on_append: value.create_stream_on_append,
1442 create_stream_on_read: value.create_stream_on_read,
1443 }
1444 }
1445}
1446
1447#[derive(Debug, Clone)]
1448#[non_exhaustive]
1449pub struct ReconfigureBasinInput {
1451 pub name: BasinName,
1453 pub config: BasinReconfiguration,
1455}
1456
1457impl ReconfigureBasinInput {
1458 pub fn new(name: BasinName, config: BasinReconfiguration) -> Self {
1460 Self { name, config }
1461 }
1462}
1463
1464#[derive(Debug, Clone, Default)]
1465#[non_exhaustive]
1466pub struct ListAccessTokensInput {
1468 pub prefix: AccessTokenIdPrefix,
1472 pub start_after: AccessTokenIdStartAfter,
1476 pub limit: Option<usize>,
1480}
1481
1482impl ListAccessTokensInput {
1483 pub fn new() -> Self {
1485 Self::default()
1486 }
1487
1488 pub fn with_prefix(self, prefix: AccessTokenIdPrefix) -> Self {
1490 Self { prefix, ..self }
1491 }
1492
1493 pub fn with_start_after(self, start_after: AccessTokenIdStartAfter) -> Self {
1496 Self {
1497 start_after,
1498 ..self
1499 }
1500 }
1501
1502 pub fn with_limit(self, limit: usize) -> Self {
1504 Self {
1505 limit: Some(limit),
1506 ..self
1507 }
1508 }
1509}
1510
1511impl From<ListAccessTokensInput> for api::access::ListAccessTokensRequest {
1512 fn from(value: ListAccessTokensInput) -> Self {
1513 Self {
1514 prefix: Some(value.prefix),
1515 start_after: Some(value.start_after),
1516 limit: value.limit,
1517 }
1518 }
1519}
1520
1521#[derive(Debug, Clone, Default)]
1522pub struct ListAllAccessTokensInput {
1524 pub prefix: AccessTokenIdPrefix,
1528 pub start_after: AccessTokenIdStartAfter,
1532}
1533
1534impl ListAllAccessTokensInput {
1535 pub fn new() -> Self {
1537 Self::default()
1538 }
1539
1540 pub fn with_prefix(self, prefix: AccessTokenIdPrefix) -> Self {
1542 Self { prefix, ..self }
1543 }
1544
1545 pub fn with_start_after(self, start_after: AccessTokenIdStartAfter) -> Self {
1548 Self {
1549 start_after,
1550 ..self
1551 }
1552 }
1553}
1554
1555#[derive(Debug, Clone, PartialEq, Eq)]
1556#[non_exhaustive]
1557pub struct LocationInfo {
1559 pub name: LocationName,
1561 pub is_private: bool,
1563}
1564
1565impl From<api::location::LocationInfo> for LocationInfo {
1566 fn from(value: api::location::LocationInfo) -> Self {
1567 Self {
1568 name: value.name,
1569 is_private: value.is_private,
1570 }
1571 }
1572}
1573
1574#[derive(Debug, Clone)]
1575#[non_exhaustive]
1576pub struct AccessTokenInfo {
1578 pub id: AccessTokenId,
1580 pub expires_at: Option<S2DateTime>,
1582 pub auto_prefix_streams: bool,
1585 pub scope: AccessTokenScope,
1587}
1588
1589impl TryFrom<api::access::AccessTokenInfo> for AccessTokenInfo {
1590 type Error = ValidationError;
1591
1592 fn try_from(value: api::access::AccessTokenInfo) -> Result<Self, Self::Error> {
1593 let expires_at = value.expires_at.map(S2DateTime::try_from).transpose()?;
1594 Ok(Self {
1595 id: value.id,
1596 expires_at,
1597 auto_prefix_streams: value.auto_prefix_streams,
1598 scope: value.scope.into(),
1599 })
1600 }
1601}
1602
1603#[derive(Debug, Clone)]
1604pub enum BasinMatcher {
1608 None,
1610 Exact(BasinName),
1612 Prefix(BasinNamePrefix),
1614}
1615
1616#[derive(Debug, Clone)]
1617pub enum StreamMatcher {
1621 None,
1623 Exact(StreamName),
1625 Prefix(StreamNamePrefix),
1627}
1628
1629#[derive(Debug, Clone)]
1630pub enum AccessTokenMatcher {
1634 None,
1636 Exact(AccessTokenId),
1638 Prefix(AccessTokenIdPrefix),
1640}
1641
1642#[derive(Debug, Clone, Default)]
1643#[non_exhaustive]
1644pub struct ReadWritePermissions {
1646 pub read: bool,
1650 pub write: bool,
1654}
1655
1656impl ReadWritePermissions {
1657 pub fn new() -> Self {
1659 Self::default()
1660 }
1661
1662 pub fn read_only() -> Self {
1664 Self {
1665 read: true,
1666 write: false,
1667 }
1668 }
1669
1670 pub fn write_only() -> Self {
1672 Self {
1673 read: false,
1674 write: true,
1675 }
1676 }
1677
1678 pub fn read_write() -> Self {
1680 Self {
1681 read: true,
1682 write: true,
1683 }
1684 }
1685}
1686
1687impl From<ReadWritePermissions> for api::access::ReadWritePermissions {
1688 fn from(value: ReadWritePermissions) -> Self {
1689 Self {
1690 read: Some(value.read),
1691 write: Some(value.write),
1692 }
1693 }
1694}
1695
1696impl From<api::access::ReadWritePermissions> for ReadWritePermissions {
1697 fn from(value: api::access::ReadWritePermissions) -> Self {
1698 Self {
1699 read: value.read.unwrap_or_default(),
1700 write: value.write.unwrap_or_default(),
1701 }
1702 }
1703}
1704
1705#[derive(Debug, Clone, Default)]
1706#[non_exhaustive]
1707pub struct OperationGroupPermissions {
1711 pub account: Option<ReadWritePermissions>,
1715 pub basin: Option<ReadWritePermissions>,
1719 pub stream: Option<ReadWritePermissions>,
1723}
1724
1725impl OperationGroupPermissions {
1726 pub fn new() -> Self {
1728 Self::default()
1729 }
1730
1731 pub fn read_only_all() -> Self {
1733 Self {
1734 account: Some(ReadWritePermissions::read_only()),
1735 basin: Some(ReadWritePermissions::read_only()),
1736 stream: Some(ReadWritePermissions::read_only()),
1737 }
1738 }
1739
1740 pub fn write_only_all() -> Self {
1742 Self {
1743 account: Some(ReadWritePermissions::write_only()),
1744 basin: Some(ReadWritePermissions::write_only()),
1745 stream: Some(ReadWritePermissions::write_only()),
1746 }
1747 }
1748
1749 pub fn read_write_all() -> Self {
1751 Self {
1752 account: Some(ReadWritePermissions::read_write()),
1753 basin: Some(ReadWritePermissions::read_write()),
1754 stream: Some(ReadWritePermissions::read_write()),
1755 }
1756 }
1757
1758 pub fn with_account(self, account: ReadWritePermissions) -> Self {
1760 Self {
1761 account: Some(account),
1762 ..self
1763 }
1764 }
1765
1766 pub fn with_basin(self, basin: ReadWritePermissions) -> Self {
1768 Self {
1769 basin: Some(basin),
1770 ..self
1771 }
1772 }
1773
1774 pub fn with_stream(self, stream: ReadWritePermissions) -> Self {
1776 Self {
1777 stream: Some(stream),
1778 ..self
1779 }
1780 }
1781}
1782
1783impl From<OperationGroupPermissions> for api::access::PermittedOperationGroups {
1784 fn from(value: OperationGroupPermissions) -> Self {
1785 Self {
1786 account: value.account.map(Into::into),
1787 basin: value.basin.map(Into::into),
1788 stream: value.stream.map(Into::into),
1789 }
1790 }
1791}
1792
1793impl From<api::access::PermittedOperationGroups> for OperationGroupPermissions {
1794 fn from(value: api::access::PermittedOperationGroups) -> Self {
1795 Self {
1796 account: value.account.map(Into::into),
1797 basin: value.basin.map(Into::into),
1798 stream: value.stream.map(Into::into),
1799 }
1800 }
1801}
1802
1803#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
1804pub enum Operation {
1808 ListBasins,
1810 CreateBasin,
1812 GetBasinConfig,
1814 DeleteBasin,
1816 ReconfigureBasin,
1818 ListAccessTokens,
1820 IssueAccessToken,
1822 RevokeAccessToken,
1824 GetAccountMetrics,
1826 GetBasinMetrics,
1828 GetStreamMetrics,
1830 ListStreams,
1832 CreateStream,
1834 GetStreamConfig,
1836 DeleteStream,
1838 ReconfigureStream,
1840 CheckTail,
1842 Append,
1844 Read,
1846 Trim,
1848 Fence,
1850 ListLocations,
1852 GetDefaultLocation,
1854 SetDefaultLocation,
1856}
1857
1858impl From<Operation> for api::access::Operation {
1859 fn from(value: Operation) -> Self {
1860 match value {
1861 Operation::ListBasins => api::access::Operation::ListBasins,
1862 Operation::CreateBasin => api::access::Operation::CreateBasin,
1863 Operation::DeleteBasin => api::access::Operation::DeleteBasin,
1864 Operation::ReconfigureBasin => api::access::Operation::ReconfigureBasin,
1865 Operation::GetBasinConfig => api::access::Operation::GetBasinConfig,
1866 Operation::IssueAccessToken => api::access::Operation::IssueAccessToken,
1867 Operation::RevokeAccessToken => api::access::Operation::RevokeAccessToken,
1868 Operation::ListAccessTokens => api::access::Operation::ListAccessTokens,
1869 Operation::ListStreams => api::access::Operation::ListStreams,
1870 Operation::CreateStream => api::access::Operation::CreateStream,
1871 Operation::DeleteStream => api::access::Operation::DeleteStream,
1872 Operation::GetStreamConfig => api::access::Operation::GetStreamConfig,
1873 Operation::ReconfigureStream => api::access::Operation::ReconfigureStream,
1874 Operation::CheckTail => api::access::Operation::CheckTail,
1875 Operation::Append => api::access::Operation::Append,
1876 Operation::Read => api::access::Operation::Read,
1877 Operation::Trim => api::access::Operation::Trim,
1878 Operation::Fence => api::access::Operation::Fence,
1879 Operation::GetAccountMetrics => api::access::Operation::AccountMetrics,
1880 Operation::GetBasinMetrics => api::access::Operation::BasinMetrics,
1881 Operation::GetStreamMetrics => api::access::Operation::StreamMetrics,
1882 Operation::ListLocations => api::access::Operation::ListLocations,
1883 Operation::GetDefaultLocation => api::access::Operation::GetDefaultLocation,
1884 Operation::SetDefaultLocation => api::access::Operation::SetDefaultLocation,
1885 }
1886 }
1887}
1888
1889impl From<api::access::Operation> for Operation {
1890 fn from(value: api::access::Operation) -> Self {
1891 match value {
1892 api::access::Operation::ListBasins => Operation::ListBasins,
1893 api::access::Operation::CreateBasin => Operation::CreateBasin,
1894 api::access::Operation::DeleteBasin => Operation::DeleteBasin,
1895 api::access::Operation::ReconfigureBasin => Operation::ReconfigureBasin,
1896 api::access::Operation::GetBasinConfig => Operation::GetBasinConfig,
1897 api::access::Operation::IssueAccessToken => Operation::IssueAccessToken,
1898 api::access::Operation::RevokeAccessToken => Operation::RevokeAccessToken,
1899 api::access::Operation::ListAccessTokens => Operation::ListAccessTokens,
1900 api::access::Operation::ListStreams => Operation::ListStreams,
1901 api::access::Operation::CreateStream => Operation::CreateStream,
1902 api::access::Operation::DeleteStream => Operation::DeleteStream,
1903 api::access::Operation::GetStreamConfig => Operation::GetStreamConfig,
1904 api::access::Operation::ReconfigureStream => Operation::ReconfigureStream,
1905 api::access::Operation::CheckTail => Operation::CheckTail,
1906 api::access::Operation::Append => Operation::Append,
1907 api::access::Operation::Read => Operation::Read,
1908 api::access::Operation::Trim => Operation::Trim,
1909 api::access::Operation::Fence => Operation::Fence,
1910 api::access::Operation::AccountMetrics => Operation::GetAccountMetrics,
1911 api::access::Operation::BasinMetrics => Operation::GetBasinMetrics,
1912 api::access::Operation::StreamMetrics => Operation::GetStreamMetrics,
1913 api::access::Operation::ListLocations => Operation::ListLocations,
1914 api::access::Operation::GetDefaultLocation => Operation::GetDefaultLocation,
1915 api::access::Operation::SetDefaultLocation => Operation::SetDefaultLocation,
1916 }
1917 }
1918}
1919
1920#[derive(Debug, Clone)]
1921#[non_exhaustive]
1922pub struct AccessTokenScopeInput {
1930 basins: Option<BasinMatcher>,
1931 streams: Option<StreamMatcher>,
1932 access_tokens: Option<AccessTokenMatcher>,
1933 op_group_perms: Option<OperationGroupPermissions>,
1934 ops: HashSet<Operation>,
1935}
1936
1937impl AccessTokenScopeInput {
1938 pub fn from_ops(ops: impl IntoIterator<Item = Operation>) -> Self {
1940 Self {
1941 basins: None,
1942 streams: None,
1943 access_tokens: None,
1944 op_group_perms: None,
1945 ops: ops.into_iter().collect(),
1946 }
1947 }
1948
1949 pub fn from_op_group_perms(op_group_perms: OperationGroupPermissions) -> Self {
1951 Self {
1952 basins: None,
1953 streams: None,
1954 access_tokens: None,
1955 op_group_perms: Some(op_group_perms),
1956 ops: HashSet::default(),
1957 }
1958 }
1959
1960 pub fn with_ops(self, ops: impl IntoIterator<Item = Operation>) -> Self {
1962 Self {
1963 ops: ops.into_iter().collect(),
1964 ..self
1965 }
1966 }
1967
1968 pub fn with_op_group_perms(self, op_group_perms: OperationGroupPermissions) -> Self {
1970 Self {
1971 op_group_perms: Some(op_group_perms),
1972 ..self
1973 }
1974 }
1975
1976 pub fn with_basins(self, basins: BasinMatcher) -> Self {
1980 Self {
1981 basins: Some(basins),
1982 ..self
1983 }
1984 }
1985
1986 pub fn with_streams(self, streams: StreamMatcher) -> Self {
1990 Self {
1991 streams: Some(streams),
1992 ..self
1993 }
1994 }
1995
1996 pub fn with_access_tokens(self, access_tokens: AccessTokenMatcher) -> Self {
2000 Self {
2001 access_tokens: Some(access_tokens),
2002 ..self
2003 }
2004 }
2005}
2006
2007#[derive(Debug, Clone)]
2008#[non_exhaustive]
2009pub struct AccessTokenScope {
2011 pub basins: Option<BasinMatcher>,
2013 pub streams: Option<StreamMatcher>,
2015 pub access_tokens: Option<AccessTokenMatcher>,
2017 pub op_group_perms: Option<OperationGroupPermissions>,
2019 pub ops: HashSet<Operation>,
2021}
2022
2023impl From<api::access::AccessTokenScope> for AccessTokenScope {
2024 fn from(value: api::access::AccessTokenScope) -> Self {
2025 Self {
2026 basins: value.basins.map(|rs| match rs {
2027 api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e)) => {
2028 BasinMatcher::Exact(e)
2029 }
2030 api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty) => {
2031 BasinMatcher::None
2032 }
2033 api::access::ResourceSet::Prefix(p) => BasinMatcher::Prefix(p),
2034 }),
2035 streams: value.streams.map(|rs| match rs {
2036 api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e)) => {
2037 StreamMatcher::Exact(e)
2038 }
2039 api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty) => {
2040 StreamMatcher::None
2041 }
2042 api::access::ResourceSet::Prefix(p) => StreamMatcher::Prefix(p),
2043 }),
2044 access_tokens: value.access_tokens.map(|rs| match rs {
2045 api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e)) => {
2046 AccessTokenMatcher::Exact(e)
2047 }
2048 api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty) => {
2049 AccessTokenMatcher::None
2050 }
2051 api::access::ResourceSet::Prefix(p) => AccessTokenMatcher::Prefix(p),
2052 }),
2053 op_group_perms: value.op_groups.map(Into::into),
2054 ops: value
2055 .ops
2056 .map(|ops| ops.into_iter().map(Into::into).collect())
2057 .unwrap_or_default(),
2058 }
2059 }
2060}
2061
2062impl From<AccessTokenScopeInput> for api::access::AccessTokenScope {
2063 fn from(value: AccessTokenScopeInput) -> Self {
2064 Self {
2065 basins: value.basins.map(|rs| match rs {
2066 BasinMatcher::None => {
2067 api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty)
2068 }
2069 BasinMatcher::Exact(e) => {
2070 api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e))
2071 }
2072 BasinMatcher::Prefix(p) => api::access::ResourceSet::Prefix(p),
2073 }),
2074 streams: value.streams.map(|rs| match rs {
2075 StreamMatcher::None => {
2076 api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty)
2077 }
2078 StreamMatcher::Exact(e) => {
2079 api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e))
2080 }
2081 StreamMatcher::Prefix(p) => api::access::ResourceSet::Prefix(p),
2082 }),
2083 access_tokens: value.access_tokens.map(|rs| match rs {
2084 AccessTokenMatcher::None => {
2085 api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty)
2086 }
2087 AccessTokenMatcher::Exact(e) => {
2088 api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e))
2089 }
2090 AccessTokenMatcher::Prefix(p) => api::access::ResourceSet::Prefix(p),
2091 }),
2092 op_groups: value.op_group_perms.map(Into::into),
2093 ops: if value.ops.is_empty() {
2094 None
2095 } else {
2096 Some(value.ops.into_iter().map(Into::into).collect())
2097 },
2098 }
2099 }
2100}
2101
2102#[derive(Debug, Clone)]
2103#[non_exhaustive]
2104pub struct IssueAccessTokenInput {
2106 pub id: AccessTokenId,
2108 pub expires_at: Option<S2DateTime>,
2113 pub auto_prefix_streams: bool,
2121 pub scope: AccessTokenScopeInput,
2123}
2124
2125impl IssueAccessTokenInput {
2126 pub fn new(id: AccessTokenId, scope: AccessTokenScopeInput) -> Self {
2128 Self {
2129 id,
2130 expires_at: None,
2131 auto_prefix_streams: false,
2132 scope,
2133 }
2134 }
2135
2136 pub fn with_expires_at(self, expires_at: S2DateTime) -> Self {
2138 Self {
2139 expires_at: Some(expires_at),
2140 ..self
2141 }
2142 }
2143
2144 pub fn with_auto_prefix_streams(self, auto_prefix_streams: bool) -> Self {
2147 Self {
2148 auto_prefix_streams,
2149 ..self
2150 }
2151 }
2152}
2153
2154impl From<IssueAccessTokenInput> for api::access::IssueAccessTokenRequest {
2155 fn from(value: IssueAccessTokenInput) -> Self {
2156 Self {
2157 id: value.id,
2158 expires_at: value.expires_at.map(Into::into),
2159 auto_prefix_streams: value.auto_prefix_streams.then_some(true),
2160 scope: value.scope.into(),
2161 }
2162 }
2163}
2164
2165#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2166pub enum TimeseriesInterval {
2168 Minute,
2170 Hour,
2172 Day,
2174}
2175
2176impl From<TimeseriesInterval> for api::metrics::TimeseriesInterval {
2177 fn from(value: TimeseriesInterval) -> Self {
2178 match value {
2179 TimeseriesInterval::Minute => api::metrics::TimeseriesInterval::Minute,
2180 TimeseriesInterval::Hour => api::metrics::TimeseriesInterval::Hour,
2181 TimeseriesInterval::Day => api::metrics::TimeseriesInterval::Day,
2182 }
2183 }
2184}
2185
2186impl From<api::metrics::TimeseriesInterval> for TimeseriesInterval {
2187 fn from(value: api::metrics::TimeseriesInterval) -> Self {
2188 match value {
2189 api::metrics::TimeseriesInterval::Minute => TimeseriesInterval::Minute,
2190 api::metrics::TimeseriesInterval::Hour => TimeseriesInterval::Hour,
2191 api::metrics::TimeseriesInterval::Day => TimeseriesInterval::Day,
2192 }
2193 }
2194}
2195
2196#[derive(Debug, Clone, Copy)]
2197#[non_exhaustive]
2198pub struct TimeRange {
2200 pub start: u32,
2202 pub end: u32,
2204}
2205
2206impl TimeRange {
2207 pub fn new(start: u32, end: u32) -> Self {
2209 Self { start, end }
2210 }
2211}
2212
2213#[derive(Debug, Clone, Copy)]
2214#[non_exhaustive]
2215pub struct TimeRangeAndInterval {
2217 pub start: u32,
2219 pub end: u32,
2221 pub interval: Option<TimeseriesInterval>,
2225}
2226
2227impl TimeRangeAndInterval {
2228 pub fn new(start: u32, end: u32) -> Self {
2230 Self {
2231 start,
2232 end,
2233 interval: None,
2234 }
2235 }
2236
2237 pub fn with_interval(self, interval: TimeseriesInterval) -> Self {
2239 Self {
2240 interval: Some(interval),
2241 ..self
2242 }
2243 }
2244}
2245
2246#[derive(Debug, Clone, Copy)]
2247pub enum AccountMetricSet {
2249 ActiveBasins(TimeRange),
2252 AccountOps(TimeRangeAndInterval),
2259}
2260
2261#[derive(Debug, Clone)]
2262#[non_exhaustive]
2263pub struct GetAccountMetricsInput {
2265 pub set: AccountMetricSet,
2267}
2268
2269impl GetAccountMetricsInput {
2270 pub fn new(set: AccountMetricSet) -> Self {
2272 Self { set }
2273 }
2274}
2275
2276impl From<GetAccountMetricsInput> for api::metrics::AccountMetricSetRequest {
2277 fn from(value: GetAccountMetricsInput) -> Self {
2278 let (set, start, end, interval) = match value.set {
2279 AccountMetricSet::ActiveBasins(args) => (
2280 api::metrics::AccountMetricSet::ActiveBasins,
2281 args.start,
2282 args.end,
2283 None,
2284 ),
2285 AccountMetricSet::AccountOps(args) => (
2286 api::metrics::AccountMetricSet::AccountOps,
2287 args.start,
2288 args.end,
2289 args.interval,
2290 ),
2291 };
2292 Self {
2293 set,
2294 start: Some(start),
2295 end: Some(end),
2296 interval: interval.map(Into::into),
2297 }
2298 }
2299}
2300
2301#[derive(Debug, Clone, Copy)]
2302pub enum BasinMetricSet {
2304 Storage(TimeRange),
2307 AppendOps(TimeRangeAndInterval),
2315 ReadOps(TimeRangeAndInterval),
2323 ReadThroughput(TimeRangeAndInterval),
2330 AppendThroughput(TimeRangeAndInterval),
2337 BasinOps(TimeRangeAndInterval),
2344}
2345
2346#[derive(Debug, Clone)]
2347#[non_exhaustive]
2348pub struct GetBasinMetricsInput {
2350 pub name: BasinName,
2352 pub set: BasinMetricSet,
2354}
2355
2356impl GetBasinMetricsInput {
2357 pub fn new(name: BasinName, set: BasinMetricSet) -> Self {
2359 Self { name, set }
2360 }
2361}
2362
2363impl From<GetBasinMetricsInput> for (BasinName, api::metrics::BasinMetricSetRequest) {
2364 fn from(value: GetBasinMetricsInput) -> Self {
2365 let (set, start, end, interval) = match value.set {
2366 BasinMetricSet::Storage(args) => (
2367 api::metrics::BasinMetricSet::Storage,
2368 args.start,
2369 args.end,
2370 None,
2371 ),
2372 BasinMetricSet::AppendOps(args) => (
2373 api::metrics::BasinMetricSet::AppendOps,
2374 args.start,
2375 args.end,
2376 args.interval,
2377 ),
2378 BasinMetricSet::ReadOps(args) => (
2379 api::metrics::BasinMetricSet::ReadOps,
2380 args.start,
2381 args.end,
2382 args.interval,
2383 ),
2384 BasinMetricSet::ReadThroughput(args) => (
2385 api::metrics::BasinMetricSet::ReadThroughput,
2386 args.start,
2387 args.end,
2388 args.interval,
2389 ),
2390 BasinMetricSet::AppendThroughput(args) => (
2391 api::metrics::BasinMetricSet::AppendThroughput,
2392 args.start,
2393 args.end,
2394 args.interval,
2395 ),
2396 BasinMetricSet::BasinOps(args) => (
2397 api::metrics::BasinMetricSet::BasinOps,
2398 args.start,
2399 args.end,
2400 args.interval,
2401 ),
2402 };
2403 (
2404 value.name,
2405 api::metrics::BasinMetricSetRequest {
2406 set,
2407 start: Some(start),
2408 end: Some(end),
2409 interval: interval.map(Into::into),
2410 },
2411 )
2412 }
2413}
2414
2415#[derive(Debug, Clone, Copy)]
2416pub enum StreamMetricSet {
2418 Storage(TimeRange),
2421}
2422
2423#[derive(Debug, Clone)]
2424#[non_exhaustive]
2425pub struct GetStreamMetricsInput {
2427 pub basin_name: BasinName,
2429 pub stream_name: StreamName,
2431 pub set: StreamMetricSet,
2433}
2434
2435impl GetStreamMetricsInput {
2436 pub fn new(basin_name: BasinName, stream_name: StreamName, set: StreamMetricSet) -> Self {
2439 Self {
2440 basin_name,
2441 stream_name,
2442 set,
2443 }
2444 }
2445}
2446
2447impl From<GetStreamMetricsInput> for (BasinName, StreamName, api::metrics::StreamMetricSetRequest) {
2448 fn from(value: GetStreamMetricsInput) -> Self {
2449 let (set, start, end, interval) = match value.set {
2450 StreamMetricSet::Storage(args) => (
2451 api::metrics::StreamMetricSet::Storage,
2452 args.start,
2453 args.end,
2454 None,
2455 ),
2456 };
2457 (
2458 value.basin_name,
2459 value.stream_name,
2460 api::metrics::StreamMetricSetRequest {
2461 set,
2462 start: Some(start),
2463 end: Some(end),
2464 interval,
2465 },
2466 )
2467 }
2468}
2469
2470#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2471pub enum MetricUnit {
2473 Bytes,
2475 Operations,
2477}
2478
2479impl From<api::metrics::MetricUnit> for MetricUnit {
2480 fn from(value: api::metrics::MetricUnit) -> Self {
2481 match value {
2482 api::metrics::MetricUnit::Bytes => MetricUnit::Bytes,
2483 api::metrics::MetricUnit::Operations => MetricUnit::Operations,
2484 }
2485 }
2486}
2487
2488#[derive(Debug, Clone)]
2489#[non_exhaustive]
2490pub struct ScalarMetric {
2492 pub name: String,
2494 pub unit: MetricUnit,
2496 pub value: f64,
2498}
2499
2500#[derive(Debug, Clone)]
2501#[non_exhaustive]
2502pub struct AccumulationMetric {
2505 pub name: String,
2507 pub unit: MetricUnit,
2509 pub interval: TimeseriesInterval,
2511 pub values: Vec<(u32, f64)>,
2515}
2516
2517#[derive(Debug, Clone)]
2518#[non_exhaustive]
2519pub struct GaugeMetric {
2521 pub name: String,
2523 pub unit: MetricUnit,
2525 pub values: Vec<(u32, f64)>,
2528}
2529
2530#[derive(Debug, Clone)]
2531#[non_exhaustive]
2532pub struct LabelMetric {
2534 pub name: String,
2536 pub values: Vec<String>,
2538}
2539
2540#[derive(Debug, Clone)]
2541pub enum Metric {
2543 Scalar(ScalarMetric),
2545 Accumulation(AccumulationMetric),
2548 Gauge(GaugeMetric),
2550 Label(LabelMetric),
2552}
2553
2554impl From<api::metrics::Metric> for Metric {
2555 fn from(value: api::metrics::Metric) -> Self {
2556 match value {
2557 api::metrics::Metric::Scalar(sm) => Metric::Scalar(ScalarMetric {
2558 name: sm.name.into(),
2559 unit: sm.unit.into(),
2560 value: sm.value,
2561 }),
2562 api::metrics::Metric::Accumulation(am) => Metric::Accumulation(AccumulationMetric {
2563 name: am.name.into(),
2564 unit: am.unit.into(),
2565 interval: am.interval.into(),
2566 values: am.values,
2567 }),
2568 api::metrics::Metric::Gauge(gm) => Metric::Gauge(GaugeMetric {
2569 name: gm.name.into(),
2570 unit: gm.unit.into(),
2571 values: gm.values,
2572 }),
2573 api::metrics::Metric::Label(lm) => Metric::Label(LabelMetric {
2574 name: lm.name.into(),
2575 values: lm.values,
2576 }),
2577 }
2578 }
2579}
2580
2581#[derive(Debug, Clone, Default)]
2582#[non_exhaustive]
2583pub struct ListStreamsInput {
2585 pub prefix: StreamNamePrefix,
2589 pub start_after: StreamNameStartAfter,
2593 pub limit: Option<usize>,
2597}
2598
2599impl ListStreamsInput {
2600 pub fn new() -> Self {
2602 Self::default()
2603 }
2604
2605 pub fn with_prefix(self, prefix: StreamNamePrefix) -> Self {
2607 Self { prefix, ..self }
2608 }
2609
2610 pub fn with_start_after(self, start_after: StreamNameStartAfter) -> Self {
2613 Self {
2614 start_after,
2615 ..self
2616 }
2617 }
2618
2619 pub fn with_limit(self, limit: usize) -> Self {
2621 Self {
2622 limit: Some(limit),
2623 ..self
2624 }
2625 }
2626}
2627
2628impl From<ListStreamsInput> for api::stream::ListStreamsRequest {
2629 fn from(value: ListStreamsInput) -> Self {
2630 Self {
2631 prefix: Some(value.prefix),
2632 start_after: Some(value.start_after),
2633 limit: value.limit,
2634 }
2635 }
2636}
2637
2638#[derive(Debug, Clone, Default)]
2639pub struct ListAllStreamsInput {
2641 pub prefix: StreamNamePrefix,
2645 pub start_after: StreamNameStartAfter,
2649 pub include_deleted: bool,
2653}
2654
2655impl ListAllStreamsInput {
2656 pub fn new() -> Self {
2658 Self::default()
2659 }
2660
2661 pub fn with_prefix(self, prefix: StreamNamePrefix) -> Self {
2663 Self { prefix, ..self }
2664 }
2665
2666 pub fn with_start_after(self, start_after: StreamNameStartAfter) -> Self {
2669 Self {
2670 start_after,
2671 ..self
2672 }
2673 }
2674
2675 pub fn with_include_deleted(self, include_deleted: bool) -> Self {
2677 Self {
2678 include_deleted,
2679 ..self
2680 }
2681 }
2682}
2683
2684#[derive(Debug, Clone, PartialEq, Eq)]
2685#[non_exhaustive]
2686pub struct StreamInfo {
2688 pub name: StreamName,
2690 pub created_at: S2DateTime,
2692 pub deleted_at: Option<S2DateTime>,
2694 pub cipher: Option<EncryptionAlgorithm>,
2696}
2697
2698impl TryFrom<api::stream::StreamInfo> for StreamInfo {
2699 type Error = ValidationError;
2700
2701 fn try_from(value: api::stream::StreamInfo) -> Result<Self, Self::Error> {
2702 Ok(Self {
2703 name: value.name,
2704 created_at: value.created_at.try_into()?,
2705 deleted_at: value.deleted_at.map(S2DateTime::try_from).transpose()?,
2706 cipher: value.cipher.map(Into::into),
2707 })
2708 }
2709}
2710
2711#[derive(Debug, Clone)]
2712#[non_exhaustive]
2713pub struct CreateStreamInput {
2715 pub name: StreamName,
2717 pub config: Option<StreamConfig>,
2721 idempotency_token: String,
2722}
2723
2724impl CreateStreamInput {
2725 pub fn new(name: StreamName) -> Self {
2727 Self {
2728 name,
2729 config: None,
2730 idempotency_token: idempotency_token(),
2731 }
2732 }
2733
2734 pub fn with_config(self, config: StreamConfig) -> Self {
2736 Self {
2737 config: Some(config),
2738 ..self
2739 }
2740 }
2741}
2742
2743impl From<CreateStreamInput> for (api::stream::CreateStreamRequest, String) {
2744 fn from(value: CreateStreamInput) -> Self {
2745 (
2746 api::stream::CreateStreamRequest {
2747 stream: value.name,
2748 config: value.config.map(Into::into),
2749 },
2750 value.idempotency_token,
2751 )
2752 }
2753}
2754
2755#[derive(Debug, Clone)]
2756#[non_exhaustive]
2757pub struct EnsureStreamInput {
2760 pub name: StreamName,
2762 pub config: Option<StreamConfig>,
2766}
2767
2768impl EnsureStreamInput {
2769 pub fn new(name: StreamName) -> Self {
2771 Self { name, config: None }
2772 }
2773
2774 pub fn with_config(self, config: StreamConfig) -> Self {
2776 Self {
2777 config: Some(config),
2778 ..self
2779 }
2780 }
2781}
2782
2783impl From<EnsureStreamInput> for (StreamName, Option<api::config::StreamConfig>) {
2784 fn from(value: EnsureStreamInput) -> Self {
2785 (value.name, value.config.map(Into::into))
2786 }
2787}
2788
2789#[derive(Debug, Clone)]
2790#[non_exhaustive]
2791pub struct DeleteStreamInput {
2793 pub name: StreamName,
2795 pub ignore_not_found: bool,
2797}
2798
2799impl DeleteStreamInput {
2800 pub fn new(name: StreamName) -> Self {
2802 Self {
2803 name,
2804 ignore_not_found: false,
2805 }
2806 }
2807
2808 pub fn with_ignore_not_found(self, ignore_not_found: bool) -> Self {
2810 Self {
2811 ignore_not_found,
2812 ..self
2813 }
2814 }
2815}
2816
2817#[derive(Debug, Clone)]
2818#[non_exhaustive]
2819pub struct ReconfigureStreamInput {
2821 pub name: StreamName,
2823 pub config: StreamReconfiguration,
2825}
2826
2827impl ReconfigureStreamInput {
2828 pub fn new(name: StreamName, config: StreamReconfiguration) -> Self {
2830 Self { name, config }
2831 }
2832}
2833
2834#[derive(Debug, Clone, PartialEq, Eq)]
2835pub struct FencingToken(String);
2841
2842impl FencingToken {
2843 pub(crate) fn from_server(value: String) -> Self {
2844 Self(value)
2845 }
2846
2847 pub fn generate(n: usize) -> Result<Self, ValidationError> {
2849 rand::rng()
2850 .sample_iter(&rand::distr::Alphanumeric)
2851 .take(n)
2852 .map(char::from)
2853 .collect::<String>()
2854 .parse()
2855 }
2856}
2857
2858impl FromStr for FencingToken {
2859 type Err = ValidationError;
2860
2861 fn from_str(s: &str) -> Result<Self, Self::Err> {
2862 if s.len() > MAX_FENCING_TOKEN_LENGTH {
2863 return Err(ValidationError(format!(
2864 "fencing token exceeds {MAX_FENCING_TOKEN_LENGTH} bytes in length",
2865 )));
2866 }
2867 Ok(FencingToken(s.to_string()))
2868 }
2869}
2870
2871impl std::fmt::Display for FencingToken {
2872 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2873 write!(f, "{}", self.0)
2874 }
2875}
2876
2877impl Deref for FencingToken {
2878 type Target = str;
2879
2880 fn deref(&self) -> &Self::Target {
2881 &self.0
2882 }
2883}
2884
2885#[derive(Debug, Clone, Copy, PartialEq)]
2886#[non_exhaustive]
2887pub struct StreamPosition {
2889 pub seq_num: u64,
2891 pub timestamp: u64,
2894}
2895
2896impl StreamPosition {
2897 pub fn new(seq_num: u64, timestamp: u64) -> Self {
2901 Self { seq_num, timestamp }
2902 }
2903}
2904
2905impl std::fmt::Display for StreamPosition {
2906 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2907 write!(f, "seq_num={}, timestamp={}", self.seq_num, self.timestamp)
2908 }
2909}
2910
2911impl From<api::stream::proto::StreamPosition> for StreamPosition {
2912 fn from(value: api::stream::proto::StreamPosition) -> Self {
2913 Self {
2914 seq_num: value.seq_num,
2915 timestamp: value.timestamp,
2916 }
2917 }
2918}
2919
2920impl From<api::stream::StreamPosition> for StreamPosition {
2921 fn from(value: api::stream::StreamPosition) -> Self {
2922 Self {
2923 seq_num: value.seq_num,
2924 timestamp: value.timestamp,
2925 }
2926 }
2927}
2928
2929#[derive(Debug, Clone, PartialEq)]
2930#[non_exhaustive]
2931pub struct Header {
2933 pub name: Bytes,
2935 pub value: Bytes,
2937}
2938
2939impl Header {
2940 pub fn new(name: impl Into<Bytes>, value: impl Into<Bytes>) -> Self {
2942 Self {
2943 name: name.into(),
2944 value: value.into(),
2945 }
2946 }
2947}
2948
2949impl From<Header> for api::stream::proto::Header {
2950 fn from(value: Header) -> Self {
2951 Self {
2952 name: value.name,
2953 value: value.value,
2954 }
2955 }
2956}
2957
2958impl From<api::stream::proto::Header> for Header {
2959 fn from(value: api::stream::proto::Header) -> Self {
2960 Self {
2961 name: value.name,
2962 value: value.value,
2963 }
2964 }
2965}
2966
2967#[derive(Debug, Clone, PartialEq)]
2968pub struct AppendRecord {
2970 body: Bytes,
2971 headers: Vec<Header>,
2972 timestamp: Option<u64>,
2973}
2974
2975impl AppendRecord {
2976 fn validate(self) -> Result<Self, ValidationError> {
2977 if self.metered_bytes() > RECORD_BATCH_MAX.bytes {
2978 Err(ValidationError(format!(
2979 "metered_bytes: {} exceeds {}",
2980 self.metered_bytes(),
2981 RECORD_BATCH_MAX.bytes
2982 )))
2983 } else {
2984 Ok(self)
2985 }
2986 }
2987
2988 pub fn new(body: impl Into<Bytes>) -> Result<Self, ValidationError> {
2990 let record = Self {
2991 body: body.into(),
2992 headers: Vec::default(),
2993 timestamp: None,
2994 };
2995 record.validate()
2996 }
2997
2998 pub fn with_headers(
3000 self,
3001 headers: impl IntoIterator<Item = Header>,
3002 ) -> Result<Self, ValidationError> {
3003 let record = Self {
3004 headers: headers.into_iter().collect(),
3005 ..self
3006 };
3007 record.validate()
3008 }
3009
3010 pub fn with_timestamp(self, timestamp: u64) -> Self {
3014 Self {
3015 timestamp: Some(timestamp),
3016 ..self
3017 }
3018 }
3019
3020 pub fn body(&self) -> &[u8] {
3022 &self.body
3023 }
3024
3025 pub fn headers(&self) -> &[Header] {
3027 &self.headers
3028 }
3029
3030 pub fn timestamp(&self) -> Option<u64> {
3032 self.timestamp
3033 }
3034}
3035
3036impl From<AppendRecord> for api::stream::proto::AppendRecord {
3037 fn from(value: AppendRecord) -> Self {
3038 Self {
3039 timestamp: value.timestamp,
3040 headers: value.headers.into_iter().map(Into::into).collect(),
3041 body: value.body,
3042 }
3043 }
3044}
3045
3046pub trait MeteredBytes {
3053 fn metered_bytes(&self) -> usize;
3055}
3056
3057macro_rules! metered_bytes_impl {
3058 ($ty:ty) => {
3059 impl MeteredBytes for $ty {
3060 fn metered_bytes(&self) -> usize {
3061 8 + (2 * self.headers.len())
3062 + self
3063 .headers
3064 .iter()
3065 .map(|h| h.name.len() + h.value.len())
3066 .sum::<usize>()
3067 + self.body.len()
3068 }
3069 }
3070 };
3071}
3072
3073metered_bytes_impl!(AppendRecord);
3074
3075impl MeteredSize for AppendRecord {
3076 fn metered_size(&self) -> usize {
3077 self.metered_bytes()
3078 }
3079}
3080
3081#[derive(Debug, Clone)]
3082pub struct AppendRecordBatch(Metered<Vec<AppendRecord>>);
3091
3092impl From<Metered<Vec<AppendRecord>>> for AppendRecordBatch {
3093 fn from(records: Metered<Vec<AppendRecord>>) -> Self {
3094 Self(records)
3095 }
3096}
3097
3098impl AppendRecordBatch {
3099 pub fn try_from_iter<I>(iter: I) -> Result<Self, ValidationError>
3101 where
3102 I: IntoIterator<Item = AppendRecord>,
3103 {
3104 let mut records = Metered::with_capacity(RECORD_BATCH_MAX.count);
3105
3106 for record in iter {
3107 records.push(Metered::from(record));
3108
3109 if records.metered_size() > RECORD_BATCH_MAX.bytes {
3110 return Err(ValidationError(format!(
3111 "batch size in metered bytes ({}) exceeds {}",
3112 records.metered_size(),
3113 RECORD_BATCH_MAX.bytes
3114 )));
3115 }
3116
3117 if records.len() > RECORD_BATCH_MAX.count {
3118 return Err(ValidationError(format!(
3119 "number of records in the batch exceeds {}",
3120 RECORD_BATCH_MAX.count
3121 )));
3122 }
3123 }
3124
3125 if records.is_empty() {
3126 return Err(ValidationError("batch is empty".into()));
3127 }
3128
3129 Ok(records.into())
3130 }
3131}
3132
3133impl Deref for AppendRecordBatch {
3134 type Target = [AppendRecord];
3135
3136 fn deref(&self) -> &Self::Target {
3137 &self.0[..]
3138 }
3139}
3140
3141impl MeteredBytes for AppendRecordBatch {
3142 fn metered_bytes(&self) -> usize {
3143 self.0.metered_size()
3144 }
3145}
3146
3147impl IntoIterator for AppendRecordBatch {
3148 type Item = AppendRecord;
3149 type IntoIter = std::vec::IntoIter<AppendRecord>;
3150
3151 fn into_iter(self) -> Self::IntoIter {
3152 self.0.into_iter()
3153 }
3154}
3155
3156impl<'a> IntoIterator for &'a AppendRecordBatch {
3157 type Item = &'a AppendRecord;
3158 type IntoIter = std::slice::Iter<'a, AppendRecord>;
3159
3160 fn into_iter(self) -> Self::IntoIter {
3161 self.0.iter()
3162 }
3163}
3164
3165#[derive(Debug, Clone)]
3166pub enum Command {
3168 Fence {
3170 fencing_token: FencingToken,
3172 },
3173 Trim {
3175 trim_point: u64,
3177 },
3178}
3179
3180#[derive(Debug, Clone)]
3181#[non_exhaustive]
3182pub struct CommandRecord {
3186 pub command: Command,
3188 pub timestamp: Option<u64>,
3190}
3191
3192impl CommandRecord {
3193 const FENCE: &[u8] = b"fence";
3194 const TRIM: &[u8] = b"trim";
3195
3196 pub fn fence(fencing_token: FencingToken) -> Self {
3201 Self {
3202 command: Command::Fence { fencing_token },
3203 timestamp: None,
3204 }
3205 }
3206
3207 pub fn trim(trim_point: u64) -> Self {
3214 Self {
3215 command: Command::Trim { trim_point },
3216 timestamp: None,
3217 }
3218 }
3219
3220 pub fn with_timestamp(self, timestamp: u64) -> Self {
3222 Self {
3223 timestamp: Some(timestamp),
3224 ..self
3225 }
3226 }
3227}
3228
3229impl From<CommandRecord> for AppendRecord {
3230 fn from(value: CommandRecord) -> Self {
3231 let (header_value, body) = match value.command {
3232 Command::Fence { fencing_token } => (
3233 CommandRecord::FENCE,
3234 Bytes::copy_from_slice(fencing_token.as_bytes()),
3235 ),
3236 Command::Trim { trim_point } => (
3237 CommandRecord::TRIM,
3238 Bytes::copy_from_slice(&trim_point.to_be_bytes()),
3239 ),
3240 };
3241 Self {
3242 body,
3243 headers: vec![Header::new("", header_value)],
3244 timestamp: value.timestamp,
3245 }
3246 }
3247}
3248
3249#[derive(Debug, Clone)]
3250#[non_exhaustive]
3251pub struct AppendInput {
3254 pub records: AppendRecordBatch,
3256 pub match_seq_num: Option<u64>,
3260 pub fencing_token: Option<FencingToken>,
3265}
3266
3267impl AppendInput {
3268 pub fn new(records: AppendRecordBatch) -> Self {
3270 Self {
3271 records,
3272 match_seq_num: None,
3273 fencing_token: None,
3274 }
3275 }
3276
3277 pub fn with_match_seq_num(self, match_seq_num: u64) -> Self {
3279 Self {
3280 match_seq_num: Some(match_seq_num),
3281 ..self
3282 }
3283 }
3284
3285 pub fn with_fencing_token(self, fencing_token: FencingToken) -> Self {
3287 Self {
3288 fencing_token: Some(fencing_token),
3289 ..self
3290 }
3291 }
3292}
3293
3294impl From<AppendInput> for api::stream::proto::AppendInput {
3295 fn from(value: AppendInput) -> Self {
3296 Self {
3297 records: value.records.iter().cloned().map(Into::into).collect(),
3298 match_seq_num: value.match_seq_num,
3299 fencing_token: value.fencing_token.map(|t| t.to_string()),
3300 }
3301 }
3302}
3303
3304#[derive(Debug, Clone, PartialEq)]
3305#[non_exhaustive]
3306pub struct AppendAck {
3308 pub start: StreamPosition,
3310 pub end: StreamPosition,
3316 pub tail: StreamPosition,
3321}
3322
3323impl AppendAck {
3324 pub fn new(start: StreamPosition, end: StreamPosition, tail: StreamPosition) -> Self {
3328 Self { start, end, tail }
3329 }
3330}
3331
3332impl From<api::stream::proto::AppendAck> for AppendAck {
3333 fn from(value: api::stream::proto::AppendAck) -> Self {
3334 Self {
3335 start: value.start.unwrap_or_default().into(),
3336 end: value.end.unwrap_or_default().into(),
3337 tail: value.tail.unwrap_or_default().into(),
3338 }
3339 }
3340}
3341
3342#[derive(Debug, Clone, Copy)]
3343pub enum ReadFrom {
3345 SeqNum(u64),
3347 Timestamp(u64),
3349 TailOffset(u64),
3351}
3352
3353impl Default for ReadFrom {
3354 fn default() -> Self {
3355 Self::SeqNum(0)
3356 }
3357}
3358
3359#[derive(Debug, Default, Clone)]
3360#[non_exhaustive]
3361pub struct ReadStart {
3363 pub from: ReadFrom,
3367 pub clamp_to_tail: bool,
3371}
3372
3373impl ReadStart {
3374 pub fn new() -> Self {
3376 Self::default()
3377 }
3378
3379 pub fn with_from(self, from: ReadFrom) -> Self {
3381 Self { from, ..self }
3382 }
3383
3384 pub fn with_clamp_to_tail(self, clamp_to_tail: bool) -> Self {
3386 Self {
3387 clamp_to_tail,
3388 ..self
3389 }
3390 }
3391}
3392
3393impl From<ReadStart> for api::stream::ReadStart {
3394 fn from(value: ReadStart) -> Self {
3395 let (seq_num, timestamp, tail_offset) = match value.from {
3396 ReadFrom::SeqNum(n) => (Some(n), None, None),
3397 ReadFrom::Timestamp(t) => (None, Some(t), None),
3398 ReadFrom::TailOffset(o) => (None, None, Some(o)),
3399 };
3400 Self {
3401 seq_num,
3402 timestamp,
3403 tail_offset,
3404 clamp: if value.clamp_to_tail {
3405 Some(true)
3406 } else {
3407 None
3408 },
3409 }
3410 }
3411}
3412
3413#[derive(Debug, Clone, Default)]
3414#[non_exhaustive]
3415pub struct ReadLimits {
3417 pub count: Option<usize>,
3421 pub bytes: Option<usize>,
3425}
3426
3427impl ReadLimits {
3428 pub fn new() -> Self {
3430 Self::default()
3431 }
3432
3433 pub fn with_count(self, count: usize) -> Self {
3435 Self {
3436 count: Some(count),
3437 ..self
3438 }
3439 }
3440
3441 pub fn with_bytes(self, bytes: usize) -> Self {
3443 Self {
3444 bytes: Some(bytes),
3445 ..self
3446 }
3447 }
3448}
3449
3450#[derive(Debug, Clone, Default)]
3451#[non_exhaustive]
3452pub struct ReadStop {
3454 pub limits: ReadLimits,
3458 pub until: Option<RangeTo<u64>>,
3462 pub wait: Option<u32>,
3472}
3473
3474impl ReadStop {
3475 pub fn new() -> Self {
3477 Self::default()
3478 }
3479
3480 pub fn with_limits(self, limits: ReadLimits) -> Self {
3482 Self { limits, ..self }
3483 }
3484
3485 pub fn with_until(self, until: RangeTo<u64>) -> Self {
3487 Self {
3488 until: Some(until),
3489 ..self
3490 }
3491 }
3492
3493 pub fn with_wait(self, wait: u32) -> Self {
3495 Self {
3496 wait: Some(wait),
3497 ..self
3498 }
3499 }
3500}
3501
3502impl From<ReadStop> for api::stream::ReadEnd {
3503 fn from(value: ReadStop) -> Self {
3504 Self {
3505 count: value.limits.count,
3506 bytes: value.limits.bytes,
3507 until: value.until.map(|r| r.end),
3508 wait: value.wait,
3509 }
3510 }
3511}
3512
3513#[derive(Debug, Clone, Default)]
3514#[non_exhaustive]
3515pub struct ReadInput {
3518 pub start: ReadStart,
3522 pub stop: ReadStop,
3526 pub ignore_command_records: bool,
3530}
3531
3532#[derive(Debug, Clone, Copy, Default, Eq, PartialEq)]
3533#[non_exhaustive]
3534pub enum ReadSessionRetryPolicy {
3536 #[default]
3538 Budgeted,
3539 Indefinite,
3545}
3546
3547#[derive(Debug, Clone, Default)]
3548#[non_exhaustive]
3549pub struct ReadSessionConfig {
3551 pub retry_policy: ReadSessionRetryPolicy,
3557}
3558
3559impl ReadSessionConfig {
3560 pub fn new() -> Self {
3562 Self::default()
3563 }
3564
3565 pub fn with_retry_policy(self, retry_policy: ReadSessionRetryPolicy) -> Self {
3567 Self {
3568 retry_policy,
3569 ..self
3570 }
3571 }
3572}
3573
3574impl ReadInput {
3575 pub fn new() -> Self {
3577 Self::default()
3578 }
3579
3580 pub fn with_start(self, start: ReadStart) -> Self {
3582 Self { start, ..self }
3583 }
3584
3585 pub fn with_stop(self, stop: ReadStop) -> Self {
3587 Self { stop, ..self }
3588 }
3589
3590 pub fn with_ignore_command_records(self, ignore_command_records: bool) -> Self {
3592 Self {
3593 ignore_command_records,
3594 ..self
3595 }
3596 }
3597}
3598
3599#[derive(Debug, Clone)]
3600#[non_exhaustive]
3601pub struct SequencedRecord {
3603 pub seq_num: u64,
3605 pub body: Bytes,
3607 pub headers: Vec<Header>,
3609 pub timestamp: u64,
3611}
3612
3613impl SequencedRecord {
3614 pub fn from_parts(
3618 seq_num: u64,
3619 timestamp: u64,
3620 headers: Vec<Header>,
3621 body: impl Into<Bytes>,
3622 ) -> Self {
3623 Self {
3624 seq_num,
3625 timestamp,
3626 body: body.into(),
3627 headers,
3628 }
3629 }
3630
3631 pub fn is_command_record(&self) -> bool {
3633 self.headers.len() == 1 && *self.headers[0].name == *b""
3634 }
3635}
3636
3637impl From<api::stream::proto::SequencedRecord> for SequencedRecord {
3638 fn from(value: api::stream::proto::SequencedRecord) -> Self {
3639 Self {
3640 seq_num: value.seq_num,
3641 body: value.body,
3642 headers: value.headers.into_iter().map(Into::into).collect(),
3643 timestamp: value.timestamp,
3644 }
3645 }
3646}
3647
3648metered_bytes_impl!(SequencedRecord);
3649
3650#[derive(Debug, Clone)]
3651#[non_exhaustive]
3652pub struct ReadBatch {
3655 pub records: Vec<SequencedRecord>,
3662 pub tail: Option<StreamPosition>,
3667}
3668
3669impl ReadBatch {
3670 pub fn new(records: Vec<SequencedRecord>, tail: Option<StreamPosition>) -> Self {
3674 Self { records, tail }
3675 }
3676
3677 pub(crate) fn from_api(batch: api::stream::proto::ReadBatch) -> Self {
3678 Self {
3679 records: batch.records.into_iter().map(Into::into).collect(),
3680 tail: batch.tail.map(Into::into),
3681 }
3682 }
3683}
3684
3685pub type Streaming<T> = Pin<Box<dyn Send + futures_core::Stream<Item = Result<T, RequestError>>>>;
3687
3688fn idempotency_token() -> String {
3689 uuid::Uuid::new_v4().simple().to_string()
3690}
3691
3692#[cfg(test)]
3693mod tests {
3694 use proptest::prelude::*;
3695 use rstest::rstest;
3696
3697 use super::*;
3698
3699 type HeaderParts = (Vec<u8>, Vec<u8>);
3700 type AppendRecordParts = (Vec<u8>, Vec<HeaderParts>);
3701
3702 fn byte_vec_strategy(max_len: usize) -> impl Strategy<Value = Vec<u8>> {
3703 prop::collection::vec(any::<u8>(), 0..=max_len)
3704 }
3705
3706 fn header_parts_strategy() -> impl Strategy<Value = HeaderParts> {
3707 (byte_vec_strategy(32), byte_vec_strategy(64))
3708 }
3709
3710 fn string_strategy(max_chars: usize) -> impl Strategy<Value = String> {
3711 prop::collection::vec(any::<char>(), 0..=max_chars)
3712 .prop_map(|chars| chars.into_iter().collect())
3713 }
3714
3715 fn read_from_strategy() -> impl Strategy<Value = ReadFrom> {
3716 prop_oneof![
3717 any::<u64>().prop_map(ReadFrom::SeqNum),
3718 any::<u64>().prop_map(ReadFrom::Timestamp),
3719 any::<u64>().prop_map(ReadFrom::TailOffset),
3720 ]
3721 }
3722
3723 fn append_record_parts_strategy() -> impl Strategy<Value = AppendRecordParts> {
3724 (
3725 byte_vec_strategy(256),
3726 prop::collection::vec(header_parts_strategy(), 0..=16),
3727 )
3728 }
3729
3730 fn proto_stream_position_strategy() -> impl Strategy<Value = api::stream::proto::StreamPosition>
3731 {
3732 (any::<u64>(), any::<u64>()).prop_map(|(seq_num, timestamp)| {
3733 api::stream::proto::StreamPosition { seq_num, timestamp }
3734 })
3735 }
3736
3737 fn headers_from_parts(headers: &[HeaderParts]) -> Vec<Header> {
3738 headers
3739 .iter()
3740 .map(|(name, value)| Header::new(Bytes::from(name.clone()), Bytes::from(value.clone())))
3741 .collect()
3742 }
3743
3744 fn expected_metered_bytes(body: &[u8], headers: &[HeaderParts]) -> usize {
3745 8 + (2 * headers.len())
3746 + headers
3747 .iter()
3748 .map(|(name, value)| name.len() + value.len())
3749 .sum::<usize>()
3750 + body.len()
3751 }
3752
3753 #[test]
3756 fn s2_datetime_parse_valid_rfc3339() {
3757 let dt: S2DateTime = "2024-01-15T12:30:00Z".parse().unwrap();
3758 assert_eq!(dt.to_string(), "2024-01-15T12:30:00Z");
3759 }
3760
3761 #[test]
3762 fn s2_datetime_parse_with_offset() {
3763 let dt: S2DateTime = "2024-06-01T08:00:00+05:30".parse().unwrap();
3764 assert_eq!(dt.to_string(), "2024-06-01T08:00:00+05:30");
3765
3766 let offset_dt: time::OffsetDateTime = dt.into();
3767 assert_eq!(
3768 offset_dt.offset(),
3769 time::UtcOffset::from_hms(5, 30, 0).unwrap()
3770 );
3771 }
3772
3773 #[test]
3774 fn s2_datetime_parse_invalid() {
3775 let err = "not-a-date".parse::<S2DateTime>();
3776 assert!(err.is_err());
3777 }
3778
3779 #[test]
3780 fn s2_datetime_roundtrip_via_offset_datetime() {
3781 let odt = time::OffsetDateTime::from_unix_timestamp(1_700_000_000).unwrap();
3782 let dt = S2DateTime::try_from(odt).unwrap();
3783 let back: time::OffsetDateTime = dt.into();
3784 assert_eq!(odt, back);
3785 }
3786
3787 #[rstest]
3790 #[case::https_with_scheme("https://aws.s2.dev", Scheme::HTTPS)]
3791 #[case::http_with_scheme("http://localhost:8080", Scheme::HTTP)]
3792 #[case::default_https("aws.s2.dev", Scheme::HTTPS)]
3793 fn account_endpoint_parse(#[case] input: &str, #[case] expected_scheme: Scheme) {
3794 let ep: AccountEndpoint = input.parse().unwrap();
3795 assert_eq!(ep.scheme, expected_scheme);
3796 }
3797
3798 #[rstest]
3801 #[case::https_parent_zone("https://{basin}.b.s2.dev", Scheme::HTTPS, true)]
3802 #[case::http_direct("http://localhost:8080", Scheme::HTTP, false)]
3803 #[case::default_https_parent_zone("{basin}.b.s2.dev", Scheme::HTTPS, true)]
3804 fn basin_endpoint_parse(
3805 #[case] input: &str,
3806 #[case] expected_scheme: Scheme,
3807 #[case] expected_parent_zone: bool,
3808 ) {
3809 let ep: BasinEndpoint = input.parse().unwrap();
3810 assert_eq!(ep.scheme, expected_scheme);
3811 assert_eq!(
3812 matches!(ep.authority, BasinAuthority::ParentZone(_)),
3813 expected_parent_zone
3814 );
3815 }
3816
3817 #[test]
3820 fn s2_endpoints_new_requires_same_scheme() {
3821 let account: AccountEndpoint = "https://aws.s2.dev".parse().unwrap();
3822 let basin: BasinEndpoint = "http://localhost:8080".parse().unwrap();
3823 let err = S2Endpoints::new(account, basin);
3824 assert!(err.is_err());
3825 }
3826
3827 #[test]
3828 fn s2_endpoints_new_same_scheme_succeeds() {
3829 let account: AccountEndpoint = "https://aws.s2.dev".parse().unwrap();
3830 let basin: BasinEndpoint = "https://{basin}.b.s2.dev".parse().unwrap();
3831 let ep = S2Endpoints::new(account, basin).unwrap();
3832 assert_eq!(ep.scheme, Scheme::HTTPS);
3833 }
3834
3835 #[test]
3836 fn s2_endpoints_for_endpoint_defaults_to_https() {
3837 let ep = S2Endpoints::for_endpoint("localhost:8080").unwrap();
3838 let authority: Authority = "localhost:8080".parse().unwrap();
3839 assert_eq!(ep.scheme, Scheme::HTTPS);
3840 assert_eq!(ep.account_authority, authority);
3841 assert_eq!(ep.basin_authority, BasinAuthority::Direct(authority));
3842 }
3843
3844 #[test]
3845 fn s2_endpoints_for_endpoint_accepts_explicit_scheme() {
3846 let ep = S2Endpoints::for_endpoint("http://localhost:8080").unwrap();
3847 assert_eq!(ep.scheme, Scheme::HTTP);
3848 }
3849
3850 #[test]
3851 fn s2_endpoints_for_endpoint_rejects_invalid_endpoint() {
3852 assert!(S2Endpoints::for_endpoint("not a valid endpoint").is_err());
3853 }
3854
3855 #[rstest]
3858 #[case::none(Compression::None, CompressionAlgorithm::None)]
3859 #[case::gzip(Compression::Gzip, CompressionAlgorithm::Gzip)]
3860 #[case::zstd(Compression::Zstd, CompressionAlgorithm::Zstd)]
3861 fn compression_conversion(#[case] sdk: Compression, #[case] api: CompressionAlgorithm) {
3862 assert_eq!(CompressionAlgorithm::from(sdk), api);
3863 }
3864
3865 #[test]
3868 fn retry_config_defaults() {
3869 let rc = RetryConfig::default();
3870 assert_eq!(rc.max_attempts.get(), 3);
3871 assert_eq!(rc.min_base_delay, Duration::from_millis(100));
3872 assert_eq!(rc.max_base_delay, Duration::from_secs(1));
3873 assert!(matches!(rc.append_retry_policy, AppendRetryPolicy::All));
3874 }
3875
3876 #[test]
3877 fn retry_config_max_retries() {
3878 let rc = RetryConfig::default();
3879 assert_eq!(rc.max_retries(), 2);
3880 }
3881
3882 #[test]
3885 fn s2_config_defaults() {
3886 let cfg = S2Config::new("test-token");
3887 assert_eq!(cfg.connection_timeout, Duration::from_secs(3));
3888 assert_eq!(cfg.request_timeout, Duration::from_secs(5));
3889 assert!(!cfg.insecure_skip_cert_verification);
3890 }
3891
3892 #[rstest]
3895 #[case::standard(StorageClass::Standard)]
3896 #[case::express(StorageClass::Express)]
3897 fn storage_class_roundtrip(#[case] sdk: StorageClass) {
3898 let api: api::config::StorageClass = sdk.into();
3899 let back: StorageClass = api.into();
3900 assert_eq!(back, sdk);
3901 }
3902
3903 #[rstest]
3906 #[case::age(RetentionPolicy::Age(3600))]
3907 #[case::infinite(RetentionPolicy::Infinite)]
3908 fn retention_policy_roundtrip(#[case] sdk: RetentionPolicy) {
3909 let api: api::config::RetentionPolicy = sdk.into();
3910 let back: RetentionPolicy = api.into();
3911 assert_eq!(back, sdk);
3912 }
3913
3914 #[rstest]
3917 #[case::client_prefer(
3918 TimestampingMode::ClientPrefer,
3919 api::config::TimestampingMode::ClientPrefer
3920 )]
3921 #[case::client_require(
3922 TimestampingMode::ClientRequire,
3923 api::config::TimestampingMode::ClientRequire
3924 )]
3925 #[case::arrival(TimestampingMode::Arrival, api::config::TimestampingMode::Arrival)]
3926 fn timestamping_mode_roundtrip(
3927 #[case] sdk: TimestampingMode,
3928 #[case] expected_api: api::config::TimestampingMode,
3929 ) {
3930 let converted: api::config::TimestampingMode = sdk.into();
3931 assert_eq!(converted, expected_api);
3932 let back: TimestampingMode = converted.into();
3933 assert_eq!(back, sdk);
3934 }
3935
3936 #[test]
3939 fn timestamping_config_roundtrip() {
3940 let sdk = TimestampingConfig {
3941 mode: Some(TimestampingMode::Arrival),
3942 uncapped: Some(true),
3943 };
3944 let api: api::config::TimestampingConfig = sdk.into();
3945 let back: TimestampingConfig = api.into();
3946 assert_eq!(back, sdk);
3947 }
3948
3949 #[test]
3952 fn delete_on_empty_config_roundtrip() {
3953 let sdk = DeleteOnEmptyConfig::new().with_min_age(Duration::from_secs(300));
3954 let api: api::config::DeleteOnEmptyConfig = sdk.into();
3955 let back: DeleteOnEmptyConfig = api.into();
3956 assert_eq!(back, sdk);
3957 }
3958
3959 #[test]
3962 fn stream_config_builder_and_roundtrip() {
3963 let sdk = StreamConfig::new()
3964 .with_storage_class(StorageClass::Express)
3965 .with_retention_policy(RetentionPolicy::Age(86400))
3966 .with_timestamping(TimestampingConfig {
3967 mode: Some(TimestampingMode::ClientPrefer),
3968 uncapped: None,
3969 })
3970 .with_delete_on_empty(DeleteOnEmptyConfig { min_age_secs: 60 });
3971 let api: api::config::StreamConfig = sdk.clone().into();
3972 let back: StreamConfig = api.into();
3973 assert_eq!(back, sdk);
3974 }
3975
3976 #[test]
3979 fn basin_config_builder_and_roundtrip() {
3980 let sdk = BasinConfig::new()
3981 .with_default_stream_config(
3982 StreamConfig::new().with_storage_class(StorageClass::Standard),
3983 )
3984 .with_create_stream_on_append(true)
3985 .with_create_stream_on_read(false);
3986 let api: api::config::BasinConfig = sdk.clone().into();
3987 let back: BasinConfig = api.into();
3988 assert_eq!(back, sdk);
3989 }
3990
3991 proptest! {
3994 #[test]
3995 fn fencing_token_parse_accepts_only_within_byte_limit(
3996 token in string_strategy(MAX_FENCING_TOKEN_LENGTH + 8),
3997 ) {
3998 let parsed = token.parse::<FencingToken>();
3999
4000 if token.len() <= MAX_FENCING_TOKEN_LENGTH {
4001 prop_assert_eq!(parsed.unwrap().to_string(), token);
4002 } else {
4003 prop_assert!(parsed.is_err());
4004 }
4005 }
4006 }
4007
4008 #[test]
4011 fn stream_position_display() {
4012 let pos = StreamPosition {
4013 seq_num: 42,
4014 timestamp: 1700000000,
4015 };
4016 assert_eq!(pos.to_string(), "seq_num=42, timestamp=1700000000");
4017 }
4018
4019 proptest! {
4020 #[test]
4021 fn stream_position_conversions_preserve_values(seq_num in any::<u64>(), timestamp in any::<u64>()) {
4022 let proto: StreamPosition = api::stream::proto::StreamPosition {
4023 seq_num,
4024 timestamp,
4025 }
4026 .into();
4027 prop_assert_eq!(proto.seq_num, seq_num);
4028 prop_assert_eq!(proto.timestamp, timestamp);
4029
4030 let api: StreamPosition = api::stream::StreamPosition {
4031 seq_num,
4032 timestamp,
4033 }
4034 .into();
4035 prop_assert_eq!(api.seq_num, seq_num);
4036 prop_assert_eq!(api.timestamp, timestamp);
4037 }
4038 }
4039
4040 proptest! {
4043 #[test]
4044 fn header_proto_roundtrip_preserves_binary_parts(
4045 name in byte_vec_strategy(64),
4046 value in byte_vec_strategy(128),
4047 ) {
4048 let header = Header::new(Bytes::from(name.clone()), Bytes::from(value.clone()));
4049 let proto: api::stream::proto::Header = header.into();
4050 let back: Header = proto.into();
4051
4052 prop_assert_eq!(back.name.as_ref(), name.as_slice());
4053 prop_assert_eq!(back.value.as_ref(), value.as_slice());
4054 }
4055 }
4056
4057 #[test]
4060 fn append_record_too_large() {
4061 let big_body = vec![0u8; RECORD_BATCH_MAX.bytes + 1];
4062 assert!(AppendRecord::new(big_body).is_err());
4063 }
4064
4065 proptest! {
4068 #[test]
4069 fn append_record_preserves_fields_and_metered_byte_formula(
4070 (body, headers) in append_record_parts_strategy(),
4071 timestamp in proptest::option::of(any::<u64>()),
4072 ) {
4073 let mut record = AppendRecord::new(body.clone())
4074 .unwrap()
4075 .with_headers(headers_from_parts(&headers))
4076 .unwrap();
4077 if let Some(timestamp) = timestamp {
4078 record = record.with_timestamp(timestamp);
4079 }
4080
4081 prop_assert_eq!(record.body(), body.as_slice());
4082 prop_assert_eq!(record.headers().len(), headers.len());
4083 prop_assert_eq!(record.timestamp(), timestamp);
4084 prop_assert_eq!(record.metered_bytes(), expected_metered_bytes(&body, &headers));
4085
4086 for (actual, (expected_name, expected_value)) in record.headers().iter().zip(headers.iter()) {
4087 prop_assert_eq!(actual.name.as_ref(), expected_name.as_slice());
4088 prop_assert_eq!(actual.value.as_ref(), expected_value.as_slice());
4089 }
4090 }
4091 }
4092
4093 #[test]
4096 fn append_record_batch_empty_is_err() {
4097 let result = AppendRecordBatch::try_from_iter(vec![]);
4098 assert!(result.is_err());
4099 }
4100
4101 #[test]
4102 fn append_record_batch_too_many_records() {
4103 let records: Vec<_> = (0..1001).map(|_| AppendRecord::new("x").unwrap()).collect();
4104 let result = AppendRecordBatch::try_from_iter(records);
4105 assert!(result.is_err());
4106 }
4107
4108 proptest! {
4109 #[test]
4110 fn append_record_batch_metered_bytes_is_sum_of_records(
4111 records in prop::collection::vec(append_record_parts_strategy(), 1..=32),
4112 ) {
4113 let expected = records
4114 .iter()
4115 .map(|(body, headers)| expected_metered_bytes(body, headers))
4116 .sum::<usize>();
4117 let records = records
4118 .into_iter()
4119 .map(|(body, headers)| {
4120 AppendRecord::new(body)
4121 .unwrap()
4122 .with_headers(headers_from_parts(&headers))
4123 .unwrap()
4124 })
4125 .collect::<Vec<_>>();
4126
4127 let batch = AppendRecordBatch::try_from_iter(records).unwrap();
4128 prop_assert_eq!(batch.metered_bytes(), expected);
4129 prop_assert_eq!(batch.iter().map(MeteredBytes::metered_bytes).sum::<usize>(), expected);
4130 }
4131 }
4132
4133 #[test]
4136 fn command_record_fence() {
4137 let token: FencingToken = "tok".parse().unwrap();
4138 let cmd = CommandRecord::fence(token);
4139 let record: AppendRecord = cmd.into();
4140 assert_eq!(record.headers().len(), 1);
4141 assert_eq!(record.headers()[0].name.as_ref(), b"");
4142 assert_eq!(record.headers()[0].value.as_ref(), b"fence");
4143 assert_eq!(record.body(), b"tok");
4144 }
4145
4146 #[test]
4147 fn command_record_trim() {
4148 let cmd = CommandRecord::trim(42);
4149 let record: AppendRecord = cmd.into();
4150 assert_eq!(record.headers().len(), 1);
4151 assert_eq!(record.headers()[0].value.as_ref(), b"trim");
4152 assert_eq!(record.body(), &42u64.to_be_bytes());
4153 }
4154
4155 #[rstest]
4158 #[case::command(vec![Header::new("", "fence")], true)]
4159 #[case::regular(vec![Header::new("key", "value")], false)]
4160 #[case::no_headers(vec![], false)]
4161 fn sequenced_record_command_detection(#[case] headers: Vec<Header>, #[case] expected: bool) {
4162 let record = SequencedRecord {
4163 seq_num: 0,
4164 body: Bytes::from("data"),
4165 headers,
4166 timestamp: 0,
4167 };
4168 assert_eq!(record.is_command_record(), expected);
4169 }
4170
4171 proptest! {
4174 #[test]
4175 fn read_start_to_api_sets_only_selected_position_field(
4176 from in read_from_strategy(),
4177 clamp_to_tail in any::<bool>(),
4178 ) {
4179 let (seq_num, timestamp, tail_offset) = match from {
4180 ReadFrom::SeqNum(value) => (Some(value), None, None),
4181 ReadFrom::Timestamp(value) => (None, Some(value), None),
4182 ReadFrom::TailOffset(value) => (None, None, Some(value)),
4183 };
4184 let api: api::stream::ReadStart = ReadStart::new()
4185 .with_from(from)
4186 .with_clamp_to_tail(clamp_to_tail)
4187 .into();
4188
4189 prop_assert_eq!(api.seq_num, seq_num);
4190 prop_assert_eq!(api.timestamp, timestamp);
4191 prop_assert_eq!(api.tail_offset, tail_offset);
4192 prop_assert_eq!(api.clamp, clamp_to_tail.then_some(true));
4193 }
4194 }
4195
4196 #[test]
4199 fn read_stop_to_api() {
4200 let stop = ReadStop::new()
4201 .with_limits(ReadLimits::new().with_count(50))
4202 .with_until(..1000)
4203 .with_wait(30);
4204 let api: api::stream::ReadEnd = stop.into();
4205 assert_eq!(api.count, Some(50));
4206 assert_eq!(api.until, Some(1000));
4207 assert_eq!(api.wait, Some(30));
4208 }
4209
4210 #[test]
4213 fn operation_roundtrip_all_variants() {
4214 let variants = [
4215 Operation::ListBasins,
4216 Operation::CreateBasin,
4217 Operation::GetBasinConfig,
4218 Operation::DeleteBasin,
4219 Operation::ReconfigureBasin,
4220 Operation::ListAccessTokens,
4221 Operation::IssueAccessToken,
4222 Operation::RevokeAccessToken,
4223 Operation::GetAccountMetrics,
4224 Operation::GetBasinMetrics,
4225 Operation::GetStreamMetrics,
4226 Operation::ListStreams,
4227 Operation::CreateStream,
4228 Operation::GetStreamConfig,
4229 Operation::DeleteStream,
4230 Operation::ReconfigureStream,
4231 Operation::CheckTail,
4232 Operation::Append,
4233 Operation::Read,
4234 Operation::Trim,
4235 Operation::Fence,
4236 Operation::ListLocations,
4237 Operation::GetDefaultLocation,
4238 Operation::SetDefaultLocation,
4239 ];
4240 for op in variants {
4241 let api_op: api::access::Operation = op.into();
4242 let back: Operation = api_op.into();
4243 assert_eq!(back, op);
4244 }
4245 }
4246
4247 #[test]
4250 fn metric_unit_conversion() {
4251 assert_eq!(
4252 MetricUnit::from(api::metrics::MetricUnit::Bytes),
4253 MetricUnit::Bytes
4254 );
4255 assert_eq!(
4256 MetricUnit::from(api::metrics::MetricUnit::Operations),
4257 MetricUnit::Operations
4258 );
4259 }
4260
4261 proptest! {
4264 #[test]
4265 fn append_ack_from_proto_preserves_present_positions_and_defaults_missing(
4266 start in proptest::option::of(proto_stream_position_strategy()),
4267 end in proptest::option::of(proto_stream_position_strategy()),
4268 tail in proptest::option::of(proto_stream_position_strategy()),
4269 ) {
4270 let expected_start = start.unwrap_or_default();
4271 let expected_end = end.unwrap_or_default();
4272 let expected_tail = tail.unwrap_or_default();
4273 let ack: AppendAck = api::stream::proto::AppendAck { start, end, tail }.into();
4274
4275 prop_assert_eq!(ack.start.seq_num, expected_start.seq_num);
4276 prop_assert_eq!(ack.start.timestamp, expected_start.timestamp);
4277 prop_assert_eq!(ack.end.seq_num, expected_end.seq_num);
4278 prop_assert_eq!(ack.end.timestamp, expected_end.timestamp);
4279 prop_assert_eq!(ack.tail.seq_num, expected_tail.seq_num);
4280 prop_assert_eq!(ack.tail.timestamp, expected_tail.timestamp);
4281 }
4282 }
4283
4284 #[test]
4287 fn read_batch_from_api() {
4288 let proto_batch = api::stream::proto::ReadBatch {
4289 records: vec![api::stream::proto::SequencedRecord {
4290 seq_num: 0,
4291 body: Bytes::from("hi"),
4292 headers: vec![api::stream::proto::Header {
4293 name: Bytes::from("k"),
4294 value: Bytes::from("v"),
4295 }],
4296 timestamp: 42,
4297 }],
4298 tail: Some(api::stream::proto::StreamPosition {
4299 seq_num: 1,
4300 timestamp: 42,
4301 }),
4302 };
4303 let batch = ReadBatch::from_api(proto_batch);
4304 assert_eq!(batch.records.len(), 1);
4305 assert_eq!(batch.records[0].seq_num, 0);
4306 assert_eq!(batch.records[0].timestamp, 42);
4307 assert_eq!(batch.records[0].body.as_ref(), b"hi");
4308 assert_eq!(batch.records[0].headers.len(), 1);
4309 assert_eq!(batch.records[0].headers[0].name.as_ref(), b"k");
4310 assert_eq!(batch.records[0].headers[0].value.as_ref(), b"v");
4311 assert_eq!(
4312 batch.tail,
4313 Some(StreamPosition {
4314 seq_num: 1,
4315 timestamp: 42,
4316 })
4317 );
4318 }
4319
4320 #[test]
4323 fn create_basin_input_to_api() {
4324 let name: BasinName = "test-basin-name".parse().unwrap();
4325 let input = CreateBasinInput::new(name.clone()).with_config(BasinConfig::new());
4326 let (req, token): (api::basin::CreateBasinRequest, String) = input.into();
4327 assert_eq!(req.basin, name);
4328 assert!(req.config.is_some());
4329 assert!(!token.is_empty());
4330 }
4331
4332 #[test]
4335 fn create_stream_input_to_api() {
4336 let name: StreamName = "my-stream".parse().unwrap();
4337 let input = CreateStreamInput::new(name.clone()).with_config(StreamConfig::new());
4338 let (req, token): (api::stream::CreateStreamRequest, String) = input.into();
4339 assert_eq!(req.stream, name);
4340 assert!(req.config.is_some());
4341 assert!(!token.is_empty());
4342 }
4343
4344 #[test]
4347 fn sequenced_record_from_proto() {
4348 let proto = api::stream::proto::SequencedRecord {
4349 seq_num: 99,
4350 body: Bytes::from("data"),
4351 headers: vec![api::stream::proto::Header {
4352 name: Bytes::from("k"),
4353 value: Bytes::from("v"),
4354 }],
4355 timestamp: 1234,
4356 };
4357 let record: SequencedRecord = proto.into();
4358 assert_eq!(record.seq_num, 99);
4359 assert_eq!(record.body.as_ref(), b"data");
4360 assert_eq!(record.headers.len(), 1);
4361 assert_eq!(record.headers[0].name.as_ref(), b"k");
4362 assert_eq!(record.headers[0].value.as_ref(), b"v");
4363 assert_eq!(record.timestamp, 1234);
4364 }
4365}