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 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 )
55 .await?;
56 let mut caught_up = session.caught_up();
57
58 loop {
59 tokio::select! {
60 tail = &mut caught_up => {
61 println!("Caught up through sequence number {}", tail?.seq_num);
62 break;
63 }
64 Some(batch) = session.next() => {
65 print_batch("Read before catching up", &batch?);
66 }
67 }
68 }
69
70 let ack = stream
71 .append(AppendInput::new(AppendRecordBatch::try_from_iter([
72 AppendRecord::new("third")?,
73 ])?))
74 .await?;
75 println!(
76 "Appended another record at sequence number {}",
77 ack.start.seq_num
78 );
79
80 while let Some(batch) = session.next().await {
81 let batch = batch?;
82 print_batch("Read after catching up", &batch);
83 if batch
84 .records
85 .iter()
86 .any(|record| record.seq_num == ack.start.seq_num)
87 {
88 break;
89 }
90 }
91 println!("Session is caught up again: {}", session.is_caught_up());
92
93 drop(session);
94 basin
95 .delete_stream(DeleteStreamInput::new(stream_name))
96 .await?;
97 s2.delete_basin(DeleteBasinInput::new(basin_name)).await?;
98
99 Ok(())
100}