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 } => {
617 self.commit_index += 1;
618 Ok(GroupWriteResponse::PurgeBucket(PurgeBucketResponse {
619 placement,
620 removed_streams,
621 group_commit_index: self.commit_index,
622 }))
623 }
624 StreamResponse::Error {
625 code,
626 message,
627 next_offset,
628 context,
629 } => Err(GroupEngineError::stream_with_context(
630 code,
631 message,
632 next_offset,
633 context,
634 )),
635 other @ (StreamResponse::BucketCreated { .. }
636 | StreamResponse::BucketAlreadyExists { .. }
637 | StreamResponse::BucketDeleted { .. }) => Err(GroupEngineError::new(format!(
638 "unexpected group write response: {other:?}"
639 ))),
640 }
641 }
642
643 pub(crate) fn cold_hot_backlog_for(
644 &self,
645 stream_id: BucketStreamId,
646 ) -> Result<ColdHotBacklog, GroupEngineError> {
647 let stream_hot_bytes = self.state_machine.hot_payload_len(&stream_id).unwrap_or(0);
648 Ok(ColdHotBacklog {
649 stream_id,
650 stream_hot_bytes,
651 group_hot_bytes: self.state_machine.total_hot_payload_bytes(),
652 })
653 }
654
655 pub fn check_cold_write_admission_bytes(
656 &self,
657 stream_id: &BucketStreamId,
658 admission: ColdWriteAdmission,
659 incoming_bytes: u64,
660 ) -> Result<(), GroupEngineError> {
661 let Some(limit) = admission.max_hot_bytes_per_group else {
662 return Ok(());
663 };
664 if incoming_bytes == 0 {
665 return Ok(());
666 }
667 let before = self.state_machine.total_hot_payload_bytes();
668 let after = before.saturating_add(incoming_bytes);
669 if after <= limit {
670 return Ok(());
671 }
672 Err(GroupEngineError::cold_backpressure(
673 stream_id.clone(),
674 before,
675 after,
676 limit,
677 ))
678 }
679
680 pub(crate) fn create_stream_with_admission_inner(
681 &mut self,
682 request: CreateStreamRequest,
683 placement: ShardPlacement,
684 admission: ColdWriteAdmission,
685 ) -> Result<CreateStreamResponse, GroupEngineError> {
686 let stream_id = request.stream_id.clone();
687 if admission.is_enabled() {
688 let mut preview = self.clone();
689 let preview_response = match preview
690 .apply_committed_write(GroupWriteCommand::from(request.clone()), placement)?
691 {
692 GroupWriteResponse::CreateStream(response) => response,
693 other => {
694 return Err(GroupEngineError::new(format!(
695 "unexpected create stream preview response: {other:?}"
696 )));
697 }
698 };
699 if !preview_response.already_exists {
700 self.check_cold_write_admission_bytes(
701 &stream_id,
702 admission,
703 u64::try_from(request.initial_payload.len()).expect("payload len fits u64"),
704 )?;
705 }
706 }
707 let response =
708 match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
709 GroupWriteResponse::CreateStream(response) => response,
710 other => {
711 return Err(GroupEngineError::new(format!(
712 "unexpected create stream write response: {other:?}"
713 )));
714 }
715 };
716 Ok(response)
717 }
718
719 pub(crate) fn append_with_admission_inner(
720 &mut self,
721 request: AppendRequest,
722 placement: ShardPlacement,
723 admission: ColdWriteAdmission,
724 ) -> Result<AppendResponse, GroupEngineError> {
725 let stream_id = request.stream_id.clone();
726 if admission.is_enabled() {
727 let mut preview = self.clone();
728 let preview_response = match preview
729 .apply_committed_write(GroupWriteCommand::from(request.clone()), placement)?
730 {
731 GroupWriteResponse::Append(response) => response,
732 other => {
733 return Err(GroupEngineError::new(format!(
734 "unexpected append preview response: {other:?}"
735 )));
736 }
737 };
738 if !preview_response.deduplicated {
739 self.check_cold_write_admission_bytes(
740 &stream_id,
741 admission,
742 u64::try_from(request.payload.len()).expect("payload len fits u64"),
743 )?;
744 }
745 }
746 let response =
747 match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
748 GroupWriteResponse::Append(response) => response,
749 other => {
750 return Err(GroupEngineError::new(format!(
751 "unexpected append write response: {other:?}"
752 )));
753 }
754 };
755 Ok(response)
756 }
757
758 pub(crate) fn append_batch_with_admission_inner(
759 &mut self,
760 request: AppendBatchRequest,
761 placement: ShardPlacement,
762 admission: ColdWriteAdmission,
763 ) -> Result<GroupAppendBatchResponse, GroupEngineError> {
764 let stream_id = request.stream_id.clone();
765 let incoming_bytes = request
766 .payloads
767 .iter()
768 .map(|payload| u64::try_from(payload.len()).expect("payload len fits u64"))
769 .sum();
770 if admission.is_enabled() {
771 let mut preview = self.clone();
772 let preview_response = match preview
773 .apply_committed_write(GroupWriteCommand::from(request.clone()), placement)?
774 {
775 GroupWriteResponse::AppendBatch(response) => response,
776 other => {
777 return Err(GroupEngineError::new(format!(
778 "unexpected append batch preview response: {other:?}"
779 )));
780 }
781 };
782 let mutates = preview_response
783 .items
784 .iter()
785 .any(|item| matches!(item, Ok(response) if !response.deduplicated));
786 if mutates {
787 self.check_cold_write_admission_bytes(&stream_id, admission, incoming_bytes)?;
788 }
789 }
790 let response =
791 match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
792 GroupWriteResponse::AppendBatch(response) => response,
793 other => {
794 return Err(GroupEngineError::new(format!(
795 "unexpected append batch write response: {other:?}"
796 )));
797 }
798 };
799 Ok(response)
800 }
801
802 pub fn access_requires_write(
803 &self,
804 stream_id: &BucketStreamId,
805 now_ms: u64,
806 renew_ttl: bool,
807 ) -> Result<bool, GroupEngineError> {
808 self.state_machine
809 .access_requires_write(stream_id, now_ms, renew_ttl)
810 .map_err(stream_response_error)
811 }
812
813 pub(crate) fn apply_access_command(
814 &mut self,
815 stream_id: BucketStreamId,
816 now_ms: u64,
817 renew_ttl: bool,
818 placement: ShardPlacement,
819 ) -> Result<TouchStreamAccessResponse, GroupEngineError> {
820 match self.apply_committed_write(
821 GroupWriteCommand::Stream(StreamCommand::TouchStreamAccess {
822 stream_id,
823 now_ms,
824 renew_ttl,
825 }),
826 placement,
827 )? {
828 GroupWriteResponse::TouchStreamAccess(response) => Ok(response),
829 other => Err(GroupEngineError::new(format!(
830 "unexpected touch stream access write response: {other:?}"
831 ))),
832 }
833 }
834
835 pub(crate) fn ensure_stream_access(
836 &mut self,
837 stream_id: &BucketStreamId,
838 now_ms: u64,
839 renew_ttl: bool,
840 placement: ShardPlacement,
841 ) -> Result<Option<TouchStreamAccessResponse>, GroupEngineError> {
842 if !self.access_requires_write(stream_id, now_ms, renew_ttl)? {
843 return Ok(None);
844 }
845 let response =
846 self.apply_access_command(stream_id.clone(), now_ms, renew_ttl, placement)?;
847 if response.expired {
848 return Err(GroupEngineError::stream(
849 StreamErrorCode::StreamNotFound,
850 format!("stream '{stream_id}' does not exist"),
851 ));
852 }
853 Ok(Some(response))
854 }
855
856 pub(crate) fn append_payload(
857 &mut self,
858 input: AppendPayloadInput<'_>,
859 placement: ShardPlacement,
860 ) -> Result<AppendResponse, GroupEngineError> {
861 let AppendPayloadInput {
862 stream_id,
863 content_type,
864 payload,
865 close_after,
866 stream_seq,
867 producer,
868 now_ms,
869 record_match,
870 } = input;
871 let stream_count_key = stream_id.clone();
872 let response = self.state_machine.append_borrowed(AppendStreamInput {
873 stream_id,
874 content_type,
875 payload,
876 close_after,
877 stream_seq,
878 producer,
879 now_ms,
880 record_match,
881 });
882 self.append_response_from_stream(stream_count_key, response, placement)
883 }
884
885 fn append_response_from_stream(
886 &mut self,
887 stream_id: BucketStreamId,
888 response: StreamResponse,
889 placement: ShardPlacement,
890 ) -> Result<AppendResponse, GroupEngineError> {
891 match response {
892 StreamResponse::Appended {
893 offset,
894 next_offset,
895 closed,
896 deduplicated,
897 producer,
898 ..
899 } => {
900 let stream_hot_bytes = self.state_machine.hot_payload_len(&stream_id).unwrap_or(0);
901 let group_hot_bytes = self.state_machine.total_hot_payload_bytes();
902 let stream_append_count = self
903 .stream_append_counts
904 .entry(stream_id.clone())
905 .or_insert(0);
906 let record_range = self
907 .state_machine
908 .record_range_for_append(&stream_id, offset, next_offset, producer.as_ref())
909 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?;
910 if !deduplicated {
911 self.commit_index += 1;
912 *stream_append_count += 1;
913 }
914 Ok(AppendResponse {
915 placement,
916 start_offset: offset,
917 next_offset,
918 stream_append_count: *stream_append_count,
919 group_commit_index: self.commit_index,
920 closed,
921 deduplicated,
922 producer,
923 record_range,
924 stream_hot_bytes,
925 group_hot_bytes,
926 })
927 }
928 StreamResponse::Error {
929 code,
930 message,
931 next_offset,
932 context,
933 } => Err(GroupEngineError::stream_with_context(
934 code,
935 message,
936 next_offset,
937 context,
938 )),
939 other => Err(GroupEngineError::new(format!(
940 "unexpected append response: {other:?}"
941 ))),
942 }
943 }
944
945 pub fn read_stream_plan(
946 &mut self,
947 request: &ReadStreamRequest,
948 placement: ShardPlacement,
949 ) -> Result<StreamReadPlan, GroupEngineError> {
950 self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
951 self.read_stream_plan_after_access(request)
952 }
953
954 pub fn read_stream_plan_after_access(
955 &self,
956 request: &ReadStreamRequest,
957 ) -> Result<StreamReadPlan, GroupEngineError> {
958 let Some(record) = request.record else {
959 let mut plan = self
960 .state_machine
961 .read_plan_at(
962 &request.stream_id,
963 request.offset,
964 request.max_len,
965 request.now_ms,
966 )
967 .map_err(stream_response_error)?;
968 plan.retained_record_range = self
969 .state_machine
970 .record_range(&request.stream_id)
971 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?;
972 return Ok(plan);
973 };
974 let retained_record_range = self
975 .state_machine
976 .record_range(&request.stream_id)
977 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?
978 .ok_or_else(|| {
979 GroupEngineError::stream(
980 StreamErrorCode::InvalidRecordBoundaries,
981 "record coordinates are inactive for this stream",
982 )
983 })?;
984 if record < retained_record_range.first_record {
985 return Err(GroupEngineError::stream(
986 StreamErrorCode::StreamGone,
987 format!(
988 "record {record} is older than first retained record {}",
989 retained_record_range.first_record
990 ),
991 ));
992 }
993 if record > retained_record_range.next_record {
994 return Err(GroupEngineError::stream(
995 StreamErrorCode::InvalidRecordBoundaries,
996 format!(
997 "record {record} is beyond record tail {}",
998 retained_record_range.next_record
999 ),
1000 ));
1001 }
1002 let next_record = request
1003 .max_records
1004 .map(|limit| record.saturating_add(limit))
1005 .unwrap_or(retained_record_range.next_record)
1006 .min(retained_record_range.next_record);
1007 let offset = self
1008 .state_machine
1009 .offset_for_record(&request.stream_id, record)
1010 .map_err(|err| GroupEngineError::new(format!("record offset: {err:?}")))?
1011 .ok_or_else(|| GroupEngineError::new("record stream disappeared"))?;
1012 let next_offset = self
1013 .state_machine
1014 .offset_for_record(&request.stream_id, next_record)
1015 .map_err(|err| GroupEngineError::new(format!("record offset: {err:?}")))?
1016 .ok_or_else(|| GroupEngineError::new("record stream disappeared"))?;
1017 let max_len = usize::try_from(next_offset.saturating_sub(offset))
1018 .map_err(|_| GroupEngineError::new("record read window exceeds usize"))?;
1019 let mut plan = self
1020 .state_machine
1021 .read_plan_at(&request.stream_id, offset, max_len, request.now_ms)
1022 .map_err(stream_response_error)?;
1023 plan.retained_record_range = Some(retained_record_range);
1024 plan.record_range = Some(ursula_stream::StreamRecordRange {
1025 first_record: record,
1026 next_record,
1027 });
1028 Ok(plan)
1029 }
1030
1031 pub fn bucket_usage_report(&self) -> Vec<ursula_stream::BucketUsageSnapshot> {
1034 self.state_machine.bucket_usage_report()
1035 }
1036
1037 pub fn head_stream_after_access(
1038 &mut self,
1039 request: &HeadStreamRequest,
1040 placement: ShardPlacement,
1041 ) -> Result<HeadStreamResponse, GroupEngineError> {
1042 let Some(metadata) = self
1043 .state_machine
1044 .head_at(&request.stream_id, request.now_ms)
1045 else {
1046 return Err(GroupEngineError::stream(
1047 StreamErrorCode::StreamNotFound,
1048 format!("stream '{}' does not exist", request.stream_id),
1049 ));
1050 };
1051 let content_type = metadata.content_type.clone();
1052 let tail_offset = metadata.tail_offset;
1053 let closed = metadata.status == ursula_stream::StreamStatus::Closed;
1054 let stream_ttl_seconds = metadata.stream_ttl_seconds;
1055 let stream_expires_at_ms = metadata.stream_expires_at_ms;
1056 let _ = metadata;
1057 let snapshot = self
1058 .state_machine
1059 .latest_snapshot(&request.stream_id)
1060 .map_err(stream_response_error)?;
1061 Ok(HeadStreamResponse {
1062 placement,
1063 content_type,
1064 tail_offset,
1065 cold_hot_start_offset: self.state_machine.hot_start_offset(&request.stream_id),
1066 closed,
1067 stream_ttl_seconds,
1068 stream_expires_at_ms,
1069 snapshot_offset: snapshot.as_ref().map(|snapshot| snapshot.offset),
1070 snapshot_digest: snapshot.map(|snapshot| snapshot.digest),
1071 retained_offset: self.state_machine.retained_offset(&request.stream_id),
1072 integrity: self
1073 .state_machine
1074 .integrity_snapshot(&request.stream_id)
1075 .map_err(stream_response_error)?,
1076 record_range: self
1077 .state_machine
1078 .record_range(&request.stream_id)
1079 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?,
1080 })
1081 }
1082
1083 pub fn get_stream_attrs_after_access(
1084 &mut self,
1085 request: &GetStreamAttrsRequest,
1086 placement: ShardPlacement,
1087 ) -> Result<GetStreamAttrsResponse, GroupEngineError> {
1088 if self
1089 .state_machine
1090 .head_at(&request.stream_id, request.now_ms)
1091 .is_none()
1092 {
1093 return Err(GroupEngineError::stream(
1094 StreamErrorCode::StreamNotFound,
1095 format!("stream '{}' does not exist", request.stream_id),
1096 ));
1097 }
1098 Ok(GetStreamAttrsResponse {
1099 placement,
1100 attrs: self.state_machine.stream_attrs(&request.stream_id).cloned(),
1101 })
1102 }
1103
1104 pub async fn read_payload_from_plan(
1105 cold_store: Option<&ColdStoreHandle>,
1106 cold_index_cache: Option<&Arc<ColdIndexPageCache<ColdStoreColdIndexPageStore>>>,
1107 stream_id: &BucketStreamId,
1108 plan: &StreamReadPlan,
1109 ) -> Result<Vec<u8>, GroupEngineError> {
1110 let mut payload = Vec::new();
1111 for segment in &plan.segments {
1112 match segment {
1113 StreamReadSegment::Hot(bytes) => payload.extend_from_slice(bytes),
1114 StreamReadSegment::ColdIndex(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 Some(cache) = cold_index_cache else {
1123 return Err(GroupEngineError::stream_with_next_offset(
1124 StreamErrorCode::InvalidColdFlush,
1125 format!("stream '{stream_id}' read requires cold index page cache"),
1126 Some(plan.next_offset),
1127 ));
1128 };
1129 let objects = cache
1130 .object_segments_for_read(stream_id, segment)
1131 .await
1132 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1133 let segment_end = segment
1134 .read_start_offset
1135 .saturating_add(u64::try_from(segment.len).expect("read len fits u64"));
1136 let mut cursor = segment.read_start_offset;
1137 for object in objects {
1138 let start = object
1143 .start_offset
1144 .max(segment.read_start_offset)
1145 .max(cursor);
1146 let end = object.end_offset.min(segment_end);
1147 if start >= end {
1148 continue;
1149 }
1150 let bytes = cold_store
1151 .read_object_range_for_stream(
1152 stream_id,
1153 &object,
1154 start,
1155 usize::try_from(end - start).expect("object read len fits usize"),
1156 )
1157 .await
1158 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1159 payload.extend_from_slice(&bytes);
1160 cursor = end;
1161 }
1162 }
1163 StreamReadSegment::Object(segment) => {
1164 let Some(cold_store) = cold_store else {
1165 return Err(GroupEngineError::stream_with_next_offset(
1166 StreamErrorCode::InvalidColdFlush,
1167 format!("stream '{stream_id}' read requires object payload store"),
1168 Some(plan.next_offset),
1169 ));
1170 };
1171 let bytes = cold_store
1172 .read_object_range_for_stream(
1173 stream_id,
1174 &segment.object,
1175 segment.read_start_offset,
1176 segment.len,
1177 )
1178 .await
1179 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1180 payload.extend_from_slice(&bytes);
1181 }
1182 }
1183 }
1184 Ok(payload)
1185 }
1186
1187 pub(crate) async fn read_own_payload_from_plan(
1188 &self,
1189 stream_id: &BucketStreamId,
1190 plan: &StreamReadPlan,
1191 ) -> Result<Vec<u8>, GroupEngineError> {
1192 Self::read_payload_from_plan(
1193 self.cold_store.as_ref(),
1194 self.cold_index_cache.as_ref(),
1195 stream_id,
1196 plan,
1197 )
1198 .await
1199 }
1200
1201 pub(crate) async fn bootstrap_updates(
1202 &self,
1203 stream_id: &BucketStreamId,
1204 records: &[StreamMessageRecord],
1205 content_type: &str,
1206 now_ms: u64,
1207 ) -> Result<Vec<BootstrapUpdate>, GroupEngineError> {
1208 let mut updates = Vec::with_capacity(records.len());
1209 for record in records {
1210 let len = usize::try_from(record.end_offset - record.start_offset).map_err(|_| {
1211 GroupEngineError::stream(
1212 StreamErrorCode::InvalidSnapshot,
1213 format!(
1214 "bootstrap message [{}..{}) for stream '{stream_id}' is too large",
1215 record.start_offset, record.end_offset
1216 ),
1217 )
1218 })?;
1219 let plan = self
1220 .state_machine
1221 .read_plan_at(stream_id, record.start_offset, len, now_ms)
1222 .map_err(stream_response_error)?;
1223 let payload = self.read_own_payload_from_plan(stream_id, &plan).await?;
1224 updates.push(BootstrapUpdate {
1225 start_offset: record.start_offset,
1226 next_offset: record.end_offset,
1227 content_type: content_type.to_owned(),
1228 payload,
1229 });
1230 }
1231 Ok(updates)
1232 }
1233
1234 pub(crate) fn build_snapshot(&self, placement: ShardPlacement) -> GroupSnapshot {
1235 let stream_snapshot = self.state_machine.snapshot();
1236 let stream_append_counts = self.stream_append_counts_snapshot(&stream_snapshot);
1237 GroupSnapshot {
1238 placement,
1239 group_commit_index: self.commit_index,
1240 stream_snapshot,
1241 stream_append_counts,
1242 }
1243 }
1244
1245 pub(crate) fn stream_append_counts_snapshot(
1246 &self,
1247 stream_snapshot: &ursula_stream::StreamSnapshot,
1248 ) -> Vec<StreamAppendCount> {
1249 let live: HashSet<&BucketStreamId> = stream_snapshot
1256 .streams
1257 .iter()
1258 .map(|entry| &entry.metadata.stream_id)
1259 .collect();
1260 let mut counts = self
1261 .stream_append_counts
1262 .iter()
1263 .filter(|(stream_id, _)| live.contains(stream_id))
1264 .map(|(stream_id, append_count)| StreamAppendCount {
1265 stream_id: stream_id.clone(),
1266 append_count: *append_count,
1267 })
1268 .collect::<Vec<_>>();
1269 counts.sort_by(|left, right| compare_stream_ids(&left.stream_id, &right.stream_id));
1270 counts
1271 }
1272
1273 pub fn stream_tail_offset(&self, stream_id: &BucketStreamId) -> Option<u64> {
1274 self.state_machine
1275 .head(stream_id)
1276 .map(|metadata| metadata.tail_offset)
1277 }
1278
1279 pub(crate) fn install_snapshot_inner(
1280 &mut self,
1281 snapshot: GroupSnapshot,
1282 ) -> Result<(), GroupEngineError> {
1283 let GroupSnapshot {
1284 placement: _,
1285 group_commit_index,
1286 stream_snapshot,
1287 stream_append_counts,
1288 } = snapshot;
1289 self.install_snapshot_parts(group_commit_index, stream_snapshot, stream_append_counts)
1290 }
1291
1292 pub(crate) fn install_snapshot_parts(
1293 &mut self,
1294 group_commit_index: u64,
1295 stream_snapshot: StreamSnapshot,
1296 stream_append_counts: Vec<StreamAppendCount>,
1297 ) -> Result<(), GroupEngineError> {
1298 let stream_ids = stream_snapshot
1299 .streams
1300 .iter()
1301 .map(|entry| entry.metadata.stream_id.clone())
1302 .collect::<HashSet<_>>();
1303 let state_machine = StreamStateMachine::restore(stream_snapshot)
1304 .map_err(|err| GroupEngineError::new(format!("restore stream snapshot: {err}")))?;
1305 let stream_append_counts = restore_stream_append_counts(stream_append_counts, &stream_ids)?;
1306
1307 self.commit_index = group_commit_index;
1308 self.state_machine = state_machine;
1309 self.stream_append_counts = stream_append_counts;
1310 Ok(())
1311 }
1312}
1313
1314impl GroupEngine for InMemoryGroupEngine {
1315 fn create_stream<'a>(
1316 &'a mut self,
1317 request: CreateStreamRequest,
1318 placement: ShardPlacement,
1319 admission: ColdWriteAdmission,
1320 ) -> GroupCreateStreamFuture<'a> {
1321 if admission.is_enabled() {
1322 return Box::pin(async move {
1323 self.create_stream_with_admission_inner(request, placement, admission)
1324 });
1325 }
1326 let command = GroupWriteCommand::from(request);
1327 Box::pin(async move {
1328 match self.apply_committed_write(command, placement)? {
1329 GroupWriteResponse::CreateStream(response) => Ok(response),
1330 other => Err(GroupEngineError::new(format!(
1331 "unexpected create stream write response: {other:?}"
1332 ))),
1333 }
1334 })
1335 }
1336
1337 fn create_stream_external<'a>(
1338 &'a mut self,
1339 request: CreateStreamExternalRequest,
1340 placement: ShardPlacement,
1341 ) -> GroupCreateStreamFuture<'a> {
1342 Box::pin(async move {
1343 if let Some(cold_store) = self.cold_store.as_ref() {
1344 let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1345 write_external_segment_index_pages(
1346 &store,
1347 &request.stream_id,
1348 0,
1349 &request.initial_payload,
1350 )
1351 .await
1352 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1353 }
1354 let command = GroupWriteCommand::from(request);
1355 match self.apply_committed_write(command, placement)? {
1356 GroupWriteResponse::CreateStream(response) => Ok(response),
1357 other => Err(GroupEngineError::new(format!(
1358 "unexpected external create stream write response: {other:?}"
1359 ))),
1360 }
1361 })
1362 }
1363
1364 fn read_stream<'a>(
1365 &'a mut self,
1366 request: ReadStreamRequest,
1367 placement: ShardPlacement,
1368 ) -> GroupReadStreamFuture<'a> {
1369 Box::pin(async move {
1370 self.read_stream_parts(request, placement)
1371 .await?
1372 .into_response()
1373 .await
1374 })
1375 }
1376
1377 fn read_stream_parts<'a>(
1378 &'a mut self,
1379 request: ReadStreamRequest,
1380 placement: ShardPlacement,
1381 ) -> GroupReadStreamPartsFuture<'a> {
1382 Box::pin(async move {
1383 let stream_id = request.stream_id.clone();
1384 let plan = self.read_stream_plan(&request, placement)?;
1385 Ok(GroupReadStreamParts::from_plan(
1386 placement,
1387 stream_id,
1388 plan,
1389 self.cold_store(),
1390 self.cold_index_cache.clone(),
1391 ))
1392 })
1393 }
1394
1395 fn publish_snapshot<'a>(
1396 &'a mut self,
1397 request: PublishSnapshotRequest,
1398 placement: ShardPlacement,
1399 ) -> GroupPublishSnapshotFuture<'a> {
1400 Box::pin(async move {
1401 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1402 let command = GroupWriteCommand::from(request);
1403 match self.apply_committed_write(command, placement)? {
1404 GroupWriteResponse::PublishSnapshot(response) => Ok(response),
1405 other => Err(GroupEngineError::new(format!(
1406 "unexpected publish snapshot write response: {other:?}"
1407 ))),
1408 }
1409 })
1410 }
1411
1412 fn advance_retention<'a>(
1413 &'a mut self,
1414 request: AdvanceRetentionRequest,
1415 placement: ShardPlacement,
1416 ) -> GroupAdvanceRetentionFuture<'a> {
1417 Box::pin(async move {
1418 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1419 let command = GroupWriteCommand::from(request);
1420 match self.apply_committed_write(command, placement)? {
1421 GroupWriteResponse::AdvanceRetention(response) => Ok(response),
1422 other => Err(GroupEngineError::new(format!(
1423 "unexpected advance retention write response: {other:?}"
1424 ))),
1425 }
1426 })
1427 }
1428
1429 fn import_group_state<'a>(
1430 &'a mut self,
1431 request: ImportGroupStateRequest,
1432 placement: ShardPlacement,
1433 ) -> crate::GroupImportGroupStateFuture<'a> {
1434 Box::pin(async move {
1435 let command = GroupWriteCommand::from(StreamCommand::from(request));
1436 match self.apply_committed_write(command, placement)? {
1437 GroupWriteResponse::ImportGroupState(response) => Ok(response),
1438 other => Err(GroupEngineError::new(format!(
1439 "unexpected group state import response: {other:?}"
1440 ))),
1441 }
1442 })
1443 }
1444
1445 fn set_bucket_quota<'a>(
1446 &'a mut self,
1447 request: SetBucketQuotaRequest,
1448 placement: ShardPlacement,
1449 ) -> GroupSetBucketQuotaFuture<'a> {
1450 Box::pin(async move {
1451 let command = GroupWriteCommand::from(request);
1452 match self.apply_committed_write(command, placement)? {
1453 GroupWriteResponse::SetBucketQuota(response) => Ok(response),
1454 other => Err(GroupEngineError::new(format!(
1455 "unexpected set bucket quota write response: {other:?}"
1456 ))),
1457 }
1458 })
1459 }
1460
1461 fn read_snapshot<'a>(
1462 &'a mut self,
1463 request: ReadSnapshotRequest,
1464 placement: ShardPlacement,
1465 ) -> GroupReadSnapshotFuture<'a> {
1466 Box::pin(async move {
1467 self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1468 let snapshot = match request.snapshot_offset {
1469 Some(offset) => self
1470 .state_machine
1471 .read_snapshot(&request.stream_id, offset)
1472 .map_err(stream_response_error)?,
1473 None => self
1474 .state_machine
1475 .latest_snapshot(&request.stream_id)
1476 .map_err(stream_response_error)?
1477 .ok_or_else(|| {
1478 GroupEngineError::stream(
1479 StreamErrorCode::SnapshotNotFound,
1480 format!("stream '{}' has no visible snapshot", request.stream_id),
1481 )
1482 })?,
1483 };
1484 let tail_offset = self
1485 .state_machine
1486 .head_at(&request.stream_id, request.now_ms)
1487 .map(|metadata| metadata.tail_offset)
1488 .unwrap_or(snapshot.offset);
1489 Ok(ReadSnapshotResponse {
1490 placement,
1491 snapshot_offset: snapshot.offset,
1492 next_offset: snapshot.offset,
1493 content_type: snapshot.content_type,
1494 snapshot_digest: snapshot.digest,
1495 payload: snapshot.payload,
1496 up_to_date: snapshot.offset == tail_offset,
1497 record_range: self
1498 .state_machine
1499 .record_range(&request.stream_id)
1500 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?,
1501 })
1502 })
1503 }
1504
1505 fn delete_snapshot<'a>(
1506 &'a mut self,
1507 request: DeleteSnapshotRequest,
1508 placement: ShardPlacement,
1509 ) -> GroupDeleteSnapshotFuture<'a> {
1510 Box::pin(async move {
1511 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1512 match self
1513 .state_machine
1514 .delete_snapshot(&request.stream_id, request.snapshot_offset)
1515 {
1516 StreamResponse::Error {
1517 code,
1518 message,
1519 next_offset,
1520 context,
1521 } => Err(GroupEngineError::stream_with_context(
1522 code,
1523 message,
1524 next_offset,
1525 context,
1526 )),
1527 other => Err(GroupEngineError::new(format!(
1528 "unexpected delete snapshot response: {other:?}"
1529 ))),
1530 }
1531 })
1532 }
1533
1534 fn bootstrap_stream<'a>(
1535 &'a mut self,
1536 request: BootstrapStreamRequest,
1537 placement: ShardPlacement,
1538 ) -> GroupBootstrapStreamFuture<'a> {
1539 Box::pin(async move {
1540 self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1541 let plan = self
1542 .state_machine
1543 .bootstrap_plan(&request.stream_id)
1544 .map_err(stream_response_error)?;
1545 let snapshot_offset = plan.snapshot.as_ref().map(|snapshot| snapshot.offset);
1546 let snapshot_content_type = plan
1547 .snapshot
1548 .as_ref()
1549 .map(|snapshot| snapshot.content_type.clone())
1550 .unwrap_or_else(|| DEFAULT_CONTENT_TYPE.to_owned());
1551 let snapshot_payload = plan
1552 .snapshot
1553 .as_ref()
1554 .map(|snapshot| snapshot.payload.clone())
1555 .unwrap_or_default();
1556 let updates = self
1557 .bootstrap_updates(
1558 &request.stream_id,
1559 &plan.updates,
1560 &plan.content_type,
1561 request.now_ms,
1562 )
1563 .await?;
1564 Ok(BootstrapStreamResponse {
1565 placement,
1566 snapshot_offset,
1567 snapshot_content_type,
1568 snapshot_payload,
1569 updates,
1570 next_offset: plan.next_offset,
1571 up_to_date: plan.up_to_date,
1572 closed: plan.closed,
1573 record_range: self
1574 .state_machine
1575 .record_range(&request.stream_id)
1576 .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?,
1577 })
1578 })
1579 }
1580
1581 fn touch_stream_access<'a>(
1582 &'a mut self,
1583 stream_id: BucketStreamId,
1584 now_ms: u64,
1585 renew_ttl: bool,
1586 placement: ShardPlacement,
1587 ) -> GroupTouchStreamAccessFuture<'a> {
1588 Box::pin(async move { self.apply_access_command(stream_id, now_ms, renew_ttl, placement) })
1589 }
1590
1591 fn head_stream<'a>(
1592 &'a mut self,
1593 request: HeadStreamRequest,
1594 placement: ShardPlacement,
1595 ) -> GroupHeadStreamFuture<'a> {
1596 Box::pin(async move {
1597 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1598 self.head_stream_after_access(&request, placement)
1599 })
1600 }
1601
1602 fn bucket_usage<'a>(&'a mut self, _placement: ShardPlacement) -> GroupBucketUsageFuture<'a> {
1603 Box::pin(async move { Ok(self.state_machine.bucket_usage_report()) })
1604 }
1605
1606 fn get_stream_attrs<'a>(
1607 &'a mut self,
1608 request: GetStreamAttrsRequest,
1609 placement: ShardPlacement,
1610 ) -> GroupGetStreamAttrsFuture<'a> {
1611 Box::pin(async move {
1612 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1613 self.get_stream_attrs_after_access(&request, placement)
1614 })
1615 }
1616
1617 fn update_stream_attrs<'a>(
1618 &'a mut self,
1619 request: UpdateStreamAttrsRequest,
1620 placement: ShardPlacement,
1621 ) -> GroupUpdateStreamAttrsFuture<'a> {
1622 Box::pin(async move {
1623 match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
1624 GroupWriteResponse::UpdateStreamAttrs(response) => Ok(response),
1625 other => Err(GroupEngineError::new(format!(
1626 "unexpected update stream attrs write response: {other:?}"
1627 ))),
1628 }
1629 })
1630 }
1631
1632 fn close_stream<'a>(
1633 &'a mut self,
1634 request: CloseStreamRequest,
1635 placement: ShardPlacement,
1636 ) -> GroupCloseStreamFuture<'a> {
1637 Box::pin(async move {
1638 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1639 let command = GroupWriteCommand::from(request);
1640 match self.apply_committed_write(command, placement)? {
1641 GroupWriteResponse::CloseStream(response) => Ok(response),
1642 other => Err(GroupEngineError::new(format!(
1643 "unexpected close stream write response: {other:?}"
1644 ))),
1645 }
1646 })
1647 }
1648
1649 fn delete_stream<'a>(
1650 &'a mut self,
1651 request: DeleteStreamRequest,
1652 placement: ShardPlacement,
1653 ) -> GroupDeleteStreamFuture<'a> {
1654 let command = GroupWriteCommand::from(request);
1655 Box::pin(async move {
1656 match self.apply_committed_write(command, placement)? {
1657 GroupWriteResponse::DeleteStream(response) => Ok(response),
1658 other => Err(GroupEngineError::new(format!(
1659 "unexpected delete stream write response: {other:?}"
1660 ))),
1661 }
1662 })
1663 }
1664
1665 fn purge_bucket<'a>(
1666 &'a mut self,
1667 bucket_id: String,
1668 placement: ShardPlacement,
1669 ) -> GroupPurgeBucketFuture<'a> {
1670 Box::pin(async move {
1671 match self.apply_committed_write(
1672 GroupWriteCommand::Stream(StreamCommand::PurgeBucket { bucket_id }),
1673 placement,
1674 )? {
1675 GroupWriteResponse::PurgeBucket(response) => Ok(response),
1676 other => Err(GroupEngineError::new(format!(
1677 "unexpected purge bucket write response: {other:?}"
1678 ))),
1679 }
1680 })
1681 }
1682
1683 fn ack_cold_gc<'a>(
1684 &'a mut self,
1685 up_to_seq: u64,
1686 placement: ShardPlacement,
1687 ) -> GroupAckColdGcFuture<'a> {
1688 Box::pin(async move {
1689 match self.apply_committed_write(
1690 GroupWriteCommand::Stream(StreamCommand::AckColdGc { up_to_seq }),
1691 placement,
1692 )? {
1693 GroupWriteResponse::AckColdGc(response) => Ok(response),
1694 other => Err(GroupEngineError::new(format!(
1695 "unexpected ack cold gc write response: {other:?}"
1696 ))),
1697 }
1698 })
1699 }
1700
1701 fn plan_cold_gc<'a>(
1702 &'a mut self,
1703 max: usize,
1704 _placement: ShardPlacement,
1705 ) -> GroupPlanColdGcFuture<'a> {
1706 let entries = self.state_machine.pending_cold_gc_batch(max);
1707 Box::pin(async move { Ok(entries) })
1708 }
1709
1710 fn append<'a>(
1711 &'a mut self,
1712 request: AppendRequest,
1713 placement: ShardPlacement,
1714 admission: ColdWriteAdmission,
1715 ) -> GroupAppendFuture<'a> {
1716 if admission.is_enabled() {
1717 return Box::pin(async move {
1718 self.append_with_admission_inner(request, placement, admission)
1719 });
1720 }
1721 Box::pin(async move {
1722 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1723 let command = GroupWriteCommand::from(request);
1724 match self.apply_committed_write(command, placement)? {
1725 GroupWriteResponse::Append(response) => Ok(response),
1726 other => Err(GroupEngineError::new(format!(
1727 "unexpected append write response: {other:?}"
1728 ))),
1729 }
1730 })
1731 }
1732
1733 fn append_transaction<'a>(
1734 &'a mut self,
1735 request: AppendTransactionRequest,
1736 placement: ShardPlacement,
1737 admission: ColdWriteAdmission,
1738 ) -> GroupAppendTransactionFuture<'a> {
1739 Box::pin(async move {
1740 let Some(first) = request.operations.first() else {
1741 return Err(GroupEngineError::new(
1742 "append transaction must contain at least one operation",
1743 ));
1744 };
1745 self.check_cold_write_admission_bytes(
1746 &first.stream_id,
1747 admission,
1748 request.payload_bytes(),
1749 )?;
1750 let command = GroupWriteCommand::Transaction {
1751 commands: request
1752 .operations
1753 .into_iter()
1754 .map(StreamCommand::from)
1755 .collect(),
1756 };
1757 let GroupWriteResponse::Batch(items) =
1758 self.apply_committed_write(command, placement)?
1759 else {
1760 return Err(GroupEngineError::new(
1761 "unexpected append transaction write response",
1762 ));
1763 };
1764 let items = items
1765 .into_iter()
1766 .map(|item| match item? {
1767 GroupWriteResponse::Append(response) => Ok(response),
1768 other => Err(GroupEngineError::new(format!(
1769 "unexpected append transaction item response: {other:?}"
1770 ))),
1771 })
1772 .collect::<Result<Vec<_>, _>>()?;
1773 Ok(AppendTransactionResponse { placement, items })
1774 })
1775 }
1776
1777 fn append_external<'a>(
1778 &'a mut self,
1779 request: AppendExternalRequest,
1780 placement: ShardPlacement,
1781 ) -> GroupAppendFuture<'a> {
1782 Box::pin(async move {
1783 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1784 if let Some(cold_store) = self.cold_store.as_ref() {
1785 let start_offset = self
1786 .state_machine
1787 .head(&request.stream_id)
1788 .map(|metadata| metadata.tail_offset)
1789 .ok_or_else(|| {
1790 GroupEngineError::stream(
1791 ursula_stream::StreamErrorCode::StreamNotFound,
1792 format!("stream '{}' does not exist", request.stream_id),
1793 )
1794 })?;
1795 let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1796 write_external_segment_index_pages(
1797 &store,
1798 &request.stream_id,
1799 start_offset,
1800 &request.payload,
1801 )
1802 .await
1803 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1804 }
1805 let command = GroupWriteCommand::from(request);
1806 match self.apply_committed_write(command, placement)? {
1807 GroupWriteResponse::Append(response) => Ok(response),
1808 other => Err(GroupEngineError::new(format!(
1809 "unexpected external append write response: {other:?}"
1810 ))),
1811 }
1812 })
1813 }
1814
1815 fn append_batch<'a>(
1816 &'a mut self,
1817 request: AppendBatchRequest,
1818 placement: ShardPlacement,
1819 admission: ColdWriteAdmission,
1820 ) -> GroupAppendBatchFuture<'a> {
1821 if admission.is_enabled() {
1822 return Box::pin(async move {
1823 self.append_batch_with_admission_inner(request, placement, admission)
1824 });
1825 }
1826 Box::pin(async move {
1827 self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1828 let command = GroupWriteCommand::from(request);
1829 match self.apply_committed_write(command, placement)? {
1830 GroupWriteResponse::AppendBatch(response) => Ok(response),
1831 other => Err(GroupEngineError::new(format!(
1832 "unexpected append batch write response: {other:?}"
1833 ))),
1834 }
1835 })
1836 }
1837
1838 fn flush_cold<'a>(
1839 &'a mut self,
1840 request: FlushColdRequest,
1841 placement: ShardPlacement,
1842 ) -> GroupFlushColdFuture<'a> {
1843 Box::pin(async move {
1844 let mut index_rollback = None;
1845 if !request.chunk.shared_object
1846 && let Some(cold_store) = self.cold_store.as_ref()
1847 {
1848 let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1849 let rollback = write_cold_chunk_index_pages_with_rollback(
1850 &store,
1851 &request.stream_id,
1852 &request.chunk,
1853 )
1854 .await
1855 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1856 index_rollback = Some((store, rollback));
1857 }
1858 let command = GroupWriteCommand::from(request);
1859 match self.apply_committed_write(command, placement) {
1860 Ok(GroupWriteResponse::FlushCold(response)) => Ok(response),
1861 Ok(other) => {
1862 if let Some((store, rollback)) = index_rollback {
1863 rollback_cold_index_pages(&store, rollback)
1864 .await
1865 .map_err(|err| GroupEngineError::new(err.to_string()))?;
1866 }
1867 Err(GroupEngineError::new(format!(
1868 "unexpected flush cold write response: {other:?}"
1869 )))
1870 }
1871 Err(err) => {
1872 if let Some((store, rollback)) = index_rollback {
1873 rollback_cold_index_pages(&store, rollback).await.map_err(
1874 |rollback_err| {
1875 GroupEngineError::new(format!(
1876 "rollback cold index after flush failure: {rollback_err}"
1877 ))
1878 },
1879 )?;
1880 }
1881 Err(err)
1882 }
1883 }
1884 })
1885 }
1886
1887 fn compact_cold<'a>(
1888 &'a mut self,
1889 request: CompactColdRequest,
1890 placement: ShardPlacement,
1891 ) -> GroupCompactColdFuture<'a> {
1892 Box::pin(async move {
1893 let mut index_rollback = None;
1894 if let Some(cold_store) = self.cold_store.as_ref() {
1895 let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1896 let Some(rollback) = replace_cold_chunk_index_pages_with_rollback(
1897 &store,
1898 &request.stream_id,
1899 &request.old_chunks,
1900 &request.replacement,
1901 )
1902 .await
1903 .map_err(|err| GroupEngineError::new(err.to_string()))?
1904 else {
1905 return Err(GroupEngineError::new(
1906 "cold compaction input no longer matches the cold index",
1907 ));
1908 };
1909 index_rollback = Some((store, rollback));
1910 }
1911 let command = GroupWriteCommand::from(request);
1912 let result = match self.apply_committed_write(command, placement) {
1913 Ok(GroupWriteResponse::CompactCold(response)) => Ok(response),
1914 Ok(other) => Err(GroupEngineError::new(format!(
1915 "unexpected compact cold write response: {other:?}"
1916 ))),
1917 Err(err) => Err(err),
1918 };
1919 if result.is_err()
1920 && let Some((store, rollback)) = index_rollback
1921 {
1922 rollback_cold_index_pages(&store, rollback)
1923 .await
1924 .map_err(|err| {
1925 GroupEngineError::new(format!(
1926 "rollback cold index after compaction failure: {err}"
1927 ))
1928 })?;
1929 }
1930 result
1931 })
1932 }
1933
1934 fn plan_cold_flush<'a>(
1935 &'a mut self,
1936 request: PlanColdFlushRequest,
1937 _placement: ShardPlacement,
1938 ) -> GroupPlanColdFlushFuture<'a> {
1939 Box::pin(async move {
1940 self.state_machine
1941 .plan_cold_flush(
1942 &request.stream_id,
1943 request.min_hot_bytes,
1944 request.max_flush_bytes,
1945 )
1946 .map_err(stream_response_error)
1947 })
1948 }
1949
1950 fn plan_next_cold_flush_batch<'a>(
1951 &'a mut self,
1952 request: PlanGroupColdFlushRequest,
1953 _placement: ShardPlacement,
1954 max_candidates: usize,
1955 ) -> GroupPlanNextColdFlushBatchFuture<'a> {
1956 Box::pin(async move {
1957 self.state_machine
1958 .plan_next_cold_flush_batch(
1959 request.min_hot_bytes,
1960 request.max_flush_bytes,
1961 request.max_batch_bytes,
1962 max_candidates,
1963 )
1964 .map_err(stream_response_error)
1965 })
1966 }
1967
1968 fn cold_hot_backlog<'a>(
1969 &'a mut self,
1970 stream_id: BucketStreamId,
1971 _placement: ShardPlacement,
1972 ) -> GroupColdHotBacklogFuture<'a> {
1973 Box::pin(async move { self.cold_hot_backlog_for(stream_id) })
1974 }
1975
1976 fn snapshot<'a>(&'a mut self, placement: ShardPlacement) -> GroupSnapshotFuture<'a> {
1977 Box::pin(async move { Ok(self.build_snapshot(placement)) })
1978 }
1979
1980 fn install_snapshot<'a>(
1981 &'a mut self,
1982 snapshot: GroupSnapshot,
1983 ) -> GroupInstallSnapshotFuture<'a> {
1984 Box::pin(async move { self.install_snapshot_inner(snapshot) })
1985 }
1986}
1987
1988#[derive(Debug, Clone, Default)]
1989pub struct InMemoryGroupEngineFactory {
1990 cold_store: Option<ColdStoreHandle>,
1991}
1992
1993impl InMemoryGroupEngineFactory {
1994 pub fn new() -> Self {
1995 Self::default()
1996 }
1997
1998 pub fn with_cold_store(cold_store: Option<ColdStoreHandle>) -> Self {
1999 Self { cold_store }
2000 }
2001}
2002
2003impl GroupEngineFactory for InMemoryGroupEngineFactory {
2004 fn create<'a>(
2005 &'a self,
2006 _placement: ShardPlacement,
2007 _metrics: GroupEngineMetrics,
2008 ) -> GroupEngineCreateFuture<'a> {
2009 Box::pin(async move {
2010 let mut engine = InMemoryGroupEngine::default();
2011 engine.set_cold_store(self.cold_store.clone());
2012 let engine: Box<dyn GroupEngine> = Box::new(engine);
2013 Ok(engine)
2014 })
2015 }
2016}
2017
2018pub(crate) fn compare_stream_ids(
2019 left: &BucketStreamId,
2020 right: &BucketStreamId,
2021) -> std::cmp::Ordering {
2022 left.bucket_id
2023 .cmp(&right.bucket_id)
2024 .then_with(|| left.stream_id.cmp(&right.stream_id))
2025}
2026pub(crate) fn ensure_bucket_exists(
2027 state_machine: &mut StreamStateMachine,
2028 stream_id: &BucketStreamId,
2029) -> Result<(), GroupEngineError> {
2030 if state_machine.bucket_exists(&stream_id.bucket_id) {
2031 return Ok(());
2032 }
2033
2034 match state_machine.apply(StreamCommand::CreateBucket {
2035 bucket_id: stream_id.bucket_id.clone(),
2036 }) {
2037 StreamResponse::BucketCreated { .. } | StreamResponse::BucketAlreadyExists { .. } => Ok(()),
2038 StreamResponse::Error {
2039 code,
2040 message,
2041 next_offset,
2042 context,
2043 } => Err(GroupEngineError::stream_with_context(
2044 code,
2045 message,
2046 next_offset,
2047 context,
2048 )),
2049 other => Err(GroupEngineError::new(format!(
2050 "unexpected create bucket response: {other:?}"
2051 ))),
2052 }
2053}
2054
2055fn command_stream_id(command: &StreamCommand) -> Option<BucketStreamId> {
2057 match command {
2058 StreamCommand::CreateBucket { .. }
2059 | StreamCommand::DeleteBucket { .. }
2060 | StreamCommand::PurgeBucket { .. }
2061 | StreamCommand::AckColdGc { .. }
2062 | StreamCommand::ImportSnapshot { .. }
2063 | StreamCommand::SetBucketQuota { .. } => None,
2064 StreamCommand::CreateStream { stream_id, .. }
2065 | StreamCommand::CreateExternal { stream_id, .. }
2066 | StreamCommand::Append { stream_id, .. }
2067 | StreamCommand::AppendExternal { stream_id, .. }
2068 | StreamCommand::AppendBatch { stream_id, .. }
2069 | StreamCommand::PublishSnapshot { stream_id, .. }
2070 | StreamCommand::AdvanceRetention { stream_id, .. }
2071 | StreamCommand::TouchStreamAccess { stream_id, .. }
2072 | StreamCommand::UpdateStreamAttrs { stream_id, .. }
2073 | StreamCommand::FlushCold { stream_id, .. }
2074 | StreamCommand::CompactCold { stream_id, .. }
2075 | StreamCommand::Close { stream_id, .. }
2076 | StreamCommand::DeleteStream { stream_id } => Some(stream_id.clone()),
2077 }
2078}
2079
2080fn command_producer(command: &StreamCommand) -> Option<ProducerRequest> {
2081 match command {
2082 StreamCommand::Close { producer, .. } => producer.clone(),
2083 _ => None,
2084 }
2085}
2086
2087fn require_response_stream_id(
2088 stream_id: Option<BucketStreamId>,
2089 response: &str,
2090) -> Result<BucketStreamId, GroupEngineError> {
2091 stream_id.ok_or_else(|| {
2092 GroupEngineError::new(format!(
2093 "{response} response for a command without a stream id"
2094 ))
2095 })
2096}
2097
2098pub(crate) fn stream_response_error(response: StreamResponse) -> GroupEngineError {
2099 match response {
2100 StreamResponse::Error {
2101 code,
2102 message,
2103 next_offset,
2104 context,
2105 } => GroupEngineError::stream_with_context(code, message, next_offset, context),
2106 other => GroupEngineError::new(format!("unexpected stream response error: {other:?}")),
2107 }
2108}
2109
2110pub(crate) fn restore_stream_append_counts(
2111 counts: Vec<StreamAppendCount>,
2112 snapshot_stream_ids: &HashSet<BucketStreamId>,
2113) -> Result<HashMap<BucketStreamId, u64>, GroupEngineError> {
2114 let mut restored = HashMap::with_capacity(counts.len());
2115 for count in counts {
2116 if !snapshot_stream_ids.contains(&count.stream_id) {
2117 return Err(GroupEngineError::new(format!(
2118 "append count references missing snapshot stream '{}'",
2119 count.stream_id
2120 )));
2121 }
2122 if restored
2123 .insert(count.stream_id.clone(), count.append_count)
2124 .is_some()
2125 {
2126 return Err(GroupEngineError::new(format!(
2127 "snapshot contains duplicate append count for stream '{}'",
2128 count.stream_id
2129 )));
2130 }
2131 }
2132 Ok(restored)
2133}