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