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 append_payload(
1015 &mut self,
1016 input: AppendPayloadInput<'_>,
1017 placement: ShardPlacement,
1018 ) -> Result<AppendResponse, GroupEngineError> {
1019 let AppendPayloadInput {
1020 stream_id,
1021 content_type,
1022 payload,
1023 close_after,
1024 stream_seq,
1025 producer,
1026 now_ms,
1027 } = input;
1028 let stream_count_key = stream_id.clone();
1029 let response = self.state_machine.append_borrowed(AppendStreamInput {
1030 stream_id,
1031 content_type,
1032 payload,
1033 close_after,
1034 stream_seq,
1035 producer,
1036 now_ms,
1037 });
1038 match response {
1039 StreamResponse::Appended {
1040 offset,
1041 next_offset,
1042 closed,
1043 deduplicated,
1044 producer,
1045 ..
1046 } => {
1047 let stream_append_count = self
1048 .stream_append_counts
1049 .entry(stream_count_key)
1050 .or_insert(0);
1051 if !deduplicated {
1052 self.commit_index += 1;
1053 *stream_append_count += 1;
1054 }
1055 Ok(AppendResponse {
1056 placement,
1057 start_offset: offset,
1058 next_offset,
1059 stream_append_count: *stream_append_count,
1060 group_commit_index: self.commit_index,
1061 closed,
1062 deduplicated,
1063 producer,
1064 })
1065 }
1066 StreamResponse::Error {
1067 code,
1068 message,
1069 next_offset,
1070 context,
1071 } => Err(GroupEngineError::stream_with_context(
1072 code,
1073 message,
1074 next_offset,
1075 context,
1076 )),
1077 other => Err(GroupEngineError::new(format!(
1078 "unexpected append response: {other:?}"
1079 ))),
1080 }
1081 }
1082
1083 pub fn read_stream_plan(
1084 &mut self,
1085 request: &ReadStreamRequest,
1086 placement: ShardPlacement,
1087 ) -> Result<StreamReadPlan, GroupEngineError> {
1088 self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1089 self.read_stream_plan_after_access(request)
1090 }
1091
1092 pub fn read_stream_plan_after_access(
1093 &self,
1094 request: &ReadStreamRequest,
1095 ) -> Result<StreamReadPlan, GroupEngineError> {
1096 self.state_machine
1097 .read_plan_at(
1098 &request.stream_id,
1099 request.offset,
1100 request.max_len,
1101 request.now_ms,
1102 )
1103 .map_err(stream_response_error)
1104 }
1105
1106 pub fn head_stream_after_access(
1107 &mut self,
1108 request: &HeadStreamRequest,
1109 placement: ShardPlacement,
1110 ) -> Result<HeadStreamResponse, GroupEngineError> {
1111 let Some(metadata) = self
1112 .state_machine
1113 .head_at(&request.stream_id, request.now_ms)
1114 else {
1115 return Err(GroupEngineError::stream(
1116 StreamErrorCode::StreamNotFound,
1117 format!("stream '{}' does not exist", request.stream_id),
1118 ));
1119 };
1120 let content_type = metadata.content_type.clone();
1121 let tail_offset = metadata.tail_offset;
1122 let closed = metadata.status == ursula_stream::StreamStatus::Closed;
1123 let stream_ttl_seconds = metadata.stream_ttl_seconds;
1124 let stream_expires_at_ms = metadata.stream_expires_at_ms;
1125 let _ = metadata;
1126 Ok(HeadStreamResponse {
1127 placement,
1128 content_type,
1129 tail_offset,
1130 cold_hot_start_offset: self.state_machine.hot_start_offset(&request.stream_id),
1131 closed,
1132 stream_ttl_seconds,
1133 stream_expires_at_ms,
1134 snapshot_offset: self
1135 .state_machine
1136 .latest_snapshot(&request.stream_id)
1137 .map_err(stream_response_error)?
1138 .map(|snapshot| snapshot.offset),
1139 integrity: self
1140 .state_machine
1141 .integrity_snapshot(&request.stream_id)
1142 .map_err(stream_response_error)?,
1143 })
1144 }
1145
1146 pub fn get_stream_attrs_after_access(
1147 &mut self,
1148 request: &GetStreamAttrsRequest,
1149 placement: ShardPlacement,
1150 ) -> Result<GetStreamAttrsResponse, GroupEngineError> {
1151 if self
1152 .state_machine
1153 .head_at(&request.stream_id, request.now_ms)
1154 .is_none()
1155 {
1156 return Err(GroupEngineError::stream(
1157 StreamErrorCode::StreamNotFound,
1158 format!("stream '{}' does not exist", request.stream_id),
1159 ));
1160 }
1161 Ok(GetStreamAttrsResponse {
1162 placement,
1163 attrs: self.state_machine.stream_attrs(&request.stream_id).cloned(),
1164 })
1165 }
1166
1167 pub async fn read_payload_from_plan(
1168 cold_store: Option<&ColdStoreHandle>,
1169 cold_index_cache: Option<&Arc<ColdIndexPageCache<ColdStoreColdIndexPageStore>>>,
1170 stream_id: &BucketStreamId,
1171 plan: &StreamReadPlan,
1172 ) -> Result<Vec<u8>, GroupEngineError> {
1173 let mut payload = Vec::new();
1174 for segment in &plan.segments {
1175 match segment {
1176 StreamReadSegment::Hot(bytes) => payload.extend_from_slice(bytes),
1177 StreamReadSegment::ColdIndex(segment) => {
1178 let Some(cold_store) = cold_store else {
1179 return Err(GroupEngineError::stream_with_next_offset(
1180 StreamErrorCode::InvalidColdFlush,
1181 format!("stream '{stream_id}' read requires object payload store"),
1182 Some(plan.next_offset),
1183 ));
1184 };
1185 let Some(cache) = cold_index_cache else {
1186 return Err(GroupEngineError::stream_with_next_offset(
1187 StreamErrorCode::InvalidColdFlush,
1188 format!("stream '{stream_id}' read requires cold index page cache"),
1189 Some(plan.next_offset),
1190 ));
1191 };
1192 let objects = cache
1193 .object_segments_for_read(stream_id, segment)
1194 .await
1195 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1196 let segment_end = segment
1197 .read_start_offset
1198 .saturating_add(u64::try_from(segment.len).expect("read len fits u64"));
1199 for object in objects {
1200 let start = object.start_offset.max(segment.read_start_offset);
1201 let end = object.end_offset.min(segment_end);
1202 if start >= end {
1203 continue;
1204 }
1205 let bytes = cold_store
1206 .read_object_range_for_stream(
1207 stream_id,
1208 &object,
1209 start,
1210 usize::try_from(end - start).expect("object read len fits usize"),
1211 )
1212 .await
1213 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1214 payload.extend_from_slice(&bytes);
1215 }
1216 }
1217 StreamReadSegment::Object(segment) => {
1218 let Some(cold_store) = cold_store else {
1219 return Err(GroupEngineError::stream_with_next_offset(
1220 StreamErrorCode::InvalidColdFlush,
1221 format!("stream '{stream_id}' read requires object payload store"),
1222 Some(plan.next_offset),
1223 ));
1224 };
1225 let bytes = cold_store
1226 .read_object_range_for_stream(
1227 stream_id,
1228 &segment.object,
1229 segment.read_start_offset,
1230 segment.len,
1231 )
1232 .await
1233 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1234 payload.extend_from_slice(&bytes);
1235 }
1236 }
1237 }
1238 Ok(payload)
1239 }
1240
1241 pub(crate) async fn read_own_payload_from_plan(
1242 &self,
1243 stream_id: &BucketStreamId,
1244 plan: &StreamReadPlan,
1245 ) -> Result<Vec<u8>, GroupEngineError> {
1246 Self::read_payload_from_plan(
1247 self.cold_store.as_ref(),
1248 self.cold_index_cache.as_ref(),
1249 stream_id,
1250 plan,
1251 )
1252 .await
1253 }
1254
1255 pub(crate) async fn bootstrap_updates(
1256 &self,
1257 stream_id: &BucketStreamId,
1258 records: &[StreamMessageRecord],
1259 content_type: &str,
1260 now_ms: u64,
1261 ) -> Result<Vec<BootstrapUpdate>, GroupEngineError> {
1262 let mut updates = Vec::with_capacity(records.len());
1263 for record in records {
1264 let len = usize::try_from(record.end_offset - record.start_offset).map_err(|_| {
1265 GroupEngineError::stream(
1266 StreamErrorCode::InvalidSnapshot,
1267 format!(
1268 "bootstrap message [{}..{}) for stream '{stream_id}' is too large",
1269 record.start_offset, record.end_offset
1270 ),
1271 )
1272 })?;
1273 let plan = self
1274 .state_machine
1275 .read_plan_at(stream_id, record.start_offset, len, now_ms)
1276 .map_err(stream_response_error)?;
1277 let payload = self.read_own_payload_from_plan(stream_id, &plan).await?;
1278 updates.push(BootstrapUpdate {
1279 start_offset: record.start_offset,
1280 next_offset: record.end_offset,
1281 content_type: content_type.to_owned(),
1282 payload,
1283 });
1284 }
1285 Ok(updates)
1286 }
1287
1288 pub(crate) fn build_snapshot(&self, placement: ShardPlacement) -> GroupSnapshot {
1289 let stream_snapshot = self.state_machine.snapshot();
1290 let stream_append_counts = self.stream_append_counts_snapshot(&stream_snapshot);
1291 GroupSnapshot {
1292 placement,
1293 group_commit_index: self.commit_index,
1294 stream_snapshot,
1295 stream_append_counts,
1296 }
1297 }
1298
1299 pub(crate) fn stream_append_counts_snapshot(
1300 &self,
1301 stream_snapshot: &ursula_stream::StreamSnapshot,
1302 ) -> Vec<StreamAppendCount> {
1303 let live: HashSet<&BucketStreamId> = stream_snapshot
1310 .streams
1311 .iter()
1312 .map(|entry| &entry.metadata.stream_id)
1313 .collect();
1314 let mut counts = self
1315 .stream_append_counts
1316 .iter()
1317 .filter(|(stream_id, _)| live.contains(stream_id))
1318 .map(|(stream_id, append_count)| StreamAppendCount {
1319 stream_id: stream_id.clone(),
1320 append_count: *append_count,
1321 })
1322 .collect::<Vec<_>>();
1323 counts.sort_by(|left, right| compare_stream_ids(&left.stream_id, &right.stream_id));
1324 counts
1325 }
1326
1327 pub fn stream_tail_offset(&self, stream_id: &BucketStreamId) -> Option<u64> {
1328 self.state_machine
1329 .head(stream_id)
1330 .map(|metadata| metadata.tail_offset)
1331 }
1332
1333 pub(crate) fn install_snapshot_inner(
1334 &mut self,
1335 snapshot: GroupSnapshot,
1336 ) -> Result<(), GroupEngineError> {
1337 let GroupSnapshot {
1338 placement: _,
1339 group_commit_index,
1340 stream_snapshot,
1341 stream_append_counts,
1342 } = snapshot;
1343 self.install_snapshot_parts(group_commit_index, stream_snapshot, stream_append_counts)
1344 }
1345
1346 pub(crate) fn install_snapshot_parts(
1347 &mut self,
1348 group_commit_index: u64,
1349 stream_snapshot: StreamSnapshot,
1350 stream_append_counts: Vec<StreamAppendCount>,
1351 ) -> Result<(), GroupEngineError> {
1352 let stream_ids = stream_snapshot
1353 .streams
1354 .iter()
1355 .map(|entry| entry.metadata.stream_id.clone())
1356 .collect::<HashSet<_>>();
1357 let state_machine = StreamStateMachine::restore(stream_snapshot)
1358 .map_err(|err| GroupEngineError::new(format!("restore stream snapshot: {err}")))?;
1359 let stream_append_counts = restore_stream_append_counts(stream_append_counts, &stream_ids)?;
1360
1361 self.commit_index = group_commit_index;
1362 self.state_machine = state_machine;
1363 self.stream_append_counts = stream_append_counts;
1364 Ok(())
1365 }
1366}
1367
1368impl GroupEngine for InMemoryGroupEngine {
1369 fn create_stream<'a>(
1370 &'a mut self,
1371 request: CreateStreamRequest,
1372 placement: ShardPlacement,
1373 ) -> GroupCreateStreamFuture<'a> {
1374 let command = GroupWriteCommand::from(request);
1375 Box::pin(async move {
1376 match self.apply_committed_write(command, placement)? {
1377 GroupWriteResponse::CreateStream(response) => Ok(response),
1378 other => Err(GroupEngineError::new(format!(
1379 "unexpected create stream write response: {other:?}"
1380 ))),
1381 }
1382 })
1383 }
1384
1385 fn create_stream_with_cold_admission<'a>(
1386 &'a mut self,
1387 request: CreateStreamRequest,
1388 placement: ShardPlacement,
1389 admission: ColdWriteAdmission,
1390 ) -> GroupCreateStreamFuture<'a> {
1391 if !admission.is_enabled() {
1392 return self.create_stream(request, placement);
1393 }
1394 Box::pin(
1395 async move { self.create_stream_with_admission_inner(request, placement, admission) },
1396 )
1397 }
1398
1399 fn create_stream_external<'a>(
1400 &'a mut self,
1401 request: CreateStreamExternalRequest,
1402 placement: ShardPlacement,
1403 ) -> GroupCreateStreamFuture<'a> {
1404 Box::pin(async move {
1405 if let Some(cold_store) = self.cold_store.as_ref() {
1406 let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1407 write_external_segment_index_pages(
1408 &store,
1409 &request.stream_id,
1410 0,
1411 &request.initial_payload,
1412 )
1413 .await
1414 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1415 }
1416 let command = GroupWriteCommand::from(request);
1417 match self.apply_committed_write(command, placement)? {
1418 GroupWriteResponse::CreateStream(response) => Ok(response),
1419 other => Err(GroupEngineError::new(format!(
1420 "unexpected external create stream write response: {other:?}"
1421 ))),
1422 }
1423 })
1424 }
1425
1426 fn read_stream<'a>(
1427 &'a mut self,
1428 request: ReadStreamRequest,
1429 placement: ShardPlacement,
1430 ) -> GroupReadStreamFuture<'a> {
1431 Box::pin(async move {
1432 self.read_stream_parts(request, placement)
1433 .await?
1434 .into_response()
1435 .await
1436 })
1437 }
1438
1439 fn read_stream_parts<'a>(
1440 &'a mut self,
1441 request: ReadStreamRequest,
1442 placement: ShardPlacement,
1443 ) -> GroupReadStreamPartsFuture<'a> {
1444 Box::pin(async move {
1445 let stream_id = request.stream_id.clone();
1446 let plan = self.read_stream_plan(&request, placement)?;
1447 Ok(GroupReadStreamParts::from_plan(
1448 placement,
1449 stream_id,
1450 plan,
1451 self.cold_store(),
1452 self.cold_index_cache.clone(),
1453 ))
1454 })
1455 }
1456
1457 fn publish_snapshot<'a>(
1458 &'a mut self,
1459 request: PublishSnapshotRequest,
1460 placement: ShardPlacement,
1461 ) -> GroupPublishSnapshotFuture<'a> {
1462 Box::pin(async move {
1463 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1464 let command = GroupWriteCommand::from(request);
1465 match self.apply_committed_write(command, placement)? {
1466 GroupWriteResponse::PublishSnapshot(response) => Ok(response),
1467 other => Err(GroupEngineError::new(format!(
1468 "unexpected publish snapshot write response: {other:?}"
1469 ))),
1470 }
1471 })
1472 }
1473
1474 fn read_snapshot<'a>(
1475 &'a mut self,
1476 request: ReadSnapshotRequest,
1477 placement: ShardPlacement,
1478 ) -> GroupReadSnapshotFuture<'a> {
1479 Box::pin(async move {
1480 self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1481 let snapshot = match request.snapshot_offset {
1482 Some(offset) => self
1483 .state_machine
1484 .read_snapshot(&request.stream_id, offset)
1485 .map_err(stream_response_error)?,
1486 None => self
1487 .state_machine
1488 .latest_snapshot(&request.stream_id)
1489 .map_err(stream_response_error)?
1490 .ok_or_else(|| {
1491 GroupEngineError::stream(
1492 StreamErrorCode::SnapshotNotFound,
1493 format!("stream '{}' has no visible snapshot", request.stream_id),
1494 )
1495 })?,
1496 };
1497 let tail_offset = self
1498 .state_machine
1499 .head_at(&request.stream_id, request.now_ms)
1500 .map(|metadata| metadata.tail_offset)
1501 .unwrap_or(snapshot.offset);
1502 Ok(ReadSnapshotResponse {
1503 placement,
1504 snapshot_offset: snapshot.offset,
1505 next_offset: snapshot.offset,
1506 content_type: snapshot.content_type,
1507 payload: snapshot.payload,
1508 up_to_date: snapshot.offset == tail_offset,
1509 })
1510 })
1511 }
1512
1513 fn delete_snapshot<'a>(
1514 &'a mut self,
1515 request: DeleteSnapshotRequest,
1516 placement: ShardPlacement,
1517 ) -> GroupDeleteSnapshotFuture<'a> {
1518 Box::pin(async move {
1519 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1520 match self
1521 .state_machine
1522 .delete_snapshot(&request.stream_id, request.snapshot_offset)
1523 {
1524 StreamResponse::Error {
1525 code,
1526 message,
1527 next_offset,
1528 context,
1529 } => Err(GroupEngineError::stream_with_context(
1530 code,
1531 message,
1532 next_offset,
1533 context,
1534 )),
1535 other => Err(GroupEngineError::new(format!(
1536 "unexpected delete snapshot response: {other:?}"
1537 ))),
1538 }
1539 })
1540 }
1541
1542 fn bootstrap_stream<'a>(
1543 &'a mut self,
1544 request: BootstrapStreamRequest,
1545 placement: ShardPlacement,
1546 ) -> GroupBootstrapStreamFuture<'a> {
1547 Box::pin(async move {
1548 self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1549 let plan = self
1550 .state_machine
1551 .bootstrap_plan(&request.stream_id)
1552 .map_err(stream_response_error)?;
1553 let snapshot_offset = plan.snapshot.as_ref().map(|snapshot| snapshot.offset);
1554 let snapshot_content_type = plan
1555 .snapshot
1556 .as_ref()
1557 .map(|snapshot| snapshot.content_type.clone())
1558 .unwrap_or_else(|| DEFAULT_CONTENT_TYPE.to_owned());
1559 let snapshot_payload = plan
1560 .snapshot
1561 .as_ref()
1562 .map(|snapshot| snapshot.payload.clone())
1563 .unwrap_or_default();
1564 let updates = self
1565 .bootstrap_updates(
1566 &request.stream_id,
1567 &plan.updates,
1568 &plan.content_type,
1569 request.now_ms,
1570 )
1571 .await?;
1572 Ok(BootstrapStreamResponse {
1573 placement,
1574 snapshot_offset,
1575 snapshot_content_type,
1576 snapshot_payload,
1577 updates,
1578 next_offset: plan.next_offset,
1579 up_to_date: plan.up_to_date,
1580 closed: plan.closed,
1581 })
1582 })
1583 }
1584
1585 fn touch_stream_access<'a>(
1586 &'a mut self,
1587 stream_id: BucketStreamId,
1588 now_ms: u64,
1589 renew_ttl: bool,
1590 placement: ShardPlacement,
1591 ) -> GroupTouchStreamAccessFuture<'a> {
1592 Box::pin(async move { self.apply_access_command(stream_id, now_ms, renew_ttl, placement) })
1593 }
1594
1595 fn add_fork_ref<'a>(
1596 &'a mut self,
1597 stream_id: BucketStreamId,
1598 now_ms: u64,
1599 placement: ShardPlacement,
1600 ) -> GroupForkRefFuture<'a> {
1601 Box::pin(async move {
1602 match self.apply_committed_write(
1603 GroupWriteCommand::AddForkRef { stream_id, now_ms },
1604 placement,
1605 )? {
1606 GroupWriteResponse::AddForkRef(response) => Ok(response),
1607 other => Err(GroupEngineError::new(format!(
1608 "unexpected add fork ref write response: {other:?}"
1609 ))),
1610 }
1611 })
1612 }
1613
1614 fn release_fork_ref<'a>(
1615 &'a mut self,
1616 stream_id: BucketStreamId,
1617 placement: ShardPlacement,
1618 ) -> GroupForkRefFuture<'a> {
1619 Box::pin(async move {
1620 match self
1621 .apply_committed_write(GroupWriteCommand::ReleaseForkRef { stream_id }, placement)?
1622 {
1623 GroupWriteResponse::ReleaseForkRef(response) => Ok(response),
1624 other => Err(GroupEngineError::new(format!(
1625 "unexpected release fork ref write response: {other:?}"
1626 ))),
1627 }
1628 })
1629 }
1630
1631 fn head_stream<'a>(
1632 &'a mut self,
1633 request: HeadStreamRequest,
1634 placement: ShardPlacement,
1635 ) -> GroupHeadStreamFuture<'a> {
1636 Box::pin(async move {
1637 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1638 self.head_stream_after_access(&request, placement)
1639 })
1640 }
1641
1642 fn get_stream_attrs<'a>(
1643 &'a mut self,
1644 request: GetStreamAttrsRequest,
1645 placement: ShardPlacement,
1646 ) -> GroupGetStreamAttrsFuture<'a> {
1647 Box::pin(async move {
1648 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1649 self.get_stream_attrs_after_access(&request, placement)
1650 })
1651 }
1652
1653 fn update_stream_attrs<'a>(
1654 &'a mut self,
1655 request: UpdateStreamAttrsRequest,
1656 placement: ShardPlacement,
1657 ) -> GroupUpdateStreamAttrsFuture<'a> {
1658 Box::pin(async move {
1659 match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
1660 GroupWriteResponse::UpdateStreamAttrs(response) => Ok(response),
1661 other => Err(GroupEngineError::new(format!(
1662 "unexpected update stream attrs write response: {other:?}"
1663 ))),
1664 }
1665 })
1666 }
1667
1668 fn close_stream<'a>(
1669 &'a mut self,
1670 request: CloseStreamRequest,
1671 placement: ShardPlacement,
1672 ) -> GroupCloseStreamFuture<'a> {
1673 Box::pin(async move {
1674 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1675 let command = GroupWriteCommand::from(request);
1676 match self.apply_committed_write(command, placement)? {
1677 GroupWriteResponse::CloseStream(response) => Ok(response),
1678 other => Err(GroupEngineError::new(format!(
1679 "unexpected close stream write response: {other:?}"
1680 ))),
1681 }
1682 })
1683 }
1684
1685 fn delete_stream<'a>(
1686 &'a mut self,
1687 request: DeleteStreamRequest,
1688 placement: ShardPlacement,
1689 ) -> GroupDeleteStreamFuture<'a> {
1690 let command = GroupWriteCommand::from(request);
1691 Box::pin(async move {
1692 match self.apply_committed_write(command, placement)? {
1693 GroupWriteResponse::DeleteStream(response) => Ok(response),
1694 other => Err(GroupEngineError::new(format!(
1695 "unexpected delete stream write response: {other:?}"
1696 ))),
1697 }
1698 })
1699 }
1700
1701 fn ack_cold_gc<'a>(
1702 &'a mut self,
1703 up_to_seq: u64,
1704 placement: ShardPlacement,
1705 ) -> GroupAckColdGcFuture<'a> {
1706 Box::pin(async move {
1707 match self
1708 .apply_committed_write(GroupWriteCommand::AckColdGc { up_to_seq }, placement)?
1709 {
1710 GroupWriteResponse::AckColdGc(response) => Ok(response),
1711 other => Err(GroupEngineError::new(format!(
1712 "unexpected ack cold gc write response: {other:?}"
1713 ))),
1714 }
1715 })
1716 }
1717
1718 fn plan_cold_gc<'a>(
1719 &'a mut self,
1720 max: usize,
1721 _placement: ShardPlacement,
1722 ) -> GroupPlanColdGcFuture<'a> {
1723 let entries = self.state_machine.pending_cold_gc_batch(max);
1724 Box::pin(async move { Ok(entries) })
1725 }
1726
1727 fn append<'a>(
1728 &'a mut self,
1729 request: AppendRequest,
1730 placement: ShardPlacement,
1731 ) -> GroupAppendFuture<'a> {
1732 Box::pin(async move {
1733 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1734 let command = GroupWriteCommand::from(request);
1735 match self.apply_committed_write(command, placement)? {
1736 GroupWriteResponse::Append(response) => Ok(response),
1737 other => Err(GroupEngineError::new(format!(
1738 "unexpected append write response: {other:?}"
1739 ))),
1740 }
1741 })
1742 }
1743
1744 fn append_with_cold_admission<'a>(
1745 &'a mut self,
1746 request: AppendRequest,
1747 placement: ShardPlacement,
1748 admission: ColdWriteAdmission,
1749 ) -> GroupAppendFuture<'a> {
1750 if !admission.is_enabled() {
1751 return self.append(request, placement);
1752 }
1753 Box::pin(async move { self.append_with_admission_inner(request, placement, admission) })
1754 }
1755
1756 fn append_external<'a>(
1757 &'a mut self,
1758 request: AppendExternalRequest,
1759 placement: ShardPlacement,
1760 ) -> GroupAppendFuture<'a> {
1761 Box::pin(async move {
1762 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1763 if let Some(cold_store) = self.cold_store.as_ref() {
1764 let start_offset = self
1765 .state_machine
1766 .head(&request.stream_id)
1767 .map(|metadata| metadata.tail_offset)
1768 .ok_or_else(|| {
1769 GroupEngineError::stream(
1770 ursula_stream::StreamErrorCode::StreamNotFound,
1771 format!("stream '{}' does not exist", request.stream_id),
1772 )
1773 })?;
1774 let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1775 write_external_segment_index_pages(
1776 &store,
1777 &request.stream_id,
1778 start_offset,
1779 &request.payload,
1780 )
1781 .await
1782 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1783 }
1784 let command = GroupWriteCommand::from(request);
1785 match self.apply_committed_write(command, placement)? {
1786 GroupWriteResponse::Append(response) => Ok(response),
1787 other => Err(GroupEngineError::new(format!(
1788 "unexpected external append write response: {other:?}"
1789 ))),
1790 }
1791 })
1792 }
1793
1794 fn append_batch<'a>(
1795 &'a mut self,
1796 request: AppendBatchRequest,
1797 placement: ShardPlacement,
1798 ) -> GroupAppendBatchFuture<'a> {
1799 Box::pin(async move {
1800 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1801 let command = GroupWriteCommand::from(request);
1802 match self.apply_committed_write(command, placement)? {
1803 GroupWriteResponse::AppendBatch(response) => Ok(response),
1804 other => Err(GroupEngineError::new(format!(
1805 "unexpected append batch write response: {other:?}"
1806 ))),
1807 }
1808 })
1809 }
1810
1811 fn append_batch_with_cold_admission<'a>(
1812 &'a mut self,
1813 request: AppendBatchRequest,
1814 placement: ShardPlacement,
1815 admission: ColdWriteAdmission,
1816 ) -> GroupAppendBatchFuture<'a> {
1817 if !admission.is_enabled() {
1818 return self.append_batch(request, placement);
1819 }
1820 Box::pin(
1821 async move { self.append_batch_with_admission_inner(request, placement, admission) },
1822 )
1823 }
1824
1825 fn flush_cold<'a>(
1826 &'a mut self,
1827 request: FlushColdRequest,
1828 placement: ShardPlacement,
1829 ) -> GroupFlushColdFuture<'a> {
1830 Box::pin(async move {
1831 let mut index_rollback = None;
1832 if let Some(cold_store) = self.cold_store.as_ref() {
1833 let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1834 let rollback = write_cold_chunk_index_pages_with_rollback(
1835 &store,
1836 &request.stream_id,
1837 &request.chunk,
1838 )
1839 .await
1840 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1841 index_rollback = Some((store, rollback));
1842 }
1843 let command = GroupWriteCommand::from(request);
1844 match self.apply_committed_write(command, placement) {
1845 Ok(GroupWriteResponse::FlushCold(response)) => Ok(response),
1846 Ok(other) => {
1847 if let Some((store, rollback)) = index_rollback {
1848 rollback_cold_index_pages(&store, rollback)
1849 .await
1850 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1851 }
1852 Err(GroupEngineError::new(format!(
1853 "unexpected flush cold write response: {other:?}"
1854 )))
1855 }
1856 Err(err) => {
1857 if let Some((store, rollback)) = index_rollback {
1858 rollback_cold_index_pages(&store, rollback).await.map_err(
1859 |rollback_err| {
1860 GroupEngineError::new(format!(
1861 "rollback cold index after flush failure: {rollback_err}"
1862 ))
1863 },
1864 )?;
1865 }
1866 Err(err)
1867 }
1868 }
1869 })
1870 }
1871
1872 fn plan_cold_flush<'a>(
1873 &'a mut self,
1874 request: PlanColdFlushRequest,
1875 _placement: ShardPlacement,
1876 ) -> GroupPlanColdFlushFuture<'a> {
1877 Box::pin(async move {
1878 self.state_machine
1879 .plan_cold_flush(
1880 &request.stream_id,
1881 request.min_hot_bytes,
1882 request.max_flush_bytes,
1883 )
1884 .map_err(stream_response_error)
1885 })
1886 }
1887
1888 fn plan_next_cold_flush<'a>(
1889 &'a mut self,
1890 request: PlanGroupColdFlushRequest,
1891 _placement: ShardPlacement,
1892 ) -> GroupPlanNextColdFlushFuture<'a> {
1893 Box::pin(async move {
1894 self.state_machine
1895 .plan_next_cold_flush_batch(request.min_hot_bytes, request.max_flush_bytes, 1)
1896 .map(|candidates| candidates.into_iter().next())
1897 .map_err(stream_response_error)
1898 })
1899 }
1900
1901 fn plan_next_cold_flush_batch<'a>(
1902 &'a mut self,
1903 request: PlanGroupColdFlushRequest,
1904 _placement: ShardPlacement,
1905 max_candidates: usize,
1906 ) -> GroupPlanNextColdFlushBatchFuture<'a> {
1907 Box::pin(async move {
1908 self.state_machine
1909 .plan_next_cold_flush_batch(
1910 request.min_hot_bytes,
1911 request.max_flush_bytes,
1912 max_candidates,
1913 )
1914 .map_err(stream_response_error)
1915 })
1916 }
1917
1918 fn cold_hot_backlog<'a>(
1919 &'a mut self,
1920 stream_id: BucketStreamId,
1921 _placement: ShardPlacement,
1922 ) -> GroupColdHotBacklogFuture<'a> {
1923 Box::pin(async move { self.cold_hot_backlog_for(stream_id) })
1924 }
1925
1926 fn snapshot<'a>(&'a mut self, placement: ShardPlacement) -> GroupSnapshotFuture<'a> {
1927 Box::pin(async move { Ok(self.build_snapshot(placement)) })
1928 }
1929
1930 fn install_snapshot<'a>(
1931 &'a mut self,
1932 snapshot: GroupSnapshot,
1933 ) -> GroupInstallSnapshotFuture<'a> {
1934 Box::pin(async move { self.install_snapshot_inner(snapshot) })
1935 }
1936}
1937
1938#[derive(Debug, Clone, Default)]
1939pub struct InMemoryGroupEngineFactory {
1940 cold_store: Option<ColdStoreHandle>,
1941}
1942
1943impl InMemoryGroupEngineFactory {
1944 pub fn new() -> Self {
1945 Self::default()
1946 }
1947
1948 pub fn with_cold_store(cold_store: Option<ColdStoreHandle>) -> Self {
1949 Self { cold_store }
1950 }
1951}
1952
1953impl GroupEngineFactory for InMemoryGroupEngineFactory {
1954 fn create<'a>(
1955 &'a self,
1956 _placement: ShardPlacement,
1957 _metrics: GroupEngineMetrics,
1958 ) -> GroupEngineCreateFuture<'a> {
1959 Box::pin(async move {
1960 let mut engine = InMemoryGroupEngine::default();
1961 engine.set_cold_store(self.cold_store.clone());
1962 let engine: Box<dyn GroupEngine> = Box::new(engine);
1963 Ok(engine)
1964 })
1965 }
1966}
1967
1968pub(crate) fn compare_stream_ids(
1969 left: &BucketStreamId,
1970 right: &BucketStreamId,
1971) -> std::cmp::Ordering {
1972 left.bucket_id
1973 .cmp(&right.bucket_id)
1974 .then_with(|| left.stream_id.cmp(&right.stream_id))
1975}
1976pub(crate) fn ensure_bucket_exists(
1977 state_machine: &mut StreamStateMachine,
1978 stream_id: &BucketStreamId,
1979) -> Result<(), GroupEngineError> {
1980 if state_machine.bucket_exists(&stream_id.bucket_id) {
1981 return Ok(());
1982 }
1983
1984 match state_machine.apply(StreamCommand::CreateBucket {
1985 bucket_id: stream_id.bucket_id.clone(),
1986 }) {
1987 StreamResponse::BucketCreated { .. } | StreamResponse::BucketAlreadyExists { .. } => Ok(()),
1988 StreamResponse::Error {
1989 code,
1990 message,
1991 next_offset,
1992 context,
1993 } => Err(GroupEngineError::stream_with_context(
1994 code,
1995 message,
1996 next_offset,
1997 context,
1998 )),
1999 other => Err(GroupEngineError::new(format!(
2000 "unexpected create bucket response: {other:?}"
2001 ))),
2002 }
2003}
2004
2005pub(crate) fn stream_response_error(response: StreamResponse) -> GroupEngineError {
2006 match response {
2007 StreamResponse::Error {
2008 code,
2009 message,
2010 next_offset,
2011 context,
2012 } => GroupEngineError::stream_with_context(code, message, next_offset, context),
2013 other => GroupEngineError::new(format!("unexpected stream response error: {other:?}")),
2014 }
2015}
2016
2017pub(crate) fn restore_stream_append_counts(
2018 counts: Vec<StreamAppendCount>,
2019 snapshot_stream_ids: &HashSet<BucketStreamId>,
2020) -> Result<HashMap<BucketStreamId, u64>, GroupEngineError> {
2021 let mut restored = HashMap::with_capacity(counts.len());
2022 for count in counts {
2023 if !snapshot_stream_ids.contains(&count.stream_id) {
2024 return Err(GroupEngineError::new(format!(
2025 "append count references missing snapshot stream '{}'",
2026 count.stream_id
2027 )));
2028 }
2029 if restored
2030 .insert(count.stream_id.clone(), count.append_count)
2031 .is_some()
2032 {
2033 return Err(GroupEngineError::new(format!(
2034 "snapshot contains duplicate append count for stream '{}'",
2035 count.stream_id
2036 )));
2037 }
2038 }
2039 Ok(restored)
2040}