Skip to main content

caught_up/
caught_up.rs

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}