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