Skip to main content

S2Endpoints

Struct S2Endpoints 

Source
#[non_exhaustive]
pub struct S2Endpoints { /* private fields */ }
Expand description

Endpoints for the S2 environment.

Implementations§

Source§

impl S2Endpoints

Source

pub fn new( account_endpoint: AccountEndpoint, basin_endpoint: BasinEndpoint, ) -> Result<Self, ValidationError>

Create a new S2Endpoints with the given account and basin endpoints.

Examples found in repository?
examples/docs_configuration.rs (lines 17-20)
12fn main() -> Result<(), Box<dyn std::error::Error>> {
13    // Example: Custom endpoints (e.g., for s2-lite local dev)
14    {
15        // ANCHOR: custom-endpoints
16        let client = S2::new(
17            S2Config::new("local-token").with_endpoints(S2Endpoints::new(
18                AccountEndpoint::new("http://localhost:8080")?,
19                BasinEndpoint::new("http://localhost:8080")?,
20            )?),
21        )?;
22        // ANCHOR_END: custom-endpoints
23        println!("Created client with custom endpoints: {:?}", client);
24    }
25
26    // Example: Custom retry configuration
27    {
28        let access_token = std::env::var("S2_ACCESS_TOKEN").unwrap_or_else(|_| "demo".into());
29        // ANCHOR: retry-config
30        let client = S2::new(
31            S2Config::new(access_token).with_retry(
32                RetryConfig::new()
33                    .with_max_attempts(NonZeroU32::new(5).unwrap())
34                    .with_min_base_delay(Duration::from_millis(100))
35                    .with_max_base_delay(Duration::from_secs(2)),
36            ),
37        )?;
38        // ANCHOR_END: retry-config
39        println!("Created client with retry config: {:?}", client);
40    }
41
42    // Example: Custom timeout configuration
43    {
44        let access_token = std::env::var("S2_ACCESS_TOKEN").unwrap_or_else(|_| "demo".into());
45        // ANCHOR: timeout-config
46        let client = S2::new(
47            S2Config::new(access_token)
48                .with_connection_timeout(Duration::from_secs(5))
49                .with_request_timeout(Duration::from_secs(10)),
50        )?;
51        // ANCHOR_END: timeout-config
52        println!("Created client with timeout config: {:?}", client);
53    }
54
55    Ok(())
56}
Source

pub fn for_endpoint(endpoint: &str) -> Result<Self, ValidationError>

Create endpoints for a single account and basin endpoint.

This is useful for S2-compatible services that expose both APIs at one endpoint.

Source

pub fn from_env() -> Result<Self, ValidationError>

Create a new S2Endpoints from environment variables.

The following environment variables are expected to be set:

  • S2_ACCOUNT_ENDPOINT - Account-level endpoint.
  • S2_BASIN_ENDPOINT - Basin-level endpoint.
Examples found in repository?
examples/caught_up.rs (line 28)
21async fn main() -> Result<(), Box<dyn std::error::Error>> {
22    let access_token =
23        std::env::var("S2_ACCESS_TOKEN").map_err(|_| "S2_ACCESS_TOKEN env var not set")?;
24    let mut config = S2Config::new(access_token);
25    if std::env::var_os("S2_ACCOUNT_ENDPOINT").is_some()
26        || std::env::var_os("S2_BASIN_ENDPOINT").is_some()
27    {
28        config = config.with_endpoints(S2Endpoints::from_env()?);
29    }
30
31    let suffix = &uuid::Uuid::new_v4().simple().to_string()[..8];
32    let basin_name: BasinName = format!("caught-up-{suffix}").parse()?;
33    let stream_name: StreamName = "example".parse()?;
34    let s2 = S2::new(config)?;
35    let basin = s2.basin(basin_name.clone());
36
37    s2.create_basin(CreateBasinInput::new(basin_name.clone()))
38        .await?;
39    basin
40        .create_stream(CreateStreamInput::new(stream_name.clone()))
41        .await?;
42    let stream = basin.stream(stream_name.clone());
43
44    stream
45        .append(AppendInput::new(AppendRecordBatch::try_from_iter([
46            AppendRecord::new("first")?,
47            AppendRecord::new("second")?,
48        ])?))
49        .await?;
50
51    let mut session = stream
52        .read_session(
53            ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::TailOffset(2))),
54            ReadSessionConfig::default(),
55        )
56        .await?;
57    let mut caught_up = session.caught_up();
58
59    loop {
60        tokio::select! {
61            tail = &mut caught_up => {
62                println!("Caught up through sequence number {}", tail?.seq_num);
63                break;
64            }
65            Some(batch) = session.next() => {
66                print_batch("Read before catching up", &batch?);
67            }
68        }
69    }
70
71    let ack = stream
72        .append(AppendInput::new(AppendRecordBatch::try_from_iter([
73            AppendRecord::new("third")?,
74        ])?))
75        .await?;
76    println!(
77        "Appended another record at sequence number {}",
78        ack.start.seq_num
79    );
80
81    while let Some(batch) = session.next().await {
82        let batch = batch?;
83        print_batch("Read after catching up", &batch);
84        if batch
85            .records
86            .iter()
87            .any(|record| record.seq_num == ack.start.seq_num)
88        {
89            break;
90        }
91    }
92    println!("Session is caught up again: {}", session.is_caught_up());
93
94    drop(session);
95    basin
96        .delete_stream(DeleteStreamInput::new(stream_name))
97        .await?;
98    s2.delete_basin(DeleteBasinInput::new(basin_name)).await?;
99
100    Ok(())
101}
Source

pub fn for_cloud() -> Self

Return the default S2 Cloud endpoints.

Trait Implementations§

Source§

impl Clone for S2Endpoints

Source§

fn clone(&self) -> S2Endpoints

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Debug for S2Endpoints

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more