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 append_payload(
1015        &mut self,
1016        input: AppendPayloadInput<'_>,
1017        placement: ShardPlacement,
1018    ) -> Result<AppendResponse, GroupEngineError> {
1019        let AppendPayloadInput {
1020            stream_id,
1021            content_type,
1022            payload,
1023            close_after,
1024            stream_seq,
1025            producer,
1026            now_ms,
1027        } = input;
1028        let stream_count_key = stream_id.clone();
1029        let response = self.state_machine.append_borrowed(AppendStreamInput {
1030            stream_id,
1031            content_type,
1032            payload,
1033            close_after,
1034            stream_seq,
1035            producer,
1036            now_ms,
1037        });
1038        match response {
1039            StreamResponse::Appended {
1040                offset,
1041                next_offset,
1042                closed,
1043                deduplicated,
1044                producer,
1045                ..
1046            } => {
1047                let stream_append_count = self
1048                    .stream_append_counts
1049                    .entry(stream_count_key)
1050                    .or_insert(0);
1051                if !deduplicated {
1052                    self.commit_index += 1;
1053                    *stream_append_count += 1;
1054                }
1055                Ok(AppendResponse {
1056                    placement,
1057                    start_offset: offset,
1058                    next_offset,
1059                    stream_append_count: *stream_append_count,
1060                    group_commit_index: self.commit_index,
1061                    closed,
1062                    deduplicated,
1063                    producer,
1064                })
1065            }
1066            StreamResponse::Error {
1067                code,
1068                message,
1069                next_offset,
1070                context,
1071            } => Err(GroupEngineError::stream_with_context(
1072                code,
1073                message,
1074                next_offset,
1075                context,
1076            )),
1077            other => Err(GroupEngineError::new(format!(
1078                "unexpected append response: {other:?}"
1079            ))),
1080        }
1081    }
1082
1083    pub fn read_stream_plan(
1084        &mut self,
1085        request: &ReadStreamRequest,
1086        placement: ShardPlacement,
1087    ) -> Result<StreamReadPlan, GroupEngineError> {
1088        self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1089        self.read_stream_plan_after_access(request)
1090    }
1091
1092    pub fn read_stream_plan_after_access(
1093        &self,
1094        request: &ReadStreamRequest,
1095    ) -> Result<StreamReadPlan, GroupEngineError> {
1096        self.state_machine
1097            .read_plan_at(
1098                &request.stream_id,
1099                request.offset,
1100                request.max_len,
1101                request.now_ms,
1102            )
1103            .map_err(stream_response_error)
1104    }
1105
1106    pub fn head_stream_after_access(
1107        &mut self,
1108        request: &HeadStreamRequest,
1109        placement: ShardPlacement,
1110    ) -> Result<HeadStreamResponse, GroupEngineError> {
1111        let Some(metadata) = self
1112            .state_machine
1113            .head_at(&request.stream_id, request.now_ms)
1114        else {
1115            return Err(GroupEngineError::stream(
1116                StreamErrorCode::StreamNotFound,
1117                format!("stream '{}' does not exist", request.stream_id),
1118            ));
1119        };
1120        let content_type = metadata.content_type.clone();
1121        let tail_offset = metadata.tail_offset;
1122        let closed = metadata.status == ursula_stream::StreamStatus::Closed;
1123        let stream_ttl_seconds = metadata.stream_ttl_seconds;
1124        let stream_expires_at_ms = metadata.stream_expires_at_ms;
1125        let _ = metadata;
1126        Ok(HeadStreamResponse {
1127            placement,
1128            content_type,
1129            tail_offset,
1130            cold_hot_start_offset: self.state_machine.hot_start_offset(&request.stream_id),
1131            closed,
1132            stream_ttl_seconds,
1133            stream_expires_at_ms,
1134            snapshot_offset: self
1135                .state_machine
1136                .latest_snapshot(&request.stream_id)
1137                .map_err(stream_response_error)?
1138                .map(|snapshot| snapshot.offset),
1139            integrity: self
1140                .state_machine
1141                .integrity_snapshot(&request.stream_id)
1142                .map_err(stream_response_error)?,
1143        })
1144    }
1145
1146    pub fn get_stream_attrs_after_access(
1147        &mut self,
1148        request: &GetStreamAttrsRequest,
1149        placement: ShardPlacement,
1150    ) -> Result<GetStreamAttrsResponse, GroupEngineError> {
1151        if self
1152            .state_machine
1153            .head_at(&request.stream_id, request.now_ms)
1154            .is_none()
1155        {
1156            return Err(GroupEngineError::stream(
1157                StreamErrorCode::StreamNotFound,
1158                format!("stream '{}' does not exist", request.stream_id),
1159            ));
1160        }
1161        Ok(GetStreamAttrsResponse {
1162            placement,
1163            attrs: self.state_machine.stream_attrs(&request.stream_id).cloned(),
1164        })
1165    }
1166
1167    pub async fn read_payload_from_plan(
1168        cold_store: Option<&ColdStoreHandle>,
1169        cold_index_cache: Option<&Arc<ColdIndexPageCache<ColdStoreColdIndexPageStore>>>,
1170        stream_id: &BucketStreamId,
1171        plan: &StreamReadPlan,
1172    ) -> Result<Vec<u8>, GroupEngineError> {
1173        let mut payload = Vec::new();
1174        for segment in &plan.segments {
1175            match segment {
1176                StreamReadSegment::Hot(bytes) => payload.extend_from_slice(bytes),
1177                StreamReadSegment::ColdIndex(segment) => {
1178                    let Some(cold_store) = cold_store else {
1179                        return Err(GroupEngineError::stream_with_next_offset(
1180                            StreamErrorCode::InvalidColdFlush,
1181                            format!("stream '{stream_id}' read requires object payload store"),
1182                            Some(plan.next_offset),
1183                        ));
1184                    };
1185                    let Some(cache) = cold_index_cache else {
1186                        return Err(GroupEngineError::stream_with_next_offset(
1187                            StreamErrorCode::InvalidColdFlush,
1188                            format!("stream '{stream_id}' read requires cold index page cache"),
1189                            Some(plan.next_offset),
1190                        ));
1191                    };
1192                    let objects = cache
1193                        .object_segments_for_read(stream_id, segment)
1194                        .await
1195                        .map_err(|err| GroupEngineError::new(err.to_string()))?;
1196                    let segment_end = segment
1197                        .read_start_offset
1198                        .saturating_add(u64::try_from(segment.len).expect("read len fits u64"));
1199                    for object in objects {
1200                        let start = object.start_offset.max(segment.read_start_offset);
1201                        let end = object.end_offset.min(segment_end);
1202                        if start >= end {
1203                            continue;
1204                        }
1205                        let bytes = cold_store
1206                            .read_object_range_for_stream(
1207                                stream_id,
1208                                &object,
1209                                start,
1210                                usize::try_from(end - start).expect("object read len fits usize"),
1211                            )
1212                            .await
1213                            .map_err(|err| GroupEngineError::new(err.to_string()))?;
1214                        payload.extend_from_slice(&bytes);
1215                    }
1216                }
1217                StreamReadSegment::Object(segment) => {
1218                    let Some(cold_store) = cold_store else {
1219                        return Err(GroupEngineError::stream_with_next_offset(
1220                            StreamErrorCode::InvalidColdFlush,
1221                            format!("stream '{stream_id}' read requires object payload store"),
1222                            Some(plan.next_offset),
1223                        ));
1224                    };
1225                    let bytes = cold_store
1226                        .read_object_range_for_stream(
1227                            stream_id,
1228                            &segment.object,
1229                            segment.read_start_offset,
1230                            segment.len,
1231                        )
1232                        .await
1233                        .map_err(|err| GroupEngineError::new(err.to_string()))?;
1234                    payload.extend_from_slice(&bytes);
1235                }
1236            }
1237        }
1238        Ok(payload)
1239    }
1240
1241    pub(crate) async fn read_own_payload_from_plan(
1242        &self,
1243        stream_id: &BucketStreamId,
1244        plan: &StreamReadPlan,
1245    ) -> Result<Vec<u8>, GroupEngineError> {
1246        Self::read_payload_from_plan(
1247            self.cold_store.as_ref(),
1248            self.cold_index_cache.as_ref(),
1249            stream_id,
1250            plan,
1251        )
1252        .await
1253    }
1254
1255    pub(crate) async fn bootstrap_updates(
1256        &self,
1257        stream_id: &BucketStreamId,
1258        records: &[StreamMessageRecord],
1259        content_type: &str,
1260        now_ms: u64,
1261    ) -> Result<Vec<BootstrapUpdate>, GroupEngineError> {
1262        let mut updates = Vec::with_capacity(records.len());
1263        for record in records {
1264            let len = usize::try_from(record.end_offset - record.start_offset).map_err(|_| {
1265                GroupEngineError::stream(
1266                    StreamErrorCode::InvalidSnapshot,
1267                    format!(
1268                        "bootstrap message [{}..{}) for stream '{stream_id}' is too large",
1269                        record.start_offset, record.end_offset
1270                    ),
1271                )
1272            })?;
1273            let plan = self
1274                .state_machine
1275                .read_plan_at(stream_id, record.start_offset, len, now_ms)
1276                .map_err(stream_response_error)?;
1277            let payload = self.read_own_payload_from_plan(stream_id, &plan).await?;
1278            updates.push(BootstrapUpdate {
1279                start_offset: record.start_offset,
1280                next_offset: record.end_offset,
1281                content_type: content_type.to_owned(),
1282                payload,
1283            });
1284        }
1285        Ok(updates)
1286    }
1287
1288    pub(crate) fn build_snapshot(&self, placement: ShardPlacement) -> GroupSnapshot {
1289        let stream_snapshot = self.state_machine.snapshot();
1290        let stream_append_counts = self.stream_append_counts_snapshot(&stream_snapshot);
1291        GroupSnapshot {
1292            placement,
1293            group_commit_index: self.commit_index,
1294            stream_snapshot,
1295            stream_append_counts,
1296        }
1297    }
1298
1299    pub(crate) fn stream_append_counts_snapshot(
1300        &self,
1301        stream_snapshot: &ursula_stream::StreamSnapshot,
1302    ) -> Vec<StreamAppendCount> {
1303        // Only emit append counts for streams actually present in the snapshot.
1304        // A deleted/expired stream can leave a stale entry in the runtime map;
1305        // emitting it would make every follower's `install_snapshot` fail the
1306        // `restore_stream_append_counts` consistency check, so a lagging node
1307        // could never catch up (and leadership transfer, which catches the
1308        // target up via a snapshot, could never complete).
1309        let live: HashSet<&BucketStreamId> = stream_snapshot
1310            .streams
1311            .iter()
1312            .map(|entry| &entry.metadata.stream_id)
1313            .collect();
1314        let mut counts = self
1315            .stream_append_counts
1316            .iter()
1317            .filter(|(stream_id, _)| live.contains(stream_id))
1318            .map(|(stream_id, append_count)| StreamAppendCount {
1319                stream_id: stream_id.clone(),
1320                append_count: *append_count,
1321            })
1322            .collect::<Vec<_>>();
1323        counts.sort_by(|left, right| compare_stream_ids(&left.stream_id, &right.stream_id));
1324        counts
1325    }
1326
1327    pub fn stream_tail_offset(&self, stream_id: &BucketStreamId) -> Option<u64> {
1328        self.state_machine
1329            .head(stream_id)
1330            .map(|metadata| metadata.tail_offset)
1331    }
1332
1333    pub(crate) fn install_snapshot_inner(
1334        &mut self,
1335        snapshot: GroupSnapshot,
1336    ) -> Result<(), GroupEngineError> {
1337        let GroupSnapshot {
1338            placement: _,
1339            group_commit_index,
1340            stream_snapshot,
1341            stream_append_counts,
1342        } = snapshot;
1343        self.install_snapshot_parts(group_commit_index, stream_snapshot, stream_append_counts)
1344    }
1345
1346    pub(crate) fn install_snapshot_parts(
1347        &mut self,
1348        group_commit_index: u64,
1349        stream_snapshot: StreamSnapshot,
1350        stream_append_counts: Vec<StreamAppendCount>,
1351    ) -> Result<(), GroupEngineError> {
1352        let stream_ids = stream_snapshot
1353            .streams
1354            .iter()
1355            .map(|entry| entry.metadata.stream_id.clone())
1356            .collect::<HashSet<_>>();
1357        let state_machine = StreamStateMachine::restore(stream_snapshot)
1358            .map_err(|err| GroupEngineError::new(format!("restore stream snapshot: {err}")))?;
1359        let stream_append_counts = restore_stream_append_counts(stream_append_counts, &stream_ids)?;
1360
1361        self.commit_index = group_commit_index;
1362        self.state_machine = state_machine;
1363        self.stream_append_counts = stream_append_counts;
1364        Ok(())
1365    }
1366}
1367
1368impl GroupEngine for InMemoryGroupEngine {
1369    fn create_stream<'a>(
1370        &'a mut self,
1371        request: CreateStreamRequest,
1372        placement: ShardPlacement,
1373    ) -> GroupCreateStreamFuture<'a> {
1374        let command = GroupWriteCommand::from(request);
1375        Box::pin(async move {
1376            match self.apply_committed_write(command, placement)? {
1377                GroupWriteResponse::CreateStream(response) => Ok(response),
1378                other => Err(GroupEngineError::new(format!(
1379                    "unexpected create stream write response: {other:?}"
1380                ))),
1381            }
1382        })
1383    }
1384
1385    fn create_stream_with_cold_admission<'a>(
1386        &'a mut self,
1387        request: CreateStreamRequest,
1388        placement: ShardPlacement,
1389        admission: ColdWriteAdmission,
1390    ) -> GroupCreateStreamFuture<'a> {
1391        if !admission.is_enabled() {
1392            return self.create_stream(request, placement);
1393        }
1394        Box::pin(
1395            async move { self.create_stream_with_admission_inner(request, placement, admission) },
1396        )
1397    }
1398
1399    fn create_stream_external<'a>(
1400        &'a mut self,
1401        request: CreateStreamExternalRequest,
1402        placement: ShardPlacement,
1403    ) -> GroupCreateStreamFuture<'a> {
1404        Box::pin(async move {
1405            if let Some(cold_store) = self.cold_store.as_ref() {
1406                let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1407                write_external_segment_index_pages(
1408                    &store,
1409                    &request.stream_id,
1410                    0,
1411                    &request.initial_payload,
1412                )
1413                .await
1414                .map_err(|err| GroupEngineError::new(err.to_string()))?;
1415            }
1416            let command = GroupWriteCommand::from(request);
1417            match self.apply_committed_write(command, placement)? {
1418                GroupWriteResponse::CreateStream(response) => Ok(response),
1419                other => Err(GroupEngineError::new(format!(
1420                    "unexpected external create stream write response: {other:?}"
1421                ))),
1422            }
1423        })
1424    }
1425
1426    fn read_stream<'a>(
1427        &'a mut self,
1428        request: ReadStreamRequest,
1429        placement: ShardPlacement,
1430    ) -> GroupReadStreamFuture<'a> {
1431        Box::pin(async move {
1432            self.read_stream_parts(request, placement)
1433                .await?
1434                .into_response()
1435                .await
1436        })
1437    }
1438
1439    fn read_stream_parts<'a>(
1440        &'a mut self,
1441        request: ReadStreamRequest,
1442        placement: ShardPlacement,
1443    ) -> GroupReadStreamPartsFuture<'a> {
1444        Box::pin(async move {
1445            let stream_id = request.stream_id.clone();
1446            let plan = self.read_stream_plan(&request, placement)?;
1447            Ok(GroupReadStreamParts::from_plan(
1448                placement,
1449                stream_id,
1450                plan,
1451                self.cold_store(),
1452                self.cold_index_cache.clone(),
1453            ))
1454        })
1455    }
1456
1457    fn publish_snapshot<'a>(
1458        &'a mut self,
1459        request: PublishSnapshotRequest,
1460        placement: ShardPlacement,
1461    ) -> GroupPublishSnapshotFuture<'a> {
1462        Box::pin(async move {
1463            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1464            let command = GroupWriteCommand::from(request);
1465            match self.apply_committed_write(command, placement)? {
1466                GroupWriteResponse::PublishSnapshot(response) => Ok(response),
1467                other => Err(GroupEngineError::new(format!(
1468                    "unexpected publish snapshot write response: {other:?}"
1469                ))),
1470            }
1471        })
1472    }
1473
1474    fn read_snapshot<'a>(
1475        &'a mut self,
1476        request: ReadSnapshotRequest,
1477        placement: ShardPlacement,
1478    ) -> GroupReadSnapshotFuture<'a> {
1479        Box::pin(async move {
1480            self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1481            let snapshot = match request.snapshot_offset {
1482                Some(offset) => self
1483                    .state_machine
1484                    .read_snapshot(&request.stream_id, offset)
1485                    .map_err(stream_response_error)?,
1486                None => self
1487                    .state_machine
1488                    .latest_snapshot(&request.stream_id)
1489                    .map_err(stream_response_error)?
1490                    .ok_or_else(|| {
1491                        GroupEngineError::stream(
1492                            StreamErrorCode::SnapshotNotFound,
1493                            format!("stream '{}' has no visible snapshot", request.stream_id),
1494                        )
1495                    })?,
1496            };
1497            let tail_offset = self
1498                .state_machine
1499                .head_at(&request.stream_id, request.now_ms)
1500                .map(|metadata| metadata.tail_offset)
1501                .unwrap_or(snapshot.offset);
1502            Ok(ReadSnapshotResponse {
1503                placement,
1504                snapshot_offset: snapshot.offset,
1505                next_offset: snapshot.offset,
1506                content_type: snapshot.content_type,
1507                payload: snapshot.payload,
1508                up_to_date: snapshot.offset == tail_offset,
1509            })
1510        })
1511    }
1512
1513    fn delete_snapshot<'a>(
1514        &'a mut self,
1515        request: DeleteSnapshotRequest,
1516        placement: ShardPlacement,
1517    ) -> GroupDeleteSnapshotFuture<'a> {
1518        Box::pin(async move {
1519            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1520            match self
1521                .state_machine
1522                .delete_snapshot(&request.stream_id, request.snapshot_offset)
1523            {
1524                StreamResponse::Error {
1525                    code,
1526                    message,
1527                    next_offset,
1528                    context,
1529                } => Err(GroupEngineError::stream_with_context(
1530                    code,
1531                    message,
1532                    next_offset,
1533                    context,
1534                )),
1535                other => Err(GroupEngineError::new(format!(
1536                    "unexpected delete snapshot response: {other:?}"
1537                ))),
1538            }
1539        })
1540    }
1541
1542    fn bootstrap_stream<'a>(
1543        &'a mut self,
1544        request: BootstrapStreamRequest,
1545        placement: ShardPlacement,
1546    ) -> GroupBootstrapStreamFuture<'a> {
1547        Box::pin(async move {
1548            self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1549            let plan = self
1550                .state_machine
1551                .bootstrap_plan(&request.stream_id)
1552                .map_err(stream_response_error)?;
1553            let snapshot_offset = plan.snapshot.as_ref().map(|snapshot| snapshot.offset);
1554            let snapshot_content_type = plan
1555                .snapshot
1556                .as_ref()
1557                .map(|snapshot| snapshot.content_type.clone())
1558                .unwrap_or_else(|| DEFAULT_CONTENT_TYPE.to_owned());
1559            let snapshot_payload = plan
1560                .snapshot
1561                .as_ref()
1562                .map(|snapshot| snapshot.payload.clone())
1563                .unwrap_or_default();
1564            let updates = self
1565                .bootstrap_updates(
1566                    &request.stream_id,
1567                    &plan.updates,
1568                    &plan.content_type,
1569                    request.now_ms,
1570                )
1571                .await?;
1572            Ok(BootstrapStreamResponse {
1573                placement,
1574                snapshot_offset,
1575                snapshot_content_type,
1576                snapshot_payload,
1577                updates,
1578                next_offset: plan.next_offset,
1579                up_to_date: plan.up_to_date,
1580                closed: plan.closed,
1581            })
1582        })
1583    }
1584
1585    fn touch_stream_access<'a>(
1586        &'a mut self,
1587        stream_id: BucketStreamId,
1588        now_ms: u64,
1589        renew_ttl: bool,
1590        placement: ShardPlacement,
1591    ) -> GroupTouchStreamAccessFuture<'a> {
1592        Box::pin(async move { self.apply_access_command(stream_id, now_ms, renew_ttl, placement) })
1593    }
1594
1595    fn add_fork_ref<'a>(
1596        &'a mut self,
1597        stream_id: BucketStreamId,
1598        now_ms: u64,
1599        placement: ShardPlacement,
1600    ) -> GroupForkRefFuture<'a> {
1601        Box::pin(async move {
1602            match self.apply_committed_write(
1603                GroupWriteCommand::AddForkRef { stream_id, now_ms },
1604                placement,
1605            )? {
1606                GroupWriteResponse::AddForkRef(response) => Ok(response),
1607                other => Err(GroupEngineError::new(format!(
1608                    "unexpected add fork ref write response: {other:?}"
1609                ))),
1610            }
1611        })
1612    }
1613
1614    fn release_fork_ref<'a>(
1615        &'a mut self,
1616        stream_id: BucketStreamId,
1617        placement: ShardPlacement,
1618    ) -> GroupForkRefFuture<'a> {
1619        Box::pin(async move {
1620            match self
1621                .apply_committed_write(GroupWriteCommand::ReleaseForkRef { stream_id }, placement)?
1622            {
1623                GroupWriteResponse::ReleaseForkRef(response) => Ok(response),
1624                other => Err(GroupEngineError::new(format!(
1625                    "unexpected release fork ref write response: {other:?}"
1626                ))),
1627            }
1628        })
1629    }
1630
1631    fn head_stream<'a>(
1632        &'a mut self,
1633        request: HeadStreamRequest,
1634        placement: ShardPlacement,
1635    ) -> GroupHeadStreamFuture<'a> {
1636        Box::pin(async move {
1637            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1638            self.head_stream_after_access(&request, placement)
1639        })
1640    }
1641
1642    fn get_stream_attrs<'a>(
1643        &'a mut self,
1644        request: GetStreamAttrsRequest,
1645        placement: ShardPlacement,
1646    ) -> GroupGetStreamAttrsFuture<'a> {
1647        Box::pin(async move {
1648            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1649            self.get_stream_attrs_after_access(&request, placement)
1650        })
1651    }
1652
1653    fn update_stream_attrs<'a>(
1654        &'a mut self,
1655        request: UpdateStreamAttrsRequest,
1656        placement: ShardPlacement,
1657    ) -> GroupUpdateStreamAttrsFuture<'a> {
1658        Box::pin(async move {
1659            match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
1660                GroupWriteResponse::UpdateStreamAttrs(response) => Ok(response),
1661                other => Err(GroupEngineError::new(format!(
1662                    "unexpected update stream attrs write response: {other:?}"
1663                ))),
1664            }
1665        })
1666    }
1667
1668    fn close_stream<'a>(
1669        &'a mut self,
1670        request: CloseStreamRequest,
1671        placement: ShardPlacement,
1672    ) -> GroupCloseStreamFuture<'a> {
1673        Box::pin(async move {
1674            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1675            let command = GroupWriteCommand::from(request);
1676            match self.apply_committed_write(command, placement)? {
1677                GroupWriteResponse::CloseStream(response) => Ok(response),
1678                other => Err(GroupEngineError::new(format!(
1679                    "unexpected close stream write response: {other:?}"
1680                ))),
1681            }
1682        })
1683    }
1684
1685    fn delete_stream<'a>(
1686        &'a mut self,
1687        request: DeleteStreamRequest,
1688        placement: ShardPlacement,
1689    ) -> GroupDeleteStreamFuture<'a> {
1690        let command = GroupWriteCommand::from(request);
1691        Box::pin(async move {
1692            match self.apply_committed_write(command, placement)? {
1693                GroupWriteResponse::DeleteStream(response) => Ok(response),
1694                other => Err(GroupEngineError::new(format!(
1695                    "unexpected delete stream write response: {other:?}"
1696                ))),
1697            }
1698        })
1699    }
1700
1701    fn ack_cold_gc<'a>(
1702        &'a mut self,
1703        up_to_seq: u64,
1704        placement: ShardPlacement,
1705    ) -> GroupAckColdGcFuture<'a> {
1706        Box::pin(async move {
1707            match self
1708                .apply_committed_write(GroupWriteCommand::AckColdGc { up_to_seq }, placement)?
1709            {
1710                GroupWriteResponse::AckColdGc(response) => Ok(response),
1711                other => Err(GroupEngineError::new(format!(
1712                    "unexpected ack cold gc write response: {other:?}"
1713                ))),
1714            }
1715        })
1716    }
1717
1718    fn plan_cold_gc<'a>(
1719        &'a mut self,
1720        max: usize,
1721        _placement: ShardPlacement,
1722    ) -> GroupPlanColdGcFuture<'a> {
1723        let entries = self.state_machine.pending_cold_gc_batch(max);
1724        Box::pin(async move { Ok(entries) })
1725    }
1726
1727    fn append<'a>(
1728        &'a mut self,
1729        request: AppendRequest,
1730        placement: ShardPlacement,
1731    ) -> GroupAppendFuture<'a> {
1732        Box::pin(async move {
1733            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1734            let command = GroupWriteCommand::from(request);
1735            match self.apply_committed_write(command, placement)? {
1736                GroupWriteResponse::Append(response) => Ok(response),
1737                other => Err(GroupEngineError::new(format!(
1738                    "unexpected append write response: {other:?}"
1739                ))),
1740            }
1741        })
1742    }
1743
1744    fn append_with_cold_admission<'a>(
1745        &'a mut self,
1746        request: AppendRequest,
1747        placement: ShardPlacement,
1748        admission: ColdWriteAdmission,
1749    ) -> GroupAppendFuture<'a> {
1750        if !admission.is_enabled() {
1751            return self.append(request, placement);
1752        }
1753        Box::pin(async move { self.append_with_admission_inner(request, placement, admission) })
1754    }
1755
1756    fn append_external<'a>(
1757        &'a mut self,
1758        request: AppendExternalRequest,
1759        placement: ShardPlacement,
1760    ) -> GroupAppendFuture<'a> {
1761        Box::pin(async move {
1762            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1763            if let Some(cold_store) = self.cold_store.as_ref() {
1764                let start_offset = self
1765                    .state_machine
1766                    .head(&request.stream_id)
1767                    .map(|metadata| metadata.tail_offset)
1768                    .ok_or_else(|| {
1769                        GroupEngineError::stream(
1770                            ursula_stream::StreamErrorCode::StreamNotFound,
1771                            format!("stream '{}' does not exist", request.stream_id),
1772                        )
1773                    })?;
1774                let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1775                write_external_segment_index_pages(
1776                    &store,
1777                    &request.stream_id,
1778                    start_offset,
1779                    &request.payload,
1780                )
1781                .await
1782                .map_err(|err| GroupEngineError::new(err.to_string()))?;
1783            }
1784            let command = GroupWriteCommand::from(request);
1785            match self.apply_committed_write(command, placement)? {
1786                GroupWriteResponse::Append(response) => Ok(response),
1787                other => Err(GroupEngineError::new(format!(
1788                    "unexpected external append write response: {other:?}"
1789                ))),
1790            }
1791        })
1792    }
1793
1794    fn append_batch<'a>(
1795        &'a mut self,
1796        request: AppendBatchRequest,
1797        placement: ShardPlacement,
1798    ) -> GroupAppendBatchFuture<'a> {
1799        Box::pin(async move {
1800            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1801            let command = GroupWriteCommand::from(request);
1802            match self.apply_committed_write(command, placement)? {
1803                GroupWriteResponse::AppendBatch(response) => Ok(response),
1804                other => Err(GroupEngineError::new(format!(
1805                    "unexpected append batch write response: {other:?}"
1806                ))),
1807            }
1808        })
1809    }
1810
1811    fn append_batch_with_cold_admission<'a>(
1812        &'a mut self,
1813        request: AppendBatchRequest,
1814        placement: ShardPlacement,
1815        admission: ColdWriteAdmission,
1816    ) -> GroupAppendBatchFuture<'a> {
1817        if !admission.is_enabled() {
1818            return self.append_batch(request, placement);
1819        }
1820        Box::pin(
1821            async move { self.append_batch_with_admission_inner(request, placement, admission) },
1822        )
1823    }
1824
1825    fn flush_cold<'a>(
1826        &'a mut self,
1827        request: FlushColdRequest,
1828        placement: ShardPlacement,
1829    ) -> GroupFlushColdFuture<'a> {
1830        Box::pin(async move {
1831            let mut index_rollback = None;
1832            if let Some(cold_store) = self.cold_store.as_ref() {
1833                let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1834                let rollback = write_cold_chunk_index_pages_with_rollback(
1835                    &store,
1836                    &request.stream_id,
1837                    &request.chunk,
1838                )
1839                .await
1840                .map_err(|err| GroupEngineError::new(err.to_string()))?;
1841                index_rollback = Some((store, rollback));
1842            }
1843            let command = GroupWriteCommand::from(request);
1844            match self.apply_committed_write(command, placement) {
1845                Ok(GroupWriteResponse::FlushCold(response)) => Ok(response),
1846                Ok(other) => {
1847                    if let Some((store, rollback)) = index_rollback {
1848                        rollback_cold_index_pages(&store, rollback)
1849                            .await
1850                            .map_err(|err| GroupEngineError::new(err.to_string()))?;
1851                    }
1852                    Err(GroupEngineError::new(format!(
1853                        "unexpected flush cold write response: {other:?}"
1854                    )))
1855                }
1856                Err(err) => {
1857                    if let Some((store, rollback)) = index_rollback {
1858                        rollback_cold_index_pages(&store, rollback).await.map_err(
1859                            |rollback_err| {
1860                                GroupEngineError::new(format!(
1861                                    "rollback cold index after flush failure: {rollback_err}"
1862                                ))
1863                            },
1864                        )?;
1865                    }
1866                    Err(err)
1867                }
1868            }
1869        })
1870    }
1871
1872    fn plan_cold_flush<'a>(
1873        &'a mut self,
1874        request: PlanColdFlushRequest,
1875        _placement: ShardPlacement,
1876    ) -> GroupPlanColdFlushFuture<'a> {
1877        Box::pin(async move {
1878            self.state_machine
1879                .plan_cold_flush(
1880                    &request.stream_id,
1881                    request.min_hot_bytes,
1882                    request.max_flush_bytes,
1883                )
1884                .map_err(stream_response_error)
1885        })
1886    }
1887
1888    fn plan_next_cold_flush<'a>(
1889        &'a mut self,
1890        request: PlanGroupColdFlushRequest,
1891        _placement: ShardPlacement,
1892    ) -> GroupPlanNextColdFlushFuture<'a> {
1893        Box::pin(async move {
1894            self.state_machine
1895                .plan_next_cold_flush_batch(request.min_hot_bytes, request.max_flush_bytes, 1)
1896                .map(|candidates| candidates.into_iter().next())
1897                .map_err(stream_response_error)
1898        })
1899    }
1900
1901    fn plan_next_cold_flush_batch<'a>(
1902        &'a mut self,
1903        request: PlanGroupColdFlushRequest,
1904        _placement: ShardPlacement,
1905        max_candidates: usize,
1906    ) -> GroupPlanNextColdFlushBatchFuture<'a> {
1907        Box::pin(async move {
1908            self.state_machine
1909                .plan_next_cold_flush_batch(
1910                    request.min_hot_bytes,
1911                    request.max_flush_bytes,
1912                    max_candidates,
1913                )
1914                .map_err(stream_response_error)
1915        })
1916    }
1917
1918    fn cold_hot_backlog<'a>(
1919        &'a mut self,
1920        stream_id: BucketStreamId,
1921        _placement: ShardPlacement,
1922    ) -> GroupColdHotBacklogFuture<'a> {
1923        Box::pin(async move { self.cold_hot_backlog_for(stream_id) })
1924    }
1925
1926    fn snapshot<'a>(&'a mut self, placement: ShardPlacement) -> GroupSnapshotFuture<'a> {
1927        Box::pin(async move { Ok(self.build_snapshot(placement)) })
1928    }
1929
1930    fn install_snapshot<'a>(
1931        &'a mut self,
1932        snapshot: GroupSnapshot,
1933    ) -> GroupInstallSnapshotFuture<'a> {
1934        Box::pin(async move { self.install_snapshot_inner(snapshot) })
1935    }
1936}
1937
1938#[derive(Debug, Clone, Default)]
1939pub struct InMemoryGroupEngineFactory {
1940    cold_store: Option<ColdStoreHandle>,
1941}
1942
1943impl InMemoryGroupEngineFactory {
1944    pub fn new() -> Self {
1945        Self::default()
1946    }
1947
1948    pub fn with_cold_store(cold_store: Option<ColdStoreHandle>) -> Self {
1949        Self { cold_store }
1950    }
1951}
1952
1953impl GroupEngineFactory for InMemoryGroupEngineFactory {
1954    fn create<'a>(
1955        &'a self,
1956        _placement: ShardPlacement,
1957        _metrics: GroupEngineMetrics,
1958    ) -> GroupEngineCreateFuture<'a> {
1959        Box::pin(async move {
1960            let mut engine = InMemoryGroupEngine::default();
1961            engine.set_cold_store(self.cold_store.clone());
1962            let engine: Box<dyn GroupEngine> = Box::new(engine);
1963            Ok(engine)
1964        })
1965    }
1966}
1967
1968pub(crate) fn compare_stream_ids(
1969    left: &BucketStreamId,
1970    right: &BucketStreamId,
1971) -> std::cmp::Ordering {
1972    left.bucket_id
1973        .cmp(&right.bucket_id)
1974        .then_with(|| left.stream_id.cmp(&right.stream_id))
1975}
1976pub(crate) fn ensure_bucket_exists(
1977    state_machine: &mut StreamStateMachine,
1978    stream_id: &BucketStreamId,
1979) -> Result<(), GroupEngineError> {
1980    if state_machine.bucket_exists(&stream_id.bucket_id) {
1981        return Ok(());
1982    }
1983
1984    match state_machine.apply(StreamCommand::CreateBucket {
1985        bucket_id: stream_id.bucket_id.clone(),
1986    }) {
1987        StreamResponse::BucketCreated { .. } | StreamResponse::BucketAlreadyExists { .. } => Ok(()),
1988        StreamResponse::Error {
1989            code,
1990            message,
1991            next_offset,
1992            context,
1993        } => Err(GroupEngineError::stream_with_context(
1994            code,
1995            message,
1996            next_offset,
1997            context,
1998        )),
1999        other => Err(GroupEngineError::new(format!(
2000            "unexpected create bucket response: {other:?}"
2001        ))),
2002    }
2003}
2004
2005pub(crate) fn stream_response_error(response: StreamResponse) -> GroupEngineError {
2006    match response {
2007        StreamResponse::Error {
2008            code,
2009            message,
2010            next_offset,
2011            context,
2012        } => GroupEngineError::stream_with_context(code, message, next_offset, context),
2013        other => GroupEngineError::new(format!("unexpected stream response error: {other:?}")),
2014    }
2015}
2016
2017pub(crate) fn restore_stream_append_counts(
2018    counts: Vec<StreamAppendCount>,
2019    snapshot_stream_ids: &HashSet<BucketStreamId>,
2020) -> Result<HashMap<BucketStreamId, u64>, GroupEngineError> {
2021    let mut restored = HashMap::with_capacity(counts.len());
2022    for count in counts {
2023        if !snapshot_stream_ids.contains(&count.stream_id) {
2024            return Err(GroupEngineError::new(format!(
2025                "append count references missing snapshot stream '{}'",
2026                count.stream_id
2027            )));
2028        }
2029        if restored
2030            .insert(count.stream_id.clone(), count.append_count)
2031            .is_some()
2032        {
2033            return Err(GroupEngineError::new(format!(
2034                "snapshot contains duplicate append count for stream '{}'",
2035                count.stream_id
2036            )));
2037        }
2038    }
2039    Ok(restored)
2040}