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}