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