Skip to main content

s2_sdk/
types.rs

1//! Types relevant to [`S2`](crate::S2), [`S2Basin`](crate::S2Basin), and
2//! [`S2Stream`](crate::S2Stream).
3use std::{
4    collections::HashSet,
5    env::VarError,
6    fmt,
7    num::NonZeroU32,
8    ops::{Deref, RangeTo},
9    pin::Pin,
10    str::FromStr,
11    sync::Arc,
12    time::Duration,
13};
14
15use bytes::Bytes;
16use http::{
17    header::HeaderValue,
18    uri::{Authority, Scheme},
19};
20use rand::RngExt;
21use s2_api::{v1 as api, v1::stream::s2s::CompressionAlgorithm};
22/// Validation error.
23pub use s2_common::ValidationError;
24/// Access token ID.
25///
26/// **Note:** It must be unique to the account and between 1 and 96 bytes in length.
27pub use s2_common::access::AccessTokenId;
28/// See [`ListAccessTokensInput::prefix`].
29pub use s2_common::access::AccessTokenIdPrefix;
30/// See [`ListAccessTokensInput::start_after`].
31pub use s2_common::access::AccessTokenIdStartAfter;
32/// Basin name.
33///
34/// **Note:** It must be globally unique and between 8 and 48 bytes in length. It can only
35/// comprise lowercase letters, numbers, and hyphens. It cannot begin or end with a hyphen.
36pub use s2_common::basin::BasinName;
37/// See [`ListBasinsInput::prefix`].
38pub use s2_common::basin::BasinNamePrefix;
39/// See [`ListBasinsInput::start_after`].
40pub use s2_common::basin::BasinNameStartAfter;
41/// Location name.
42///
43/// **Note:** It must be between 1 and 64 characters in length and can only comprise ASCII
44/// letters, numbers, colons, hyphens, and periods.
45pub use s2_common::location::LocationName;
46/// Stream name.
47///
48/// **Note:** It must be unique to the basin and between 1 and 512 bytes in length.
49pub use s2_common::stream::StreamName;
50/// See [`ListStreamsInput::prefix`].
51pub use s2_common::stream::StreamNamePrefix;
52/// See [`ListStreamsInput::start_after`].
53pub use s2_common::stream::StreamNameStartAfter;
54pub use s2_common::{
55    caps::RECORD_BATCH_MAX,
56    encryption::{EncryptionAlgorithm, EncryptionKey},
57};
58
59pub(crate) const ONE_MIB: u32 = 1024 * 1024;
60
61use s2_common::{
62    maybe::Maybe,
63    record::{MAX_FENCING_TOKEN_LENGTH, Metered, MeteredSize},
64    resources::ProvisionResult,
65};
66use secrecy::SecretString;
67
68use crate::error::RequestError;
69
70/// An RFC 3339 datetime.
71///
72/// It can be created in either of the following ways:
73/// - Parse an RFC 3339 datetime string using [`FromStr`] or [`str::parse`].
74/// - Convert from [`time::OffsetDateTime`] using [`TryFrom`]/[`TryInto`].
75#[derive(Debug, Clone, Copy, PartialEq, Eq)]
76pub struct S2DateTime(time::OffsetDateTime);
77
78impl TryFrom<time::OffsetDateTime> for S2DateTime {
79    type Error = ValidationError;
80
81    fn try_from(dt: time::OffsetDateTime) -> Result<Self, Self::Error> {
82        dt.format(&time::format_description::well_known::Rfc3339)
83            .map_err(|e| ValidationError(format!("not a valid RFC 3339 datetime: {e}")))?;
84        Ok(Self(dt))
85    }
86}
87
88impl From<S2DateTime> for time::OffsetDateTime {
89    fn from(dt: S2DateTime) -> Self {
90        dt.0
91    }
92}
93
94impl FromStr for S2DateTime {
95    type Err = ValidationError;
96
97    fn from_str(s: &str) -> Result<Self, Self::Err> {
98        time::OffsetDateTime::parse(s, &time::format_description::well_known::Rfc3339)
99            .map(Self)
100            .map_err(|e| ValidationError(format!("not a valid RFC 3339 datetime: {e}")))
101    }
102}
103
104impl fmt::Display for S2DateTime {
105    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
106        write!(
107            f,
108            "{}",
109            self.0
110                .format(&time::format_description::well_known::Rfc3339)
111                .expect("RFC3339 formatting should not fail for S2DateTime")
112        )
113    }
114}
115
116/// Authority for connecting to an S2 basin.
117#[derive(Debug, Clone, PartialEq)]
118pub(crate) enum BasinAuthority {
119    /// Parent zone for basins. DNS is used to route to the correct cell for the basin.
120    ParentZone(Authority),
121    /// Direct cell authority. Basin is expected to be hosted by this cell.
122    Direct(Authority),
123}
124
125/// Account endpoint.
126#[derive(Debug, Clone)]
127pub struct AccountEndpoint {
128    scheme: Scheme,
129    authority: Authority,
130}
131
132impl AccountEndpoint {
133    /// Create a new [`AccountEndpoint`] with the given endpoint.
134    pub fn new(endpoint: &str) -> Result<Self, ValidationError> {
135        endpoint.parse()
136    }
137}
138
139impl FromStr for AccountEndpoint {
140    type Err = ValidationError;
141
142    fn from_str(s: &str) -> Result<Self, Self::Err> {
143        let (scheme, authority) = match s.find("://") {
144            Some(idx) => {
145                let scheme: Scheme = s[..idx]
146                    .parse()
147                    .map_err(|_| "invalid account endpoint scheme".to_string())?;
148                (scheme, &s[idx + 3..])
149            }
150            None => (Scheme::HTTPS, s),
151        };
152        Ok(Self {
153            scheme,
154            authority: authority
155                .parse()
156                .map_err(|e| format!("invalid account endpoint authority: {e}"))?,
157        })
158    }
159}
160
161/// Basin endpoint.
162#[derive(Debug, Clone)]
163pub struct BasinEndpoint {
164    scheme: Scheme,
165    authority: BasinAuthority,
166}
167
168impl BasinEndpoint {
169    /// Create a new [`BasinEndpoint`] with the given endpoint.
170    pub fn new(endpoint: &str) -> Result<Self, ValidationError> {
171        endpoint.parse()
172    }
173}
174
175impl FromStr for BasinEndpoint {
176    type Err = ValidationError;
177
178    fn from_str(s: &str) -> Result<Self, Self::Err> {
179        let (scheme, authority) = match s.find("://") {
180            Some(idx) => {
181                let scheme: Scheme = s[..idx]
182                    .parse()
183                    .map_err(|_| "invalid basin endpoint scheme".to_string())?;
184                (scheme, &s[idx + 3..])
185            }
186            None => (Scheme::HTTPS, s),
187        };
188        let authority = if let Some(authority) = authority.strip_prefix("{basin}.") {
189            BasinAuthority::ParentZone(
190                authority
191                    .parse()
192                    .map_err(|e| format!("invalid basin endpoint authority: {e}"))?,
193            )
194        } else {
195            BasinAuthority::Direct(
196                authority
197                    .parse()
198                    .map_err(|e| format!("invalid basin endpoint authority: {e}"))?,
199            )
200        };
201        Ok(Self { scheme, authority })
202    }
203}
204
205#[derive(Debug, Clone)]
206#[non_exhaustive]
207/// Endpoints for the S2 environment.
208pub struct S2Endpoints {
209    pub(crate) scheme: Scheme,
210    pub(crate) account_authority: Authority,
211    pub(crate) basin_authority: BasinAuthority,
212}
213
214impl S2Endpoints {
215    /// Create a new [`S2Endpoints`] with the given account and basin endpoints.
216    pub fn new(
217        account_endpoint: AccountEndpoint,
218        basin_endpoint: BasinEndpoint,
219    ) -> Result<Self, ValidationError> {
220        if account_endpoint.scheme != basin_endpoint.scheme {
221            return Err("account and basin endpoints must have the same scheme".into());
222        }
223        Ok(Self {
224            scheme: account_endpoint.scheme,
225            account_authority: account_endpoint.authority,
226            basin_authority: basin_endpoint.authority,
227        })
228    }
229
230    /// Create endpoints for a single account and basin endpoint.
231    ///
232    /// This is useful for S2-compatible services that expose both APIs at one endpoint.
233    pub fn for_endpoint(endpoint: &str) -> Result<Self, ValidationError> {
234        Self::new(
235            AccountEndpoint::new(endpoint)?,
236            BasinEndpoint::new(endpoint)?,
237        )
238    }
239
240    /// Create a new [`S2Endpoints`] from environment variables.
241    ///
242    /// The following environment variables are expected to be set:
243    /// - `S2_ACCOUNT_ENDPOINT` - Account-level endpoint.
244    /// - `S2_BASIN_ENDPOINT` - Basin-level endpoint.
245    pub fn from_env() -> Result<Self, ValidationError> {
246        let account_endpoint: AccountEndpoint = match std::env::var("S2_ACCOUNT_ENDPOINT") {
247            Ok(endpoint) => endpoint.parse()?,
248            Err(VarError::NotPresent) => return Err("S2_ACCOUNT_ENDPOINT env var not set".into()),
249            Err(VarError::NotUnicode(_)) => {
250                return Err("S2_ACCOUNT_ENDPOINT is not valid unicode".into());
251            }
252        };
253
254        let basin_endpoint: BasinEndpoint = match std::env::var("S2_BASIN_ENDPOINT") {
255            Ok(endpoint) => endpoint.parse()?,
256            Err(VarError::NotPresent) => return Err("S2_BASIN_ENDPOINT env var not set".into()),
257            Err(VarError::NotUnicode(_)) => {
258                return Err("S2_BASIN_ENDPOINT is not valid unicode".into());
259            }
260        };
261
262        if account_endpoint.scheme != basin_endpoint.scheme {
263            return Err(
264                "S2_ACCOUNT_ENDPOINT and S2_BASIN_ENDPOINT must have the same scheme".into(),
265            );
266        }
267
268        Ok(Self {
269            scheme: account_endpoint.scheme,
270            account_authority: account_endpoint.authority,
271            basin_authority: basin_endpoint.authority,
272        })
273    }
274
275    /// Return the default S2 Cloud endpoints.
276    pub fn for_cloud() -> Self {
277        Self {
278            scheme: Scheme::HTTPS,
279            account_authority: "a.s2.dev".try_into().expect("valid authority"),
280            basin_authority: BasinAuthority::ParentZone(
281                "b.s2.dev".try_into().expect("valid authority"),
282            ),
283        }
284    }
285}
286
287#[derive(Debug, Clone, Copy)]
288/// Compression algorithm for request and response bodies.
289pub enum Compression {
290    /// No compression.
291    None,
292    /// Gzip compression.
293    Gzip,
294    /// Zstd compression.
295    Zstd,
296}
297
298impl From<Compression> for CompressionAlgorithm {
299    fn from(value: Compression) -> Self {
300        match value {
301            Compression::None => CompressionAlgorithm::None,
302            Compression::Gzip => CompressionAlgorithm::Gzip,
303            Compression::Zstd => CompressionAlgorithm::Zstd,
304        }
305    }
306}
307
308#[derive(Debug, Clone, Copy, PartialEq)]
309#[non_exhaustive]
310/// Retry policy for [`append`](crate::S2Stream::append) and
311/// [`append_session`](crate::S2Stream::append_session) operations.
312pub enum AppendRetryPolicy {
313    /// Retry all appends. Use when duplicate records on the stream are acceptable.
314    All,
315    /// Retry when it can be determined that the request had no side effects.
316    ///
317    /// Uses a frame-level signal to detect whether any body frames were consumed
318    /// by the HTTP transport. If no frames were sent, the server never saw the
319    /// request, so retry is safe and will not cause duplicate records.
320    ///
321    /// Certain server errors (`rate_limited`, `hot_server`) are also safe to
322    /// retry regardless of frame signal state, since they guarantee no mutation
323    /// occurred.
324    NoSideEffects,
325}
326
327#[derive(Debug, Clone)]
328#[non_exhaustive]
329/// Configuration for retrying requests in case of transient failures.
330///
331/// Exponential backoff with jitter is the retry strategy. Below is the pseudocode for the strategy:
332/// ```text
333/// base_delay = min(min_base_delay · 2ⁿ, max_base_delay)    (n = retry attempt, starting from 0)
334///     jitter = rand[0, base_delay]
335///     delay  = base_delay + jitter
336/// ````
337pub struct RetryConfig {
338    /// Total number of attempts including the initial try. A value of `1` means no retries.
339    ///
340    /// Defaults to `3`.
341    pub max_attempts: NonZeroU32,
342    /// Minimum base delay for retries.
343    ///
344    /// Defaults to `100ms`.
345    pub min_base_delay: Duration,
346    /// Maximum base delay for retries.
347    ///
348    /// Defaults to `1s`.
349    pub max_base_delay: Duration,
350    /// Retry policy for [`append`](crate::S2Stream::append) and
351    /// [`append_session`](crate::S2Stream::append_session) operations.
352    ///
353    /// Defaults to `All`.
354    pub append_retry_policy: AppendRetryPolicy,
355}
356
357impl Default for RetryConfig {
358    fn default() -> Self {
359        Self {
360            max_attempts: NonZeroU32::new(3).expect("valid non-zero u32"),
361            min_base_delay: Duration::from_millis(100),
362            max_base_delay: Duration::from_secs(1),
363            append_retry_policy: AppendRetryPolicy::All,
364        }
365    }
366}
367
368impl RetryConfig {
369    /// Create a new [`RetryConfig`] with default settings.
370    pub fn new() -> Self {
371        Self::default()
372    }
373
374    pub(crate) fn max_retries(&self) -> u32 {
375        self.max_attempts.get() - 1
376    }
377
378    /// Set the total number of attempts including the initial try.
379    pub fn with_max_attempts(self, max_attempts: NonZeroU32) -> Self {
380        Self {
381            max_attempts,
382            ..self
383        }
384    }
385
386    /// Set the minimum base delay for retries.
387    pub fn with_min_base_delay(self, min_base_delay: Duration) -> Self {
388        Self {
389            min_base_delay,
390            ..self
391        }
392    }
393
394    /// Set the maximum base delay for retries.
395    pub fn with_max_base_delay(self, max_base_delay: Duration) -> Self {
396        Self {
397            max_base_delay,
398            ..self
399        }
400    }
401
402    /// Set the retry policy for [`append`](crate::S2Stream::append) and
403    /// [`append_session`](crate::S2Stream::append_session) operations.
404    pub fn with_append_retry_policy(self, append_retry_policy: AppendRetryPolicy) -> Self {
405        Self {
406            append_retry_policy,
407            ..self
408        }
409    }
410}
411
412#[derive(Debug, Clone)]
413#[non_exhaustive]
414/// Configuration for [`S2`](crate::S2).
415pub struct S2Config {
416    pub(crate) access_token: SecretString,
417    pub(crate) endpoints: S2Endpoints,
418    pub(crate) connection_timeout: Duration,
419    pub(crate) request_timeout: Duration,
420    pub(crate) retry: RetryConfig,
421    pub(crate) compression: Compression,
422    pub(crate) user_agent: HeaderValue,
423    pub(crate) insecure_skip_cert_verification: bool,
424    pub(crate) rustls_crypto_provider: Option<Arc<rustls::crypto::CryptoProvider>>,
425}
426
427impl S2Config {
428    /// Create a new [`S2Config`] with the given access token and default settings.
429    pub fn new(access_token: impl Into<String>) -> Self {
430        Self {
431            access_token: access_token.into().into(),
432            endpoints: S2Endpoints::for_cloud(),
433            connection_timeout: Duration::from_secs(3),
434            request_timeout: Duration::from_secs(5),
435            retry: RetryConfig::new(),
436            compression: Compression::None,
437            user_agent: concat!("s2-sdk-rust/", env!("CARGO_PKG_VERSION"))
438                .parse()
439                .expect("valid user agent"),
440            insecure_skip_cert_verification: false,
441            rustls_crypto_provider: default_rustls_crypto_provider(),
442        }
443    }
444
445    /// Set the S2 endpoints to connect to.
446    pub fn with_endpoints(self, endpoints: S2Endpoints) -> Self {
447        Self { endpoints, ..self }
448    }
449
450    /// Set the timeout for establishing a connection to the server.
451    ///
452    /// Defaults to `3s`.
453    pub fn with_connection_timeout(self, connection_timeout: Duration) -> Self {
454        Self {
455            connection_timeout,
456            ..self
457        }
458    }
459
460    /// Set the timeout for requests.
461    ///
462    /// Defaults to `5s`.
463    pub fn with_request_timeout(self, request_timeout: Duration) -> Self {
464        Self {
465            request_timeout,
466            ..self
467        }
468    }
469
470    /// Set the retry configuration for requests.
471    ///
472    /// See [`RetryConfig`] for defaults.
473    pub fn with_retry(self, retry: RetryConfig) -> Self {
474        Self { retry, ..self }
475    }
476
477    /// Set the compression algorithm for requests and responses.
478    ///
479    /// Defaults to no compression.
480    pub fn with_compression(self, compression: Compression) -> Self {
481        Self {
482            compression,
483            ..self
484        }
485    }
486
487    /// Skip TLS certificate verification (insecure).
488    ///
489    /// This is useful for connecting to endpoints with self-signed certificates
490    /// or certificates that don't match the hostname (similar to `curl -k`).
491    ///
492    /// # Warning
493    ///
494    /// This disables certificate verification and should only be used for
495    /// testing or development purposes. **Never use this in production.**
496    ///
497    /// Defaults to `false`.
498    pub fn with_insecure_skip_cert_verification(self, skip: bool) -> Self {
499        Self {
500            insecure_skip_cert_verification: skip,
501            ..self
502        }
503    }
504
505    /// Use a specific rustls crypto provider for SDK TLS connections.
506    ///
507    /// With default features enabled, the SDK uses the `aws-lc-rs` provider.
508    /// With default features disabled, the SDK uses rustls's process-global
509    /// provider if one has been installed, or returns an error otherwise.
510    ///
511    /// Use this when your application needs a specific rustls provider, such as
512    /// `ring` or a custom [`rustls::crypto::CryptoProvider`]. The corresponding
513    /// rustls provider feature must be enabled in the dependency graph.
514    pub fn with_rustls_crypto_provider(
515        self,
516        provider: impl Into<Arc<rustls::crypto::CryptoProvider>>,
517    ) -> Self {
518        Self {
519            rustls_crypto_provider: Some(provider.into()),
520            ..self
521        }
522    }
523
524    /// Use rustls's `aws-lc-rs` crypto provider.
525    ///
526    /// Requires the `rustls-aws-lc-rs` crate feature.
527    #[cfg(feature = "rustls-aws-lc-rs")]
528    pub fn with_rustls_aws_lc_rs_crypto_provider(self) -> Self {
529        self.with_rustls_crypto_provider(rustls::crypto::aws_lc_rs::default_provider())
530    }
531
532    /// Use rustls's `ring` crypto provider.
533    ///
534    /// Requires the `rustls-ring` crate feature.
535    #[cfg(feature = "rustls-ring")]
536    pub fn with_rustls_ring_crypto_provider(self) -> Self {
537        self.with_rustls_crypto_provider(rustls::crypto::ring::default_provider())
538    }
539
540    #[doc(hidden)]
541    #[cfg(feature = "_hidden")]
542    pub fn with_user_agent(self, user_agent: impl Into<String>) -> Result<Self, ValidationError> {
543        let user_agent = user_agent
544            .into()
545            .parse()
546            .map_err(|e| ValidationError(format!("invalid user agent: {e}")))?;
547        Ok(Self { user_agent, ..self })
548    }
549}
550
551#[cfg(feature = "rustls-aws-lc-rs")]
552fn default_rustls_crypto_provider() -> Option<Arc<rustls::crypto::CryptoProvider>> {
553    Some(Arc::new(rustls::crypto::aws_lc_rs::default_provider()))
554}
555
556#[cfg(all(not(feature = "rustls-aws-lc-rs"), feature = "rustls-ring"))]
557fn default_rustls_crypto_provider() -> Option<Arc<rustls::crypto::CryptoProvider>> {
558    Some(Arc::new(rustls::crypto::ring::default_provider()))
559}
560
561#[cfg(all(not(feature = "rustls-aws-lc-rs"), not(feature = "rustls-ring")))]
562fn default_rustls_crypto_provider() -> Option<Arc<rustls::crypto::CryptoProvider>> {
563    None
564}
565
566#[derive(Debug, Default, Clone, PartialEq, Eq)]
567#[non_exhaustive]
568/// A page of values.
569pub struct Page<T> {
570    /// Values in this page.
571    pub values: Vec<T>,
572    /// Whether there are more pages.
573    pub has_more: bool,
574}
575
576impl<T> Page<T> {
577    pub(crate) fn new(values: impl Into<Vec<T>>, has_more: bool) -> Self {
578        Self {
579            values: values.into(),
580            has_more,
581        }
582    }
583}
584
585#[derive(Debug, Clone, Copy, PartialEq, Eq)]
586/// Storage class for recent appends.
587pub enum StorageClass {
588    /// Standard storage class that offers append latencies under `500ms`.
589    Standard,
590    /// Express storage class that offers append latencies under `50ms`.
591    Express,
592}
593
594impl From<api::config::StorageClass> for StorageClass {
595    fn from(value: api::config::StorageClass) -> Self {
596        match value {
597            api::config::StorageClass::Standard => StorageClass::Standard,
598            api::config::StorageClass::Express => StorageClass::Express,
599        }
600    }
601}
602
603impl From<StorageClass> for api::config::StorageClass {
604    fn from(value: StorageClass) -> Self {
605        match value {
606            StorageClass::Standard => api::config::StorageClass::Standard,
607            StorageClass::Express => api::config::StorageClass::Express,
608        }
609    }
610}
611
612#[derive(Debug, Clone, Copy, PartialEq, Eq)]
613/// Retention policy for records in a stream.
614pub enum RetentionPolicy {
615    /// Age in seconds. Records older than this age are automatically trimmed.
616    Age(u64),
617    /// Records are retained indefinitely unless explicitly trimmed.
618    Infinite,
619}
620
621impl From<api::config::RetentionPolicy> for RetentionPolicy {
622    fn from(value: api::config::RetentionPolicy) -> Self {
623        match value {
624            api::config::RetentionPolicy::Age(secs) => RetentionPolicy::Age(secs),
625            api::config::RetentionPolicy::Infinite(_) => RetentionPolicy::Infinite,
626        }
627    }
628}
629
630impl From<RetentionPolicy> for api::config::RetentionPolicy {
631    fn from(value: RetentionPolicy) -> Self {
632        match value {
633            RetentionPolicy::Age(secs) => api::config::RetentionPolicy::Age(secs),
634            RetentionPolicy::Infinite => {
635                api::config::RetentionPolicy::Infinite(api::config::InfiniteRetention {})
636            }
637        }
638    }
639}
640
641#[derive(Debug, Clone, Copy, PartialEq, Eq)]
642/// Timestamping mode for appends that influences how timestamps are handled.
643pub enum TimestampingMode {
644    /// Prefer client-specified timestamp if present otherwise use arrival time.
645    ClientPrefer,
646    /// Require a client-specified timestamp and reject the append if it is missing.
647    ClientRequire,
648    /// Use the arrival time and ignore any client-specified timestamp.
649    Arrival,
650}
651
652impl From<api::config::TimestampingMode> for TimestampingMode {
653    fn from(value: api::config::TimestampingMode) -> Self {
654        match value {
655            api::config::TimestampingMode::ClientPrefer => TimestampingMode::ClientPrefer,
656            api::config::TimestampingMode::ClientRequire => TimestampingMode::ClientRequire,
657            api::config::TimestampingMode::Arrival => TimestampingMode::Arrival,
658        }
659    }
660}
661
662impl From<TimestampingMode> for api::config::TimestampingMode {
663    fn from(value: TimestampingMode) -> Self {
664        match value {
665            TimestampingMode::ClientPrefer => api::config::TimestampingMode::ClientPrefer,
666            TimestampingMode::ClientRequire => api::config::TimestampingMode::ClientRequire,
667            TimestampingMode::Arrival => api::config::TimestampingMode::Arrival,
668        }
669    }
670}
671
672#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
673#[non_exhaustive]
674/// Configuration for timestamping behavior.
675pub struct TimestampingConfig {
676    /// Timestamping mode for appends that influences how timestamps are handled.
677    ///
678    /// Defaults to [`ClientPrefer`](TimestampingMode::ClientPrefer).
679    pub mode: Option<TimestampingMode>,
680    /// Whether client-specified timestamps are allowed to exceed the arrival time.
681    ///
682    /// Defaults to `false` (client timestamps are capped at the arrival time).
683    pub uncapped: Option<bool>,
684}
685
686impl TimestampingConfig {
687    /// Create a new [`TimestampingConfig`] with default settings.
688    pub fn new() -> Self {
689        Self::default()
690    }
691
692    /// Set the timestamping mode for appends that influences how timestamps are handled.
693    pub fn with_mode(self, mode: TimestampingMode) -> Self {
694        Self {
695            mode: Some(mode),
696            ..self
697        }
698    }
699
700    /// Set whether client-specified timestamps are allowed to exceed the arrival time.
701    pub fn with_uncapped(self, uncapped: bool) -> Self {
702        Self {
703            uncapped: Some(uncapped),
704            ..self
705        }
706    }
707}
708
709impl From<api::config::TimestampingConfig> for TimestampingConfig {
710    fn from(value: api::config::TimestampingConfig) -> Self {
711        Self {
712            mode: value.mode.map(Into::into),
713            uncapped: value.uncapped,
714        }
715    }
716}
717
718impl From<TimestampingConfig> for api::config::TimestampingConfig {
719    fn from(value: TimestampingConfig) -> Self {
720        Self {
721            mode: value.mode.map(Into::into),
722            uncapped: value.uncapped,
723        }
724    }
725}
726
727#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
728#[non_exhaustive]
729/// Configuration for automatically deleting a stream when it becomes empty.
730pub struct DeleteOnEmptyConfig {
731    /// Minimum age in seconds before an empty stream can be deleted.
732    ///
733    /// Defaults to `0` (disables automatic deletion).
734    pub min_age_secs: u64,
735}
736
737impl DeleteOnEmptyConfig {
738    /// Create a new [`DeleteOnEmptyConfig`] with default settings.
739    pub fn new() -> Self {
740        Self::default()
741    }
742
743    /// Set the minimum age in seconds before an empty stream can be deleted.
744    pub fn with_min_age(self, min_age: Duration) -> Self {
745        Self {
746            min_age_secs: min_age.as_secs(),
747        }
748    }
749}
750
751impl From<api::config::DeleteOnEmptyConfig> for DeleteOnEmptyConfig {
752    fn from(value: api::config::DeleteOnEmptyConfig) -> Self {
753        Self {
754            min_age_secs: value.min_age_secs,
755        }
756    }
757}
758
759impl From<DeleteOnEmptyConfig> for api::config::DeleteOnEmptyConfig {
760    fn from(value: DeleteOnEmptyConfig) -> Self {
761        Self {
762            min_age_secs: value.min_age_secs,
763        }
764    }
765}
766
767#[derive(Debug, Clone, Default, PartialEq, Eq)]
768#[non_exhaustive]
769/// Configuration for a stream.
770pub struct StreamConfig {
771    /// Storage class for the stream.
772    ///
773    /// Defaults to [`Express`](StorageClass::Express).
774    pub storage_class: Option<StorageClass>,
775    /// Retention policy for records in the stream.
776    ///
777    /// Defaults to `7 days` of retention.
778    pub retention_policy: Option<RetentionPolicy>,
779    /// Configuration for timestamping behavior.
780    ///
781    /// See [`TimestampingConfig`] for defaults.
782    pub timestamping: Option<TimestampingConfig>,
783    /// Configuration for automatically deleting the stream when it becomes empty.
784    ///
785    /// See [`DeleteOnEmptyConfig`] for defaults.
786    pub delete_on_empty: Option<DeleteOnEmptyConfig>,
787}
788
789impl StreamConfig {
790    /// Create a new [`StreamConfig`] with default settings.
791    pub fn new() -> Self {
792        Self::default()
793    }
794
795    /// Set the storage class for the stream.
796    pub fn with_storage_class(self, storage_class: StorageClass) -> Self {
797        Self {
798            storage_class: Some(storage_class),
799            ..self
800        }
801    }
802
803    /// Set the retention policy for records in the stream.
804    pub fn with_retention_policy(self, retention_policy: RetentionPolicy) -> Self {
805        Self {
806            retention_policy: Some(retention_policy),
807            ..self
808        }
809    }
810
811    /// Set the configuration for timestamping behavior.
812    pub fn with_timestamping(self, timestamping: TimestampingConfig) -> Self {
813        Self {
814            timestamping: Some(timestamping),
815            ..self
816        }
817    }
818
819    /// Set the configuration for automatically deleting the stream when it becomes empty.
820    pub fn with_delete_on_empty(self, delete_on_empty: DeleteOnEmptyConfig) -> Self {
821        Self {
822            delete_on_empty: Some(delete_on_empty),
823            ..self
824        }
825    }
826}
827
828impl From<api::config::StreamConfig> for StreamConfig {
829    fn from(value: api::config::StreamConfig) -> Self {
830        Self {
831            storage_class: value.storage_class.map(Into::into),
832            retention_policy: value.retention_policy.map(Into::into),
833            timestamping: value.timestamping.map(Into::into),
834            delete_on_empty: value.delete_on_empty.map(Into::into),
835        }
836    }
837}
838
839impl From<StreamConfig> for api::config::StreamConfig {
840    fn from(value: StreamConfig) -> Self {
841        Self {
842            storage_class: value.storage_class.map(Into::into),
843            retention_policy: value.retention_policy.map(Into::into),
844            timestamping: value.timestamping.map(Into::into),
845            delete_on_empty: value.delete_on_empty.map(Into::into),
846        }
847    }
848}
849
850#[derive(Debug, Clone, Default, PartialEq, Eq)]
851#[non_exhaustive]
852/// Configuration for a basin.
853pub struct BasinConfig {
854    /// Default configuration for all streams in the basin.
855    ///
856    /// See [`StreamConfig`] for defaults.
857    pub default_stream_config: Option<StreamConfig>,
858    /// Encryption algorithm to apply to newly created streams in the basin.
859    pub stream_cipher: Option<EncryptionAlgorithm>,
860    /// Whether to create stream on append if it doesn't exist using default stream configuration.
861    ///
862    /// Defaults to `false`.
863    pub create_stream_on_append: bool,
864    /// Whether to create stream on read if it doesn't exist using default stream configuration.
865    ///
866    /// Defaults to `false`.
867    pub create_stream_on_read: bool,
868}
869
870impl BasinConfig {
871    /// Create a new [`BasinConfig`] with default settings.
872    pub fn new() -> Self {
873        Self::default()
874    }
875
876    /// Set the default configuration for all streams in the basin.
877    pub fn with_default_stream_config(self, config: StreamConfig) -> Self {
878        Self {
879            default_stream_config: Some(config),
880            ..self
881        }
882    }
883
884    /// Set the encryption algorithm to apply to newly created streams in the basin.
885    pub fn with_stream_cipher(self, stream_cipher: EncryptionAlgorithm) -> Self {
886        Self {
887            stream_cipher: Some(stream_cipher),
888            ..self
889        }
890    }
891
892    /// Set whether to create stream on append if it doesn't exist using default stream
893    /// configuration.
894    pub fn with_create_stream_on_append(self, create_stream_on_append: bool) -> Self {
895        Self {
896            create_stream_on_append,
897            ..self
898        }
899    }
900
901    /// Set whether to create stream on read if it doesn't exist using default stream configuration.
902    pub fn with_create_stream_on_read(self, create_stream_on_read: bool) -> Self {
903        Self {
904            create_stream_on_read,
905            ..self
906        }
907    }
908}
909
910impl From<api::config::BasinConfig> for BasinConfig {
911    fn from(value: api::config::BasinConfig) -> Self {
912        Self {
913            default_stream_config: value.default_stream_config.map(Into::into),
914            stream_cipher: value.stream_cipher.map(Into::into),
915            create_stream_on_append: value.create_stream_on_append,
916            create_stream_on_read: value.create_stream_on_read,
917        }
918    }
919}
920
921impl From<BasinConfig> for api::config::BasinConfig {
922    fn from(value: BasinConfig) -> Self {
923        Self {
924            default_stream_config: value.default_stream_config.map(Into::into),
925            stream_cipher: value.stream_cipher.map(Into::into),
926            create_stream_on_append: value.create_stream_on_append,
927            create_stream_on_read: value.create_stream_on_read,
928        }
929    }
930}
931
932#[derive(Debug, Clone)]
933#[non_exhaustive]
934/// Input for [`create_basin`](crate::S2::create_basin) operation.
935pub struct CreateBasinInput {
936    /// Basin name.
937    pub name: BasinName,
938    /// Configuration for the basin.
939    ///
940    /// See [`BasinConfig`] for defaults.
941    pub config: Option<BasinConfig>,
942    /// Location of the basin.
943    ///
944    /// If omitted when creating, uses the default location for the account.
945    pub location: Option<LocationName>,
946    idempotency_token: String,
947}
948
949impl CreateBasinInput {
950    /// Create a new [`CreateBasinInput`] with the given basin name.
951    pub fn new(name: BasinName) -> Self {
952        Self {
953            name,
954            config: None,
955            location: None,
956            idempotency_token: idempotency_token(),
957        }
958    }
959
960    /// Set the configuration for the basin.
961    pub fn with_config(self, config: BasinConfig) -> Self {
962        Self {
963            config: Some(config),
964            ..self
965        }
966    }
967
968    /// Set the location of the basin.
969    pub fn with_location<S>(self, location: S) -> Result<Self, ValidationError>
970    where
971        S: TryInto<LocationName>,
972        S::Error: fmt::Display,
973    {
974        let location = location
975            .try_into()
976            .map_err(|e| ValidationError(e.to_string()))?;
977        Ok(Self {
978            location: Some(location),
979            ..self
980        })
981    }
982}
983
984impl From<CreateBasinInput> for (api::basin::CreateBasinRequest, String) {
985    fn from(value: CreateBasinInput) -> Self {
986        (
987            api::basin::CreateBasinRequest {
988                basin: value.name,
989                config: value.config.map(Into::into),
990                location: value.location,
991            },
992            value.idempotency_token,
993        )
994    }
995}
996
997#[derive(Debug, Clone)]
998#[non_exhaustive]
999/// Input for [`ensure_basin`](crate::S2::ensure_basin) operation.
1000pub struct EnsureBasinInput {
1001    /// Basin name.
1002    pub name: BasinName,
1003    /// Configuration for the basin.
1004    ///
1005    /// See [`BasinConfig`] for defaults.
1006    pub config: Option<BasinConfig>,
1007    /// Location of the basin.
1008    ///
1009    /// If omitted when creating, uses the default location for the account. Cannot be changed once
1010    /// set.
1011    pub location: Option<LocationName>,
1012}
1013
1014impl EnsureBasinInput {
1015    /// Create a new [`EnsureBasinInput`] with the given basin name.
1016    pub fn new(name: BasinName) -> Self {
1017        Self {
1018            name,
1019            config: None,
1020            location: None,
1021        }
1022    }
1023
1024    /// Set the configuration for the basin.
1025    pub fn with_config(self, config: BasinConfig) -> Self {
1026        Self {
1027            config: Some(config),
1028            ..self
1029        }
1030    }
1031
1032    /// Set the location of the basin.
1033    pub fn with_location<S>(self, location: S) -> Result<Self, ValidationError>
1034    where
1035        S: TryInto<LocationName>,
1036        S::Error: fmt::Display,
1037    {
1038        let location = location
1039            .try_into()
1040            .map_err(|e| ValidationError(e.to_string()))?;
1041        Ok(Self {
1042            location: Some(location),
1043            ..self
1044        })
1045    }
1046}
1047
1048impl From<EnsureBasinInput> for (BasinName, Option<api::basin::EnsureBasinRequest>) {
1049    fn from(value: EnsureBasinInput) -> Self {
1050        let config = value.config;
1051        let request = if config.is_some() || value.location.is_some() {
1052            Some(api::basin::EnsureBasinRequest {
1053                config: config.map(Into::into),
1054                location: value.location,
1055            })
1056        } else {
1057            None
1058        };
1059        (value.name, request)
1060    }
1061}
1062
1063#[derive(Debug, Clone)]
1064/// Output for `ensure` operations ([`ensure_basin`](crate::S2::ensure_basin),
1065/// [`ensure_stream`](crate::S2Basin::ensure_stream)).
1066pub enum EnsureOutput<T> {
1067    /// Resource created.
1068    Created(T),
1069    /// Resource already existed, and its config was updated.
1070    ConfigUpdated(T),
1071    /// Resource already existed, and its config is unchanged.
1072    ConfigUnchanged(T),
1073}
1074
1075impl<T> From<ProvisionResult<T>> for EnsureOutput<T> {
1076    fn from(result: ProvisionResult<T>) -> Self {
1077        match result {
1078            ProvisionResult::Created(info) => EnsureOutput::Created(info),
1079            ProvisionResult::Updated(info) => EnsureOutput::ConfigUpdated(info),
1080            ProvisionResult::Noop(info) => EnsureOutput::ConfigUnchanged(info),
1081        }
1082    }
1083}
1084
1085#[derive(Debug, Clone, Default)]
1086#[non_exhaustive]
1087/// Input for [`list_basins`](crate::S2::list_basins) operation.
1088pub struct ListBasinsInput {
1089    /// Filter basins whose names begin with this value.
1090    ///
1091    /// Defaults to `""`.
1092    pub prefix: BasinNamePrefix,
1093    /// Filter basins whose names are lexicographically greater than this value.
1094    ///
1095    /// Defaults to `""`.
1096    pub start_after: BasinNameStartAfter,
1097    /// Number of basins to return in a page. Will be clamped to a maximum of `1000`.
1098    ///
1099    /// Defaults to `1000`.
1100    pub limit: Option<usize>,
1101}
1102
1103impl ListBasinsInput {
1104    /// Create a new [`ListBasinsInput`] with default values.
1105    pub fn new() -> Self {
1106        Self::default()
1107    }
1108
1109    /// Set the prefix used to filter basins whose names begin with this value.
1110    pub fn with_prefix(self, prefix: BasinNamePrefix) -> Self {
1111        Self { prefix, ..self }
1112    }
1113
1114    /// Set the value used to filter basins whose names are lexicographically greater than this
1115    /// value.
1116    pub fn with_start_after(self, start_after: BasinNameStartAfter) -> Self {
1117        Self {
1118            start_after,
1119            ..self
1120        }
1121    }
1122
1123    /// Set the limit on number of basins to return in a page.
1124    pub fn with_limit(self, limit: usize) -> Self {
1125        Self {
1126            limit: Some(limit),
1127            ..self
1128        }
1129    }
1130}
1131
1132impl From<ListBasinsInput> for api::basin::ListBasinsRequest {
1133    fn from(value: ListBasinsInput) -> Self {
1134        Self {
1135            prefix: Some(value.prefix),
1136            start_after: Some(value.start_after),
1137            limit: value.limit,
1138        }
1139    }
1140}
1141
1142#[derive(Debug, Clone, Default)]
1143/// Input for [`list_all_basins`](crate::S2::list_all_basins) operation.
1144pub struct ListAllBasinsInput {
1145    /// Filter basins whose names begin with this value.
1146    ///
1147    /// Defaults to `""`.
1148    pub prefix: BasinNamePrefix,
1149    /// Filter basins whose names are lexicographically greater than this value.
1150    ///
1151    /// Defaults to `""`.
1152    pub start_after: BasinNameStartAfter,
1153    /// Whether to include basins that are being deleted.
1154    ///
1155    /// Defaults to `false`.
1156    pub include_deleted: bool,
1157}
1158
1159impl ListAllBasinsInput {
1160    /// Create a new [`ListAllBasinsInput`] with default values.
1161    pub fn new() -> Self {
1162        Self::default()
1163    }
1164
1165    /// Set the prefix used to filter basins whose names begin with this value.
1166    pub fn with_prefix(self, prefix: BasinNamePrefix) -> Self {
1167        Self { prefix, ..self }
1168    }
1169
1170    /// Set the value used to filter basins whose names are lexicographically greater than this
1171    /// value.
1172    pub fn with_start_after(self, start_after: BasinNameStartAfter) -> Self {
1173        Self {
1174            start_after,
1175            ..self
1176        }
1177    }
1178
1179    /// Set whether to include basins that are being deleted.
1180    pub fn with_include_deleted(self, include_deleted: bool) -> Self {
1181        Self {
1182            include_deleted,
1183            ..self
1184        }
1185    }
1186}
1187
1188#[derive(Debug, Clone, PartialEq, Eq)]
1189#[non_exhaustive]
1190/// Basin information.
1191pub struct BasinInfo {
1192    /// Basin name.
1193    pub name: BasinName,
1194    /// Location of the basin.
1195    pub location: Option<LocationName>,
1196    /// Creation time.
1197    pub created_at: S2DateTime,
1198    /// Deletion time if the basin is being deleted.
1199    pub deleted_at: Option<S2DateTime>,
1200}
1201
1202impl TryFrom<api::basin::BasinInfo> for BasinInfo {
1203    type Error = ValidationError;
1204
1205    fn try_from(value: api::basin::BasinInfo) -> Result<Self, Self::Error> {
1206        Ok(Self {
1207            name: value.name,
1208            location: value.location,
1209            created_at: value.created_at.try_into()?,
1210            deleted_at: value.deleted_at.map(S2DateTime::try_from).transpose()?,
1211        })
1212    }
1213}
1214
1215#[derive(Debug, Clone)]
1216#[non_exhaustive]
1217/// Input for [`delete_basin`](crate::S2::delete_basin) operation.
1218pub struct DeleteBasinInput {
1219    /// Basin name.
1220    pub name: BasinName,
1221    /// Whether to ignore `Not Found` error if the basin doesn't exist.
1222    pub ignore_not_found: bool,
1223}
1224
1225impl DeleteBasinInput {
1226    /// Create a new [`DeleteBasinInput`] with the given basin name.
1227    pub fn new(name: BasinName) -> Self {
1228        Self {
1229            name,
1230            ignore_not_found: false,
1231        }
1232    }
1233
1234    /// Set whether to ignore `Not Found` error if the basin is not existing.
1235    pub fn with_ignore_not_found(self, ignore_not_found: bool) -> Self {
1236        Self {
1237            ignore_not_found,
1238            ..self
1239        }
1240    }
1241}
1242
1243#[derive(Debug, Clone, Default)]
1244#[non_exhaustive]
1245/// Reconfiguration for [`TimestampingConfig`].
1246pub struct TimestampingReconfiguration {
1247    /// Override for the existing [`mode`](TimestampingConfig::mode).
1248    pub mode: Maybe<Option<TimestampingMode>>,
1249    /// Override for the existing [`uncapped`](TimestampingConfig::uncapped) setting.
1250    pub uncapped: Maybe<Option<bool>>,
1251}
1252
1253impl TimestampingReconfiguration {
1254    /// Create a new [`TimestampingReconfiguration`].
1255    pub fn new() -> Self {
1256        Self::default()
1257    }
1258
1259    /// Set the override for the existing [`mode`](TimestampingConfig::mode).
1260    pub fn with_mode(self, mode: TimestampingMode) -> Self {
1261        Self {
1262            mode: Maybe::Specified(Some(mode)),
1263            ..self
1264        }
1265    }
1266
1267    /// Set the override for the existing [`uncapped`](TimestampingConfig::uncapped).
1268    pub fn with_uncapped(self, uncapped: bool) -> Self {
1269        Self {
1270            uncapped: Maybe::Specified(Some(uncapped)),
1271            ..self
1272        }
1273    }
1274}
1275
1276impl From<TimestampingReconfiguration> for api::config::TimestampingReconfiguration {
1277    fn from(value: TimestampingReconfiguration) -> Self {
1278        Self {
1279            mode: value.mode.map(|m| m.map(Into::into)),
1280            uncapped: value.uncapped,
1281        }
1282    }
1283}
1284
1285#[derive(Debug, Clone, Default)]
1286#[non_exhaustive]
1287/// Reconfiguration for [`DeleteOnEmptyConfig`].
1288pub struct DeleteOnEmptyReconfiguration {
1289    /// Override for the existing [`min_age_secs`](DeleteOnEmptyConfig::min_age_secs).
1290    pub min_age_secs: Maybe<Option<u64>>,
1291}
1292
1293impl DeleteOnEmptyReconfiguration {
1294    /// Create a new [`DeleteOnEmptyReconfiguration`].
1295    pub fn new() -> Self {
1296        Self::default()
1297    }
1298
1299    /// Set the override for the existing [`min_age_secs`](DeleteOnEmptyConfig::min_age_secs).
1300    pub fn with_min_age(self, min_age: Duration) -> Self {
1301        Self {
1302            min_age_secs: Maybe::Specified(Some(min_age.as_secs())),
1303        }
1304    }
1305}
1306
1307impl From<DeleteOnEmptyReconfiguration> for api::config::DeleteOnEmptyReconfiguration {
1308    fn from(value: DeleteOnEmptyReconfiguration) -> Self {
1309        Self {
1310            min_age_secs: value.min_age_secs,
1311        }
1312    }
1313}
1314
1315#[derive(Debug, Clone, Default)]
1316#[non_exhaustive]
1317/// Reconfiguration for [`StreamConfig`].
1318pub struct StreamReconfiguration {
1319    /// Override for the existing [`storage_class`](StreamConfig::storage_class).
1320    pub storage_class: Maybe<Option<StorageClass>>,
1321    /// Override for the existing [`retention_policy`](StreamConfig::retention_policy).
1322    pub retention_policy: Maybe<Option<RetentionPolicy>>,
1323    /// Override for the existing [`timestamping`](StreamConfig::timestamping).
1324    pub timestamping: Maybe<Option<TimestampingReconfiguration>>,
1325    /// Override for the existing [`delete_on_empty`](StreamConfig::delete_on_empty).
1326    pub delete_on_empty: Maybe<Option<DeleteOnEmptyReconfiguration>>,
1327}
1328
1329impl StreamReconfiguration {
1330    /// Create a new [`StreamReconfiguration`].
1331    pub fn new() -> Self {
1332        Self::default()
1333    }
1334
1335    /// Set the override for the existing [`storage_class`](StreamConfig::storage_class).
1336    pub fn with_storage_class(self, storage_class: StorageClass) -> Self {
1337        Self {
1338            storage_class: Maybe::Specified(Some(storage_class)),
1339            ..self
1340        }
1341    }
1342
1343    /// Set the override for the existing [`retention_policy`](StreamConfig::retention_policy).
1344    pub fn with_retention_policy(self, retention_policy: RetentionPolicy) -> Self {
1345        Self {
1346            retention_policy: Maybe::Specified(Some(retention_policy)),
1347            ..self
1348        }
1349    }
1350
1351    /// Set the override for the existing [`timestamping`](StreamConfig::timestamping).
1352    pub fn with_timestamping(self, timestamping: TimestampingReconfiguration) -> Self {
1353        Self {
1354            timestamping: Maybe::Specified(Some(timestamping)),
1355            ..self
1356        }
1357    }
1358
1359    /// Set the override for the existing [`delete_on_empty`](StreamConfig::delete_on_empty).
1360    pub fn with_delete_on_empty(self, delete_on_empty: DeleteOnEmptyReconfiguration) -> Self {
1361        Self {
1362            delete_on_empty: Maybe::Specified(Some(delete_on_empty)),
1363            ..self
1364        }
1365    }
1366}
1367
1368impl From<StreamReconfiguration> for api::config::StreamReconfiguration {
1369    fn from(value: StreamReconfiguration) -> Self {
1370        Self {
1371            storage_class: value.storage_class.map(|m| m.map(Into::into)),
1372            retention_policy: value.retention_policy.map(|m| m.map(Into::into)),
1373            timestamping: value.timestamping.map(|m| m.map(Into::into)),
1374            delete_on_empty: value.delete_on_empty.map(|m| m.map(Into::into)),
1375        }
1376    }
1377}
1378
1379#[derive(Debug, Clone, Default)]
1380#[non_exhaustive]
1381/// Reconfiguration for [`BasinConfig`].
1382pub struct BasinReconfiguration {
1383    /// Override for the existing [`default_stream_config`](BasinConfig::default_stream_config).
1384    pub default_stream_config: Maybe<Option<StreamReconfiguration>>,
1385    /// Override for the existing [`stream_cipher`](BasinConfig::stream_cipher).
1386    pub stream_cipher: Maybe<Option<EncryptionAlgorithm>>,
1387    /// Override for the existing
1388    /// [`create_stream_on_append`](BasinConfig::create_stream_on_append).
1389    pub create_stream_on_append: Maybe<bool>,
1390    /// Override for the existing [`create_stream_on_read`](BasinConfig::create_stream_on_read).
1391    pub create_stream_on_read: Maybe<bool>,
1392}
1393
1394impl BasinReconfiguration {
1395    /// Create a new [`BasinReconfiguration`].
1396    pub fn new() -> Self {
1397        Self::default()
1398    }
1399
1400    /// Set the override for the existing
1401    /// [`default_stream_config`](BasinConfig::default_stream_config).
1402    pub fn with_default_stream_config(self, config: StreamReconfiguration) -> Self {
1403        Self {
1404            default_stream_config: Maybe::Specified(Some(config)),
1405            ..self
1406        }
1407    }
1408
1409    /// Set the override for the existing [`stream_cipher`](BasinConfig::stream_cipher).
1410    pub fn with_stream_cipher(self, stream_cipher: EncryptionAlgorithm) -> Self {
1411        Self {
1412            stream_cipher: Maybe::Specified(Some(stream_cipher)),
1413            ..self
1414        }
1415    }
1416
1417    /// Set the override for the existing
1418    /// [`create_stream_on_append`](BasinConfig::create_stream_on_append).
1419    pub fn with_create_stream_on_append(self, create_stream_on_append: bool) -> Self {
1420        Self {
1421            create_stream_on_append: Maybe::Specified(create_stream_on_append),
1422            ..self
1423        }
1424    }
1425
1426    /// Set the override for the existing
1427    /// [`create_stream_on_read`](BasinConfig::create_stream_on_read).
1428    pub fn with_create_stream_on_read(self, create_stream_on_read: bool) -> Self {
1429        Self {
1430            create_stream_on_read: Maybe::Specified(create_stream_on_read),
1431            ..self
1432        }
1433    }
1434}
1435
1436impl From<BasinReconfiguration> for api::config::BasinReconfiguration {
1437    fn from(value: BasinReconfiguration) -> Self {
1438        Self {
1439            default_stream_config: value.default_stream_config.map(|m| m.map(Into::into)),
1440            stream_cipher: value.stream_cipher.map(|m| m.map(Into::into)),
1441            create_stream_on_append: value.create_stream_on_append,
1442            create_stream_on_read: value.create_stream_on_read,
1443        }
1444    }
1445}
1446
1447#[derive(Debug, Clone)]
1448#[non_exhaustive]
1449/// Input for [`reconfigure_basin`](crate::S2::reconfigure_basin) operation.
1450pub struct ReconfigureBasinInput {
1451    /// Basin name.
1452    pub name: BasinName,
1453    /// Reconfiguration for [`BasinConfig`].
1454    pub config: BasinReconfiguration,
1455}
1456
1457impl ReconfigureBasinInput {
1458    /// Create a new [`ReconfigureBasinInput`] with the given basin name and reconfiguration.
1459    pub fn new(name: BasinName, config: BasinReconfiguration) -> Self {
1460        Self { name, config }
1461    }
1462}
1463
1464#[derive(Debug, Clone, Default)]
1465#[non_exhaustive]
1466/// Input for [`list_access_tokens`](crate::S2::list_access_tokens) operation.
1467pub struct ListAccessTokensInput {
1468    /// Filter access tokens whose IDs begin with this value.
1469    ///
1470    /// Defaults to `""`.
1471    pub prefix: AccessTokenIdPrefix,
1472    /// Filter access tokens whose IDs are lexicographically greater than this value.
1473    ///
1474    /// Defaults to `""`.
1475    pub start_after: AccessTokenIdStartAfter,
1476    /// Number of access tokens to return in a page. Will be clamped to a maximum of `1000`.
1477    ///
1478    /// Defaults to `1000`.
1479    pub limit: Option<usize>,
1480}
1481
1482impl ListAccessTokensInput {
1483    /// Create a new [`ListAccessTokensInput`] with default values.
1484    pub fn new() -> Self {
1485        Self::default()
1486    }
1487
1488    /// Set the prefix used to filter access tokens whose IDs begin with this value.
1489    pub fn with_prefix(self, prefix: AccessTokenIdPrefix) -> Self {
1490        Self { prefix, ..self }
1491    }
1492
1493    /// Set the value used to filter access tokens whose IDs are lexicographically greater than this
1494    /// value.
1495    pub fn with_start_after(self, start_after: AccessTokenIdStartAfter) -> Self {
1496        Self {
1497            start_after,
1498            ..self
1499        }
1500    }
1501
1502    /// Set the limit on number of access tokens to return in a page.
1503    pub fn with_limit(self, limit: usize) -> Self {
1504        Self {
1505            limit: Some(limit),
1506            ..self
1507        }
1508    }
1509}
1510
1511impl From<ListAccessTokensInput> for api::access::ListAccessTokensRequest {
1512    fn from(value: ListAccessTokensInput) -> Self {
1513        Self {
1514            prefix: Some(value.prefix),
1515            start_after: Some(value.start_after),
1516            limit: value.limit,
1517        }
1518    }
1519}
1520
1521#[derive(Debug, Clone, Default)]
1522/// Input for [`list_all_access_tokens`](crate::S2::list_all_access_tokens) operation.
1523pub struct ListAllAccessTokensInput {
1524    /// Filter access tokens whose IDs begin with this value.
1525    ///
1526    /// Defaults to `""`.
1527    pub prefix: AccessTokenIdPrefix,
1528    /// Filter access tokens whose IDs are lexicographically greater than this value.
1529    ///
1530    /// Defaults to `""`.
1531    pub start_after: AccessTokenIdStartAfter,
1532}
1533
1534impl ListAllAccessTokensInput {
1535    /// Create a new [`ListAllAccessTokensInput`] with default values.
1536    pub fn new() -> Self {
1537        Self::default()
1538    }
1539
1540    /// Set the prefix used to filter access tokens whose IDs begin with this value.
1541    pub fn with_prefix(self, prefix: AccessTokenIdPrefix) -> Self {
1542        Self { prefix, ..self }
1543    }
1544
1545    /// Set the value used to filter access tokens whose IDs are lexicographically greater than
1546    /// this value.
1547    pub fn with_start_after(self, start_after: AccessTokenIdStartAfter) -> Self {
1548        Self {
1549            start_after,
1550            ..self
1551        }
1552    }
1553}
1554
1555#[derive(Debug, Clone, PartialEq, Eq)]
1556#[non_exhaustive]
1557/// Location information.
1558pub struct LocationInfo {
1559    /// Location name.
1560    pub name: LocationName,
1561    /// Location represents a private placement, limited by account.
1562    pub is_private: bool,
1563}
1564
1565impl From<api::location::LocationInfo> for LocationInfo {
1566    fn from(value: api::location::LocationInfo) -> Self {
1567        Self {
1568            name: value.name,
1569            is_private: value.is_private,
1570        }
1571    }
1572}
1573
1574#[derive(Debug, Clone)]
1575#[non_exhaustive]
1576/// Access token information.
1577pub struct AccessTokenInfo {
1578    /// Access token ID.
1579    pub id: AccessTokenId,
1580    /// Expiration time, or `None` if the token does not expire.
1581    pub expires_at: Option<S2DateTime>,
1582    /// Whether to automatically prefix stream names during creation and strip the prefix during
1583    /// listing.
1584    pub auto_prefix_streams: bool,
1585    /// Scope of the access token.
1586    pub scope: AccessTokenScope,
1587}
1588
1589impl TryFrom<api::access::AccessTokenInfo> for AccessTokenInfo {
1590    type Error = ValidationError;
1591
1592    fn try_from(value: api::access::AccessTokenInfo) -> Result<Self, Self::Error> {
1593        let expires_at = value.expires_at.map(S2DateTime::try_from).transpose()?;
1594        Ok(Self {
1595            id: value.id,
1596            expires_at,
1597            auto_prefix_streams: value.auto_prefix_streams,
1598            scope: value.scope.into(),
1599        })
1600    }
1601}
1602
1603#[derive(Debug, Clone)]
1604/// Pattern for matching basins.
1605///
1606/// See [`AccessTokenScope::basins`].
1607pub enum BasinMatcher {
1608    /// Match no basins.
1609    None,
1610    /// Match exactly this basin.
1611    Exact(BasinName),
1612    /// Match all basins with this prefix.
1613    Prefix(BasinNamePrefix),
1614}
1615
1616#[derive(Debug, Clone)]
1617/// Pattern for matching streams.
1618///
1619/// See [`AccessTokenScope::streams`].
1620pub enum StreamMatcher {
1621    /// Match no streams.
1622    None,
1623    /// Match exactly this stream.
1624    Exact(StreamName),
1625    /// Match all streams with this prefix.
1626    Prefix(StreamNamePrefix),
1627}
1628
1629#[derive(Debug, Clone)]
1630/// Pattern for matching access tokens.
1631///
1632/// See [`AccessTokenScope::access_tokens`].
1633pub enum AccessTokenMatcher {
1634    /// Match no access tokens.
1635    None,
1636    /// Match exactly this access token.
1637    Exact(AccessTokenId),
1638    /// Match all access tokens with this prefix.
1639    Prefix(AccessTokenIdPrefix),
1640}
1641
1642#[derive(Debug, Clone, Default)]
1643#[non_exhaustive]
1644/// Permissions indicating allowed operations.
1645pub struct ReadWritePermissions {
1646    /// Read permission.
1647    ///
1648    /// Defaults to `false`.
1649    pub read: bool,
1650    /// Write permission.
1651    ///
1652    /// Defaults to `false`.
1653    pub write: bool,
1654}
1655
1656impl ReadWritePermissions {
1657    /// Create a new [`ReadWritePermissions`] with default values.
1658    pub fn new() -> Self {
1659        Self::default()
1660    }
1661
1662    /// Create read-only permissions.
1663    pub fn read_only() -> Self {
1664        Self {
1665            read: true,
1666            write: false,
1667        }
1668    }
1669
1670    /// Create write-only permissions.
1671    pub fn write_only() -> Self {
1672        Self {
1673            read: false,
1674            write: true,
1675        }
1676    }
1677
1678    /// Create read-write permissions.
1679    pub fn read_write() -> Self {
1680        Self {
1681            read: true,
1682            write: true,
1683        }
1684    }
1685}
1686
1687impl From<ReadWritePermissions> for api::access::ReadWritePermissions {
1688    fn from(value: ReadWritePermissions) -> Self {
1689        Self {
1690            read: Some(value.read),
1691            write: Some(value.write),
1692        }
1693    }
1694}
1695
1696impl From<api::access::ReadWritePermissions> for ReadWritePermissions {
1697    fn from(value: api::access::ReadWritePermissions) -> Self {
1698        Self {
1699            read: value.read.unwrap_or_default(),
1700            write: value.write.unwrap_or_default(),
1701        }
1702    }
1703}
1704
1705#[derive(Debug, Clone, Default)]
1706#[non_exhaustive]
1707/// Permissions at the operation group level.
1708///
1709/// See [`AccessTokenScope::op_group_perms`].
1710pub struct OperationGroupPermissions {
1711    /// Account-level access permissions.
1712    ///
1713    /// Defaults to `None`.
1714    pub account: Option<ReadWritePermissions>,
1715    /// Basin-level access permissions.
1716    ///
1717    /// Defaults to `None`.
1718    pub basin: Option<ReadWritePermissions>,
1719    /// Stream-level access permissions.
1720    ///
1721    /// Defaults to `None`.
1722    pub stream: Option<ReadWritePermissions>,
1723}
1724
1725impl OperationGroupPermissions {
1726    /// Create a new [`OperationGroupPermissions`] with default values.
1727    pub fn new() -> Self {
1728        Self::default()
1729    }
1730
1731    /// Create read-only permissions for all groups.
1732    pub fn read_only_all() -> Self {
1733        Self {
1734            account: Some(ReadWritePermissions::read_only()),
1735            basin: Some(ReadWritePermissions::read_only()),
1736            stream: Some(ReadWritePermissions::read_only()),
1737        }
1738    }
1739
1740    /// Create write-only permissions for all groups.
1741    pub fn write_only_all() -> Self {
1742        Self {
1743            account: Some(ReadWritePermissions::write_only()),
1744            basin: Some(ReadWritePermissions::write_only()),
1745            stream: Some(ReadWritePermissions::write_only()),
1746        }
1747    }
1748
1749    /// Create read-write permissions for all groups.
1750    pub fn read_write_all() -> Self {
1751        Self {
1752            account: Some(ReadWritePermissions::read_write()),
1753            basin: Some(ReadWritePermissions::read_write()),
1754            stream: Some(ReadWritePermissions::read_write()),
1755        }
1756    }
1757
1758    /// Set account-level access permissions.
1759    pub fn with_account(self, account: ReadWritePermissions) -> Self {
1760        Self {
1761            account: Some(account),
1762            ..self
1763        }
1764    }
1765
1766    /// Set basin-level access permissions.
1767    pub fn with_basin(self, basin: ReadWritePermissions) -> Self {
1768        Self {
1769            basin: Some(basin),
1770            ..self
1771        }
1772    }
1773
1774    /// Set stream-level access permissions.
1775    pub fn with_stream(self, stream: ReadWritePermissions) -> Self {
1776        Self {
1777            stream: Some(stream),
1778            ..self
1779        }
1780    }
1781}
1782
1783impl From<OperationGroupPermissions> for api::access::PermittedOperationGroups {
1784    fn from(value: OperationGroupPermissions) -> Self {
1785        Self {
1786            account: value.account.map(Into::into),
1787            basin: value.basin.map(Into::into),
1788            stream: value.stream.map(Into::into),
1789        }
1790    }
1791}
1792
1793impl From<api::access::PermittedOperationGroups> for OperationGroupPermissions {
1794    fn from(value: api::access::PermittedOperationGroups) -> Self {
1795        Self {
1796            account: value.account.map(Into::into),
1797            basin: value.basin.map(Into::into),
1798            stream: value.stream.map(Into::into),
1799        }
1800    }
1801}
1802
1803#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
1804/// Individual operation that can be permitted.
1805///
1806/// See [`AccessTokenScope::ops`].
1807pub enum Operation {
1808    /// List basins.
1809    ListBasins,
1810    /// Create a basin.
1811    CreateBasin,
1812    /// Get basin configuration.
1813    GetBasinConfig,
1814    /// Delete a basin.
1815    DeleteBasin,
1816    /// Reconfigure a basin.
1817    ReconfigureBasin,
1818    /// List access tokens.
1819    ListAccessTokens,
1820    /// Issue an access token.
1821    IssueAccessToken,
1822    /// Revoke an access token.
1823    RevokeAccessToken,
1824    /// Get account metrics.
1825    GetAccountMetrics,
1826    /// Get basin metrics.
1827    GetBasinMetrics,
1828    /// Get stream metrics.
1829    GetStreamMetrics,
1830    /// List streams.
1831    ListStreams,
1832    /// Create a stream.
1833    CreateStream,
1834    /// Get stream configuration.
1835    GetStreamConfig,
1836    /// Delete a stream.
1837    DeleteStream,
1838    /// Reconfigure a stream.
1839    ReconfigureStream,
1840    /// Check the tail of a stream.
1841    CheckTail,
1842    /// Append records to a stream.
1843    Append,
1844    /// Read records from a stream.
1845    Read,
1846    /// Trim records on a stream.
1847    Trim,
1848    /// Set the fencing token on a stream.
1849    Fence,
1850    /// List locations.
1851    ListLocations,
1852    /// Get the default location.
1853    GetDefaultLocation,
1854    /// Set the default location.
1855    SetDefaultLocation,
1856}
1857
1858impl From<Operation> for api::access::Operation {
1859    fn from(value: Operation) -> Self {
1860        match value {
1861            Operation::ListBasins => api::access::Operation::ListBasins,
1862            Operation::CreateBasin => api::access::Operation::CreateBasin,
1863            Operation::DeleteBasin => api::access::Operation::DeleteBasin,
1864            Operation::ReconfigureBasin => api::access::Operation::ReconfigureBasin,
1865            Operation::GetBasinConfig => api::access::Operation::GetBasinConfig,
1866            Operation::IssueAccessToken => api::access::Operation::IssueAccessToken,
1867            Operation::RevokeAccessToken => api::access::Operation::RevokeAccessToken,
1868            Operation::ListAccessTokens => api::access::Operation::ListAccessTokens,
1869            Operation::ListStreams => api::access::Operation::ListStreams,
1870            Operation::CreateStream => api::access::Operation::CreateStream,
1871            Operation::DeleteStream => api::access::Operation::DeleteStream,
1872            Operation::GetStreamConfig => api::access::Operation::GetStreamConfig,
1873            Operation::ReconfigureStream => api::access::Operation::ReconfigureStream,
1874            Operation::CheckTail => api::access::Operation::CheckTail,
1875            Operation::Append => api::access::Operation::Append,
1876            Operation::Read => api::access::Operation::Read,
1877            Operation::Trim => api::access::Operation::Trim,
1878            Operation::Fence => api::access::Operation::Fence,
1879            Operation::GetAccountMetrics => api::access::Operation::AccountMetrics,
1880            Operation::GetBasinMetrics => api::access::Operation::BasinMetrics,
1881            Operation::GetStreamMetrics => api::access::Operation::StreamMetrics,
1882            Operation::ListLocations => api::access::Operation::ListLocations,
1883            Operation::GetDefaultLocation => api::access::Operation::GetDefaultLocation,
1884            Operation::SetDefaultLocation => api::access::Operation::SetDefaultLocation,
1885        }
1886    }
1887}
1888
1889impl From<api::access::Operation> for Operation {
1890    fn from(value: api::access::Operation) -> Self {
1891        match value {
1892            api::access::Operation::ListBasins => Operation::ListBasins,
1893            api::access::Operation::CreateBasin => Operation::CreateBasin,
1894            api::access::Operation::DeleteBasin => Operation::DeleteBasin,
1895            api::access::Operation::ReconfigureBasin => Operation::ReconfigureBasin,
1896            api::access::Operation::GetBasinConfig => Operation::GetBasinConfig,
1897            api::access::Operation::IssueAccessToken => Operation::IssueAccessToken,
1898            api::access::Operation::RevokeAccessToken => Operation::RevokeAccessToken,
1899            api::access::Operation::ListAccessTokens => Operation::ListAccessTokens,
1900            api::access::Operation::ListStreams => Operation::ListStreams,
1901            api::access::Operation::CreateStream => Operation::CreateStream,
1902            api::access::Operation::DeleteStream => Operation::DeleteStream,
1903            api::access::Operation::GetStreamConfig => Operation::GetStreamConfig,
1904            api::access::Operation::ReconfigureStream => Operation::ReconfigureStream,
1905            api::access::Operation::CheckTail => Operation::CheckTail,
1906            api::access::Operation::Append => Operation::Append,
1907            api::access::Operation::Read => Operation::Read,
1908            api::access::Operation::Trim => Operation::Trim,
1909            api::access::Operation::Fence => Operation::Fence,
1910            api::access::Operation::AccountMetrics => Operation::GetAccountMetrics,
1911            api::access::Operation::BasinMetrics => Operation::GetBasinMetrics,
1912            api::access::Operation::StreamMetrics => Operation::GetStreamMetrics,
1913            api::access::Operation::ListLocations => Operation::ListLocations,
1914            api::access::Operation::GetDefaultLocation => Operation::GetDefaultLocation,
1915            api::access::Operation::SetDefaultLocation => Operation::SetDefaultLocation,
1916        }
1917    }
1918}
1919
1920#[derive(Debug, Clone)]
1921#[non_exhaustive]
1922/// Scope of an access token.
1923///
1924/// **Note:** The final set of permitted operations is the union of [`ops`](AccessTokenScope::ops)
1925/// and the operations permitted by [`op_group_perms`](AccessTokenScope::op_group_perms). Also, the
1926/// final set must not be empty.
1927///
1928/// See [`IssueAccessTokenInput::scope`].
1929pub struct AccessTokenScopeInput {
1930    basins: Option<BasinMatcher>,
1931    streams: Option<StreamMatcher>,
1932    access_tokens: Option<AccessTokenMatcher>,
1933    op_group_perms: Option<OperationGroupPermissions>,
1934    ops: HashSet<Operation>,
1935}
1936
1937impl AccessTokenScopeInput {
1938    /// Create a new [`AccessTokenScopeInput`] with the given permitted operations.
1939    pub fn from_ops(ops: impl IntoIterator<Item = Operation>) -> Self {
1940        Self {
1941            basins: None,
1942            streams: None,
1943            access_tokens: None,
1944            op_group_perms: None,
1945            ops: ops.into_iter().collect(),
1946        }
1947    }
1948
1949    /// Create a new [`AccessTokenScopeInput`] with the given operation group permissions.
1950    pub fn from_op_group_perms(op_group_perms: OperationGroupPermissions) -> Self {
1951        Self {
1952            basins: None,
1953            streams: None,
1954            access_tokens: None,
1955            op_group_perms: Some(op_group_perms),
1956            ops: HashSet::default(),
1957        }
1958    }
1959
1960    /// Set the permitted operations.
1961    pub fn with_ops(self, ops: impl IntoIterator<Item = Operation>) -> Self {
1962        Self {
1963            ops: ops.into_iter().collect(),
1964            ..self
1965        }
1966    }
1967
1968    /// Set the access permissions at the operation group level.
1969    pub fn with_op_group_perms(self, op_group_perms: OperationGroupPermissions) -> Self {
1970        Self {
1971            op_group_perms: Some(op_group_perms),
1972            ..self
1973        }
1974    }
1975
1976    /// Set the permitted basins.
1977    ///
1978    /// Defaults to no basins.
1979    pub fn with_basins(self, basins: BasinMatcher) -> Self {
1980        Self {
1981            basins: Some(basins),
1982            ..self
1983        }
1984    }
1985
1986    /// Set the permitted streams.
1987    ///
1988    /// Defaults to no streams.
1989    pub fn with_streams(self, streams: StreamMatcher) -> Self {
1990        Self {
1991            streams: Some(streams),
1992            ..self
1993        }
1994    }
1995
1996    /// Set the permitted access tokens.
1997    ///
1998    /// Defaults to no access tokens.
1999    pub fn with_access_tokens(self, access_tokens: AccessTokenMatcher) -> Self {
2000        Self {
2001            access_tokens: Some(access_tokens),
2002            ..self
2003        }
2004    }
2005}
2006
2007#[derive(Debug, Clone)]
2008#[non_exhaustive]
2009/// Scope of an access token.
2010pub struct AccessTokenScope {
2011    /// Permitted basins.
2012    pub basins: Option<BasinMatcher>,
2013    /// Permitted streams.
2014    pub streams: Option<StreamMatcher>,
2015    /// Permitted access tokens.
2016    pub access_tokens: Option<AccessTokenMatcher>,
2017    /// Permissions at the operation group level.
2018    pub op_group_perms: Option<OperationGroupPermissions>,
2019    /// Permitted operations.
2020    pub ops: HashSet<Operation>,
2021}
2022
2023impl From<api::access::AccessTokenScope> for AccessTokenScope {
2024    fn from(value: api::access::AccessTokenScope) -> Self {
2025        Self {
2026            basins: value.basins.map(|rs| match rs {
2027                api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e)) => {
2028                    BasinMatcher::Exact(e)
2029                }
2030                api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty) => {
2031                    BasinMatcher::None
2032                }
2033                api::access::ResourceSet::Prefix(p) => BasinMatcher::Prefix(p),
2034            }),
2035            streams: value.streams.map(|rs| match rs {
2036                api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e)) => {
2037                    StreamMatcher::Exact(e)
2038                }
2039                api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty) => {
2040                    StreamMatcher::None
2041                }
2042                api::access::ResourceSet::Prefix(p) => StreamMatcher::Prefix(p),
2043            }),
2044            access_tokens: value.access_tokens.map(|rs| match rs {
2045                api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e)) => {
2046                    AccessTokenMatcher::Exact(e)
2047                }
2048                api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty) => {
2049                    AccessTokenMatcher::None
2050                }
2051                api::access::ResourceSet::Prefix(p) => AccessTokenMatcher::Prefix(p),
2052            }),
2053            op_group_perms: value.op_groups.map(Into::into),
2054            ops: value
2055                .ops
2056                .map(|ops| ops.into_iter().map(Into::into).collect())
2057                .unwrap_or_default(),
2058        }
2059    }
2060}
2061
2062impl From<AccessTokenScopeInput> for api::access::AccessTokenScope {
2063    fn from(value: AccessTokenScopeInput) -> Self {
2064        Self {
2065            basins: value.basins.map(|rs| match rs {
2066                BasinMatcher::None => {
2067                    api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty)
2068                }
2069                BasinMatcher::Exact(e) => {
2070                    api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e))
2071                }
2072                BasinMatcher::Prefix(p) => api::access::ResourceSet::Prefix(p),
2073            }),
2074            streams: value.streams.map(|rs| match rs {
2075                StreamMatcher::None => {
2076                    api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty)
2077                }
2078                StreamMatcher::Exact(e) => {
2079                    api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e))
2080                }
2081                StreamMatcher::Prefix(p) => api::access::ResourceSet::Prefix(p),
2082            }),
2083            access_tokens: value.access_tokens.map(|rs| match rs {
2084                AccessTokenMatcher::None => {
2085                    api::access::ResourceSet::Exact(api::access::MaybeEmpty::Empty)
2086                }
2087                AccessTokenMatcher::Exact(e) => {
2088                    api::access::ResourceSet::Exact(api::access::MaybeEmpty::NonEmpty(e))
2089                }
2090                AccessTokenMatcher::Prefix(p) => api::access::ResourceSet::Prefix(p),
2091            }),
2092            op_groups: value.op_group_perms.map(Into::into),
2093            ops: if value.ops.is_empty() {
2094                None
2095            } else {
2096                Some(value.ops.into_iter().map(Into::into).collect())
2097            },
2098        }
2099    }
2100}
2101
2102#[derive(Debug, Clone)]
2103#[non_exhaustive]
2104/// Input for [`issue_access_token`](crate::S2::issue_access_token).
2105pub struct IssueAccessTokenInput {
2106    /// Access token ID.
2107    pub id: AccessTokenId,
2108    /// Expiration time.
2109    ///
2110    /// Defaults to the expiration time of requestor's access token passed via
2111    /// [`S2Config`](S2Config::new).
2112    pub expires_at: Option<S2DateTime>,
2113    /// Whether to automatically prefix stream names during creation and strip the prefix during
2114    /// listing.
2115    ///
2116    /// **Note:** [`scope.streams`](AccessTokenScopeInput::with_streams) must be set with the
2117    /// prefix.
2118    ///
2119    /// Defaults to `false`.
2120    pub auto_prefix_streams: bool,
2121    /// Scope of the token.
2122    pub scope: AccessTokenScopeInput,
2123}
2124
2125impl IssueAccessTokenInput {
2126    /// Create a new [`IssueAccessTokenInput`] with the given id and scope.
2127    pub fn new(id: AccessTokenId, scope: AccessTokenScopeInput) -> Self {
2128        Self {
2129            id,
2130            expires_at: None,
2131            auto_prefix_streams: false,
2132            scope,
2133        }
2134    }
2135
2136    /// Set the expiration time.
2137    pub fn with_expires_at(self, expires_at: S2DateTime) -> Self {
2138        Self {
2139            expires_at: Some(expires_at),
2140            ..self
2141        }
2142    }
2143
2144    /// Set whether to automatically prefix stream names during creation and strip the prefix during
2145    /// listing.
2146    pub fn with_auto_prefix_streams(self, auto_prefix_streams: bool) -> Self {
2147        Self {
2148            auto_prefix_streams,
2149            ..self
2150        }
2151    }
2152}
2153
2154impl From<IssueAccessTokenInput> for api::access::IssueAccessTokenRequest {
2155    fn from(value: IssueAccessTokenInput) -> Self {
2156        Self {
2157            id: value.id,
2158            expires_at: value.expires_at.map(Into::into),
2159            auto_prefix_streams: value.auto_prefix_streams.then_some(true),
2160            scope: value.scope.into(),
2161        }
2162    }
2163}
2164
2165#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2166/// Interval to accumulate over for timeseries metric sets.
2167pub enum TimeseriesInterval {
2168    /// Minute.
2169    Minute,
2170    /// Hour.
2171    Hour,
2172    /// Day.
2173    Day,
2174}
2175
2176impl From<TimeseriesInterval> for api::metrics::TimeseriesInterval {
2177    fn from(value: TimeseriesInterval) -> Self {
2178        match value {
2179            TimeseriesInterval::Minute => api::metrics::TimeseriesInterval::Minute,
2180            TimeseriesInterval::Hour => api::metrics::TimeseriesInterval::Hour,
2181            TimeseriesInterval::Day => api::metrics::TimeseriesInterval::Day,
2182        }
2183    }
2184}
2185
2186impl From<api::metrics::TimeseriesInterval> for TimeseriesInterval {
2187    fn from(value: api::metrics::TimeseriesInterval) -> Self {
2188        match value {
2189            api::metrics::TimeseriesInterval::Minute => TimeseriesInterval::Minute,
2190            api::metrics::TimeseriesInterval::Hour => TimeseriesInterval::Hour,
2191            api::metrics::TimeseriesInterval::Day => TimeseriesInterval::Day,
2192        }
2193    }
2194}
2195
2196#[derive(Debug, Clone, Copy)]
2197#[non_exhaustive]
2198/// Time range as Unix epoch seconds.
2199pub struct TimeRange {
2200    /// Start timestamp (inclusive).
2201    pub start: u32,
2202    /// End timestamp (exclusive).
2203    pub end: u32,
2204}
2205
2206impl TimeRange {
2207    /// Create a new [`TimeRange`] with the given start and end timestamps.
2208    pub fn new(start: u32, end: u32) -> Self {
2209        Self { start, end }
2210    }
2211}
2212
2213#[derive(Debug, Clone, Copy)]
2214#[non_exhaustive]
2215/// Time range as Unix epoch seconds and accumulation interval.
2216pub struct TimeRangeAndInterval {
2217    /// Start timestamp (inclusive).
2218    pub start: u32,
2219    /// End timestamp (exclusive).
2220    pub end: u32,
2221    /// Interval to accumulate over for timeseries metric sets.
2222    ///
2223    /// Default is dependent on the requested metric set.
2224    pub interval: Option<TimeseriesInterval>,
2225}
2226
2227impl TimeRangeAndInterval {
2228    /// Create a new [`TimeRangeAndInterval`] with the given start and end timestamps.
2229    pub fn new(start: u32, end: u32) -> Self {
2230        Self {
2231            start,
2232            end,
2233            interval: None,
2234        }
2235    }
2236
2237    /// Set the interval to accumulate over for timeseries metric sets.
2238    pub fn with_interval(self, interval: TimeseriesInterval) -> Self {
2239        Self {
2240            interval: Some(interval),
2241            ..self
2242        }
2243    }
2244}
2245
2246#[derive(Debug, Clone, Copy)]
2247/// Account metric set to return.
2248pub enum AccountMetricSet {
2249    /// Returns a [`LabelMetric`] representing all basins which had at least one stream within the
2250    /// specified time range.
2251    ActiveBasins(TimeRange),
2252    /// Returns [`AccumulationMetric`]s, one per account operation type.
2253    ///
2254    /// Each metric represents a timeseries of the number of operations, with one accumulated value
2255    /// per interval over the requested time range.
2256    ///
2257    /// [`interval`](TimeRangeAndInterval::interval) defaults to [`hour`](TimeseriesInterval::Hour).
2258    AccountOps(TimeRangeAndInterval),
2259}
2260
2261#[derive(Debug, Clone)]
2262#[non_exhaustive]
2263/// Input for [`get_account_metrics`](crate::S2::get_account_metrics) operation.
2264pub struct GetAccountMetricsInput {
2265    /// Metric set to return.
2266    pub set: AccountMetricSet,
2267}
2268
2269impl GetAccountMetricsInput {
2270    /// Create a new [`GetAccountMetricsInput`] with the given account metric set.
2271    pub fn new(set: AccountMetricSet) -> Self {
2272        Self { set }
2273    }
2274}
2275
2276impl From<GetAccountMetricsInput> for api::metrics::AccountMetricSetRequest {
2277    fn from(value: GetAccountMetricsInput) -> Self {
2278        let (set, start, end, interval) = match value.set {
2279            AccountMetricSet::ActiveBasins(args) => (
2280                api::metrics::AccountMetricSet::ActiveBasins,
2281                args.start,
2282                args.end,
2283                None,
2284            ),
2285            AccountMetricSet::AccountOps(args) => (
2286                api::metrics::AccountMetricSet::AccountOps,
2287                args.start,
2288                args.end,
2289                args.interval,
2290            ),
2291        };
2292        Self {
2293            set,
2294            start: Some(start),
2295            end: Some(end),
2296            interval: interval.map(Into::into),
2297        }
2298    }
2299}
2300
2301#[derive(Debug, Clone, Copy)]
2302/// Basin metric set to return.
2303pub enum BasinMetricSet {
2304    /// Returns a [`GaugeMetric`] representing a timeseries of total stored bytes across all streams
2305    /// in the basin, with one observed value for each hour over the requested time range.
2306    Storage(TimeRange),
2307    /// Returns [`AccumulationMetric`]s, one per storage class (standard, express).
2308    ///
2309    /// Each metric represents a timeseries of the number of append operations across all streams
2310    /// in the basin, with one accumulated value per interval over the requested time range.
2311    ///
2312    /// [`interval`](TimeRangeAndInterval::interval) defaults to
2313    /// [`minute`](TimeseriesInterval::Minute).
2314    AppendOps(TimeRangeAndInterval),
2315    /// Returns [`AccumulationMetric`]s, one per read type (unary, streaming).
2316    ///
2317    /// Each metric represents a timeseries of the number of read operations across all streams
2318    /// in the basin, with one accumulated value per interval over the requested time range.
2319    ///
2320    /// [`interval`](TimeRangeAndInterval::interval) defaults to
2321    /// [`minute`](TimeseriesInterval::Minute).
2322    ReadOps(TimeRangeAndInterval),
2323    /// Returns an [`AccumulationMetric`] representing a timeseries of total read bytes
2324    /// across all streams in the basin, with one accumulated value per interval
2325    /// over the requested time range.
2326    ///
2327    /// [`interval`](TimeRangeAndInterval::interval) defaults to
2328    /// [`minute`](TimeseriesInterval::Minute).
2329    ReadThroughput(TimeRangeAndInterval),
2330    /// Returns an [`AccumulationMetric`] representing a timeseries of total appended bytes
2331    /// across all streams in the basin, with one accumulated value per interval
2332    /// over the requested time range.
2333    ///
2334    /// [`interval`](TimeRangeAndInterval::interval) defaults to
2335    /// [`minute`](TimeseriesInterval::Minute).
2336    AppendThroughput(TimeRangeAndInterval),
2337    /// Returns [`AccumulationMetric`]s, one per basin operation type.
2338    ///
2339    /// Each metric represents a timeseries of the number of operations, with one accumulated value
2340    /// per interval over the requested time range.
2341    ///
2342    /// [`interval`](TimeRangeAndInterval::interval) defaults to [`hour`](TimeseriesInterval::Hour).
2343    BasinOps(TimeRangeAndInterval),
2344}
2345
2346#[derive(Debug, Clone)]
2347#[non_exhaustive]
2348/// Input for [`get_basin_metrics`](crate::S2::get_basin_metrics) operation.
2349pub struct GetBasinMetricsInput {
2350    /// Basin name.
2351    pub name: BasinName,
2352    /// Metric set to return.
2353    pub set: BasinMetricSet,
2354}
2355
2356impl GetBasinMetricsInput {
2357    /// Create a new [`GetBasinMetricsInput`] with the given basin name and metric set.
2358    pub fn new(name: BasinName, set: BasinMetricSet) -> Self {
2359        Self { name, set }
2360    }
2361}
2362
2363impl From<GetBasinMetricsInput> for (BasinName, api::metrics::BasinMetricSetRequest) {
2364    fn from(value: GetBasinMetricsInput) -> Self {
2365        let (set, start, end, interval) = match value.set {
2366            BasinMetricSet::Storage(args) => (
2367                api::metrics::BasinMetricSet::Storage,
2368                args.start,
2369                args.end,
2370                None,
2371            ),
2372            BasinMetricSet::AppendOps(args) => (
2373                api::metrics::BasinMetricSet::AppendOps,
2374                args.start,
2375                args.end,
2376                args.interval,
2377            ),
2378            BasinMetricSet::ReadOps(args) => (
2379                api::metrics::BasinMetricSet::ReadOps,
2380                args.start,
2381                args.end,
2382                args.interval,
2383            ),
2384            BasinMetricSet::ReadThroughput(args) => (
2385                api::metrics::BasinMetricSet::ReadThroughput,
2386                args.start,
2387                args.end,
2388                args.interval,
2389            ),
2390            BasinMetricSet::AppendThroughput(args) => (
2391                api::metrics::BasinMetricSet::AppendThroughput,
2392                args.start,
2393                args.end,
2394                args.interval,
2395            ),
2396            BasinMetricSet::BasinOps(args) => (
2397                api::metrics::BasinMetricSet::BasinOps,
2398                args.start,
2399                args.end,
2400                args.interval,
2401            ),
2402        };
2403        (
2404            value.name,
2405            api::metrics::BasinMetricSetRequest {
2406                set,
2407                start: Some(start),
2408                end: Some(end),
2409                interval: interval.map(Into::into),
2410            },
2411        )
2412    }
2413}
2414
2415#[derive(Debug, Clone, Copy)]
2416/// Stream metric set to return.
2417pub enum StreamMetricSet {
2418    /// Returns a [`GaugeMetric`] representing a timeseries of total stored bytes for the stream,
2419    /// with one observed value for each minute over the requested time range.
2420    Storage(TimeRange),
2421}
2422
2423#[derive(Debug, Clone)]
2424#[non_exhaustive]
2425/// Input for [`get_stream_metrics`](crate::S2::get_stream_metrics) operation.
2426pub struct GetStreamMetricsInput {
2427    /// Basin name.
2428    pub basin_name: BasinName,
2429    /// Stream name.
2430    pub stream_name: StreamName,
2431    /// Metric set to return.
2432    pub set: StreamMetricSet,
2433}
2434
2435impl GetStreamMetricsInput {
2436    /// Create a new [`GetStreamMetricsInput`] with the given basin name, stream name and metric
2437    /// set.
2438    pub fn new(basin_name: BasinName, stream_name: StreamName, set: StreamMetricSet) -> Self {
2439        Self {
2440            basin_name,
2441            stream_name,
2442            set,
2443        }
2444    }
2445}
2446
2447impl From<GetStreamMetricsInput> for (BasinName, StreamName, api::metrics::StreamMetricSetRequest) {
2448    fn from(value: GetStreamMetricsInput) -> Self {
2449        let (set, start, end, interval) = match value.set {
2450            StreamMetricSet::Storage(args) => (
2451                api::metrics::StreamMetricSet::Storage,
2452                args.start,
2453                args.end,
2454                None,
2455            ),
2456        };
2457        (
2458            value.basin_name,
2459            value.stream_name,
2460            api::metrics::StreamMetricSetRequest {
2461                set,
2462                start: Some(start),
2463                end: Some(end),
2464                interval,
2465            },
2466        )
2467    }
2468}
2469
2470#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2471/// Unit in which metric values are measured.
2472pub enum MetricUnit {
2473    /// Size in bytes.
2474    Bytes,
2475    /// Number of operations.
2476    Operations,
2477}
2478
2479impl From<api::metrics::MetricUnit> for MetricUnit {
2480    fn from(value: api::metrics::MetricUnit) -> Self {
2481        match value {
2482            api::metrics::MetricUnit::Bytes => MetricUnit::Bytes,
2483            api::metrics::MetricUnit::Operations => MetricUnit::Operations,
2484        }
2485    }
2486}
2487
2488#[derive(Debug, Clone)]
2489#[non_exhaustive]
2490/// Single named value.
2491pub struct ScalarMetric {
2492    /// Metric name.
2493    pub name: String,
2494    /// Unit for the metric value.
2495    pub unit: MetricUnit,
2496    /// Metric value.
2497    pub value: f64,
2498}
2499
2500#[derive(Debug, Clone)]
2501#[non_exhaustive]
2502/// Named series of `(timestamp, value)` datapoints, each representing an accumulation over a
2503/// specified interval.
2504pub struct AccumulationMetric {
2505    /// Timeseries name.
2506    pub name: String,
2507    /// Unit for the accumulated values.
2508    pub unit: MetricUnit,
2509    /// The interval at which datapoints are accumulated.
2510    pub interval: TimeseriesInterval,
2511    /// Series of `(timestamp, value)` datapoints. Each datapoint represents the accumulated
2512    /// `value` for the time period starting at the `timestamp` (in Unix epoch seconds), spanning
2513    /// one `interval`.
2514    pub values: Vec<(u32, f64)>,
2515}
2516
2517#[derive(Debug, Clone)]
2518#[non_exhaustive]
2519/// Named series of `(timestamp, value)` datapoints, each representing an instantaneous value.
2520pub struct GaugeMetric {
2521    /// Timeseries name.
2522    pub name: String,
2523    /// Unit for the instantaneous values.
2524    pub unit: MetricUnit,
2525    /// Series of `(timestamp, value)` datapoints. Each datapoint represents the `value` at the
2526    /// instant of the `timestamp` (in Unix epoch seconds).
2527    pub values: Vec<(u32, f64)>,
2528}
2529
2530#[derive(Debug, Clone)]
2531#[non_exhaustive]
2532/// Set of string labels.
2533pub struct LabelMetric {
2534    /// Label name.
2535    pub name: String,
2536    /// Label values.
2537    pub values: Vec<String>,
2538}
2539
2540#[derive(Debug, Clone)]
2541/// Individual metric in a returned metric set.
2542pub enum Metric {
2543    /// Single named value.
2544    Scalar(ScalarMetric),
2545    /// Named series of `(timestamp, value)` datapoints, each representing an accumulation over a
2546    /// specified interval.
2547    Accumulation(AccumulationMetric),
2548    /// Named series of `(timestamp, value)` datapoints, each representing an instantaneous value.
2549    Gauge(GaugeMetric),
2550    /// Set of string labels.
2551    Label(LabelMetric),
2552}
2553
2554impl From<api::metrics::Metric> for Metric {
2555    fn from(value: api::metrics::Metric) -> Self {
2556        match value {
2557            api::metrics::Metric::Scalar(sm) => Metric::Scalar(ScalarMetric {
2558                name: sm.name.into(),
2559                unit: sm.unit.into(),
2560                value: sm.value,
2561            }),
2562            api::metrics::Metric::Accumulation(am) => Metric::Accumulation(AccumulationMetric {
2563                name: am.name.into(),
2564                unit: am.unit.into(),
2565                interval: am.interval.into(),
2566                values: am.values,
2567            }),
2568            api::metrics::Metric::Gauge(gm) => Metric::Gauge(GaugeMetric {
2569                name: gm.name.into(),
2570                unit: gm.unit.into(),
2571                values: gm.values,
2572            }),
2573            api::metrics::Metric::Label(lm) => Metric::Label(LabelMetric {
2574                name: lm.name.into(),
2575                values: lm.values,
2576            }),
2577        }
2578    }
2579}
2580
2581#[derive(Debug, Clone, Default)]
2582#[non_exhaustive]
2583/// Input for [`list_streams`](crate::S2Basin::list_streams) operation.
2584pub struct ListStreamsInput {
2585    /// Filter streams whose names begin with this value.
2586    ///
2587    /// Defaults to `""`.
2588    pub prefix: StreamNamePrefix,
2589    /// Filter streams whose names are lexicographically greater than this value.
2590    ///
2591    /// Defaults to `""`.
2592    pub start_after: StreamNameStartAfter,
2593    /// Number of streams to return in a page. Will be clamped to a maximum of `1000`.
2594    ///
2595    /// Defaults to `1000`.
2596    pub limit: Option<usize>,
2597}
2598
2599impl ListStreamsInput {
2600    /// Create a new [`ListStreamsInput`] with default values.
2601    pub fn new() -> Self {
2602        Self::default()
2603    }
2604
2605    /// Set the prefix used to filter streams whose names begin with this value.
2606    pub fn with_prefix(self, prefix: StreamNamePrefix) -> Self {
2607        Self { prefix, ..self }
2608    }
2609
2610    /// Set the value used to filter streams whose names are lexicographically greater than this
2611    /// value.
2612    pub fn with_start_after(self, start_after: StreamNameStartAfter) -> Self {
2613        Self {
2614            start_after,
2615            ..self
2616        }
2617    }
2618
2619    /// Set the limit on number of streams to return in a page.
2620    pub fn with_limit(self, limit: usize) -> Self {
2621        Self {
2622            limit: Some(limit),
2623            ..self
2624        }
2625    }
2626}
2627
2628impl From<ListStreamsInput> for api::stream::ListStreamsRequest {
2629    fn from(value: ListStreamsInput) -> Self {
2630        Self {
2631            prefix: Some(value.prefix),
2632            start_after: Some(value.start_after),
2633            limit: value.limit,
2634        }
2635    }
2636}
2637
2638#[derive(Debug, Clone, Default)]
2639/// Input for [`list_all_streams`](crate::S2Basin::list_all_streams) operation.
2640pub struct ListAllStreamsInput {
2641    /// Filter streams whose names begin with this value.
2642    ///
2643    /// Defaults to `""`.
2644    pub prefix: StreamNamePrefix,
2645    /// Filter streams whose names are lexicographically greater than this value.
2646    ///
2647    /// Defaults to `""`.
2648    pub start_after: StreamNameStartAfter,
2649    /// Whether to include streams that are being deleted.
2650    ///
2651    /// Defaults to `false`.
2652    pub include_deleted: bool,
2653}
2654
2655impl ListAllStreamsInput {
2656    /// Create a new [`ListAllStreamsInput`] with default values.
2657    pub fn new() -> Self {
2658        Self::default()
2659    }
2660
2661    /// Set the prefix used to filter streams whose names begin with this value.
2662    pub fn with_prefix(self, prefix: StreamNamePrefix) -> Self {
2663        Self { prefix, ..self }
2664    }
2665
2666    /// Set the value used to filter streams whose names are lexicographically greater than this
2667    /// value.
2668    pub fn with_start_after(self, start_after: StreamNameStartAfter) -> Self {
2669        Self {
2670            start_after,
2671            ..self
2672        }
2673    }
2674
2675    /// Set whether to include streams that are being deleted.
2676    pub fn with_include_deleted(self, include_deleted: bool) -> Self {
2677        Self {
2678            include_deleted,
2679            ..self
2680        }
2681    }
2682}
2683
2684#[derive(Debug, Clone, PartialEq, Eq)]
2685#[non_exhaustive]
2686/// Stream information.
2687pub struct StreamInfo {
2688    /// Stream name.
2689    pub name: StreamName,
2690    /// Creation time.
2691    pub created_at: S2DateTime,
2692    /// Deletion time if the stream is being deleted.
2693    pub deleted_at: Option<S2DateTime>,
2694    /// Encryption algorithm for this stream, if encryption is enabled.
2695    pub cipher: Option<EncryptionAlgorithm>,
2696}
2697
2698impl TryFrom<api::stream::StreamInfo> for StreamInfo {
2699    type Error = ValidationError;
2700
2701    fn try_from(value: api::stream::StreamInfo) -> Result<Self, Self::Error> {
2702        Ok(Self {
2703            name: value.name,
2704            created_at: value.created_at.try_into()?,
2705            deleted_at: value.deleted_at.map(S2DateTime::try_from).transpose()?,
2706            cipher: value.cipher.map(Into::into),
2707        })
2708    }
2709}
2710
2711#[derive(Debug, Clone)]
2712#[non_exhaustive]
2713/// Input for [`create_stream`](crate::S2Basin::create_stream) operation.
2714pub struct CreateStreamInput {
2715    /// Stream name.
2716    pub name: StreamName,
2717    /// Configuration for the stream.
2718    ///
2719    /// See [`StreamConfig`] for defaults.
2720    pub config: Option<StreamConfig>,
2721    idempotency_token: String,
2722}
2723
2724impl CreateStreamInput {
2725    /// Create a new [`CreateStreamInput`] with the given stream name.
2726    pub fn new(name: StreamName) -> Self {
2727        Self {
2728            name,
2729            config: None,
2730            idempotency_token: idempotency_token(),
2731        }
2732    }
2733
2734    /// Set the configuration for the stream.
2735    pub fn with_config(self, config: StreamConfig) -> Self {
2736        Self {
2737            config: Some(config),
2738            ..self
2739        }
2740    }
2741}
2742
2743impl From<CreateStreamInput> for (api::stream::CreateStreamRequest, String) {
2744    fn from(value: CreateStreamInput) -> Self {
2745        (
2746            api::stream::CreateStreamRequest {
2747                stream: value.name,
2748                config: value.config.map(Into::into),
2749            },
2750            value.idempotency_token,
2751        )
2752    }
2753}
2754
2755#[derive(Debug, Clone)]
2756#[non_exhaustive]
2757/// Input for [`ensure_stream`](crate::S2Basin::ensure_stream)
2758/// operation.
2759pub struct EnsureStreamInput {
2760    /// Stream name.
2761    pub name: StreamName,
2762    /// Configuration for the stream.
2763    ///
2764    /// See [`StreamConfig`] for defaults.
2765    pub config: Option<StreamConfig>,
2766}
2767
2768impl EnsureStreamInput {
2769    /// Create a new [`EnsureStreamInput`] with the given stream name.
2770    pub fn new(name: StreamName) -> Self {
2771        Self { name, config: None }
2772    }
2773
2774    /// Set the configuration for the stream.
2775    pub fn with_config(self, config: StreamConfig) -> Self {
2776        Self {
2777            config: Some(config),
2778            ..self
2779        }
2780    }
2781}
2782
2783impl From<EnsureStreamInput> for (StreamName, Option<api::config::StreamConfig>) {
2784    fn from(value: EnsureStreamInput) -> Self {
2785        (value.name, value.config.map(Into::into))
2786    }
2787}
2788
2789#[derive(Debug, Clone)]
2790#[non_exhaustive]
2791/// Input of [`delete_stream`](crate::S2Basin::delete_stream) operation.
2792pub struct DeleteStreamInput {
2793    /// Stream name.
2794    pub name: StreamName,
2795    /// Whether to ignore `Not Found` error if the stream doesn't exist.
2796    pub ignore_not_found: bool,
2797}
2798
2799impl DeleteStreamInput {
2800    /// Create a new [`DeleteStreamInput`] with the given stream name.
2801    pub fn new(name: StreamName) -> Self {
2802        Self {
2803            name,
2804            ignore_not_found: false,
2805        }
2806    }
2807
2808    /// Set whether to ignore `Not Found` error if the stream doesn't exist.
2809    pub fn with_ignore_not_found(self, ignore_not_found: bool) -> Self {
2810        Self {
2811            ignore_not_found,
2812            ..self
2813        }
2814    }
2815}
2816
2817#[derive(Debug, Clone)]
2818#[non_exhaustive]
2819/// Input for [`reconfigure_stream`](crate::S2Basin::reconfigure_stream) operation.
2820pub struct ReconfigureStreamInput {
2821    /// Stream name.
2822    pub name: StreamName,
2823    /// Reconfiguration for [`StreamConfig`].
2824    pub config: StreamReconfiguration,
2825}
2826
2827impl ReconfigureStreamInput {
2828    /// Create a new [`ReconfigureStreamInput`] with the given stream name and reconfiguration.
2829    pub fn new(name: StreamName, config: StreamReconfiguration) -> Self {
2830        Self { name, config }
2831    }
2832}
2833
2834#[derive(Debug, Clone, PartialEq, Eq)]
2835/// Token for fencing appends to a stream.
2836///
2837/// **Note:** It must not exceed 36 bytes in length.
2838///
2839/// See [`CommandRecord::fence`] and [`AppendInput::fencing_token`].
2840pub struct FencingToken(String);
2841
2842impl FencingToken {
2843    pub(crate) fn from_server(value: String) -> Self {
2844        Self(value)
2845    }
2846
2847    /// Generate a random alphanumeric fencing token of `n` bytes.
2848    pub fn generate(n: usize) -> Result<Self, ValidationError> {
2849        rand::rng()
2850            .sample_iter(&rand::distr::Alphanumeric)
2851            .take(n)
2852            .map(char::from)
2853            .collect::<String>()
2854            .parse()
2855    }
2856}
2857
2858impl FromStr for FencingToken {
2859    type Err = ValidationError;
2860
2861    fn from_str(s: &str) -> Result<Self, Self::Err> {
2862        if s.len() > MAX_FENCING_TOKEN_LENGTH {
2863            return Err(ValidationError(format!(
2864                "fencing token exceeds {MAX_FENCING_TOKEN_LENGTH} bytes in length",
2865            )));
2866        }
2867        Ok(FencingToken(s.to_string()))
2868    }
2869}
2870
2871impl std::fmt::Display for FencingToken {
2872    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2873        write!(f, "{}", self.0)
2874    }
2875}
2876
2877impl Deref for FencingToken {
2878    type Target = str;
2879
2880    fn deref(&self) -> &Self::Target {
2881        &self.0
2882    }
2883}
2884
2885#[derive(Debug, Clone, Copy, PartialEq)]
2886#[non_exhaustive]
2887/// A position in a stream.
2888pub struct StreamPosition {
2889    /// Sequence number assigned by the service.
2890    pub seq_num: u64,
2891    /// Timestamp. When assigned by the service, represents milliseconds since Unix epoch.
2892    /// User-specified timestamps are passed through as-is.
2893    pub timestamp: u64,
2894}
2895
2896impl StreamPosition {
2897    /// Construct a stream position.
2898    ///
2899    /// This is intended for building fixtures in downstream tests.
2900    pub fn new(seq_num: u64, timestamp: u64) -> Self {
2901        Self { seq_num, timestamp }
2902    }
2903}
2904
2905impl std::fmt::Display for StreamPosition {
2906    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2907        write!(f, "seq_num={}, timestamp={}", self.seq_num, self.timestamp)
2908    }
2909}
2910
2911impl From<api::stream::proto::StreamPosition> for StreamPosition {
2912    fn from(value: api::stream::proto::StreamPosition) -> Self {
2913        Self {
2914            seq_num: value.seq_num,
2915            timestamp: value.timestamp,
2916        }
2917    }
2918}
2919
2920impl From<api::stream::StreamPosition> for StreamPosition {
2921    fn from(value: api::stream::StreamPosition) -> Self {
2922        Self {
2923            seq_num: value.seq_num,
2924            timestamp: value.timestamp,
2925        }
2926    }
2927}
2928
2929#[derive(Debug, Clone, PartialEq)]
2930#[non_exhaustive]
2931/// A name-value pair.
2932pub struct Header {
2933    /// Name.
2934    pub name: Bytes,
2935    /// Value.
2936    pub value: Bytes,
2937}
2938
2939impl Header {
2940    /// Create a new [`Header`] with the given name and value.
2941    pub fn new(name: impl Into<Bytes>, value: impl Into<Bytes>) -> Self {
2942        Self {
2943            name: name.into(),
2944            value: value.into(),
2945        }
2946    }
2947}
2948
2949impl From<Header> for api::stream::proto::Header {
2950    fn from(value: Header) -> Self {
2951        Self {
2952            name: value.name,
2953            value: value.value,
2954        }
2955    }
2956}
2957
2958impl From<api::stream::proto::Header> for Header {
2959    fn from(value: api::stream::proto::Header) -> Self {
2960        Self {
2961            name: value.name,
2962            value: value.value,
2963        }
2964    }
2965}
2966
2967#[derive(Debug, Clone, PartialEq)]
2968/// A record to append.
2969pub struct AppendRecord {
2970    body: Bytes,
2971    headers: Vec<Header>,
2972    timestamp: Option<u64>,
2973}
2974
2975impl AppendRecord {
2976    fn validate(self) -> Result<Self, ValidationError> {
2977        if self.metered_bytes() > RECORD_BATCH_MAX.bytes {
2978            Err(ValidationError(format!(
2979                "metered_bytes: {} exceeds {}",
2980                self.metered_bytes(),
2981                RECORD_BATCH_MAX.bytes
2982            )))
2983        } else {
2984            Ok(self)
2985        }
2986    }
2987
2988    /// Create a new [`AppendRecord`] with the given record body.
2989    pub fn new(body: impl Into<Bytes>) -> Result<Self, ValidationError> {
2990        let record = Self {
2991            body: body.into(),
2992            headers: Vec::default(),
2993            timestamp: None,
2994        };
2995        record.validate()
2996    }
2997
2998    /// Set the headers for this record.
2999    pub fn with_headers(
3000        self,
3001        headers: impl IntoIterator<Item = Header>,
3002    ) -> Result<Self, ValidationError> {
3003        let record = Self {
3004            headers: headers.into_iter().collect(),
3005            ..self
3006        };
3007        record.validate()
3008    }
3009
3010    /// Set the timestamp for this record.
3011    ///
3012    /// Precise semantics depend on [`StreamConfig::timestamping`].
3013    pub fn with_timestamp(self, timestamp: u64) -> Self {
3014        Self {
3015            timestamp: Some(timestamp),
3016            ..self
3017        }
3018    }
3019
3020    /// Get the body of this record.
3021    pub fn body(&self) -> &[u8] {
3022        &self.body
3023    }
3024
3025    /// Get the headers of this record.
3026    pub fn headers(&self) -> &[Header] {
3027        &self.headers
3028    }
3029
3030    /// Get the timestamp of this record.
3031    pub fn timestamp(&self) -> Option<u64> {
3032        self.timestamp
3033    }
3034}
3035
3036impl From<AppendRecord> for api::stream::proto::AppendRecord {
3037    fn from(value: AppendRecord) -> Self {
3038        Self {
3039            timestamp: value.timestamp,
3040            headers: value.headers.into_iter().map(Into::into).collect(),
3041            body: value.body,
3042        }
3043    }
3044}
3045
3046/// Metered byte size calculation.
3047///
3048/// Formula for a record:
3049/// ```text
3050/// 8 + 2 * len(headers) + sum(len(h.name) + len(h.value) for h in headers) + len(body)
3051/// ```
3052pub trait MeteredBytes {
3053    /// Returns the metered byte size.
3054    fn metered_bytes(&self) -> usize;
3055}
3056
3057macro_rules! metered_bytes_impl {
3058    ($ty:ty) => {
3059        impl MeteredBytes for $ty {
3060            fn metered_bytes(&self) -> usize {
3061                8 + (2 * self.headers.len())
3062                    + self
3063                        .headers
3064                        .iter()
3065                        .map(|h| h.name.len() + h.value.len())
3066                        .sum::<usize>()
3067                    + self.body.len()
3068            }
3069        }
3070    };
3071}
3072
3073metered_bytes_impl!(AppendRecord);
3074
3075impl MeteredSize for AppendRecord {
3076    fn metered_size(&self) -> usize {
3077        self.metered_bytes()
3078    }
3079}
3080
3081#[derive(Debug, Clone)]
3082/// A batch of records to append atomically.
3083///
3084/// **Note:** It must contain at least `1` record and no more than `1000`.
3085/// The total size of the batch must not exceed `1MiB` in metered bytes.
3086///
3087/// See [`AppendRecordBatches`](crate::batching::AppendRecordBatches) and
3088/// [`AppendInputs`](crate::batching::AppendInputs) for convenient and automatic batching of records
3089/// that takes care of the abovementioned constraints.
3090pub struct AppendRecordBatch(Metered<Vec<AppendRecord>>);
3091
3092impl From<Metered<Vec<AppendRecord>>> for AppendRecordBatch {
3093    fn from(records: Metered<Vec<AppendRecord>>) -> Self {
3094        Self(records)
3095    }
3096}
3097
3098impl AppendRecordBatch {
3099    /// Try to create an [`AppendRecordBatch`] from an iterator of [`AppendRecord`]s.
3100    pub fn try_from_iter<I>(iter: I) -> Result<Self, ValidationError>
3101    where
3102        I: IntoIterator<Item = AppendRecord>,
3103    {
3104        let mut records = Metered::with_capacity(RECORD_BATCH_MAX.count);
3105
3106        for record in iter {
3107            records.push(Metered::from(record));
3108
3109            if records.metered_size() > RECORD_BATCH_MAX.bytes {
3110                return Err(ValidationError(format!(
3111                    "batch size in metered bytes ({}) exceeds {}",
3112                    records.metered_size(),
3113                    RECORD_BATCH_MAX.bytes
3114                )));
3115            }
3116
3117            if records.len() > RECORD_BATCH_MAX.count {
3118                return Err(ValidationError(format!(
3119                    "number of records in the batch exceeds {}",
3120                    RECORD_BATCH_MAX.count
3121                )));
3122            }
3123        }
3124
3125        if records.is_empty() {
3126            return Err(ValidationError("batch is empty".into()));
3127        }
3128
3129        Ok(records.into())
3130    }
3131}
3132
3133impl Deref for AppendRecordBatch {
3134    type Target = [AppendRecord];
3135
3136    fn deref(&self) -> &Self::Target {
3137        &self.0[..]
3138    }
3139}
3140
3141impl MeteredBytes for AppendRecordBatch {
3142    fn metered_bytes(&self) -> usize {
3143        self.0.metered_size()
3144    }
3145}
3146
3147impl IntoIterator for AppendRecordBatch {
3148    type Item = AppendRecord;
3149    type IntoIter = std::vec::IntoIter<AppendRecord>;
3150
3151    fn into_iter(self) -> Self::IntoIter {
3152        self.0.into_iter()
3153    }
3154}
3155
3156impl<'a> IntoIterator for &'a AppendRecordBatch {
3157    type Item = &'a AppendRecord;
3158    type IntoIter = std::slice::Iter<'a, AppendRecord>;
3159
3160    fn into_iter(self) -> Self::IntoIter {
3161        self.0.iter()
3162    }
3163}
3164
3165#[derive(Debug, Clone)]
3166/// Command to signal an operation.
3167pub enum Command {
3168    /// Fence operation.
3169    Fence {
3170        /// Fencing token.
3171        fencing_token: FencingToken,
3172    },
3173    /// Trim operation.
3174    Trim {
3175        /// Trim point.
3176        trim_point: u64,
3177    },
3178}
3179
3180#[derive(Debug, Clone)]
3181#[non_exhaustive]
3182/// Command record for signaling operations to the service.
3183///
3184/// See [here](https://s2.dev/docs/rest/records/overview#command-records) for more information.
3185pub struct CommandRecord {
3186    /// Command to signal an operation.
3187    pub command: Command,
3188    /// Timestamp for this record.
3189    pub timestamp: Option<u64>,
3190}
3191
3192impl CommandRecord {
3193    const FENCE: &[u8] = b"fence";
3194    const TRIM: &[u8] = b"trim";
3195
3196    /// Create a fence command record with the given fencing token.
3197    ///
3198    /// Fencing is strongly consistent, and subsequent appends that specify a
3199    /// fencing token will fail if it does not match.
3200    pub fn fence(fencing_token: FencingToken) -> Self {
3201        Self {
3202            command: Command::Fence { fencing_token },
3203            timestamp: None,
3204        }
3205    }
3206
3207    /// Create a trim command record with the given trim point.
3208    ///
3209    /// Trim point is the desired earliest sequence number for the stream.
3210    ///
3211    /// Trimming is eventually consistent, and trimmed records may be visible
3212    /// for a brief period.
3213    pub fn trim(trim_point: u64) -> Self {
3214        Self {
3215            command: Command::Trim { trim_point },
3216            timestamp: None,
3217        }
3218    }
3219
3220    /// Set the timestamp for this record.
3221    pub fn with_timestamp(self, timestamp: u64) -> Self {
3222        Self {
3223            timestamp: Some(timestamp),
3224            ..self
3225        }
3226    }
3227}
3228
3229impl From<CommandRecord> for AppendRecord {
3230    fn from(value: CommandRecord) -> Self {
3231        let (header_value, body) = match value.command {
3232            Command::Fence { fencing_token } => (
3233                CommandRecord::FENCE,
3234                Bytes::copy_from_slice(fencing_token.as_bytes()),
3235            ),
3236            Command::Trim { trim_point } => (
3237                CommandRecord::TRIM,
3238                Bytes::copy_from_slice(&trim_point.to_be_bytes()),
3239            ),
3240        };
3241        Self {
3242            body,
3243            headers: vec![Header::new("", header_value)],
3244            timestamp: value.timestamp,
3245        }
3246    }
3247}
3248
3249#[derive(Debug, Clone)]
3250#[non_exhaustive]
3251/// Input for [`append`](crate::S2Stream::append) operation and
3252/// [`AppendSession::submit`](crate::append_session::AppendSession::submit).
3253pub struct AppendInput {
3254    /// Batch of records to append atomically.
3255    pub records: AppendRecordBatch,
3256    /// Expected sequence number for the first record in the batch.
3257    ///
3258    /// If unspecified, no matching is performed. If specified and mismatched, the append fails.
3259    pub match_seq_num: Option<u64>,
3260    /// Fencing token to match against the stream's current fencing token.
3261    ///
3262    /// If unspecified, no matching is performed. If specified and mismatched,
3263    /// the append fails. A stream defaults to `""` as its fencing token.
3264    pub fencing_token: Option<FencingToken>,
3265}
3266
3267impl AppendInput {
3268    /// Create a new [`AppendInput`] with the given batch of records.
3269    pub fn new(records: AppendRecordBatch) -> Self {
3270        Self {
3271            records,
3272            match_seq_num: None,
3273            fencing_token: None,
3274        }
3275    }
3276
3277    /// Set the expected sequence number for the first record in the batch.
3278    pub fn with_match_seq_num(self, match_seq_num: u64) -> Self {
3279        Self {
3280            match_seq_num: Some(match_seq_num),
3281            ..self
3282        }
3283    }
3284
3285    /// Set the fencing token to match against the stream's current fencing token.
3286    pub fn with_fencing_token(self, fencing_token: FencingToken) -> Self {
3287        Self {
3288            fencing_token: Some(fencing_token),
3289            ..self
3290        }
3291    }
3292}
3293
3294impl From<AppendInput> for api::stream::proto::AppendInput {
3295    fn from(value: AppendInput) -> Self {
3296        Self {
3297            records: value.records.iter().cloned().map(Into::into).collect(),
3298            match_seq_num: value.match_seq_num,
3299            fencing_token: value.fencing_token.map(|t| t.to_string()),
3300        }
3301    }
3302}
3303
3304#[derive(Debug, Clone, PartialEq)]
3305#[non_exhaustive]
3306/// Acknowledgement for an [`AppendInput`].
3307pub struct AppendAck {
3308    /// Sequence number and timestamp of the first record that was appended.
3309    pub start: StreamPosition,
3310    /// Sequence number of the last record that was appended + 1, and timestamp of the last record
3311    /// that was appended.
3312    ///
3313    /// The difference between `end.seq_num` and `start.seq_num` will be the number of records
3314    /// appended.
3315    pub end: StreamPosition,
3316    /// Sequence number that will be assigned to the next record on the stream, and timestamp of
3317    /// the last record on the stream.
3318    ///
3319    /// This can be greater than the `end` position in case of concurrent appends.
3320    pub tail: StreamPosition,
3321}
3322
3323impl AppendAck {
3324    /// Construct an append acknowledgement.
3325    ///
3326    /// This is intended for building fixtures in downstream tests.
3327    pub fn new(start: StreamPosition, end: StreamPosition, tail: StreamPosition) -> Self {
3328        Self { start, end, tail }
3329    }
3330}
3331
3332impl From<api::stream::proto::AppendAck> for AppendAck {
3333    fn from(value: api::stream::proto::AppendAck) -> Self {
3334        Self {
3335            start: value.start.unwrap_or_default().into(),
3336            end: value.end.unwrap_or_default().into(),
3337            tail: value.tail.unwrap_or_default().into(),
3338        }
3339    }
3340}
3341
3342#[derive(Debug, Clone, Copy)]
3343/// Starting position for reading from a stream.
3344pub enum ReadFrom {
3345    /// Read from this sequence number.
3346    SeqNum(u64),
3347    /// Read from this timestamp.
3348    Timestamp(u64),
3349    /// Read from N records before the tail.
3350    TailOffset(u64),
3351}
3352
3353impl Default for ReadFrom {
3354    fn default() -> Self {
3355        Self::SeqNum(0)
3356    }
3357}
3358
3359#[derive(Debug, Default, Clone)]
3360#[non_exhaustive]
3361/// Where to start reading.
3362pub struct ReadStart {
3363    /// Starting position.
3364    ///
3365    /// Defaults to reading from sequence number `0`.
3366    pub from: ReadFrom,
3367    /// Whether to start from tail if the requested starting position is beyond it.
3368    ///
3369    /// Defaults to `false` (errors if position is beyond tail).
3370    pub clamp_to_tail: bool,
3371}
3372
3373impl ReadStart {
3374    /// Create a new [`ReadStart`] with default values.
3375    pub fn new() -> Self {
3376        Self::default()
3377    }
3378
3379    /// Set the starting position.
3380    pub fn with_from(self, from: ReadFrom) -> Self {
3381        Self { from, ..self }
3382    }
3383
3384    /// Set whether to start from tail if the requested starting position is beyond it.
3385    pub fn with_clamp_to_tail(self, clamp_to_tail: bool) -> Self {
3386        Self {
3387            clamp_to_tail,
3388            ..self
3389        }
3390    }
3391}
3392
3393impl From<ReadStart> for api::stream::ReadStart {
3394    fn from(value: ReadStart) -> Self {
3395        let (seq_num, timestamp, tail_offset) = match value.from {
3396            ReadFrom::SeqNum(n) => (Some(n), None, None),
3397            ReadFrom::Timestamp(t) => (None, Some(t), None),
3398            ReadFrom::TailOffset(o) => (None, None, Some(o)),
3399        };
3400        Self {
3401            seq_num,
3402            timestamp,
3403            tail_offset,
3404            clamp: if value.clamp_to_tail {
3405                Some(true)
3406            } else {
3407                None
3408            },
3409        }
3410    }
3411}
3412
3413#[derive(Debug, Clone, Default)]
3414#[non_exhaustive]
3415/// Limits on how much to read.
3416pub struct ReadLimits {
3417    /// Limit on number of records.
3418    ///
3419    /// Defaults to `1000` for non-streaming read.
3420    pub count: Option<usize>,
3421    /// Limit on total metered bytes of records.
3422    ///
3423    /// Defaults to `1MiB` for non-streaming read.
3424    pub bytes: Option<usize>,
3425}
3426
3427impl ReadLimits {
3428    /// Create a new [`ReadLimits`] with default values.
3429    pub fn new() -> Self {
3430        Self::default()
3431    }
3432
3433    /// Set the limit on number of records.
3434    pub fn with_count(self, count: usize) -> Self {
3435        Self {
3436            count: Some(count),
3437            ..self
3438        }
3439    }
3440
3441    /// Set the limit on total metered bytes of records.
3442    pub fn with_bytes(self, bytes: usize) -> Self {
3443        Self {
3444            bytes: Some(bytes),
3445            ..self
3446        }
3447    }
3448}
3449
3450#[derive(Debug, Clone, Default)]
3451#[non_exhaustive]
3452/// When to stop reading.
3453pub struct ReadStop {
3454    /// Limits on how much to read.
3455    ///
3456    /// See [`ReadLimits`] for defaults.
3457    pub limits: ReadLimits,
3458    /// Timestamp at which to stop (exclusive).
3459    ///
3460    /// Defaults to `None`.
3461    pub until: Option<RangeTo<u64>>,
3462    /// Duration in seconds to wait for new records before stopping. Will be clamped to `60`
3463    /// seconds for [`read`](crate::S2Stream::read).
3464    ///
3465    /// Defaults to:
3466    /// - `0` (no wait) for [`read`](crate::S2Stream::read).
3467    /// - `0` (no wait) for [`read_session`](crate::S2Stream::read_session) if `limits` or `until`
3468    ///   is specified.
3469    /// - Infinite wait for [`read_session`](crate::S2Stream::read_session) if neither `limits` nor
3470    ///   `until` is specified.
3471    pub wait: Option<u32>,
3472}
3473
3474impl ReadStop {
3475    /// Create a new [`ReadStop`] with default values.
3476    pub fn new() -> Self {
3477        Self::default()
3478    }
3479
3480    /// Set the limits on how much to read.
3481    pub fn with_limits(self, limits: ReadLimits) -> Self {
3482        Self { limits, ..self }
3483    }
3484
3485    /// Set the timestamp at which to stop (exclusive).
3486    pub fn with_until(self, until: RangeTo<u64>) -> Self {
3487        Self {
3488            until: Some(until),
3489            ..self
3490        }
3491    }
3492
3493    /// Set the duration in seconds to wait for new records before stopping.
3494    pub fn with_wait(self, wait: u32) -> Self {
3495        Self {
3496            wait: Some(wait),
3497            ..self
3498        }
3499    }
3500}
3501
3502impl From<ReadStop> for api::stream::ReadEnd {
3503    fn from(value: ReadStop) -> Self {
3504        Self {
3505            count: value.limits.count,
3506            bytes: value.limits.bytes,
3507            until: value.until.map(|r| r.end),
3508            wait: value.wait,
3509        }
3510    }
3511}
3512
3513#[derive(Debug, Clone, Default)]
3514#[non_exhaustive]
3515/// Input for [`read`](crate::S2Stream::read) and [`read_session`](crate::S2Stream::read_session)
3516/// operations.
3517pub struct ReadInput {
3518    /// Where to start reading.
3519    ///
3520    /// See [`ReadStart`] for defaults.
3521    pub start: ReadStart,
3522    /// When to stop reading.
3523    ///
3524    /// See [`ReadStop`] for defaults.
3525    pub stop: ReadStop,
3526    /// Whether to filter out command records from the stream when reading.
3527    ///
3528    /// Defaults to `false`.
3529    pub ignore_command_records: bool,
3530}
3531
3532#[derive(Debug, Clone, Copy, Default, Eq, PartialEq)]
3533#[non_exhaustive]
3534/// Retry policy for a continuous read session.
3535pub enum ReadSessionRetryPolicy {
3536    /// Stop after the retry budget configured by [`RetryConfig`] is exhausted.
3537    #[default]
3538    Budgeted,
3539    /// Keep retrying retryable failures after the configured retry budget is exhausted.
3540    ///
3541    /// This also applies while establishing the initial session, so
3542    /// [`read_session`](crate::S2Stream::read_session) may remain pending through retryable
3543    /// failures until it connects or the future is cancelled.
3544    Indefinite,
3545}
3546
3547#[derive(Debug, Clone, Default)]
3548#[non_exhaustive]
3549/// Configuration for a continuous read session.
3550pub struct ReadSessionConfig {
3551    /// Policy for retrying retryable failures.
3552    ///
3553    /// Clean stream ends and non-retryable failures always terminate the session.
3554    ///
3555    /// Defaults to [`ReadSessionRetryPolicy::Budgeted`].
3556    pub retry_policy: ReadSessionRetryPolicy,
3557}
3558
3559impl ReadSessionConfig {
3560    /// Create a new [`ReadSessionConfig`] with default settings.
3561    pub fn new() -> Self {
3562        Self::default()
3563    }
3564
3565    /// Set the policy for retrying retryable failures.
3566    pub fn with_retry_policy(self, retry_policy: ReadSessionRetryPolicy) -> Self {
3567        Self {
3568            retry_policy,
3569            ..self
3570        }
3571    }
3572}
3573
3574impl ReadInput {
3575    /// Create a new [`ReadInput`] with default values.
3576    pub fn new() -> Self {
3577        Self::default()
3578    }
3579
3580    /// Set where to start reading.
3581    pub fn with_start(self, start: ReadStart) -> Self {
3582        Self { start, ..self }
3583    }
3584
3585    /// Set when to stop reading.
3586    pub fn with_stop(self, stop: ReadStop) -> Self {
3587        Self { stop, ..self }
3588    }
3589
3590    /// Set whether to filter out command records from the stream when reading.
3591    pub fn with_ignore_command_records(self, ignore_command_records: bool) -> Self {
3592        Self {
3593            ignore_command_records,
3594            ..self
3595        }
3596    }
3597}
3598
3599#[derive(Debug, Clone)]
3600#[non_exhaustive]
3601/// Record that is durably sequenced on a stream.
3602pub struct SequencedRecord {
3603    /// Sequence number assigned to this record.
3604    pub seq_num: u64,
3605    /// Body of this record.
3606    pub body: Bytes,
3607    /// Headers for this record.
3608    pub headers: Vec<Header>,
3609    /// Timestamp for this record.
3610    pub timestamp: u64,
3611}
3612
3613impl SequencedRecord {
3614    /// Construct a sequenced record from its plain-data fields.
3615    ///
3616    /// This is intended for building fixtures in downstream tests.
3617    pub fn from_parts(
3618        seq_num: u64,
3619        timestamp: u64,
3620        headers: Vec<Header>,
3621        body: impl Into<Bytes>,
3622    ) -> Self {
3623        Self {
3624            seq_num,
3625            timestamp,
3626            body: body.into(),
3627            headers,
3628        }
3629    }
3630
3631    /// Whether this is a command record.
3632    pub fn is_command_record(&self) -> bool {
3633        self.headers.len() == 1 && *self.headers[0].name == *b""
3634    }
3635}
3636
3637impl From<api::stream::proto::SequencedRecord> for SequencedRecord {
3638    fn from(value: api::stream::proto::SequencedRecord) -> Self {
3639        Self {
3640            seq_num: value.seq_num,
3641            body: value.body,
3642            headers: value.headers.into_iter().map(Into::into).collect(),
3643            timestamp: value.timestamp,
3644        }
3645    }
3646}
3647
3648metered_bytes_impl!(SequencedRecord);
3649
3650#[derive(Debug, Clone)]
3651#[non_exhaustive]
3652/// Batch of records returned by [`read`](crate::S2Stream::read) or streamed by
3653/// [`read_session`](crate::S2Stream::read_session).
3654pub struct ReadBatch {
3655    /// Records that are durably sequenced on the stream.
3656    ///
3657    /// It can be empty only for a [`read`](crate::S2Stream::read) operation when:
3658    /// - the [`stop condition`](ReadInput::stop) was already met, or
3659    /// - all records in the batch were command records and
3660    ///   [`ignore_command_records`](ReadInput::ignore_command_records) was set to `true`.
3661    pub records: Vec<SequencedRecord>,
3662    /// Sequence number that will be assigned to the next record on the stream, and timestamp of
3663    /// the last record.
3664    ///
3665    /// It will only be present when reading recent records.
3666    pub tail: Option<StreamPosition>,
3667}
3668
3669impl ReadBatch {
3670    /// Construct a read batch.
3671    ///
3672    /// This is intended for building fixtures in downstream tests.
3673    pub fn new(records: Vec<SequencedRecord>, tail: Option<StreamPosition>) -> Self {
3674        Self { records, tail }
3675    }
3676
3677    pub(crate) fn from_api(batch: api::stream::proto::ReadBatch) -> Self {
3678        Self {
3679            records: batch.records.into_iter().map(Into::into).collect(),
3680            tail: batch.tail.map(Into::into),
3681        }
3682    }
3683}
3684
3685/// A stream of values of type `Result<T, RequestError>`.
3686pub type Streaming<T> = Pin<Box<dyn Send + futures_core::Stream<Item = Result<T, RequestError>>>>;
3687
3688fn idempotency_token() -> String {
3689    uuid::Uuid::new_v4().simple().to_string()
3690}
3691
3692#[cfg(test)]
3693mod tests {
3694    use proptest::prelude::*;
3695    use rstest::rstest;
3696
3697    use super::*;
3698
3699    type HeaderParts = (Vec<u8>, Vec<u8>);
3700    type AppendRecordParts = (Vec<u8>, Vec<HeaderParts>);
3701
3702    fn byte_vec_strategy(max_len: usize) -> impl Strategy<Value = Vec<u8>> {
3703        prop::collection::vec(any::<u8>(), 0..=max_len)
3704    }
3705
3706    fn header_parts_strategy() -> impl Strategy<Value = HeaderParts> {
3707        (byte_vec_strategy(32), byte_vec_strategy(64))
3708    }
3709
3710    fn string_strategy(max_chars: usize) -> impl Strategy<Value = String> {
3711        prop::collection::vec(any::<char>(), 0..=max_chars)
3712            .prop_map(|chars| chars.into_iter().collect())
3713    }
3714
3715    fn read_from_strategy() -> impl Strategy<Value = ReadFrom> {
3716        prop_oneof![
3717            any::<u64>().prop_map(ReadFrom::SeqNum),
3718            any::<u64>().prop_map(ReadFrom::Timestamp),
3719            any::<u64>().prop_map(ReadFrom::TailOffset),
3720        ]
3721    }
3722
3723    fn append_record_parts_strategy() -> impl Strategy<Value = AppendRecordParts> {
3724        (
3725            byte_vec_strategy(256),
3726            prop::collection::vec(header_parts_strategy(), 0..=16),
3727        )
3728    }
3729
3730    fn proto_stream_position_strategy() -> impl Strategy<Value = api::stream::proto::StreamPosition>
3731    {
3732        (any::<u64>(), any::<u64>()).prop_map(|(seq_num, timestamp)| {
3733            api::stream::proto::StreamPosition { seq_num, timestamp }
3734        })
3735    }
3736
3737    fn headers_from_parts(headers: &[HeaderParts]) -> Vec<Header> {
3738        headers
3739            .iter()
3740            .map(|(name, value)| Header::new(Bytes::from(name.clone()), Bytes::from(value.clone())))
3741            .collect()
3742    }
3743
3744    fn expected_metered_bytes(body: &[u8], headers: &[HeaderParts]) -> usize {
3745        8 + (2 * headers.len())
3746            + headers
3747                .iter()
3748                .map(|(name, value)| name.len() + value.len())
3749                .sum::<usize>()
3750            + body.len()
3751    }
3752
3753    // -- S2DateTime --
3754
3755    #[test]
3756    fn s2_datetime_parse_valid_rfc3339() {
3757        let dt: S2DateTime = "2024-01-15T12:30:00Z".parse().unwrap();
3758        assert_eq!(dt.to_string(), "2024-01-15T12:30:00Z");
3759    }
3760
3761    #[test]
3762    fn s2_datetime_parse_with_offset() {
3763        let dt: S2DateTime = "2024-06-01T08:00:00+05:30".parse().unwrap();
3764        assert_eq!(dt.to_string(), "2024-06-01T08:00:00+05:30");
3765
3766        let offset_dt: time::OffsetDateTime = dt.into();
3767        assert_eq!(
3768            offset_dt.offset(),
3769            time::UtcOffset::from_hms(5, 30, 0).unwrap()
3770        );
3771    }
3772
3773    #[test]
3774    fn s2_datetime_parse_invalid() {
3775        let err = "not-a-date".parse::<S2DateTime>();
3776        assert!(err.is_err());
3777    }
3778
3779    #[test]
3780    fn s2_datetime_roundtrip_via_offset_datetime() {
3781        let odt = time::OffsetDateTime::from_unix_timestamp(1_700_000_000).unwrap();
3782        let dt = S2DateTime::try_from(odt).unwrap();
3783        let back: time::OffsetDateTime = dt.into();
3784        assert_eq!(odt, back);
3785    }
3786
3787    // -- AccountEndpoint --
3788
3789    #[rstest]
3790    #[case::https_with_scheme("https://aws.s2.dev", Scheme::HTTPS)]
3791    #[case::http_with_scheme("http://localhost:8080", Scheme::HTTP)]
3792    #[case::default_https("aws.s2.dev", Scheme::HTTPS)]
3793    fn account_endpoint_parse(#[case] input: &str, #[case] expected_scheme: Scheme) {
3794        let ep: AccountEndpoint = input.parse().unwrap();
3795        assert_eq!(ep.scheme, expected_scheme);
3796    }
3797
3798    // -- BasinEndpoint --
3799
3800    #[rstest]
3801    #[case::https_parent_zone("https://{basin}.b.s2.dev", Scheme::HTTPS, true)]
3802    #[case::http_direct("http://localhost:8080", Scheme::HTTP, false)]
3803    #[case::default_https_parent_zone("{basin}.b.s2.dev", Scheme::HTTPS, true)]
3804    fn basin_endpoint_parse(
3805        #[case] input: &str,
3806        #[case] expected_scheme: Scheme,
3807        #[case] expected_parent_zone: bool,
3808    ) {
3809        let ep: BasinEndpoint = input.parse().unwrap();
3810        assert_eq!(ep.scheme, expected_scheme);
3811        assert_eq!(
3812            matches!(ep.authority, BasinAuthority::ParentZone(_)),
3813            expected_parent_zone
3814        );
3815    }
3816
3817    // -- S2Endpoints --
3818
3819    #[test]
3820    fn s2_endpoints_new_requires_same_scheme() {
3821        let account: AccountEndpoint = "https://aws.s2.dev".parse().unwrap();
3822        let basin: BasinEndpoint = "http://localhost:8080".parse().unwrap();
3823        let err = S2Endpoints::new(account, basin);
3824        assert!(err.is_err());
3825    }
3826
3827    #[test]
3828    fn s2_endpoints_new_same_scheme_succeeds() {
3829        let account: AccountEndpoint = "https://aws.s2.dev".parse().unwrap();
3830        let basin: BasinEndpoint = "https://{basin}.b.s2.dev".parse().unwrap();
3831        let ep = S2Endpoints::new(account, basin).unwrap();
3832        assert_eq!(ep.scheme, Scheme::HTTPS);
3833    }
3834
3835    #[test]
3836    fn s2_endpoints_for_endpoint_defaults_to_https() {
3837        let ep = S2Endpoints::for_endpoint("localhost:8080").unwrap();
3838        let authority: Authority = "localhost:8080".parse().unwrap();
3839        assert_eq!(ep.scheme, Scheme::HTTPS);
3840        assert_eq!(ep.account_authority, authority);
3841        assert_eq!(ep.basin_authority, BasinAuthority::Direct(authority));
3842    }
3843
3844    #[test]
3845    fn s2_endpoints_for_endpoint_accepts_explicit_scheme() {
3846        let ep = S2Endpoints::for_endpoint("http://localhost:8080").unwrap();
3847        assert_eq!(ep.scheme, Scheme::HTTP);
3848    }
3849
3850    #[test]
3851    fn s2_endpoints_for_endpoint_rejects_invalid_endpoint() {
3852        assert!(S2Endpoints::for_endpoint("not a valid endpoint").is_err());
3853    }
3854
3855    // -- Compression --
3856
3857    #[rstest]
3858    #[case::none(Compression::None, CompressionAlgorithm::None)]
3859    #[case::gzip(Compression::Gzip, CompressionAlgorithm::Gzip)]
3860    #[case::zstd(Compression::Zstd, CompressionAlgorithm::Zstd)]
3861    fn compression_conversion(#[case] sdk: Compression, #[case] api: CompressionAlgorithm) {
3862        assert_eq!(CompressionAlgorithm::from(sdk), api);
3863    }
3864
3865    // -- RetryConfig --
3866
3867    #[test]
3868    fn retry_config_defaults() {
3869        let rc = RetryConfig::default();
3870        assert_eq!(rc.max_attempts.get(), 3);
3871        assert_eq!(rc.min_base_delay, Duration::from_millis(100));
3872        assert_eq!(rc.max_base_delay, Duration::from_secs(1));
3873        assert!(matches!(rc.append_retry_policy, AppendRetryPolicy::All));
3874    }
3875
3876    #[test]
3877    fn retry_config_max_retries() {
3878        let rc = RetryConfig::default();
3879        assert_eq!(rc.max_retries(), 2);
3880    }
3881
3882    // -- S2Config --
3883
3884    #[test]
3885    fn s2_config_defaults() {
3886        let cfg = S2Config::new("test-token");
3887        assert_eq!(cfg.connection_timeout, Duration::from_secs(3));
3888        assert_eq!(cfg.request_timeout, Duration::from_secs(5));
3889        assert!(!cfg.insecure_skip_cert_verification);
3890    }
3891
3892    // -- StorageClass --
3893
3894    #[rstest]
3895    #[case::standard(StorageClass::Standard)]
3896    #[case::express(StorageClass::Express)]
3897    fn storage_class_roundtrip(#[case] sdk: StorageClass) {
3898        let api: api::config::StorageClass = sdk.into();
3899        let back: StorageClass = api.into();
3900        assert_eq!(back, sdk);
3901    }
3902
3903    // -- RetentionPolicy --
3904
3905    #[rstest]
3906    #[case::age(RetentionPolicy::Age(3600))]
3907    #[case::infinite(RetentionPolicy::Infinite)]
3908    fn retention_policy_roundtrip(#[case] sdk: RetentionPolicy) {
3909        let api: api::config::RetentionPolicy = sdk.into();
3910        let back: RetentionPolicy = api.into();
3911        assert_eq!(back, sdk);
3912    }
3913
3914    // -- TimestampingMode --
3915
3916    #[rstest]
3917    #[case::client_prefer(
3918        TimestampingMode::ClientPrefer,
3919        api::config::TimestampingMode::ClientPrefer
3920    )]
3921    #[case::client_require(
3922        TimestampingMode::ClientRequire,
3923        api::config::TimestampingMode::ClientRequire
3924    )]
3925    #[case::arrival(TimestampingMode::Arrival, api::config::TimestampingMode::Arrival)]
3926    fn timestamping_mode_roundtrip(
3927        #[case] sdk: TimestampingMode,
3928        #[case] expected_api: api::config::TimestampingMode,
3929    ) {
3930        let converted: api::config::TimestampingMode = sdk.into();
3931        assert_eq!(converted, expected_api);
3932        let back: TimestampingMode = converted.into();
3933        assert_eq!(back, sdk);
3934    }
3935
3936    // -- TimestampingConfig --
3937
3938    #[test]
3939    fn timestamping_config_roundtrip() {
3940        let sdk = TimestampingConfig {
3941            mode: Some(TimestampingMode::Arrival),
3942            uncapped: Some(true),
3943        };
3944        let api: api::config::TimestampingConfig = sdk.into();
3945        let back: TimestampingConfig = api.into();
3946        assert_eq!(back, sdk);
3947    }
3948
3949    // -- DeleteOnEmptyConfig --
3950
3951    #[test]
3952    fn delete_on_empty_config_roundtrip() {
3953        let sdk = DeleteOnEmptyConfig::new().with_min_age(Duration::from_secs(300));
3954        let api: api::config::DeleteOnEmptyConfig = sdk.into();
3955        let back: DeleteOnEmptyConfig = api.into();
3956        assert_eq!(back, sdk);
3957    }
3958
3959    // -- StreamConfig --
3960
3961    #[test]
3962    fn stream_config_builder_and_roundtrip() {
3963        let sdk = StreamConfig::new()
3964            .with_storage_class(StorageClass::Express)
3965            .with_retention_policy(RetentionPolicy::Age(86400))
3966            .with_timestamping(TimestampingConfig {
3967                mode: Some(TimestampingMode::ClientPrefer),
3968                uncapped: None,
3969            })
3970            .with_delete_on_empty(DeleteOnEmptyConfig { min_age_secs: 60 });
3971        let api: api::config::StreamConfig = sdk.clone().into();
3972        let back: StreamConfig = api.into();
3973        assert_eq!(back, sdk);
3974    }
3975
3976    // -- BasinConfig --
3977
3978    #[test]
3979    fn basin_config_builder_and_roundtrip() {
3980        let sdk = BasinConfig::new()
3981            .with_default_stream_config(
3982                StreamConfig::new().with_storage_class(StorageClass::Standard),
3983            )
3984            .with_create_stream_on_append(true)
3985            .with_create_stream_on_read(false);
3986        let api: api::config::BasinConfig = sdk.clone().into();
3987        let back: BasinConfig = api.into();
3988        assert_eq!(back, sdk);
3989    }
3990
3991    // -- FencingToken --
3992
3993    proptest! {
3994        #[test]
3995        fn fencing_token_parse_accepts_only_within_byte_limit(
3996            token in string_strategy(MAX_FENCING_TOKEN_LENGTH + 8),
3997        ) {
3998            let parsed = token.parse::<FencingToken>();
3999
4000            if token.len() <= MAX_FENCING_TOKEN_LENGTH {
4001                prop_assert_eq!(parsed.unwrap().to_string(), token);
4002            } else {
4003                prop_assert!(parsed.is_err());
4004            }
4005        }
4006    }
4007
4008    // -- StreamPosition --
4009
4010    #[test]
4011    fn stream_position_display() {
4012        let pos = StreamPosition {
4013            seq_num: 42,
4014            timestamp: 1700000000,
4015        };
4016        assert_eq!(pos.to_string(), "seq_num=42, timestamp=1700000000");
4017    }
4018
4019    proptest! {
4020        #[test]
4021        fn stream_position_conversions_preserve_values(seq_num in any::<u64>(), timestamp in any::<u64>()) {
4022            let proto: StreamPosition = api::stream::proto::StreamPosition {
4023                seq_num,
4024                timestamp,
4025            }
4026            .into();
4027            prop_assert_eq!(proto.seq_num, seq_num);
4028            prop_assert_eq!(proto.timestamp, timestamp);
4029
4030            let api: StreamPosition = api::stream::StreamPosition {
4031                seq_num,
4032                timestamp,
4033            }
4034            .into();
4035            prop_assert_eq!(api.seq_num, seq_num);
4036            prop_assert_eq!(api.timestamp, timestamp);
4037        }
4038    }
4039
4040    // -- Header --
4041
4042    proptest! {
4043        #[test]
4044        fn header_proto_roundtrip_preserves_binary_parts(
4045            name in byte_vec_strategy(64),
4046            value in byte_vec_strategy(128),
4047        ) {
4048            let header = Header::new(Bytes::from(name.clone()), Bytes::from(value.clone()));
4049            let proto: api::stream::proto::Header = header.into();
4050            let back: Header = proto.into();
4051
4052            prop_assert_eq!(back.name.as_ref(), name.as_slice());
4053            prop_assert_eq!(back.value.as_ref(), value.as_slice());
4054        }
4055    }
4056
4057    // -- AppendRecord --
4058
4059    #[test]
4060    fn append_record_too_large() {
4061        let big_body = vec![0u8; RECORD_BATCH_MAX.bytes + 1];
4062        assert!(AppendRecord::new(big_body).is_err());
4063    }
4064
4065    // -- MeteredBytes --
4066
4067    proptest! {
4068        #[test]
4069        fn append_record_preserves_fields_and_metered_byte_formula(
4070            (body, headers) in append_record_parts_strategy(),
4071            timestamp in proptest::option::of(any::<u64>()),
4072        ) {
4073            let mut record = AppendRecord::new(body.clone())
4074                .unwrap()
4075                .with_headers(headers_from_parts(&headers))
4076                .unwrap();
4077            if let Some(timestamp) = timestamp {
4078                record = record.with_timestamp(timestamp);
4079            }
4080
4081            prop_assert_eq!(record.body(), body.as_slice());
4082            prop_assert_eq!(record.headers().len(), headers.len());
4083            prop_assert_eq!(record.timestamp(), timestamp);
4084            prop_assert_eq!(record.metered_bytes(), expected_metered_bytes(&body, &headers));
4085
4086            for (actual, (expected_name, expected_value)) in record.headers().iter().zip(headers.iter()) {
4087                prop_assert_eq!(actual.name.as_ref(), expected_name.as_slice());
4088                prop_assert_eq!(actual.value.as_ref(), expected_value.as_slice());
4089            }
4090        }
4091    }
4092
4093    // -- AppendRecordBatch --
4094
4095    #[test]
4096    fn append_record_batch_empty_is_err() {
4097        let result = AppendRecordBatch::try_from_iter(vec![]);
4098        assert!(result.is_err());
4099    }
4100
4101    #[test]
4102    fn append_record_batch_too_many_records() {
4103        let records: Vec<_> = (0..1001).map(|_| AppendRecord::new("x").unwrap()).collect();
4104        let result = AppendRecordBatch::try_from_iter(records);
4105        assert!(result.is_err());
4106    }
4107
4108    proptest! {
4109        #[test]
4110        fn append_record_batch_metered_bytes_is_sum_of_records(
4111            records in prop::collection::vec(append_record_parts_strategy(), 1..=32),
4112        ) {
4113            let expected = records
4114                .iter()
4115                .map(|(body, headers)| expected_metered_bytes(body, headers))
4116                .sum::<usize>();
4117            let records = records
4118                .into_iter()
4119                .map(|(body, headers)| {
4120                    AppendRecord::new(body)
4121                        .unwrap()
4122                        .with_headers(headers_from_parts(&headers))
4123                        .unwrap()
4124                })
4125                .collect::<Vec<_>>();
4126
4127            let batch = AppendRecordBatch::try_from_iter(records).unwrap();
4128            prop_assert_eq!(batch.metered_bytes(), expected);
4129            prop_assert_eq!(batch.iter().map(MeteredBytes::metered_bytes).sum::<usize>(), expected);
4130        }
4131    }
4132
4133    // -- CommandRecord --
4134
4135    #[test]
4136    fn command_record_fence() {
4137        let token: FencingToken = "tok".parse().unwrap();
4138        let cmd = CommandRecord::fence(token);
4139        let record: AppendRecord = cmd.into();
4140        assert_eq!(record.headers().len(), 1);
4141        assert_eq!(record.headers()[0].name.as_ref(), b"");
4142        assert_eq!(record.headers()[0].value.as_ref(), b"fence");
4143        assert_eq!(record.body(), b"tok");
4144    }
4145
4146    #[test]
4147    fn command_record_trim() {
4148        let cmd = CommandRecord::trim(42);
4149        let record: AppendRecord = cmd.into();
4150        assert_eq!(record.headers().len(), 1);
4151        assert_eq!(record.headers()[0].value.as_ref(), b"trim");
4152        assert_eq!(record.body(), &42u64.to_be_bytes());
4153    }
4154
4155    // -- SequencedRecord --
4156
4157    #[rstest]
4158    #[case::command(vec![Header::new("", "fence")], true)]
4159    #[case::regular(vec![Header::new("key", "value")], false)]
4160    #[case::no_headers(vec![], false)]
4161    fn sequenced_record_command_detection(#[case] headers: Vec<Header>, #[case] expected: bool) {
4162        let record = SequencedRecord {
4163            seq_num: 0,
4164            body: Bytes::from("data"),
4165            headers,
4166            timestamp: 0,
4167        };
4168        assert_eq!(record.is_command_record(), expected);
4169    }
4170
4171    // -- ReadStart --
4172
4173    proptest! {
4174        #[test]
4175        fn read_start_to_api_sets_only_selected_position_field(
4176            from in read_from_strategy(),
4177            clamp_to_tail in any::<bool>(),
4178        ) {
4179            let (seq_num, timestamp, tail_offset) = match from {
4180                ReadFrom::SeqNum(value) => (Some(value), None, None),
4181                ReadFrom::Timestamp(value) => (None, Some(value), None),
4182                ReadFrom::TailOffset(value) => (None, None, Some(value)),
4183            };
4184            let api: api::stream::ReadStart = ReadStart::new()
4185                .with_from(from)
4186                .with_clamp_to_tail(clamp_to_tail)
4187                .into();
4188
4189            prop_assert_eq!(api.seq_num, seq_num);
4190            prop_assert_eq!(api.timestamp, timestamp);
4191            prop_assert_eq!(api.tail_offset, tail_offset);
4192            prop_assert_eq!(api.clamp, clamp_to_tail.then_some(true));
4193        }
4194    }
4195
4196    // -- ReadStop --
4197
4198    #[test]
4199    fn read_stop_to_api() {
4200        let stop = ReadStop::new()
4201            .with_limits(ReadLimits::new().with_count(50))
4202            .with_until(..1000)
4203            .with_wait(30);
4204        let api: api::stream::ReadEnd = stop.into();
4205        assert_eq!(api.count, Some(50));
4206        assert_eq!(api.until, Some(1000));
4207        assert_eq!(api.wait, Some(30));
4208    }
4209
4210    // -- Operation roundtrip --
4211
4212    #[test]
4213    fn operation_roundtrip_all_variants() {
4214        let variants = [
4215            Operation::ListBasins,
4216            Operation::CreateBasin,
4217            Operation::GetBasinConfig,
4218            Operation::DeleteBasin,
4219            Operation::ReconfigureBasin,
4220            Operation::ListAccessTokens,
4221            Operation::IssueAccessToken,
4222            Operation::RevokeAccessToken,
4223            Operation::GetAccountMetrics,
4224            Operation::GetBasinMetrics,
4225            Operation::GetStreamMetrics,
4226            Operation::ListStreams,
4227            Operation::CreateStream,
4228            Operation::GetStreamConfig,
4229            Operation::DeleteStream,
4230            Operation::ReconfigureStream,
4231            Operation::CheckTail,
4232            Operation::Append,
4233            Operation::Read,
4234            Operation::Trim,
4235            Operation::Fence,
4236            Operation::ListLocations,
4237            Operation::GetDefaultLocation,
4238            Operation::SetDefaultLocation,
4239        ];
4240        for op in variants {
4241            let api_op: api::access::Operation = op.into();
4242            let back: Operation = api_op.into();
4243            assert_eq!(back, op);
4244        }
4245    }
4246
4247    // -- MetricUnit --
4248
4249    #[test]
4250    fn metric_unit_conversion() {
4251        assert_eq!(
4252            MetricUnit::from(api::metrics::MetricUnit::Bytes),
4253            MetricUnit::Bytes
4254        );
4255        assert_eq!(
4256            MetricUnit::from(api::metrics::MetricUnit::Operations),
4257            MetricUnit::Operations
4258        );
4259    }
4260
4261    // -- AppendAck --
4262
4263    proptest! {
4264        #[test]
4265        fn append_ack_from_proto_preserves_present_positions_and_defaults_missing(
4266            start in proptest::option::of(proto_stream_position_strategy()),
4267            end in proptest::option::of(proto_stream_position_strategy()),
4268            tail in proptest::option::of(proto_stream_position_strategy()),
4269        ) {
4270            let expected_start = start.unwrap_or_default();
4271            let expected_end = end.unwrap_or_default();
4272            let expected_tail = tail.unwrap_or_default();
4273            let ack: AppendAck = api::stream::proto::AppendAck { start, end, tail }.into();
4274
4275            prop_assert_eq!(ack.start.seq_num, expected_start.seq_num);
4276            prop_assert_eq!(ack.start.timestamp, expected_start.timestamp);
4277            prop_assert_eq!(ack.end.seq_num, expected_end.seq_num);
4278            prop_assert_eq!(ack.end.timestamp, expected_end.timestamp);
4279            prop_assert_eq!(ack.tail.seq_num, expected_tail.seq_num);
4280            prop_assert_eq!(ack.tail.timestamp, expected_tail.timestamp);
4281        }
4282    }
4283
4284    // -- ReadBatch --
4285
4286    #[test]
4287    fn read_batch_from_api() {
4288        let proto_batch = api::stream::proto::ReadBatch {
4289            records: vec![api::stream::proto::SequencedRecord {
4290                seq_num: 0,
4291                body: Bytes::from("hi"),
4292                headers: vec![api::stream::proto::Header {
4293                    name: Bytes::from("k"),
4294                    value: Bytes::from("v"),
4295                }],
4296                timestamp: 42,
4297            }],
4298            tail: Some(api::stream::proto::StreamPosition {
4299                seq_num: 1,
4300                timestamp: 42,
4301            }),
4302        };
4303        let batch = ReadBatch::from_api(proto_batch);
4304        assert_eq!(batch.records.len(), 1);
4305        assert_eq!(batch.records[0].seq_num, 0);
4306        assert_eq!(batch.records[0].timestamp, 42);
4307        assert_eq!(batch.records[0].body.as_ref(), b"hi");
4308        assert_eq!(batch.records[0].headers.len(), 1);
4309        assert_eq!(batch.records[0].headers[0].name.as_ref(), b"k");
4310        assert_eq!(batch.records[0].headers[0].value.as_ref(), b"v");
4311        assert_eq!(
4312            batch.tail,
4313            Some(StreamPosition {
4314                seq_num: 1,
4315                timestamp: 42,
4316            })
4317        );
4318    }
4319
4320    // -- CreateBasinInput --
4321
4322    #[test]
4323    fn create_basin_input_to_api() {
4324        let name: BasinName = "test-basin-name".parse().unwrap();
4325        let input = CreateBasinInput::new(name.clone()).with_config(BasinConfig::new());
4326        let (req, token): (api::basin::CreateBasinRequest, String) = input.into();
4327        assert_eq!(req.basin, name);
4328        assert!(req.config.is_some());
4329        assert!(!token.is_empty());
4330    }
4331
4332    // -- CreateStreamInput --
4333
4334    #[test]
4335    fn create_stream_input_to_api() {
4336        let name: StreamName = "my-stream".parse().unwrap();
4337        let input = CreateStreamInput::new(name.clone()).with_config(StreamConfig::new());
4338        let (req, token): (api::stream::CreateStreamRequest, String) = input.into();
4339        assert_eq!(req.stream, name);
4340        assert!(req.config.is_some());
4341        assert!(!token.is_empty());
4342    }
4343
4344    // -- SequencedRecord from proto --
4345
4346    #[test]
4347    fn sequenced_record_from_proto() {
4348        let proto = api::stream::proto::SequencedRecord {
4349            seq_num: 99,
4350            body: Bytes::from("data"),
4351            headers: vec![api::stream::proto::Header {
4352                name: Bytes::from("k"),
4353                value: Bytes::from("v"),
4354            }],
4355            timestamp: 1234,
4356        };
4357        let record: SequencedRecord = proto.into();
4358        assert_eq!(record.seq_num, 99);
4359        assert_eq!(record.body.as_ref(), b"data");
4360        assert_eq!(record.headers.len(), 1);
4361        assert_eq!(record.headers[0].name.as_ref(), b"k");
4362        assert_eq!(record.headers[0].value.as_ref(), b"v");
4363        assert_eq!(record.timestamp, 1234);
4364    }
4365}