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