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        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}