Skip to main content

ursula_runtime/engine/
in_memory.rs

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