#[non_exhaustive]pub struct S2Endpoints { /* private fields */ }Expand description
Endpoints for the S2 environment.
Implementations§
Source§impl S2Endpoints
impl S2Endpoints
Sourcepub fn new(
account_endpoint: AccountEndpoint,
basin_endpoint: BasinEndpoint,
) -> Result<Self, ValidationError>
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}Sourcepub fn for_endpoint(endpoint: &str) -> Result<Self, ValidationError>
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.
Sourcepub fn from_env() -> Result<Self, ValidationError>
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}Trait Implementations§
Source§impl Clone for S2Endpoints
impl Clone for S2Endpoints
Source§fn clone(&self) -> S2Endpoints
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)
fn clone_from(&mut self, source: &Self)
Performs copy-assignment from
source. Read moreAuto Trait Implementations§
impl !Freeze for S2Endpoints
impl RefUnwindSafe for S2Endpoints
impl Send for S2Endpoints
impl Sync for S2Endpoints
impl Unpin for S2Endpoints
impl UnsafeUnpin for S2Endpoints
impl UnwindSafe for S2Endpoints
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self>
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 moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
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