docs_streams/
docs_streams.rs1use std::time::Duration;
6
7use futures_util::StreamExt;
8use s2_sdk::{
9 S2,
10 append_session::AppendSessionConfig,
11 batching::BatchingConfig,
12 producer::ProducerConfig,
13 types::{
14 AppendInput, AppendRecord, AppendRecordBatch, BasinName, ReadFrom, ReadInput, ReadLimits,
15 ReadSessionConfig, ReadStart, ReadStop, S2Config, StreamName,
16 },
17};
18
19#[tokio::main]
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 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 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 println!(
51 "Wrote records {} through {}",
52 ack.start.seq_num,
53 ack.end.seq_num - 1
54 );
55 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 let session = stream.append_session(AppendSessionConfig::new());
73
74 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 let ack = ticket.await?;
83 println!("Durable at seqNum {}", ack.start.seq_num);
84
85 session.close().await?;
86 let producer = stream.producer(
90 ProducerConfig::new()
91 .with_batching(BatchingConfig::new().with_linger(Duration::from_millis(5))),
92 );
93
94 let ticket = producer.submit(AppendRecord::new("my event")?).await?;
96
97 let ack = ticket.await?;
99 println!("Record durable at seqNum {}", ack.seq_num);
100
101 producer.close().await?;
102 let tail = stream.check_tail().await?;
106 println!("Stream has {} records", tail.seq_num);
107 basin
111 .delete_stream(s2_sdk::types::DeleteStreamInput::new(stream_name))
112 .await?;
113
114 println!("Streams examples completed");
115
116 if std::env::var("RUN_READ_SESSIONS").is_err() {
119 return Ok(());
120 }
121
122 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 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 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 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 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 Ok(())
221}