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