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