Skip to main content

ReadInput

Struct ReadInput 

Source
#[non_exhaustive]
pub struct ReadInput { pub start: ReadStart, pub stop: ReadStop, pub ignore_command_records: bool, pub stream_config: Option<StreamConfig>, }
Expand description

Input for read and read_session operations.

Fields (Non-exhaustive)§

This struct is marked as non-exhaustive
Non-exhaustive structs could have additional fields added in future. Therefore, non-exhaustive structs cannot be constructed in external crates using the traditional Struct { .. } syntax; cannot be matched against without a wildcard ..; and struct update syntax will not work.
§start: ReadStart

Where to start reading.

See ReadStart for defaults.

§stop: ReadStop

When to stop reading.

See ReadStop for defaults.

§ignore_command_records: bool

Whether to filter out command records from the stream when reading.

Defaults to false.

§stream_config: Option<StreamConfig>

Stream configuration to apply if the stream is created on read.

Unset fields inherit the basin’s default stream configuration. Ignored if the stream already exists.

Implementations§

Source§

impl ReadInput

Source

pub fn new() -> Self

Create a new ReadInput with default values.

Examples found in repository?
examples/consumer.rs (line 17)
9async fn main() -> Result<(), Box<dyn std::error::Error>> {
10    let access_token = std::env::var("S2_ACCESS_TOKEN")?;
11    let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
12    let stream_name: StreamName = std::env::var("S2_STREAM")?.parse()?;
13
14    let s2 = S2::new(S2Config::new(access_token))?;
15    let stream = s2.basin(basin_name).stream(stream_name);
16
17    let input = ReadInput::new();
18    let mut batches = stream
19        .read_session(input, ReadSessionConfig::default())
20        .await?;
21    loop {
22        select! {
23            batch = batches.next() => {
24                let Some(batch) = batch else { break };
25                let batch = batch?;
26                println!("{batch:?}");
27            }
28            _ = tokio::signal::ctrl_c() => break,
29        }
30    }
31
32    Ok(())
33}
More examples
Hide additional examples
examples/get_latest_record.rs (line 22)
9async fn main() -> Result<(), Box<dyn std::error::Error>> {
10    let access_token =
11        std::env::var("S2_ACCESS_TOKEN").map_err(|_| "S2_ACCESS_TOKEN env var not set")?;
12    let basin_name: BasinName = std::env::var("S2_BASIN")
13        .map_err(|_| "S2_BASIN env var not set")?
14        .parse()?;
15    let stream_name: StreamName = std::env::var("S2_STREAM")
16        .map_err(|_| "S2_STREAM env var not set")?
17        .parse()?;
18
19    let s2 = S2::new(S2Config::new(access_token))?;
20    let stream = s2.basin(basin_name).stream(stream_name);
21
22    let input = ReadInput::new()
23        .with_start(ReadStart::new().with_from(ReadFrom::TailOffset(1)))
24        .with_stop(ReadStop::new().with_limits(ReadLimits::new().with_count(1)));
25    let batch = stream.read(input).await?;
26    println!("{batch:#?}");
27
28    Ok(())
29}
examples/docs_encryption.rs (line 62)
15async fn main() -> Result<(), Box<dyn std::error::Error>> {
16    let access_token = std::env::var("S2_ACCESS_TOKEN")?;
17    let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
18    let stream_name: StreamName = format!(
19        "docs-encryption-{}",
20        std::time::SystemTime::now()
21            .duration_since(std::time::UNIX_EPOCH)?
22            .as_millis()
23    )
24    .parse()?;
25
26    let client = S2::new(S2Config::new(access_token))?;
27
28    // ANCHOR: basin-cipher
29    client
30        .create_basin(
31            CreateBasinInput::new(basin_name.clone())
32                .with_config(BasinConfig::new().with_stream_cipher(EncryptionAlgorithm::Aegis256)),
33        )
34        .await?;
35
36    client
37        .reconfigure_basin(ReconfigureBasinInput::new(
38            basin_name.clone(),
39            BasinReconfiguration::new().with_stream_cipher(EncryptionAlgorithm::Aes256Gcm),
40        ))
41        .await?;
42    // ANCHOR_END: basin-cipher
43
44    let basin = client.basin(basin_name.clone());
45    basin
46        .create_stream(CreateStreamInput::new(stream_name.clone()))
47        .await?;
48
49    // ANCHOR: append-read
50    let stream = basin
51        .stream(stream_name.clone())
52        .with_encryption_key(std::env::var("S2_ENCRYPTION_KEY")?.parse()?);
53
54    stream
55        .append(AppendInput::new(AppendRecordBatch::try_from_iter([
56            AppendRecord::new("top secret")?,
57        ])?))
58        .await?;
59
60    let batch = stream
61        .read(
62            ReadInput::new()
63                .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
64                .with_stop(ReadStop::new().with_limits(ReadLimits::new().with_count(10))),
65        )
66        .await?;
67    // ANCHOR_END: append-read
68
69    println!("Read {} encrypted record(s)", batch.records.len());
70
71    basin
72        .delete_stream(DeleteStreamInput::new(stream_name))
73        .await?;
74
75    Ok(())
76}
examples/caught_up.rs (line 53)
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}
examples/docs_streams.rs (line 60)
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    // Create a temporary stream for examples
28    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    // ANCHOR: simple-append
40    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    // ack tells us where the records landed
50    println!(
51        "Wrote records {} through {}",
52        ack.start.seq_num,
53        ack.end.seq_num - 1
54    );
55    // ANCHOR_END: simple-append
56
57    // ANCHOR: simple-read
58    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    // ANCHOR_END: simple-read
70
71    // ANCHOR: append-session
72    let session = stream.append_session(AppendSessionConfig::new());
73
74    // Submit a batch - this enqueues it and returns a ticket
75    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    // Wait for durability
82    let ack = ticket.await?;
83    println!("Durable at seqNum {}", ack.start.seq_num);
84
85    session.close().await?;
86    // ANCHOR_END: append-session
87
88    // ANCHOR: producer
89    let producer = stream.producer(
90        ProducerConfig::new()
91            .with_batching(BatchingConfig::new().with_linger(Duration::from_millis(5))),
92    );
93
94    // Submit individual records
95    let ticket = producer.submit(AppendRecord::new("my event")?).await?;
96
97    // Force the partial batch and wait for durability without closing the producer
98    producer.flush().await?;
99
100    // Get the exact sequence number
101    let ack = ticket.await?;
102    println!("Record durable at seqNum {}", ack.seq_num);
103
104    producer.close().await?;
105    // ANCHOR_END: producer
106
107    // ANCHOR: check-tail
108    let tail = stream.check_tail().await?;
109    println!("Stream has {} records", tail.seq_num);
110    // ANCHOR_END: check-tail
111
112    // Cleanup
113    basin
114        .delete_stream(s2_sdk::types::DeleteStreamInput::new(stream_name))
115        .await?;
116
117    println!("Streams examples completed");
118
119    // The following read session examples are for documentation snippets only.
120    // They are not executed because they would block waiting for new records.
121    if std::env::var("RUN_READ_SESSIONS").is_err() {
122        return Ok(());
123    }
124
125    // ANCHOR: read-session
126    let mut session = stream
127        .read_session(
128            ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0))),
129            ReadSessionConfig::default(),
130        )
131        .await?;
132
133    while let Some(batch) = session.next().await {
134        let batch = batch?;
135        for record in batch.records {
136            println!("[{}] {:?}", record.seq_num, record.body);
137        }
138    }
139    // ANCHOR_END: read-session
140
141    // ANCHOR: read-session-tail-offset
142    // Start reading from 10 records before the current tail
143    let mut session = stream
144        .read_session(
145            ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::TailOffset(10))),
146            ReadSessionConfig::default(),
147        )
148        .await?;
149
150    while let Some(batch) = session.next().await {
151        let batch = batch?;
152        for record in batch.records {
153            println!("[{}] {:?}", record.seq_num, record.body);
154        }
155    }
156    // ANCHOR_END: read-session-tail-offset
157
158    // ANCHOR: read-session-timestamp
159    // Start reading from a specific timestamp
160    let one_hour_ago = std::time::SystemTime::now()
161        .duration_since(std::time::UNIX_EPOCH)?
162        .as_millis() as u64
163        - 3600 * 1000;
164    let mut session = stream
165        .read_session(
166            ReadInput::new()
167                .with_start(ReadStart::new().with_from(ReadFrom::Timestamp(one_hour_ago))),
168            ReadSessionConfig::default(),
169        )
170        .await?;
171
172    while let Some(batch) = session.next().await {
173        let batch = batch?;
174        for record in batch.records {
175            println!("[{}] {:?}", record.seq_num, record.body);
176        }
177    }
178    // ANCHOR_END: read-session-timestamp
179
180    // ANCHOR: read-session-until
181    // Read records until a specific timestamp
182    let one_hour_ago = std::time::SystemTime::now()
183        .duration_since(std::time::UNIX_EPOCH)?
184        .as_millis() as u64
185        - 3600 * 1000;
186    let mut session = stream
187        .read_session(
188            ReadInput::new()
189                .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
190                .with_stop(ReadStop::new().with_until(..one_hour_ago)),
191            ReadSessionConfig::default(),
192        )
193        .await?;
194
195    while let Some(batch) = session.next().await {
196        let batch = batch?;
197        for record in batch.records {
198            println!("[{}] {:?}", record.seq_num, record.body);
199        }
200    }
201    // ANCHOR_END: read-session-until
202
203    // ANCHOR: read-session-wait
204    // Read all available records, and once reaching the current tail, wait an additional 30 seconds
205    // for new ones
206    let mut session = stream
207        .read_session(
208            ReadInput::new()
209                .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
210                .with_stop(ReadStop::new().with_wait(30)),
211            ReadSessionConfig::default(),
212        )
213        .await?;
214
215    while let Some(batch) = session.next().await {
216        let batch = batch?;
217        for record in batch.records {
218            println!("[{}] {:?}", record.seq_num, record.body);
219        }
220    }
221    // ANCHOR_END: read-session-wait
222
223    Ok(())
224}
Source

pub fn with_start(self, start: ReadStart) -> Self

Set where to start reading.

Examples found in repository?
examples/get_latest_record.rs (line 23)
9async fn main() -> Result<(), Box<dyn std::error::Error>> {
10    let access_token =
11        std::env::var("S2_ACCESS_TOKEN").map_err(|_| "S2_ACCESS_TOKEN env var not set")?;
12    let basin_name: BasinName = std::env::var("S2_BASIN")
13        .map_err(|_| "S2_BASIN env var not set")?
14        .parse()?;
15    let stream_name: StreamName = std::env::var("S2_STREAM")
16        .map_err(|_| "S2_STREAM env var not set")?
17        .parse()?;
18
19    let s2 = S2::new(S2Config::new(access_token))?;
20    let stream = s2.basin(basin_name).stream(stream_name);
21
22    let input = ReadInput::new()
23        .with_start(ReadStart::new().with_from(ReadFrom::TailOffset(1)))
24        .with_stop(ReadStop::new().with_limits(ReadLimits::new().with_count(1)));
25    let batch = stream.read(input).await?;
26    println!("{batch:#?}");
27
28    Ok(())
29}
More examples
Hide additional examples
examples/docs_encryption.rs (line 63)
15async fn main() -> Result<(), Box<dyn std::error::Error>> {
16    let access_token = std::env::var("S2_ACCESS_TOKEN")?;
17    let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
18    let stream_name: StreamName = format!(
19        "docs-encryption-{}",
20        std::time::SystemTime::now()
21            .duration_since(std::time::UNIX_EPOCH)?
22            .as_millis()
23    )
24    .parse()?;
25
26    let client = S2::new(S2Config::new(access_token))?;
27
28    // ANCHOR: basin-cipher
29    client
30        .create_basin(
31            CreateBasinInput::new(basin_name.clone())
32                .with_config(BasinConfig::new().with_stream_cipher(EncryptionAlgorithm::Aegis256)),
33        )
34        .await?;
35
36    client
37        .reconfigure_basin(ReconfigureBasinInput::new(
38            basin_name.clone(),
39            BasinReconfiguration::new().with_stream_cipher(EncryptionAlgorithm::Aes256Gcm),
40        ))
41        .await?;
42    // ANCHOR_END: basin-cipher
43
44    let basin = client.basin(basin_name.clone());
45    basin
46        .create_stream(CreateStreamInput::new(stream_name.clone()))
47        .await?;
48
49    // ANCHOR: append-read
50    let stream = basin
51        .stream(stream_name.clone())
52        .with_encryption_key(std::env::var("S2_ENCRYPTION_KEY")?.parse()?);
53
54    stream
55        .append(AppendInput::new(AppendRecordBatch::try_from_iter([
56            AppendRecord::new("top secret")?,
57        ])?))
58        .await?;
59
60    let batch = stream
61        .read(
62            ReadInput::new()
63                .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
64                .with_stop(ReadStop::new().with_limits(ReadLimits::new().with_count(10))),
65        )
66        .await?;
67    // ANCHOR_END: append-read
68
69    println!("Read {} encrypted record(s)", batch.records.len());
70
71    basin
72        .delete_stream(DeleteStreamInput::new(stream_name))
73        .await?;
74
75    Ok(())
76}
examples/caught_up.rs (line 53)
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}
examples/docs_streams.rs (line 61)
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    // Create a temporary stream for examples
28    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    // ANCHOR: simple-append
40    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    // ack tells us where the records landed
50    println!(
51        "Wrote records {} through {}",
52        ack.start.seq_num,
53        ack.end.seq_num - 1
54    );
55    // ANCHOR_END: simple-append
56
57    // ANCHOR: simple-read
58    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    // ANCHOR_END: simple-read
70
71    // ANCHOR: append-session
72    let session = stream.append_session(AppendSessionConfig::new());
73
74    // Submit a batch - this enqueues it and returns a ticket
75    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    // Wait for durability
82    let ack = ticket.await?;
83    println!("Durable at seqNum {}", ack.start.seq_num);
84
85    session.close().await?;
86    // ANCHOR_END: append-session
87
88    // ANCHOR: producer
89    let producer = stream.producer(
90        ProducerConfig::new()
91            .with_batching(BatchingConfig::new().with_linger(Duration::from_millis(5))),
92    );
93
94    // Submit individual records
95    let ticket = producer.submit(AppendRecord::new("my event")?).await?;
96
97    // Force the partial batch and wait for durability without closing the producer
98    producer.flush().await?;
99
100    // Get the exact sequence number
101    let ack = ticket.await?;
102    println!("Record durable at seqNum {}", ack.seq_num);
103
104    producer.close().await?;
105    // ANCHOR_END: producer
106
107    // ANCHOR: check-tail
108    let tail = stream.check_tail().await?;
109    println!("Stream has {} records", tail.seq_num);
110    // ANCHOR_END: check-tail
111
112    // Cleanup
113    basin
114        .delete_stream(s2_sdk::types::DeleteStreamInput::new(stream_name))
115        .await?;
116
117    println!("Streams examples completed");
118
119    // The following read session examples are for documentation snippets only.
120    // They are not executed because they would block waiting for new records.
121    if std::env::var("RUN_READ_SESSIONS").is_err() {
122        return Ok(());
123    }
124
125    // ANCHOR: read-session
126    let mut session = stream
127        .read_session(
128            ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0))),
129            ReadSessionConfig::default(),
130        )
131        .await?;
132
133    while let Some(batch) = session.next().await {
134        let batch = batch?;
135        for record in batch.records {
136            println!("[{}] {:?}", record.seq_num, record.body);
137        }
138    }
139    // ANCHOR_END: read-session
140
141    // ANCHOR: read-session-tail-offset
142    // Start reading from 10 records before the current tail
143    let mut session = stream
144        .read_session(
145            ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::TailOffset(10))),
146            ReadSessionConfig::default(),
147        )
148        .await?;
149
150    while let Some(batch) = session.next().await {
151        let batch = batch?;
152        for record in batch.records {
153            println!("[{}] {:?}", record.seq_num, record.body);
154        }
155    }
156    // ANCHOR_END: read-session-tail-offset
157
158    // ANCHOR: read-session-timestamp
159    // Start reading from a specific timestamp
160    let one_hour_ago = std::time::SystemTime::now()
161        .duration_since(std::time::UNIX_EPOCH)?
162        .as_millis() as u64
163        - 3600 * 1000;
164    let mut session = stream
165        .read_session(
166            ReadInput::new()
167                .with_start(ReadStart::new().with_from(ReadFrom::Timestamp(one_hour_ago))),
168            ReadSessionConfig::default(),
169        )
170        .await?;
171
172    while let Some(batch) = session.next().await {
173        let batch = batch?;
174        for record in batch.records {
175            println!("[{}] {:?}", record.seq_num, record.body);
176        }
177    }
178    // ANCHOR_END: read-session-timestamp
179
180    // ANCHOR: read-session-until
181    // Read records until a specific timestamp
182    let one_hour_ago = std::time::SystemTime::now()
183        .duration_since(std::time::UNIX_EPOCH)?
184        .as_millis() as u64
185        - 3600 * 1000;
186    let mut session = stream
187        .read_session(
188            ReadInput::new()
189                .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
190                .with_stop(ReadStop::new().with_until(..one_hour_ago)),
191            ReadSessionConfig::default(),
192        )
193        .await?;
194
195    while let Some(batch) = session.next().await {
196        let batch = batch?;
197        for record in batch.records {
198            println!("[{}] {:?}", record.seq_num, record.body);
199        }
200    }
201    // ANCHOR_END: read-session-until
202
203    // ANCHOR: read-session-wait
204    // Read all available records, and once reaching the current tail, wait an additional 30 seconds
205    // for new ones
206    let mut session = stream
207        .read_session(
208            ReadInput::new()
209                .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
210                .with_stop(ReadStop::new().with_wait(30)),
211            ReadSessionConfig::default(),
212        )
213        .await?;
214
215    while let Some(batch) = session.next().await {
216        let batch = batch?;
217        for record in batch.records {
218            println!("[{}] {:?}", record.seq_num, record.body);
219        }
220    }
221    // ANCHOR_END: read-session-wait
222
223    Ok(())
224}
Source

pub fn with_stop(self, stop: ReadStop) -> Self

Set when to stop reading.

Examples found in repository?
examples/get_latest_record.rs (line 24)
9async fn main() -> Result<(), Box<dyn std::error::Error>> {
10    let access_token =
11        std::env::var("S2_ACCESS_TOKEN").map_err(|_| "S2_ACCESS_TOKEN env var not set")?;
12    let basin_name: BasinName = std::env::var("S2_BASIN")
13        .map_err(|_| "S2_BASIN env var not set")?
14        .parse()?;
15    let stream_name: StreamName = std::env::var("S2_STREAM")
16        .map_err(|_| "S2_STREAM env var not set")?
17        .parse()?;
18
19    let s2 = S2::new(S2Config::new(access_token))?;
20    let stream = s2.basin(basin_name).stream(stream_name);
21
22    let input = ReadInput::new()
23        .with_start(ReadStart::new().with_from(ReadFrom::TailOffset(1)))
24        .with_stop(ReadStop::new().with_limits(ReadLimits::new().with_count(1)));
25    let batch = stream.read(input).await?;
26    println!("{batch:#?}");
27
28    Ok(())
29}
More examples
Hide additional examples
examples/docs_encryption.rs (line 64)
15async fn main() -> Result<(), Box<dyn std::error::Error>> {
16    let access_token = std::env::var("S2_ACCESS_TOKEN")?;
17    let basin_name: BasinName = std::env::var("S2_BASIN")?.parse()?;
18    let stream_name: StreamName = format!(
19        "docs-encryption-{}",
20        std::time::SystemTime::now()
21            .duration_since(std::time::UNIX_EPOCH)?
22            .as_millis()
23    )
24    .parse()?;
25
26    let client = S2::new(S2Config::new(access_token))?;
27
28    // ANCHOR: basin-cipher
29    client
30        .create_basin(
31            CreateBasinInput::new(basin_name.clone())
32                .with_config(BasinConfig::new().with_stream_cipher(EncryptionAlgorithm::Aegis256)),
33        )
34        .await?;
35
36    client
37        .reconfigure_basin(ReconfigureBasinInput::new(
38            basin_name.clone(),
39            BasinReconfiguration::new().with_stream_cipher(EncryptionAlgorithm::Aes256Gcm),
40        ))
41        .await?;
42    // ANCHOR_END: basin-cipher
43
44    let basin = client.basin(basin_name.clone());
45    basin
46        .create_stream(CreateStreamInput::new(stream_name.clone()))
47        .await?;
48
49    // ANCHOR: append-read
50    let stream = basin
51        .stream(stream_name.clone())
52        .with_encryption_key(std::env::var("S2_ENCRYPTION_KEY")?.parse()?);
53
54    stream
55        .append(AppendInput::new(AppendRecordBatch::try_from_iter([
56            AppendRecord::new("top secret")?,
57        ])?))
58        .await?;
59
60    let batch = stream
61        .read(
62            ReadInput::new()
63                .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
64                .with_stop(ReadStop::new().with_limits(ReadLimits::new().with_count(10))),
65        )
66        .await?;
67    // ANCHOR_END: append-read
68
69    println!("Read {} encrypted record(s)", batch.records.len());
70
71    basin
72        .delete_stream(DeleteStreamInput::new(stream_name))
73        .await?;
74
75    Ok(())
76}
examples/docs_streams.rs (line 62)
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    // Create a temporary stream for examples
28    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    // ANCHOR: simple-append
40    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    // ack tells us where the records landed
50    println!(
51        "Wrote records {} through {}",
52        ack.start.seq_num,
53        ack.end.seq_num - 1
54    );
55    // ANCHOR_END: simple-append
56
57    // ANCHOR: simple-read
58    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    // ANCHOR_END: simple-read
70
71    // ANCHOR: append-session
72    let session = stream.append_session(AppendSessionConfig::new());
73
74    // Submit a batch - this enqueues it and returns a ticket
75    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    // Wait for durability
82    let ack = ticket.await?;
83    println!("Durable at seqNum {}", ack.start.seq_num);
84
85    session.close().await?;
86    // ANCHOR_END: append-session
87
88    // ANCHOR: producer
89    let producer = stream.producer(
90        ProducerConfig::new()
91            .with_batching(BatchingConfig::new().with_linger(Duration::from_millis(5))),
92    );
93
94    // Submit individual records
95    let ticket = producer.submit(AppendRecord::new("my event")?).await?;
96
97    // Force the partial batch and wait for durability without closing the producer
98    producer.flush().await?;
99
100    // Get the exact sequence number
101    let ack = ticket.await?;
102    println!("Record durable at seqNum {}", ack.seq_num);
103
104    producer.close().await?;
105    // ANCHOR_END: producer
106
107    // ANCHOR: check-tail
108    let tail = stream.check_tail().await?;
109    println!("Stream has {} records", tail.seq_num);
110    // ANCHOR_END: check-tail
111
112    // Cleanup
113    basin
114        .delete_stream(s2_sdk::types::DeleteStreamInput::new(stream_name))
115        .await?;
116
117    println!("Streams examples completed");
118
119    // The following read session examples are for documentation snippets only.
120    // They are not executed because they would block waiting for new records.
121    if std::env::var("RUN_READ_SESSIONS").is_err() {
122        return Ok(());
123    }
124
125    // ANCHOR: read-session
126    let mut session = stream
127        .read_session(
128            ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0))),
129            ReadSessionConfig::default(),
130        )
131        .await?;
132
133    while let Some(batch) = session.next().await {
134        let batch = batch?;
135        for record in batch.records {
136            println!("[{}] {:?}", record.seq_num, record.body);
137        }
138    }
139    // ANCHOR_END: read-session
140
141    // ANCHOR: read-session-tail-offset
142    // Start reading from 10 records before the current tail
143    let mut session = stream
144        .read_session(
145            ReadInput::new().with_start(ReadStart::new().with_from(ReadFrom::TailOffset(10))),
146            ReadSessionConfig::default(),
147        )
148        .await?;
149
150    while let Some(batch) = session.next().await {
151        let batch = batch?;
152        for record in batch.records {
153            println!("[{}] {:?}", record.seq_num, record.body);
154        }
155    }
156    // ANCHOR_END: read-session-tail-offset
157
158    // ANCHOR: read-session-timestamp
159    // Start reading from a specific timestamp
160    let one_hour_ago = std::time::SystemTime::now()
161        .duration_since(std::time::UNIX_EPOCH)?
162        .as_millis() as u64
163        - 3600 * 1000;
164    let mut session = stream
165        .read_session(
166            ReadInput::new()
167                .with_start(ReadStart::new().with_from(ReadFrom::Timestamp(one_hour_ago))),
168            ReadSessionConfig::default(),
169        )
170        .await?;
171
172    while let Some(batch) = session.next().await {
173        let batch = batch?;
174        for record in batch.records {
175            println!("[{}] {:?}", record.seq_num, record.body);
176        }
177    }
178    // ANCHOR_END: read-session-timestamp
179
180    // ANCHOR: read-session-until
181    // Read records until a specific timestamp
182    let one_hour_ago = std::time::SystemTime::now()
183        .duration_since(std::time::UNIX_EPOCH)?
184        .as_millis() as u64
185        - 3600 * 1000;
186    let mut session = stream
187        .read_session(
188            ReadInput::new()
189                .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
190                .with_stop(ReadStop::new().with_until(..one_hour_ago)),
191            ReadSessionConfig::default(),
192        )
193        .await?;
194
195    while let Some(batch) = session.next().await {
196        let batch = batch?;
197        for record in batch.records {
198            println!("[{}] {:?}", record.seq_num, record.body);
199        }
200    }
201    // ANCHOR_END: read-session-until
202
203    // ANCHOR: read-session-wait
204    // Read all available records, and once reaching the current tail, wait an additional 30 seconds
205    // for new ones
206    let mut session = stream
207        .read_session(
208            ReadInput::new()
209                .with_start(ReadStart::new().with_from(ReadFrom::SeqNum(0)))
210                .with_stop(ReadStop::new().with_wait(30)),
211            ReadSessionConfig::default(),
212        )
213        .await?;
214
215    while let Some(batch) = session.next().await {
216        let batch = batch?;
217        for record in batch.records {
218            println!("[{}] {:?}", record.seq_num, record.body);
219        }
220    }
221    // ANCHOR_END: read-session-wait
222
223    Ok(())
224}
Source

pub fn with_ignore_command_records(self, ignore_command_records: bool) -> Self

Set whether to filter out command records from the stream when reading.

Source

pub fn with_stream_config(self, stream_config: StreamConfig) -> Self

Set the stream configuration to apply if the stream is created on read.

Trait Implementations§

Source§

impl Clone for ReadInput

Source§

fn clone(&self) -> ReadInput

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Debug for ReadInput

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl Default for ReadInput

Source§

fn default() -> ReadInput

Returns the “default value” for a type. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more