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