Skip to main content

ursula_runtime/
request.rs

1use std::sync::Arc;
2
3use bytes::Bytes;
4use serde::Deserialize;
5use serde::Serialize;
6use ursula_shard::BucketStreamId;
7use ursula_shard::ShardPlacement;
8use ursula_stream::ColdChunkRef;
9use ursula_stream::ExternalPayloadRef;
10use ursula_stream::ProducerRequest;
11use ursula_stream::StreamIntegritySnapshot;
12use ursula_stream::StreamReadPlan;
13use ursula_stream::StreamReadSegment;
14
15use crate::cold_index::ColdIndexPageCache;
16use crate::cold_index::ColdStoreColdIndexPageStore;
17use crate::cold_store::ColdStoreHandle;
18use crate::cold_store::DEFAULT_CONTENT_TYPE;
19use crate::engine::GroupEngineError;
20use crate::engine::in_memory::InMemoryGroupEngine;
21use crate::error::RuntimeError;
22
23#[derive(Debug, Clone, PartialEq, Eq)]
24pub struct CreateStreamRequest {
25    pub stream_id: BucketStreamId,
26    pub content_type: String,
27    pub content_type_explicit: bool,
28    pub initial_payload: Bytes,
29    pub close_after: bool,
30    pub stream_seq: Option<String>,
31    pub producer: Option<ProducerRequest>,
32    pub stream_ttl_seconds: Option<u64>,
33    pub stream_expires_at_ms: Option<u64>,
34    pub forked_from: Option<BucketStreamId>,
35    pub fork_offset: Option<u64>,
36    pub now_ms: u64,
37}
38
39#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
40pub struct CreateStreamExternalRequest {
41    pub stream_id: BucketStreamId,
42    pub content_type: String,
43    pub initial_payload: ExternalPayloadRef,
44    pub close_after: bool,
45    pub stream_seq: Option<String>,
46    pub producer: Option<ProducerRequest>,
47    pub stream_ttl_seconds: Option<u64>,
48    pub stream_expires_at_ms: Option<u64>,
49    pub forked_from: Option<BucketStreamId>,
50    pub fork_offset: Option<u64>,
51    pub now_ms: u64,
52}
53
54impl CreateStreamExternalRequest {
55    pub fn from_create_request(
56        request: CreateStreamRequest,
57        initial_payload: ExternalPayloadRef,
58    ) -> Self {
59        Self {
60            stream_id: request.stream_id,
61            content_type: request.content_type,
62            initial_payload,
63            close_after: request.close_after,
64            stream_seq: request.stream_seq,
65            producer: request.producer,
66            stream_ttl_seconds: request.stream_ttl_seconds,
67            stream_expires_at_ms: request.stream_expires_at_ms,
68            forked_from: request.forked_from,
69            fork_offset: request.fork_offset,
70            now_ms: request.now_ms,
71        }
72    }
73}
74
75impl CreateStreamRequest {
76    pub fn new(stream_id: BucketStreamId, content_type: impl Into<String>) -> Self {
77        Self {
78            stream_id,
79            content_type: content_type.into(),
80            content_type_explicit: true,
81            initial_payload: Bytes::new(),
82            close_after: false,
83            stream_seq: None,
84            producer: None,
85            stream_ttl_seconds: None,
86            stream_expires_at_ms: None,
87            forked_from: None,
88            fork_offset: None,
89            now_ms: 0,
90        }
91    }
92}
93
94#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
95pub struct CreateStreamResponse {
96    pub placement: ShardPlacement,
97    pub next_offset: u64,
98    pub closed: bool,
99    pub already_exists: bool,
100    pub group_commit_index: u64,
101}
102
103#[derive(Debug, Clone, PartialEq, Eq)]
104pub struct HeadStreamRequest {
105    pub stream_id: BucketStreamId,
106    pub now_ms: u64,
107}
108
109#[derive(Debug, Clone, PartialEq, Eq)]
110pub struct HeadStreamResponse {
111    pub placement: ShardPlacement,
112    pub content_type: String,
113    pub tail_offset: u64,
114    pub cold_hot_start_offset: u64,
115    pub closed: bool,
116    pub stream_ttl_seconds: Option<u64>,
117    pub stream_expires_at_ms: Option<u64>,
118    pub snapshot_offset: Option<u64>,
119    pub integrity: StreamIntegritySnapshot,
120}
121
122#[derive(Debug, Clone, PartialEq, Eq)]
123pub struct ReadStreamRequest {
124    pub stream_id: BucketStreamId,
125    pub offset: u64,
126    pub max_len: usize,
127    pub now_ms: u64,
128}
129
130#[derive(Debug, Clone, PartialEq, Eq)]
131pub struct ReadStreamResponse {
132    pub placement: ShardPlacement,
133    pub offset: u64,
134    pub next_offset: u64,
135    pub content_type: String,
136    pub payload: Vec<u8>,
137    pub up_to_date: bool,
138    pub closed: bool,
139}
140
141pub enum GroupReadStreamBody {
142    Materialized(Vec<u8>),
143    Planned {
144        stream_id: BucketStreamId,
145        plan: StreamReadPlan,
146        cold_store: Option<ColdStoreHandle>,
147        cold_index_cache: Option<Arc<ColdIndexPageCache<ColdStoreColdIndexPageStore>>>,
148    },
149    #[cfg(test)]
150    Blocking {
151        entered: Arc<crate::rt::sync::Notify>,
152        materialized: Arc<crate::rt::sync::Notify>,
153        release: Arc<crate::rt::sync::Notify>,
154        payload: Vec<u8>,
155    },
156}
157
158pub struct GroupReadStreamParts {
159    pub placement: ShardPlacement,
160    pub offset: u64,
161    pub next_offset: u64,
162    pub content_type: String,
163    pub up_to_date: bool,
164    pub closed: bool,
165    pub body: GroupReadStreamBody,
166}
167
168impl GroupReadStreamParts {
169    pub fn from_response(response: ReadStreamResponse) -> Self {
170        Self {
171            placement: response.placement,
172            offset: response.offset,
173            next_offset: response.next_offset,
174            content_type: response.content_type,
175            up_to_date: response.up_to_date,
176            closed: response.closed,
177            body: GroupReadStreamBody::Materialized(response.payload),
178        }
179    }
180
181    pub fn from_plan(
182        placement: ShardPlacement,
183        stream_id: BucketStreamId,
184        plan: StreamReadPlan,
185        cold_store: Option<ColdStoreHandle>,
186        cold_index_cache: Option<Arc<ColdIndexPageCache<ColdStoreColdIndexPageStore>>>,
187    ) -> Self {
188        Self {
189            placement,
190            offset: plan.offset,
191            next_offset: plan.next_offset,
192            content_type: plan.content_type.clone(),
193            up_to_date: plan.up_to_date,
194            closed: plan.closed,
195            body: GroupReadStreamBody::Planned {
196                stream_id,
197                plan,
198                cold_store,
199                cold_index_cache,
200            },
201        }
202    }
203
204    pub async fn into_response(self) -> Result<ReadStreamResponse, GroupEngineError> {
205        let payload = match &self.body {
206            GroupReadStreamBody::Materialized(payload) => payload.clone(),
207            GroupReadStreamBody::Planned {
208                stream_id,
209                plan,
210                cold_store,
211                cold_index_cache,
212            } => {
213                InMemoryGroupEngine::read_payload_from_plan(
214                    cold_store.as_ref(),
215                    cold_index_cache.as_ref(),
216                    stream_id,
217                    plan,
218                )
219                .await?
220            }
221            #[cfg(test)]
222            GroupReadStreamBody::Blocking {
223                entered,
224                materialized,
225                release,
226                payload,
227            } => {
228                entered.notify_one();
229                materialized.notify_one();
230                release.notified().await;
231                payload.clone()
232            }
233        };
234        Ok(ReadStreamResponse {
235            placement: self.placement,
236            offset: self.offset,
237            next_offset: self.next_offset,
238            content_type: self.content_type,
239            payload,
240            up_to_date: self.up_to_date,
241            closed: self.closed,
242        })
243    }
244
245    pub fn payload_is_empty(&self) -> bool {
246        match &self.body {
247            GroupReadStreamBody::Materialized(payload) => payload.is_empty(),
248            GroupReadStreamBody::Planned { plan, .. } => {
249                plan.segments.iter().all(|segment| match segment {
250                    StreamReadSegment::Hot(payload) => payload.is_empty(),
251                    StreamReadSegment::ColdIndex(segment) => segment.len == 0,
252                    StreamReadSegment::Object(segment) => segment.len == 0,
253                })
254            }
255            #[cfg(test)]
256            GroupReadStreamBody::Blocking { payload, .. } => payload.is_empty(),
257        }
258    }
259}
260
261#[derive(Debug, Clone, PartialEq, Eq)]
262pub struct PublishSnapshotRequest {
263    pub stream_id: BucketStreamId,
264    pub snapshot_offset: u64,
265    pub content_type: String,
266    pub payload: Bytes,
267    pub now_ms: u64,
268}
269
270#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
271pub struct PublishSnapshotResponse {
272    pub placement: ShardPlacement,
273    pub snapshot_offset: u64,
274    pub group_commit_index: u64,
275}
276
277#[derive(Debug, Clone, PartialEq, Eq)]
278pub struct ReadSnapshotRequest {
279    pub stream_id: BucketStreamId,
280    pub snapshot_offset: Option<u64>,
281    pub now_ms: u64,
282}
283
284#[derive(Debug, Clone, PartialEq, Eq)]
285pub struct ReadSnapshotResponse {
286    pub placement: ShardPlacement,
287    pub snapshot_offset: u64,
288    pub next_offset: u64,
289    pub content_type: String,
290    pub payload: Vec<u8>,
291    pub up_to_date: bool,
292}
293
294#[derive(Debug, Clone, PartialEq, Eq)]
295pub struct DeleteSnapshotRequest {
296    pub stream_id: BucketStreamId,
297    pub snapshot_offset: u64,
298    pub now_ms: u64,
299}
300
301#[derive(Debug, Clone, PartialEq, Eq)]
302pub struct BootstrapStreamRequest {
303    pub stream_id: BucketStreamId,
304    pub now_ms: u64,
305}
306
307#[derive(Debug, Clone, PartialEq, Eq)]
308pub struct BootstrapUpdate {
309    pub start_offset: u64,
310    pub next_offset: u64,
311    pub content_type: String,
312    pub payload: Vec<u8>,
313}
314
315#[derive(Debug, Clone, PartialEq, Eq)]
316pub struct BootstrapStreamResponse {
317    pub placement: ShardPlacement,
318    pub snapshot_offset: Option<u64>,
319    pub snapshot_content_type: String,
320    pub snapshot_payload: Vec<u8>,
321    pub updates: Vec<BootstrapUpdate>,
322    pub next_offset: u64,
323    pub up_to_date: bool,
324    pub closed: bool,
325}
326
327#[derive(Debug, Clone, PartialEq, Eq)]
328pub struct CloseStreamRequest {
329    pub stream_id: BucketStreamId,
330    pub stream_seq: Option<String>,
331    pub producer: Option<ProducerRequest>,
332    pub now_ms: u64,
333}
334
335#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
336pub struct CloseStreamResponse {
337    pub placement: ShardPlacement,
338    pub next_offset: u64,
339    pub group_commit_index: u64,
340    pub deduplicated: bool,
341}
342
343#[derive(Debug, Clone, PartialEq, Eq)]
344pub struct DeleteStreamRequest {
345    pub stream_id: BucketStreamId,
346}
347
348#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
349pub struct DeleteStreamResponse {
350    pub placement: ShardPlacement,
351    pub group_commit_index: u64,
352    pub hard_deleted: bool,
353    pub parent_to_release: Option<BucketStreamId>,
354}
355
356#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
357pub struct AckColdGcResponse {
358    pub placement: ShardPlacement,
359    pub removed: u64,
360    pub group_commit_index: u64,
361}
362
363#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
364pub struct ForkRefResponse {
365    pub placement: ShardPlacement,
366    pub fork_ref_count: u64,
367    pub hard_deleted: bool,
368    pub parent_to_release: Option<BucketStreamId>,
369    pub group_commit_index: u64,
370}
371
372#[derive(Debug, Clone, PartialEq, Eq)]
373pub struct FlushColdRequest {
374    pub stream_id: BucketStreamId,
375    pub chunk: ColdChunkRef,
376}
377
378#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
379pub struct FlushColdResponse {
380    pub placement: ShardPlacement,
381    pub hot_start_offset: u64,
382    pub group_commit_index: u64,
383}
384
385#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
386pub struct TouchStreamAccessResponse {
387    pub placement: ShardPlacement,
388    pub changed: bool,
389    pub expired: bool,
390    pub group_commit_index: u64,
391}
392
393#[derive(Debug, Clone, PartialEq, Eq)]
394pub struct PlanColdFlushRequest {
395    pub stream_id: BucketStreamId,
396    pub min_hot_bytes: usize,
397    pub max_flush_bytes: usize,
398}
399
400#[derive(Debug, Clone, PartialEq, Eq)]
401pub struct PlanGroupColdFlushRequest {
402    pub min_hot_bytes: usize,
403    pub max_flush_bytes: usize,
404}
405
406#[derive(Debug, Clone, PartialEq, Eq)]
407pub struct ColdHotBacklog {
408    pub stream_id: BucketStreamId,
409    pub stream_hot_bytes: u64,
410    pub group_hot_bytes: u64,
411}
412
413#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
414pub struct ColdWriteAdmission {
415    pub max_hot_bytes_per_group: Option<u64>,
416}
417
418impl ColdWriteAdmission {
419    pub(crate) fn is_enabled(self) -> bool {
420        self.max_hot_bytes_per_group.is_some()
421    }
422}
423
424#[derive(Debug, Clone, PartialEq, Eq)]
425pub struct AppendRequest {
426    pub stream_id: BucketStreamId,
427    pub content_type: String,
428    pub payload: Bytes,
429    pub close_after: bool,
430    pub stream_seq: Option<String>,
431    pub producer: Option<ProducerRequest>,
432    pub now_ms: u64,
433}
434
435#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
436pub struct AppendExternalRequest {
437    pub stream_id: BucketStreamId,
438    pub content_type: String,
439    pub payload: ExternalPayloadRef,
440    pub close_after: bool,
441    pub stream_seq: Option<String>,
442    pub producer: Option<ProducerRequest>,
443    pub now_ms: u64,
444}
445
446impl AppendExternalRequest {
447    pub fn from_append_request(request: AppendRequest, payload: ExternalPayloadRef) -> Self {
448        Self {
449            stream_id: request.stream_id,
450            content_type: request.content_type,
451            payload,
452            close_after: request.close_after,
453            stream_seq: request.stream_seq,
454            producer: request.producer,
455            now_ms: request.now_ms,
456        }
457    }
458}
459
460impl AppendRequest {
461    pub fn new(stream_id: BucketStreamId, payload_len: u64) -> Self {
462        Self {
463            stream_id,
464            content_type: DEFAULT_CONTENT_TYPE.to_owned(),
465            payload: Bytes::from(vec![
466                0;
467                usize::try_from(payload_len)
468                    .expect("payload_len fits usize")
469            ]),
470            close_after: false,
471            stream_seq: None,
472            producer: None,
473            now_ms: 0,
474        }
475    }
476
477    pub fn from_bytes(stream_id: BucketStreamId, payload: impl Into<Bytes>) -> Self {
478        Self {
479            stream_id,
480            content_type: DEFAULT_CONTENT_TYPE.to_owned(),
481            payload: payload.into(),
482            close_after: false,
483            stream_seq: None,
484            producer: None,
485            now_ms: 0,
486        }
487    }
488
489    pub fn payload_len(&self) -> u64 {
490        u64::try_from(self.payload.len()).expect("payload len fits u64")
491    }
492}
493
494#[derive(Debug, Clone, PartialEq, Eq)]
495pub struct AppendBatchRequest {
496    pub stream_id: BucketStreamId,
497    pub content_type: String,
498    pub payloads: Vec<Bytes>,
499    pub producer: Option<ProducerRequest>,
500    pub now_ms: u64,
501}
502
503impl AppendBatchRequest {
504    pub fn new<P>(stream_id: BucketStreamId, payloads: Vec<P>) -> Self
505    where P: Into<Bytes> {
506        Self {
507            stream_id,
508            content_type: DEFAULT_CONTENT_TYPE.to_owned(),
509            payloads: payloads.into_iter().map(Into::into).collect(),
510            producer: None,
511            now_ms: 0,
512        }
513    }
514}
515
516#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
517pub struct AppendResponse {
518    pub placement: ShardPlacement,
519    pub start_offset: u64,
520    pub next_offset: u64,
521    pub stream_append_count: u64,
522    pub group_commit_index: u64,
523    pub closed: bool,
524    pub deduplicated: bool,
525    pub producer: Option<ProducerRequest>,
526}
527
528#[derive(Debug, Clone, PartialEq, Eq)]
529pub struct AppendBatchResponse {
530    pub placement: ShardPlacement,
531    pub items: Vec<Result<AppendResponse, RuntimeError>>,
532}
533
534#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
535pub struct StreamAppendCount {
536    pub stream_id: BucketStreamId,
537    pub append_count: u64,
538}