1use std::collections::HashMap;
2use std::collections::HashSet;
3use std::sync::Arc;
4
5use bytes::Bytes;
6use ursula_shard::BucketStreamId;
7use ursula_shard::ShardPlacement;
8use ursula_stream::AppendStreamInput;
9use ursula_stream::ProducerRequest;
10use ursula_stream::StreamCommand;
11use ursula_stream::StreamErrorCode;
12use ursula_stream::StreamMessageRecord;
13use ursula_stream::StreamReadPlan;
14use ursula_stream::StreamReadSegment;
15use ursula_stream::StreamResponse;
16use ursula_stream::StreamSnapshot;
17use ursula_stream::StreamStateMachine;
18
19use super::GroupAckColdGcFuture;
20use super::GroupAppendBatchFuture;
21use super::GroupAppendBatchResponse;
22use super::GroupAppendFuture;
23use super::GroupBootstrapStreamFuture;
24use super::GroupCloseStreamFuture;
25use super::GroupColdHotBacklogFuture;
26use super::GroupCreateStreamFuture;
27use super::GroupDeleteSnapshotFuture;
28use super::GroupDeleteStreamFuture;
29use super::GroupEngine;
30use super::GroupEngineCreateFuture;
31use super::GroupEngineError;
32use super::GroupEngineFactory;
33use super::GroupEngineMetrics;
34use super::GroupFlushColdFuture;
35use super::GroupGetStreamAttrsFuture;
36use super::GroupHeadStreamFuture;
37use super::GroupInstallSnapshotFuture;
38use super::GroupPlanColdFlushFuture;
39use super::GroupPlanColdGcFuture;
40use super::GroupPlanNextColdFlushBatchFuture;
41use super::GroupPublishSnapshotFuture;
42use super::GroupReadSnapshotFuture;
43use super::GroupReadStreamFuture;
44use super::GroupReadStreamPartsFuture;
45use super::GroupSnapshotFuture;
46use super::GroupTouchStreamAccessFuture;
47use super::GroupUpdateStreamAttrsFuture;
48use super::GroupWriteResponse;
49use crate::cold_index::ColdIndexPageCache;
50use crate::cold_index::ColdStoreColdIndexPageStore;
51use crate::cold_index::rollback_cold_index_pages;
52use crate::cold_index::write_cold_chunk_index_pages_with_rollback;
53use crate::cold_index::write_external_segment_index_pages;
54use crate::cold_store::ColdStoreHandle;
55use crate::cold_store::DEFAULT_CONTENT_TYPE;
56use crate::command::GroupSnapshot;
57use crate::command::GroupWriteCommand;
58use crate::request::AckColdGcResponse;
59use crate::request::AppendBatchRequest;
60use crate::request::AppendExternalRequest;
61use crate::request::AppendRequest;
62use crate::request::AppendResponse;
63use crate::request::BootstrapStreamRequest;
64use crate::request::BootstrapStreamResponse;
65use crate::request::BootstrapUpdate;
66use crate::request::CloseStreamRequest;
67use crate::request::CloseStreamResponse;
68use crate::request::ColdHotBacklog;
69use crate::request::ColdWriteAdmission;
70use crate::request::CreateStreamExternalRequest;
71use crate::request::CreateStreamRequest;
72use crate::request::CreateStreamResponse;
73use crate::request::DeleteSnapshotRequest;
74use crate::request::DeleteStreamRequest;
75use crate::request::DeleteStreamResponse;
76use crate::request::FlushColdRequest;
77use crate::request::FlushColdResponse;
78use crate::request::GetStreamAttrsRequest;
79use crate::request::GetStreamAttrsResponse;
80use crate::request::GroupReadStreamParts;
81use crate::request::HeadStreamRequest;
82use crate::request::HeadStreamResponse;
83use crate::request::PlanColdFlushRequest;
84use crate::request::PlanGroupColdFlushRequest;
85use crate::request::PublishSnapshotRequest;
86use crate::request::PublishSnapshotResponse;
87use crate::request::ReadSnapshotRequest;
88use crate::request::ReadSnapshotResponse;
89use crate::request::ReadStreamRequest;
90use crate::request::StreamAppendCount;
91use crate::request::TouchStreamAccessResponse;
92use crate::request::UpdateStreamAttrsRequest;
93use crate::request::UpdateStreamAttrsResponse;
94
95pub(crate) struct AppendPayloadInput<'a> {
96 stream_id: BucketStreamId,
97 content_type: Option<&'a str>,
98 payload: &'a [u8],
99 close_after: bool,
100 stream_seq: Option<String>,
101 producer: Option<ProducerRequest>,
102 now_ms: u64,
103 record_match: Option<u64>,
104}
105
106#[derive(Debug, Clone, Default)]
107pub struct InMemoryGroupEngine {
108 pub(crate) commit_index: u64,
109 pub(crate) state_machine: StreamStateMachine,
110 pub(crate) stream_append_counts: HashMap<BucketStreamId, u64>,
111 pub(crate) cold_store: Option<ColdStoreHandle>,
112 pub(crate) cold_index_cache: Option<Arc<ColdIndexPageCache<ColdStoreColdIndexPageStore>>>,
113}
114
115impl InMemoryGroupEngine {
116 pub fn with_cold_store(cold_store: ColdStoreHandle) -> Self {
117 let mut engine = Self::default();
118 engine.set_cold_store(Some(cold_store));
119 engine
120 }
121
122 pub fn cold_store(&self) -> Option<ColdStoreHandle> {
123 self.cold_store.clone()
124 }
125
126 pub(crate) fn set_cold_store(&mut self, cold_store: Option<ColdStoreHandle>) {
127 self.cold_index_cache = cold_store.as_ref().map(|cold_store| {
128 Arc::new(ColdIndexPageCache::new(
129 Arc::new(ColdStoreColdIndexPageStore::new(cold_store.clone())),
130 1024,
131 ))
132 });
133 self.cold_store = cold_store;
134 }
135
136 pub fn apply_committed_write(
137 &mut self,
138 command: GroupWriteCommand,
139 placement: ShardPlacement,
140 ) -> Result<GroupWriteResponse, GroupEngineError> {
141 match command {
142 GroupWriteCommand::Stream(command) => self.apply_stream_command(command, placement),
143 GroupWriteCommand::Batch { commands } => Ok(GroupWriteResponse::Batch(
144 commands
145 .into_iter()
146 .map(|command| self.apply_stream_command(command, placement))
147 .collect(),
148 )),
149 }
150 }
151
152 pub fn apply_stream_command(
156 &mut self,
157 command: StreamCommand,
158 placement: ShardPlacement,
159 ) -> Result<GroupWriteResponse, GroupEngineError> {
160 match command {
161 StreamCommand::Append {
164 stream_id,
165 content_type,
166 payload,
167 close_after,
168 stream_seq,
169 producer,
170 now_ms,
171 record_match,
172 } => self
173 .append_payload(
174 AppendPayloadInput {
175 stream_id,
176 content_type: content_type.as_deref(),
177 payload: &payload,
178 close_after,
179 stream_seq,
180 producer,
181 now_ms,
182 record_match,
183 },
184 placement,
185 )
186 .map(GroupWriteResponse::Append),
187 StreamCommand::AppendBatch {
188 stream_id,
189 content_type,
190 payloads,
191 producer,
192 now_ms,
193 } => self.apply_append_batch(
194 stream_id,
195 content_type,
196 payloads,
197 producer,
198 now_ms,
199 placement,
200 ),
201 command => {
202 let stream_id = command_stream_id(&command);
203 let command_producer = command_producer(&command);
204 if let StreamCommand::CreateStream { stream_id, .. }
205 | StreamCommand::CreateExternal { stream_id, .. } = &command
206 {
207 ensure_bucket_exists(&mut self.state_machine, stream_id)?;
208 }
209 let response = self.state_machine.apply(command);
210 self.group_response_from_stream(response, stream_id, command_producer, placement)
211 }
212 }
213 }
214
215 fn apply_append_batch(
216 &mut self,
217 stream_id: BucketStreamId,
218 content_type: Option<String>,
219 payloads: Vec<Bytes>,
220 producer: Option<ProducerRequest>,
221 now_ms: u64,
222 placement: ShardPlacement,
223 ) -> Result<GroupWriteResponse, GroupEngineError> {
224 if let Some(producer) = producer {
225 let payload_refs = payloads.iter().map(Bytes::as_ref).collect::<Vec<_>>();
226 let batch = self
227 .state_machine
228 .append_batch_borrowed(
229 stream_id.clone(),
230 content_type.as_deref(),
231 &payload_refs,
232 Some(producer.clone()),
233 now_ms,
234 )
235 .map_err(stream_response_error)?;
236 let old_commit_index = self.commit_index;
237 let old_append_count = *self.stream_append_counts.get(&stream_id).unwrap_or(&0);
238 if !batch.deduplicated {
239 let count = u64::try_from(batch.items.len()).expect("item count fits u64");
240 self.commit_index += count;
241 *self
242 .stream_append_counts
243 .entry(stream_id.clone())
244 .or_insert(0) += count;
245 }
246 let items = batch
247 .items
248 .into_iter()
249 .enumerate()
250 .map(|(index, item)| {
251 let item_index = u64::try_from(index + 1).expect("item index fits u64");
252 Ok(AppendResponse {
253 placement,
254 start_offset: item.offset,
255 next_offset: item.next_offset,
256 stream_append_count: if item.deduplicated {
257 old_append_count
258 } else {
259 old_append_count + item_index
260 },
261 group_commit_index: if item.deduplicated {
262 old_commit_index
263 } else {
264 old_commit_index + item_index
265 },
266 closed: item.closed,
267 deduplicated: item.deduplicated,
268 producer: None,
269 record_range: self
270 .state_machine
271 .record_range_for_append(
272 &stream_id,
273 item.offset,
274 item.next_offset,
275 Some(&producer),
276 )
277 .map_err(|err| {
278 GroupEngineError::new(format!("record range: {err:?}"))
279 })?,
280 })
281 })
282 .collect();
283 return Ok(GroupWriteResponse::AppendBatch(GroupAppendBatchResponse {
284 placement,
285 items,
286 }));
287 }
288
289 let mut items = Vec::with_capacity(payloads.len());
290 for payload in payloads {
291 if payload.is_empty() {
292 items.push(Err(GroupEngineError::stream(
293 StreamErrorCode::EmptyAppend,
294 "append payload must be non-empty",
295 )));
296 continue;
297 }
298 items.push(self.append_payload(
299 AppendPayloadInput {
300 stream_id: stream_id.clone(),
301 content_type: content_type.as_deref(),
302 payload: &payload,
303 close_after: false,
304 stream_seq: None,
305 producer: None,
306 now_ms,
307 record_match: None,
308 },
309 placement,
310 ));
311 }
312 Ok(GroupWriteResponse::AppendBatch(GroupAppendBatchResponse {
313 placement,
314 items,
315 }))
316 }
317
318 fn group_response_from_stream(
321 &mut self,
322 response: StreamResponse,
323 stream_id: Option<BucketStreamId>,
324 command_producer: Option<ProducerRequest>,
325 placement: ShardPlacement,
326 ) -> Result<GroupWriteResponse, GroupEngineError> {
327 match response {
328 StreamResponse::Created {
329 next_offset,
330 closed,
331 ..
332 } => {
333 let stream_id = require_response_stream_id(stream_id, "created")?;
334 self.commit_index += 1;
335 Ok(GroupWriteResponse::CreateStream(CreateStreamResponse {
336 placement,
337 next_offset,
338 closed,
339 already_exists: false,
340 group_commit_index: self.commit_index,
341 record_range: self
342 .state_machine
343 .record_range(&stream_id)
344 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?,
345 }))
346 }
347 StreamResponse::AlreadyExists {
348 next_offset,
349 closed,
350 ..
351 } => Ok(GroupWriteResponse::CreateStream(CreateStreamResponse {
352 placement,
353 next_offset,
354 closed,
355 already_exists: true,
356 group_commit_index: self.commit_index,
357 record_range: None,
358 })),
359 StreamResponse::Appended {
360 offset,
361 next_offset,
362 closed,
363 deduplicated,
364 producer,
365 } => {
366 let stream_id = require_response_stream_id(stream_id, "appended")?;
367 let record_range = self
368 .state_machine
369 .record_range_for_append(&stream_id, offset, next_offset, producer.as_ref())
370 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?;
371 let stream_append_count = self.stream_append_counts.entry(stream_id).or_insert(0);
372 if !deduplicated {
373 self.commit_index += 1;
374 *stream_append_count += 1;
375 }
376 Ok(GroupWriteResponse::Append(AppendResponse {
377 placement,
378 start_offset: offset,
379 next_offset,
380 stream_append_count: *stream_append_count,
381 group_commit_index: self.commit_index,
382 closed,
383 deduplicated,
384 producer,
385 record_range,
386 }))
387 }
388 StreamResponse::SnapshotPublished {
389 snapshot_offset,
390 record_range,
391 } => {
392 self.commit_index += 1;
393 Ok(GroupWriteResponse::PublishSnapshot(
394 PublishSnapshotResponse {
395 placement,
396 snapshot_offset,
397 group_commit_index: self.commit_index,
398 record_range,
399 },
400 ))
401 }
402 StreamResponse::Accessed { changed, expired } => {
403 if changed || expired {
404 self.commit_index += 1;
405 }
406 Ok(GroupWriteResponse::TouchStreamAccess(
407 TouchStreamAccessResponse {
408 placement,
409 changed,
410 expired,
411 group_commit_index: self.commit_index,
412 },
413 ))
414 }
415 StreamResponse::AttrsUpdated { changed } => {
416 if changed {
417 self.commit_index += 1;
418 }
419 Ok(GroupWriteResponse::UpdateStreamAttrs(
420 UpdateStreamAttrsResponse {
421 placement,
422 changed,
423 group_commit_index: self.commit_index,
424 },
425 ))
426 }
427 StreamResponse::ColdFlushed { hot_start_offset } => {
428 self.commit_index += 1;
429 Ok(GroupWriteResponse::FlushCold(FlushColdResponse {
430 placement,
431 hot_start_offset,
432 group_commit_index: self.commit_index,
433 }))
434 }
435 StreamResponse::Closed {
436 next_offset,
437 deduplicated,
438 ..
439 } => {
440 let stream_id = require_response_stream_id(stream_id, "closed")?;
441 let record_range = self
442 .state_machine
443 .record_range_for_append(
444 &stream_id,
445 next_offset,
446 next_offset,
447 command_producer.as_ref(),
448 )
449 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?;
450 if !deduplicated {
451 self.commit_index += 1;
452 }
453 Ok(GroupWriteResponse::CloseStream(CloseStreamResponse {
454 placement,
455 next_offset,
456 group_commit_index: self.commit_index,
457 deduplicated,
458 record_range,
459 }))
460 }
461 StreamResponse::Deleted => {
462 let stream_id = require_response_stream_id(stream_id, "deleted")?;
463 self.commit_index += 1;
464 self.stream_append_counts.remove(&stream_id);
467 Ok(GroupWriteResponse::DeleteStream(DeleteStreamResponse {
468 placement,
469 group_commit_index: self.commit_index,
470 }))
471 }
472 StreamResponse::ColdGcAcked { removed } => {
473 self.commit_index += 1;
474 Ok(GroupWriteResponse::AckColdGc(AckColdGcResponse {
475 placement,
476 removed,
477 group_commit_index: self.commit_index,
478 }))
479 }
480 StreamResponse::Error {
481 code,
482 message,
483 next_offset,
484 context,
485 } => Err(GroupEngineError::stream_with_context(
486 code,
487 message,
488 next_offset,
489 context,
490 )),
491 other @ (StreamResponse::BucketCreated { .. }
492 | StreamResponse::BucketAlreadyExists { .. }
493 | StreamResponse::BucketDeleted { .. }) => Err(GroupEngineError::new(format!(
494 "unexpected group write response: {other:?}"
495 ))),
496 }
497 }
498
499 pub(crate) fn cold_hot_backlog_for(
500 &self,
501 stream_id: BucketStreamId,
502 ) -> Result<ColdHotBacklog, GroupEngineError> {
503 let stream_hot_bytes = self.state_machine.hot_payload_len(&stream_id).unwrap_or(0);
504 Ok(ColdHotBacklog {
505 stream_id,
506 stream_hot_bytes,
507 group_hot_bytes: self.state_machine.total_hot_payload_bytes(),
508 })
509 }
510
511 pub fn check_cold_write_admission_bytes(
512 &self,
513 stream_id: &BucketStreamId,
514 admission: ColdWriteAdmission,
515 incoming_bytes: u64,
516 ) -> Result<(), GroupEngineError> {
517 let Some(limit) = admission.max_hot_bytes_per_group else {
518 return Ok(());
519 };
520 if incoming_bytes == 0 {
521 return Ok(());
522 }
523 let before = self.state_machine.total_hot_payload_bytes();
524 let after = before.saturating_add(incoming_bytes);
525 if after <= limit {
526 return Ok(());
527 }
528 Err(GroupEngineError::cold_backpressure(
529 stream_id.clone(),
530 before,
531 after,
532 limit,
533 ))
534 }
535
536 pub(crate) fn create_stream_with_admission_inner(
537 &mut self,
538 request: CreateStreamRequest,
539 placement: ShardPlacement,
540 admission: ColdWriteAdmission,
541 ) -> Result<CreateStreamResponse, GroupEngineError> {
542 let stream_id = request.stream_id.clone();
543 if admission.is_enabled() {
544 let mut preview = self.clone();
545 let preview_response = match preview
546 .apply_committed_write(GroupWriteCommand::from(request.clone()), placement)?
547 {
548 GroupWriteResponse::CreateStream(response) => response,
549 other => {
550 return Err(GroupEngineError::new(format!(
551 "unexpected create stream preview response: {other:?}"
552 )));
553 }
554 };
555 if !preview_response.already_exists {
556 self.check_cold_write_admission_bytes(
557 &stream_id,
558 admission,
559 u64::try_from(request.initial_payload.len()).expect("payload len fits u64"),
560 )?;
561 }
562 }
563 let response =
564 match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
565 GroupWriteResponse::CreateStream(response) => response,
566 other => {
567 return Err(GroupEngineError::new(format!(
568 "unexpected create stream write response: {other:?}"
569 )));
570 }
571 };
572 Ok(response)
573 }
574
575 pub(crate) fn append_with_admission_inner(
576 &mut self,
577 request: AppendRequest,
578 placement: ShardPlacement,
579 admission: ColdWriteAdmission,
580 ) -> Result<AppendResponse, GroupEngineError> {
581 let stream_id = request.stream_id.clone();
582 if admission.is_enabled() {
583 let mut preview = self.clone();
584 let preview_response = match preview
585 .apply_committed_write(GroupWriteCommand::from(request.clone()), placement)?
586 {
587 GroupWriteResponse::Append(response) => response,
588 other => {
589 return Err(GroupEngineError::new(format!(
590 "unexpected append preview response: {other:?}"
591 )));
592 }
593 };
594 if !preview_response.deduplicated {
595 self.check_cold_write_admission_bytes(
596 &stream_id,
597 admission,
598 u64::try_from(request.payload.len()).expect("payload len fits u64"),
599 )?;
600 }
601 }
602 let response =
603 match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
604 GroupWriteResponse::Append(response) => response,
605 other => {
606 return Err(GroupEngineError::new(format!(
607 "unexpected append write response: {other:?}"
608 )));
609 }
610 };
611 Ok(response)
612 }
613
614 pub(crate) fn append_batch_with_admission_inner(
615 &mut self,
616 request: AppendBatchRequest,
617 placement: ShardPlacement,
618 admission: ColdWriteAdmission,
619 ) -> Result<GroupAppendBatchResponse, GroupEngineError> {
620 let stream_id = request.stream_id.clone();
621 let incoming_bytes = request
622 .payloads
623 .iter()
624 .map(|payload| u64::try_from(payload.len()).expect("payload len fits u64"))
625 .sum();
626 if admission.is_enabled() {
627 let mut preview = self.clone();
628 let preview_response = match preview
629 .apply_committed_write(GroupWriteCommand::from(request.clone()), placement)?
630 {
631 GroupWriteResponse::AppendBatch(response) => response,
632 other => {
633 return Err(GroupEngineError::new(format!(
634 "unexpected append batch preview response: {other:?}"
635 )));
636 }
637 };
638 let mutates = preview_response
639 .items
640 .iter()
641 .any(|item| matches!(item, Ok(response) if !response.deduplicated));
642 if mutates {
643 self.check_cold_write_admission_bytes(&stream_id, admission, incoming_bytes)?;
644 }
645 }
646 let response =
647 match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
648 GroupWriteResponse::AppendBatch(response) => response,
649 other => {
650 return Err(GroupEngineError::new(format!(
651 "unexpected append batch write response: {other:?}"
652 )));
653 }
654 };
655 Ok(response)
656 }
657
658 pub fn access_requires_write(
659 &self,
660 stream_id: &BucketStreamId,
661 now_ms: u64,
662 renew_ttl: bool,
663 ) -> Result<bool, GroupEngineError> {
664 self.state_machine
665 .access_requires_write(stream_id, now_ms, renew_ttl)
666 .map_err(stream_response_error)
667 }
668
669 pub(crate) fn apply_access_command(
670 &mut self,
671 stream_id: BucketStreamId,
672 now_ms: u64,
673 renew_ttl: bool,
674 placement: ShardPlacement,
675 ) -> Result<TouchStreamAccessResponse, GroupEngineError> {
676 match self.apply_committed_write(
677 GroupWriteCommand::Stream(StreamCommand::TouchStreamAccess {
678 stream_id,
679 now_ms,
680 renew_ttl,
681 }),
682 placement,
683 )? {
684 GroupWriteResponse::TouchStreamAccess(response) => Ok(response),
685 other => Err(GroupEngineError::new(format!(
686 "unexpected touch stream access write response: {other:?}"
687 ))),
688 }
689 }
690
691 pub(crate) fn ensure_stream_access(
692 &mut self,
693 stream_id: &BucketStreamId,
694 now_ms: u64,
695 renew_ttl: bool,
696 placement: ShardPlacement,
697 ) -> Result<Option<TouchStreamAccessResponse>, GroupEngineError> {
698 if !self.access_requires_write(stream_id, now_ms, renew_ttl)? {
699 return Ok(None);
700 }
701 let response =
702 self.apply_access_command(stream_id.clone(), now_ms, renew_ttl, placement)?;
703 if response.expired {
704 return Err(GroupEngineError::stream(
705 StreamErrorCode::StreamNotFound,
706 format!("stream '{stream_id}' does not exist"),
707 ));
708 }
709 Ok(Some(response))
710 }
711
712 pub(crate) fn append_payload(
713 &mut self,
714 input: AppendPayloadInput<'_>,
715 placement: ShardPlacement,
716 ) -> Result<AppendResponse, GroupEngineError> {
717 let AppendPayloadInput {
718 stream_id,
719 content_type,
720 payload,
721 close_after,
722 stream_seq,
723 producer,
724 now_ms,
725 record_match,
726 } = input;
727 let stream_count_key = stream_id.clone();
728 let response = self.state_machine.append_borrowed(AppendStreamInput {
729 stream_id,
730 content_type,
731 payload,
732 close_after,
733 stream_seq,
734 producer,
735 now_ms,
736 record_match,
737 });
738 match response {
739 StreamResponse::Appended {
740 offset,
741 next_offset,
742 closed,
743 deduplicated,
744 producer,
745 ..
746 } => {
747 let stream_append_count = self
748 .stream_append_counts
749 .entry(stream_count_key.clone())
750 .or_insert(0);
751 let record_range = self
752 .state_machine
753 .record_range_for_append(
754 &stream_count_key,
755 offset,
756 next_offset,
757 producer.as_ref(),
758 )
759 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?;
760 if !deduplicated {
761 self.commit_index += 1;
762 *stream_append_count += 1;
763 }
764 Ok(AppendResponse {
765 placement,
766 start_offset: offset,
767 next_offset,
768 stream_append_count: *stream_append_count,
769 group_commit_index: self.commit_index,
770 closed,
771 deduplicated,
772 producer,
773 record_range,
774 })
775 }
776 StreamResponse::Error {
777 code,
778 message,
779 next_offset,
780 context,
781 } => Err(GroupEngineError::stream_with_context(
782 code,
783 message,
784 next_offset,
785 context,
786 )),
787 other => Err(GroupEngineError::new(format!(
788 "unexpected append response: {other:?}"
789 ))),
790 }
791 }
792
793 pub fn read_stream_plan(
794 &mut self,
795 request: &ReadStreamRequest,
796 placement: ShardPlacement,
797 ) -> Result<StreamReadPlan, GroupEngineError> {
798 self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
799 self.read_stream_plan_after_access(request)
800 }
801
802 pub fn read_stream_plan_after_access(
803 &self,
804 request: &ReadStreamRequest,
805 ) -> Result<StreamReadPlan, GroupEngineError> {
806 let Some(record) = request.record else {
807 let mut plan = self
808 .state_machine
809 .read_plan_at(
810 &request.stream_id,
811 request.offset,
812 request.max_len,
813 request.now_ms,
814 )
815 .map_err(stream_response_error)?;
816 plan.retained_record_range = self
817 .state_machine
818 .record_range(&request.stream_id)
819 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?;
820 return Ok(plan);
821 };
822 let retained_record_range = self
823 .state_machine
824 .record_range(&request.stream_id)
825 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?
826 .ok_or_else(|| {
827 GroupEngineError::stream(
828 StreamErrorCode::InvalidRecordBoundaries,
829 "record coordinates are inactive for this stream",
830 )
831 })?;
832 if record < retained_record_range.first_record {
833 return Err(GroupEngineError::stream(
834 StreamErrorCode::StreamGone,
835 format!(
836 "record {record} is older than first retained record {}",
837 retained_record_range.first_record
838 ),
839 ));
840 }
841 if record > retained_record_range.next_record {
842 return Err(GroupEngineError::stream(
843 StreamErrorCode::InvalidRecordBoundaries,
844 format!(
845 "record {record} is beyond record tail {}",
846 retained_record_range.next_record
847 ),
848 ));
849 }
850 let next_record = request
851 .max_records
852 .map(|limit| record.saturating_add(limit))
853 .unwrap_or(retained_record_range.next_record)
854 .min(retained_record_range.next_record);
855 let offset = self
856 .state_machine
857 .offset_for_record(&request.stream_id, record)
858 .map_err(|err| GroupEngineError::new(format!("record offset: {err:?}")))?
859 .ok_or_else(|| GroupEngineError::new("record stream disappeared"))?;
860 let next_offset = self
861 .state_machine
862 .offset_for_record(&request.stream_id, next_record)
863 .map_err(|err| GroupEngineError::new(format!("record offset: {err:?}")))?
864 .ok_or_else(|| GroupEngineError::new("record stream disappeared"))?;
865 let max_len = usize::try_from(next_offset.saturating_sub(offset))
866 .map_err(|_| GroupEngineError::new("record read window exceeds usize"))?;
867 let mut plan = self
868 .state_machine
869 .read_plan_at(&request.stream_id, offset, max_len, request.now_ms)
870 .map_err(stream_response_error)?;
871 plan.retained_record_range = Some(retained_record_range);
872 plan.record_range = Some(ursula_stream::StreamRecordRange {
873 first_record: record,
874 next_record,
875 });
876 Ok(plan)
877 }
878
879 pub fn head_stream_after_access(
880 &mut self,
881 request: &HeadStreamRequest,
882 placement: ShardPlacement,
883 ) -> Result<HeadStreamResponse, GroupEngineError> {
884 let Some(metadata) = self
885 .state_machine
886 .head_at(&request.stream_id, request.now_ms)
887 else {
888 return Err(GroupEngineError::stream(
889 StreamErrorCode::StreamNotFound,
890 format!("stream '{}' does not exist", request.stream_id),
891 ));
892 };
893 let content_type = metadata.content_type.clone();
894 let tail_offset = metadata.tail_offset;
895 let closed = metadata.status == ursula_stream::StreamStatus::Closed;
896 let stream_ttl_seconds = metadata.stream_ttl_seconds;
897 let stream_expires_at_ms = metadata.stream_expires_at_ms;
898 let _ = metadata;
899 Ok(HeadStreamResponse {
900 placement,
901 content_type,
902 tail_offset,
903 cold_hot_start_offset: self.state_machine.hot_start_offset(&request.stream_id),
904 closed,
905 stream_ttl_seconds,
906 stream_expires_at_ms,
907 snapshot_offset: self
908 .state_machine
909 .latest_snapshot(&request.stream_id)
910 .map_err(stream_response_error)?
911 .map(|snapshot| snapshot.offset),
912 integrity: self
913 .state_machine
914 .integrity_snapshot(&request.stream_id)
915 .map_err(stream_response_error)?,
916 record_range: self
917 .state_machine
918 .record_range(&request.stream_id)
919 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?,
920 })
921 }
922
923 pub fn get_stream_attrs_after_access(
924 &mut self,
925 request: &GetStreamAttrsRequest,
926 placement: ShardPlacement,
927 ) -> Result<GetStreamAttrsResponse, GroupEngineError> {
928 if self
929 .state_machine
930 .head_at(&request.stream_id, request.now_ms)
931 .is_none()
932 {
933 return Err(GroupEngineError::stream(
934 StreamErrorCode::StreamNotFound,
935 format!("stream '{}' does not exist", request.stream_id),
936 ));
937 }
938 Ok(GetStreamAttrsResponse {
939 placement,
940 attrs: self.state_machine.stream_attrs(&request.stream_id).cloned(),
941 })
942 }
943
944 pub async fn read_payload_from_plan(
945 cold_store: Option<&ColdStoreHandle>,
946 cold_index_cache: Option<&Arc<ColdIndexPageCache<ColdStoreColdIndexPageStore>>>,
947 stream_id: &BucketStreamId,
948 plan: &StreamReadPlan,
949 ) -> Result<Vec<u8>, GroupEngineError> {
950 let mut payload = Vec::new();
951 for segment in &plan.segments {
952 match segment {
953 StreamReadSegment::Hot(bytes) => payload.extend_from_slice(bytes),
954 StreamReadSegment::ColdIndex(segment) => {
955 let Some(cold_store) = cold_store else {
956 return Err(GroupEngineError::stream_with_next_offset(
957 StreamErrorCode::InvalidColdFlush,
958 format!("stream '{stream_id}' read requires object payload store"),
959 Some(plan.next_offset),
960 ));
961 };
962 let Some(cache) = cold_index_cache else {
963 return Err(GroupEngineError::stream_with_next_offset(
964 StreamErrorCode::InvalidColdFlush,
965 format!("stream '{stream_id}' read requires cold index page cache"),
966 Some(plan.next_offset),
967 ));
968 };
969 let objects = cache
970 .object_segments_for_read(stream_id, segment)
971 .await
972 .map_err(|err| GroupEngineError::new(err.to_string()))?;
973 let segment_end = segment
974 .read_start_offset
975 .saturating_add(u64::try_from(segment.len).expect("read len fits u64"));
976 for object in objects {
977 let start = object.start_offset.max(segment.read_start_offset);
978 let end = object.end_offset.min(segment_end);
979 if start >= end {
980 continue;
981 }
982 let bytes = cold_store
983 .read_object_range_for_stream(
984 stream_id,
985 &object,
986 start,
987 usize::try_from(end - start).expect("object read len fits usize"),
988 )
989 .await
990 .map_err(|err| GroupEngineError::new(err.to_string()))?;
991 payload.extend_from_slice(&bytes);
992 }
993 }
994 StreamReadSegment::Object(segment) => {
995 let Some(cold_store) = cold_store else {
996 return Err(GroupEngineError::stream_with_next_offset(
997 StreamErrorCode::InvalidColdFlush,
998 format!("stream '{stream_id}' read requires object payload store"),
999 Some(plan.next_offset),
1000 ));
1001 };
1002 let bytes = cold_store
1003 .read_object_range_for_stream(
1004 stream_id,
1005 &segment.object,
1006 segment.read_start_offset,
1007 segment.len,
1008 )
1009 .await
1010 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1011 payload.extend_from_slice(&bytes);
1012 }
1013 }
1014 }
1015 Ok(payload)
1016 }
1017
1018 pub(crate) async fn read_own_payload_from_plan(
1019 &self,
1020 stream_id: &BucketStreamId,
1021 plan: &StreamReadPlan,
1022 ) -> Result<Vec<u8>, GroupEngineError> {
1023 Self::read_payload_from_plan(
1024 self.cold_store.as_ref(),
1025 self.cold_index_cache.as_ref(),
1026 stream_id,
1027 plan,
1028 )
1029 .await
1030 }
1031
1032 pub(crate) async fn bootstrap_updates(
1033 &self,
1034 stream_id: &BucketStreamId,
1035 records: &[StreamMessageRecord],
1036 content_type: &str,
1037 now_ms: u64,
1038 ) -> Result<Vec<BootstrapUpdate>, GroupEngineError> {
1039 let mut updates = Vec::with_capacity(records.len());
1040 for record in records {
1041 let len = usize::try_from(record.end_offset - record.start_offset).map_err(|_| {
1042 GroupEngineError::stream(
1043 StreamErrorCode::InvalidSnapshot,
1044 format!(
1045 "bootstrap message [{}..{}) for stream '{stream_id}' is too large",
1046 record.start_offset, record.end_offset
1047 ),
1048 )
1049 })?;
1050 let plan = self
1051 .state_machine
1052 .read_plan_at(stream_id, record.start_offset, len, now_ms)
1053 .map_err(stream_response_error)?;
1054 let payload = self.read_own_payload_from_plan(stream_id, &plan).await?;
1055 updates.push(BootstrapUpdate {
1056 start_offset: record.start_offset,
1057 next_offset: record.end_offset,
1058 content_type: content_type.to_owned(),
1059 payload,
1060 });
1061 }
1062 Ok(updates)
1063 }
1064
1065 pub(crate) fn build_snapshot(&self, placement: ShardPlacement) -> GroupSnapshot {
1066 let stream_snapshot = self.state_machine.snapshot();
1067 let stream_append_counts = self.stream_append_counts_snapshot(&stream_snapshot);
1068 GroupSnapshot {
1069 placement,
1070 group_commit_index: self.commit_index,
1071 stream_snapshot,
1072 stream_append_counts,
1073 }
1074 }
1075
1076 pub(crate) fn stream_append_counts_snapshot(
1077 &self,
1078 stream_snapshot: &ursula_stream::StreamSnapshot,
1079 ) -> Vec<StreamAppendCount> {
1080 let live: HashSet<&BucketStreamId> = stream_snapshot
1087 .streams
1088 .iter()
1089 .map(|entry| &entry.metadata.stream_id)
1090 .collect();
1091 let mut counts = self
1092 .stream_append_counts
1093 .iter()
1094 .filter(|(stream_id, _)| live.contains(stream_id))
1095 .map(|(stream_id, append_count)| StreamAppendCount {
1096 stream_id: stream_id.clone(),
1097 append_count: *append_count,
1098 })
1099 .collect::<Vec<_>>();
1100 counts.sort_by(|left, right| compare_stream_ids(&left.stream_id, &right.stream_id));
1101 counts
1102 }
1103
1104 pub fn stream_tail_offset(&self, stream_id: &BucketStreamId) -> Option<u64> {
1105 self.state_machine
1106 .head(stream_id)
1107 .map(|metadata| metadata.tail_offset)
1108 }
1109
1110 pub(crate) fn install_snapshot_inner(
1111 &mut self,
1112 snapshot: GroupSnapshot,
1113 ) -> Result<(), GroupEngineError> {
1114 let GroupSnapshot {
1115 placement: _,
1116 group_commit_index,
1117 stream_snapshot,
1118 stream_append_counts,
1119 } = snapshot;
1120 self.install_snapshot_parts(group_commit_index, stream_snapshot, stream_append_counts)
1121 }
1122
1123 pub(crate) fn install_snapshot_parts(
1124 &mut self,
1125 group_commit_index: u64,
1126 stream_snapshot: StreamSnapshot,
1127 stream_append_counts: Vec<StreamAppendCount>,
1128 ) -> Result<(), GroupEngineError> {
1129 let stream_ids = stream_snapshot
1130 .streams
1131 .iter()
1132 .map(|entry| entry.metadata.stream_id.clone())
1133 .collect::<HashSet<_>>();
1134 let state_machine = StreamStateMachine::restore(stream_snapshot)
1135 .map_err(|err| GroupEngineError::new(format!("restore stream snapshot: {err}")))?;
1136 let stream_append_counts = restore_stream_append_counts(stream_append_counts, &stream_ids)?;
1137
1138 self.commit_index = group_commit_index;
1139 self.state_machine = state_machine;
1140 self.stream_append_counts = stream_append_counts;
1141 Ok(())
1142 }
1143}
1144
1145impl GroupEngine for InMemoryGroupEngine {
1146 fn create_stream<'a>(
1147 &'a mut self,
1148 request: CreateStreamRequest,
1149 placement: ShardPlacement,
1150 admission: ColdWriteAdmission,
1151 ) -> GroupCreateStreamFuture<'a> {
1152 if admission.is_enabled() {
1153 return Box::pin(async move {
1154 self.create_stream_with_admission_inner(request, placement, admission)
1155 });
1156 }
1157 let command = GroupWriteCommand::from(request);
1158 Box::pin(async move {
1159 match self.apply_committed_write(command, placement)? {
1160 GroupWriteResponse::CreateStream(response) => Ok(response),
1161 other => Err(GroupEngineError::new(format!(
1162 "unexpected create stream write response: {other:?}"
1163 ))),
1164 }
1165 })
1166 }
1167
1168 fn create_stream_external<'a>(
1169 &'a mut self,
1170 request: CreateStreamExternalRequest,
1171 placement: ShardPlacement,
1172 ) -> GroupCreateStreamFuture<'a> {
1173 Box::pin(async move {
1174 if let Some(cold_store) = self.cold_store.as_ref() {
1175 let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1176 write_external_segment_index_pages(
1177 &store,
1178 &request.stream_id,
1179 0,
1180 &request.initial_payload,
1181 )
1182 .await
1183 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1184 }
1185 let command = GroupWriteCommand::from(request);
1186 match self.apply_committed_write(command, placement)? {
1187 GroupWriteResponse::CreateStream(response) => Ok(response),
1188 other => Err(GroupEngineError::new(format!(
1189 "unexpected external create stream write response: {other:?}"
1190 ))),
1191 }
1192 })
1193 }
1194
1195 fn read_stream<'a>(
1196 &'a mut self,
1197 request: ReadStreamRequest,
1198 placement: ShardPlacement,
1199 ) -> GroupReadStreamFuture<'a> {
1200 Box::pin(async move {
1201 self.read_stream_parts(request, placement)
1202 .await?
1203 .into_response()
1204 .await
1205 })
1206 }
1207
1208 fn read_stream_parts<'a>(
1209 &'a mut self,
1210 request: ReadStreamRequest,
1211 placement: ShardPlacement,
1212 ) -> GroupReadStreamPartsFuture<'a> {
1213 Box::pin(async move {
1214 let stream_id = request.stream_id.clone();
1215 let plan = self.read_stream_plan(&request, placement)?;
1216 Ok(GroupReadStreamParts::from_plan(
1217 placement,
1218 stream_id,
1219 plan,
1220 self.cold_store(),
1221 self.cold_index_cache.clone(),
1222 ))
1223 })
1224 }
1225
1226 fn publish_snapshot<'a>(
1227 &'a mut self,
1228 request: PublishSnapshotRequest,
1229 placement: ShardPlacement,
1230 ) -> GroupPublishSnapshotFuture<'a> {
1231 Box::pin(async move {
1232 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1233 let command = GroupWriteCommand::from(request);
1234 match self.apply_committed_write(command, placement)? {
1235 GroupWriteResponse::PublishSnapshot(response) => Ok(response),
1236 other => Err(GroupEngineError::new(format!(
1237 "unexpected publish snapshot write response: {other:?}"
1238 ))),
1239 }
1240 })
1241 }
1242
1243 fn read_snapshot<'a>(
1244 &'a mut self,
1245 request: ReadSnapshotRequest,
1246 placement: ShardPlacement,
1247 ) -> GroupReadSnapshotFuture<'a> {
1248 Box::pin(async move {
1249 self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1250 let snapshot = match request.snapshot_offset {
1251 Some(offset) => self
1252 .state_machine
1253 .read_snapshot(&request.stream_id, offset)
1254 .map_err(stream_response_error)?,
1255 None => self
1256 .state_machine
1257 .latest_snapshot(&request.stream_id)
1258 .map_err(stream_response_error)?
1259 .ok_or_else(|| {
1260 GroupEngineError::stream(
1261 StreamErrorCode::SnapshotNotFound,
1262 format!("stream '{}' has no visible snapshot", request.stream_id),
1263 )
1264 })?,
1265 };
1266 let tail_offset = self
1267 .state_machine
1268 .head_at(&request.stream_id, request.now_ms)
1269 .map(|metadata| metadata.tail_offset)
1270 .unwrap_or(snapshot.offset);
1271 Ok(ReadSnapshotResponse {
1272 placement,
1273 snapshot_offset: snapshot.offset,
1274 next_offset: snapshot.offset,
1275 content_type: snapshot.content_type,
1276 payload: snapshot.payload,
1277 up_to_date: snapshot.offset == tail_offset,
1278 record_range: self
1279 .state_machine
1280 .record_range(&request.stream_id)
1281 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?,
1282 })
1283 })
1284 }
1285
1286 fn delete_snapshot<'a>(
1287 &'a mut self,
1288 request: DeleteSnapshotRequest,
1289 placement: ShardPlacement,
1290 ) -> GroupDeleteSnapshotFuture<'a> {
1291 Box::pin(async move {
1292 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1293 match self
1294 .state_machine
1295 .delete_snapshot(&request.stream_id, request.snapshot_offset)
1296 {
1297 StreamResponse::Error {
1298 code,
1299 message,
1300 next_offset,
1301 context,
1302 } => Err(GroupEngineError::stream_with_context(
1303 code,
1304 message,
1305 next_offset,
1306 context,
1307 )),
1308 other => Err(GroupEngineError::new(format!(
1309 "unexpected delete snapshot response: {other:?}"
1310 ))),
1311 }
1312 })
1313 }
1314
1315 fn bootstrap_stream<'a>(
1316 &'a mut self,
1317 request: BootstrapStreamRequest,
1318 placement: ShardPlacement,
1319 ) -> GroupBootstrapStreamFuture<'a> {
1320 Box::pin(async move {
1321 self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1322 let plan = self
1323 .state_machine
1324 .bootstrap_plan(&request.stream_id)
1325 .map_err(stream_response_error)?;
1326 let snapshot_offset = plan.snapshot.as_ref().map(|snapshot| snapshot.offset);
1327 let snapshot_content_type = plan
1328 .snapshot
1329 .as_ref()
1330 .map(|snapshot| snapshot.content_type.clone())
1331 .unwrap_or_else(|| DEFAULT_CONTENT_TYPE.to_owned());
1332 let snapshot_payload = plan
1333 .snapshot
1334 .as_ref()
1335 .map(|snapshot| snapshot.payload.clone())
1336 .unwrap_or_default();
1337 let updates = self
1338 .bootstrap_updates(
1339 &request.stream_id,
1340 &plan.updates,
1341 &plan.content_type,
1342 request.now_ms,
1343 )
1344 .await?;
1345 Ok(BootstrapStreamResponse {
1346 placement,
1347 snapshot_offset,
1348 snapshot_content_type,
1349 snapshot_payload,
1350 updates,
1351 next_offset: plan.next_offset,
1352 up_to_date: plan.up_to_date,
1353 closed: plan.closed,
1354 record_range: self
1355 .state_machine
1356 .record_range(&request.stream_id)
1357 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?,
1358 })
1359 })
1360 }
1361
1362 fn touch_stream_access<'a>(
1363 &'a mut self,
1364 stream_id: BucketStreamId,
1365 now_ms: u64,
1366 renew_ttl: bool,
1367 placement: ShardPlacement,
1368 ) -> GroupTouchStreamAccessFuture<'a> {
1369 Box::pin(async move { self.apply_access_command(stream_id, now_ms, renew_ttl, placement) })
1370 }
1371
1372 fn head_stream<'a>(
1373 &'a mut self,
1374 request: HeadStreamRequest,
1375 placement: ShardPlacement,
1376 ) -> GroupHeadStreamFuture<'a> {
1377 Box::pin(async move {
1378 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1379 self.head_stream_after_access(&request, placement)
1380 })
1381 }
1382
1383 fn get_stream_attrs<'a>(
1384 &'a mut self,
1385 request: GetStreamAttrsRequest,
1386 placement: ShardPlacement,
1387 ) -> GroupGetStreamAttrsFuture<'a> {
1388 Box::pin(async move {
1389 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1390 self.get_stream_attrs_after_access(&request, placement)
1391 })
1392 }
1393
1394 fn update_stream_attrs<'a>(
1395 &'a mut self,
1396 request: UpdateStreamAttrsRequest,
1397 placement: ShardPlacement,
1398 ) -> GroupUpdateStreamAttrsFuture<'a> {
1399 Box::pin(async move {
1400 match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
1401 GroupWriteResponse::UpdateStreamAttrs(response) => Ok(response),
1402 other => Err(GroupEngineError::new(format!(
1403 "unexpected update stream attrs write response: {other:?}"
1404 ))),
1405 }
1406 })
1407 }
1408
1409 fn close_stream<'a>(
1410 &'a mut self,
1411 request: CloseStreamRequest,
1412 placement: ShardPlacement,
1413 ) -> GroupCloseStreamFuture<'a> {
1414 Box::pin(async move {
1415 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1416 let command = GroupWriteCommand::from(request);
1417 match self.apply_committed_write(command, placement)? {
1418 GroupWriteResponse::CloseStream(response) => Ok(response),
1419 other => Err(GroupEngineError::new(format!(
1420 "unexpected close stream write response: {other:?}"
1421 ))),
1422 }
1423 })
1424 }
1425
1426 fn delete_stream<'a>(
1427 &'a mut self,
1428 request: DeleteStreamRequest,
1429 placement: ShardPlacement,
1430 ) -> GroupDeleteStreamFuture<'a> {
1431 let command = GroupWriteCommand::from(request);
1432 Box::pin(async move {
1433 match self.apply_committed_write(command, placement)? {
1434 GroupWriteResponse::DeleteStream(response) => Ok(response),
1435 other => Err(GroupEngineError::new(format!(
1436 "unexpected delete stream write response: {other:?}"
1437 ))),
1438 }
1439 })
1440 }
1441
1442 fn ack_cold_gc<'a>(
1443 &'a mut self,
1444 up_to_seq: u64,
1445 placement: ShardPlacement,
1446 ) -> GroupAckColdGcFuture<'a> {
1447 Box::pin(async move {
1448 match self.apply_committed_write(
1449 GroupWriteCommand::Stream(StreamCommand::AckColdGc { up_to_seq }),
1450 placement,
1451 )? {
1452 GroupWriteResponse::AckColdGc(response) => Ok(response),
1453 other => Err(GroupEngineError::new(format!(
1454 "unexpected ack cold gc write response: {other:?}"
1455 ))),
1456 }
1457 })
1458 }
1459
1460 fn plan_cold_gc<'a>(
1461 &'a mut self,
1462 max: usize,
1463 _placement: ShardPlacement,
1464 ) -> GroupPlanColdGcFuture<'a> {
1465 let entries = self.state_machine.pending_cold_gc_batch(max);
1466 Box::pin(async move { Ok(entries) })
1467 }
1468
1469 fn append<'a>(
1470 &'a mut self,
1471 request: AppendRequest,
1472 placement: ShardPlacement,
1473 admission: ColdWriteAdmission,
1474 ) -> GroupAppendFuture<'a> {
1475 if admission.is_enabled() {
1476 return Box::pin(async move {
1477 self.append_with_admission_inner(request, placement, admission)
1478 });
1479 }
1480 Box::pin(async move {
1481 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1482 let command = GroupWriteCommand::from(request);
1483 match self.apply_committed_write(command, placement)? {
1484 GroupWriteResponse::Append(response) => Ok(response),
1485 other => Err(GroupEngineError::new(format!(
1486 "unexpected append write response: {other:?}"
1487 ))),
1488 }
1489 })
1490 }
1491
1492 fn append_external<'a>(
1493 &'a mut self,
1494 request: AppendExternalRequest,
1495 placement: ShardPlacement,
1496 ) -> GroupAppendFuture<'a> {
1497 Box::pin(async move {
1498 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1499 if let Some(cold_store) = self.cold_store.as_ref() {
1500 let start_offset = self
1501 .state_machine
1502 .head(&request.stream_id)
1503 .map(|metadata| metadata.tail_offset)
1504 .ok_or_else(|| {
1505 GroupEngineError::stream(
1506 ursula_stream::StreamErrorCode::StreamNotFound,
1507 format!("stream '{}' does not exist", request.stream_id),
1508 )
1509 })?;
1510 let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1511 write_external_segment_index_pages(
1512 &store,
1513 &request.stream_id,
1514 start_offset,
1515 &request.payload,
1516 )
1517 .await
1518 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1519 }
1520 let command = GroupWriteCommand::from(request);
1521 match self.apply_committed_write(command, placement)? {
1522 GroupWriteResponse::Append(response) => Ok(response),
1523 other => Err(GroupEngineError::new(format!(
1524 "unexpected external append write response: {other:?}"
1525 ))),
1526 }
1527 })
1528 }
1529
1530 fn append_batch<'a>(
1531 &'a mut self,
1532 request: AppendBatchRequest,
1533 placement: ShardPlacement,
1534 admission: ColdWriteAdmission,
1535 ) -> GroupAppendBatchFuture<'a> {
1536 if admission.is_enabled() {
1537 return Box::pin(async move {
1538 self.append_batch_with_admission_inner(request, placement, admission)
1539 });
1540 }
1541 Box::pin(async move {
1542 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1543 let command = GroupWriteCommand::from(request);
1544 match self.apply_committed_write(command, placement)? {
1545 GroupWriteResponse::AppendBatch(response) => Ok(response),
1546 other => Err(GroupEngineError::new(format!(
1547 "unexpected append batch write response: {other:?}"
1548 ))),
1549 }
1550 })
1551 }
1552
1553 fn flush_cold<'a>(
1554 &'a mut self,
1555 request: FlushColdRequest,
1556 placement: ShardPlacement,
1557 ) -> GroupFlushColdFuture<'a> {
1558 Box::pin(async move {
1559 let mut index_rollback = None;
1560 if let Some(cold_store) = self.cold_store.as_ref() {
1561 let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1562 let rollback = write_cold_chunk_index_pages_with_rollback(
1563 &store,
1564 &request.stream_id,
1565 &request.chunk,
1566 )
1567 .await
1568 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1569 index_rollback = Some((store, rollback));
1570 }
1571 let command = GroupWriteCommand::from(request);
1572 match self.apply_committed_write(command, placement) {
1573 Ok(GroupWriteResponse::FlushCold(response)) => Ok(response),
1574 Ok(other) => {
1575 if let Some((store, rollback)) = index_rollback {
1576 rollback_cold_index_pages(&store, rollback)
1577 .await
1578 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1579 }
1580 Err(GroupEngineError::new(format!(
1581 "unexpected flush cold write response: {other:?}"
1582 )))
1583 }
1584 Err(err) => {
1585 if let Some((store, rollback)) = index_rollback {
1586 rollback_cold_index_pages(&store, rollback).await.map_err(
1587 |rollback_err| {
1588 GroupEngineError::new(format!(
1589 "rollback cold index after flush failure: {rollback_err}"
1590 ))
1591 },
1592 )?;
1593 }
1594 Err(err)
1595 }
1596 }
1597 })
1598 }
1599
1600 fn plan_cold_flush<'a>(
1601 &'a mut self,
1602 request: PlanColdFlushRequest,
1603 _placement: ShardPlacement,
1604 ) -> GroupPlanColdFlushFuture<'a> {
1605 Box::pin(async move {
1606 self.state_machine
1607 .plan_cold_flush(
1608 &request.stream_id,
1609 request.min_hot_bytes,
1610 request.max_flush_bytes,
1611 )
1612 .map_err(stream_response_error)
1613 })
1614 }
1615
1616 fn plan_next_cold_flush_batch<'a>(
1617 &'a mut self,
1618 request: PlanGroupColdFlushRequest,
1619 _placement: ShardPlacement,
1620 max_candidates: usize,
1621 ) -> GroupPlanNextColdFlushBatchFuture<'a> {
1622 Box::pin(async move {
1623 self.state_machine
1624 .plan_next_cold_flush_batch(
1625 request.min_hot_bytes,
1626 request.max_flush_bytes,
1627 max_candidates,
1628 )
1629 .map_err(stream_response_error)
1630 })
1631 }
1632
1633 fn cold_hot_backlog<'a>(
1634 &'a mut self,
1635 stream_id: BucketStreamId,
1636 _placement: ShardPlacement,
1637 ) -> GroupColdHotBacklogFuture<'a> {
1638 Box::pin(async move { self.cold_hot_backlog_for(stream_id) })
1639 }
1640
1641 fn snapshot<'a>(&'a mut self, placement: ShardPlacement) -> GroupSnapshotFuture<'a> {
1642 Box::pin(async move { Ok(self.build_snapshot(placement)) })
1643 }
1644
1645 fn install_snapshot<'a>(
1646 &'a mut self,
1647 snapshot: GroupSnapshot,
1648 ) -> GroupInstallSnapshotFuture<'a> {
1649 Box::pin(async move { self.install_snapshot_inner(snapshot) })
1650 }
1651}
1652
1653#[derive(Debug, Clone, Default)]
1654pub struct InMemoryGroupEngineFactory {
1655 cold_store: Option<ColdStoreHandle>,
1656}
1657
1658impl InMemoryGroupEngineFactory {
1659 pub fn new() -> Self {
1660 Self::default()
1661 }
1662
1663 pub fn with_cold_store(cold_store: Option<ColdStoreHandle>) -> Self {
1664 Self { cold_store }
1665 }
1666}
1667
1668impl GroupEngineFactory for InMemoryGroupEngineFactory {
1669 fn create<'a>(
1670 &'a self,
1671 _placement: ShardPlacement,
1672 _metrics: GroupEngineMetrics,
1673 ) -> GroupEngineCreateFuture<'a> {
1674 Box::pin(async move {
1675 let mut engine = InMemoryGroupEngine::default();
1676 engine.set_cold_store(self.cold_store.clone());
1677 let engine: Box<dyn GroupEngine> = Box::new(engine);
1678 Ok(engine)
1679 })
1680 }
1681}
1682
1683pub(crate) fn compare_stream_ids(
1684 left: &BucketStreamId,
1685 right: &BucketStreamId,
1686) -> std::cmp::Ordering {
1687 left.bucket_id
1688 .cmp(&right.bucket_id)
1689 .then_with(|| left.stream_id.cmp(&right.stream_id))
1690}
1691pub(crate) fn ensure_bucket_exists(
1692 state_machine: &mut StreamStateMachine,
1693 stream_id: &BucketStreamId,
1694) -> Result<(), GroupEngineError> {
1695 if state_machine.bucket_exists(&stream_id.bucket_id) {
1696 return Ok(());
1697 }
1698
1699 match state_machine.apply(StreamCommand::CreateBucket {
1700 bucket_id: stream_id.bucket_id.clone(),
1701 }) {
1702 StreamResponse::BucketCreated { .. } | StreamResponse::BucketAlreadyExists { .. } => Ok(()),
1703 StreamResponse::Error {
1704 code,
1705 message,
1706 next_offset,
1707 context,
1708 } => Err(GroupEngineError::stream_with_context(
1709 code,
1710 message,
1711 next_offset,
1712 context,
1713 )),
1714 other => Err(GroupEngineError::new(format!(
1715 "unexpected create bucket response: {other:?}"
1716 ))),
1717 }
1718}
1719
1720fn command_stream_id(command: &StreamCommand) -> Option<BucketStreamId> {
1722 match command {
1723 StreamCommand::CreateBucket { .. }
1724 | StreamCommand::DeleteBucket { .. }
1725 | StreamCommand::AckColdGc { .. } => None,
1726 StreamCommand::CreateStream { stream_id, .. }
1727 | StreamCommand::CreateExternal { stream_id, .. }
1728 | StreamCommand::Append { stream_id, .. }
1729 | StreamCommand::AppendExternal { stream_id, .. }
1730 | StreamCommand::AppendBatch { stream_id, .. }
1731 | StreamCommand::PublishSnapshot { stream_id, .. }
1732 | StreamCommand::TouchStreamAccess { stream_id, .. }
1733 | StreamCommand::UpdateStreamAttrs { stream_id, .. }
1734 | StreamCommand::FlushCold { stream_id, .. }
1735 | StreamCommand::Close { stream_id, .. }
1736 | StreamCommand::DeleteStream { stream_id } => Some(stream_id.clone()),
1737 }
1738}
1739
1740fn command_producer(command: &StreamCommand) -> Option<ProducerRequest> {
1741 match command {
1742 StreamCommand::Close { producer, .. } => producer.clone(),
1743 _ => None,
1744 }
1745}
1746
1747fn require_response_stream_id(
1748 stream_id: Option<BucketStreamId>,
1749 response: &str,
1750) -> Result<BucketStreamId, GroupEngineError> {
1751 stream_id.ok_or_else(|| {
1752 GroupEngineError::new(format!(
1753 "{response} response for a command without a stream id"
1754 ))
1755 })
1756}
1757
1758pub(crate) fn stream_response_error(response: StreamResponse) -> GroupEngineError {
1759 match response {
1760 StreamResponse::Error {
1761 code,
1762 message,
1763 next_offset,
1764 context,
1765 } => GroupEngineError::stream_with_context(code, message, next_offset, context),
1766 other => GroupEngineError::new(format!("unexpected stream response error: {other:?}")),
1767 }
1768}
1769
1770pub(crate) fn restore_stream_append_counts(
1771 counts: Vec<StreamAppendCount>,
1772 snapshot_stream_ids: &HashSet<BucketStreamId>,
1773) -> Result<HashMap<BucketStreamId, u64>, GroupEngineError> {
1774 let mut restored = HashMap::with_capacity(counts.len());
1775 for count in counts {
1776 if !snapshot_stream_ids.contains(&count.stream_id) {
1777 return Err(GroupEngineError::new(format!(
1778 "append count references missing snapshot stream '{}'",
1779 count.stream_id
1780 )));
1781 }
1782 if restored
1783 .insert(count.stream_id.clone(), count.append_count)
1784 .is_some()
1785 {
1786 return Err(GroupEngineError::new(format!(
1787 "snapshot contains duplicate append count for stream '{}'",
1788 count.stream_id
1789 )));
1790 }
1791 }
1792 Ok(restored)
1793}