pub struct S2Basin { /* private fields */ }Expand description
A basin in an S2 account.
See S2::basin.
Implementations§
Source§impl S2Basin
impl S2Basin
Sourcepub fn stream(&self, name: StreamName) -> S2Stream
pub fn stream(&self, name: StreamName) -> S2Stream
Get an S2Stream.
Examples found in repository?
examples/docs_overview.rs (line 12)
5fn main() -> Result<(), Box<dyn std::error::Error>> {
6 // ANCHOR: create-client
7 use s2_sdk::{S2, types::S2Config};
8
9 let client = S2::new(S2Config::new(std::env::var("S2_ACCESS_TOKEN")?))?;
10
11 let basin = client.basin("my-basin".parse()?);
12 let stream = basin.stream("my-stream".parse()?);
13 // ANCHOR_END: create-client
14
15 println!("Created client for stream: {:?}", stream);
16 Ok(())
17}More examples
examples/consumer.rs (line 15)
9async fn main() -> Result<(), Box<dyn std::error::Error>> {
10 let access_token = std::env::var("S2_ACCESS_TOKEN")?;
11 let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
12 let stream_name: StreamName = std::env::var("S2_STREAM")?.parse()?;
13
14 let s2 = S2::new(S2Config::new(access_token))?;
15 let stream = s2.basin(basin_name).stream(stream_name);
16
17 let input = ReadInput::new();
18 let mut batches = stream
19 .read_session(input, ReadSessionConfig::default())
20 .await?;
21 loop {
22 select! {
23 batch = batches.next() => {
24 let Some(batch) = batch else { break };
25 let batch = batch?;
26 println!("{batch:?}");
27 }
28 _ = tokio::signal::ctrl_c() => break,
29 }
30 }
31
32 Ok(())
33}examples/get_latest_record.rs (line 20)
9async fn main() -> Result<(), Box<dyn std::error::Error>> {
10 let access_token =
11 std::env::var("S2_ACCESS_TOKEN").map_err(|_| "S2_ACCESS_TOKEN env var not set")?;
12 let basin_name: BasinName = std::env::var("S2_BASIN")
13 .map_err(|_| "S2_BASIN env var not set")?
14 .parse()?;
15 let stream_name: StreamName = std::env::var("S2_STREAM")
16 .map_err(|_| "S2_STREAM env var not set")?
17 .parse()?;
18
19 let s2 = S2::new(S2Config::new(access_token))?;
20 let stream = s2.basin(basin_name).stream(stream_name);
21
22 let input = ReadInput::new()
23 .with_start(ReadStart::new().with_from(ReadFrom::TailOffset(1)))
24 .with_stop(ReadStop::new().with_limits(ReadLimits::new().with_count(1)));
25 let batch = stream.read(input).await?;
26 println!("{batch:#?}");
27
28 Ok(())
29}examples/explicit_trim.rs (line 18)
7async fn main() -> Result<(), Box<dyn std::error::Error>> {
8 let access_token =
9 std::env::var("S2_ACCESS_TOKEN").map_err(|_| "S2_ACCESS_TOKEN env var not set")?;
10 let basin_name: BasinName = std::env::var("S2_BASIN")
11 .map_err(|_| "S2_BASIN env var not set")?
12 .parse()?;
13 let stream_name: StreamName = std::env::var("S2_STREAM")
14 .map_err(|_| "S2_STREAM env var not set")?
15 .parse()?;
16
17 let s2 = S2::new(S2Config::new(access_token))?;
18 let stream = s2.basin(basin_name).stream(stream_name);
19
20 let tail = stream.check_tail().await?;
21 if tail.seq_num == 0 {
22 println!("Empty stream");
23 return Ok(());
24 }
25
26 let input = AppendInput::new(AppendRecordBatch::try_from_iter([CommandRecord::trim(
27 tail.seq_num - 1,
28 )
29 .into()])?);
30 stream.append(input).await?;
31 println!("Trim requested");
32
33 Ok(())
34}examples/producer.rs (line 19)
8async fn main() -> Result<(), Box<dyn std::error::Error>> {
9 let access_token =
10 std::env::var("S2_ACCESS_TOKEN").map_err(|_| "S2_ACCESS_TOKEN env var not set")?;
11 let basin_name: BasinName = std::env::var("S2_BASIN")
12 .map_err(|_| "S2_BASIN env var not set")?
13 .parse()?;
14 let stream_name: StreamName = std::env::var("S2_STREAM")
15 .map_err(|_| "S2_STREAM env var not set")?
16 .parse()?;
17
18 let s2 = S2::new(S2Config::new(access_token))?;
19 let stream = s2.basin(basin_name).stream(stream_name);
20
21 let producer = stream.producer(ProducerConfig::new());
22
23 let ticket1 = producer.submit(AppendRecord::new("lorem")?).await?;
24 let ticket2 = producer.submit(AppendRecord::new("ipsum")?).await?;
25
26 let ack1 = ticket1.await?;
27 let ack2 = ticket2.await?;
28 println!("Record 1 seq_num: {}", ack1.seq_num);
29 println!("Record 2 seq_num: {}", ack2.seq_num);
30
31 producer.close().await?;
32
33 Ok(())
34}examples/docs_encryption.rs (line 51)
15async fn main() -> Result<(), Box<dyn std::error::Error>> {
16 let access_token = std::env::var("S2_ACCESS_TOKEN")?;
17 let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
18 let stream_name: StreamName = format!(
19 "docs-encryption-{}",
20 std::time::SystemTime::now()
21 .duration_since(std::time::UNIX_EPOCH)?
22 .as_millis()
23 )
24 .parse()?;
25
26 let client = S2::new(S2Config::new(access_token))?;
27
28 // ANCHOR: basin-cipher
29 client
30 .create_basin(
31 CreateBasinInput::new(basin_name.clone())
32 .with_config(BasinConfig::new().with_stream_cipher(EncryptionAlgorithm::Aegis256)),
33 )
34 .await?;
35
36 client
37 .reconfigure_basin(ReconfigureBasinInput::new(
38 basin_name.clone(),
39 BasinReconfiguration::new().with_stream_cipher(EncryptionAlgorithm::Aes256Gcm),
40 ))
41 .await?;
42 // ANCHOR_END: basin-cipher
43
44 let basin = client.basin(basin_name.clone());
45 basin
46 .create_stream(CreateStreamInput::new(stream_name.clone()))
47 .await?;
48
49 // ANCHOR: append-read
50 let stream = basin
51 .stream(stream_name.clone())
52 .with_encryption_key(std::env::var("S2_ENCRYPTION_KEY")?.parse()?);
53
54 stream
55 .append(AppendInput::new(AppendRecordBatch::try_from_iter([
56 AppendRecord::new("top secret")?,
57 ])?))
58 .await?;
59
60 let batch = stream
61 .read(
62 ReadInput::new()
63 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
64 .with_stop(ReadStop::new().with_limits(ReadLimits::new().with_count(10))),
65 )
66 .await?;
67 // ANCHOR_END: append-read
68
69 println!("Read {} encrypted record(s)", batch.records.len());
70
71 basin
72 .delete_stream(DeleteStreamInput::new(stream_name))
73 .await?;
74
75 Ok(())
76}Additional examples can be found in:
Sourcepub async fn list_streams(
&self,
input: ListStreamsInput,
) -> Result<Page<StreamInfo>, RequestError>
pub async fn list_streams( &self, input: ListStreamsInput, ) -> Result<Page<StreamInfo>, RequestError>
List a page of streams.
See list_all_streams for automatic pagination.
Examples found in repository?
examples/list_streams.rs (line 18)
7async fn main() -> Result<(), Box<dyn std::error::Error>> {
8 let access_token =
9 std::env::var("S2_ACCESS_TOKEN").map_err(|_| "S2_ACCESS_TOKEN env var not set")?;
10 let basin_name: BasinName = std::env::var("S2_BASIN")
11 .map_err(|_| "S2_BASIN env var not set")?
12 .parse()?;
13
14 let s2 = S2::new(S2Config::new(access_token))?;
15 let basin = s2.basin(basin_name);
16
17 let input = ListStreamsInput::new().with_prefix("my-".parse()?);
18 let page = basin.list_streams(input).await?;
19 println!("{page:#?}");
20
21 Ok(())
22}More examples
examples/docs_account_and_basins.rs (line 47)
17async fn main() -> Result<(), Box<dyn std::error::Error>> {
18 let access_token = std::env::var("S2_ACCESS_TOKEN")?;
19 let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
20
21 let client = S2::new(S2Config::new(access_token))?;
22
23 // ANCHOR: basin-operations
24 // List basins
25 let basins = client.list_basins(ListBasinsInput::new()).await?;
26
27 // Create a basin
28 client
29 .create_basin(CreateBasinInput::new("my-events".parse()?))
30 .await?;
31
32 // Get configuration
33 let config = client.get_basin_config("my-events".parse()?).await?;
34
35 // Delete
36 client
37 .delete_basin(DeleteBasinInput::new("my-events".parse()?))
38 .await?;
39 // ANCHOR_END: basin-operations
40 println!("Basins: {:?}, config: {:?}", basins, config);
41
42 let basin = client.basin(basin_name);
43
44 // ANCHOR: stream-operations
45 // List streams
46 let streams = basin
47 .list_streams(ListStreamsInput::new().with_prefix("user-".parse()?))
48 .await?;
49
50 // Create a stream
51 // Optionally, pass `.with_config(StreamConfig { .. })` to CreateStreamInput.
52 basin
53 .create_stream(CreateStreamInput::new("user-actions".parse()?))
54 .await?;
55
56 // Get configuration
57 let config = basin.get_stream_config("user-actions".parse()?).await?;
58
59 // Delete
60 basin
61 .delete_stream(DeleteStreamInput::new("user-actions".parse()?))
62 .await?;
63 // ANCHOR_END: stream-operations
64 println!("Streams: {:?}, config: {:?}", streams, config);
65
66 // ANCHOR: access-token-basic
67 // List tokens (returns metadata, not the secret)
68 let tokens = client.list_access_tokens(Default::default()).await?;
69
70 // Issue a token scoped to streams under "users/1234/"
71 let result = client
72 .issue_access_token(
73 IssueAccessTokenInput::new(
74 "user-1234-rw-token".parse()?,
75 AccessTokenScopeInput::from_op_group_perms(
76 OperationGroupPermissions::new()
77 .with_stream(ReadWritePermissions::read_write()),
78 )
79 .with_basins(BasinMatcher::Prefix("".parse()?)) // all basins
80 .with_streams(StreamMatcher::Prefix("users/1234/".parse()?)),
81 )
82 .with_expires_at("2027-01-01T00:00:00Z".parse()?),
83 )
84 .await?;
85
86 // Revoke a token
87 client
88 .revoke_access_token("user-1234-rw-token".parse()?)
89 .await?;
90 // ANCHOR_END: access-token-basic
91 println!("Tokens: {:?}, issued: {:?}", tokens, result);
92
93 // ANCHOR: access-token-restricted
94 client
95 .issue_access_token(IssueAccessTokenInput::new(
96 "restricted-token".parse()?,
97 AccessTokenScopeInput::from_op_group_perms(
98 OperationGroupPermissions::new().with_stream(ReadWritePermissions::read_only()),
99 )
100 .with_basins(BasinMatcher::Exact("production".parse()?))
101 .with_streams(StreamMatcher::Prefix("logs/".parse()?)),
102 ))
103 .await?;
104 // ANCHOR_END: access-token-restricted
105
106 // Pagination examples - not executed by default
107 if false {
108 // ANCHOR: pagination
109 // Iterate through all streams with automatic pagination
110 let mut stream = basin.list_all_streams(ListAllStreamsInput::new());
111 while let Some(info) = stream.next().await {
112 let info = info?;
113 println!("{}", info.name);
114 }
115 // ANCHOR_END: pagination
116
117 // ANCHOR: pagination-filtering
118 // List streams with a prefix filter
119 let input = ListAllStreamsInput::new().with_prefix("events/".parse()?);
120 let mut stream = basin.list_all_streams(input);
121 while let Some(info) = stream.next().await {
122 println!("{}", info?.name);
123 }
124 // ANCHOR_END: pagination-filtering
125
126 // ANCHOR: pagination-deleted
127 // Include streams that are being deleted
128 let input = ListAllStreamsInput::new().with_include_deleted(true);
129 let mut stream = basin.list_all_streams(input);
130 while let Some(info) = stream.next().await {
131 let info = info?;
132 println!("{} {:?}", info.name, info.deleted_at);
133 }
134 // ANCHOR_END: pagination-deleted
135 }
136
137 Ok(())
138}Sourcepub fn list_all_streams(
&self,
input: ListAllStreamsInput,
) -> Streaming<StreamInfo>
pub fn list_all_streams( &self, input: ListAllStreamsInput, ) -> Streaming<StreamInfo>
List all streams, paginating automatically.
Examples found in repository?
examples/docs_account_and_basins.rs (line 110)
17async fn main() -> Result<(), Box<dyn std::error::Error>> {
18 let access_token = std::env::var("S2_ACCESS_TOKEN")?;
19 let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
20
21 let client = S2::new(S2Config::new(access_token))?;
22
23 // ANCHOR: basin-operations
24 // List basins
25 let basins = client.list_basins(ListBasinsInput::new()).await?;
26
27 // Create a basin
28 client
29 .create_basin(CreateBasinInput::new("my-events".parse()?))
30 .await?;
31
32 // Get configuration
33 let config = client.get_basin_config("my-events".parse()?).await?;
34
35 // Delete
36 client
37 .delete_basin(DeleteBasinInput::new("my-events".parse()?))
38 .await?;
39 // ANCHOR_END: basin-operations
40 println!("Basins: {:?}, config: {:?}", basins, config);
41
42 let basin = client.basin(basin_name);
43
44 // ANCHOR: stream-operations
45 // List streams
46 let streams = basin
47 .list_streams(ListStreamsInput::new().with_prefix("user-".parse()?))
48 .await?;
49
50 // Create a stream
51 // Optionally, pass `.with_config(StreamConfig { .. })` to CreateStreamInput.
52 basin
53 .create_stream(CreateStreamInput::new("user-actions".parse()?))
54 .await?;
55
56 // Get configuration
57 let config = basin.get_stream_config("user-actions".parse()?).await?;
58
59 // Delete
60 basin
61 .delete_stream(DeleteStreamInput::new("user-actions".parse()?))
62 .await?;
63 // ANCHOR_END: stream-operations
64 println!("Streams: {:?}, config: {:?}", streams, config);
65
66 // ANCHOR: access-token-basic
67 // List tokens (returns metadata, not the secret)
68 let tokens = client.list_access_tokens(Default::default()).await?;
69
70 // Issue a token scoped to streams under "users/1234/"
71 let result = client
72 .issue_access_token(
73 IssueAccessTokenInput::new(
74 "user-1234-rw-token".parse()?,
75 AccessTokenScopeInput::from_op_group_perms(
76 OperationGroupPermissions::new()
77 .with_stream(ReadWritePermissions::read_write()),
78 )
79 .with_basins(BasinMatcher::Prefix("".parse()?)) // all basins
80 .with_streams(StreamMatcher::Prefix("users/1234/".parse()?)),
81 )
82 .with_expires_at("2027-01-01T00:00:00Z".parse()?),
83 )
84 .await?;
85
86 // Revoke a token
87 client
88 .revoke_access_token("user-1234-rw-token".parse()?)
89 .await?;
90 // ANCHOR_END: access-token-basic
91 println!("Tokens: {:?}, issued: {:?}", tokens, result);
92
93 // ANCHOR: access-token-restricted
94 client
95 .issue_access_token(IssueAccessTokenInput::new(
96 "restricted-token".parse()?,
97 AccessTokenScopeInput::from_op_group_perms(
98 OperationGroupPermissions::new().with_stream(ReadWritePermissions::read_only()),
99 )
100 .with_basins(BasinMatcher::Exact("production".parse()?))
101 .with_streams(StreamMatcher::Prefix("logs/".parse()?)),
102 ))
103 .await?;
104 // ANCHOR_END: access-token-restricted
105
106 // Pagination examples - not executed by default
107 if false {
108 // ANCHOR: pagination
109 // Iterate through all streams with automatic pagination
110 let mut stream = basin.list_all_streams(ListAllStreamsInput::new());
111 while let Some(info) = stream.next().await {
112 let info = info?;
113 println!("{}", info.name);
114 }
115 // ANCHOR_END: pagination
116
117 // ANCHOR: pagination-filtering
118 // List streams with a prefix filter
119 let input = ListAllStreamsInput::new().with_prefix("events/".parse()?);
120 let mut stream = basin.list_all_streams(input);
121 while let Some(info) = stream.next().await {
122 println!("{}", info?.name);
123 }
124 // ANCHOR_END: pagination-filtering
125
126 // ANCHOR: pagination-deleted
127 // Include streams that are being deleted
128 let input = ListAllStreamsInput::new().with_include_deleted(true);
129 let mut stream = basin.list_all_streams(input);
130 while let Some(info) = stream.next().await {
131 let info = info?;
132 println!("{} {:?}", info.name, info.deleted_at);
133 }
134 // ANCHOR_END: pagination-deleted
135 }
136
137 Ok(())
138}Sourcepub async fn create_stream(
&self,
input: CreateStreamInput,
) -> Result<StreamInfo, RequestError>
pub async fn create_stream( &self, input: CreateStreamInput, ) -> Result<StreamInfo, RequestError>
Create a stream.
Examples found in repository?
examples/create_stream.rs (line 28)
10async fn main() -> Result<(), Box<dyn std::error::Error>> {
11 let access_token =
12 std::env::var("S2_ACCESS_TOKEN").map_err(|_| "S2_ACCESS_TOKEN env var not set")?;
13 let basin_name: BasinName = std::env::var("S2_BASIN")
14 .map_err(|_| "S2_BASIN env var not set")?
15 .parse()?;
16 let stream_name: StreamName = std::env::var("S2_STREAM")
17 .map_err(|_| "S2_STREAM env var not set")?
18 .parse()?;
19
20 let s2 = S2::new(S2Config::new(access_token))?;
21 let basin = s2.basin(basin_name);
22
23 let input = CreateStreamInput::new(stream_name.clone()).with_config(
24 StreamConfig::new().with_timestamping(
25 TimestampingConfig::new().with_mode(TimestampingMode::ClientRequire),
26 ),
27 );
28 let stream_info = basin.create_stream(input).await?;
29 println!("{stream_info:#?}");
30
31 let stream_config = basin.get_stream_config(stream_name).await?;
32 println!("{stream_config:#?}");
33
34 Ok(())
35}More examples
examples/docs_encryption.rs (line 46)
15async fn main() -> Result<(), Box<dyn std::error::Error>> {
16 let access_token = std::env::var("S2_ACCESS_TOKEN")?;
17 let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
18 let stream_name: StreamName = format!(
19 "docs-encryption-{}",
20 std::time::SystemTime::now()
21 .duration_since(std::time::UNIX_EPOCH)?
22 .as_millis()
23 )
24 .parse()?;
25
26 let client = S2::new(S2Config::new(access_token))?;
27
28 // ANCHOR: basin-cipher
29 client
30 .create_basin(
31 CreateBasinInput::new(basin_name.clone())
32 .with_config(BasinConfig::new().with_stream_cipher(EncryptionAlgorithm::Aegis256)),
33 )
34 .await?;
35
36 client
37 .reconfigure_basin(ReconfigureBasinInput::new(
38 basin_name.clone(),
39 BasinReconfiguration::new().with_stream_cipher(EncryptionAlgorithm::Aes256Gcm),
40 ))
41 .await?;
42 // ANCHOR_END: basin-cipher
43
44 let basin = client.basin(basin_name.clone());
45 basin
46 .create_stream(CreateStreamInput::new(stream_name.clone()))
47 .await?;
48
49 // ANCHOR: append-read
50 let stream = basin
51 .stream(stream_name.clone())
52 .with_encryption_key(std::env::var("S2_ENCRYPTION_KEY")?.parse()?);
53
54 stream
55 .append(AppendInput::new(AppendRecordBatch::try_from_iter([
56 AppendRecord::new("top secret")?,
57 ])?))
58 .await?;
59
60 let batch = stream
61 .read(
62 ReadInput::new()
63 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
64 .with_stop(ReadStop::new().with_limits(ReadLimits::new().with_count(10))),
65 )
66 .await?;
67 // ANCHOR_END: append-read
68
69 println!("Read {} encrypted record(s)", batch.records.len());
70
71 basin
72 .delete_stream(DeleteStreamInput::new(stream_name))
73 .await?;
74
75 Ok(())
76}examples/caught_up.rs (line 40)
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}examples/docs_account_and_basins.rs (line 53)
17async fn main() -> Result<(), Box<dyn std::error::Error>> {
18 let access_token = std::env::var("S2_ACCESS_TOKEN")?;
19 let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
20
21 let client = S2::new(S2Config::new(access_token))?;
22
23 // ANCHOR: basin-operations
24 // List basins
25 let basins = client.list_basins(ListBasinsInput::new()).await?;
26
27 // Create a basin
28 client
29 .create_basin(CreateBasinInput::new("my-events".parse()?))
30 .await?;
31
32 // Get configuration
33 let config = client.get_basin_config("my-events".parse()?).await?;
34
35 // Delete
36 client
37 .delete_basin(DeleteBasinInput::new("my-events".parse()?))
38 .await?;
39 // ANCHOR_END: basin-operations
40 println!("Basins: {:?}, config: {:?}", basins, config);
41
42 let basin = client.basin(basin_name);
43
44 // ANCHOR: stream-operations
45 // List streams
46 let streams = basin
47 .list_streams(ListStreamsInput::new().with_prefix("user-".parse()?))
48 .await?;
49
50 // Create a stream
51 // Optionally, pass `.with_config(StreamConfig { .. })` to CreateStreamInput.
52 basin
53 .create_stream(CreateStreamInput::new("user-actions".parse()?))
54 .await?;
55
56 // Get configuration
57 let config = basin.get_stream_config("user-actions".parse()?).await?;
58
59 // Delete
60 basin
61 .delete_stream(DeleteStreamInput::new("user-actions".parse()?))
62 .await?;
63 // ANCHOR_END: stream-operations
64 println!("Streams: {:?}, config: {:?}", streams, config);
65
66 // ANCHOR: access-token-basic
67 // List tokens (returns metadata, not the secret)
68 let tokens = client.list_access_tokens(Default::default()).await?;
69
70 // Issue a token scoped to streams under "users/1234/"
71 let result = client
72 .issue_access_token(
73 IssueAccessTokenInput::new(
74 "user-1234-rw-token".parse()?,
75 AccessTokenScopeInput::from_op_group_perms(
76 OperationGroupPermissions::new()
77 .with_stream(ReadWritePermissions::read_write()),
78 )
79 .with_basins(BasinMatcher::Prefix("".parse()?)) // all basins
80 .with_streams(StreamMatcher::Prefix("users/1234/".parse()?)),
81 )
82 .with_expires_at("2027-01-01T00:00:00Z".parse()?),
83 )
84 .await?;
85
86 // Revoke a token
87 client
88 .revoke_access_token("user-1234-rw-token".parse()?)
89 .await?;
90 // ANCHOR_END: access-token-basic
91 println!("Tokens: {:?}, issued: {:?}", tokens, result);
92
93 // ANCHOR: access-token-restricted
94 client
95 .issue_access_token(IssueAccessTokenInput::new(
96 "restricted-token".parse()?,
97 AccessTokenScopeInput::from_op_group_perms(
98 OperationGroupPermissions::new().with_stream(ReadWritePermissions::read_only()),
99 )
100 .with_basins(BasinMatcher::Exact("production".parse()?))
101 .with_streams(StreamMatcher::Prefix("logs/".parse()?)),
102 ))
103 .await?;
104 // ANCHOR_END: access-token-restricted
105
106 // Pagination examples - not executed by default
107 if false {
108 // ANCHOR: pagination
109 // Iterate through all streams with automatic pagination
110 let mut stream = basin.list_all_streams(ListAllStreamsInput::new());
111 while let Some(info) = stream.next().await {
112 let info = info?;
113 println!("{}", info.name);
114 }
115 // ANCHOR_END: pagination
116
117 // ANCHOR: pagination-filtering
118 // List streams with a prefix filter
119 let input = ListAllStreamsInput::new().with_prefix("events/".parse()?);
120 let mut stream = basin.list_all_streams(input);
121 while let Some(info) = stream.next().await {
122 println!("{}", info?.name);
123 }
124 // ANCHOR_END: pagination-filtering
125
126 // ANCHOR: pagination-deleted
127 // Include streams that are being deleted
128 let input = ListAllStreamsInput::new().with_include_deleted(true);
129 let mut stream = basin.list_all_streams(input);
130 while let Some(info) = stream.next().await {
131 let info = info?;
132 println!("{} {:?}", info.name, info.deleted_at);
133 }
134 // ANCHOR_END: pagination-deleted
135 }
136
137 Ok(())
138}examples/docs_streams.rs (line 36)
20async fn main() -> Result<(), Box<dyn std::error::Error>> {
21 let access_token = std::env::var("S2_ACCESS_TOKEN")?;
22 let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
23
24 let client = S2::new(S2Config::new(access_token))?;
25 let basin = client.basin(basin_name);
26
27 // Create a temporary stream for examples
28 let stream_name: StreamName = format!(
29 "docs-streams-{}",
30 std::time::SystemTime::now()
31 .duration_since(std::time::UNIX_EPOCH)?
32 .as_millis()
33 )
34 .parse()?;
35 basin
36 .create_stream(s2_sdk::types::CreateStreamInput::new(stream_name.clone()))
37 .await?;
38
39 // ANCHOR: simple-append
40 let stream = basin.stream(stream_name.clone());
41
42 let ack = stream
43 .append(AppendInput::new(AppendRecordBatch::try_from_iter([
44 AppendRecord::new("first event")?,
45 AppendRecord::new("second event")?,
46 ])?))
47 .await?;
48
49 // ack tells us where the records landed
50 println!(
51 "Wrote records {} through {}",
52 ack.start.seq_num,
53 ack.end.seq_num - 1
54 );
55 // ANCHOR_END: simple-append
56
57 // ANCHOR: simple-read
58 let batch = stream
59 .read(
60 ReadInput::new()
61 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
62 .with_stop(ReadStop::new().with_limits(ReadLimits::new().with_count(100))),
63 )
64 .await?;
65
66 for record in batch.records {
67 println!("[{}] {:?}", record.seq_num, record.body);
68 }
69 // ANCHOR_END: simple-read
70
71 // ANCHOR: append-session
72 let session = stream.append_session(AppendSessionConfig::new());
73
74 // Submit a batch - this enqueues it and returns a ticket
75 let records = AppendRecordBatch::try_from_iter([
76 AppendRecord::new("event-1")?,
77 AppendRecord::new("event-2")?,
78 ])?;
79 let ticket = session.submit(AppendInput::new(records)).await?;
80
81 // Wait for durability
82 let ack = ticket.await?;
83 println!("Durable at seqNum {}", ack.start.seq_num);
84
85 session.close().await?;
86 // ANCHOR_END: append-session
87
88 // ANCHOR: producer
89 let producer = stream.producer(
90 ProducerConfig::new()
91 .with_batching(BatchingConfig::new().with_linger(Duration::from_millis(5))),
92 );
93
94 // Submit individual records
95 let ticket = producer.submit(AppendRecord::new("my event")?).await?;
96
97 // Get the exact sequence number
98 let ack = ticket.await?;
99 println!("Record durable at seqNum {}", ack.seq_num);
100
101 producer.close().await?;
102 // ANCHOR_END: producer
103
104 // ANCHOR: check-tail
105 let tail = stream.check_tail().await?;
106 println!("Stream has {} records", tail.seq_num);
107 // ANCHOR_END: check-tail
108
109 // Cleanup
110 basin
111 .delete_stream(s2_sdk::types::DeleteStreamInput::new(stream_name))
112 .await?;
113
114 println!("Streams examples completed");
115
116 // The following read session examples are for documentation snippets only.
117 // They are not executed because they would block waiting for new records.
118 if std::env::var("RUN_READ_SESSIONS").is_err() {
119 return Ok(());
120 }
121
122 // ANCHOR: read-session
123 let mut session = stream
124 .read_session(
125 ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0))),
126 ReadSessionConfig::default(),
127 )
128 .await?;
129
130 while let Some(batch) = session.next().await {
131 let batch = batch?;
132 for record in batch.records {
133 println!("[{}] {:?}", record.seq_num, record.body);
134 }
135 }
136 // ANCHOR_END: read-session
137
138 // ANCHOR: read-session-tail-offset
139 // Start reading from 10 records before the current tail
140 let mut session = stream
141 .read_session(
142 ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::TailOffset(10))),
143 ReadSessionConfig::default(),
144 )
145 .await?;
146
147 while let Some(batch) = session.next().await {
148 let batch = batch?;
149 for record in batch.records {
150 println!("[{}] {:?}", record.seq_num, record.body);
151 }
152 }
153 // ANCHOR_END: read-session-tail-offset
154
155 // ANCHOR: read-session-timestamp
156 // Start reading from a specific timestamp
157 let one_hour_ago = std::time::SystemTime::now()
158 .duration_since(std::time::UNIX_EPOCH)?
159 .as_millis() as u64
160 - 3600 * 1000;
161 let mut session = stream
162 .read_session(
163 ReadInput::new()
164 .with_start(ReadStart::new().with_from(ReadFrom::Timestamp(one_hour_ago))),
165 ReadSessionConfig::default(),
166 )
167 .await?;
168
169 while let Some(batch) = session.next().await {
170 let batch = batch?;
171 for record in batch.records {
172 println!("[{}] {:?}", record.seq_num, record.body);
173 }
174 }
175 // ANCHOR_END: read-session-timestamp
176
177 // ANCHOR: read-session-until
178 // Read records until a specific timestamp
179 let one_hour_ago = std::time::SystemTime::now()
180 .duration_since(std::time::UNIX_EPOCH)?
181 .as_millis() as u64
182 - 3600 * 1000;
183 let mut session = stream
184 .read_session(
185 ReadInput::new()
186 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
187 .with_stop(ReadStop::new().with_until(..one_hour_ago)),
188 ReadSessionConfig::default(),
189 )
190 .await?;
191
192 while let Some(batch) = session.next().await {
193 let batch = batch?;
194 for record in batch.records {
195 println!("[{}] {:?}", record.seq_num, record.body);
196 }
197 }
198 // ANCHOR_END: read-session-until
199
200 // ANCHOR: read-session-wait
201 // Read all available records, and once reaching the current tail, wait an additional 30 seconds
202 // for new ones
203 let mut session = stream
204 .read_session(
205 ReadInput::new()
206 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
207 .with_stop(ReadStop::new().with_wait(30)),
208 ReadSessionConfig::default(),
209 )
210 .await?;
211
212 while let Some(batch) = session.next().await {
213 let batch = batch?;
214 for record in batch.records {
215 println!("[{}] {:?}", record.seq_num, record.body);
216 }
217 }
218 // ANCHOR_END: read-session-wait
219
220 Ok(())
221}Sourcepub async fn ensure_stream(
&self,
input: EnsureStreamInput,
) -> Result<EnsureOutput<StreamInfo>, RequestError>
pub async fn ensure_stream( &self, input: EnsureStreamInput, ) -> Result<EnsureOutput<StreamInfo>, RequestError>
Ensure a stream.
If the stream doesn’t exist, creates the stream with specified configuration.
If the stream already exists:
- Its configuration is updated to the specified configuration, if different.
- Its configuration is unchanged, if the specified configuration is same.
Sourcepub async fn get_stream_config(
&self,
name: StreamName,
) -> Result<StreamConfig, RequestError>
pub async fn get_stream_config( &self, name: StreamName, ) -> Result<StreamConfig, RequestError>
Get stream configuration.
Examples found in repository?
examples/create_stream.rs (line 31)
10async fn main() -> Result<(), Box<dyn std::error::Error>> {
11 let access_token =
12 std::env::var("S2_ACCESS_TOKEN").map_err(|_| "S2_ACCESS_TOKEN env var not set")?;
13 let basin_name: BasinName = std::env::var("S2_BASIN")
14 .map_err(|_| "S2_BASIN env var not set")?
15 .parse()?;
16 let stream_name: StreamName = std::env::var("S2_STREAM")
17 .map_err(|_| "S2_STREAM env var not set")?
18 .parse()?;
19
20 let s2 = S2::new(S2Config::new(access_token))?;
21 let basin = s2.basin(basin_name);
22
23 let input = CreateStreamInput::new(stream_name.clone()).with_config(
24 StreamConfig::new().with_timestamping(
25 TimestampingConfig::new().with_mode(TimestampingMode::ClientRequire),
26 ),
27 );
28 let stream_info = basin.create_stream(input).await?;
29 println!("{stream_info:#?}");
30
31 let stream_config = basin.get_stream_config(stream_name).await?;
32 println!("{stream_config:#?}");
33
34 Ok(())
35}More examples
examples/docs_account_and_basins.rs (line 57)
17async fn main() -> Result<(), Box<dyn std::error::Error>> {
18 let access_token = std::env::var("S2_ACCESS_TOKEN")?;
19 let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
20
21 let client = S2::new(S2Config::new(access_token))?;
22
23 // ANCHOR: basin-operations
24 // List basins
25 let basins = client.list_basins(ListBasinsInput::new()).await?;
26
27 // Create a basin
28 client
29 .create_basin(CreateBasinInput::new("my-events".parse()?))
30 .await?;
31
32 // Get configuration
33 let config = client.get_basin_config("my-events".parse()?).await?;
34
35 // Delete
36 client
37 .delete_basin(DeleteBasinInput::new("my-events".parse()?))
38 .await?;
39 // ANCHOR_END: basin-operations
40 println!("Basins: {:?}, config: {:?}", basins, config);
41
42 let basin = client.basin(basin_name);
43
44 // ANCHOR: stream-operations
45 // List streams
46 let streams = basin
47 .list_streams(ListStreamsInput::new().with_prefix("user-".parse()?))
48 .await?;
49
50 // Create a stream
51 // Optionally, pass `.with_config(StreamConfig { .. })` to CreateStreamInput.
52 basin
53 .create_stream(CreateStreamInput::new("user-actions".parse()?))
54 .await?;
55
56 // Get configuration
57 let config = basin.get_stream_config("user-actions".parse()?).await?;
58
59 // Delete
60 basin
61 .delete_stream(DeleteStreamInput::new("user-actions".parse()?))
62 .await?;
63 // ANCHOR_END: stream-operations
64 println!("Streams: {:?}, config: {:?}", streams, config);
65
66 // ANCHOR: access-token-basic
67 // List tokens (returns metadata, not the secret)
68 let tokens = client.list_access_tokens(Default::default()).await?;
69
70 // Issue a token scoped to streams under "users/1234/"
71 let result = client
72 .issue_access_token(
73 IssueAccessTokenInput::new(
74 "user-1234-rw-token".parse()?,
75 AccessTokenScopeInput::from_op_group_perms(
76 OperationGroupPermissions::new()
77 .with_stream(ReadWritePermissions::read_write()),
78 )
79 .with_basins(BasinMatcher::Prefix("".parse()?)) // all basins
80 .with_streams(StreamMatcher::Prefix("users/1234/".parse()?)),
81 )
82 .with_expires_at("2027-01-01T00:00:00Z".parse()?),
83 )
84 .await?;
85
86 // Revoke a token
87 client
88 .revoke_access_token("user-1234-rw-token".parse()?)
89 .await?;
90 // ANCHOR_END: access-token-basic
91 println!("Tokens: {:?}, issued: {:?}", tokens, result);
92
93 // ANCHOR: access-token-restricted
94 client
95 .issue_access_token(IssueAccessTokenInput::new(
96 "restricted-token".parse()?,
97 AccessTokenScopeInput::from_op_group_perms(
98 OperationGroupPermissions::new().with_stream(ReadWritePermissions::read_only()),
99 )
100 .with_basins(BasinMatcher::Exact("production".parse()?))
101 .with_streams(StreamMatcher::Prefix("logs/".parse()?)),
102 ))
103 .await?;
104 // ANCHOR_END: access-token-restricted
105
106 // Pagination examples - not executed by default
107 if false {
108 // ANCHOR: pagination
109 // Iterate through all streams with automatic pagination
110 let mut stream = basin.list_all_streams(ListAllStreamsInput::new());
111 while let Some(info) = stream.next().await {
112 let info = info?;
113 println!("{}", info.name);
114 }
115 // ANCHOR_END: pagination
116
117 // ANCHOR: pagination-filtering
118 // List streams with a prefix filter
119 let input = ListAllStreamsInput::new().with_prefix("events/".parse()?);
120 let mut stream = basin.list_all_streams(input);
121 while let Some(info) = stream.next().await {
122 println!("{}", info?.name);
123 }
124 // ANCHOR_END: pagination-filtering
125
126 // ANCHOR: pagination-deleted
127 // Include streams that are being deleted
128 let input = ListAllStreamsInput::new().with_include_deleted(true);
129 let mut stream = basin.list_all_streams(input);
130 while let Some(info) = stream.next().await {
131 let info = info?;
132 println!("{} {:?}", info.name, info.deleted_at);
133 }
134 // ANCHOR_END: pagination-deleted
135 }
136
137 Ok(())
138}Sourcepub async fn delete_stream(
&self,
input: DeleteStreamInput,
) -> Result<(), RequestError>
pub async fn delete_stream( &self, input: DeleteStreamInput, ) -> Result<(), RequestError>
Delete a stream.
Examples found in repository?
examples/delete_stream.rs (line 21)
7async fn main() -> Result<(), Box<dyn std::error::Error>> {
8 let access_token =
9 std::env::var("S2_ACCESS_TOKEN").map_err(|_| "S2_ACCESS_TOKEN env var not set")?;
10 let basin_name: BasinName = std::env::var("S2_BASIN")
11 .map_err(|_| "S2_BASIN env var not set")?
12 .parse()?;
13 let stream_name: StreamName = std::env::var("S2_STREAM")
14 .map_err(|_| "S2_STREAM env var not set")?
15 .parse()?;
16
17 let s2 = S2::new(S2Config::new(access_token))?;
18 let basin = s2.basin(basin_name);
19
20 let input = DeleteStreamInput::new(stream_name);
21 basin.delete_stream(input).await?;
22 println!("Deletion requested");
23
24 Ok(())
25}More examples
examples/docs_encryption.rs (line 72)
15async fn main() -> Result<(), Box<dyn std::error::Error>> {
16 let access_token = std::env::var("S2_ACCESS_TOKEN")?;
17 let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
18 let stream_name: StreamName = format!(
19 "docs-encryption-{}",
20 std::time::SystemTime::now()
21 .duration_since(std::time::UNIX_EPOCH)?
22 .as_millis()
23 )
24 .parse()?;
25
26 let client = S2::new(S2Config::new(access_token))?;
27
28 // ANCHOR: basin-cipher
29 client
30 .create_basin(
31 CreateBasinInput::new(basin_name.clone())
32 .with_config(BasinConfig::new().with_stream_cipher(EncryptionAlgorithm::Aegis256)),
33 )
34 .await?;
35
36 client
37 .reconfigure_basin(ReconfigureBasinInput::new(
38 basin_name.clone(),
39 BasinReconfiguration::new().with_stream_cipher(EncryptionAlgorithm::Aes256Gcm),
40 ))
41 .await?;
42 // ANCHOR_END: basin-cipher
43
44 let basin = client.basin(basin_name.clone());
45 basin
46 .create_stream(CreateStreamInput::new(stream_name.clone()))
47 .await?;
48
49 // ANCHOR: append-read
50 let stream = basin
51 .stream(stream_name.clone())
52 .with_encryption_key(std::env::var("S2_ENCRYPTION_KEY")?.parse()?);
53
54 stream
55 .append(AppendInput::new(AppendRecordBatch::try_from_iter([
56 AppendRecord::new("top secret")?,
57 ])?))
58 .await?;
59
60 let batch = stream
61 .read(
62 ReadInput::new()
63 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
64 .with_stop(ReadStop::new().with_limits(ReadLimits::new().with_count(10))),
65 )
66 .await?;
67 // ANCHOR_END: append-read
68
69 println!("Read {} encrypted record(s)", batch.records.len());
70
71 basin
72 .delete_stream(DeleteStreamInput::new(stream_name))
73 .await?;
74
75 Ok(())
76}examples/caught_up.rs (line 96)
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}examples/docs_account_and_basins.rs (line 61)
17async fn main() -> Result<(), Box<dyn std::error::Error>> {
18 let access_token = std::env::var("S2_ACCESS_TOKEN")?;
19 let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
20
21 let client = S2::new(S2Config::new(access_token))?;
22
23 // ANCHOR: basin-operations
24 // List basins
25 let basins = client.list_basins(ListBasinsInput::new()).await?;
26
27 // Create a basin
28 client
29 .create_basin(CreateBasinInput::new("my-events".parse()?))
30 .await?;
31
32 // Get configuration
33 let config = client.get_basin_config("my-events".parse()?).await?;
34
35 // Delete
36 client
37 .delete_basin(DeleteBasinInput::new("my-events".parse()?))
38 .await?;
39 // ANCHOR_END: basin-operations
40 println!("Basins: {:?}, config: {:?}", basins, config);
41
42 let basin = client.basin(basin_name);
43
44 // ANCHOR: stream-operations
45 // List streams
46 let streams = basin
47 .list_streams(ListStreamsInput::new().with_prefix("user-".parse()?))
48 .await?;
49
50 // Create a stream
51 // Optionally, pass `.with_config(StreamConfig { .. })` to CreateStreamInput.
52 basin
53 .create_stream(CreateStreamInput::new("user-actions".parse()?))
54 .await?;
55
56 // Get configuration
57 let config = basin.get_stream_config("user-actions".parse()?).await?;
58
59 // Delete
60 basin
61 .delete_stream(DeleteStreamInput::new("user-actions".parse()?))
62 .await?;
63 // ANCHOR_END: stream-operations
64 println!("Streams: {:?}, config: {:?}", streams, config);
65
66 // ANCHOR: access-token-basic
67 // List tokens (returns metadata, not the secret)
68 let tokens = client.list_access_tokens(Default::default()).await?;
69
70 // Issue a token scoped to streams under "users/1234/"
71 let result = client
72 .issue_access_token(
73 IssueAccessTokenInput::new(
74 "user-1234-rw-token".parse()?,
75 AccessTokenScopeInput::from_op_group_perms(
76 OperationGroupPermissions::new()
77 .with_stream(ReadWritePermissions::read_write()),
78 )
79 .with_basins(BasinMatcher::Prefix("".parse()?)) // all basins
80 .with_streams(StreamMatcher::Prefix("users/1234/".parse()?)),
81 )
82 .with_expires_at("2027-01-01T00:00:00Z".parse()?),
83 )
84 .await?;
85
86 // Revoke a token
87 client
88 .revoke_access_token("user-1234-rw-token".parse()?)
89 .await?;
90 // ANCHOR_END: access-token-basic
91 println!("Tokens: {:?}, issued: {:?}", tokens, result);
92
93 // ANCHOR: access-token-restricted
94 client
95 .issue_access_token(IssueAccessTokenInput::new(
96 "restricted-token".parse()?,
97 AccessTokenScopeInput::from_op_group_perms(
98 OperationGroupPermissions::new().with_stream(ReadWritePermissions::read_only()),
99 )
100 .with_basins(BasinMatcher::Exact("production".parse()?))
101 .with_streams(StreamMatcher::Prefix("logs/".parse()?)),
102 ))
103 .await?;
104 // ANCHOR_END: access-token-restricted
105
106 // Pagination examples - not executed by default
107 if false {
108 // ANCHOR: pagination
109 // Iterate through all streams with automatic pagination
110 let mut stream = basin.list_all_streams(ListAllStreamsInput::new());
111 while let Some(info) = stream.next().await {
112 let info = info?;
113 println!("{}", info.name);
114 }
115 // ANCHOR_END: pagination
116
117 // ANCHOR: pagination-filtering
118 // List streams with a prefix filter
119 let input = ListAllStreamsInput::new().with_prefix("events/".parse()?);
120 let mut stream = basin.list_all_streams(input);
121 while let Some(info) = stream.next().await {
122 println!("{}", info?.name);
123 }
124 // ANCHOR_END: pagination-filtering
125
126 // ANCHOR: pagination-deleted
127 // Include streams that are being deleted
128 let input = ListAllStreamsInput::new().with_include_deleted(true);
129 let mut stream = basin.list_all_streams(input);
130 while let Some(info) = stream.next().await {
131 let info = info?;
132 println!("{} {:?}", info.name, info.deleted_at);
133 }
134 // ANCHOR_END: pagination-deleted
135 }
136
137 Ok(())
138}examples/docs_streams.rs (line 111)
20async fn main() -> Result<(), Box<dyn std::error::Error>> {
21 let access_token = std::env::var("S2_ACCESS_TOKEN")?;
22 let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
23
24 let client = S2::new(S2Config::new(access_token))?;
25 let basin = client.basin(basin_name);
26
27 // Create a temporary stream for examples
28 let stream_name: StreamName = format!(
29 "docs-streams-{}",
30 std::time::SystemTime::now()
31 .duration_since(std::time::UNIX_EPOCH)?
32 .as_millis()
33 )
34 .parse()?;
35 basin
36 .create_stream(s2_sdk::types::CreateStreamInput::new(stream_name.clone()))
37 .await?;
38
39 // ANCHOR: simple-append
40 let stream = basin.stream(stream_name.clone());
41
42 let ack = stream
43 .append(AppendInput::new(AppendRecordBatch::try_from_iter([
44 AppendRecord::new("first event")?,
45 AppendRecord::new("second event")?,
46 ])?))
47 .await?;
48
49 // ack tells us where the records landed
50 println!(
51 "Wrote records {} through {}",
52 ack.start.seq_num,
53 ack.end.seq_num - 1
54 );
55 // ANCHOR_END: simple-append
56
57 // ANCHOR: simple-read
58 let batch = stream
59 .read(
60 ReadInput::new()
61 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
62 .with_stop(ReadStop::new().with_limits(ReadLimits::new().with_count(100))),
63 )
64 .await?;
65
66 for record in batch.records {
67 println!("[{}] {:?}", record.seq_num, record.body);
68 }
69 // ANCHOR_END: simple-read
70
71 // ANCHOR: append-session
72 let session = stream.append_session(AppendSessionConfig::new());
73
74 // Submit a batch - this enqueues it and returns a ticket
75 let records = AppendRecordBatch::try_from_iter([
76 AppendRecord::new("event-1")?,
77 AppendRecord::new("event-2")?,
78 ])?;
79 let ticket = session.submit(AppendInput::new(records)).await?;
80
81 // Wait for durability
82 let ack = ticket.await?;
83 println!("Durable at seqNum {}", ack.start.seq_num);
84
85 session.close().await?;
86 // ANCHOR_END: append-session
87
88 // ANCHOR: producer
89 let producer = stream.producer(
90 ProducerConfig::new()
91 .with_batching(BatchingConfig::new().with_linger(Duration::from_millis(5))),
92 );
93
94 // Submit individual records
95 let ticket = producer.submit(AppendRecord::new("my event")?).await?;
96
97 // Get the exact sequence number
98 let ack = ticket.await?;
99 println!("Record durable at seqNum {}", ack.seq_num);
100
101 producer.close().await?;
102 // ANCHOR_END: producer
103
104 // ANCHOR: check-tail
105 let tail = stream.check_tail().await?;
106 println!("Stream has {} records", tail.seq_num);
107 // ANCHOR_END: check-tail
108
109 // Cleanup
110 basin
111 .delete_stream(s2_sdk::types::DeleteStreamInput::new(stream_name))
112 .await?;
113
114 println!("Streams examples completed");
115
116 // The following read session examples are for documentation snippets only.
117 // They are not executed because they would block waiting for new records.
118 if std::env::var("RUN_READ_SESSIONS").is_err() {
119 return Ok(());
120 }
121
122 // ANCHOR: read-session
123 let mut session = stream
124 .read_session(
125 ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0))),
126 ReadSessionConfig::default(),
127 )
128 .await?;
129
130 while let Some(batch) = session.next().await {
131 let batch = batch?;
132 for record in batch.records {
133 println!("[{}] {:?}", record.seq_num, record.body);
134 }
135 }
136 // ANCHOR_END: read-session
137
138 // ANCHOR: read-session-tail-offset
139 // Start reading from 10 records before the current tail
140 let mut session = stream
141 .read_session(
142 ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::TailOffset(10))),
143 ReadSessionConfig::default(),
144 )
145 .await?;
146
147 while let Some(batch) = session.next().await {
148 let batch = batch?;
149 for record in batch.records {
150 println!("[{}] {:?}", record.seq_num, record.body);
151 }
152 }
153 // ANCHOR_END: read-session-tail-offset
154
155 // ANCHOR: read-session-timestamp
156 // Start reading from a specific timestamp
157 let one_hour_ago = std::time::SystemTime::now()
158 .duration_since(std::time::UNIX_EPOCH)?
159 .as_millis() as u64
160 - 3600 * 1000;
161 let mut session = stream
162 .read_session(
163 ReadInput::new()
164 .with_start(ReadStart::new().with_from(ReadFrom::Timestamp(one_hour_ago))),
165 ReadSessionConfig::default(),
166 )
167 .await?;
168
169 while let Some(batch) = session.next().await {
170 let batch = batch?;
171 for record in batch.records {
172 println!("[{}] {:?}", record.seq_num, record.body);
173 }
174 }
175 // ANCHOR_END: read-session-timestamp
176
177 // ANCHOR: read-session-until
178 // Read records until a specific timestamp
179 let one_hour_ago = std::time::SystemTime::now()
180 .duration_since(std::time::UNIX_EPOCH)?
181 .as_millis() as u64
182 - 3600 * 1000;
183 let mut session = stream
184 .read_session(
185 ReadInput::new()
186 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
187 .with_stop(ReadStop::new().with_until(..one_hour_ago)),
188 ReadSessionConfig::default(),
189 )
190 .await?;
191
192 while let Some(batch) = session.next().await {
193 let batch = batch?;
194 for record in batch.records {
195 println!("[{}] {:?}", record.seq_num, record.body);
196 }
197 }
198 // ANCHOR_END: read-session-until
199
200 // ANCHOR: read-session-wait
201 // Read all available records, and once reaching the current tail, wait an additional 30 seconds
202 // for new ones
203 let mut session = stream
204 .read_session(
205 ReadInput::new()
206 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
207 .with_stop(ReadStop::new().with_wait(30)),
208 ReadSessionConfig::default(),
209 )
210 .await?;
211
212 while let Some(batch) = session.next().await {
213 let batch = batch?;
214 for record in batch.records {
215 println!("[{}] {:?}", record.seq_num, record.body);
216 }
217 }
218 // ANCHOR_END: read-session-wait
219
220 Ok(())
221}Sourcepub async fn reconfigure_stream(
&self,
input: ReconfigureStreamInput,
) -> Result<StreamConfig, RequestError>
pub async fn reconfigure_stream( &self, input: ReconfigureStreamInput, ) -> Result<StreamConfig, RequestError>
Reconfigure a stream.
Examples found in repository?
examples/reconfigure_stream.rs (line 27)
10async fn main() -> Result<(), Box<dyn std::error::Error>> {
11 let access_token =
12 std::env::var("S2_ACCESS_TOKEN").map_err(|_| "S2_ACCESS_TOKEN env var not set")?;
13 let basin_name: BasinName = std::env::var("S2_BASIN")
14 .map_err(|_| "S2_BASIN env var not set")?
15 .parse()?;
16 let stream_name: StreamName = std::env::var("S2_STREAM")
17 .map_err(|_| "S2_STREAM env var not set")?
18 .parse()?;
19
20 let s2 = S2::new(S2Config::new(access_token))?;
21 let basin = s2.basin(basin_name);
22
23 let input = ReconfigureStreamInput::new(
24 stream_name,
25 StreamReconfiguration::new().with_retention_policy(RetentionPolicy::Age(10 * 24 * 60 * 60)),
26 );
27 let config = basin.reconfigure_stream(input).await?;
28 println!("{config:#?}");
29
30 Ok(())
31}Trait Implementations§
Auto Trait Implementations§
impl !Freeze for S2Basin
impl !RefUnwindSafe for S2Basin
impl !UnwindSafe for S2Basin
impl Send for S2Basin
impl Sync for S2Basin
impl Unpin for S2Basin
impl UnsafeUnpin for S2Basin
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