1use std::collections::HashMap;
2use std::collections::HashSet;
3use std::sync::Arc;
4
5use bytes::Bytes;
6use ursula_shard::BucketStreamId;
7use ursula_shard::CoreId;
8use ursula_shard::RaftGroupId;
9use ursula_shard::ShardId;
10use ursula_shard::ShardPlacement;
11use ursula_stream::AppendStreamInput;
12use ursula_stream::ProducerRequest;
13use ursula_stream::StreamCommand;
14use ursula_stream::StreamErrorCode;
15use ursula_stream::StreamMessageRecord;
16use ursula_stream::StreamReadPlan;
17use ursula_stream::StreamReadSegment;
18use ursula_stream::StreamResponse;
19use ursula_stream::StreamSnapshot;
20use ursula_stream::StreamStateMachine;
21
22use super::GroupAckColdGcFuture;
23use super::GroupAppendBatchFuture;
24use super::GroupAppendBatchResponse;
25use super::GroupAppendFuture;
26use super::GroupBootstrapStreamFuture;
27use super::GroupCloseStreamFuture;
28use super::GroupColdHotBacklogFuture;
29use super::GroupCreateStreamFuture;
30use super::GroupDeleteSnapshotFuture;
31use super::GroupDeleteStreamFuture;
32use super::GroupEngine;
33use super::GroupEngineCreateFuture;
34use super::GroupEngineError;
35use super::GroupEngineFactory;
36use super::GroupEngineMetrics;
37use super::GroupFlushColdFuture;
38use super::GroupForkRefFuture;
39use super::GroupHeadStreamFuture;
40use super::GroupInstallSnapshotFuture;
41use super::GroupPlanColdFlushFuture;
42use super::GroupPlanColdGcFuture;
43use super::GroupPlanNextColdFlushBatchFuture;
44use super::GroupPlanNextColdFlushFuture;
45use super::GroupPublishSnapshotFuture;
46use super::GroupReadSnapshotFuture;
47use super::GroupReadStreamFuture;
48use super::GroupReadStreamPartsFuture;
49use super::GroupSnapshotFuture;
50use super::GroupTouchStreamAccessFuture;
51use super::GroupWriteResponse;
52use crate::cold_index::ColdIndexPageCache;
53use crate::cold_index::ColdStoreColdIndexPageStore;
54use crate::cold_index::rollback_cold_index_pages;
55use crate::cold_index::write_cold_chunk_index_pages_with_rollback;
56use crate::cold_index::write_external_segment_index_pages;
57use crate::cold_store::ColdStoreHandle;
58use crate::cold_store::DEFAULT_CONTENT_TYPE;
59use crate::command::GroupSnapshot;
60use crate::command::GroupWriteCommand;
61use crate::request::AckColdGcResponse;
62use crate::request::AppendBatchRequest;
63use crate::request::AppendExternalRequest;
64use crate::request::AppendRequest;
65use crate::request::AppendResponse;
66use crate::request::BootstrapStreamRequest;
67use crate::request::BootstrapStreamResponse;
68use crate::request::BootstrapUpdate;
69use crate::request::CloseStreamRequest;
70use crate::request::CloseStreamResponse;
71use crate::request::ColdHotBacklog;
72use crate::request::ColdWriteAdmission;
73use crate::request::CreateStreamExternalRequest;
74use crate::request::CreateStreamRequest;
75use crate::request::CreateStreamResponse;
76use crate::request::DeleteSnapshotRequest;
77use crate::request::DeleteStreamRequest;
78use crate::request::DeleteStreamResponse;
79use crate::request::FlushColdRequest;
80use crate::request::FlushColdResponse;
81use crate::request::ForkRefResponse;
82use crate::request::GroupReadStreamParts;
83use crate::request::HeadStreamRequest;
84use crate::request::HeadStreamResponse;
85use crate::request::PlanColdFlushRequest;
86use crate::request::PlanGroupColdFlushRequest;
87use crate::request::PublishSnapshotRequest;
88use crate::request::PublishSnapshotResponse;
89use crate::request::ReadSnapshotRequest;
90use crate::request::ReadSnapshotResponse;
91use crate::request::ReadStreamRequest;
92use crate::request::StreamAppendCount;
93use crate::request::TouchStreamAccessResponse;
94
95pub(crate) struct AppendPayloadInput<'a> {
96 stream_id: BucketStreamId,
97 content_type: Option<&'a str>,
98 payload: &'a [u8],
99 close_after: bool,
100 stream_seq: Option<String>,
101 producer: Option<ProducerRequest>,
102 now_ms: u64,
103}
104
105#[derive(Debug, Clone, Default)]
106pub struct InMemoryGroupEngine {
107 pub(crate) commit_index: u64,
108 pub(crate) state_machine: StreamStateMachine,
109 pub(crate) stream_append_counts: HashMap<BucketStreamId, u64>,
110 pub(crate) cold_store: Option<ColdStoreHandle>,
111 pub(crate) cold_index_cache: Option<Arc<ColdIndexPageCache<ColdStoreColdIndexPageStore>>>,
112}
113
114impl InMemoryGroupEngine {
115 pub fn with_cold_store(cold_store: ColdStoreHandle) -> Self {
116 let mut engine = Self::default();
117 engine.set_cold_store(Some(cold_store));
118 engine
119 }
120
121 pub fn cold_store(&self) -> Option<ColdStoreHandle> {
122 self.cold_store.clone()
123 }
124
125 pub(crate) fn set_cold_store(&mut self, cold_store: Option<ColdStoreHandle>) {
126 self.cold_index_cache = cold_store.as_ref().map(|cold_store| {
127 Arc::new(ColdIndexPageCache::new(
128 Arc::new(ColdStoreColdIndexPageStore::new(cold_store.clone())),
129 1024,
130 ))
131 });
132 self.cold_store = cold_store;
133 }
134
135 pub fn apply_committed_write(
136 &mut self,
137 command: GroupWriteCommand,
138 placement: ShardPlacement,
139 ) -> Result<GroupWriteResponse, GroupEngineError> {
140 match command {
141 GroupWriteCommand::CreateStream {
142 stream_id,
143 content_type,
144 initial_payload,
145 close_after,
146 stream_seq,
147 producer,
148 stream_ttl_seconds,
149 stream_expires_at_ms,
150 forked_from,
151 fork_offset,
152 now_ms,
153 } => {
154 ensure_bucket_exists(&mut self.state_machine, &stream_id)?;
155 let response = self.state_machine.apply(StreamCommand::CreateStream {
156 stream_id,
157 content_type,
158 initial_payload: initial_payload.to_vec(),
159 close_after,
160 stream_seq,
161 producer,
162 stream_ttl_seconds,
163 stream_expires_at_ms,
164 forked_from,
165 fork_offset,
166 now_ms,
167 });
168 match response {
169 StreamResponse::Created {
170 next_offset,
171 closed,
172 ..
173 } => {
174 self.commit_index += 1;
175 Ok(GroupWriteResponse::CreateStream(CreateStreamResponse {
176 placement,
177 next_offset,
178 closed,
179 already_exists: false,
180 group_commit_index: self.commit_index,
181 }))
182 }
183 StreamResponse::AlreadyExists {
184 next_offset,
185 closed,
186 ..
187 } => Ok(GroupWriteResponse::CreateStream(CreateStreamResponse {
188 placement,
189 next_offset,
190 closed,
191 already_exists: true,
192 group_commit_index: self.commit_index,
193 })),
194 StreamResponse::Error {
195 code,
196 message,
197 next_offset,
198 context,
199 } => Err(GroupEngineError::stream_with_context(
200 code,
201 message,
202 next_offset,
203 context,
204 )),
205 other => Err(GroupEngineError::new(format!(
206 "unexpected create stream response: {other:?}"
207 ))),
208 }
209 }
210 GroupWriteCommand::CreateExternal {
211 stream_id,
212 content_type,
213 initial_payload,
214 close_after,
215 stream_seq,
216 producer,
217 stream_ttl_seconds,
218 stream_expires_at_ms,
219 forked_from,
220 fork_offset,
221 now_ms,
222 } => {
223 ensure_bucket_exists(&mut self.state_machine, &stream_id)?;
224 let response = self.state_machine.apply(StreamCommand::CreateExternal {
225 stream_id,
226 content_type,
227 initial_payload,
228 close_after,
229 stream_seq,
230 producer,
231 stream_ttl_seconds,
232 stream_expires_at_ms,
233 forked_from,
234 fork_offset,
235 now_ms,
236 });
237 match response {
238 StreamResponse::Created {
239 next_offset,
240 closed,
241 ..
242 } => {
243 self.commit_index += 1;
244 Ok(GroupWriteResponse::CreateStream(CreateStreamResponse {
245 placement,
246 next_offset,
247 closed,
248 already_exists: false,
249 group_commit_index: self.commit_index,
250 }))
251 }
252 StreamResponse::AlreadyExists {
253 next_offset,
254 closed,
255 ..
256 } => Ok(GroupWriteResponse::CreateStream(CreateStreamResponse {
257 placement,
258 next_offset,
259 closed,
260 already_exists: true,
261 group_commit_index: self.commit_index,
262 })),
263 StreamResponse::Error {
264 code,
265 message,
266 next_offset,
267 context,
268 } => Err(GroupEngineError::stream_with_context(
269 code,
270 message,
271 next_offset,
272 context,
273 )),
274 other => Err(GroupEngineError::new(format!(
275 "unexpected create external stream response: {other:?}"
276 ))),
277 }
278 }
279 GroupWriteCommand::Append {
280 stream_id,
281 content_type,
282 payload,
283 close_after,
284 stream_seq,
285 producer,
286 now_ms,
287 } => self
288 .append_payload(
289 AppendPayloadInput {
290 stream_id,
291 content_type: Some(&content_type),
292 payload: &payload,
293 close_after,
294 stream_seq,
295 producer,
296 now_ms,
297 },
298 placement,
299 )
300 .map(GroupWriteResponse::Append),
301 GroupWriteCommand::AppendExternal {
302 stream_id,
303 content_type,
304 payload,
305 close_after,
306 stream_seq,
307 producer,
308 now_ms,
309 } => {
310 let response = self.state_machine.apply(StreamCommand::AppendExternal {
311 stream_id: stream_id.clone(),
312 content_type: Some(content_type),
313 payload,
314 close_after,
315 stream_seq,
316 producer,
317 now_ms,
318 });
319 match response {
320 StreamResponse::Appended {
321 offset,
322 next_offset,
323 closed,
324 deduplicated,
325 producer,
326 ..
327 } => {
328 let stream_append_count =
329 self.stream_append_counts.entry(stream_id).or_insert(0);
330 if !deduplicated {
331 self.commit_index += 1;
332 *stream_append_count += 1;
333 }
334 Ok(GroupWriteResponse::Append(AppendResponse {
335 placement,
336 start_offset: offset,
337 next_offset,
338 stream_append_count: *stream_append_count,
339 group_commit_index: self.commit_index,
340 closed,
341 deduplicated,
342 producer,
343 }))
344 }
345 StreamResponse::Error {
346 code,
347 message,
348 next_offset,
349 context,
350 } => Err(GroupEngineError::stream_with_context(
351 code,
352 message,
353 next_offset,
354 context,
355 )),
356 other => Err(GroupEngineError::new(format!(
357 "unexpected append external response: {other:?}"
358 ))),
359 }
360 }
361 GroupWriteCommand::AppendBatch {
362 stream_id,
363 content_type,
364 payloads,
365 producer,
366 now_ms,
367 } => {
368 if producer.is_some() {
369 let payload_refs = payloads.iter().map(Bytes::as_ref).collect::<Vec<_>>();
370 let batch = self
371 .state_machine
372 .append_batch_borrowed(
373 stream_id.clone(),
374 Some(&content_type),
375 &payload_refs,
376 producer,
377 now_ms,
378 )
379 .map_err(stream_response_error)?;
380 let old_commit_index = self.commit_index;
381 let old_append_count = *self.stream_append_counts.get(&stream_id).unwrap_or(&0);
382 if !batch.deduplicated {
383 let count = u64::try_from(batch.items.len()).expect("item count fits u64");
384 self.commit_index += count;
385 *self.stream_append_counts.entry(stream_id).or_insert(0) += count;
386 }
387 let items = batch
388 .items
389 .into_iter()
390 .enumerate()
391 .map(|(index, item)| {
392 let item_index = u64::try_from(index + 1).expect("item index fits u64");
393 Ok(AppendResponse {
394 placement,
395 start_offset: item.offset,
396 next_offset: item.next_offset,
397 stream_append_count: if item.deduplicated {
398 old_append_count
399 } else {
400 old_append_count + item_index
401 },
402 group_commit_index: if item.deduplicated {
403 old_commit_index
404 } else {
405 old_commit_index + item_index
406 },
407 closed: item.closed,
408 deduplicated: item.deduplicated,
409 producer: None,
410 })
411 })
412 .collect();
413 return Ok(GroupWriteResponse::AppendBatch(GroupAppendBatchResponse {
414 placement,
415 items,
416 }));
417 }
418
419 let mut items = Vec::with_capacity(payloads.len());
420 for payload in payloads {
421 if payload.is_empty() {
422 items.push(Err(GroupEngineError::stream(
423 StreamErrorCode::EmptyAppend,
424 "append payload must be non-empty",
425 )));
426 continue;
427 }
428 items.push(self.append_payload(
429 AppendPayloadInput {
430 stream_id: stream_id.clone(),
431 content_type: Some(&content_type),
432 payload: &payload,
433 close_after: false,
434 stream_seq: None,
435 producer: None,
436 now_ms,
437 },
438 placement,
439 ));
440 }
441 Ok(GroupWriteResponse::AppendBatch(GroupAppendBatchResponse {
442 placement,
443 items,
444 }))
445 }
446 GroupWriteCommand::PublishSnapshot {
447 stream_id,
448 snapshot_offset,
449 content_type,
450 payload,
451 now_ms,
452 } => {
453 let response = self.state_machine.apply(StreamCommand::PublishSnapshot {
454 stream_id,
455 snapshot_offset,
456 content_type,
457 payload: payload.to_vec(),
458 now_ms,
459 });
460 match response {
461 StreamResponse::SnapshotPublished { snapshot_offset } => {
462 self.commit_index += 1;
463 Ok(GroupWriteResponse::PublishSnapshot(
464 PublishSnapshotResponse {
465 placement,
466 snapshot_offset,
467 group_commit_index: self.commit_index,
468 },
469 ))
470 }
471 StreamResponse::Error {
472 code,
473 message,
474 next_offset,
475 context,
476 } => Err(GroupEngineError::stream_with_context(
477 code,
478 message,
479 next_offset,
480 context,
481 )),
482 other => Err(GroupEngineError::new(format!(
483 "unexpected publish snapshot response: {other:?}"
484 ))),
485 }
486 }
487 GroupWriteCommand::TouchStreamAccess {
488 stream_id,
489 now_ms,
490 renew_ttl,
491 } => {
492 let response = self.state_machine.apply(StreamCommand::TouchStreamAccess {
493 stream_id,
494 now_ms,
495 renew_ttl,
496 });
497 match response {
498 StreamResponse::Accessed { changed, expired } => {
499 if changed || expired {
500 self.commit_index += 1;
501 }
502 Ok(GroupWriteResponse::TouchStreamAccess(
503 TouchStreamAccessResponse {
504 placement,
505 changed,
506 expired,
507 group_commit_index: self.commit_index,
508 },
509 ))
510 }
511 StreamResponse::Error {
512 code,
513 message,
514 next_offset,
515 context,
516 } => Err(GroupEngineError::stream_with_context(
517 code,
518 message,
519 next_offset,
520 context,
521 )),
522 other => Err(GroupEngineError::new(format!(
523 "unexpected touch stream access response: {other:?}"
524 ))),
525 }
526 }
527 GroupWriteCommand::AddForkRef { stream_id, now_ms } => {
528 let response = self
529 .state_machine
530 .apply(StreamCommand::AddForkRef { stream_id, now_ms });
531 match response {
532 StreamResponse::ForkRefAdded { fork_ref_count } => {
533 self.commit_index += 1;
534 Ok(GroupWriteResponse::AddForkRef(ForkRefResponse {
535 placement,
536 fork_ref_count,
537 hard_deleted: false,
538 parent_to_release: None,
539 group_commit_index: self.commit_index,
540 }))
541 }
542 StreamResponse::Error {
543 code,
544 message,
545 next_offset,
546 context,
547 } => Err(GroupEngineError::stream_with_context(
548 code,
549 message,
550 next_offset,
551 context,
552 )),
553 other => Err(GroupEngineError::new(format!(
554 "unexpected add fork ref response: {other:?}"
555 ))),
556 }
557 }
558 GroupWriteCommand::ReleaseForkRef { stream_id } => {
559 let response = self
560 .state_machine
561 .apply(StreamCommand::ReleaseForkRef { stream_id });
562 match response {
563 StreamResponse::ForkRefReleased {
564 hard_deleted,
565 fork_ref_count,
566 parent_to_release,
567 } => {
568 self.commit_index += 1;
569 Ok(GroupWriteResponse::ReleaseForkRef(ForkRefResponse {
570 placement,
571 fork_ref_count,
572 hard_deleted,
573 parent_to_release,
574 group_commit_index: self.commit_index,
575 }))
576 }
577 StreamResponse::Error {
578 code,
579 message,
580 next_offset,
581 context,
582 } => Err(GroupEngineError::stream_with_context(
583 code,
584 message,
585 next_offset,
586 context,
587 )),
588 other => Err(GroupEngineError::new(format!(
589 "unexpected release fork ref response: {other:?}"
590 ))),
591 }
592 }
593 GroupWriteCommand::FlushCold { stream_id, chunk } => {
594 let response = self
595 .state_machine
596 .apply(StreamCommand::FlushCold { stream_id, chunk });
597 match response {
598 StreamResponse::ColdFlushed { hot_start_offset } => {
599 self.commit_index += 1;
600 Ok(GroupWriteResponse::FlushCold(FlushColdResponse {
601 placement,
602 hot_start_offset,
603 group_commit_index: self.commit_index,
604 }))
605 }
606 StreamResponse::Error {
607 code,
608 message,
609 next_offset,
610 context,
611 } => Err(GroupEngineError::stream_with_context(
612 code,
613 message,
614 next_offset,
615 context,
616 )),
617 other => Err(GroupEngineError::new(format!(
618 "unexpected flush cold response: {other:?}"
619 ))),
620 }
621 }
622 GroupWriteCommand::CloseStream {
623 stream_id,
624 stream_seq,
625 producer,
626 now_ms,
627 } => {
628 let response = self.state_machine.apply(StreamCommand::Close {
629 stream_id,
630 stream_seq,
631 producer,
632 now_ms,
633 });
634 match response {
635 StreamResponse::Closed {
636 next_offset,
637 deduplicated,
638 ..
639 } => {
640 if !deduplicated {
641 self.commit_index += 1;
642 }
643 Ok(GroupWriteResponse::CloseStream(CloseStreamResponse {
644 placement,
645 next_offset,
646 group_commit_index: self.commit_index,
647 deduplicated,
648 }))
649 }
650 StreamResponse::Error {
651 code,
652 message,
653 next_offset,
654 context,
655 } => Err(GroupEngineError::stream_with_context(
656 code,
657 message,
658 next_offset,
659 context,
660 )),
661 other => Err(GroupEngineError::new(format!(
662 "unexpected close stream response: {other:?}"
663 ))),
664 }
665 }
666 GroupWriteCommand::DeleteStream { stream_id } => {
667 let response = self.state_machine.apply(StreamCommand::DeleteStream {
668 stream_id: stream_id.clone(),
669 });
670 match response {
671 StreamResponse::Deleted {
672 hard_deleted,
673 parent_to_release,
674 } => {
675 self.commit_index += 1;
676 if hard_deleted {
677 self.stream_append_counts.remove(&stream_id);
681 }
682 Ok(GroupWriteResponse::DeleteStream(DeleteStreamResponse {
683 placement,
684 group_commit_index: self.commit_index,
685 hard_deleted,
686 parent_to_release,
687 }))
688 }
689 StreamResponse::Error {
690 code,
691 message,
692 next_offset,
693 context,
694 } => Err(GroupEngineError::stream_with_context(
695 code,
696 message,
697 next_offset,
698 context,
699 )),
700 other => Err(GroupEngineError::new(format!(
701 "unexpected delete stream response: {other:?}"
702 ))),
703 }
704 }
705 GroupWriteCommand::AckColdGc { up_to_seq } => {
706 let response = self
707 .state_machine
708 .apply(StreamCommand::AckColdGc { up_to_seq });
709 match response {
710 StreamResponse::ColdGcAcked { removed } => {
711 self.commit_index += 1;
712 Ok(GroupWriteResponse::AckColdGc(AckColdGcResponse {
713 placement,
714 removed,
715 group_commit_index: self.commit_index,
716 }))
717 }
718 other => Err(GroupEngineError::new(format!(
719 "unexpected ack cold gc response: {other:?}"
720 ))),
721 }
722 }
723 GroupWriteCommand::Batch { commands } => Ok(GroupWriteResponse::Batch(
724 self.apply_committed_write_batch(commands, placement),
725 )),
726 }
727 }
728
729 pub(crate) fn cold_hot_backlog_for(
730 &self,
731 stream_id: BucketStreamId,
732 ) -> Result<ColdHotBacklog, GroupEngineError> {
733 let stream_hot_bytes = self.state_machine.hot_payload_len(&stream_id).unwrap_or(0);
734 Ok(ColdHotBacklog {
735 stream_id,
736 stream_hot_bytes,
737 group_hot_bytes: self.state_machine.total_hot_payload_bytes(),
738 })
739 }
740
741 pub fn check_cold_write_admission_bytes(
742 &self,
743 stream_id: &BucketStreamId,
744 admission: ColdWriteAdmission,
745 incoming_bytes: u64,
746 ) -> Result<(), GroupEngineError> {
747 let Some(limit) = admission.max_hot_bytes_per_group else {
748 return Ok(());
749 };
750 if incoming_bytes == 0 {
751 return Ok(());
752 }
753 let before = self.state_machine.total_hot_payload_bytes();
754 let after = before.saturating_add(incoming_bytes);
755 if after <= limit {
756 return Ok(());
757 }
758 Err(GroupEngineError::cold_backpressure(
759 stream_id.clone(),
760 before,
761 after,
762 limit,
763 ))
764 }
765
766 pub(crate) fn create_stream_with_admission_inner(
767 &mut self,
768 request: CreateStreamRequest,
769 placement: ShardPlacement,
770 admission: ColdWriteAdmission,
771 ) -> Result<CreateStreamResponse, GroupEngineError> {
772 let stream_id = request.stream_id.clone();
773 if admission.is_enabled() {
774 let mut preview = self.clone();
775 let preview_response = match preview
776 .apply_committed_write(GroupWriteCommand::from(request.clone()), placement)?
777 {
778 GroupWriteResponse::CreateStream(response) => response,
779 other => {
780 return Err(GroupEngineError::new(format!(
781 "unexpected create stream preview response: {other:?}"
782 )));
783 }
784 };
785 if !preview_response.already_exists {
786 self.check_cold_write_admission_bytes(
787 &stream_id,
788 admission,
789 u64::try_from(request.initial_payload.len()).expect("payload len fits u64"),
790 )?;
791 }
792 }
793 let response =
794 match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
795 GroupWriteResponse::CreateStream(response) => response,
796 other => {
797 return Err(GroupEngineError::new(format!(
798 "unexpected create stream write response: {other:?}"
799 )));
800 }
801 };
802 Ok(response)
803 }
804
805 pub(crate) fn append_with_admission_inner(
806 &mut self,
807 request: AppendRequest,
808 placement: ShardPlacement,
809 admission: ColdWriteAdmission,
810 ) -> Result<AppendResponse, GroupEngineError> {
811 let stream_id = request.stream_id.clone();
812 if admission.is_enabled() {
813 let mut preview = self.clone();
814 let preview_response = match preview
815 .apply_committed_write(GroupWriteCommand::from(request.clone()), placement)?
816 {
817 GroupWriteResponse::Append(response) => response,
818 other => {
819 return Err(GroupEngineError::new(format!(
820 "unexpected append preview response: {other:?}"
821 )));
822 }
823 };
824 if !preview_response.deduplicated {
825 self.check_cold_write_admission_bytes(
826 &stream_id,
827 admission,
828 u64::try_from(request.payload.len()).expect("payload len fits u64"),
829 )?;
830 }
831 }
832 let response =
833 match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
834 GroupWriteResponse::Append(response) => response,
835 other => {
836 return Err(GroupEngineError::new(format!(
837 "unexpected append write response: {other:?}"
838 )));
839 }
840 };
841 Ok(response)
842 }
843
844 pub(crate) fn append_batch_with_admission_inner(
845 &mut self,
846 request: AppendBatchRequest,
847 placement: ShardPlacement,
848 admission: ColdWriteAdmission,
849 ) -> Result<GroupAppendBatchResponse, GroupEngineError> {
850 let stream_id = request.stream_id.clone();
851 let incoming_bytes = request
852 .payloads
853 .iter()
854 .map(|payload| u64::try_from(payload.len()).expect("payload len fits u64"))
855 .sum();
856 if admission.is_enabled() {
857 let mut preview = self.clone();
858 let preview_response = match preview
859 .apply_committed_write(GroupWriteCommand::from(request.clone()), placement)?
860 {
861 GroupWriteResponse::AppendBatch(response) => response,
862 other => {
863 return Err(GroupEngineError::new(format!(
864 "unexpected append batch preview response: {other:?}"
865 )));
866 }
867 };
868 let mutates = preview_response
869 .items
870 .iter()
871 .any(|item| matches!(item, Ok(response) if !response.deduplicated));
872 if mutates {
873 self.check_cold_write_admission_bytes(&stream_id, admission, incoming_bytes)?;
874 }
875 }
876 let response =
877 match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
878 GroupWriteResponse::AppendBatch(response) => response,
879 other => {
880 return Err(GroupEngineError::new(format!(
881 "unexpected append batch write response: {other:?}"
882 )));
883 }
884 };
885 Ok(response)
886 }
887
888 pub fn access_requires_write(
889 &self,
890 stream_id: &BucketStreamId,
891 now_ms: u64,
892 renew_ttl: bool,
893 ) -> Result<bool, GroupEngineError> {
894 self.state_machine
895 .access_requires_write(stream_id, now_ms, renew_ttl)
896 .map_err(stream_response_error)
897 }
898
899 pub(crate) fn apply_access_command(
900 &mut self,
901 stream_id: BucketStreamId,
902 now_ms: u64,
903 renew_ttl: bool,
904 placement: ShardPlacement,
905 ) -> Result<TouchStreamAccessResponse, GroupEngineError> {
906 match self.apply_committed_write(
907 GroupWriteCommand::TouchStreamAccess {
908 stream_id,
909 now_ms,
910 renew_ttl,
911 },
912 placement,
913 )? {
914 GroupWriteResponse::TouchStreamAccess(response) => Ok(response),
915 other => Err(GroupEngineError::new(format!(
916 "unexpected touch stream access write response: {other:?}"
917 ))),
918 }
919 }
920
921 pub(crate) fn ensure_stream_access(
922 &mut self,
923 stream_id: &BucketStreamId,
924 now_ms: u64,
925 renew_ttl: bool,
926 placement: ShardPlacement,
927 ) -> Result<Option<TouchStreamAccessResponse>, GroupEngineError> {
928 if !self.access_requires_write(stream_id, now_ms, renew_ttl)? {
929 return Ok(None);
930 }
931 let response =
932 self.apply_access_command(stream_id.clone(), now_ms, renew_ttl, placement)?;
933 if response.expired {
934 return Err(GroupEngineError::stream(
935 StreamErrorCode::StreamNotFound,
936 format!("stream '{stream_id}' does not exist"),
937 ));
938 }
939 Ok(Some(response))
940 }
941
942 pub fn apply_committed_write_batch(
943 &mut self,
944 commands: Vec<GroupWriteCommand>,
945 placement: ShardPlacement,
946 ) -> Vec<Result<GroupWriteResponse, GroupEngineError>> {
947 commands
948 .into_iter()
949 .map(|command| self.apply_committed_write(command, placement))
950 .collect()
951 }
952
953 pub(crate) fn apply_replayed_write_command(
954 &mut self,
955 command: GroupWriteCommand,
956 ) -> Result<(), GroupEngineError> {
957 let placement = ShardPlacement {
958 core_id: CoreId(0),
959 shard_id: ShardId(0),
960 raft_group_id: RaftGroupId(0),
961 };
962 self.apply_committed_write(command, placement).map(|_| ())
963 }
964
965 pub(crate) fn apply_replayed_command(
966 &mut self,
967 command: StreamCommand,
968 ) -> Result<(), GroupEngineError> {
969 match command {
970 StreamCommand::CreateBucket { bucket_id } => {
971 match self
972 .state_machine
973 .apply(StreamCommand::CreateBucket { bucket_id })
974 {
975 StreamResponse::BucketCreated { .. } => {
976 self.commit_index += 1;
977 Ok(())
978 }
979 StreamResponse::BucketAlreadyExists { .. } => Ok(()),
980 StreamResponse::Error {
981 code,
982 message,
983 next_offset,
984 context,
985 } => Err(GroupEngineError::stream_with_context(
986 code,
987 message,
988 next_offset,
989 context,
990 )),
991 other => Err(GroupEngineError::new(format!(
992 "unexpected replay create bucket response: {other:?}"
993 ))),
994 }
995 }
996 StreamCommand::DeleteBucket { bucket_id } => {
997 match self
998 .state_machine
999 .apply(StreamCommand::DeleteBucket { bucket_id })
1000 {
1001 StreamResponse::BucketDeleted { .. } => {
1002 self.commit_index += 1;
1003 Ok(())
1004 }
1005 StreamResponse::Error {
1006 code,
1007 message,
1008 next_offset,
1009 context,
1010 } => Err(GroupEngineError::stream_with_context(
1011 code,
1012 message,
1013 next_offset,
1014 context,
1015 )),
1016 other => Err(GroupEngineError::new(format!(
1017 "unexpected replay delete bucket response: {other:?}"
1018 ))),
1019 }
1020 }
1021 StreamCommand::CreateStream {
1022 stream_id,
1023 content_type,
1024 initial_payload,
1025 close_after,
1026 stream_seq,
1027 producer,
1028 stream_ttl_seconds,
1029 stream_expires_at_ms,
1030 forked_from,
1031 fork_offset,
1032 now_ms,
1033 } => {
1034 ensure_bucket_exists(&mut self.state_machine, &stream_id)?;
1035 let response = self.state_machine.apply(StreamCommand::CreateStream {
1036 stream_id,
1037 content_type,
1038 initial_payload,
1039 close_after,
1040 stream_seq,
1041 producer,
1042 stream_ttl_seconds,
1043 stream_expires_at_ms,
1044 forked_from,
1045 fork_offset,
1046 now_ms,
1047 });
1048 match response {
1049 StreamResponse::Created { .. } => {
1050 self.commit_index += 1;
1051 Ok(())
1052 }
1053 StreamResponse::AlreadyExists { .. } => Ok(()),
1054 StreamResponse::Error {
1055 code,
1056 message,
1057 next_offset,
1058 context,
1059 } => Err(GroupEngineError::stream_with_context(
1060 code,
1061 message,
1062 next_offset,
1063 context,
1064 )),
1065 other => Err(GroupEngineError::new(format!(
1066 "unexpected replay create stream response: {other:?}"
1067 ))),
1068 }
1069 }
1070 StreamCommand::CreateExternal {
1071 stream_id,
1072 content_type,
1073 initial_payload,
1074 close_after,
1075 stream_seq,
1076 producer,
1077 stream_ttl_seconds,
1078 stream_expires_at_ms,
1079 forked_from,
1080 fork_offset,
1081 now_ms,
1082 } => {
1083 ensure_bucket_exists(&mut self.state_machine, &stream_id)?;
1084 let response = self.state_machine.apply(StreamCommand::CreateExternal {
1085 stream_id,
1086 content_type,
1087 initial_payload,
1088 close_after,
1089 stream_seq,
1090 producer,
1091 stream_ttl_seconds,
1092 stream_expires_at_ms,
1093 forked_from,
1094 fork_offset,
1095 now_ms,
1096 });
1097 match response {
1098 StreamResponse::Created { .. } => {
1099 self.commit_index += 1;
1100 Ok(())
1101 }
1102 StreamResponse::AlreadyExists { .. } => Ok(()),
1103 StreamResponse::Error {
1104 code,
1105 message,
1106 next_offset,
1107 context,
1108 } => Err(GroupEngineError::stream_with_context(
1109 code,
1110 message,
1111 next_offset,
1112 context,
1113 )),
1114 other => Err(GroupEngineError::new(format!(
1115 "unexpected replay external create stream response: {other:?}"
1116 ))),
1117 }
1118 }
1119 StreamCommand::Append {
1120 stream_id,
1121 content_type,
1122 payload,
1123 close_after,
1124 stream_seq,
1125 producer,
1126 now_ms,
1127 } => {
1128 let stream_count_key = stream_id.clone();
1129 let response = self.state_machine.apply(StreamCommand::Append {
1130 stream_id,
1131 content_type,
1132 payload,
1133 close_after,
1134 stream_seq,
1135 producer,
1136 now_ms,
1137 });
1138 match response {
1139 StreamResponse::Appended { deduplicated, .. } => {
1140 if !deduplicated {
1141 self.commit_index += 1;
1142 *self
1143 .stream_append_counts
1144 .entry(stream_count_key)
1145 .or_insert(0) += 1;
1146 }
1147 Ok(())
1148 }
1149 StreamResponse::Closed { deduplicated, .. } => {
1150 if !deduplicated {
1151 self.commit_index += 1;
1152 }
1153 Ok(())
1154 }
1155 StreamResponse::Error {
1156 code,
1157 message,
1158 next_offset,
1159 context,
1160 } => Err(GroupEngineError::stream_with_context(
1161 code,
1162 message,
1163 next_offset,
1164 context,
1165 )),
1166 other => Err(GroupEngineError::new(format!(
1167 "unexpected replay append response: {other:?}"
1168 ))),
1169 }
1170 }
1171 StreamCommand::AppendExternal {
1172 stream_id,
1173 content_type,
1174 payload,
1175 close_after,
1176 stream_seq,
1177 producer,
1178 now_ms,
1179 } => {
1180 let stream_count_key = stream_id.clone();
1181 let response = self.state_machine.apply(StreamCommand::AppendExternal {
1182 stream_id,
1183 content_type,
1184 payload,
1185 close_after,
1186 stream_seq,
1187 producer,
1188 now_ms,
1189 });
1190 match response {
1191 StreamResponse::Appended { deduplicated, .. } => {
1192 if !deduplicated {
1193 self.commit_index += 1;
1194 *self
1195 .stream_append_counts
1196 .entry(stream_count_key)
1197 .or_insert(0) += 1;
1198 }
1199 Ok(())
1200 }
1201 StreamResponse::Error {
1202 code,
1203 message,
1204 next_offset,
1205 context,
1206 } => Err(GroupEngineError::stream_with_context(
1207 code,
1208 message,
1209 next_offset,
1210 context,
1211 )),
1212 other => Err(GroupEngineError::new(format!(
1213 "unexpected replay external append response: {other:?}"
1214 ))),
1215 }
1216 }
1217 StreamCommand::AppendBatch {
1218 stream_id,
1219 content_type,
1220 payloads,
1221 producer,
1222 now_ms,
1223 } => {
1224 let stream_count_key = stream_id.clone();
1225 let payload_refs = payloads.iter().map(Vec::as_slice).collect::<Vec<_>>();
1226 let response = self
1227 .state_machine
1228 .append_batch_borrowed(
1229 stream_id,
1230 content_type.as_deref(),
1231 &payload_refs,
1232 producer,
1233 now_ms,
1234 )
1235 .map_err(stream_response_error)?;
1236 if !response.deduplicated {
1237 let count = u64::try_from(response.items.len()).expect("item count fits u64");
1238 self.commit_index += count;
1239 *self
1240 .stream_append_counts
1241 .entry(stream_count_key)
1242 .or_insert(0) += count;
1243 }
1244 Ok(())
1245 }
1246 StreamCommand::PublishSnapshot {
1247 stream_id,
1248 snapshot_offset,
1249 content_type,
1250 payload,
1251 now_ms,
1252 } => {
1253 let response = self.state_machine.apply(StreamCommand::PublishSnapshot {
1254 stream_id,
1255 snapshot_offset,
1256 content_type,
1257 payload,
1258 now_ms,
1259 });
1260 match response {
1261 StreamResponse::SnapshotPublished { .. } => {
1262 self.commit_index += 1;
1263 Ok(())
1264 }
1265 StreamResponse::Error {
1266 code,
1267 message,
1268 next_offset,
1269 context,
1270 } => Err(GroupEngineError::stream_with_context(
1271 code,
1272 message,
1273 next_offset,
1274 context,
1275 )),
1276 other => Err(GroupEngineError::new(format!(
1277 "unexpected replay publish snapshot response: {other:?}"
1278 ))),
1279 }
1280 }
1281 StreamCommand::TouchStreamAccess {
1282 stream_id,
1283 now_ms,
1284 renew_ttl,
1285 } => {
1286 let response = self.state_machine.apply(StreamCommand::TouchStreamAccess {
1287 stream_id,
1288 now_ms,
1289 renew_ttl,
1290 });
1291 match response {
1292 StreamResponse::Accessed { changed, expired } => {
1293 if changed || expired {
1294 self.commit_index += 1;
1295 }
1296 Ok(())
1297 }
1298 StreamResponse::Error {
1299 code,
1300 message,
1301 next_offset,
1302 context,
1303 } => Err(GroupEngineError::stream_with_context(
1304 code,
1305 message,
1306 next_offset,
1307 context,
1308 )),
1309 other => Err(GroupEngineError::new(format!(
1310 "unexpected replay touch stream access response: {other:?}"
1311 ))),
1312 }
1313 }
1314 StreamCommand::AddForkRef { stream_id, now_ms } => {
1315 let response = self
1316 .state_machine
1317 .apply(StreamCommand::AddForkRef { stream_id, now_ms });
1318 match response {
1319 StreamResponse::ForkRefAdded { .. } => {
1320 self.commit_index += 1;
1321 Ok(())
1322 }
1323 StreamResponse::Error {
1324 code,
1325 message,
1326 next_offset,
1327 context,
1328 } => Err(GroupEngineError::stream_with_context(
1329 code,
1330 message,
1331 next_offset,
1332 context,
1333 )),
1334 other => Err(GroupEngineError::new(format!(
1335 "unexpected replay add fork ref response: {other:?}"
1336 ))),
1337 }
1338 }
1339 StreamCommand::ReleaseForkRef { stream_id } => {
1340 let response = self
1341 .state_machine
1342 .apply(StreamCommand::ReleaseForkRef { stream_id });
1343 match response {
1344 StreamResponse::ForkRefReleased { .. } => {
1345 self.commit_index += 1;
1346 Ok(())
1347 }
1348 StreamResponse::Error {
1349 code,
1350 message,
1351 next_offset,
1352 context,
1353 } => Err(GroupEngineError::stream_with_context(
1354 code,
1355 message,
1356 next_offset,
1357 context,
1358 )),
1359 other => Err(GroupEngineError::new(format!(
1360 "unexpected replay release fork ref response: {other:?}"
1361 ))),
1362 }
1363 }
1364 StreamCommand::FlushCold { stream_id, chunk } => {
1365 let response = self
1366 .state_machine
1367 .apply(StreamCommand::FlushCold { stream_id, chunk });
1368 match response {
1369 StreamResponse::ColdFlushed { .. } => {
1370 self.commit_index += 1;
1371 Ok(())
1372 }
1373 StreamResponse::Error {
1374 code,
1375 message,
1376 next_offset,
1377 context,
1378 } => Err(GroupEngineError::stream_with_context(
1379 code,
1380 message,
1381 next_offset,
1382 context,
1383 )),
1384 other => Err(GroupEngineError::new(format!(
1385 "unexpected replay flush cold response: {other:?}"
1386 ))),
1387 }
1388 }
1389 StreamCommand::Close {
1390 stream_id,
1391 stream_seq,
1392 producer,
1393 now_ms,
1394 } => {
1395 let response = self.state_machine.apply(StreamCommand::Close {
1396 stream_id,
1397 stream_seq,
1398 producer,
1399 now_ms,
1400 });
1401 match response {
1402 StreamResponse::Closed { deduplicated, .. } => {
1403 if !deduplicated {
1404 self.commit_index += 1;
1405 }
1406 Ok(())
1407 }
1408 StreamResponse::Error {
1409 code,
1410 message,
1411 next_offset,
1412 context,
1413 } => Err(GroupEngineError::stream_with_context(
1414 code,
1415 message,
1416 next_offset,
1417 context,
1418 )),
1419 other => Err(GroupEngineError::new(format!(
1420 "unexpected replay close stream response: {other:?}"
1421 ))),
1422 }
1423 }
1424 StreamCommand::DeleteStream { stream_id } => {
1425 let response = self
1426 .state_machine
1427 .apply(StreamCommand::DeleteStream { stream_id });
1428 match response {
1429 StreamResponse::Deleted { .. } => {
1430 self.commit_index += 1;
1431 Ok(())
1432 }
1433 StreamResponse::Error {
1434 code,
1435 message,
1436 next_offset,
1437 context,
1438 } => Err(GroupEngineError::stream_with_context(
1439 code,
1440 message,
1441 next_offset,
1442 context,
1443 )),
1444 other => Err(GroupEngineError::new(format!(
1445 "unexpected replay delete stream response: {other:?}"
1446 ))),
1447 }
1448 }
1449 StreamCommand::AckColdGc { up_to_seq } => {
1450 match self
1451 .state_machine
1452 .apply(StreamCommand::AckColdGc { up_to_seq })
1453 {
1454 StreamResponse::ColdGcAcked { .. } => {
1455 self.commit_index += 1;
1456 Ok(())
1457 }
1458 other => Err(GroupEngineError::new(format!(
1459 "unexpected replay ack cold gc response: {other:?}"
1460 ))),
1461 }
1462 }
1463 }
1464 }
1465
1466 pub(crate) fn append_payload(
1467 &mut self,
1468 input: AppendPayloadInput<'_>,
1469 placement: ShardPlacement,
1470 ) -> Result<AppendResponse, GroupEngineError> {
1471 let AppendPayloadInput {
1472 stream_id,
1473 content_type,
1474 payload,
1475 close_after,
1476 stream_seq,
1477 producer,
1478 now_ms,
1479 } = input;
1480 let stream_count_key = stream_id.clone();
1481 let response = self.state_machine.append_borrowed(AppendStreamInput {
1482 stream_id,
1483 content_type,
1484 payload,
1485 close_after,
1486 stream_seq,
1487 producer,
1488 now_ms,
1489 });
1490 match response {
1491 StreamResponse::Appended {
1492 offset,
1493 next_offset,
1494 closed,
1495 deduplicated,
1496 producer,
1497 ..
1498 } => {
1499 let stream_append_count = self
1500 .stream_append_counts
1501 .entry(stream_count_key)
1502 .or_insert(0);
1503 if !deduplicated {
1504 self.commit_index += 1;
1505 *stream_append_count += 1;
1506 }
1507 Ok(AppendResponse {
1508 placement,
1509 start_offset: offset,
1510 next_offset,
1511 stream_append_count: *stream_append_count,
1512 group_commit_index: self.commit_index,
1513 closed,
1514 deduplicated,
1515 producer,
1516 })
1517 }
1518 StreamResponse::Error {
1519 code,
1520 message,
1521 next_offset,
1522 context,
1523 } => Err(GroupEngineError::stream_with_context(
1524 code,
1525 message,
1526 next_offset,
1527 context,
1528 )),
1529 other => Err(GroupEngineError::new(format!(
1530 "unexpected append response: {other:?}"
1531 ))),
1532 }
1533 }
1534
1535 pub fn read_stream_plan(
1536 &mut self,
1537 request: &ReadStreamRequest,
1538 placement: ShardPlacement,
1539 ) -> Result<StreamReadPlan, GroupEngineError> {
1540 self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1541 self.read_stream_plan_after_access(request)
1542 }
1543
1544 pub fn read_stream_plan_after_access(
1545 &self,
1546 request: &ReadStreamRequest,
1547 ) -> Result<StreamReadPlan, GroupEngineError> {
1548 self.state_machine
1549 .read_plan_at(
1550 &request.stream_id,
1551 request.offset,
1552 request.max_len,
1553 request.now_ms,
1554 )
1555 .map_err(stream_response_error)
1556 }
1557
1558 pub fn head_stream_after_access(
1559 &mut self,
1560 request: &HeadStreamRequest,
1561 placement: ShardPlacement,
1562 ) -> Result<HeadStreamResponse, GroupEngineError> {
1563 let Some(metadata) = self
1564 .state_machine
1565 .head_at(&request.stream_id, request.now_ms)
1566 else {
1567 return Err(GroupEngineError::stream(
1568 StreamErrorCode::StreamNotFound,
1569 format!("stream '{}' does not exist", request.stream_id),
1570 ));
1571 };
1572 let content_type = metadata.content_type.clone();
1573 let tail_offset = metadata.tail_offset;
1574 let closed = metadata.status == ursula_stream::StreamStatus::Closed;
1575 let stream_ttl_seconds = metadata.stream_ttl_seconds;
1576 let stream_expires_at_ms = metadata.stream_expires_at_ms;
1577 let _ = metadata;
1578 Ok(HeadStreamResponse {
1579 placement,
1580 content_type,
1581 tail_offset,
1582 cold_hot_start_offset: self.state_machine.hot_start_offset(&request.stream_id),
1583 closed,
1584 stream_ttl_seconds,
1585 stream_expires_at_ms,
1586 snapshot_offset: self
1587 .state_machine
1588 .latest_snapshot(&request.stream_id)
1589 .map_err(stream_response_error)?
1590 .map(|snapshot| snapshot.offset),
1591 integrity: self
1592 .state_machine
1593 .integrity_snapshot(&request.stream_id)
1594 .map_err(stream_response_error)?,
1595 })
1596 }
1597
1598 pub async fn read_payload_from_plan(
1599 cold_store: Option<&ColdStoreHandle>,
1600 cold_index_cache: Option<&Arc<ColdIndexPageCache<ColdStoreColdIndexPageStore>>>,
1601 stream_id: &BucketStreamId,
1602 plan: &StreamReadPlan,
1603 ) -> Result<Vec<u8>, GroupEngineError> {
1604 let mut payload = Vec::new();
1605 for segment in &plan.segments {
1606 match segment {
1607 StreamReadSegment::Hot(bytes) => payload.extend_from_slice(bytes),
1608 StreamReadSegment::ColdIndex(segment) => {
1609 let Some(cold_store) = cold_store else {
1610 return Err(GroupEngineError::stream_with_next_offset(
1611 StreamErrorCode::InvalidColdFlush,
1612 format!("stream '{stream_id}' read requires object payload store"),
1613 Some(plan.next_offset),
1614 ));
1615 };
1616 let Some(cache) = cold_index_cache else {
1617 return Err(GroupEngineError::stream_with_next_offset(
1618 StreamErrorCode::InvalidColdFlush,
1619 format!("stream '{stream_id}' read requires cold index page cache"),
1620 Some(plan.next_offset),
1621 ));
1622 };
1623 let objects = cache
1624 .object_segments_for_read(stream_id, segment)
1625 .await
1626 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1627 let segment_end = segment
1628 .read_start_offset
1629 .saturating_add(u64::try_from(segment.len).expect("read len fits u64"));
1630 for object in objects {
1631 let start = object.start_offset.max(segment.read_start_offset);
1632 let end = object.end_offset.min(segment_end);
1633 if start >= end {
1634 continue;
1635 }
1636 let bytes = cold_store
1637 .read_object_range_for_stream(
1638 stream_id,
1639 &object,
1640 start,
1641 usize::try_from(end - start).expect("object read len fits usize"),
1642 )
1643 .await
1644 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1645 payload.extend_from_slice(&bytes);
1646 }
1647 }
1648 StreamReadSegment::Object(segment) => {
1649 let Some(cold_store) = cold_store else {
1650 return Err(GroupEngineError::stream_with_next_offset(
1651 StreamErrorCode::InvalidColdFlush,
1652 format!("stream '{stream_id}' read requires object payload store"),
1653 Some(plan.next_offset),
1654 ));
1655 };
1656 let bytes = cold_store
1657 .read_object_range_for_stream(
1658 stream_id,
1659 &segment.object,
1660 segment.read_start_offset,
1661 segment.len,
1662 )
1663 .await
1664 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1665 payload.extend_from_slice(&bytes);
1666 }
1667 }
1668 }
1669 Ok(payload)
1670 }
1671
1672 pub(crate) async fn read_own_payload_from_plan(
1673 &self,
1674 stream_id: &BucketStreamId,
1675 plan: &StreamReadPlan,
1676 ) -> Result<Vec<u8>, GroupEngineError> {
1677 Self::read_payload_from_plan(
1678 self.cold_store.as_ref(),
1679 self.cold_index_cache.as_ref(),
1680 stream_id,
1681 plan,
1682 )
1683 .await
1684 }
1685
1686 pub(crate) async fn bootstrap_updates(
1687 &self,
1688 stream_id: &BucketStreamId,
1689 records: &[StreamMessageRecord],
1690 content_type: &str,
1691 now_ms: u64,
1692 ) -> Result<Vec<BootstrapUpdate>, GroupEngineError> {
1693 let mut updates = Vec::with_capacity(records.len());
1694 for record in records {
1695 let len = usize::try_from(record.end_offset - record.start_offset).map_err(|_| {
1696 GroupEngineError::stream(
1697 StreamErrorCode::InvalidSnapshot,
1698 format!(
1699 "bootstrap message [{}..{}) for stream '{stream_id}' is too large",
1700 record.start_offset, record.end_offset
1701 ),
1702 )
1703 })?;
1704 let plan = self
1705 .state_machine
1706 .read_plan_at(stream_id, record.start_offset, len, now_ms)
1707 .map_err(stream_response_error)?;
1708 let payload = self.read_own_payload_from_plan(stream_id, &plan).await?;
1709 updates.push(BootstrapUpdate {
1710 start_offset: record.start_offset,
1711 next_offset: record.end_offset,
1712 content_type: content_type.to_owned(),
1713 payload,
1714 });
1715 }
1716 Ok(updates)
1717 }
1718
1719 pub(crate) fn build_snapshot(&self, placement: ShardPlacement) -> GroupSnapshot {
1720 let stream_snapshot = self.state_machine.snapshot();
1721 let stream_append_counts = self.stream_append_counts_snapshot(&stream_snapshot);
1722 GroupSnapshot {
1723 placement,
1724 group_commit_index: self.commit_index,
1725 stream_snapshot,
1726 stream_append_counts,
1727 }
1728 }
1729
1730 pub(crate) fn stream_append_counts_snapshot(
1731 &self,
1732 stream_snapshot: &ursula_stream::StreamSnapshot,
1733 ) -> Vec<StreamAppendCount> {
1734 let live: HashSet<&BucketStreamId> = stream_snapshot
1741 .streams
1742 .iter()
1743 .map(|entry| &entry.metadata.stream_id)
1744 .collect();
1745 let mut counts = self
1746 .stream_append_counts
1747 .iter()
1748 .filter(|(stream_id, _)| live.contains(stream_id))
1749 .map(|(stream_id, append_count)| StreamAppendCount {
1750 stream_id: stream_id.clone(),
1751 append_count: *append_count,
1752 })
1753 .collect::<Vec<_>>();
1754 counts.sort_by(|left, right| compare_stream_ids(&left.stream_id, &right.stream_id));
1755 counts
1756 }
1757
1758 pub fn stream_tail_offset(&self, stream_id: &BucketStreamId) -> Option<u64> {
1759 self.state_machine
1760 .head(stream_id)
1761 .map(|metadata| metadata.tail_offset)
1762 }
1763
1764 pub(crate) fn install_snapshot_inner(
1765 &mut self,
1766 snapshot: GroupSnapshot,
1767 ) -> Result<(), GroupEngineError> {
1768 let GroupSnapshot {
1769 placement: _,
1770 group_commit_index,
1771 stream_snapshot,
1772 stream_append_counts,
1773 } = snapshot;
1774 self.install_snapshot_parts(group_commit_index, stream_snapshot, stream_append_counts)
1775 }
1776
1777 pub(crate) fn install_snapshot_parts(
1778 &mut self,
1779 group_commit_index: u64,
1780 stream_snapshot: StreamSnapshot,
1781 stream_append_counts: Vec<StreamAppendCount>,
1782 ) -> Result<(), GroupEngineError> {
1783 let stream_ids = stream_snapshot
1784 .streams
1785 .iter()
1786 .map(|entry| entry.metadata.stream_id.clone())
1787 .collect::<HashSet<_>>();
1788 let state_machine = StreamStateMachine::restore(stream_snapshot)
1789 .map_err(|err| GroupEngineError::new(format!("restore stream snapshot: {err}")))?;
1790 let stream_append_counts = restore_stream_append_counts(stream_append_counts, &stream_ids)?;
1791
1792 self.commit_index = group_commit_index;
1793 self.state_machine = state_machine;
1794 self.stream_append_counts = stream_append_counts;
1795 Ok(())
1796 }
1797}
1798
1799impl GroupEngine for InMemoryGroupEngine {
1800 fn create_stream<'a>(
1801 &'a mut self,
1802 request: CreateStreamRequest,
1803 placement: ShardPlacement,
1804 ) -> GroupCreateStreamFuture<'a> {
1805 let command = GroupWriteCommand::from(request);
1806 Box::pin(async move {
1807 match self.apply_committed_write(command, placement)? {
1808 GroupWriteResponse::CreateStream(response) => Ok(response),
1809 other => Err(GroupEngineError::new(format!(
1810 "unexpected create stream write response: {other:?}"
1811 ))),
1812 }
1813 })
1814 }
1815
1816 fn create_stream_with_cold_admission<'a>(
1817 &'a mut self,
1818 request: CreateStreamRequest,
1819 placement: ShardPlacement,
1820 admission: ColdWriteAdmission,
1821 ) -> GroupCreateStreamFuture<'a> {
1822 if !admission.is_enabled() {
1823 return self.create_stream(request, placement);
1824 }
1825 Box::pin(
1826 async move { self.create_stream_with_admission_inner(request, placement, admission) },
1827 )
1828 }
1829
1830 fn create_stream_external<'a>(
1831 &'a mut self,
1832 request: CreateStreamExternalRequest,
1833 placement: ShardPlacement,
1834 ) -> GroupCreateStreamFuture<'a> {
1835 Box::pin(async move {
1836 if let Some(cold_store) = self.cold_store.as_ref() {
1837 let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1838 write_external_segment_index_pages(
1839 &store,
1840 &request.stream_id,
1841 0,
1842 &request.initial_payload,
1843 )
1844 .await
1845 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1846 }
1847 let command = GroupWriteCommand::from(request);
1848 match self.apply_committed_write(command, placement)? {
1849 GroupWriteResponse::CreateStream(response) => Ok(response),
1850 other => Err(GroupEngineError::new(format!(
1851 "unexpected external create stream write response: {other:?}"
1852 ))),
1853 }
1854 })
1855 }
1856
1857 fn read_stream<'a>(
1858 &'a mut self,
1859 request: ReadStreamRequest,
1860 placement: ShardPlacement,
1861 ) -> GroupReadStreamFuture<'a> {
1862 Box::pin(async move {
1863 self.read_stream_parts(request, placement)
1864 .await?
1865 .into_response()
1866 .await
1867 })
1868 }
1869
1870 fn read_stream_parts<'a>(
1871 &'a mut self,
1872 request: ReadStreamRequest,
1873 placement: ShardPlacement,
1874 ) -> GroupReadStreamPartsFuture<'a> {
1875 Box::pin(async move {
1876 let stream_id = request.stream_id.clone();
1877 let plan = self.read_stream_plan(&request, placement)?;
1878 Ok(GroupReadStreamParts::from_plan(
1879 placement,
1880 stream_id,
1881 plan,
1882 self.cold_store(),
1883 self.cold_index_cache.clone(),
1884 ))
1885 })
1886 }
1887
1888 fn publish_snapshot<'a>(
1889 &'a mut self,
1890 request: PublishSnapshotRequest,
1891 placement: ShardPlacement,
1892 ) -> GroupPublishSnapshotFuture<'a> {
1893 Box::pin(async move {
1894 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1895 let command = GroupWriteCommand::from(request);
1896 match self.apply_committed_write(command, placement)? {
1897 GroupWriteResponse::PublishSnapshot(response) => Ok(response),
1898 other => Err(GroupEngineError::new(format!(
1899 "unexpected publish snapshot write response: {other:?}"
1900 ))),
1901 }
1902 })
1903 }
1904
1905 fn read_snapshot<'a>(
1906 &'a mut self,
1907 request: ReadSnapshotRequest,
1908 placement: ShardPlacement,
1909 ) -> GroupReadSnapshotFuture<'a> {
1910 Box::pin(async move {
1911 self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1912 let snapshot = match request.snapshot_offset {
1913 Some(offset) => self
1914 .state_machine
1915 .read_snapshot(&request.stream_id, offset)
1916 .map_err(stream_response_error)?,
1917 None => self
1918 .state_machine
1919 .latest_snapshot(&request.stream_id)
1920 .map_err(stream_response_error)?
1921 .ok_or_else(|| {
1922 GroupEngineError::stream(
1923 StreamErrorCode::SnapshotNotFound,
1924 format!("stream '{}' has no visible snapshot", request.stream_id),
1925 )
1926 })?,
1927 };
1928 let tail_offset = self
1929 .state_machine
1930 .head_at(&request.stream_id, request.now_ms)
1931 .map(|metadata| metadata.tail_offset)
1932 .unwrap_or(snapshot.offset);
1933 Ok(ReadSnapshotResponse {
1934 placement,
1935 snapshot_offset: snapshot.offset,
1936 next_offset: snapshot.offset,
1937 content_type: snapshot.content_type,
1938 payload: snapshot.payload,
1939 up_to_date: snapshot.offset == tail_offset,
1940 })
1941 })
1942 }
1943
1944 fn delete_snapshot<'a>(
1945 &'a mut self,
1946 request: DeleteSnapshotRequest,
1947 placement: ShardPlacement,
1948 ) -> GroupDeleteSnapshotFuture<'a> {
1949 Box::pin(async move {
1950 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1951 match self
1952 .state_machine
1953 .delete_snapshot(&request.stream_id, request.snapshot_offset)
1954 {
1955 StreamResponse::Error {
1956 code,
1957 message,
1958 next_offset,
1959 context,
1960 } => Err(GroupEngineError::stream_with_context(
1961 code,
1962 message,
1963 next_offset,
1964 context,
1965 )),
1966 other => Err(GroupEngineError::new(format!(
1967 "unexpected delete snapshot response: {other:?}"
1968 ))),
1969 }
1970 })
1971 }
1972
1973 fn bootstrap_stream<'a>(
1974 &'a mut self,
1975 request: BootstrapStreamRequest,
1976 placement: ShardPlacement,
1977 ) -> GroupBootstrapStreamFuture<'a> {
1978 Box::pin(async move {
1979 self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1980 let plan = self
1981 .state_machine
1982 .bootstrap_plan(&request.stream_id)
1983 .map_err(stream_response_error)?;
1984 let snapshot_offset = plan.snapshot.as_ref().map(|snapshot| snapshot.offset);
1985 let snapshot_content_type = plan
1986 .snapshot
1987 .as_ref()
1988 .map(|snapshot| snapshot.content_type.clone())
1989 .unwrap_or_else(|| DEFAULT_CONTENT_TYPE.to_owned());
1990 let snapshot_payload = plan
1991 .snapshot
1992 .as_ref()
1993 .map(|snapshot| snapshot.payload.clone())
1994 .unwrap_or_default();
1995 let updates = self
1996 .bootstrap_updates(
1997 &request.stream_id,
1998 &plan.updates,
1999 &plan.content_type,
2000 request.now_ms,
2001 )
2002 .await?;
2003 Ok(BootstrapStreamResponse {
2004 placement,
2005 snapshot_offset,
2006 snapshot_content_type,
2007 snapshot_payload,
2008 updates,
2009 next_offset: plan.next_offset,
2010 up_to_date: plan.up_to_date,
2011 closed: plan.closed,
2012 })
2013 })
2014 }
2015
2016 fn touch_stream_access<'a>(
2017 &'a mut self,
2018 stream_id: BucketStreamId,
2019 now_ms: u64,
2020 renew_ttl: bool,
2021 placement: ShardPlacement,
2022 ) -> GroupTouchStreamAccessFuture<'a> {
2023 Box::pin(async move { self.apply_access_command(stream_id, now_ms, renew_ttl, placement) })
2024 }
2025
2026 fn add_fork_ref<'a>(
2027 &'a mut self,
2028 stream_id: BucketStreamId,
2029 now_ms: u64,
2030 placement: ShardPlacement,
2031 ) -> GroupForkRefFuture<'a> {
2032 Box::pin(async move {
2033 match self.apply_committed_write(
2034 GroupWriteCommand::AddForkRef { stream_id, now_ms },
2035 placement,
2036 )? {
2037 GroupWriteResponse::AddForkRef(response) => Ok(response),
2038 other => Err(GroupEngineError::new(format!(
2039 "unexpected add fork ref write response: {other:?}"
2040 ))),
2041 }
2042 })
2043 }
2044
2045 fn release_fork_ref<'a>(
2046 &'a mut self,
2047 stream_id: BucketStreamId,
2048 placement: ShardPlacement,
2049 ) -> GroupForkRefFuture<'a> {
2050 Box::pin(async move {
2051 match self
2052 .apply_committed_write(GroupWriteCommand::ReleaseForkRef { stream_id }, placement)?
2053 {
2054 GroupWriteResponse::ReleaseForkRef(response) => Ok(response),
2055 other => Err(GroupEngineError::new(format!(
2056 "unexpected release fork ref write response: {other:?}"
2057 ))),
2058 }
2059 })
2060 }
2061
2062 fn head_stream<'a>(
2063 &'a mut self,
2064 request: HeadStreamRequest,
2065 placement: ShardPlacement,
2066 ) -> GroupHeadStreamFuture<'a> {
2067 Box::pin(async move {
2068 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
2069 self.head_stream_after_access(&request, placement)
2070 })
2071 }
2072
2073 fn close_stream<'a>(
2074 &'a mut self,
2075 request: CloseStreamRequest,
2076 placement: ShardPlacement,
2077 ) -> GroupCloseStreamFuture<'a> {
2078 Box::pin(async move {
2079 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
2080 let command = GroupWriteCommand::from(request);
2081 match self.apply_committed_write(command, placement)? {
2082 GroupWriteResponse::CloseStream(response) => Ok(response),
2083 other => Err(GroupEngineError::new(format!(
2084 "unexpected close stream write response: {other:?}"
2085 ))),
2086 }
2087 })
2088 }
2089
2090 fn delete_stream<'a>(
2091 &'a mut self,
2092 request: DeleteStreamRequest,
2093 placement: ShardPlacement,
2094 ) -> GroupDeleteStreamFuture<'a> {
2095 let command = GroupWriteCommand::from(request);
2096 Box::pin(async move {
2097 match self.apply_committed_write(command, placement)? {
2098 GroupWriteResponse::DeleteStream(response) => Ok(response),
2099 other => Err(GroupEngineError::new(format!(
2100 "unexpected delete stream write response: {other:?}"
2101 ))),
2102 }
2103 })
2104 }
2105
2106 fn ack_cold_gc<'a>(
2107 &'a mut self,
2108 up_to_seq: u64,
2109 placement: ShardPlacement,
2110 ) -> GroupAckColdGcFuture<'a> {
2111 Box::pin(async move {
2112 match self
2113 .apply_committed_write(GroupWriteCommand::AckColdGc { up_to_seq }, placement)?
2114 {
2115 GroupWriteResponse::AckColdGc(response) => Ok(response),
2116 other => Err(GroupEngineError::new(format!(
2117 "unexpected ack cold gc write response: {other:?}"
2118 ))),
2119 }
2120 })
2121 }
2122
2123 fn plan_cold_gc<'a>(
2124 &'a mut self,
2125 max: usize,
2126 _placement: ShardPlacement,
2127 ) -> GroupPlanColdGcFuture<'a> {
2128 let entries = self.state_machine.pending_cold_gc_batch(max);
2129 Box::pin(async move { Ok(entries) })
2130 }
2131
2132 fn append<'a>(
2133 &'a mut self,
2134 request: AppendRequest,
2135 placement: ShardPlacement,
2136 ) -> GroupAppendFuture<'a> {
2137 Box::pin(async move {
2138 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
2139 let command = GroupWriteCommand::from(request);
2140 match self.apply_committed_write(command, placement)? {
2141 GroupWriteResponse::Append(response) => Ok(response),
2142 other => Err(GroupEngineError::new(format!(
2143 "unexpected append write response: {other:?}"
2144 ))),
2145 }
2146 })
2147 }
2148
2149 fn append_with_cold_admission<'a>(
2150 &'a mut self,
2151 request: AppendRequest,
2152 placement: ShardPlacement,
2153 admission: ColdWriteAdmission,
2154 ) -> GroupAppendFuture<'a> {
2155 if !admission.is_enabled() {
2156 return self.append(request, placement);
2157 }
2158 Box::pin(async move { self.append_with_admission_inner(request, placement, admission) })
2159 }
2160
2161 fn append_external<'a>(
2162 &'a mut self,
2163 request: AppendExternalRequest,
2164 placement: ShardPlacement,
2165 ) -> GroupAppendFuture<'a> {
2166 Box::pin(async move {
2167 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
2168 if let Some(cold_store) = self.cold_store.as_ref() {
2169 let start_offset = self
2170 .state_machine
2171 .head(&request.stream_id)
2172 .map(|metadata| metadata.tail_offset)
2173 .ok_or_else(|| {
2174 GroupEngineError::stream(
2175 ursula_stream::StreamErrorCode::StreamNotFound,
2176 format!("stream '{}' does not exist", request.stream_id),
2177 )
2178 })?;
2179 let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
2180 write_external_segment_index_pages(
2181 &store,
2182 &request.stream_id,
2183 start_offset,
2184 &request.payload,
2185 )
2186 .await
2187 .map_err(|err| GroupEngineError::new(err.to_string()))?;
2188 }
2189 let command = GroupWriteCommand::from(request);
2190 match self.apply_committed_write(command, placement)? {
2191 GroupWriteResponse::Append(response) => Ok(response),
2192 other => Err(GroupEngineError::new(format!(
2193 "unexpected external append write response: {other:?}"
2194 ))),
2195 }
2196 })
2197 }
2198
2199 fn append_batch<'a>(
2200 &'a mut self,
2201 request: AppendBatchRequest,
2202 placement: ShardPlacement,
2203 ) -> GroupAppendBatchFuture<'a> {
2204 Box::pin(async move {
2205 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
2206 let command = GroupWriteCommand::from(request);
2207 match self.apply_committed_write(command, placement)? {
2208 GroupWriteResponse::AppendBatch(response) => Ok(response),
2209 other => Err(GroupEngineError::new(format!(
2210 "unexpected append batch write response: {other:?}"
2211 ))),
2212 }
2213 })
2214 }
2215
2216 fn append_batch_with_cold_admission<'a>(
2217 &'a mut self,
2218 request: AppendBatchRequest,
2219 placement: ShardPlacement,
2220 admission: ColdWriteAdmission,
2221 ) -> GroupAppendBatchFuture<'a> {
2222 if !admission.is_enabled() {
2223 return self.append_batch(request, placement);
2224 }
2225 Box::pin(
2226 async move { self.append_batch_with_admission_inner(request, placement, admission) },
2227 )
2228 }
2229
2230 fn flush_cold<'a>(
2231 &'a mut self,
2232 request: FlushColdRequest,
2233 placement: ShardPlacement,
2234 ) -> GroupFlushColdFuture<'a> {
2235 Box::pin(async move {
2236 let mut index_rollback = None;
2237 if let Some(cold_store) = self.cold_store.as_ref() {
2238 let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
2239 let rollback = write_cold_chunk_index_pages_with_rollback(
2240 &store,
2241 &request.stream_id,
2242 &request.chunk,
2243 )
2244 .await
2245 .map_err(|err| GroupEngineError::new(err.to_string()))?;
2246 index_rollback = Some((store, rollback));
2247 }
2248 let command = GroupWriteCommand::from(request);
2249 match self.apply_committed_write(command, placement) {
2250 Ok(GroupWriteResponse::FlushCold(response)) => Ok(response),
2251 Ok(other) => {
2252 if let Some((store, rollback)) = index_rollback {
2253 rollback_cold_index_pages(&store, rollback)
2254 .await
2255 .map_err(|err| GroupEngineError::new(err.to_string()))?;
2256 }
2257 Err(GroupEngineError::new(format!(
2258 "unexpected flush cold write response: {other:?}"
2259 )))
2260 }
2261 Err(err) => {
2262 if let Some((store, rollback)) = index_rollback {
2263 rollback_cold_index_pages(&store, rollback).await.map_err(
2264 |rollback_err| {
2265 GroupEngineError::new(format!(
2266 "rollback cold index after flush failure: {rollback_err}"
2267 ))
2268 },
2269 )?;
2270 }
2271 Err(err)
2272 }
2273 }
2274 })
2275 }
2276
2277 fn plan_cold_flush<'a>(
2278 &'a mut self,
2279 request: PlanColdFlushRequest,
2280 _placement: ShardPlacement,
2281 ) -> GroupPlanColdFlushFuture<'a> {
2282 Box::pin(async move {
2283 self.state_machine
2284 .plan_cold_flush(
2285 &request.stream_id,
2286 request.min_hot_bytes,
2287 request.max_flush_bytes,
2288 )
2289 .map_err(stream_response_error)
2290 })
2291 }
2292
2293 fn plan_next_cold_flush<'a>(
2294 &'a mut self,
2295 request: PlanGroupColdFlushRequest,
2296 _placement: ShardPlacement,
2297 ) -> GroupPlanNextColdFlushFuture<'a> {
2298 Box::pin(async move {
2299 self.state_machine
2300 .plan_next_cold_flush(request.min_hot_bytes, request.max_flush_bytes)
2301 .map_err(stream_response_error)
2302 })
2303 }
2304
2305 fn plan_next_cold_flush_batch<'a>(
2306 &'a mut self,
2307 request: PlanGroupColdFlushRequest,
2308 _placement: ShardPlacement,
2309 max_candidates: usize,
2310 ) -> GroupPlanNextColdFlushBatchFuture<'a> {
2311 Box::pin(async move {
2312 self.state_machine
2313 .plan_next_cold_flush_batch(
2314 request.min_hot_bytes,
2315 request.max_flush_bytes,
2316 max_candidates,
2317 )
2318 .map_err(stream_response_error)
2319 })
2320 }
2321
2322 fn cold_hot_backlog<'a>(
2323 &'a mut self,
2324 stream_id: BucketStreamId,
2325 _placement: ShardPlacement,
2326 ) -> GroupColdHotBacklogFuture<'a> {
2327 Box::pin(async move { self.cold_hot_backlog_for(stream_id) })
2328 }
2329
2330 fn snapshot<'a>(&'a mut self, placement: ShardPlacement) -> GroupSnapshotFuture<'a> {
2331 Box::pin(async move { Ok(self.build_snapshot(placement)) })
2332 }
2333
2334 fn install_snapshot<'a>(
2335 &'a mut self,
2336 snapshot: GroupSnapshot,
2337 ) -> GroupInstallSnapshotFuture<'a> {
2338 Box::pin(async move { self.install_snapshot_inner(snapshot) })
2339 }
2340}
2341
2342#[derive(Debug, Clone, Default)]
2343pub struct InMemoryGroupEngineFactory {
2344 cold_store: Option<ColdStoreHandle>,
2345}
2346
2347impl InMemoryGroupEngineFactory {
2348 pub fn new() -> Self {
2349 Self::default()
2350 }
2351
2352 pub fn with_cold_store(cold_store: Option<ColdStoreHandle>) -> Self {
2353 Self { cold_store }
2354 }
2355}
2356
2357impl GroupEngineFactory for InMemoryGroupEngineFactory {
2358 fn create<'a>(
2359 &'a self,
2360 _placement: ShardPlacement,
2361 _metrics: GroupEngineMetrics,
2362 ) -> GroupEngineCreateFuture<'a> {
2363 Box::pin(async move {
2364 let mut engine = InMemoryGroupEngine::default();
2365 engine.set_cold_store(self.cold_store.clone());
2366 let engine: Box<dyn GroupEngine> = Box::new(engine);
2367 Ok(engine)
2368 })
2369 }
2370}
2371
2372pub(crate) fn compare_stream_ids(
2373 left: &BucketStreamId,
2374 right: &BucketStreamId,
2375) -> std::cmp::Ordering {
2376 left.bucket_id
2377 .cmp(&right.bucket_id)
2378 .then_with(|| left.stream_id.cmp(&right.stream_id))
2379}
2380pub(crate) fn ensure_bucket_exists(
2381 state_machine: &mut StreamStateMachine,
2382 stream_id: &BucketStreamId,
2383) -> Result<(), GroupEngineError> {
2384 if state_machine.bucket_exists(&stream_id.bucket_id) {
2385 return Ok(());
2386 }
2387
2388 match state_machine.apply(StreamCommand::CreateBucket {
2389 bucket_id: stream_id.bucket_id.clone(),
2390 }) {
2391 StreamResponse::BucketCreated { .. } | StreamResponse::BucketAlreadyExists { .. } => Ok(()),
2392 StreamResponse::Error {
2393 code,
2394 message,
2395 next_offset,
2396 context,
2397 } => Err(GroupEngineError::stream_with_context(
2398 code,
2399 message,
2400 next_offset,
2401 context,
2402 )),
2403 other => Err(GroupEngineError::new(format!(
2404 "unexpected create bucket response: {other:?}"
2405 ))),
2406 }
2407}
2408
2409pub(crate) fn stream_response_error(response: StreamResponse) -> GroupEngineError {
2410 match response {
2411 StreamResponse::Error {
2412 code,
2413 message,
2414 next_offset,
2415 context,
2416 } => GroupEngineError::stream_with_context(code, message, next_offset, context),
2417 other => GroupEngineError::new(format!("unexpected stream response error: {other:?}")),
2418 }
2419}
2420
2421pub(crate) fn restore_stream_append_counts(
2422 counts: Vec<StreamAppendCount>,
2423 snapshot_stream_ids: &HashSet<BucketStreamId>,
2424) -> Result<HashMap<BucketStreamId, u64>, GroupEngineError> {
2425 let mut restored = HashMap::with_capacity(counts.len());
2426 for count in counts {
2427 if !snapshot_stream_ids.contains(&count.stream_id) {
2428 return Err(GroupEngineError::new(format!(
2429 "append count references missing snapshot stream '{}'",
2430 count.stream_id
2431 )));
2432 }
2433 if restored
2434 .insert(count.stream_id.clone(), count.append_count)
2435 .is_some()
2436 {
2437 return Err(GroupEngineError::new(format!(
2438 "snapshot contains duplicate append count for stream '{}'",
2439 count.stream_id
2440 )));
2441 }
2442 }
2443 Ok(restored)
2444}