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}