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