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
15#[cfg(feature = "_hidden")]
16use async_trait::async_trait;
17use bytes::Bytes;
18use compact_str::CompactString;
19use http::{
20 HeaderMap,
21 header::HeaderValue,
22 uri::{Authority, Scheme},
23};
24use rand::RngExt;
25use s2_api::{v1 as api, v1::stream::s2s::CompressionAlgorithm};
26pub use s2_common::ValidationError;
28pub use s2_common::access::AccessTokenId;
33pub use s2_common::access::AccessTokenIdPrefix;
35pub use s2_common::access::AccessTokenIdStartAfter;
37pub use s2_common::basin::BasinName;
42pub use s2_common::basin::BasinNamePrefix;
44pub use s2_common::basin::BasinNameStartAfter;
46pub use s2_common::location::LocationName;
51pub use s2_common::stream::StreamName;
56pub use s2_common::stream::StreamNamePrefix;
58pub use s2_common::stream::StreamNameStartAfter;
60pub use s2_common::{
61 caps::RECORD_BATCH_MAX,
62 encryption::{EncryptionAlgorithm, EncryptionKey},
63};
64
65pub(crate) const ONE_MIB: u32 = 1024 * 1024;
66
67use s2_common::{
68 maybe::Maybe,
69 record::{MAX_FENCING_TOKEN_LENGTH, Metered, MeteredSize},
70 resources::ProvisionResult,
71};
72use secrecy::SecretString;
73
74use crate::error::RequestError;
75
76#[cfg(feature = "_hidden")]
77#[derive(Debug, Clone, thiserror::Error)]
78#[error("{message}")]
79#[doc(hidden)]
80pub struct AccessTokenProviderError {
81 message: String,
82 retryable: bool,
83}
84
85#[cfg(feature = "_hidden")]
86impl AccessTokenProviderError {
87 pub fn transient(message: impl Into<String>) -> Self {
89 Self {
90 message: message.into(),
91 retryable: true,
92 }
93 }
94
95 pub fn permanent(message: impl Into<String>) -> Self {
97 Self {
98 message: message.into(),
99 retryable: false,
100 }
101 }
102
103 pub(crate) fn is_retryable(&self) -> bool {
104 self.retryable
105 }
106}
107
108#[cfg(feature = "_hidden")]
109#[async_trait]
110#[doc(hidden)]
111pub trait AccessTokenProvider: fmt::Debug + Send + Sync {
112 async fn access_token(&self) -> Result<String, AccessTokenProviderError>;
114
115 fn invalidate_access_token(&self, _rejected_access_token: &str) {}
117}
118
119#[derive(Clone)]
120pub(crate) enum AccessToken {
121 Static(SecretString),
122 #[cfg(feature = "_hidden")]
123 Provider(Arc<dyn AccessTokenProvider>),
124}
125
126#[derive(Debug, Clone, Copy, PartialEq, Eq)]
127pub(crate) enum AccessTokenMode {
128 Static,
129 #[cfg(feature = "_hidden")]
130 Refreshable,
131}
132
133impl AccessTokenMode {
134 pub(crate) fn is_refreshable(self) -> bool {
135 match self {
136 Self::Static => false,
137 #[cfg(feature = "_hidden")]
138 Self::Refreshable => true,
139 }
140 }
141}
142
143impl AccessToken {
144 pub(crate) fn mode(&self) -> AccessTokenMode {
145 match self {
146 Self::Static(_) => AccessTokenMode::Static,
147 #[cfg(feature = "_hidden")]
148 Self::Provider(_) => AccessTokenMode::Refreshable,
149 }
150 }
151}
152
153impl fmt::Debug for AccessToken {
154 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
155 match self {
156 Self::Static(_) => formatter.write_str("Static(<redacted>)"),
157 #[cfg(feature = "_hidden")]
158 Self::Provider(_) => formatter.write_str("Provider(<redacted>)"),
159 }
160 }
161}
162
163#[derive(Debug, Clone, Copy, PartialEq, Eq)]
169pub struct S2DateTime(time::OffsetDateTime);
170
171impl TryFrom<time::OffsetDateTime> for S2DateTime {
172 type Error = ValidationError;
173
174 fn try_from(dt: time::OffsetDateTime) -> Result<Self, Self::Error> {
175 dt.format(&time::format_description::well_known::Rfc3339)
176 .map_err(|e| ValidationError(format!("not a valid RFC 3339 datetime: {e}")))?;
177 Ok(Self(dt))
178 }
179}
180
181impl From<S2DateTime> for time::OffsetDateTime {
182 fn from(dt: S2DateTime) -> Self {
183 dt.0
184 }
185}
186
187impl FromStr for S2DateTime {
188 type Err = ValidationError;
189
190 fn from_str(s: &str) -> Result<Self, Self::Err> {
191 time::OffsetDateTime::parse(s, &time::format_description::well_known::Rfc3339)
192 .map(Self)
193 .map_err(|e| ValidationError(format!("not a valid RFC 3339 datetime: {e}")))
194 }
195}
196
197impl fmt::Display for S2DateTime {
198 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
199 write!(
200 f,
201 "{}",
202 self.0
203 .format(&time::format_description::well_known::Rfc3339)
204 .expect("RFC3339 formatting should not fail for S2DateTime")
205 )
206 }
207}
208
209#[derive(Debug, Clone, PartialEq)]
211pub(crate) enum BasinAuthority {
212 ParentZone(Authority),
214 Direct(Authority),
216}
217
218#[derive(Debug, Clone)]
220pub struct AccountEndpoint {
221 scheme: Scheme,
222 authority: Authority,
223}
224
225impl AccountEndpoint {
226 pub fn new(endpoint: &str) -> Result<Self, ValidationError> {
228 endpoint.parse()
229 }
230}
231
232impl FromStr for AccountEndpoint {
233 type Err = ValidationError;
234
235 fn from_str(s: &str) -> Result<Self, Self::Err> {
236 let (scheme, authority) = match s.find("://") {
237 Some(idx) => {
238 let scheme: Scheme = s[..idx]
239 .parse()
240 .map_err(|_| "invalid account endpoint scheme".to_string())?;
241 (scheme, &s[idx + 3..])
242 }
243 None => (Scheme::HTTPS, s),
244 };
245 Ok(Self {
246 scheme,
247 authority: authority
248 .parse()
249 .map_err(|e| format!("invalid account endpoint authority: {e}"))?,
250 })
251 }
252}
253
254#[derive(Debug, Clone)]
256pub struct BasinEndpoint {
257 scheme: Scheme,
258 authority: BasinAuthority,
259}
260
261impl BasinEndpoint {
262 pub fn new(endpoint: &str) -> Result<Self, ValidationError> {
264 endpoint.parse()
265 }
266}
267
268impl FromStr for BasinEndpoint {
269 type Err = ValidationError;
270
271 fn from_str(s: &str) -> Result<Self, Self::Err> {
272 let (scheme, authority) = match s.find("://") {
273 Some(idx) => {
274 let scheme: Scheme = s[..idx]
275 .parse()
276 .map_err(|_| "invalid basin endpoint scheme".to_string())?;
277 (scheme, &s[idx + 3..])
278 }
279 None => (Scheme::HTTPS, s),
280 };
281 let authority = if let Some(authority) = authority.strip_prefix("{basin}.") {
282 BasinAuthority::ParentZone(
283 authority
284 .parse()
285 .map_err(|e| format!("invalid basin endpoint authority: {e}"))?,
286 )
287 } else {
288 BasinAuthority::Direct(
289 authority
290 .parse()
291 .map_err(|e| format!("invalid basin endpoint authority: {e}"))?,
292 )
293 };
294 Ok(Self { scheme, authority })
295 }
296}
297
298#[derive(Debug, Clone)]
299#[non_exhaustive]
300pub struct S2Endpoints {
302 pub(crate) scheme: Scheme,
303 pub(crate) account_authority: Authority,
304 pub(crate) basin_authority: BasinAuthority,
305}
306
307impl S2Endpoints {
308 pub fn new(
310 account_endpoint: AccountEndpoint,
311 basin_endpoint: BasinEndpoint,
312 ) -> Result<Self, ValidationError> {
313 if account_endpoint.scheme != basin_endpoint.scheme {
314 return Err("account and basin endpoints must have the same scheme".into());
315 }
316 Ok(Self {
317 scheme: account_endpoint.scheme,
318 account_authority: account_endpoint.authority,
319 basin_authority: basin_endpoint.authority,
320 })
321 }
322
323 pub fn for_endpoint(endpoint: &str) -> Result<Self, ValidationError> {
327 Self::new(
328 AccountEndpoint::new(endpoint)?,
329 BasinEndpoint::new(endpoint)?,
330 )
331 }
332
333 pub fn from_env() -> Result<Self, ValidationError> {
339 let account_endpoint: AccountEndpoint = match std::env::var("S2_ACCOUNT_ENDPOINT") {
340 Ok(endpoint) => endpoint.parse()?,
341 Err(VarError::NotPresent) => return Err("S2_ACCOUNT_ENDPOINT env var not set".into()),
342 Err(VarError::NotUnicode(_)) => {
343 return Err("S2_ACCOUNT_ENDPOINT is not valid unicode".into());
344 }
345 };
346
347 let basin_endpoint: BasinEndpoint = match std::env::var("S2_BASIN_ENDPOINT") {
348 Ok(endpoint) => endpoint.parse()?,
349 Err(VarError::NotPresent) => return Err("S2_BASIN_ENDPOINT env var not set".into()),
350 Err(VarError::NotUnicode(_)) => {
351 return Err("S2_BASIN_ENDPOINT is not valid unicode".into());
352 }
353 };
354
355 if account_endpoint.scheme != basin_endpoint.scheme {
356 return Err(
357 "S2_ACCOUNT_ENDPOINT and S2_BASIN_ENDPOINT must have the same scheme".into(),
358 );
359 }
360
361 Ok(Self {
362 scheme: account_endpoint.scheme,
363 account_authority: account_endpoint.authority,
364 basin_authority: basin_endpoint.authority,
365 })
366 }
367
368 pub fn for_cloud() -> Self {
370 Self {
371 scheme: Scheme::HTTPS,
372 account_authority: "a.s2.dev".try_into().expect("valid authority"),
373 basin_authority: BasinAuthority::ParentZone(
374 "b.s2.dev".try_into().expect("valid authority"),
375 ),
376 }
377 }
378}
379
380#[derive(Debug, Clone, Copy)]
381pub enum Compression {
383 None,
385 Gzip,
387 Zstd,
389}
390
391impl From<Compression> for CompressionAlgorithm {
392 fn from(value: Compression) -> Self {
393 match value {
394 Compression::None => CompressionAlgorithm::None,
395 Compression::Gzip => CompressionAlgorithm::Gzip,
396 Compression::Zstd => CompressionAlgorithm::Zstd,
397 }
398 }
399}
400
401#[derive(Debug, Clone, Copy, PartialEq)]
402#[non_exhaustive]
403pub enum AppendRetryPolicy {
406 All,
408 NoSideEffects,
418}
419
420#[derive(Debug, Clone)]
421#[non_exhaustive]
422pub struct RetryConfig {
431 pub max_attempts: NonZeroU32,
435 pub min_base_delay: Duration,
439 pub max_base_delay: Duration,
443 pub append_retry_policy: AppendRetryPolicy,
448}
449
450impl Default for RetryConfig {
451 fn default() -> Self {
452 Self {
453 max_attempts: NonZeroU32::new(3).expect("valid non-zero u32"),
454 min_base_delay: Duration::from_millis(100),
455 max_base_delay: Duration::from_secs(1),
456 append_retry_policy: AppendRetryPolicy::All,
457 }
458 }
459}
460
461impl RetryConfig {
462 pub fn new() -> Self {
464 Self::default()
465 }
466
467 pub(crate) fn max_retries(&self) -> u32 {
468 self.max_attempts.get() - 1
469 }
470
471 pub fn with_max_attempts(self, max_attempts: NonZeroU32) -> Self {
473 Self {
474 max_attempts,
475 ..self
476 }
477 }
478
479 pub fn with_min_base_delay(self, min_base_delay: Duration) -> Self {
481 Self {
482 min_base_delay,
483 ..self
484 }
485 }
486
487 pub fn with_max_base_delay(self, max_base_delay: Duration) -> Self {
489 Self {
490 max_base_delay,
491 ..self
492 }
493 }
494
495 pub fn with_append_retry_policy(self, append_retry_policy: AppendRetryPolicy) -> Self {
498 Self {
499 append_retry_policy,
500 ..self
501 }
502 }
503}
504
505#[derive(Debug, Clone)]
506#[non_exhaustive]
507pub struct S2Config {
509 pub(crate) access_token: AccessToken,
510 pub(crate) endpoints: S2Endpoints,
511 pub(crate) connection_timeout: Duration,
512 pub(crate) request_timeout: Duration,
513 pub(crate) retry: RetryConfig,
514 pub(crate) compression: Compression,
515 pub(crate) user_agent: HeaderValue,
516 pub(crate) default_headers: HeaderMap,
517 pub(crate) insecure_skip_cert_verification: bool,
518 pub(crate) rustls_crypto_provider: Option<Arc<rustls::crypto::CryptoProvider>>,
519}
520
521impl S2Config {
522 pub fn new(access_token: impl Into<String>) -> Self {
524 Self {
525 access_token: AccessToken::Static(access_token.into().into()),
526 endpoints: S2Endpoints::for_cloud(),
527 connection_timeout: Duration::from_secs(3),
528 request_timeout: Duration::from_secs(5),
529 retry: RetryConfig::new(),
530 compression: Compression::None,
531 user_agent: concat!("s2-sdk-rust/", env!("CARGO_PKG_VERSION"))
532 .parse()
533 .expect("valid user agent"),
534 default_headers: HeaderMap::new(),
535 insecure_skip_cert_verification: false,
536 rustls_crypto_provider: default_rustls_crypto_provider(),
537 }
538 }
539
540 #[cfg(feature = "_hidden")]
541 #[doc(hidden)]
542 pub fn with_access_token_provider(self, provider: impl AccessTokenProvider + 'static) -> Self {
543 Self {
544 access_token: AccessToken::Provider(Arc::new(provider)),
545 ..self
546 }
547 }
548
549 pub fn with_endpoints(self, endpoints: S2Endpoints) -> Self {
551 Self { endpoints, ..self }
552 }
553
554 #[cfg(feature = "_hidden")]
577 #[doc(hidden)]
578 pub fn with_default_headers(self, default_headers: HeaderMap) -> Result<Self, ValidationError> {
579 if default_headers.contains_key(http::header::CONTENT_ENCODING) {
580 return Err(ValidationError(
581 "Content-Encoding cannot be set in default headers; use S2Config::with_compression instead"
582 .into(),
583 ));
584 }
585 for name in [
586 http::header::CONTENT_TYPE,
587 http::header::CONTENT_LENGTH,
588 http::header::TRANSFER_ENCODING,
589 ] {
590 if default_headers.contains_key(&name) {
591 return Err(ValidationError(format!(
592 "{name} cannot be set in default headers; the SDK controls request format and body framing"
593 )));
594 }
595 }
596 Ok(Self {
597 default_headers,
598 ..self
599 })
600 }
601
602 pub fn with_connection_timeout(self, connection_timeout: Duration) -> Self {
606 Self {
607 connection_timeout,
608 ..self
609 }
610 }
611
612 pub fn with_request_timeout(self, request_timeout: Duration) -> Self {
616 Self {
617 request_timeout,
618 ..self
619 }
620 }
621
622 pub fn with_retry(self, retry: RetryConfig) -> Self {
626 Self { retry, ..self }
627 }
628
629 pub fn with_compression(self, compression: Compression) -> Self {
633 Self {
634 compression,
635 ..self
636 }
637 }
638
639 pub fn with_insecure_skip_cert_verification(self, skip: bool) -> Self {
651 Self {
652 insecure_skip_cert_verification: skip,
653 ..self
654 }
655 }
656
657 pub fn with_rustls_crypto_provider(
667 self,
668 provider: impl Into<Arc<rustls::crypto::CryptoProvider>>,
669 ) -> Self {
670 Self {
671 rustls_crypto_provider: Some(provider.into()),
672 ..self
673 }
674 }
675
676 #[cfg(feature = "rustls-aws-lc-rs")]
680 pub fn with_rustls_aws_lc_rs_crypto_provider(self) -> Self {
681 self.with_rustls_crypto_provider(rustls::crypto::aws_lc_rs::default_provider())
682 }
683
684 #[cfg(feature = "rustls-ring")]
688 pub fn with_rustls_ring_crypto_provider(self) -> Self {
689 self.with_rustls_crypto_provider(rustls::crypto::ring::default_provider())
690 }
691
692 #[doc(hidden)]
693 #[cfg(feature = "_hidden")]
694 pub fn with_user_agent(self, user_agent: impl Into<String>) -> Result<Self, ValidationError> {
695 let user_agent = user_agent
696 .into()
697 .parse()
698 .map_err(|e| ValidationError(format!("invalid user agent: {e}")))?;
699 Ok(Self { user_agent, ..self })
700 }
701}
702
703#[cfg(feature = "rustls-aws-lc-rs")]
704fn default_rustls_crypto_provider() -> Option<Arc<rustls::crypto::CryptoProvider>> {
705 Some(Arc::new(rustls::crypto::aws_lc_rs::default_provider()))
706}
707
708#[cfg(all(not(feature = "rustls-aws-lc-rs"), feature = "rustls-ring"))]
709fn default_rustls_crypto_provider() -> Option<Arc<rustls::crypto::CryptoProvider>> {
710 Some(Arc::new(rustls::crypto::ring::default_provider()))
711}
712
713#[cfg(all(not(feature = "rustls-aws-lc-rs"), not(feature = "rustls-ring")))]
714fn default_rustls_crypto_provider() -> Option<Arc<rustls::crypto::CryptoProvider>> {
715 None
716}
717
718#[derive(Debug, Default, Clone, PartialEq, Eq)]
719#[non_exhaustive]
720pub struct Page<T> {
722 pub values: Vec<T>,
724 pub has_more: bool,
726}
727
728impl<T> Page<T> {
729 pub(crate) fn new(values: impl Into<Vec<T>>, has_more: bool) -> Self {
730 Self {
731 values: values.into(),
732 has_more,
733 }
734 }
735}
736
737#[derive(Debug, Clone, Copy, PartialEq, Eq)]
738pub enum RetentionPolicy {
740 Age(u64),
742 Infinite,
744}
745
746impl From<api::config::RetentionPolicy> for RetentionPolicy {
747 fn from(value: api::config::RetentionPolicy) -> Self {
748 match value {
749 api::config::RetentionPolicy::Age(secs) => RetentionPolicy::Age(secs),
750 api::config::RetentionPolicy::Infinite(_) => RetentionPolicy::Infinite,
751 }
752 }
753}
754
755impl From<RetentionPolicy> for api::config::RetentionPolicy {
756 fn from(value: RetentionPolicy) -> Self {
757 match value {
758 RetentionPolicy::Age(secs) => api::config::RetentionPolicy::Age(secs),
759 RetentionPolicy::Infinite => {
760 api::config::RetentionPolicy::Infinite(api::config::InfiniteRetention {})
761 }
762 }
763 }
764}
765
766#[derive(Debug, Clone, Copy, PartialEq, Eq)]
767pub enum TimestampingMode {
769 ClientPrefer,
771 ClientRequire,
773 Arrival,
775}
776
777impl From<api::config::TimestampingMode> for TimestampingMode {
778 fn from(value: api::config::TimestampingMode) -> Self {
779 match value {
780 api::config::TimestampingMode::ClientPrefer => TimestampingMode::ClientPrefer,
781 api::config::TimestampingMode::ClientRequire => TimestampingMode::ClientRequire,
782 api::config::TimestampingMode::Arrival => TimestampingMode::Arrival,
783 }
784 }
785}
786
787impl From<TimestampingMode> for api::config::TimestampingMode {
788 fn from(value: TimestampingMode) -> Self {
789 match value {
790 TimestampingMode::ClientPrefer => api::config::TimestampingMode::ClientPrefer,
791 TimestampingMode::ClientRequire => api::config::TimestampingMode::ClientRequire,
792 TimestampingMode::Arrival => api::config::TimestampingMode::Arrival,
793 }
794 }
795}
796
797#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
798#[non_exhaustive]
799pub struct TimestampingConfig {
801 pub mode: Option<TimestampingMode>,
805 pub uncapped: Option<bool>,
809}
810
811impl TimestampingConfig {
812 pub fn new() -> Self {
814 Self::default()
815 }
816
817 pub fn with_mode(self, mode: TimestampingMode) -> Self {
819 Self {
820 mode: Some(mode),
821 ..self
822 }
823 }
824
825 pub fn with_uncapped(self, uncapped: bool) -> Self {
827 Self {
828 uncapped: Some(uncapped),
829 ..self
830 }
831 }
832}
833
834impl From<api::config::TimestampingConfig> for TimestampingConfig {
835 fn from(value: api::config::TimestampingConfig) -> Self {
836 Self {
837 mode: value.mode.map(Into::into),
838 uncapped: value.uncapped,
839 }
840 }
841}
842
843impl From<TimestampingConfig> for api::config::TimestampingConfig {
844 fn from(value: TimestampingConfig) -> Self {
845 Self {
846 mode: value.mode.map(Into::into),
847 uncapped: value.uncapped,
848 }
849 }
850}
851
852#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
853#[non_exhaustive]
854pub struct DeleteOnEmptyConfig {
856 pub min_age_secs: u64,
860}
861
862impl DeleteOnEmptyConfig {
863 pub fn new() -> Self {
865 Self::default()
866 }
867
868 pub fn with_min_age(self, min_age: Duration) -> Self {
870 Self {
871 min_age_secs: min_age.as_secs(),
872 }
873 }
874}
875
876impl From<api::config::DeleteOnEmptyConfig> for DeleteOnEmptyConfig {
877 fn from(value: api::config::DeleteOnEmptyConfig) -> Self {
878 Self {
879 min_age_secs: value.min_age_secs,
880 }
881 }
882}
883
884impl From<DeleteOnEmptyConfig> for api::config::DeleteOnEmptyConfig {
885 fn from(value: DeleteOnEmptyConfig) -> Self {
886 Self {
887 min_age_secs: value.min_age_secs,
888 }
889 }
890}
891
892#[derive(Debug, Clone, Default, PartialEq, Eq)]
893#[non_exhaustive]
894pub struct StreamConfig {
896 pub storage_class: Option<CompactString>,
898 pub retention_policy: Option<RetentionPolicy>,
902 pub timestamping: Option<TimestampingConfig>,
906 pub delete_on_empty: Option<DeleteOnEmptyConfig>,
910}
911
912impl StreamConfig {
913 pub fn new() -> Self {
915 Self::default()
916 }
917
918 pub fn with_storage_class(self, storage_class: impl Into<CompactString>) -> Self {
920 Self {
921 storage_class: Some(storage_class.into()),
922 ..self
923 }
924 }
925
926 pub fn with_retention_policy(self, retention_policy: RetentionPolicy) -> Self {
928 Self {
929 retention_policy: Some(retention_policy),
930 ..self
931 }
932 }
933
934 pub fn with_timestamping(self, timestamping: TimestampingConfig) -> Self {
936 Self {
937 timestamping: Some(timestamping),
938 ..self
939 }
940 }
941
942 pub fn with_delete_on_empty(self, delete_on_empty: DeleteOnEmptyConfig) -> Self {
944 Self {
945 delete_on_empty: Some(delete_on_empty),
946 ..self
947 }
948 }
949}
950
951impl From<api::config::StreamConfig> for StreamConfig {
952 fn from(value: api::config::StreamConfig) -> Self {
953 Self {
954 storage_class: value.storage_class,
955 retention_policy: value.retention_policy.map(Into::into),
956 timestamping: value.timestamping.map(Into::into),
957 delete_on_empty: value.delete_on_empty.map(Into::into),
958 }
959 }
960}
961
962impl From<StreamConfig> for api::config::StreamConfig {
963 fn from(value: StreamConfig) -> Self {
964 Self {
965 storage_class: value.storage_class,
966 retention_policy: value.retention_policy.map(Into::into),
967 timestamping: value.timestamping.map(Into::into),
968 delete_on_empty: value.delete_on_empty.map(Into::into),
969 }
970 }
971}
972
973#[derive(Debug, Clone, Default, PartialEq, Eq)]
974#[non_exhaustive]
975pub struct BasinConfig {
977 pub default_stream_config: Option<StreamConfig>,
981 pub stream_cipher: Option<EncryptionAlgorithm>,
983 pub create_stream_on_append: bool,
987 pub create_stream_on_read: bool,
991}
992
993impl BasinConfig {
994 pub fn new() -> Self {
996 Self::default()
997 }
998
999 pub fn with_default_stream_config(self, config: StreamConfig) -> Self {
1001 Self {
1002 default_stream_config: Some(config),
1003 ..self
1004 }
1005 }
1006
1007 pub fn with_stream_cipher(self, stream_cipher: EncryptionAlgorithm) -> Self {
1009 Self {
1010 stream_cipher: Some(stream_cipher),
1011 ..self
1012 }
1013 }
1014
1015 pub fn with_create_stream_on_append(self, create_stream_on_append: bool) -> Self {
1018 Self {
1019 create_stream_on_append,
1020 ..self
1021 }
1022 }
1023
1024 pub fn with_create_stream_on_read(self, create_stream_on_read: bool) -> Self {
1026 Self {
1027 create_stream_on_read,
1028 ..self
1029 }
1030 }
1031}
1032
1033impl From<api::config::BasinConfig> for BasinConfig {
1034 fn from(value: api::config::BasinConfig) -> Self {
1035 Self {
1036 default_stream_config: value.default_stream_config.map(Into::into),
1037 stream_cipher: value.stream_cipher.map(Into::into),
1038 create_stream_on_append: value.create_stream_on_append,
1039 create_stream_on_read: value.create_stream_on_read,
1040 }
1041 }
1042}
1043
1044impl From<BasinConfig> for api::config::BasinConfig {
1045 fn from(value: BasinConfig) -> Self {
1046 Self {
1047 default_stream_config: value.default_stream_config.map(Into::into),
1048 stream_cipher: value.stream_cipher.map(Into::into),
1049 create_stream_on_append: value.create_stream_on_append,
1050 create_stream_on_read: value.create_stream_on_read,
1051 }
1052 }
1053}
1054
1055#[derive(Debug, Clone)]
1056#[non_exhaustive]
1057pub struct CreateBasinInput {
1059 pub name: BasinName,
1061 pub config: Option<BasinConfig>,
1065 pub location: Option<LocationName>,
1069 idempotency_token: String,
1070}
1071
1072impl CreateBasinInput {
1073 pub fn new(name: BasinName) -> Self {
1075 Self {
1076 name,
1077 config: None,
1078 location: None,
1079 idempotency_token: idempotency_token(),
1080 }
1081 }
1082
1083 pub fn with_config(self, config: BasinConfig) -> Self {
1085 Self {
1086 config: Some(config),
1087 ..self
1088 }
1089 }
1090
1091 pub fn with_location<S>(self, location: S) -> Result<Self, ValidationError>
1093 where
1094 S: TryInto<LocationName>,
1095 S::Error: fmt::Display,
1096 {
1097 let location = location
1098 .try_into()
1099 .map_err(|e| ValidationError(e.to_string()))?;
1100 Ok(Self {
1101 location: Some(location),
1102 ..self
1103 })
1104 }
1105}
1106
1107impl From<CreateBasinInput> for (api::basin::CreateBasinRequest, String) {
1108 fn from(value: CreateBasinInput) -> Self {
1109 (
1110 api::basin::CreateBasinRequest {
1111 basin: value.name,
1112 config: value.config.map(Into::into),
1113 location: value.location,
1114 },
1115 value.idempotency_token,
1116 )
1117 }
1118}
1119
1120#[derive(Debug, Clone)]
1121#[non_exhaustive]
1122pub struct EnsureBasinInput {
1124 pub name: BasinName,
1126 pub config: Option<BasinConfig>,
1130 pub location: Option<LocationName>,
1135}
1136
1137impl EnsureBasinInput {
1138 pub fn new(name: BasinName) -> Self {
1140 Self {
1141 name,
1142 config: None,
1143 location: None,
1144 }
1145 }
1146
1147 pub fn with_config(self, config: BasinConfig) -> Self {
1149 Self {
1150 config: Some(config),
1151 ..self
1152 }
1153 }
1154
1155 pub fn with_location<S>(self, location: S) -> Result<Self, ValidationError>
1157 where
1158 S: TryInto<LocationName>,
1159 S::Error: fmt::Display,
1160 {
1161 let location = location
1162 .try_into()
1163 .map_err(|e| ValidationError(e.to_string()))?;
1164 Ok(Self {
1165 location: Some(location),
1166 ..self
1167 })
1168 }
1169}
1170
1171impl From<EnsureBasinInput> for (BasinName, Option<api::basin::EnsureBasinRequest>) {
1172 fn from(value: EnsureBasinInput) -> Self {
1173 let config = value.config;
1174 let request = if config.is_some() || value.location.is_some() {
1175 Some(api::basin::EnsureBasinRequest {
1176 config: config.map(Into::into),
1177 location: value.location,
1178 })
1179 } else {
1180 None
1181 };
1182 (value.name, request)
1183 }
1184}
1185
1186#[derive(Debug, Clone)]
1187pub enum EnsureOutput<T> {
1190 Created(T),
1192 ConfigUpdated(T),
1194 ConfigUnchanged(T),
1196}
1197
1198impl<T> From<ProvisionResult<T>> for EnsureOutput<T> {
1199 fn from(result: ProvisionResult<T>) -> Self {
1200 match result {
1201 ProvisionResult::Created(info) => EnsureOutput::Created(info),
1202 ProvisionResult::Updated(info) => EnsureOutput::ConfigUpdated(info),
1203 ProvisionResult::Noop(info) => EnsureOutput::ConfigUnchanged(info),
1204 }
1205 }
1206}
1207
1208#[derive(Debug, Clone, Default)]
1209#[non_exhaustive]
1210pub struct ListBasinsInput {
1212 pub prefix: BasinNamePrefix,
1216 pub start_after: BasinNameStartAfter,
1220 pub limit: Option<usize>,
1224}
1225
1226impl ListBasinsInput {
1227 pub fn new() -> Self {
1229 Self::default()
1230 }
1231
1232 pub fn with_prefix(self, prefix: BasinNamePrefix) -> Self {
1234 Self { prefix, ..self }
1235 }
1236
1237 pub fn with_start_after(self, start_after: BasinNameStartAfter) -> Self {
1240 Self {
1241 start_after,
1242 ..self
1243 }
1244 }
1245
1246 pub fn with_limit(self, limit: usize) -> Self {
1248 Self {
1249 limit: Some(limit),
1250 ..self
1251 }
1252 }
1253}
1254
1255impl From<ListBasinsInput> for api::basin::ListBasinsRequest {
1256 fn from(value: ListBasinsInput) -> Self {
1257 Self {
1258 prefix: Some(value.prefix),
1259 start_after: Some(value.start_after),
1260 limit: value.limit,
1261 }
1262 }
1263}
1264
1265#[derive(Debug, Clone, Default)]
1266pub struct ListAllBasinsInput {
1268 pub prefix: BasinNamePrefix,
1272 pub start_after: BasinNameStartAfter,
1276 pub include_deleted: bool,
1280}
1281
1282impl ListAllBasinsInput {
1283 pub fn new() -> Self {
1285 Self::default()
1286 }
1287
1288 pub fn with_prefix(self, prefix: BasinNamePrefix) -> Self {
1290 Self { prefix, ..self }
1291 }
1292
1293 pub fn with_start_after(self, start_after: BasinNameStartAfter) -> Self {
1296 Self {
1297 start_after,
1298 ..self
1299 }
1300 }
1301
1302 pub fn with_include_deleted(self, include_deleted: bool) -> Self {
1304 Self {
1305 include_deleted,
1306 ..self
1307 }
1308 }
1309}
1310
1311#[derive(Debug, Clone, PartialEq, Eq)]
1312#[non_exhaustive]
1313pub struct BasinInfo {
1315 pub name: BasinName,
1317 pub location: Option<LocationName>,
1319 pub created_at: S2DateTime,
1321 pub deleted_at: Option<S2DateTime>,
1323}
1324
1325impl TryFrom<api::basin::BasinInfo> for BasinInfo {
1326 type Error = ValidationError;
1327
1328 fn try_from(value: api::basin::BasinInfo) -> Result<Self, Self::Error> {
1329 Ok(Self {
1330 name: value.name,
1331 location: value.location,
1332 created_at: value.created_at.try_into()?,
1333 deleted_at: value.deleted_at.map(S2DateTime::try_from).transpose()?,
1334 })
1335 }
1336}
1337
1338#[derive(Debug, Clone)]
1339#[non_exhaustive]
1340pub struct DeleteBasinInput {
1342 pub name: BasinName,
1344 pub ignore_not_found: bool,
1346}
1347
1348impl DeleteBasinInput {
1349 pub fn new(name: BasinName) -> Self {
1351 Self {
1352 name,
1353 ignore_not_found: false,
1354 }
1355 }
1356
1357 pub fn with_ignore_not_found(self, ignore_not_found: bool) -> Self {
1359 Self {
1360 ignore_not_found,
1361 ..self
1362 }
1363 }
1364}
1365
1366#[derive(Debug, Clone, Default)]
1367#[non_exhaustive]
1368pub struct TimestampingReconfiguration {
1370 pub mode: Maybe<Option<TimestampingMode>>,
1372 pub uncapped: Maybe<Option<bool>>,
1374}
1375
1376impl TimestampingReconfiguration {
1377 pub fn new() -> Self {
1379 Self::default()
1380 }
1381
1382 pub fn with_mode(self, mode: TimestampingMode) -> Self {
1384 Self {
1385 mode: Maybe::Specified(Some(mode)),
1386 ..self
1387 }
1388 }
1389
1390 pub fn with_uncapped(self, uncapped: bool) -> Self {
1392 Self {
1393 uncapped: Maybe::Specified(Some(uncapped)),
1394 ..self
1395 }
1396 }
1397}
1398
1399impl From<TimestampingReconfiguration> for api::config::TimestampingReconfiguration {
1400 fn from(value: TimestampingReconfiguration) -> Self {
1401 Self {
1402 mode: value.mode.map(|m| m.map(Into::into)),
1403 uncapped: value.uncapped,
1404 }
1405 }
1406}
1407
1408#[derive(Debug, Clone, Default)]
1409#[non_exhaustive]
1410pub struct DeleteOnEmptyReconfiguration {
1412 pub min_age_secs: Maybe<Option<u64>>,
1414}
1415
1416impl DeleteOnEmptyReconfiguration {
1417 pub fn new() -> Self {
1419 Self::default()
1420 }
1421
1422 pub fn with_min_age(self, min_age: Duration) -> Self {
1424 Self {
1425 min_age_secs: Maybe::Specified(Some(min_age.as_secs())),
1426 }
1427 }
1428}
1429
1430impl From<DeleteOnEmptyReconfiguration> for api::config::DeleteOnEmptyReconfiguration {
1431 fn from(value: DeleteOnEmptyReconfiguration) -> Self {
1432 Self {
1433 min_age_secs: value.min_age_secs,
1434 }
1435 }
1436}
1437
1438#[derive(Debug, Clone, Default)]
1439#[non_exhaustive]
1440pub struct StreamReconfiguration {
1442 pub storage_class: Maybe<Option<CompactString>>,
1444 pub retention_policy: Maybe<Option<RetentionPolicy>>,
1446 pub timestamping: Maybe<Option<TimestampingReconfiguration>>,
1448 pub delete_on_empty: Maybe<Option<DeleteOnEmptyReconfiguration>>,
1450}
1451
1452impl StreamReconfiguration {
1453 pub fn new() -> Self {
1455 Self::default()
1456 }
1457
1458 pub fn with_storage_class(self, storage_class: impl Into<CompactString>) -> Self {
1460 Self {
1461 storage_class: Maybe::Specified(Some(storage_class.into())),
1462 ..self
1463 }
1464 }
1465
1466 pub fn with_retention_policy(self, retention_policy: RetentionPolicy) -> Self {
1468 Self {
1469 retention_policy: Maybe::Specified(Some(retention_policy)),
1470 ..self
1471 }
1472 }
1473
1474 pub fn with_timestamping(self, timestamping: TimestampingReconfiguration) -> Self {
1476 Self {
1477 timestamping: Maybe::Specified(Some(timestamping)),
1478 ..self
1479 }
1480 }
1481
1482 pub fn with_delete_on_empty(self, delete_on_empty: DeleteOnEmptyReconfiguration) -> Self {
1484 Self {
1485 delete_on_empty: Maybe::Specified(Some(delete_on_empty)),
1486 ..self
1487 }
1488 }
1489}
1490
1491impl From<StreamReconfiguration> for api::config::StreamReconfiguration {
1492 fn from(value: StreamReconfiguration) -> Self {
1493 Self {
1494 storage_class: value.storage_class,
1495 retention_policy: value.retention_policy.map(|m| m.map(Into::into)),
1496 timestamping: value.timestamping.map(|m| m.map(Into::into)),
1497 delete_on_empty: value.delete_on_empty.map(|m| m.map(Into::into)),
1498 }
1499 }
1500}
1501
1502#[derive(Debug, Clone, Default)]
1503#[non_exhaustive]
1504pub struct BasinReconfiguration {
1506 pub default_stream_config: Maybe<Option<StreamReconfiguration>>,
1508 pub stream_cipher: Maybe<Option<EncryptionAlgorithm>>,
1510 pub create_stream_on_append: Maybe<bool>,
1513 pub create_stream_on_read: Maybe<bool>,
1515}
1516
1517impl BasinReconfiguration {
1518 pub fn new() -> Self {
1520 Self::default()
1521 }
1522
1523 pub fn with_default_stream_config(self, config: StreamReconfiguration) -> Self {
1526 Self {
1527 default_stream_config: Maybe::Specified(Some(config)),
1528 ..self
1529 }
1530 }
1531
1532 pub fn with_stream_cipher(self, stream_cipher: EncryptionAlgorithm) -> Self {
1534 Self {
1535 stream_cipher: Maybe::Specified(Some(stream_cipher)),
1536 ..self
1537 }
1538 }
1539
1540 pub fn with_create_stream_on_append(self, create_stream_on_append: bool) -> Self {
1543 Self {
1544 create_stream_on_append: Maybe::Specified(create_stream_on_append),
1545 ..self
1546 }
1547 }
1548
1549 pub fn with_create_stream_on_read(self, create_stream_on_read: bool) -> Self {
1552 Self {
1553 create_stream_on_read: Maybe::Specified(create_stream_on_read),
1554 ..self
1555 }
1556 }
1557}
1558
1559impl From<BasinReconfiguration> for api::config::BasinReconfiguration {
1560 fn from(value: BasinReconfiguration) -> Self {
1561 Self {
1562 default_stream_config: value.default_stream_config.map(|m| m.map(Into::into)),
1563 stream_cipher: value.stream_cipher.map(|m| m.map(Into::into)),
1564 create_stream_on_append: value.create_stream_on_append,
1565 create_stream_on_read: value.create_stream_on_read,
1566 }
1567 }
1568}
1569
1570#[derive(Debug, Clone)]
1571#[non_exhaustive]
1572pub struct ReconfigureBasinInput {
1574 pub name: BasinName,
1576 pub config: BasinReconfiguration,
1578}
1579
1580impl ReconfigureBasinInput {
1581 pub fn new(name: BasinName, config: BasinReconfiguration) -> Self {
1583 Self { name, config }
1584 }
1585}
1586
1587#[derive(Debug, Clone, Default)]
1588#[non_exhaustive]
1589pub struct ListAccessTokensInput {
1591 pub prefix: AccessTokenIdPrefix,
1595 pub start_after: AccessTokenIdStartAfter,
1599 pub limit: Option<usize>,
1603}
1604
1605impl ListAccessTokensInput {
1606 pub fn new() -> Self {
1608 Self::default()
1609 }
1610
1611 pub fn with_prefix(self, prefix: AccessTokenIdPrefix) -> Self {
1613 Self { prefix, ..self }
1614 }
1615
1616 pub fn with_start_after(self, start_after: AccessTokenIdStartAfter) -> Self {
1619 Self {
1620 start_after,
1621 ..self
1622 }
1623 }
1624
1625 pub fn with_limit(self, limit: usize) -> Self {
1627 Self {
1628 limit: Some(limit),
1629 ..self
1630 }
1631 }
1632}
1633
1634impl From<ListAccessTokensInput> for api::access::ListAccessTokensRequest {
1635 fn from(value: ListAccessTokensInput) -> Self {
1636 Self {
1637 prefix: Some(value.prefix),
1638 start_after: Some(value.start_after),
1639 limit: value.limit,
1640 }
1641 }
1642}
1643
1644#[derive(Debug, Clone, Default)]
1645pub struct ListAllAccessTokensInput {
1647 pub prefix: AccessTokenIdPrefix,
1651 pub start_after: AccessTokenIdStartAfter,
1655}
1656
1657impl ListAllAccessTokensInput {
1658 pub fn new() -> Self {
1660 Self::default()
1661 }
1662
1663 pub fn with_prefix(self, prefix: AccessTokenIdPrefix) -> Self {
1665 Self { prefix, ..self }
1666 }
1667
1668 pub fn with_start_after(self, start_after: AccessTokenIdStartAfter) -> Self {
1671 Self {
1672 start_after,
1673 ..self
1674 }
1675 }
1676}
1677
1678#[derive(Debug, Clone, PartialEq, Eq)]
1679#[non_exhaustive]
1680pub struct LocationInfo {
1682 pub name: LocationName,
1684 pub is_private: bool,
1686 pub storage_classes: Option<Vec<CompactString>>,
1688 pub default_storage_class: Option<CompactString>,
1690}
1691
1692impl From<api::location::LocationInfo> for LocationInfo {
1693 fn from(value: api::location::LocationInfo) -> Self {
1694 Self {
1695 name: value.name,
1696 is_private: value.is_private,
1697 storage_classes: value.storage_classes,
1698 default_storage_class: value.default_storage_class,
1699 }
1700 }
1701}
1702
1703#[derive(Debug, Clone)]
1704#[non_exhaustive]
1705pub struct AccessTokenInfo {
1707 pub id: AccessTokenId,
1709 pub expires_at: Option<S2DateTime>,
1711 pub auto_prefix_streams: bool,
1714 pub scope: AccessTokenScope,
1716}
1717
1718impl TryFrom<api::access::AccessTokenInfo> for AccessTokenInfo {
1719 type Error = ValidationError;
1720
1721 fn try_from(value: api::access::AccessTokenInfo) -> Result<Self, Self::Error> {
1722 let expires_at = value.expires_at.map(S2DateTime::try_from).transpose()?;
1723 Ok(Self {
1724 id: value.id,
1725 expires_at,
1726 auto_prefix_streams: value.auto_prefix_streams,
1727 scope: value.scope.into(),
1728 })
1729 }
1730}
1731
1732#[derive(Debug, Clone)]
1733pub enum BasinMatcher {
1737 None,
1739 Exact(BasinName),
1741 Prefix(BasinNamePrefix),
1743}
1744
1745#[derive(Debug, Clone)]
1746pub enum StreamMatcher {
1750 None,
1752 Exact(StreamName),
1754 Prefix(StreamNamePrefix),
1756}
1757
1758#[derive(Debug, Clone)]
1759pub enum AccessTokenMatcher {
1763 None,
1765 Exact(AccessTokenId),
1767 Prefix(AccessTokenIdPrefix),
1769}
1770
1771#[derive(Debug, Clone, Default)]
1772#[non_exhaustive]
1773pub struct ReadWritePermissions {
1775 pub read: bool,
1779 pub write: bool,
1783}
1784
1785impl ReadWritePermissions {
1786 pub fn new() -> Self {
1788 Self::default()
1789 }
1790
1791 pub fn read_only() -> Self {
1793 Self {
1794 read: true,
1795 write: false,
1796 }
1797 }
1798
1799 pub fn write_only() -> Self {
1801 Self {
1802 read: false,
1803 write: true,
1804 }
1805 }
1806
1807 pub fn read_write() -> Self {
1809 Self {
1810 read: true,
1811 write: true,
1812 }
1813 }
1814}
1815
1816impl From<ReadWritePermissions> for api::access::ReadWritePermissions {
1817 fn from(value: ReadWritePermissions) -> Self {
1818 Self {
1819 read: Some(value.read),
1820 write: Some(value.write),
1821 }
1822 }
1823}
1824
1825impl From<api::access::ReadWritePermissions> for ReadWritePermissions {
1826 fn from(value: api::access::ReadWritePermissions) -> Self {
1827 Self {
1828 read: value.read.unwrap_or_default(),
1829 write: value.write.unwrap_or_default(),
1830 }
1831 }
1832}
1833
1834#[derive(Debug, Clone, Default)]
1835#[non_exhaustive]
1836pub struct OperationGroupPermissions {
1840 pub account: Option<ReadWritePermissions>,
1844 pub basin: Option<ReadWritePermissions>,
1848 pub stream: Option<ReadWritePermissions>,
1852}
1853
1854impl OperationGroupPermissions {
1855 pub fn new() -> Self {
1857 Self::default()
1858 }
1859
1860 pub fn read_only_all() -> Self {
1862 Self {
1863 account: Some(ReadWritePermissions::read_only()),
1864 basin: Some(ReadWritePermissions::read_only()),
1865 stream: Some(ReadWritePermissions::read_only()),
1866 }
1867 }
1868
1869 pub fn write_only_all() -> Self {
1871 Self {
1872 account: Some(ReadWritePermissions::write_only()),
1873 basin: Some(ReadWritePermissions::write_only()),
1874 stream: Some(ReadWritePermissions::write_only()),
1875 }
1876 }
1877
1878 pub fn read_write_all() -> Self {
1880 Self {
1881 account: Some(ReadWritePermissions::read_write()),
1882 basin: Some(ReadWritePermissions::read_write()),
1883 stream: Some(ReadWritePermissions::read_write()),
1884 }
1885 }
1886
1887 pub fn with_account(self, account: ReadWritePermissions) -> Self {
1889 Self {
1890 account: Some(account),
1891 ..self
1892 }
1893 }
1894
1895 pub fn with_basin(self, basin: ReadWritePermissions) -> Self {
1897 Self {
1898 basin: Some(basin),
1899 ..self
1900 }
1901 }
1902
1903 pub fn with_stream(self, stream: ReadWritePermissions) -> Self {
1905 Self {
1906 stream: Some(stream),
1907 ..self
1908 }
1909 }
1910}
1911
1912impl From<OperationGroupPermissions> for api::access::PermittedOperationGroups {
1913 fn from(value: OperationGroupPermissions) -> Self {
1914 Self {
1915 account: value.account.map(Into::into),
1916 basin: value.basin.map(Into::into),
1917 stream: value.stream.map(Into::into),
1918 }
1919 }
1920}
1921
1922impl From<api::access::PermittedOperationGroups> for OperationGroupPermissions {
1923 fn from(value: api::access::PermittedOperationGroups) -> Self {
1924 Self {
1925 account: value.account.map(Into::into),
1926 basin: value.basin.map(Into::into),
1927 stream: value.stream.map(Into::into),
1928 }
1929 }
1930}
1931
1932#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
1933pub enum Operation {
1937 ListBasins,
1939 CreateBasin,
1941 GetBasinConfig,
1943 DeleteBasin,
1945 ReconfigureBasin,
1947 ListAccessTokens,
1949 IssueAccessToken,
1951 RevokeAccessToken,
1953 GetAccountMetrics,
1955 GetBasinMetrics,
1957 GetStreamMetrics,
1959 ListStreams,
1961 CreateStream,
1963 GetStreamConfig,
1965 DeleteStream,
1967 ReconfigureStream,
1969 CheckTail,
1971 Append,
1973 Read,
1975 Trim,
1977 Fence,
1979 ListLocations,
1981 GetDefaultLocation,
1983 SetDefaultLocation,
1985}
1986
1987impl From<Operation> for api::access::Operation {
1988 fn from(value: Operation) -> Self {
1989 match value {
1990 Operation::ListBasins => api::access::Operation::ListBasins,
1991 Operation::CreateBasin => api::access::Operation::CreateBasin,
1992 Operation::DeleteBasin => api::access::Operation::DeleteBasin,
1993 Operation::ReconfigureBasin => api::access::Operation::ReconfigureBasin,
1994 Operation::GetBasinConfig => api::access::Operation::GetBasinConfig,
1995 Operation::IssueAccessToken => api::access::Operation::IssueAccessToken,
1996 Operation::RevokeAccessToken => api::access::Operation::RevokeAccessToken,
1997 Operation::ListAccessTokens => api::access::Operation::ListAccessTokens,
1998 Operation::ListStreams => api::access::Operation::ListStreams,
1999 Operation::CreateStream => api::access::Operation::CreateStream,
2000 Operation::DeleteStream => api::access::Operation::DeleteStream,
2001 Operation::GetStreamConfig => api::access::Operation::GetStreamConfig,
2002 Operation::ReconfigureStream => api::access::Operation::ReconfigureStream,
2003 Operation::CheckTail => api::access::Operation::CheckTail,
2004 Operation::Append => api::access::Operation::Append,
2005 Operation::Read => api::access::Operation::Read,
2006 Operation::Trim => api::access::Operation::Trim,
2007 Operation::Fence => api::access::Operation::Fence,
2008 Operation::GetAccountMetrics => api::access::Operation::AccountMetrics,
2009 Operation::GetBasinMetrics => api::access::Operation::BasinMetrics,
2010 Operation::GetStreamMetrics => api::access::Operation::StreamMetrics,
2011 Operation::ListLocations => api::access::Operation::ListLocations,
2012 Operation::GetDefaultLocation => api::access::Operation::GetDefaultLocation,
2013 Operation::SetDefaultLocation => api::access::Operation::SetDefaultLocation,
2014 }
2015 }
2016}
2017
2018impl From<api::access::Operation> for Operation {
2019 fn from(value: api::access::Operation) -> Self {
2020 match value {
2021 api::access::Operation::ListBasins => Operation::ListBasins,
2022 api::access::Operation::CreateBasin => Operation::CreateBasin,
2023 api::access::Operation::DeleteBasin => Operation::DeleteBasin,
2024 api::access::Operation::ReconfigureBasin => Operation::ReconfigureBasin,
2025 api::access::Operation::GetBasinConfig => Operation::GetBasinConfig,
2026 api::access::Operation::IssueAccessToken => Operation::IssueAccessToken,
2027 api::access::Operation::RevokeAccessToken => Operation::RevokeAccessToken,
2028 api::access::Operation::ListAccessTokens => Operation::ListAccessTokens,
2029 api::access::Operation::ListStreams => Operation::ListStreams,
2030 api::access::Operation::CreateStream => Operation::CreateStream,
2031 api::access::Operation::DeleteStream => Operation::DeleteStream,
2032 api::access::Operation::GetStreamConfig => Operation::GetStreamConfig,
2033 api::access::Operation::ReconfigureStream => Operation::ReconfigureStream,
2034 api::access::Operation::CheckTail => Operation::CheckTail,
2035 api::access::Operation::Append => Operation::Append,
2036 api::access::Operation::Read => Operation::Read,
2037 api::access::Operation::Trim => Operation::Trim,
2038 api::access::Operation::Fence => Operation::Fence,
2039 api::access::Operation::AccountMetrics => Operation::GetAccountMetrics,
2040 api::access::Operation::BasinMetrics => Operation::GetBasinMetrics,
2041 api::access::Operation::StreamMetrics => Operation::GetStreamMetrics,
2042 api::access::Operation::ListLocations => Operation::ListLocations,
2043 api::access::Operation::GetDefaultLocation => Operation::GetDefaultLocation,
2044 api::access::Operation::SetDefaultLocation => Operation::SetDefaultLocation,
2045 }
2046 }
2047}
2048
2049#[derive(Debug, Clone)]
2050#[non_exhaustive]
2051pub struct AccessTokenScopeInput {
2059 basins: Option<BasinMatcher>,
2060 streams: Option<StreamMatcher>,
2061 access_tokens: Option<AccessTokenMatcher>,
2062 op_group_perms: Option<OperationGroupPermissions>,
2063 ops: HashSet<Operation>,
2064}
2065
2066impl AccessTokenScopeInput {
2067 pub fn from_ops(ops: impl IntoIterator<Item = Operation>) -> Self {
2069 Self {
2070 basins: None,
2071 streams: None,
2072 access_tokens: None,
2073 op_group_perms: None,
2074 ops: ops.into_iter().collect(),
2075 }
2076 }
2077
2078 pub fn from_op_group_perms(op_group_perms: OperationGroupPermissions) -> Self {
2080 Self {
2081 basins: None,
2082 streams: None,
2083 access_tokens: None,
2084 op_group_perms: Some(op_group_perms),
2085 ops: HashSet::default(),
2086 }
2087 }
2088
2089 pub fn with_ops(self, ops: impl IntoIterator<Item = Operation>) -> Self {
2091 Self {
2092 ops: ops.into_iter().collect(),
2093 ..self
2094 }
2095 }
2096
2097 pub fn with_op_group_perms(self, op_group_perms: OperationGroupPermissions) -> Self {
2099 Self {
2100 op_group_perms: Some(op_group_perms),
2101 ..self
2102 }
2103 }
2104
2105 pub fn with_basins(self, basins: BasinMatcher) -> Self {
2109 Self {
2110 basins: Some(basins),
2111 ..self
2112 }
2113 }
2114
2115 pub fn with_streams(self, streams: StreamMatcher) -> Self {
2119 Self {
2120 streams: Some(streams),
2121 ..self
2122 }
2123 }
2124
2125 pub fn with_access_tokens(self, access_tokens: AccessTokenMatcher) -> Self {
2129 Self {
2130 access_tokens: Some(access_tokens),
2131 ..self
2132 }
2133 }
2134}
2135
2136#[derive(Debug, Clone)]
2137#[non_exhaustive]
2138pub struct AccessTokenScope {
2140 pub basins: Option<BasinMatcher>,
2142 pub streams: Option<StreamMatcher>,
2144 pub access_tokens: Option<AccessTokenMatcher>,
2146 pub op_group_perms: Option<OperationGroupPermissions>,
2148 pub ops: HashSet<Operation>,
2150}
2151
2152impl From<api::access::AccessTokenScope> for AccessTokenScope {
2153 fn from(value: api::access::AccessTokenScope) -> Self {
2154 Self {
2155 basins: value.basins.map(|rs| match rs {
2156 api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e)) => {
2157 BasinMatcher::Exact(e)
2158 }
2159 api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty) => {
2160 BasinMatcher::None
2161 }
2162 api::access::ResourceSet::Prefix(p) => BasinMatcher::Prefix(p),
2163 }),
2164 streams: value.streams.map(|rs| match rs {
2165 api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e)) => {
2166 StreamMatcher::Exact(e)
2167 }
2168 api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty) => {
2169 StreamMatcher::None
2170 }
2171 api::access::ResourceSet::Prefix(p) => StreamMatcher::Prefix(p),
2172 }),
2173 access_tokens: value.access_tokens.map(|rs| match rs {
2174 api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e)) => {
2175 AccessTokenMatcher::Exact(e)
2176 }
2177 api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty) => {
2178 AccessTokenMatcher::None
2179 }
2180 api::access::ResourceSet::Prefix(p) => AccessTokenMatcher::Prefix(p),
2181 }),
2182 op_group_perms: value.op_groups.map(Into::into),
2183 ops: value
2184 .ops
2185 .map(|ops| ops.into_iter().map(Into::into).collect())
2186 .unwrap_or_default(),
2187 }
2188 }
2189}
2190
2191impl From<AccessTokenScopeInput> for api::access::AccessTokenScope {
2192 fn from(value: AccessTokenScopeInput) -> Self {
2193 Self {
2194 basins: value.basins.map(|rs| match rs {
2195 BasinMatcher::None => {
2196 api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty)
2197 }
2198 BasinMatcher::Exact(e) => {
2199 api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e))
2200 }
2201 BasinMatcher::Prefix(p) => api::access::ResourceSet::Prefix(p),
2202 }),
2203 streams: value.streams.map(|rs| match rs {
2204 StreamMatcher::None => {
2205 api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty)
2206 }
2207 StreamMatcher::Exact(e) => {
2208 api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e))
2209 }
2210 StreamMatcher::Prefix(p) => api::access::ResourceSet::Prefix(p),
2211 }),
2212 access_tokens: value.access_tokens.map(|rs| match rs {
2213 AccessTokenMatcher::None => {
2214 api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty)
2215 }
2216 AccessTokenMatcher::Exact(e) => {
2217 api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e))
2218 }
2219 AccessTokenMatcher::Prefix(p) => api::access::ResourceSet::Prefix(p),
2220 }),
2221 op_groups: value.op_group_perms.map(Into::into),
2222 ops: if value.ops.is_empty() {
2223 None
2224 } else {
2225 Some(value.ops.into_iter().map(Into::into).collect())
2226 },
2227 }
2228 }
2229}
2230
2231#[derive(Debug, Clone)]
2232#[non_exhaustive]
2233pub struct IssueAccessTokenInput {
2235 pub id: AccessTokenId,
2237 pub expires_at: Option<S2DateTime>,
2242 pub auto_prefix_streams: bool,
2250 pub scope: AccessTokenScopeInput,
2252}
2253
2254impl IssueAccessTokenInput {
2255 pub fn new(id: AccessTokenId, scope: AccessTokenScopeInput) -> Self {
2257 Self {
2258 id,
2259 expires_at: None,
2260 auto_prefix_streams: false,
2261 scope,
2262 }
2263 }
2264
2265 pub fn with_expires_at(self, expires_at: S2DateTime) -> Self {
2267 Self {
2268 expires_at: Some(expires_at),
2269 ..self
2270 }
2271 }
2272
2273 pub fn with_auto_prefix_streams(self, auto_prefix_streams: bool) -> Self {
2276 Self {
2277 auto_prefix_streams,
2278 ..self
2279 }
2280 }
2281}
2282
2283impl From<IssueAccessTokenInput> for api::access::IssueAccessTokenRequest {
2284 fn from(value: IssueAccessTokenInput) -> Self {
2285 Self {
2286 id: value.id,
2287 expires_at: value.expires_at.map(Into::into),
2288 auto_prefix_streams: value.auto_prefix_streams.then_some(true),
2289 scope: value.scope.into(),
2290 }
2291 }
2292}
2293
2294#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2295pub enum TimeseriesInterval {
2297 Minute,
2299 Hour,
2301 Day,
2303}
2304
2305impl From<TimeseriesInterval> for api::metrics::TimeseriesInterval {
2306 fn from(value: TimeseriesInterval) -> Self {
2307 match value {
2308 TimeseriesInterval::Minute => api::metrics::TimeseriesInterval::Minute,
2309 TimeseriesInterval::Hour => api::metrics::TimeseriesInterval::Hour,
2310 TimeseriesInterval::Day => api::metrics::TimeseriesInterval::Day,
2311 }
2312 }
2313}
2314
2315impl From<api::metrics::TimeseriesInterval> for TimeseriesInterval {
2316 fn from(value: api::metrics::TimeseriesInterval) -> Self {
2317 match value {
2318 api::metrics::TimeseriesInterval::Minute => TimeseriesInterval::Minute,
2319 api::metrics::TimeseriesInterval::Hour => TimeseriesInterval::Hour,
2320 api::metrics::TimeseriesInterval::Day => TimeseriesInterval::Day,
2321 }
2322 }
2323}
2324
2325#[derive(Debug, Clone, Copy)]
2326#[non_exhaustive]
2327pub struct TimeRange {
2329 pub start: u32,
2331 pub end: u32,
2333}
2334
2335impl TimeRange {
2336 pub fn new(start: u32, end: u32) -> Self {
2338 Self { start, end }
2339 }
2340}
2341
2342#[derive(Debug, Clone, Copy)]
2343#[non_exhaustive]
2344pub struct TimeRangeAndInterval {
2346 pub start: u32,
2348 pub end: u32,
2350 pub interval: Option<TimeseriesInterval>,
2354}
2355
2356impl TimeRangeAndInterval {
2357 pub fn new(start: u32, end: u32) -> Self {
2359 Self {
2360 start,
2361 end,
2362 interval: None,
2363 }
2364 }
2365
2366 pub fn with_interval(self, interval: TimeseriesInterval) -> Self {
2368 Self {
2369 interval: Some(interval),
2370 ..self
2371 }
2372 }
2373}
2374
2375#[derive(Debug, Clone, Copy)]
2376pub enum AccountMetricSet {
2378 ActiveBasins(TimeRange),
2381 AccountOps(TimeRangeAndInterval),
2388}
2389
2390#[derive(Debug, Clone)]
2391#[non_exhaustive]
2392pub struct GetAccountMetricsInput {
2394 pub set: AccountMetricSet,
2396}
2397
2398impl GetAccountMetricsInput {
2399 pub fn new(set: AccountMetricSet) -> Self {
2401 Self { set }
2402 }
2403}
2404
2405impl From<GetAccountMetricsInput> for api::metrics::AccountMetricSetRequest {
2406 fn from(value: GetAccountMetricsInput) -> Self {
2407 let (set, start, end, interval) = match value.set {
2408 AccountMetricSet::ActiveBasins(args) => (
2409 api::metrics::AccountMetricSet::ActiveBasins,
2410 args.start,
2411 args.end,
2412 None,
2413 ),
2414 AccountMetricSet::AccountOps(args) => (
2415 api::metrics::AccountMetricSet::AccountOps,
2416 args.start,
2417 args.end,
2418 args.interval,
2419 ),
2420 };
2421 Self {
2422 set,
2423 start: Some(start),
2424 end: Some(end),
2425 interval: interval.map(Into::into),
2426 }
2427 }
2428}
2429
2430#[derive(Debug, Clone, Copy)]
2431pub enum BasinMetricSet {
2433 Storage(TimeRange),
2436 AppendOps(TimeRangeAndInterval),
2444 ReadOps(TimeRangeAndInterval),
2452 ReadThroughput(TimeRangeAndInterval),
2459 AppendThroughput(TimeRangeAndInterval),
2466 BasinOps(TimeRangeAndInterval),
2473}
2474
2475#[derive(Debug, Clone)]
2476#[non_exhaustive]
2477pub struct GetBasinMetricsInput {
2479 pub name: BasinName,
2481 pub set: BasinMetricSet,
2483}
2484
2485impl GetBasinMetricsInput {
2486 pub fn new(name: BasinName, set: BasinMetricSet) -> Self {
2488 Self { name, set }
2489 }
2490}
2491
2492impl From<GetBasinMetricsInput> for (BasinName, api::metrics::BasinMetricSetRequest) {
2493 fn from(value: GetBasinMetricsInput) -> Self {
2494 let (set, start, end, interval) = match value.set {
2495 BasinMetricSet::Storage(args) => (
2496 api::metrics::BasinMetricSet::Storage,
2497 args.start,
2498 args.end,
2499 None,
2500 ),
2501 BasinMetricSet::AppendOps(args) => (
2502 api::metrics::BasinMetricSet::AppendOps,
2503 args.start,
2504 args.end,
2505 args.interval,
2506 ),
2507 BasinMetricSet::ReadOps(args) => (
2508 api::metrics::BasinMetricSet::ReadOps,
2509 args.start,
2510 args.end,
2511 args.interval,
2512 ),
2513 BasinMetricSet::ReadThroughput(args) => (
2514 api::metrics::BasinMetricSet::ReadThroughput,
2515 args.start,
2516 args.end,
2517 args.interval,
2518 ),
2519 BasinMetricSet::AppendThroughput(args) => (
2520 api::metrics::BasinMetricSet::AppendThroughput,
2521 args.start,
2522 args.end,
2523 args.interval,
2524 ),
2525 BasinMetricSet::BasinOps(args) => (
2526 api::metrics::BasinMetricSet::BasinOps,
2527 args.start,
2528 args.end,
2529 args.interval,
2530 ),
2531 };
2532 (
2533 value.name,
2534 api::metrics::BasinMetricSetRequest {
2535 set,
2536 start: Some(start),
2537 end: Some(end),
2538 interval: interval.map(Into::into),
2539 },
2540 )
2541 }
2542}
2543
2544#[derive(Debug, Clone, Copy)]
2545pub enum StreamMetricSet {
2547 Storage(TimeRange),
2550}
2551
2552#[derive(Debug, Clone)]
2553#[non_exhaustive]
2554pub struct GetStreamMetricsInput {
2556 pub basin_name: BasinName,
2558 pub stream_name: StreamName,
2560 pub set: StreamMetricSet,
2562}
2563
2564impl GetStreamMetricsInput {
2565 pub fn new(basin_name: BasinName, stream_name: StreamName, set: StreamMetricSet) -> Self {
2568 Self {
2569 basin_name,
2570 stream_name,
2571 set,
2572 }
2573 }
2574}
2575
2576impl From<GetStreamMetricsInput> for (BasinName, StreamName, api::metrics::StreamMetricSetRequest) {
2577 fn from(value: GetStreamMetricsInput) -> Self {
2578 let (set, start, end, interval) = match value.set {
2579 StreamMetricSet::Storage(args) => (
2580 api::metrics::StreamMetricSet::Storage,
2581 args.start,
2582 args.end,
2583 None,
2584 ),
2585 };
2586 (
2587 value.basin_name,
2588 value.stream_name,
2589 api::metrics::StreamMetricSetRequest {
2590 set,
2591 start: Some(start),
2592 end: Some(end),
2593 interval,
2594 },
2595 )
2596 }
2597}
2598
2599#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2600pub enum MetricUnit {
2602 Bytes,
2604 Operations,
2606}
2607
2608impl From<api::metrics::MetricUnit> for MetricUnit {
2609 fn from(value: api::metrics::MetricUnit) -> Self {
2610 match value {
2611 api::metrics::MetricUnit::Bytes => MetricUnit::Bytes,
2612 api::metrics::MetricUnit::Operations => MetricUnit::Operations,
2613 }
2614 }
2615}
2616
2617#[derive(Debug, Clone)]
2618#[non_exhaustive]
2619pub struct ScalarMetric {
2621 pub name: String,
2623 pub unit: MetricUnit,
2625 pub value: f64,
2627}
2628
2629#[derive(Debug, Clone)]
2630#[non_exhaustive]
2631pub struct AccumulationMetric {
2634 pub name: String,
2636 pub unit: MetricUnit,
2638 pub interval: TimeseriesInterval,
2640 pub values: Vec<(u32, f64)>,
2644}
2645
2646#[derive(Debug, Clone)]
2647#[non_exhaustive]
2648pub struct GaugeMetric {
2650 pub name: String,
2652 pub unit: MetricUnit,
2654 pub values: Vec<(u32, f64)>,
2657}
2658
2659#[derive(Debug, Clone)]
2660#[non_exhaustive]
2661pub struct LabelMetric {
2663 pub name: String,
2665 pub values: Vec<String>,
2667}
2668
2669#[derive(Debug, Clone)]
2670pub enum Metric {
2672 Scalar(ScalarMetric),
2674 Accumulation(AccumulationMetric),
2677 Gauge(GaugeMetric),
2679 Label(LabelMetric),
2681}
2682
2683impl From<api::metrics::Metric> for Metric {
2684 fn from(value: api::metrics::Metric) -> Self {
2685 match value {
2686 api::metrics::Metric::Scalar(sm) => Metric::Scalar(ScalarMetric {
2687 name: sm.name.into(),
2688 unit: sm.unit.into(),
2689 value: sm.value,
2690 }),
2691 api::metrics::Metric::Accumulation(am) => Metric::Accumulation(AccumulationMetric {
2692 name: am.name.into(),
2693 unit: am.unit.into(),
2694 interval: am.interval.into(),
2695 values: am.values,
2696 }),
2697 api::metrics::Metric::Gauge(gm) => Metric::Gauge(GaugeMetric {
2698 name: gm.name.into(),
2699 unit: gm.unit.into(),
2700 values: gm.values,
2701 }),
2702 api::metrics::Metric::Label(lm) => Metric::Label(LabelMetric {
2703 name: lm.name.into(),
2704 values: lm.values,
2705 }),
2706 }
2707 }
2708}
2709
2710#[derive(Debug, Clone, Default)]
2711#[non_exhaustive]
2712pub struct ListStreamsInput {
2714 pub prefix: StreamNamePrefix,
2718 pub start_after: StreamNameStartAfter,
2722 pub limit: Option<usize>,
2726}
2727
2728impl ListStreamsInput {
2729 pub fn new() -> Self {
2731 Self::default()
2732 }
2733
2734 pub fn with_prefix(self, prefix: StreamNamePrefix) -> Self {
2736 Self { prefix, ..self }
2737 }
2738
2739 pub fn with_start_after(self, start_after: StreamNameStartAfter) -> Self {
2742 Self {
2743 start_after,
2744 ..self
2745 }
2746 }
2747
2748 pub fn with_limit(self, limit: usize) -> Self {
2750 Self {
2751 limit: Some(limit),
2752 ..self
2753 }
2754 }
2755}
2756
2757impl From<ListStreamsInput> for api::stream::ListStreamsRequest {
2758 fn from(value: ListStreamsInput) -> Self {
2759 Self {
2760 prefix: Some(value.prefix),
2761 start_after: Some(value.start_after),
2762 limit: value.limit,
2763 }
2764 }
2765}
2766
2767#[derive(Debug, Clone, Default)]
2768pub struct ListAllStreamsInput {
2770 pub prefix: StreamNamePrefix,
2774 pub start_after: StreamNameStartAfter,
2778 pub include_deleted: bool,
2782}
2783
2784impl ListAllStreamsInput {
2785 pub fn new() -> Self {
2787 Self::default()
2788 }
2789
2790 pub fn with_prefix(self, prefix: StreamNamePrefix) -> Self {
2792 Self { prefix, ..self }
2793 }
2794
2795 pub fn with_start_after(self, start_after: StreamNameStartAfter) -> Self {
2798 Self {
2799 start_after,
2800 ..self
2801 }
2802 }
2803
2804 pub fn with_include_deleted(self, include_deleted: bool) -> Self {
2806 Self {
2807 include_deleted,
2808 ..self
2809 }
2810 }
2811}
2812
2813#[derive(Debug, Clone, PartialEq, Eq)]
2814#[non_exhaustive]
2815pub struct StreamInfo {
2817 pub name: StreamName,
2819 pub created_at: S2DateTime,
2821 pub deleted_at: Option<S2DateTime>,
2823 pub cipher: Option<EncryptionAlgorithm>,
2825}
2826
2827impl TryFrom<api::stream::StreamInfo> for StreamInfo {
2828 type Error = ValidationError;
2829
2830 fn try_from(value: api::stream::StreamInfo) -> Result<Self, Self::Error> {
2831 Ok(Self {
2832 name: value.name,
2833 created_at: value.created_at.try_into()?,
2834 deleted_at: value.deleted_at.map(S2DateTime::try_from).transpose()?,
2835 cipher: value.cipher.map(Into::into),
2836 })
2837 }
2838}
2839
2840#[derive(Debug, Clone)]
2841#[non_exhaustive]
2842pub struct CreateStreamInput {
2844 pub name: StreamName,
2846 pub config: Option<StreamConfig>,
2850 idempotency_token: String,
2851}
2852
2853impl CreateStreamInput {
2854 pub fn new(name: StreamName) -> Self {
2856 Self {
2857 name,
2858 config: None,
2859 idempotency_token: idempotency_token(),
2860 }
2861 }
2862
2863 pub fn with_config(self, config: StreamConfig) -> Self {
2865 Self {
2866 config: Some(config),
2867 ..self
2868 }
2869 }
2870}
2871
2872impl From<CreateStreamInput> for (api::stream::CreateStreamRequest, String) {
2873 fn from(value: CreateStreamInput) -> Self {
2874 (
2875 api::stream::CreateStreamRequest {
2876 stream: value.name,
2877 config: value.config.map(Into::into),
2878 },
2879 value.idempotency_token,
2880 )
2881 }
2882}
2883
2884#[derive(Debug, Clone)]
2885#[non_exhaustive]
2886pub struct EnsureStreamInput {
2889 pub name: StreamName,
2891 pub config: Option<StreamConfig>,
2895}
2896
2897impl EnsureStreamInput {
2898 pub fn new(name: StreamName) -> Self {
2900 Self { name, config: None }
2901 }
2902
2903 pub fn with_config(self, config: StreamConfig) -> Self {
2905 Self {
2906 config: Some(config),
2907 ..self
2908 }
2909 }
2910}
2911
2912impl From<EnsureStreamInput> for (StreamName, Option<api::config::StreamConfig>) {
2913 fn from(value: EnsureStreamInput) -> Self {
2914 (value.name, value.config.map(Into::into))
2915 }
2916}
2917
2918#[derive(Debug, Clone)]
2919#[non_exhaustive]
2920pub struct DeleteStreamInput {
2922 pub name: StreamName,
2924 pub ignore_not_found: bool,
2926}
2927
2928impl DeleteStreamInput {
2929 pub fn new(name: StreamName) -> Self {
2931 Self {
2932 name,
2933 ignore_not_found: false,
2934 }
2935 }
2936
2937 pub fn with_ignore_not_found(self, ignore_not_found: bool) -> Self {
2939 Self {
2940 ignore_not_found,
2941 ..self
2942 }
2943 }
2944}
2945
2946#[derive(Debug, Clone)]
2947#[non_exhaustive]
2948pub struct ReconfigureStreamInput {
2950 pub name: StreamName,
2952 pub config: StreamReconfiguration,
2954}
2955
2956impl ReconfigureStreamInput {
2957 pub fn new(name: StreamName, config: StreamReconfiguration) -> Self {
2959 Self { name, config }
2960 }
2961}
2962
2963#[derive(Debug, Clone, PartialEq, Eq)]
2964pub struct FencingToken(String);
2970
2971impl FencingToken {
2972 pub(crate) fn from_server(value: String) -> Self {
2973 Self(value)
2974 }
2975
2976 pub fn generate(n: usize) -> Result<Self, ValidationError> {
2978 rand::rng()
2979 .sample_iter(&rand::distr::Alphanumeric)
2980 .take(n)
2981 .map(char::from)
2982 .collect::<String>()
2983 .parse()
2984 }
2985}
2986
2987impl FromStr for FencingToken {
2988 type Err = ValidationError;
2989
2990 fn from_str(s: &str) -> Result<Self, Self::Err> {
2991 if s.len() > MAX_FENCING_TOKEN_LENGTH {
2992 return Err(ValidationError(format!(
2993 "fencing token exceeds {MAX_FENCING_TOKEN_LENGTH} bytes in length",
2994 )));
2995 }
2996 Ok(FencingToken(s.to_string()))
2997 }
2998}
2999
3000impl std::fmt::Display for FencingToken {
3001 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
3002 write!(f, "{}", self.0)
3003 }
3004}
3005
3006impl Deref for FencingToken {
3007 type Target = str;
3008
3009 fn deref(&self) -> &Self::Target {
3010 &self.0
3011 }
3012}
3013
3014#[derive(Debug, Clone, Copy, PartialEq)]
3015#[non_exhaustive]
3016pub struct StreamPosition {
3018 pub seq_num: u64,
3020 pub timestamp: u64,
3023}
3024
3025impl StreamPosition {
3026 pub fn new(seq_num: u64, timestamp: u64) -> Self {
3030 Self { seq_num, timestamp }
3031 }
3032}
3033
3034impl std::fmt::Display for StreamPosition {
3035 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
3036 write!(f, "seq_num={}, timestamp={}", self.seq_num, self.timestamp)
3037 }
3038}
3039
3040impl From<api::stream::proto::StreamPosition> for StreamPosition {
3041 fn from(value: api::stream::proto::StreamPosition) -> Self {
3042 Self {
3043 seq_num: value.seq_num,
3044 timestamp: value.timestamp,
3045 }
3046 }
3047}
3048
3049impl From<api::stream::StreamPosition> for StreamPosition {
3050 fn from(value: api::stream::StreamPosition) -> Self {
3051 Self {
3052 seq_num: value.seq_num,
3053 timestamp: value.timestamp,
3054 }
3055 }
3056}
3057
3058#[derive(Debug, Clone, PartialEq)]
3059#[non_exhaustive]
3060pub struct Header {
3062 pub name: Bytes,
3064 pub value: Bytes,
3066}
3067
3068impl Header {
3069 pub fn new(name: impl Into<Bytes>, value: impl Into<Bytes>) -> Self {
3071 Self {
3072 name: name.into(),
3073 value: value.into(),
3074 }
3075 }
3076}
3077
3078impl From<Header> for api::stream::proto::Header {
3079 fn from(value: Header) -> Self {
3080 Self {
3081 name: value.name,
3082 value: value.value,
3083 }
3084 }
3085}
3086
3087impl From<api::stream::proto::Header> for Header {
3088 fn from(value: api::stream::proto::Header) -> Self {
3089 Self {
3090 name: value.name,
3091 value: value.value,
3092 }
3093 }
3094}
3095
3096#[derive(Debug, Clone, PartialEq)]
3097pub struct AppendRecord {
3099 body: Bytes,
3100 headers: Vec<Header>,
3101 timestamp: Option<u64>,
3102}
3103
3104impl AppendRecord {
3105 fn validate(self) -> Result<Self, ValidationError> {
3106 if self.metered_bytes() > RECORD_BATCH_MAX.bytes {
3107 Err(ValidationError(format!(
3108 "metered_bytes: {} exceeds {}",
3109 self.metered_bytes(),
3110 RECORD_BATCH_MAX.bytes
3111 )))
3112 } else {
3113 Ok(self)
3114 }
3115 }
3116
3117 pub fn new(body: impl Into<Bytes>) -> Result<Self, ValidationError> {
3119 let record = Self {
3120 body: body.into(),
3121 headers: Vec::default(),
3122 timestamp: None,
3123 };
3124 record.validate()
3125 }
3126
3127 pub fn with_headers(
3129 self,
3130 headers: impl IntoIterator<Item = Header>,
3131 ) -> Result<Self, ValidationError> {
3132 let record = Self {
3133 headers: headers.into_iter().collect(),
3134 ..self
3135 };
3136 record.validate()
3137 }
3138
3139 pub fn with_timestamp(self, timestamp: u64) -> Self {
3143 Self {
3144 timestamp: Some(timestamp),
3145 ..self
3146 }
3147 }
3148
3149 pub fn body(&self) -> &[u8] {
3151 &self.body
3152 }
3153
3154 pub fn headers(&self) -> &[Header] {
3156 &self.headers
3157 }
3158
3159 pub fn timestamp(&self) -> Option<u64> {
3161 self.timestamp
3162 }
3163}
3164
3165impl From<AppendRecord> for api::stream::proto::AppendRecord {
3166 fn from(value: AppendRecord) -> Self {
3167 Self {
3168 timestamp: value.timestamp,
3169 headers: value.headers.into_iter().map(Into::into).collect(),
3170 body: value.body,
3171 }
3172 }
3173}
3174
3175pub trait MeteredBytes {
3182 fn metered_bytes(&self) -> usize;
3184}
3185
3186macro_rules! metered_bytes_impl {
3187 ($ty:ty) => {
3188 impl MeteredBytes for $ty {
3189 fn metered_bytes(&self) -> usize {
3190 8 + (2 * self.headers.len())
3191 + self
3192 .headers
3193 .iter()
3194 .map(|h| h.name.len() + h.value.len())
3195 .sum::<usize>()
3196 + self.body.len()
3197 }
3198 }
3199 };
3200}
3201
3202metered_bytes_impl!(AppendRecord);
3203
3204impl MeteredSize for AppendRecord {
3205 fn metered_size(&self) -> usize {
3206 self.metered_bytes()
3207 }
3208}
3209
3210#[derive(Debug, Clone)]
3211pub struct AppendRecordBatch(Metered<Vec<AppendRecord>>);
3220
3221impl From<Metered<Vec<AppendRecord>>> for AppendRecordBatch {
3222 fn from(records: Metered<Vec<AppendRecord>>) -> Self {
3223 Self(records)
3224 }
3225}
3226
3227impl AppendRecordBatch {
3228 pub fn try_from_iter<I>(iter: I) -> Result<Self, ValidationError>
3230 where
3231 I: IntoIterator<Item = AppendRecord>,
3232 {
3233 let mut records = Metered::with_capacity(RECORD_BATCH_MAX.count);
3234
3235 for record in iter {
3236 records.push(Metered::from(record));
3237
3238 if records.metered_size() > RECORD_BATCH_MAX.bytes {
3239 return Err(ValidationError(format!(
3240 "batch size in metered bytes ({}) exceeds {}",
3241 records.metered_size(),
3242 RECORD_BATCH_MAX.bytes
3243 )));
3244 }
3245
3246 if records.len() > RECORD_BATCH_MAX.count {
3247 return Err(ValidationError(format!(
3248 "number of records in the batch exceeds {}",
3249 RECORD_BATCH_MAX.count
3250 )));
3251 }
3252 }
3253
3254 if records.is_empty() {
3255 return Err(ValidationError("batch is empty".into()));
3256 }
3257
3258 Ok(records.into())
3259 }
3260}
3261
3262impl Deref for AppendRecordBatch {
3263 type Target = [AppendRecord];
3264
3265 fn deref(&self) -> &Self::Target {
3266 &self.0[..]
3267 }
3268}
3269
3270impl MeteredBytes for AppendRecordBatch {
3271 fn metered_bytes(&self) -> usize {
3272 self.0.metered_size()
3273 }
3274}
3275
3276impl IntoIterator for AppendRecordBatch {
3277 type Item = AppendRecord;
3278 type IntoIter = std::vec::IntoIter<AppendRecord>;
3279
3280 fn into_iter(self) -> Self::IntoIter {
3281 self.0.into_iter()
3282 }
3283}
3284
3285impl<'a> IntoIterator for &'a AppendRecordBatch {
3286 type Item = &'a AppendRecord;
3287 type IntoIter = std::slice::Iter<'a, AppendRecord>;
3288
3289 fn into_iter(self) -> Self::IntoIter {
3290 self.0.iter()
3291 }
3292}
3293
3294#[derive(Debug, Clone)]
3295pub enum Command {
3297 Fence {
3299 fencing_token: FencingToken,
3301 },
3302 Trim {
3304 trim_point: u64,
3306 },
3307}
3308
3309#[derive(Debug, Clone)]
3310#[non_exhaustive]
3311pub struct CommandRecord {
3315 pub command: Command,
3317 pub timestamp: Option<u64>,
3319}
3320
3321impl CommandRecord {
3322 const FENCE: &[u8] = b"fence";
3323 const TRIM: &[u8] = b"trim";
3324
3325 pub fn fence(fencing_token: FencingToken) -> Self {
3330 Self {
3331 command: Command::Fence { fencing_token },
3332 timestamp: None,
3333 }
3334 }
3335
3336 pub fn trim(trim_point: u64) -> Self {
3343 Self {
3344 command: Command::Trim { trim_point },
3345 timestamp: None,
3346 }
3347 }
3348
3349 pub fn with_timestamp(self, timestamp: u64) -> Self {
3351 Self {
3352 timestamp: Some(timestamp),
3353 ..self
3354 }
3355 }
3356}
3357
3358impl From<CommandRecord> for AppendRecord {
3359 fn from(value: CommandRecord) -> Self {
3360 let (header_value, body) = match value.command {
3361 Command::Fence { fencing_token } => (
3362 CommandRecord::FENCE,
3363 Bytes::copy_from_slice(fencing_token.as_bytes()),
3364 ),
3365 Command::Trim { trim_point } => (
3366 CommandRecord::TRIM,
3367 Bytes::copy_from_slice(&trim_point.to_be_bytes()),
3368 ),
3369 };
3370 Self {
3371 body,
3372 headers: vec![Header::new("", header_value)],
3373 timestamp: value.timestamp,
3374 }
3375 }
3376}
3377
3378#[derive(Debug, Clone)]
3379#[non_exhaustive]
3380pub struct AppendInput {
3383 pub records: AppendRecordBatch,
3385 pub match_seq_num: Option<u64>,
3389 pub fencing_token: Option<FencingToken>,
3394 pub stream_config: Option<StreamConfig>,
3404}
3405
3406impl AppendInput {
3407 pub fn new(records: AppendRecordBatch) -> Self {
3409 Self {
3410 records,
3411 match_seq_num: None,
3412 fencing_token: None,
3413 stream_config: None,
3414 }
3415 }
3416
3417 pub fn with_stream_config(self, stream_config: StreamConfig) -> Self {
3419 Self {
3420 stream_config: Some(stream_config),
3421 ..self
3422 }
3423 }
3424
3425 pub fn with_match_seq_num(self, match_seq_num: u64) -> Self {
3427 Self {
3428 match_seq_num: Some(match_seq_num),
3429 ..self
3430 }
3431 }
3432
3433 pub fn with_fencing_token(self, fencing_token: FencingToken) -> Self {
3435 Self {
3436 fencing_token: Some(fencing_token),
3437 ..self
3438 }
3439 }
3440}
3441
3442impl From<AppendInput> for api::stream::proto::AppendInput {
3443 fn from(value: AppendInput) -> Self {
3444 Self {
3445 records: value.records.iter().cloned().map(Into::into).collect(),
3446 match_seq_num: value.match_seq_num,
3447 fencing_token: value.fencing_token.map(|t| t.to_string()),
3448 }
3449 }
3450}
3451
3452#[derive(Debug, Clone, PartialEq)]
3453#[non_exhaustive]
3454pub struct AppendAck {
3456 pub start: StreamPosition,
3458 pub end: StreamPosition,
3464 pub tail: StreamPosition,
3469}
3470
3471impl AppendAck {
3472 pub fn new(start: StreamPosition, end: StreamPosition, tail: StreamPosition) -> Self {
3476 Self { start, end, tail }
3477 }
3478}
3479
3480impl From<api::stream::proto::AppendAck> for AppendAck {
3481 fn from(value: api::stream::proto::AppendAck) -> Self {
3482 Self {
3483 start: value.start.unwrap_or_default().into(),
3484 end: value.end.unwrap_or_default().into(),
3485 tail: value.tail.unwrap_or_default().into(),
3486 }
3487 }
3488}
3489
3490#[derive(Debug, Clone, Copy)]
3491pub enum ReadFrom {
3493 SeqNum(u64),
3495 Timestamp(u64),
3497 TailOffset(u64),
3499}
3500
3501impl Default for ReadFrom {
3502 fn default() -> Self {
3503 Self::SeqNum(0)
3504 }
3505}
3506
3507#[derive(Debug, Default, Clone)]
3508#[non_exhaustive]
3509pub struct ReadStart {
3511 pub from: ReadFrom,
3515 pub clamp_to_tail: bool,
3519}
3520
3521impl ReadStart {
3522 pub fn new() -> Self {
3524 Self::default()
3525 }
3526
3527 pub fn with_from(self, from: ReadFrom) -> Self {
3529 Self { from, ..self }
3530 }
3531
3532 pub fn with_clamp_to_tail(self, clamp_to_tail: bool) -> Self {
3534 Self {
3535 clamp_to_tail,
3536 ..self
3537 }
3538 }
3539}
3540
3541impl From<ReadStart> for api::stream::ReadStart {
3542 fn from(value: ReadStart) -> Self {
3543 let (seq_num, timestamp, tail_offset) = match value.from {
3544 ReadFrom::SeqNum(n) => (Some(n), None, None),
3545 ReadFrom::Timestamp(t) => (None, Some(t), None),
3546 ReadFrom::TailOffset(o) => (None, None, Some(o)),
3547 };
3548 Self {
3549 seq_num,
3550 timestamp,
3551 tail_offset,
3552 clamp: if value.clamp_to_tail {
3553 Some(true)
3554 } else {
3555 None
3556 },
3557 }
3558 }
3559}
3560
3561#[derive(Debug, Clone, Default)]
3562#[non_exhaustive]
3563pub struct ReadLimits {
3565 pub count: Option<usize>,
3569 pub bytes: Option<usize>,
3573}
3574
3575impl ReadLimits {
3576 pub fn new() -> Self {
3578 Self::default()
3579 }
3580
3581 pub fn with_count(self, count: usize) -> Self {
3583 Self {
3584 count: Some(count),
3585 ..self
3586 }
3587 }
3588
3589 pub fn with_bytes(self, bytes: usize) -> Self {
3591 Self {
3592 bytes: Some(bytes),
3593 ..self
3594 }
3595 }
3596}
3597
3598#[derive(Debug, Clone, Default)]
3599#[non_exhaustive]
3600pub struct ReadStop {
3602 pub limits: ReadLimits,
3606 pub until: Option<RangeTo<u64>>,
3610 pub wait: Option<u32>,
3620}
3621
3622impl ReadStop {
3623 pub fn new() -> Self {
3625 Self::default()
3626 }
3627
3628 pub fn with_limits(self, limits: ReadLimits) -> Self {
3630 Self { limits, ..self }
3631 }
3632
3633 pub fn with_until(self, until: RangeTo<u64>) -> Self {
3635 Self {
3636 until: Some(until),
3637 ..self
3638 }
3639 }
3640
3641 pub fn with_wait(self, wait: u32) -> Self {
3643 Self {
3644 wait: Some(wait),
3645 ..self
3646 }
3647 }
3648}
3649
3650impl From<ReadStop> for api::stream::ReadEnd {
3651 fn from(value: ReadStop) -> Self {
3652 Self {
3653 count: value.limits.count,
3654 bytes: value.limits.bytes,
3655 until: value.until.map(|r| r.end),
3656 wait: value.wait,
3657 }
3658 }
3659}
3660
3661#[derive(Debug, Clone, Default)]
3662#[non_exhaustive]
3663pub struct ReadInput {
3666 pub start: ReadStart,
3670 pub stop: ReadStop,
3674 pub ignore_command_records: bool,
3678 pub stream_config: Option<StreamConfig>,
3683}
3684
3685#[derive(Debug, Clone, Copy, Default, Eq, PartialEq)]
3686#[non_exhaustive]
3687pub enum ReadSessionRetryPolicy {
3689 #[default]
3691 Budgeted,
3692 Indefinite,
3698}
3699
3700#[derive(Debug, Clone, Default)]
3701#[non_exhaustive]
3702pub struct ReadSessionConfig {
3704 pub retry_policy: ReadSessionRetryPolicy,
3710}
3711
3712impl ReadSessionConfig {
3713 pub fn new() -> Self {
3715 Self::default()
3716 }
3717
3718 pub fn with_retry_policy(self, retry_policy: ReadSessionRetryPolicy) -> Self {
3720 Self {
3721 retry_policy,
3722 ..self
3723 }
3724 }
3725}
3726
3727impl ReadInput {
3728 pub fn new() -> Self {
3730 Self::default()
3731 }
3732
3733 pub fn with_start(self, start: ReadStart) -> Self {
3735 Self { start, ..self }
3736 }
3737
3738 pub fn with_stop(self, stop: ReadStop) -> Self {
3740 Self { stop, ..self }
3741 }
3742
3743 pub fn with_ignore_command_records(self, ignore_command_records: bool) -> Self {
3745 Self {
3746 ignore_command_records,
3747 ..self
3748 }
3749 }
3750
3751 pub fn with_stream_config(self, stream_config: StreamConfig) -> Self {
3753 Self {
3754 stream_config: Some(stream_config),
3755 ..self
3756 }
3757 }
3758}
3759
3760#[derive(Debug, Clone)]
3761#[non_exhaustive]
3762pub struct SequencedRecord {
3764 pub seq_num: u64,
3766 pub body: Bytes,
3768 pub headers: Vec<Header>,
3770 pub timestamp: u64,
3772}
3773
3774impl SequencedRecord {
3775 pub fn from_parts(
3779 seq_num: u64,
3780 timestamp: u64,
3781 headers: Vec<Header>,
3782 body: impl Into<Bytes>,
3783 ) -> Self {
3784 Self {
3785 seq_num,
3786 timestamp,
3787 body: body.into(),
3788 headers,
3789 }
3790 }
3791
3792 pub fn is_command_record(&self) -> bool {
3794 self.headers.len() == 1 && *self.headers[0].name == *b""
3795 }
3796}
3797
3798impl From<api::stream::proto::SequencedRecord> for SequencedRecord {
3799 fn from(value: api::stream::proto::SequencedRecord) -> Self {
3800 Self {
3801 seq_num: value.seq_num,
3802 body: value.body,
3803 headers: value.headers.into_iter().map(Into::into).collect(),
3804 timestamp: value.timestamp,
3805 }
3806 }
3807}
3808
3809metered_bytes_impl!(SequencedRecord);
3810
3811#[derive(Debug, Clone)]
3812#[non_exhaustive]
3813pub struct ReadBatch {
3816 pub records: Vec<SequencedRecord>,
3823 pub tail: Option<StreamPosition>,
3828}
3829
3830impl ReadBatch {
3831 pub fn new(records: Vec<SequencedRecord>, tail: Option<StreamPosition>) -> Self {
3835 Self { records, tail }
3836 }
3837
3838 pub(crate) fn from_api(batch: api::stream::proto::ReadBatch) -> Self {
3839 Self {
3840 records: batch.records.into_iter().map(Into::into).collect(),
3841 tail: batch.tail.map(Into::into),
3842 }
3843 }
3844}
3845
3846pub type Streaming<T> = Pin<Box<dyn Send + futures_core::Stream<Item = Result<T, RequestError>>>>;
3848
3849fn idempotency_token() -> String {
3850 uuid::Uuid::new_v4().simple().to_string()
3851}
3852
3853#[cfg(test)]
3854mod tests {
3855 use proptest::prelude::*;
3856 use rstest::rstest;
3857
3858 use super::*;
3859
3860 type HeaderParts = (Vec<u8>, Vec<u8>);
3861 type AppendRecordParts = (Vec<u8>, Vec<HeaderParts>);
3862
3863 fn byte_vec_strategy(max_len: usize) -> impl Strategy<Value = Vec<u8>> {
3864 prop::collection::vec(any::<u8>(), 0..=max_len)
3865 }
3866
3867 fn header_parts_strategy() -> impl Strategy<Value = HeaderParts> {
3868 (byte_vec_strategy(32), byte_vec_strategy(64))
3869 }
3870
3871 fn string_strategy(max_chars: usize) -> impl Strategy<Value = String> {
3872 prop::collection::vec(any::<char>(), 0..=max_chars)
3873 .prop_map(|chars| chars.into_iter().collect())
3874 }
3875
3876 fn read_from_strategy() -> impl Strategy<Value = ReadFrom> {
3877 prop_oneof![
3878 any::<u64>().prop_map(ReadFrom::SeqNum),
3879 any::<u64>().prop_map(ReadFrom::Timestamp),
3880 any::<u64>().prop_map(ReadFrom::TailOffset),
3881 ]
3882 }
3883
3884 fn append_record_parts_strategy() -> impl Strategy<Value = AppendRecordParts> {
3885 (
3886 byte_vec_strategy(256),
3887 prop::collection::vec(header_parts_strategy(), 0..=16),
3888 )
3889 }
3890
3891 fn proto_stream_position_strategy() -> impl Strategy<Value = api::stream::proto::StreamPosition>
3892 {
3893 (any::<u64>(), any::<u64>()).prop_map(|(seq_num, timestamp)| {
3894 api::stream::proto::StreamPosition { seq_num, timestamp }
3895 })
3896 }
3897
3898 fn headers_from_parts(headers: &[HeaderParts]) -> Vec<Header> {
3899 headers
3900 .iter()
3901 .map(|(name, value)| Header::new(Bytes::from(name.clone()), Bytes::from(value.clone())))
3902 .collect()
3903 }
3904
3905 fn expected_metered_bytes(body: &[u8], headers: &[HeaderParts]) -> usize {
3906 8 + (2 * headers.len())
3907 + headers
3908 .iter()
3909 .map(|(name, value)| name.len() + value.len())
3910 .sum::<usize>()
3911 + body.len()
3912 }
3913
3914 #[test]
3917 fn s2_datetime_parse_valid_rfc3339() {
3918 let dt: S2DateTime = "2024-01-15T12:30:00Z".parse().unwrap();
3919 assert_eq!(dt.to_string(), "2024-01-15T12:30:00Z");
3920 }
3921
3922 #[test]
3923 fn s2_datetime_parse_with_offset() {
3924 let dt: S2DateTime = "2024-06-01T08:00:00+05:30".parse().unwrap();
3925 assert_eq!(dt.to_string(), "2024-06-01T08:00:00+05:30");
3926
3927 let offset_dt: time::OffsetDateTime = dt.into();
3928 assert_eq!(
3929 offset_dt.offset(),
3930 time::UtcOffset::from_hms(5, 30, 0).unwrap()
3931 );
3932 }
3933
3934 #[test]
3935 fn s2_datetime_parse_invalid() {
3936 let err = "not-a-date".parse::<S2DateTime>();
3937 assert!(err.is_err());
3938 }
3939
3940 #[test]
3941 fn s2_datetime_roundtrip_via_offset_datetime() {
3942 let odt = time::OffsetDateTime::from_unix_timestamp(1_700_000_000).unwrap();
3943 let dt = S2DateTime::try_from(odt).unwrap();
3944 let back: time::OffsetDateTime = dt.into();
3945 assert_eq!(odt, back);
3946 }
3947
3948 #[rstest]
3951 #[case::https_with_scheme("https://aws.s2.dev", Scheme::HTTPS)]
3952 #[case::http_with_scheme("http://localhost:8080", Scheme::HTTP)]
3953 #[case::default_https("aws.s2.dev", Scheme::HTTPS)]
3954 fn account_endpoint_parse(#[case] input: &str, #[case] expected_scheme: Scheme) {
3955 let ep: AccountEndpoint = input.parse().unwrap();
3956 assert_eq!(ep.scheme, expected_scheme);
3957 }
3958
3959 #[rstest]
3962 #[case::https_parent_zone("https://{basin}.b.s2.dev", Scheme::HTTPS, true)]
3963 #[case::http_direct("http://localhost:8080", Scheme::HTTP, false)]
3964 #[case::default_https_parent_zone("{basin}.b.s2.dev", Scheme::HTTPS, true)]
3965 fn basin_endpoint_parse(
3966 #[case] input: &str,
3967 #[case] expected_scheme: Scheme,
3968 #[case] expected_parent_zone: bool,
3969 ) {
3970 let ep: BasinEndpoint = input.parse().unwrap();
3971 assert_eq!(ep.scheme, expected_scheme);
3972 assert_eq!(
3973 matches!(ep.authority, BasinAuthority::ParentZone(_)),
3974 expected_parent_zone
3975 );
3976 }
3977
3978 #[test]
3981 fn s2_endpoints_new_requires_same_scheme() {
3982 let account: AccountEndpoint = "https://aws.s2.dev".parse().unwrap();
3983 let basin: BasinEndpoint = "http://localhost:8080".parse().unwrap();
3984 let err = S2Endpoints::new(account, basin);
3985 assert!(err.is_err());
3986 }
3987
3988 #[test]
3989 fn s2_endpoints_new_same_scheme_succeeds() {
3990 let account: AccountEndpoint = "https://aws.s2.dev".parse().unwrap();
3991 let basin: BasinEndpoint = "https://{basin}.b.s2.dev".parse().unwrap();
3992 let ep = S2Endpoints::new(account, basin).unwrap();
3993 assert_eq!(ep.scheme, Scheme::HTTPS);
3994 }
3995
3996 #[test]
3997 fn s2_endpoints_for_endpoint_defaults_to_https() {
3998 let ep = S2Endpoints::for_endpoint("localhost:8080").unwrap();
3999 let authority: Authority = "localhost:8080".parse().unwrap();
4000 assert_eq!(ep.scheme, Scheme::HTTPS);
4001 assert_eq!(ep.account_authority, authority);
4002 assert_eq!(ep.basin_authority, BasinAuthority::Direct(authority));
4003 }
4004
4005 #[test]
4006 fn s2_endpoints_for_endpoint_accepts_explicit_scheme() {
4007 let ep = S2Endpoints::for_endpoint("http://localhost:8080").unwrap();
4008 assert_eq!(ep.scheme, Scheme::HTTP);
4009 }
4010
4011 #[test]
4012 fn s2_endpoints_for_endpoint_rejects_invalid_endpoint() {
4013 assert!(S2Endpoints::for_endpoint("not a valid endpoint").is_err());
4014 }
4015
4016 #[rstest]
4019 #[case::none(Compression::None, CompressionAlgorithm::None)]
4020 #[case::gzip(Compression::Gzip, CompressionAlgorithm::Gzip)]
4021 #[case::zstd(Compression::Zstd, CompressionAlgorithm::Zstd)]
4022 fn compression_conversion(#[case] sdk: Compression, #[case] api: CompressionAlgorithm) {
4023 assert_eq!(CompressionAlgorithm::from(sdk), api);
4024 }
4025
4026 #[test]
4029 fn retry_config_defaults() {
4030 let rc = RetryConfig::default();
4031 assert_eq!(rc.max_attempts.get(), 3);
4032 assert_eq!(rc.min_base_delay, Duration::from_millis(100));
4033 assert_eq!(rc.max_base_delay, Duration::from_secs(1));
4034 assert!(matches!(rc.append_retry_policy, AppendRetryPolicy::All));
4035 }
4036
4037 #[test]
4038 fn retry_config_max_retries() {
4039 let rc = RetryConfig::default();
4040 assert_eq!(rc.max_retries(), 2);
4041 }
4042
4043 #[test]
4046 fn s2_config_defaults() {
4047 let cfg = S2Config::new("test-token");
4048 assert_eq!(cfg.connection_timeout, Duration::from_secs(3));
4049 assert_eq!(cfg.request_timeout, Duration::from_secs(5));
4050 assert!(!cfg.insecure_skip_cert_verification);
4051 }
4052
4053 #[cfg(feature = "_hidden")]
4054 #[rstest]
4055 #[case::matching_compression("content-encoding", "gzip", Compression::Gzip)]
4056 #[case::mixed_case("Content-Encoding", "identity", Compression::None)]
4057 #[case::empty_value("content-encoding", "", Compression::None)]
4058 fn default_headers_reject_content_encoding(
4059 #[case] name: &str,
4060 #[case] value: &str,
4061 #[case] compression: Compression,
4062 ) {
4063 let headers = HeaderMap::from_iter([(
4064 name.parse::<http::header::HeaderName>().unwrap(),
4065 HeaderValue::from_str(value).unwrap(),
4066 )]);
4067 let error = S2Config::new("token")
4068 .with_compression(compression)
4069 .with_default_headers(headers)
4070 .unwrap_err();
4071 assert!(error.0.contains("Content-Encoding"));
4072 assert!(error.0.contains("with_compression"));
4073 }
4074
4075 #[cfg(feature = "_hidden")]
4076 #[rstest]
4077 #[case::content_type_s2s("content-type", "s2s/proto")]
4078 #[case::content_type_protobuf("content-type", "application/protobuf")]
4079 #[case::content_type_json("content-type", "application/json")]
4080 #[case::content_type_mixed_case("Content-Type", "s2s/proto")]
4081 #[case::content_type_empty("content-type", "")]
4082 #[case::content_length("content-length", "123")]
4083 #[case::content_length_mixed_case("Content-Length", "0")]
4084 #[case::content_length_empty("content-length", "")]
4085 #[case::transfer_encoding("transfer-encoding", "chunked")]
4086 #[case::transfer_encoding_mixed_case("Transfer-Encoding", "chunked")]
4087 #[case::transfer_encoding_empty("transfer-encoding", "")]
4088 fn default_headers_reject_framing_headers(#[case] name: &str, #[case] value: &str) {
4089 let headers = HeaderMap::from_iter([(
4090 name.parse::<http::header::HeaderName>().unwrap(),
4091 HeaderValue::from_str(value).unwrap(),
4092 )]);
4093 let error = S2Config::new("token")
4094 .with_default_headers(headers)
4095 .unwrap_err();
4096 assert!(error.0.contains(&name.to_ascii_lowercase()));
4097 assert!(error.0.contains("framing"));
4098 }
4099
4100 #[rstest]
4103 #[case::age(RetentionPolicy::Age(3600))]
4104 #[case::infinite(RetentionPolicy::Infinite)]
4105 fn retention_policy_roundtrip(#[case] sdk: RetentionPolicy) {
4106 let api: api::config::RetentionPolicy = sdk.into();
4107 let back: RetentionPolicy = api.into();
4108 assert_eq!(back, sdk);
4109 }
4110
4111 #[rstest]
4114 #[case::client_prefer(
4115 TimestampingMode::ClientPrefer,
4116 api::config::TimestampingMode::ClientPrefer
4117 )]
4118 #[case::client_require(
4119 TimestampingMode::ClientRequire,
4120 api::config::TimestampingMode::ClientRequire
4121 )]
4122 #[case::arrival(TimestampingMode::Arrival, api::config::TimestampingMode::Arrival)]
4123 fn timestamping_mode_roundtrip(
4124 #[case] sdk: TimestampingMode,
4125 #[case] expected_api: api::config::TimestampingMode,
4126 ) {
4127 let converted: api::config::TimestampingMode = sdk.into();
4128 assert_eq!(converted, expected_api);
4129 let back: TimestampingMode = converted.into();
4130 assert_eq!(back, sdk);
4131 }
4132
4133 #[test]
4136 fn timestamping_config_roundtrip() {
4137 let sdk = TimestampingConfig {
4138 mode: Some(TimestampingMode::Arrival),
4139 uncapped: Some(true),
4140 };
4141 let api: api::config::TimestampingConfig = sdk.into();
4142 let back: TimestampingConfig = api.into();
4143 assert_eq!(back, sdk);
4144 }
4145
4146 #[test]
4149 fn delete_on_empty_config_roundtrip() {
4150 let sdk = DeleteOnEmptyConfig::new().with_min_age(Duration::from_secs(300));
4151 let api: api::config::DeleteOnEmptyConfig = sdk.into();
4152 let back: DeleteOnEmptyConfig = api.into();
4153 assert_eq!(back, sdk);
4154 }
4155
4156 #[test]
4159 fn stream_config_builder_and_roundtrip() {
4160 let sdk = StreamConfig::new()
4161 .with_storage_class("express")
4162 .with_retention_policy(RetentionPolicy::Age(86400))
4163 .with_timestamping(TimestampingConfig {
4164 mode: Some(TimestampingMode::ClientPrefer),
4165 uncapped: None,
4166 })
4167 .with_delete_on_empty(DeleteOnEmptyConfig { min_age_secs: 60 });
4168 let api: api::config::StreamConfig = sdk.clone().into();
4169 let back: StreamConfig = api.into();
4170 assert_eq!(back, sdk);
4171 }
4172
4173 #[test]
4176 fn basin_config_builder_and_roundtrip() {
4177 let sdk = BasinConfig::new()
4178 .with_default_stream_config(StreamConfig::new().with_storage_class("standard"))
4179 .with_create_stream_on_append(true)
4180 .with_create_stream_on_read(false);
4181 let api: api::config::BasinConfig = sdk.clone().into();
4182 let back: BasinConfig = api.into();
4183 assert_eq!(back, sdk);
4184 }
4185
4186 proptest! {
4189 #[test]
4190 fn fencing_token_parse_accepts_only_within_byte_limit(
4191 token in string_strategy(MAX_FENCING_TOKEN_LENGTH + 8),
4192 ) {
4193 let parsed = token.parse::<FencingToken>();
4194
4195 if token.len() <= MAX_FENCING_TOKEN_LENGTH {
4196 prop_assert_eq!(parsed.unwrap().to_string(), token);
4197 } else {
4198 prop_assert!(parsed.is_err());
4199 }
4200 }
4201 }
4202
4203 #[test]
4206 fn stream_position_display() {
4207 let pos = StreamPosition {
4208 seq_num: 42,
4209 timestamp: 1700000000,
4210 };
4211 assert_eq!(pos.to_string(), "seq_num=42, timestamp=1700000000");
4212 }
4213
4214 proptest! {
4215 #[test]
4216 fn stream_position_conversions_preserve_values(seq_num in any::<u64>(), timestamp in any::<u64>()) {
4217 let proto: StreamPosition = api::stream::proto::StreamPosition {
4218 seq_num,
4219 timestamp,
4220 }
4221 .into();
4222 prop_assert_eq!(proto.seq_num, seq_num);
4223 prop_assert_eq!(proto.timestamp, timestamp);
4224
4225 let api: StreamPosition = api::stream::StreamPosition {
4226 seq_num,
4227 timestamp,
4228 }
4229 .into();
4230 prop_assert_eq!(api.seq_num, seq_num);
4231 prop_assert_eq!(api.timestamp, timestamp);
4232 }
4233 }
4234
4235 proptest! {
4238 #[test]
4239 fn header_proto_roundtrip_preserves_binary_parts(
4240 name in byte_vec_strategy(64),
4241 value in byte_vec_strategy(128),
4242 ) {
4243 let header = Header::new(Bytes::from(name.clone()), Bytes::from(value.clone()));
4244 let proto: api::stream::proto::Header = header.into();
4245 let back: Header = proto.into();
4246
4247 prop_assert_eq!(back.name.as_ref(), name.as_slice());
4248 prop_assert_eq!(back.value.as_ref(), value.as_slice());
4249 }
4250 }
4251
4252 #[test]
4255 fn append_record_too_large() {
4256 let big_body = vec![0u8; RECORD_BATCH_MAX.bytes + 1];
4257 assert!(AppendRecord::new(big_body).is_err());
4258 }
4259
4260 proptest! {
4263 #[test]
4264 fn append_record_preserves_fields_and_metered_byte_formula(
4265 (body, headers) in append_record_parts_strategy(),
4266 timestamp in proptest::option::of(any::<u64>()),
4267 ) {
4268 let mut record = AppendRecord::new(body.clone())
4269 .unwrap()
4270 .with_headers(headers_from_parts(&headers))
4271 .unwrap();
4272 if let Some(timestamp) = timestamp {
4273 record = record.with_timestamp(timestamp);
4274 }
4275
4276 prop_assert_eq!(record.body(), body.as_slice());
4277 prop_assert_eq!(record.headers().len(), headers.len());
4278 prop_assert_eq!(record.timestamp(), timestamp);
4279 prop_assert_eq!(record.metered_bytes(), expected_metered_bytes(&body, &headers));
4280
4281 for (actual, (expected_name, expected_value)) in record.headers().iter().zip(headers.iter()) {
4282 prop_assert_eq!(actual.name.as_ref(), expected_name.as_slice());
4283 prop_assert_eq!(actual.value.as_ref(), expected_value.as_slice());
4284 }
4285 }
4286 }
4287
4288 #[test]
4291 fn append_record_batch_empty_is_err() {
4292 let result = AppendRecordBatch::try_from_iter(vec![]);
4293 assert!(result.is_err());
4294 }
4295
4296 #[test]
4297 fn append_record_batch_too_many_records() {
4298 let records: Vec<_> = (0..1001).map(|_| AppendRecord::new("x").unwrap()).collect();
4299 let result = AppendRecordBatch::try_from_iter(records);
4300 assert!(result.is_err());
4301 }
4302
4303 proptest! {
4304 #[test]
4305 fn append_record_batch_metered_bytes_is_sum_of_records(
4306 records in prop::collection::vec(append_record_parts_strategy(), 1..=32),
4307 ) {
4308 let expected = records
4309 .iter()
4310 .map(|(body, headers)| expected_metered_bytes(body, headers))
4311 .sum::<usize>();
4312 let records = records
4313 .into_iter()
4314 .map(|(body, headers)| {
4315 AppendRecord::new(body)
4316 .unwrap()
4317 .with_headers(headers_from_parts(&headers))
4318 .unwrap()
4319 })
4320 .collect::<Vec<_>>();
4321
4322 let batch = AppendRecordBatch::try_from_iter(records).unwrap();
4323 prop_assert_eq!(batch.metered_bytes(), expected);
4324 prop_assert_eq!(batch.iter().map(MeteredBytes::metered_bytes).sum::<usize>(), expected);
4325 }
4326 }
4327
4328 #[test]
4331 fn command_record_fence() {
4332 let token: FencingToken = "tok".parse().unwrap();
4333 let cmd = CommandRecord::fence(token);
4334 let record: AppendRecord = cmd.into();
4335 assert_eq!(record.headers().len(), 1);
4336 assert_eq!(record.headers()[0].name.as_ref(), b"");
4337 assert_eq!(record.headers()[0].value.as_ref(), b"fence");
4338 assert_eq!(record.body(), b"tok");
4339 }
4340
4341 #[test]
4342 fn command_record_trim() {
4343 let cmd = CommandRecord::trim(42);
4344 let record: AppendRecord = cmd.into();
4345 assert_eq!(record.headers().len(), 1);
4346 assert_eq!(record.headers()[0].value.as_ref(), b"trim");
4347 assert_eq!(record.body(), &42u64.to_be_bytes());
4348 }
4349
4350 #[rstest]
4353 #[case::command(vec![Header::new("", "fence")], true)]
4354 #[case::regular(vec![Header::new("key", "value")], false)]
4355 #[case::no_headers(vec![], false)]
4356 fn sequenced_record_command_detection(#[case] headers: Vec<Header>, #[case] expected: bool) {
4357 let record = SequencedRecord {
4358 seq_num: 0,
4359 body: Bytes::from("data"),
4360 headers,
4361 timestamp: 0,
4362 };
4363 assert_eq!(record.is_command_record(), expected);
4364 }
4365
4366 proptest! {
4369 #[test]
4370 fn read_start_to_api_sets_only_selected_position_field(
4371 from in read_from_strategy(),
4372 clamp_to_tail in any::<bool>(),
4373 ) {
4374 let (seq_num, timestamp, tail_offset) = match from {
4375 ReadFrom::SeqNum(value) => (Some(value), None, None),
4376 ReadFrom::Timestamp(value) => (None, Some(value), None),
4377 ReadFrom::TailOffset(value) => (None, None, Some(value)),
4378 };
4379 let api: api::stream::ReadStart = ReadStart::new()
4380 .with_from(from)
4381 .with_clamp_to_tail(clamp_to_tail)
4382 .into();
4383
4384 prop_assert_eq!(api.seq_num, seq_num);
4385 prop_assert_eq!(api.timestamp, timestamp);
4386 prop_assert_eq!(api.tail_offset, tail_offset);
4387 prop_assert_eq!(api.clamp, clamp_to_tail.then_some(true));
4388 }
4389 }
4390
4391 #[test]
4394 fn read_stop_to_api() {
4395 let stop = ReadStop::new()
4396 .with_limits(ReadLimits::new().with_count(50))
4397 .with_until(..1000)
4398 .with_wait(30);
4399 let api: api::stream::ReadEnd = stop.into();
4400 assert_eq!(api.count, Some(50));
4401 assert_eq!(api.until, Some(1000));
4402 assert_eq!(api.wait, Some(30));
4403 }
4404
4405 #[test]
4408 fn operation_roundtrip_all_variants() {
4409 let variants = [
4410 Operation::ListBasins,
4411 Operation::CreateBasin,
4412 Operation::GetBasinConfig,
4413 Operation::DeleteBasin,
4414 Operation::ReconfigureBasin,
4415 Operation::ListAccessTokens,
4416 Operation::IssueAccessToken,
4417 Operation::RevokeAccessToken,
4418 Operation::GetAccountMetrics,
4419 Operation::GetBasinMetrics,
4420 Operation::GetStreamMetrics,
4421 Operation::ListStreams,
4422 Operation::CreateStream,
4423 Operation::GetStreamConfig,
4424 Operation::DeleteStream,
4425 Operation::ReconfigureStream,
4426 Operation::CheckTail,
4427 Operation::Append,
4428 Operation::Read,
4429 Operation::Trim,
4430 Operation::Fence,
4431 Operation::ListLocations,
4432 Operation::GetDefaultLocation,
4433 Operation::SetDefaultLocation,
4434 ];
4435 for op in variants {
4436 let api_op: api::access::Operation = op.into();
4437 let back: Operation = api_op.into();
4438 assert_eq!(back, op);
4439 }
4440 }
4441
4442 #[test]
4445 fn metric_unit_conversion() {
4446 assert_eq!(
4447 MetricUnit::from(api::metrics::MetricUnit::Bytes),
4448 MetricUnit::Bytes
4449 );
4450 assert_eq!(
4451 MetricUnit::from(api::metrics::MetricUnit::Operations),
4452 MetricUnit::Operations
4453 );
4454 }
4455
4456 proptest! {
4459 #[test]
4460 fn append_ack_from_proto_preserves_present_positions_and_defaults_missing(
4461 start in proptest::option::of(proto_stream_position_strategy()),
4462 end in proptest::option::of(proto_stream_position_strategy()),
4463 tail in proptest::option::of(proto_stream_position_strategy()),
4464 ) {
4465 let expected_start = start.unwrap_or_default();
4466 let expected_end = end.unwrap_or_default();
4467 let expected_tail = tail.unwrap_or_default();
4468 let ack: AppendAck = api::stream::proto::AppendAck { start, end, tail }.into();
4469
4470 prop_assert_eq!(ack.start.seq_num, expected_start.seq_num);
4471 prop_assert_eq!(ack.start.timestamp, expected_start.timestamp);
4472 prop_assert_eq!(ack.end.seq_num, expected_end.seq_num);
4473 prop_assert_eq!(ack.end.timestamp, expected_end.timestamp);
4474 prop_assert_eq!(ack.tail.seq_num, expected_tail.seq_num);
4475 prop_assert_eq!(ack.tail.timestamp, expected_tail.timestamp);
4476 }
4477 }
4478
4479 #[test]
4482 fn read_batch_from_api() {
4483 let proto_batch = api::stream::proto::ReadBatch {
4484 records: vec![api::stream::proto::SequencedRecord {
4485 seq_num: 0,
4486 body: Bytes::from("hi"),
4487 headers: vec![api::stream::proto::Header {
4488 name: Bytes::from("k"),
4489 value: Bytes::from("v"),
4490 }],
4491 timestamp: 42,
4492 }],
4493 tail: Some(api::stream::proto::StreamPosition {
4494 seq_num: 1,
4495 timestamp: 42,
4496 }),
4497 };
4498 let batch = ReadBatch::from_api(proto_batch);
4499 assert_eq!(batch.records.len(), 1);
4500 assert_eq!(batch.records[0].seq_num, 0);
4501 assert_eq!(batch.records[0].timestamp, 42);
4502 assert_eq!(batch.records[0].body.as_ref(), b"hi");
4503 assert_eq!(batch.records[0].headers.len(), 1);
4504 assert_eq!(batch.records[0].headers[0].name.as_ref(), b"k");
4505 assert_eq!(batch.records[0].headers[0].value.as_ref(), b"v");
4506 assert_eq!(
4507 batch.tail,
4508 Some(StreamPosition {
4509 seq_num: 1,
4510 timestamp: 42,
4511 })
4512 );
4513 }
4514
4515 #[test]
4518 fn create_basin_input_to_api() {
4519 let name: BasinName = "test-basin-name".parse().unwrap();
4520 let input = CreateBasinInput::new(name.clone()).with_config(BasinConfig::new());
4521 let (req, token): (api::basin::CreateBasinRequest, String) = input.into();
4522 assert_eq!(req.basin, name);
4523 assert!(req.config.is_some());
4524 assert!(!token.is_empty());
4525 }
4526
4527 #[test]
4530 fn create_stream_input_to_api() {
4531 let name: StreamName = "my-stream".parse().unwrap();
4532 let input = CreateStreamInput::new(name.clone()).with_config(StreamConfig::new());
4533 let (req, token): (api::stream::CreateStreamRequest, String) = input.into();
4534 assert_eq!(req.stream, name);
4535 assert!(req.config.is_some());
4536 assert!(!token.is_empty());
4537 }
4538
4539 #[test]
4542 fn sequenced_record_from_proto() {
4543 let proto = api::stream::proto::SequencedRecord {
4544 seq_num: 99,
4545 body: Bytes::from("data"),
4546 headers: vec![api::stream::proto::Header {
4547 name: Bytes::from("k"),
4548 value: Bytes::from("v"),
4549 }],
4550 timestamp: 1234,
4551 };
4552 let record: SequencedRecord = proto.into();
4553 assert_eq!(record.seq_num, 99);
4554 assert_eq!(record.body.as_ref(), b"data");
4555 assert_eq!(record.headers.len(), 1);
4556 assert_eq!(record.headers[0].name.as_ref(), b"k");
4557 assert_eq!(record.headers[0].value.as_ref(), b"v");
4558 assert_eq!(record.timestamp, 1234);
4559 }
4560}