1use futures_util::StreamExt;
2use s2_sdk::{
3 S2,
4 types::{
5 AppendInput, AppendRecord, AppendRecordBatch, BasinName, CreateBasinInput,
6 CreateStreamInput, DeleteBasinInput, DeleteStreamInput, ReadBatch, ReadFrom, ReadInput,
7 ReadSessionConfig, ReadStart, S2Config, S2Endpoints, StreamName,
8 },
9};
10
11fn print_batch(label: &str, batch: &ReadBatch) {
12 let seq_nums = batch
13 .records
14 .iter()
15 .map(|record| record.seq_num)
16 .collect::<Vec<_>>();
17 println!("{label}: {seq_nums:?}");
18}
19
20#[tokio::main]
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}