Skip to main content

s2_lite/backend/
streams.rs

1use s2_common::{
2    basin::BasinName,
3    config::{OptionalStreamConfig, StreamConfig, StreamReconfiguration},
4    record::StreamPosition,
5    resources::{Page, ProvisionMode, ProvisionResult, RequestToken},
6    stream::{ListStreamsRequest, StreamInfo, StreamName},
7};
8use s2_storage::bash::Bash;
9use slatedb::{
10    IsolationLevel,
11    config::{DurabilityLevel, ScanOptions},
12};
13use time::OffsetDateTime;
14use tracing::instrument;
15
16use super::{
17    Backend,
18    store::db_txn_get,
19    streamer::{TerminalTrimCondition, TerminalTrimOutcome, doe_arm_delay},
20};
21use crate::{
22    backend::{
23        error::{
24            BasinDeletionPendingError, BasinNotFoundError, DeleteStreamError, GetStreamConfigError,
25            ListStreamsError, ProvisionStreamError, ReconfigureStreamError, StorageError,
26            StreamAlreadyExistsError, StreamDeletionPendingError, StreamNotFoundError,
27            StreamerError,
28        },
29        kv,
30    },
31    stream_id::StreamId,
32};
33
34impl Backend {
35    pub async fn list_streams(
36        &self,
37        basin: BasinName,
38        request: ListStreamsRequest,
39    ) -> Result<Page<StreamInfo>, ListStreamsError> {
40        let ListStreamsRequest {
41            prefix,
42            start_after,
43            limit,
44        } = request;
45
46        let key_range = kv::stream_meta::ser_key_range(&basin, &prefix, &start_after);
47        if key_range.is_empty() {
48            return Ok(Page::new_empty());
49        }
50
51        let scan_opts = ScanOptions {
52            durability_filter: DurabilityLevel::Remote,
53            ..Default::default()
54        };
55        let mut it = self.db.scan_with_options(key_range, &scan_opts).await?;
56
57        let mut streams = Vec::with_capacity(limit.as_usize());
58        let mut has_more = false;
59        while let Some(kv) = it.next().await? {
60            let (deser_basin, stream) = kv::stream_meta::deser_key(kv.key)?;
61            assert_eq!(deser_basin.as_ref(), basin.as_ref());
62            assert!(stream.as_ref() > start_after.as_ref());
63            assert!(stream.as_ref() >= prefix.as_ref());
64            if streams.len() == limit.as_usize() {
65                has_more = true;
66                break;
67            }
68            let meta = kv::stream_meta::deser_value(kv.value)?;
69            streams.push(StreamInfo {
70                name: stream,
71                created_at: meta.created_at,
72                deleted_at: meta.deleted_at,
73                cipher: meta.cipher,
74            });
75        }
76        Ok(Page::new(streams, has_more))
77    }
78
79    /// Invariant: any outcome asserting the stream exists — `Created`,
80    /// `Updated`, `Noop`, or `StreamAlreadyExists` — is durably visible
81    /// (readable at `DurabilityLevel::Remote`) by the time this returns.
82    pub async fn provision_stream(
83        &self,
84        basin: BasinName,
85        stream: StreamName,
86        config: OptionalStreamConfig,
87        mode: ProvisionMode,
88    ) -> Result<ProvisionResult<StreamInfo>, ProvisionStreamError> {
89        let txn = self.db.begin(IsolationLevel::SerializableSnapshot).await?;
90
91        let Some(basin_meta) = db_txn_get(
92            &txn,
93            kv::basin_meta::ser_key(&basin),
94            kv::basin_meta::deser_value,
95        )
96        .await?
97        else {
98            return Err(BasinNotFoundError { basin }.into());
99        };
100
101        if basin_meta.deleted_at.is_some() {
102            return Err(BasinDeletionPendingError { basin }.into());
103        }
104
105        let stream_meta_key = kv::stream_meta::ser_key(&basin, &stream);
106
107        // Existence is decided from a Memory-level read; capture the row's
108        // commit seq so exists-outcomes can await durability (see fn doc).
109        let existing_entry = txn.get_key_value(&stream_meta_key).await?;
110        let existing_seq = existing_entry.as_ref().map(|kv| kv.seq);
111        let existing_meta = existing_entry
112            .map(|kv| kv::stream_meta::deser_value(kv.value))
113            .transpose()
114            .map_err(StorageError::from)?;
115        if let Some(existing_meta) = &existing_meta
116            && existing_meta.deleted_at.is_some()
117        {
118            return Err(ProvisionStreamError::StreamDeletionPending(
119                StreamDeletionPendingError,
120            ));
121        }
122
123        let basin_defaults = basin_meta.config.default_stream_config;
124        let (outcome, prior_doe_min_age) = match (existing_meta, mode) {
125            (Some(existing), ProvisionMode::CreateOnly { request_token }) => {
126                let new_creation_idempotency_key = request_token
127                    .as_ref()
128                    .map(|req_token| creation_idempotency_key(req_token, &config));
129                let result = if new_creation_idempotency_key.is_some()
130                    && existing.creation_idempotency_key == new_creation_idempotency_key
131                {
132                    Ok(ProvisionResult::Noop(StreamInfo {
133                        name: stream,
134                        created_at: existing.created_at,
135                        deleted_at: None,
136                        cipher: existing.cipher,
137                    }))
138                } else {
139                    Err(StreamAlreadyExistsError { basin, stream }.into())
140                };
141                drop(txn);
142                self.await_durable_seq(existing_seq.expect("existing meta was read"))
143                    .await?;
144                return result;
145            }
146            (Some(existing), ProvisionMode::Ensure) => {
147                let desired_config = config.merge(basin_defaults);
148                let config_unchanged = existing.config == desired_config;
149                let meta = kv::stream_meta::StreamMeta {
150                    config: desired_config,
151                    cipher: existing.cipher,
152                    created_at: existing.created_at,
153                    deleted_at: None,
154                    creation_idempotency_key: existing.creation_idempotency_key,
155                };
156                (
157                    if config_unchanged {
158                        ProvisionResult::Noop(meta)
159                    } else {
160                        ProvisionResult::Updated(meta)
161                    },
162                    existing.config.delete_on_empty.min_age(),
163                )
164            }
165            (None, ProvisionMode::CreateOnly { request_token }) => {
166                let new_creation_idempotency_key = request_token
167                    .as_ref()
168                    .map(|req_token| creation_idempotency_key(req_token, &config));
169                (
170                    ProvisionResult::Created(kv::stream_meta::StreamMeta {
171                        config: config.merge(basin_defaults),
172                        cipher: basin_meta.config.stream_cipher,
173                        created_at: OffsetDateTime::now_utc(),
174                        deleted_at: None,
175                        creation_idempotency_key: new_creation_idempotency_key,
176                    }),
177                    None,
178                )
179            }
180            (None, ProvisionMode::Ensure) => (
181                ProvisionResult::Created(kv::stream_meta::StreamMeta {
182                    config: config.merge(basin_defaults),
183                    cipher: basin_meta.config.stream_cipher,
184                    created_at: OffsetDateTime::now_utc(),
185                    deleted_at: None,
186                    creation_idempotency_key: None,
187                }),
188                None,
189            ),
190        };
191
192        if matches!(&outcome, ProvisionResult::Noop(_)) {
193            drop(txn);
194            self.await_durable_seq(existing_seq.expect("noop implies existing meta"))
195                .await?;
196        } else {
197            let meta = outcome.inner();
198
199            txn.put(&stream_meta_key, kv::stream_meta::ser_value(meta))?;
200
201            let stream_id = StreamId::new(&basin, &stream);
202
203            if matches!(&outcome, ProvisionResult::Created(_)) {
204                txn.put(
205                    kv::stream_id_mapping::ser_key(stream_id),
206                    kv::stream_id_mapping::ser_value(&basin, &stream),
207                )?;
208                txn.put(
209                    kv::stream_tail_position::ser_key(stream_id),
210                    kv::stream_tail_position::ser_value(StreamPosition::MIN),
211                )?;
212            }
213
214            if let Some(min_age) = meta.config.delete_on_empty.min_age()
215                && (matches!(&outcome, ProvisionResult::Created(_)) || prior_doe_min_age.is_none())
216            {
217                txn.put(
218                    kv::stream_doe_deadline::ser_key(
219                        kv::timestamp::TimestampSecs::after(doe_arm_delay(
220                            meta.config.retention_policy.age().unwrap_or_default(),
221                            min_age,
222                        )),
223                        stream_id,
224                    ),
225                    kv::stream_doe_deadline::ser_value(min_age),
226                )?;
227            }
228
229            txn.commit().await?;
230        }
231
232        if let ProvisionResult::Updated(meta) = &outcome
233            && let Some(client) = self.streamer_client_if_active(&basin, &stream)
234        {
235            client.advise_reconfig(meta.config.clone());
236        }
237
238        Ok(outcome.map(|meta| StreamInfo {
239            name: stream,
240            created_at: meta.created_at,
241            deleted_at: None,
242            cipher: meta.cipher,
243        }))
244    }
245
246    pub(super) async fn stream_id_mapping(
247        &self,
248        stream_id: StreamId,
249    ) -> Result<Option<(BasinName, StreamName)>, StorageError> {
250        self.db_get(
251            kv::stream_id_mapping::ser_key(stream_id),
252            kv::stream_id_mapping::deser_value,
253        )
254        .await
255    }
256
257    pub async fn get_stream_config(
258        &self,
259        basin: BasinName,
260        stream: StreamName,
261    ) -> Result<StreamConfig, GetStreamConfigError> {
262        let meta = self
263            .db_get(
264                kv::stream_meta::ser_key(&basin, &stream),
265                kv::stream_meta::deser_value,
266            )
267            .await?
268            .ok_or_else(|| StreamNotFoundError {
269                basin: basin.clone(),
270                stream: stream.clone(),
271            })?;
272        if meta.deleted_at.is_some() {
273            return Err(StreamDeletionPendingError.into());
274        }
275        Ok(meta.config)
276    }
277
278    pub async fn reconfigure_stream(
279        &self,
280        basin: BasinName,
281        stream: StreamName,
282        reconfig: StreamReconfiguration,
283    ) -> Result<StreamConfig, ReconfigureStreamError> {
284        let txn = self.db.begin(IsolationLevel::SerializableSnapshot).await?;
285
286        let meta_key = kv::stream_meta::ser_key(&basin, &stream);
287        let (basin_meta, meta) = tokio::try_join!(
288            db_txn_get(
289                &txn,
290                kv::basin_meta::ser_key(&basin),
291                kv::basin_meta::deser_value,
292            ),
293            db_txn_get(&txn, &meta_key, kv::stream_meta::deser_value),
294        )?;
295
296        let basin_meta = basin_meta.ok_or_else(|| BasinNotFoundError {
297            basin: basin.clone(),
298        })?;
299        if basin_meta.deleted_at.is_some() {
300            return Err(BasinDeletionPendingError { basin }.into());
301        }
302
303        let mut meta = meta.ok_or_else(|| StreamNotFoundError {
304            basin: basin.clone(),
305            stream: stream.clone(),
306        })?;
307
308        if meta.deleted_at.is_some() {
309            return Err(StreamDeletionPendingError.into());
310        }
311
312        let prior_doe_min_age = meta.config.delete_on_empty.min_age();
313
314        meta.config = OptionalStreamConfig::from(meta.config)
315            .reconfigure(reconfig)
316            .merge(basin_meta.config.default_stream_config);
317
318        txn.put(&meta_key, kv::stream_meta::ser_value(&meta))?;
319
320        let stream_id = StreamId::new(&basin, &stream);
321        if let Some(min_age) = meta.config.delete_on_empty.min_age()
322            && prior_doe_min_age.is_none()
323        {
324            txn.put(
325                kv::stream_doe_deadline::ser_key(
326                    kv::timestamp::TimestampSecs::after(doe_arm_delay(
327                        meta.config.retention_policy.age().unwrap_or_default(),
328                        min_age,
329                    )),
330                    stream_id,
331                ),
332                kv::stream_doe_deadline::ser_value(min_age),
333            )?;
334        }
335
336        txn.commit().await?;
337
338        if let Some(client) = self.streamer_client_if_active(&basin, &stream) {
339            client.advise_reconfig(meta.config.clone());
340        }
341
342        Ok(meta.config)
343    }
344
345    #[instrument(ret, err, skip(self))]
346    pub async fn delete_stream(
347        &self,
348        basin: BasinName,
349        stream: StreamName,
350    ) -> Result<(), DeleteStreamError> {
351        self.delete_stream_with_condition(basin, stream, TerminalTrimCondition::Always)
352            .await
353    }
354
355    pub(super) async fn delete_stream_with_condition(
356        &self,
357        basin: BasinName,
358        stream: StreamName,
359        condition: TerminalTrimCondition,
360    ) -> Result<(), DeleteStreamError> {
361        let outcome = match self.streamer_client_guarded(&basin, &stream).await {
362            Ok(client) => client.terminal_trim(condition).await?,
363            Err(StreamerError::Storage(e)) => {
364                return Err(DeleteStreamError::Storage(e));
365            }
366            Err(StreamerError::StreamNotFound(e)) => {
367                return Err(DeleteStreamError::StreamNotFound(e));
368            }
369            Err(StreamerError::StreamDeletionPending(_)) => TerminalTrimOutcome::DeletionPending,
370        };
371        match outcome {
372            TerminalTrimOutcome::DeletionPending => self.mark_stream_deleted(basin, stream).await,
373            TerminalTrimOutcome::Ineligible => Ok(()),
374        }
375    }
376
377    async fn mark_stream_deleted(
378        &self,
379        basin: BasinName,
380        stream: StreamName,
381    ) -> Result<(), DeleteStreamError> {
382        let txn = self.db.begin(IsolationLevel::SerializableSnapshot).await?;
383        let meta_key = kv::stream_meta::ser_key(&basin, &stream);
384        let mut meta = db_txn_get(&txn, &meta_key, kv::stream_meta::deser_value)
385            .await?
386            .ok_or_else(|| StreamNotFoundError {
387                basin,
388                stream: stream.clone(),
389            })?;
390        if meta.deleted_at.is_none() {
391            meta.deleted_at = Some(OffsetDateTime::now_utc());
392            txn.put(&meta_key, kv::stream_meta::ser_value(&meta))?;
393            txn.commit().await?;
394        }
395        Ok(())
396    }
397}
398
399fn creation_idempotency_key(req_token: &RequestToken, config: &OptionalStreamConfig) -> Bash {
400    Bash::length_prefixed(&[
401        req_token.as_bytes(),
402        &s2_api::v1::config::StreamConfig::to_opt(config.clone())
403            .as_ref()
404            .map(|v| serde_json::to_vec(v).expect("serializable"))
405            .unwrap_or_default(),
406    ])
407}