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