#[non_exhaustive]pub struct ReadInput {
pub start: ReadStart,
pub stop: ReadStop,
pub ignore_command_records: bool,
}Expand description
Input for read and read_session
operations.
Fields (Non-exhaustive)§
This struct is marked as non-exhaustive
Non-exhaustive structs could have additional fields added in future. Therefore, non-exhaustive structs cannot be constructed in external crates using the traditional
Struct { .. } syntax; cannot be matched against without a wildcard ..; and struct update syntax will not work.start: ReadStartWhere to start reading.
See ReadStart for defaults.
stop: ReadStopWhen to stop reading.
See ReadStop for defaults.
ignore_command_records: boolWhether to filter out command records from the stream when reading.
Defaults to false.
Implementations§
Source§impl ReadInput
impl ReadInput
Sourcepub fn new() -> Self
pub fn new() -> Self
Create a new ReadInput with default values.
Examples found in repository?
examples/consumer.rs (line 17)
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}More examples
examples/get_latest_record.rs (line 22)
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/docs_encryption.rs (line 62)
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 53)
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_streams.rs (line 60)
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 // Force the partial batch and wait for durability without closing the producer
98 producer.flush().await?;
99
100 // Get the exact sequence number
101 let ack = ticket.await?;
102 println!("Record durable at seqNum {}", ack.seq_num);
103
104 producer.close().await?;
105 // ANCHOR_END: producer
106
107 // ANCHOR: check-tail
108 let tail = stream.check_tail().await?;
109 println!("Stream has {} records", tail.seq_num);
110 // ANCHOR_END: check-tail
111
112 // Cleanup
113 basin
114 .delete_stream(s2_sdk::types::DeleteStreamInput::new(stream_name))
115 .await?;
116
117 println!("Streams examples completed");
118
119 // The following read session examples are for documentation snippets only.
120 // They are not executed because they would block waiting for new records.
121 if std::env::var("RUN_READ_SESSIONS").is_err() {
122 return Ok(());
123 }
124
125 // ANCHOR: read-session
126 let mut session = stream
127 .read_session(
128 ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0))),
129 ReadSessionConfig::default(),
130 )
131 .await?;
132
133 while let Some(batch) = session.next().await {
134 let batch = batch?;
135 for record in batch.records {
136 println!("[{}] {:?}", record.seq_num, record.body);
137 }
138 }
139 // ANCHOR_END: read-session
140
141 // ANCHOR: read-session-tail-offset
142 // Start reading from 10 records before the current tail
143 let mut session = stream
144 .read_session(
145 ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::TailOffset(10))),
146 ReadSessionConfig::default(),
147 )
148 .await?;
149
150 while let Some(batch) = session.next().await {
151 let batch = batch?;
152 for record in batch.records {
153 println!("[{}] {:?}", record.seq_num, record.body);
154 }
155 }
156 // ANCHOR_END: read-session-tail-offset
157
158 // ANCHOR: read-session-timestamp
159 // Start reading from a specific timestamp
160 let one_hour_ago = std::time::SystemTime::now()
161 .duration_since(std::time::UNIX_EPOCH)?
162 .as_millis() as u64
163 - 3600 * 1000;
164 let mut session = stream
165 .read_session(
166 ReadInput::new()
167 .with_start(ReadStart::new().with_from(ReadFrom::Timestamp(one_hour_ago))),
168 ReadSessionConfig::default(),
169 )
170 .await?;
171
172 while let Some(batch) = session.next().await {
173 let batch = batch?;
174 for record in batch.records {
175 println!("[{}] {:?}", record.seq_num, record.body);
176 }
177 }
178 // ANCHOR_END: read-session-timestamp
179
180 // ANCHOR: read-session-until
181 // Read records until a specific timestamp
182 let one_hour_ago = std::time::SystemTime::now()
183 .duration_since(std::time::UNIX_EPOCH)?
184 .as_millis() as u64
185 - 3600 * 1000;
186 let mut session = stream
187 .read_session(
188 ReadInput::new()
189 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
190 .with_stop(ReadStop::new().with_until(..one_hour_ago)),
191 ReadSessionConfig::default(),
192 )
193 .await?;
194
195 while let Some(batch) = session.next().await {
196 let batch = batch?;
197 for record in batch.records {
198 println!("[{}] {:?}", record.seq_num, record.body);
199 }
200 }
201 // ANCHOR_END: read-session-until
202
203 // ANCHOR: read-session-wait
204 // Read all available records, and once reaching the current tail, wait an additional 30 seconds
205 // for new ones
206 let mut session = stream
207 .read_session(
208 ReadInput::new()
209 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
210 .with_stop(ReadStop::new().with_wait(30)),
211 ReadSessionConfig::default(),
212 )
213 .await?;
214
215 while let Some(batch) = session.next().await {
216 let batch = batch?;
217 for record in batch.records {
218 println!("[{}] {:?}", record.seq_num, record.body);
219 }
220 }
221 // ANCHOR_END: read-session-wait
222
223 Ok(())
224}Sourcepub fn with_start(self, start: ReadStart) -> Self
pub fn with_start(self, start: ReadStart) -> Self
Set where to start reading.
Examples found in repository?
examples/get_latest_record.rs (line 23)
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}More examples
examples/docs_encryption.rs (line 63)
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 53)
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_streams.rs (line 61)
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 // Force the partial batch and wait for durability without closing the producer
98 producer.flush().await?;
99
100 // Get the exact sequence number
101 let ack = ticket.await?;
102 println!("Record durable at seqNum {}", ack.seq_num);
103
104 producer.close().await?;
105 // ANCHOR_END: producer
106
107 // ANCHOR: check-tail
108 let tail = stream.check_tail().await?;
109 println!("Stream has {} records", tail.seq_num);
110 // ANCHOR_END: check-tail
111
112 // Cleanup
113 basin
114 .delete_stream(s2_sdk::types::DeleteStreamInput::new(stream_name))
115 .await?;
116
117 println!("Streams examples completed");
118
119 // The following read session examples are for documentation snippets only.
120 // They are not executed because they would block waiting for new records.
121 if std::env::var("RUN_READ_SESSIONS").is_err() {
122 return Ok(());
123 }
124
125 // ANCHOR: read-session
126 let mut session = stream
127 .read_session(
128 ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0))),
129 ReadSessionConfig::default(),
130 )
131 .await?;
132
133 while let Some(batch) = session.next().await {
134 let batch = batch?;
135 for record in batch.records {
136 println!("[{}] {:?}", record.seq_num, record.body);
137 }
138 }
139 // ANCHOR_END: read-session
140
141 // ANCHOR: read-session-tail-offset
142 // Start reading from 10 records before the current tail
143 let mut session = stream
144 .read_session(
145 ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::TailOffset(10))),
146 ReadSessionConfig::default(),
147 )
148 .await?;
149
150 while let Some(batch) = session.next().await {
151 let batch = batch?;
152 for record in batch.records {
153 println!("[{}] {:?}", record.seq_num, record.body);
154 }
155 }
156 // ANCHOR_END: read-session-tail-offset
157
158 // ANCHOR: read-session-timestamp
159 // Start reading from a specific timestamp
160 let one_hour_ago = std::time::SystemTime::now()
161 .duration_since(std::time::UNIX_EPOCH)?
162 .as_millis() as u64
163 - 3600 * 1000;
164 let mut session = stream
165 .read_session(
166 ReadInput::new()
167 .with_start(ReadStart::new().with_from(ReadFrom::Timestamp(one_hour_ago))),
168 ReadSessionConfig::default(),
169 )
170 .await?;
171
172 while let Some(batch) = session.next().await {
173 let batch = batch?;
174 for record in batch.records {
175 println!("[{}] {:?}", record.seq_num, record.body);
176 }
177 }
178 // ANCHOR_END: read-session-timestamp
179
180 // ANCHOR: read-session-until
181 // Read records until a specific timestamp
182 let one_hour_ago = std::time::SystemTime::now()
183 .duration_since(std::time::UNIX_EPOCH)?
184 .as_millis() as u64
185 - 3600 * 1000;
186 let mut session = stream
187 .read_session(
188 ReadInput::new()
189 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
190 .with_stop(ReadStop::new().with_until(..one_hour_ago)),
191 ReadSessionConfig::default(),
192 )
193 .await?;
194
195 while let Some(batch) = session.next().await {
196 let batch = batch?;
197 for record in batch.records {
198 println!("[{}] {:?}", record.seq_num, record.body);
199 }
200 }
201 // ANCHOR_END: read-session-until
202
203 // ANCHOR: read-session-wait
204 // Read all available records, and once reaching the current tail, wait an additional 30 seconds
205 // for new ones
206 let mut session = stream
207 .read_session(
208 ReadInput::new()
209 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
210 .with_stop(ReadStop::new().with_wait(30)),
211 ReadSessionConfig::default(),
212 )
213 .await?;
214
215 while let Some(batch) = session.next().await {
216 let batch = batch?;
217 for record in batch.records {
218 println!("[{}] {:?}", record.seq_num, record.body);
219 }
220 }
221 // ANCHOR_END: read-session-wait
222
223 Ok(())
224}Sourcepub fn with_stop(self, stop: ReadStop) -> Self
pub fn with_stop(self, stop: ReadStop) -> Self
Set when to stop reading.
Examples found in repository?
examples/get_latest_record.rs (line 24)
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}More examples
examples/docs_encryption.rs (line 64)
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/docs_streams.rs (line 62)
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 // Force the partial batch and wait for durability without closing the producer
98 producer.flush().await?;
99
100 // Get the exact sequence number
101 let ack = ticket.await?;
102 println!("Record durable at seqNum {}", ack.seq_num);
103
104 producer.close().await?;
105 // ANCHOR_END: producer
106
107 // ANCHOR: check-tail
108 let tail = stream.check_tail().await?;
109 println!("Stream has {} records", tail.seq_num);
110 // ANCHOR_END: check-tail
111
112 // Cleanup
113 basin
114 .delete_stream(s2_sdk::types::DeleteStreamInput::new(stream_name))
115 .await?;
116
117 println!("Streams examples completed");
118
119 // The following read session examples are for documentation snippets only.
120 // They are not executed because they would block waiting for new records.
121 if std::env::var("RUN_READ_SESSIONS").is_err() {
122 return Ok(());
123 }
124
125 // ANCHOR: read-session
126 let mut session = stream
127 .read_session(
128 ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0))),
129 ReadSessionConfig::default(),
130 )
131 .await?;
132
133 while let Some(batch) = session.next().await {
134 let batch = batch?;
135 for record in batch.records {
136 println!("[{}] {:?}", record.seq_num, record.body);
137 }
138 }
139 // ANCHOR_END: read-session
140
141 // ANCHOR: read-session-tail-offset
142 // Start reading from 10 records before the current tail
143 let mut session = stream
144 .read_session(
145 ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::TailOffset(10))),
146 ReadSessionConfig::default(),
147 )
148 .await?;
149
150 while let Some(batch) = session.next().await {
151 let batch = batch?;
152 for record in batch.records {
153 println!("[{}] {:?}", record.seq_num, record.body);
154 }
155 }
156 // ANCHOR_END: read-session-tail-offset
157
158 // ANCHOR: read-session-timestamp
159 // Start reading from a specific timestamp
160 let one_hour_ago = std::time::SystemTime::now()
161 .duration_since(std::time::UNIX_EPOCH)?
162 .as_millis() as u64
163 - 3600 * 1000;
164 let mut session = stream
165 .read_session(
166 ReadInput::new()
167 .with_start(ReadStart::new().with_from(ReadFrom::Timestamp(one_hour_ago))),
168 ReadSessionConfig::default(),
169 )
170 .await?;
171
172 while let Some(batch) = session.next().await {
173 let batch = batch?;
174 for record in batch.records {
175 println!("[{}] {:?}", record.seq_num, record.body);
176 }
177 }
178 // ANCHOR_END: read-session-timestamp
179
180 // ANCHOR: read-session-until
181 // Read records until a specific timestamp
182 let one_hour_ago = std::time::SystemTime::now()
183 .duration_since(std::time::UNIX_EPOCH)?
184 .as_millis() as u64
185 - 3600 * 1000;
186 let mut session = stream
187 .read_session(
188 ReadInput::new()
189 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
190 .with_stop(ReadStop::new().with_until(..one_hour_ago)),
191 ReadSessionConfig::default(),
192 )
193 .await?;
194
195 while let Some(batch) = session.next().await {
196 let batch = batch?;
197 for record in batch.records {
198 println!("[{}] {:?}", record.seq_num, record.body);
199 }
200 }
201 // ANCHOR_END: read-session-until
202
203 // ANCHOR: read-session-wait
204 // Read all available records, and once reaching the current tail, wait an additional 30 seconds
205 // for new ones
206 let mut session = stream
207 .read_session(
208 ReadInput::new()
209 .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
210 .with_stop(ReadStop::new().with_wait(30)),
211 ReadSessionConfig::default(),
212 )
213 .await?;
214
215 while let Some(batch) = session.next().await {
216 let batch = batch?;
217 for record in batch.records {
218 println!("[{}] {:?}", record.seq_num, record.body);
219 }
220 }
221 // ANCHOR_END: read-session-wait
222
223 Ok(())
224}Sourcepub fn with_ignore_command_records(self, ignore_command_records: bool) -> Self
pub fn with_ignore_command_records(self, ignore_command_records: bool) -> Self
Set whether to filter out command records from the stream when reading.
Trait Implementations§
Auto Trait Implementations§
impl Freeze for ReadInput
impl RefUnwindSafe for ReadInput
impl Send for ReadInput
impl Sync for ReadInput
impl Unpin for ReadInput
impl UnsafeUnpin for ReadInput
impl UnwindSafe for ReadInput
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