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                pending_cold_gc_entries,
617            } => {
618                self.commit_index += 1;
619                Ok(GroupWriteResponse::PurgeBucket(PurgeBucketResponse {
620                    placement,
621                    removed_streams,
622                    pending_cold_gc_entries,
623                    group_commit_index: self.commit_index,
624                }))
625            }
626            StreamResponse::Error {
627                code,
628                message,
629                next_offset,
630                context,
631            } => Err(GroupEngineError::stream_with_context(
632                code,
633                message,
634                next_offset,
635                context,
636            )),
637            other @ (StreamResponse::BucketCreated { .. }
638            | StreamResponse::BucketAlreadyExists { .. }
639            | StreamResponse::BucketDeleted { .. }) => Err(GroupEngineError::new(format!(
640                "unexpected group write response: {other:?}"
641            ))),
642        }
643    }
644
645    pub(crate) fn cold_hot_backlog_for(
646        &self,
647        stream_id: BucketStreamId,
648    ) -> Result<ColdHotBacklog, GroupEngineError> {
649        let stream_hot_bytes = self.state_machine.hot_payload_len(&stream_id).unwrap_or(0);
650        Ok(ColdHotBacklog {
651            stream_id,
652            stream_hot_bytes,
653            group_hot_bytes: self.state_machine.total_hot_payload_bytes(),
654        })
655    }
656
657    pub fn check_cold_write_admission_bytes(
658        &self,
659        stream_id: &BucketStreamId,
660        admission: ColdWriteAdmission,
661        incoming_bytes: u64,
662    ) -> Result<(), GroupEngineError> {
663        let Some(limit) = admission.max_hot_bytes_per_group else {
664            return Ok(());
665        };
666        if incoming_bytes == 0 {
667            return Ok(());
668        }
669        let before = self.state_machine.total_hot_payload_bytes();
670        let after = before.saturating_add(incoming_bytes);
671        if after <= limit {
672            return Ok(());
673        }
674        Err(GroupEngineError::cold_backpressure(
675            stream_id.clone(),
676            before,
677            after,
678            limit,
679        ))
680    }
681
682    pub(crate) fn create_stream_with_admission_inner(
683        &mut self,
684        request: CreateStreamRequest,
685        placement: ShardPlacement,
686        admission: ColdWriteAdmission,
687    ) -> Result<CreateStreamResponse, GroupEngineError> {
688        let stream_id = request.stream_id.clone();
689        if admission.is_enabled() {
690            let mut preview = self.clone();
691            let preview_response = match preview
692                .apply_committed_write(GroupWriteCommand::from(request.clone()), placement)?
693            {
694                GroupWriteResponse::CreateStream(response) => response,
695                other => {
696                    return Err(GroupEngineError::new(format!(
697                        "unexpected create stream preview response: {other:?}"
698                    )));
699                }
700            };
701            if !preview_response.already_exists {
702                self.check_cold_write_admission_bytes(
703                    &stream_id,
704                    admission,
705                    u64::try_from(request.initial_payload.len()).expect("payload len fits u64"),
706                )?;
707            }
708        }
709        let response =
710            match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
711                GroupWriteResponse::CreateStream(response) => response,
712                other => {
713                    return Err(GroupEngineError::new(format!(
714                        "unexpected create stream write response: {other:?}"
715                    )));
716                }
717            };
718        Ok(response)
719    }
720
721    pub(crate) fn append_with_admission_inner(
722        &mut self,
723        request: AppendRequest,
724        placement: ShardPlacement,
725        admission: ColdWriteAdmission,
726    ) -> Result<AppendResponse, GroupEngineError> {
727        let stream_id = request.stream_id.clone();
728        if admission.is_enabled() {
729            let mut preview = self.clone();
730            let preview_response = match preview
731                .apply_committed_write(GroupWriteCommand::from(request.clone()), placement)?
732            {
733                GroupWriteResponse::Append(response) => response,
734                other => {
735                    return Err(GroupEngineError::new(format!(
736                        "unexpected append preview response: {other:?}"
737                    )));
738                }
739            };
740            if !preview_response.deduplicated {
741                self.check_cold_write_admission_bytes(
742                    &stream_id,
743                    admission,
744                    u64::try_from(request.payload.len()).expect("payload len fits u64"),
745                )?;
746            }
747        }
748        let response =
749            match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
750                GroupWriteResponse::Append(response) => response,
751                other => {
752                    return Err(GroupEngineError::new(format!(
753                        "unexpected append write response: {other:?}"
754                    )));
755                }
756            };
757        Ok(response)
758    }
759
760    pub(crate) fn append_batch_with_admission_inner(
761        &mut self,
762        request: AppendBatchRequest,
763        placement: ShardPlacement,
764        admission: ColdWriteAdmission,
765    ) -> Result<GroupAppendBatchResponse, GroupEngineError> {
766        let stream_id = request.stream_id.clone();
767        let incoming_bytes = request
768            .payloads
769            .iter()
770            .map(|payload| u64::try_from(payload.len()).expect("payload len fits u64"))
771            .sum();
772        if admission.is_enabled() {
773            let mut preview = self.clone();
774            let preview_response = match preview
775                .apply_committed_write(GroupWriteCommand::from(request.clone()), placement)?
776            {
777                GroupWriteResponse::AppendBatch(response) => response,
778                other => {
779                    return Err(GroupEngineError::new(format!(
780                        "unexpected append batch preview response: {other:?}"
781                    )));
782                }
783            };
784            let mutates = preview_response
785                .items
786                .iter()
787                .any(|item| matches!(item, Ok(response) if !response.deduplicated));
788            if mutates {
789                self.check_cold_write_admission_bytes(&stream_id, admission, incoming_bytes)?;
790            }
791        }
792        let response =
793            match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
794                GroupWriteResponse::AppendBatch(response) => response,
795                other => {
796                    return Err(GroupEngineError::new(format!(
797                        "unexpected append batch write response: {other:?}"
798                    )));
799                }
800            };
801        Ok(response)
802    }
803
804    pub fn access_requires_write(
805        &self,
806        stream_id: &BucketStreamId,
807        now_ms: u64,
808        renew_ttl: bool,
809    ) -> Result<bool, GroupEngineError> {
810        self.state_machine
811            .access_requires_write(stream_id, now_ms, renew_ttl)
812            .map_err(stream_response_error)
813    }
814
815    pub(crate) fn apply_access_command(
816        &mut self,
817        stream_id: BucketStreamId,
818        now_ms: u64,
819        renew_ttl: bool,
820        placement: ShardPlacement,
821    ) -> Result<TouchStreamAccessResponse, GroupEngineError> {
822        match self.apply_committed_write(
823            GroupWriteCommand::Stream(StreamCommand::TouchStreamAccess {
824                stream_id,
825                now_ms,
826                renew_ttl,
827            }),
828            placement,
829        )? {
830            GroupWriteResponse::TouchStreamAccess(response) => Ok(response),
831            other => Err(GroupEngineError::new(format!(
832                "unexpected touch stream access write response: {other:?}"
833            ))),
834        }
835    }
836
837    pub(crate) fn ensure_stream_access(
838        &mut self,
839        stream_id: &BucketStreamId,
840        now_ms: u64,
841        renew_ttl: bool,
842        placement: ShardPlacement,
843    ) -> Result<Option<TouchStreamAccessResponse>, GroupEngineError> {
844        if !self.access_requires_write(stream_id, now_ms, renew_ttl)? {
845            return Ok(None);
846        }
847        let response =
848            self.apply_access_command(stream_id.clone(), now_ms, renew_ttl, placement)?;
849        if response.expired {
850            return Err(GroupEngineError::stream(
851                StreamErrorCode::StreamNotFound,
852                format!("stream '{stream_id}' does not exist"),
853            ));
854        }
855        Ok(Some(response))
856    }
857
858    pub(crate) fn append_payload(
859        &mut self,
860        input: AppendPayloadInput<'_>,
861        placement: ShardPlacement,
862    ) -> Result<AppendResponse, GroupEngineError> {
863        let AppendPayloadInput {
864            stream_id,
865            content_type,
866            payload,
867            close_after,
868            stream_seq,
869            producer,
870            now_ms,
871            record_match,
872        } = input;
873        let stream_count_key = stream_id.clone();
874        let response = self.state_machine.append_borrowed(AppendStreamInput {
875            stream_id,
876            content_type,
877            payload,
878            close_after,
879            stream_seq,
880            producer,
881            now_ms,
882            record_match,
883        });
884        self.append_response_from_stream(stream_count_key, response, placement)
885    }
886
887    fn append_response_from_stream(
888        &mut self,
889        stream_id: BucketStreamId,
890        response: StreamResponse,
891        placement: ShardPlacement,
892    ) -> Result<AppendResponse, GroupEngineError> {
893        match response {
894            StreamResponse::Appended {
895                offset,
896                next_offset,
897                closed,
898                deduplicated,
899                producer,
900                ..
901            } => {
902                let stream_hot_bytes = self.state_machine.hot_payload_len(&stream_id).unwrap_or(0);
903                let group_hot_bytes = self.state_machine.total_hot_payload_bytes();
904                let stream_append_count = self
905                    .stream_append_counts
906                    .entry(stream_id.clone())
907                    .or_insert(0);
908                let record_range = self
909                    .state_machine
910                    .record_range_for_append(&stream_id, offset, next_offset, producer.as_ref())
911                    .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?;
912                if !deduplicated {
913                    self.commit_index += 1;
914                    *stream_append_count += 1;
915                }
916                Ok(AppendResponse {
917                    placement,
918                    start_offset: offset,
919                    next_offset,
920                    stream_append_count: *stream_append_count,
921                    group_commit_index: self.commit_index,
922                    closed,
923                    deduplicated,
924                    producer,
925                    record_range,
926                    stream_hot_bytes,
927                    group_hot_bytes,
928                })
929            }
930            StreamResponse::Error {
931                code,
932                message,
933                next_offset,
934                context,
935            } => Err(GroupEngineError::stream_with_context(
936                code,
937                message,
938                next_offset,
939                context,
940            )),
941            other => Err(GroupEngineError::new(format!(
942                "unexpected append response: {other:?}"
943            ))),
944        }
945    }
946
947    pub fn read_stream_plan(
948        &mut self,
949        request: &ReadStreamRequest,
950        placement: ShardPlacement,
951    ) -> Result<StreamReadPlan, GroupEngineError> {
952        self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
953        self.read_stream_plan_after_access(request)
954    }
955
956    pub fn read_stream_plan_after_access(
957        &self,
958        request: &ReadStreamRequest,
959    ) -> Result<StreamReadPlan, GroupEngineError> {
960        let Some(record) = request.record else {
961            let mut plan = self
962                .state_machine
963                .read_plan_at(
964                    &request.stream_id,
965                    request.offset,
966                    request.max_len,
967                    request.now_ms,
968                )
969                .map_err(stream_response_error)?;
970            plan.retained_record_range = self
971                .state_machine
972                .record_range(&request.stream_id)
973                .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?;
974            return Ok(plan);
975        };
976        let retained_record_range = self
977            .state_machine
978            .record_range(&request.stream_id)
979            .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?
980            .ok_or_else(|| {
981                GroupEngineError::stream(
982                    StreamErrorCode::InvalidRecordBoundaries,
983                    "record coordinates are inactive for this stream",
984                )
985            })?;
986        if record < retained_record_range.first_record {
987            return Err(GroupEngineError::stream(
988                StreamErrorCode::StreamGone,
989                format!(
990                    "record {record} is older than first retained record {}",
991                    retained_record_range.first_record
992                ),
993            ));
994        }
995        if record > retained_record_range.next_record {
996            return Err(GroupEngineError::stream(
997                StreamErrorCode::InvalidRecordBoundaries,
998                format!(
999                    "record {record} is beyond record tail {}",
1000                    retained_record_range.next_record
1001                ),
1002            ));
1003        }
1004        let next_record = request
1005            .max_records
1006            .map(|limit| record.saturating_add(limit))
1007            .unwrap_or(retained_record_range.next_record)
1008            .min(retained_record_range.next_record);
1009        let offset = self
1010            .state_machine
1011            .offset_for_record(&request.stream_id, record)
1012            .map_err(|err| GroupEngineError::new(format!("record offset: {err:?}")))?
1013            .ok_or_else(|| GroupEngineError::new("record stream disappeared"))?;
1014        let next_offset = self
1015            .state_machine
1016            .offset_for_record(&request.stream_id, next_record)
1017            .map_err(|err| GroupEngineError::new(format!("record offset: {err:?}")))?
1018            .ok_or_else(|| GroupEngineError::new("record stream disappeared"))?;
1019        let max_len = usize::try_from(next_offset.saturating_sub(offset))
1020            .map_err(|_| GroupEngineError::new("record read window exceeds usize"))?;
1021        let mut plan = self
1022            .state_machine
1023            .read_plan_at(&request.stream_id, offset, max_len, request.now_ms)
1024            .map_err(stream_response_error)?;
1025        plan.retained_record_range = Some(retained_record_range);
1026        plan.record_range = Some(ursula_stream::StreamRecordRange {
1027            first_record: record,
1028            next_record,
1029        });
1030        Ok(plan)
1031    }
1032
1033    /// Per-bucket usage held by this group's state machine. Public so the
1034    /// Raft engine can serve usage reads from its applied state machine.
1035    pub fn bucket_usage_report(&self) -> Vec<ursula_stream::BucketUsageSnapshot> {
1036        self.state_machine.bucket_usage_report()
1037    }
1038
1039    pub fn head_stream_after_access(
1040        &mut self,
1041        request: &HeadStreamRequest,
1042        placement: ShardPlacement,
1043    ) -> Result<HeadStreamResponse, GroupEngineError> {
1044        let Some(metadata) = self
1045            .state_machine
1046            .head_at(&request.stream_id, request.now_ms)
1047        else {
1048            return Err(GroupEngineError::stream(
1049                StreamErrorCode::StreamNotFound,
1050                format!("stream '{}' does not exist", request.stream_id),
1051            ));
1052        };
1053        let content_type = metadata.content_type.clone();
1054        let tail_offset = metadata.tail_offset;
1055        let closed = metadata.status == ursula_stream::StreamStatus::Closed;
1056        let stream_ttl_seconds = metadata.stream_ttl_seconds;
1057        let stream_expires_at_ms = metadata.stream_expires_at_ms;
1058        let _ = metadata;
1059        let snapshot = self
1060            .state_machine
1061            .latest_snapshot(&request.stream_id)
1062            .map_err(stream_response_error)?;
1063        Ok(HeadStreamResponse {
1064            placement,
1065            content_type,
1066            tail_offset,
1067            cold_hot_start_offset: self.state_machine.hot_start_offset(&request.stream_id),
1068            closed,
1069            stream_ttl_seconds,
1070            stream_expires_at_ms,
1071            snapshot_offset: snapshot.as_ref().map(|snapshot| snapshot.offset),
1072            snapshot_digest: snapshot.map(|snapshot| snapshot.digest),
1073            retained_offset: self.state_machine.retained_offset(&request.stream_id),
1074            integrity: self
1075                .state_machine
1076                .integrity_snapshot(&request.stream_id)
1077                .map_err(stream_response_error)?,
1078            record_range: self
1079                .state_machine
1080                .record_range(&request.stream_id)
1081                .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?,
1082        })
1083    }
1084
1085    pub fn get_stream_attrs_after_access(
1086        &mut self,
1087        request: &GetStreamAttrsRequest,
1088        placement: ShardPlacement,
1089    ) -> Result<GetStreamAttrsResponse, GroupEngineError> {
1090        if self
1091            .state_machine
1092            .head_at(&request.stream_id, request.now_ms)
1093            .is_none()
1094        {
1095            return Err(GroupEngineError::stream(
1096                StreamErrorCode::StreamNotFound,
1097                format!("stream '{}' does not exist", request.stream_id),
1098            ));
1099        }
1100        Ok(GetStreamAttrsResponse {
1101            placement,
1102            attrs: self.state_machine.stream_attrs(&request.stream_id).cloned(),
1103        })
1104    }
1105
1106    pub async fn read_payload_from_plan(
1107        cold_store: Option<&ColdStoreHandle>,
1108        cold_index_cache: Option<&Arc<ColdIndexPageCache<ColdStoreColdIndexPageStore>>>,
1109        stream_id: &BucketStreamId,
1110        plan: &StreamReadPlan,
1111    ) -> Result<Vec<u8>, GroupEngineError> {
1112        let mut payload = Vec::new();
1113        for segment in &plan.segments {
1114            match segment {
1115                StreamReadSegment::Hot(bytes) => payload.extend_from_slice(bytes),
1116                StreamReadSegment::ColdIndex(segment) => {
1117                    let Some(cold_store) = cold_store else {
1118                        return Err(GroupEngineError::stream_with_next_offset(
1119                            StreamErrorCode::InvalidColdFlush,
1120                            format!("stream '{stream_id}' read requires object payload store"),
1121                            Some(plan.next_offset),
1122                        ));
1123                    };
1124                    let Some(cache) = cold_index_cache else {
1125                        return Err(GroupEngineError::stream_with_next_offset(
1126                            StreamErrorCode::InvalidColdFlush,
1127                            format!("stream '{stream_id}' read requires cold index page cache"),
1128                            Some(plan.next_offset),
1129                        ));
1130                    };
1131                    let objects = cache
1132                        .object_segments_for_read(stream_id, segment)
1133                        .await
1134                        .map_err(|err| GroupEngineError::new(err.to_string()))?;
1135                    let segment_end = segment
1136                        .read_start_offset
1137                        .saturating_add(u64::try_from(segment.len).expect("read len fits u64"));
1138                    let mut cursor = segment.read_start_offset;
1139                    for object in objects {
1140                        // A retried cold flush can leave overlapping objects in
1141                        // an index page. They describe the same byte range, so
1142                        // materialize each byte once while still using the
1143                        // original object start for range addressing.
1144                        let start = object
1145                            .start_offset
1146                            .max(segment.read_start_offset)
1147                            .max(cursor);
1148                        let end = object.end_offset.min(segment_end);
1149                        if start >= end {
1150                            continue;
1151                        }
1152                        let bytes = cold_store
1153                            .read_object_range_for_stream(
1154                                stream_id,
1155                                &object,
1156                                start,
1157                                usize::try_from(end - start).expect("object read len fits usize"),
1158                            )
1159                            .await
1160                            .map_err(|err| GroupEngineError::new(err.to_string()))?;
1161                        payload.extend_from_slice(&bytes);
1162                        cursor = end;
1163                    }
1164                }
1165                StreamReadSegment::Object(segment) => {
1166                    let Some(cold_store) = cold_store else {
1167                        return Err(GroupEngineError::stream_with_next_offset(
1168                            StreamErrorCode::InvalidColdFlush,
1169                            format!("stream '{stream_id}' read requires object payload store"),
1170                            Some(plan.next_offset),
1171                        ));
1172                    };
1173                    let bytes = cold_store
1174                        .read_object_range_for_stream(
1175                            stream_id,
1176                            &segment.object,
1177                            segment.read_start_offset,
1178                            segment.len,
1179                        )
1180                        .await
1181                        .map_err(|err| GroupEngineError::new(err.to_string()))?;
1182                    payload.extend_from_slice(&bytes);
1183                }
1184            }
1185        }
1186        Ok(payload)
1187    }
1188
1189    pub(crate) async fn read_own_payload_from_plan(
1190        &self,
1191        stream_id: &BucketStreamId,
1192        plan: &StreamReadPlan,
1193    ) -> Result<Vec<u8>, GroupEngineError> {
1194        Self::read_payload_from_plan(
1195            self.cold_store.as_ref(),
1196            self.cold_index_cache.as_ref(),
1197            stream_id,
1198            plan,
1199        )
1200        .await
1201    }
1202
1203    pub(crate) async fn bootstrap_updates(
1204        &self,
1205        stream_id: &BucketStreamId,
1206        records: &[StreamMessageRecord],
1207        content_type: &str,
1208        now_ms: u64,
1209    ) -> Result<Vec<BootstrapUpdate>, GroupEngineError> {
1210        let mut updates = Vec::with_capacity(records.len());
1211        for record in records {
1212            let len = usize::try_from(record.end_offset - record.start_offset).map_err(|_| {
1213                GroupEngineError::stream(
1214                    StreamErrorCode::InvalidSnapshot,
1215                    format!(
1216                        "bootstrap message [{}..{}) for stream '{stream_id}' is too large",
1217                        record.start_offset, record.end_offset
1218                    ),
1219                )
1220            })?;
1221            let plan = self
1222                .state_machine
1223                .read_plan_at(stream_id, record.start_offset, len, now_ms)
1224                .map_err(stream_response_error)?;
1225            let payload = self.read_own_payload_from_plan(stream_id, &plan).await?;
1226            updates.push(BootstrapUpdate {
1227                start_offset: record.start_offset,
1228                next_offset: record.end_offset,
1229                content_type: content_type.to_owned(),
1230                payload,
1231            });
1232        }
1233        Ok(updates)
1234    }
1235
1236    pub(crate) fn build_snapshot(&self, placement: ShardPlacement) -> GroupSnapshot {
1237        let stream_snapshot = self.state_machine.snapshot();
1238        let stream_append_counts = self.stream_append_counts_snapshot(&stream_snapshot);
1239        GroupSnapshot {
1240            placement,
1241            group_commit_index: self.commit_index,
1242            stream_snapshot,
1243            stream_append_counts,
1244        }
1245    }
1246
1247    pub(crate) fn stream_append_counts_snapshot(
1248        &self,
1249        stream_snapshot: &ursula_stream::StreamSnapshot,
1250    ) -> Vec<StreamAppendCount> {
1251        // Only emit append counts for streams actually present in the snapshot.
1252        // A deleted/expired stream can leave a stale entry in the runtime map;
1253        // emitting it would make every follower's `install_snapshot` fail the
1254        // `restore_stream_append_counts` consistency check, so a lagging node
1255        // could never catch up (and leadership transfer, which catches the
1256        // target up via a snapshot, could never complete).
1257        let live: HashSet<&BucketStreamId> = stream_snapshot
1258            .streams
1259            .iter()
1260            .map(|entry| &entry.metadata.stream_id)
1261            .collect();
1262        let mut counts = self
1263            .stream_append_counts
1264            .iter()
1265            .filter(|(stream_id, _)| live.contains(stream_id))
1266            .map(|(stream_id, append_count)| StreamAppendCount {
1267                stream_id: stream_id.clone(),
1268                append_count: *append_count,
1269            })
1270            .collect::<Vec<_>>();
1271        counts.sort_by(|left, right| compare_stream_ids(&left.stream_id, &right.stream_id));
1272        counts
1273    }
1274
1275    pub fn stream_tail_offset(&self, stream_id: &BucketStreamId) -> Option<u64> {
1276        self.state_machine
1277            .head(stream_id)
1278            .map(|metadata| metadata.tail_offset)
1279    }
1280
1281    pub(crate) fn install_snapshot_inner(
1282        &mut self,
1283        snapshot: GroupSnapshot,
1284    ) -> Result<(), GroupEngineError> {
1285        let GroupSnapshot {
1286            placement: _,
1287            group_commit_index,
1288            stream_snapshot,
1289            stream_append_counts,
1290        } = snapshot;
1291        self.install_snapshot_parts(group_commit_index, stream_snapshot, stream_append_counts)
1292    }
1293
1294    pub(crate) fn install_snapshot_parts(
1295        &mut self,
1296        group_commit_index: u64,
1297        stream_snapshot: StreamSnapshot,
1298        stream_append_counts: Vec<StreamAppendCount>,
1299    ) -> Result<(), GroupEngineError> {
1300        let stream_ids = stream_snapshot
1301            .streams
1302            .iter()
1303            .map(|entry| entry.metadata.stream_id.clone())
1304            .collect::<HashSet<_>>();
1305        let state_machine = StreamStateMachine::restore(stream_snapshot)
1306            .map_err(|err| GroupEngineError::new(format!("restore stream snapshot: {err}")))?;
1307        let stream_append_counts = restore_stream_append_counts(stream_append_counts, &stream_ids)?;
1308
1309        self.commit_index = group_commit_index;
1310        self.state_machine = state_machine;
1311        self.stream_append_counts = stream_append_counts;
1312        Ok(())
1313    }
1314}
1315
1316impl GroupEngine for InMemoryGroupEngine {
1317    fn create_stream<'a>(
1318        &'a mut self,
1319        request: CreateStreamRequest,
1320        placement: ShardPlacement,
1321        admission: ColdWriteAdmission,
1322    ) -> GroupCreateStreamFuture<'a> {
1323        if admission.is_enabled() {
1324            return Box::pin(async move {
1325                self.create_stream_with_admission_inner(request, placement, admission)
1326            });
1327        }
1328        let command = GroupWriteCommand::from(request);
1329        Box::pin(async move {
1330            match self.apply_committed_write(command, placement)? {
1331                GroupWriteResponse::CreateStream(response) => Ok(response),
1332                other => Err(GroupEngineError::new(format!(
1333                    "unexpected create stream write response: {other:?}"
1334                ))),
1335            }
1336        })
1337    }
1338
1339    fn create_stream_external<'a>(
1340        &'a mut self,
1341        request: CreateStreamExternalRequest,
1342        placement: ShardPlacement,
1343    ) -> GroupCreateStreamFuture<'a> {
1344        Box::pin(async move {
1345            if let Some(cold_store) = self.cold_store.as_ref() {
1346                let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1347                write_external_segment_index_pages(
1348                    &store,
1349                    &request.stream_id,
1350                    0,
1351                    &request.initial_payload,
1352                )
1353                .await
1354                .map_err(|err| GroupEngineError::new(err.to_string()))?;
1355            }
1356            let command = GroupWriteCommand::from(request);
1357            match self.apply_committed_write(command, placement)? {
1358                GroupWriteResponse::CreateStream(response) => Ok(response),
1359                other => Err(GroupEngineError::new(format!(
1360                    "unexpected external create stream write response: {other:?}"
1361                ))),
1362            }
1363        })
1364    }
1365
1366    fn read_stream<'a>(
1367        &'a mut self,
1368        request: ReadStreamRequest,
1369        placement: ShardPlacement,
1370    ) -> GroupReadStreamFuture<'a> {
1371        Box::pin(async move {
1372            self.read_stream_parts(request, placement)
1373                .await?
1374                .into_response()
1375                .await
1376        })
1377    }
1378
1379    fn read_stream_parts<'a>(
1380        &'a mut self,
1381        request: ReadStreamRequest,
1382        placement: ShardPlacement,
1383    ) -> GroupReadStreamPartsFuture<'a> {
1384        Box::pin(async move {
1385            let stream_id = request.stream_id.clone();
1386            let plan = self.read_stream_plan(&request, placement)?;
1387            Ok(GroupReadStreamParts::from_plan(
1388                placement,
1389                stream_id,
1390                plan,
1391                self.cold_store(),
1392                self.cold_index_cache.clone(),
1393            ))
1394        })
1395    }
1396
1397    fn publish_snapshot<'a>(
1398        &'a mut self,
1399        request: PublishSnapshotRequest,
1400        placement: ShardPlacement,
1401    ) -> GroupPublishSnapshotFuture<'a> {
1402        Box::pin(async move {
1403            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1404            let command = GroupWriteCommand::from(request);
1405            match self.apply_committed_write(command, placement)? {
1406                GroupWriteResponse::PublishSnapshot(response) => Ok(response),
1407                other => Err(GroupEngineError::new(format!(
1408                    "unexpected publish snapshot write response: {other:?}"
1409                ))),
1410            }
1411        })
1412    }
1413
1414    fn advance_retention<'a>(
1415        &'a mut self,
1416        request: AdvanceRetentionRequest,
1417        placement: ShardPlacement,
1418    ) -> GroupAdvanceRetentionFuture<'a> {
1419        Box::pin(async move {
1420            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1421            let command = GroupWriteCommand::from(request);
1422            match self.apply_committed_write(command, placement)? {
1423                GroupWriteResponse::AdvanceRetention(response) => Ok(response),
1424                other => Err(GroupEngineError::new(format!(
1425                    "unexpected advance retention write response: {other:?}"
1426                ))),
1427            }
1428        })
1429    }
1430
1431    fn import_group_state<'a>(
1432        &'a mut self,
1433        request: ImportGroupStateRequest,
1434        placement: ShardPlacement,
1435    ) -> crate::GroupImportGroupStateFuture<'a> {
1436        Box::pin(async move {
1437            let command = GroupWriteCommand::from(StreamCommand::from(request));
1438            match self.apply_committed_write(command, placement)? {
1439                GroupWriteResponse::ImportGroupState(response) => Ok(response),
1440                other => Err(GroupEngineError::new(format!(
1441                    "unexpected group state import response: {other:?}"
1442                ))),
1443            }
1444        })
1445    }
1446
1447    fn set_bucket_quota<'a>(
1448        &'a mut self,
1449        request: SetBucketQuotaRequest,
1450        placement: ShardPlacement,
1451    ) -> GroupSetBucketQuotaFuture<'a> {
1452        Box::pin(async move {
1453            let command = GroupWriteCommand::from(request);
1454            match self.apply_committed_write(command, placement)? {
1455                GroupWriteResponse::SetBucketQuota(response) => Ok(response),
1456                other => Err(GroupEngineError::new(format!(
1457                    "unexpected set bucket quota write response: {other:?}"
1458                ))),
1459            }
1460        })
1461    }
1462
1463    fn read_snapshot<'a>(
1464        &'a mut self,
1465        request: ReadSnapshotRequest,
1466        placement: ShardPlacement,
1467    ) -> GroupReadSnapshotFuture<'a> {
1468        Box::pin(async move {
1469            self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1470            let snapshot = match request.snapshot_offset {
1471                Some(offset) => self
1472                    .state_machine
1473                    .read_snapshot(&request.stream_id, offset)
1474                    .map_err(stream_response_error)?,
1475                None => self
1476                    .state_machine
1477                    .latest_snapshot(&request.stream_id)
1478                    .map_err(stream_response_error)?
1479                    .ok_or_else(|| {
1480                        GroupEngineError::stream(
1481                            StreamErrorCode::SnapshotNotFound,
1482                            format!("stream '{}' has no visible snapshot", request.stream_id),
1483                        )
1484                    })?,
1485            };
1486            let tail_offset = self
1487                .state_machine
1488                .head_at(&request.stream_id, request.now_ms)
1489                .map(|metadata| metadata.tail_offset)
1490                .unwrap_or(snapshot.offset);
1491            Ok(ReadSnapshotResponse {
1492                placement,
1493                snapshot_offset: snapshot.offset,
1494                next_offset: snapshot.offset,
1495                content_type: snapshot.content_type,
1496                snapshot_digest: snapshot.digest,
1497                payload: snapshot.payload,
1498                up_to_date: snapshot.offset == tail_offset,
1499                record_range: self
1500                    .state_machine
1501                    .record_range(&request.stream_id)
1502                    .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?,
1503            })
1504        })
1505    }
1506
1507    fn delete_snapshot<'a>(
1508        &'a mut self,
1509        request: DeleteSnapshotRequest,
1510        placement: ShardPlacement,
1511    ) -> GroupDeleteSnapshotFuture<'a> {
1512        Box::pin(async move {
1513            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1514            match self
1515                .state_machine
1516                .delete_snapshot(&request.stream_id, request.snapshot_offset)
1517            {
1518                StreamResponse::Error {
1519                    code,
1520                    message,
1521                    next_offset,
1522                    context,
1523                } => Err(GroupEngineError::stream_with_context(
1524                    code,
1525                    message,
1526                    next_offset,
1527                    context,
1528                )),
1529                other => Err(GroupEngineError::new(format!(
1530                    "unexpected delete snapshot response: {other:?}"
1531                ))),
1532            }
1533        })
1534    }
1535
1536    fn bootstrap_stream<'a>(
1537        &'a mut self,
1538        request: BootstrapStreamRequest,
1539        placement: ShardPlacement,
1540    ) -> GroupBootstrapStreamFuture<'a> {
1541        Box::pin(async move {
1542            self.ensure_stream_access(&request.stream_id, request.now_ms, true, placement)?;
1543            let plan = self
1544                .state_machine
1545                .bootstrap_plan(&request.stream_id)
1546                .map_err(stream_response_error)?;
1547            let snapshot_offset = plan.snapshot.as_ref().map(|snapshot| snapshot.offset);
1548            let snapshot_content_type = plan
1549                .snapshot
1550                .as_ref()
1551                .map(|snapshot| snapshot.content_type.clone())
1552                .unwrap_or_else(|| DEFAULT_CONTENT_TYPE.to_owned());
1553            let snapshot_payload = plan
1554                .snapshot
1555                .as_ref()
1556                .map(|snapshot| snapshot.payload.clone())
1557                .unwrap_or_default();
1558            let updates = self
1559                .bootstrap_updates(
1560                    &request.stream_id,
1561                    &plan.updates,
1562                    &plan.content_type,
1563                    request.now_ms,
1564                )
1565                .await?;
1566            Ok(BootstrapStreamResponse {
1567                placement,
1568                snapshot_offset,
1569                snapshot_content_type,
1570                snapshot_payload,
1571                updates,
1572                next_offset: plan.next_offset,
1573                up_to_date: plan.up_to_date,
1574                closed: plan.closed,
1575                record_range: self
1576                    .state_machine
1577                    .record_range(&request.stream_id)
1578                    .map_err(|err| GroupEngineError::new(format!("record range: {err:?}")))?,
1579            })
1580        })
1581    }
1582
1583    fn touch_stream_access<'a>(
1584        &'a mut self,
1585        stream_id: BucketStreamId,
1586        now_ms: u64,
1587        renew_ttl: bool,
1588        placement: ShardPlacement,
1589    ) -> GroupTouchStreamAccessFuture<'a> {
1590        Box::pin(async move { self.apply_access_command(stream_id, now_ms, renew_ttl, placement) })
1591    }
1592
1593    fn head_stream<'a>(
1594        &'a mut self,
1595        request: HeadStreamRequest,
1596        placement: ShardPlacement,
1597    ) -> GroupHeadStreamFuture<'a> {
1598        Box::pin(async move {
1599            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1600            self.head_stream_after_access(&request, placement)
1601        })
1602    }
1603
1604    fn bucket_usage<'a>(&'a mut self, _placement: ShardPlacement) -> GroupBucketUsageFuture<'a> {
1605        Box::pin(async move { Ok(self.state_machine.bucket_usage_report()) })
1606    }
1607
1608    fn get_stream_attrs<'a>(
1609        &'a mut self,
1610        request: GetStreamAttrsRequest,
1611        placement: ShardPlacement,
1612    ) -> GroupGetStreamAttrsFuture<'a> {
1613        Box::pin(async move {
1614            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1615            self.get_stream_attrs_after_access(&request, placement)
1616        })
1617    }
1618
1619    fn update_stream_attrs<'a>(
1620        &'a mut self,
1621        request: UpdateStreamAttrsRequest,
1622        placement: ShardPlacement,
1623    ) -> GroupUpdateStreamAttrsFuture<'a> {
1624        Box::pin(async move {
1625            match self.apply_committed_write(GroupWriteCommand::from(request), placement)? {
1626                GroupWriteResponse::UpdateStreamAttrs(response) => Ok(response),
1627                other => Err(GroupEngineError::new(format!(
1628                    "unexpected update stream attrs write response: {other:?}"
1629                ))),
1630            }
1631        })
1632    }
1633
1634    fn close_stream<'a>(
1635        &'a mut self,
1636        request: CloseStreamRequest,
1637        placement: ShardPlacement,
1638    ) -> GroupCloseStreamFuture<'a> {
1639        Box::pin(async move {
1640            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1641            let command = GroupWriteCommand::from(request);
1642            match self.apply_committed_write(command, placement)? {
1643                GroupWriteResponse::CloseStream(response) => Ok(response),
1644                other => Err(GroupEngineError::new(format!(
1645                    "unexpected close stream write response: {other:?}"
1646                ))),
1647            }
1648        })
1649    }
1650
1651    fn delete_stream<'a>(
1652        &'a mut self,
1653        request: DeleteStreamRequest,
1654        placement: ShardPlacement,
1655    ) -> GroupDeleteStreamFuture<'a> {
1656        let command = GroupWriteCommand::from(request);
1657        Box::pin(async move {
1658            match self.apply_committed_write(command, placement)? {
1659                GroupWriteResponse::DeleteStream(response) => Ok(response),
1660                other => Err(GroupEngineError::new(format!(
1661                    "unexpected delete stream write response: {other:?}"
1662                ))),
1663            }
1664        })
1665    }
1666
1667    fn purge_bucket<'a>(
1668        &'a mut self,
1669        bucket_id: String,
1670        placement: ShardPlacement,
1671    ) -> GroupPurgeBucketFuture<'a> {
1672        Box::pin(async move {
1673            match self.apply_committed_write(
1674                GroupWriteCommand::Stream(StreamCommand::PurgeBucket { bucket_id }),
1675                placement,
1676            )? {
1677                GroupWriteResponse::PurgeBucket(response) => Ok(response),
1678                other => Err(GroupEngineError::new(format!(
1679                    "unexpected purge bucket write response: {other:?}"
1680                ))),
1681            }
1682        })
1683    }
1684
1685    fn ack_cold_gc<'a>(
1686        &'a mut self,
1687        up_to_seq: u64,
1688        placement: ShardPlacement,
1689    ) -> GroupAckColdGcFuture<'a> {
1690        Box::pin(async move {
1691            match self.apply_committed_write(
1692                GroupWriteCommand::Stream(StreamCommand::AckColdGc { up_to_seq }),
1693                placement,
1694            )? {
1695                GroupWriteResponse::AckColdGc(response) => Ok(response),
1696                other => Err(GroupEngineError::new(format!(
1697                    "unexpected ack cold gc write response: {other:?}"
1698                ))),
1699            }
1700        })
1701    }
1702
1703    fn plan_cold_gc<'a>(
1704        &'a mut self,
1705        max: usize,
1706        _placement: ShardPlacement,
1707    ) -> GroupPlanColdGcFuture<'a> {
1708        let entries = self.state_machine.pending_cold_gc_batch(max);
1709        Box::pin(async move { Ok(entries) })
1710    }
1711
1712    fn append<'a>(
1713        &'a mut self,
1714        request: AppendRequest,
1715        placement: ShardPlacement,
1716        admission: ColdWriteAdmission,
1717    ) -> GroupAppendFuture<'a> {
1718        if admission.is_enabled() {
1719            return Box::pin(async move {
1720                self.append_with_admission_inner(request, placement, admission)
1721            });
1722        }
1723        Box::pin(async move {
1724            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1725            let command = GroupWriteCommand::from(request);
1726            match self.apply_committed_write(command, placement)? {
1727                GroupWriteResponse::Append(response) => Ok(response),
1728                other => Err(GroupEngineError::new(format!(
1729                    "unexpected append write response: {other:?}"
1730                ))),
1731            }
1732        })
1733    }
1734
1735    fn append_transaction<'a>(
1736        &'a mut self,
1737        request: AppendTransactionRequest,
1738        placement: ShardPlacement,
1739        admission: ColdWriteAdmission,
1740    ) -> GroupAppendTransactionFuture<'a> {
1741        Box::pin(async move {
1742            let Some(first) = request.operations.first() else {
1743                return Err(GroupEngineError::new(
1744                    "append transaction must contain at least one operation",
1745                ));
1746            };
1747            self.check_cold_write_admission_bytes(
1748                &first.stream_id,
1749                admission,
1750                request.payload_bytes(),
1751            )?;
1752            let command = GroupWriteCommand::Transaction {
1753                commands: request
1754                    .operations
1755                    .into_iter()
1756                    .map(StreamCommand::from)
1757                    .collect(),
1758            };
1759            let GroupWriteResponse::Batch(items) =
1760                self.apply_committed_write(command, placement)?
1761            else {
1762                return Err(GroupEngineError::new(
1763                    "unexpected append transaction write response",
1764                ));
1765            };
1766            let items = items
1767                .into_iter()
1768                .map(|item| match item? {
1769                    GroupWriteResponse::Append(response) => Ok(response),
1770                    other => Err(GroupEngineError::new(format!(
1771                        "unexpected append transaction item response: {other:?}"
1772                    ))),
1773                })
1774                .collect::<Result<Vec<_>, _>>()?;
1775            Ok(AppendTransactionResponse { placement, items })
1776        })
1777    }
1778
1779    fn append_external<'a>(
1780        &'a mut self,
1781        request: AppendExternalRequest,
1782        placement: ShardPlacement,
1783    ) -> GroupAppendFuture<'a> {
1784        Box::pin(async move {
1785            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1786            if let Some(cold_store) = self.cold_store.as_ref() {
1787                let start_offset = self
1788                    .state_machine
1789                    .head(&request.stream_id)
1790                    .map(|metadata| metadata.tail_offset)
1791                    .ok_or_else(|| {
1792                        GroupEngineError::stream(
1793                            ursula_stream::StreamErrorCode::StreamNotFound,
1794                            format!("stream '{}' does not exist", request.stream_id),
1795                        )
1796                    })?;
1797                let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1798                write_external_segment_index_pages(
1799                    &store,
1800                    &request.stream_id,
1801                    start_offset,
1802                    &request.payload,
1803                )
1804                .await
1805                .map_err(|err| GroupEngineError::new(err.to_string()))?;
1806            }
1807            let command = GroupWriteCommand::from(request);
1808            match self.apply_committed_write(command, placement)? {
1809                GroupWriteResponse::Append(response) => Ok(response),
1810                other => Err(GroupEngineError::new(format!(
1811                    "unexpected external append write response: {other:?}"
1812                ))),
1813            }
1814        })
1815    }
1816
1817    fn append_batch<'a>(
1818        &'a mut self,
1819        request: AppendBatchRequest,
1820        placement: ShardPlacement,
1821        admission: ColdWriteAdmission,
1822    ) -> GroupAppendBatchFuture<'a> {
1823        if admission.is_enabled() {
1824            return Box::pin(async move {
1825                self.append_batch_with_admission_inner(request, placement, admission)
1826            });
1827        }
1828        Box::pin(async move {
1829            self.ensure_stream_access(&request.stream_id, request.now_ms, false, placement)?;
1830            let command = GroupWriteCommand::from(request);
1831            match self.apply_committed_write(command, placement)? {
1832                GroupWriteResponse::AppendBatch(response) => Ok(response),
1833                other => Err(GroupEngineError::new(format!(
1834                    "unexpected append batch write response: {other:?}"
1835                ))),
1836            }
1837        })
1838    }
1839
1840    fn flush_cold<'a>(
1841        &'a mut self,
1842        request: FlushColdRequest,
1843        placement: ShardPlacement,
1844    ) -> GroupFlushColdFuture<'a> {
1845        Box::pin(async move {
1846            let mut index_rollback = None;
1847            if !request.chunk.shared_object
1848                && let Some(cold_store) = self.cold_store.as_ref()
1849            {
1850                let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1851                let rollback = write_cold_chunk_index_pages_with_rollback(
1852                    &store,
1853                    &request.stream_id,
1854                    &request.chunk,
1855                )
1856                .await
1857                .map_err(|err| GroupEngineError::new(err.to_string()))?;
1858                index_rollback = Some((store, rollback));
1859            }
1860            let command = GroupWriteCommand::from(request);
1861            match self.apply_committed_write(command, placement) {
1862                Ok(GroupWriteResponse::FlushCold(response)) => Ok(response),
1863                Ok(other) => {
1864                    if let Some((store, rollback)) = index_rollback {
1865                        rollback_cold_index_pages(&store, rollback)
1866                            .await
1867                            .map_err(|err| GroupEngineError::new(err.to_string()))?;
1868                    }
1869                    Err(GroupEngineError::new(format!(
1870                        "unexpected flush cold write response: {other:?}"
1871                    )))
1872                }
1873                Err(err) => {
1874                    if let Some((store, rollback)) = index_rollback {
1875                        rollback_cold_index_pages(&store, rollback).await.map_err(
1876                            |rollback_err| {
1877                                GroupEngineError::new(format!(
1878                                    "rollback cold index after flush failure: {rollback_err}"
1879                                ))
1880                            },
1881                        )?;
1882                    }
1883                    Err(err)
1884                }
1885            }
1886        })
1887    }
1888
1889    fn compact_cold<'a>(
1890        &'a mut self,
1891        request: CompactColdRequest,
1892        placement: ShardPlacement,
1893    ) -> GroupCompactColdFuture<'a> {
1894        Box::pin(async move {
1895            let mut index_rollback = None;
1896            if let Some(cold_store) = self.cold_store.as_ref() {
1897                let store = ColdStoreColdIndexPageStore::new(cold_store.clone());
1898                let rollback = if request.old_chunks.iter().all(|chunk| chunk.shared_object) {
1899                    write_cold_chunk_index_pages_with_rollback(
1900                        &store,
1901                        &request.stream_id,
1902                        &request.replacement,
1903                    )
1904                    .await
1905                    .map_err(|err| GroupEngineError::new(err.to_string()))?
1906                } else {
1907                    let Some(rollback) = replace_cold_chunk_index_pages_with_rollback(
1908                        &store,
1909                        &request.stream_id,
1910                        &request.old_chunks,
1911                        &request.replacement,
1912                    )
1913                    .await
1914                    .map_err(|err| GroupEngineError::new(err.to_string()))?
1915                    else {
1916                        return Err(GroupEngineError::new(
1917                            "cold compaction input no longer matches the cold index",
1918                        ));
1919                    };
1920                    rollback
1921                };
1922                index_rollback = Some((store, rollback));
1923            }
1924            let command = GroupWriteCommand::from(request);
1925            let result = match self.apply_committed_write(command, placement) {
1926                Ok(GroupWriteResponse::CompactCold(response)) => Ok(response),
1927                Ok(other) => Err(GroupEngineError::new(format!(
1928                    "unexpected compact cold write response: {other:?}"
1929                ))),
1930                Err(err) => Err(err),
1931            };
1932            if result.is_err()
1933                && let Some((store, rollback)) = index_rollback
1934            {
1935                rollback_cold_index_pages(&store, rollback)
1936                    .await
1937                    .map_err(|err| {
1938                        GroupEngineError::new(format!(
1939                            "rollback cold index after compaction failure: {err}"
1940                        ))
1941                    })?;
1942            }
1943            result
1944        })
1945    }
1946
1947    fn plan_cold_flush<'a>(
1948        &'a mut self,
1949        request: PlanColdFlushRequest,
1950        _placement: ShardPlacement,
1951    ) -> GroupPlanColdFlushFuture<'a> {
1952        Box::pin(async move {
1953            self.state_machine
1954                .plan_cold_flush(
1955                    &request.stream_id,
1956                    request.min_hot_bytes,
1957                    request.max_flush_bytes,
1958                )
1959                .map_err(stream_response_error)
1960        })
1961    }
1962
1963    fn plan_next_cold_flush_batch<'a>(
1964        &'a mut self,
1965        request: PlanGroupColdFlushRequest,
1966        _placement: ShardPlacement,
1967        max_candidates: usize,
1968    ) -> GroupPlanNextColdFlushBatchFuture<'a> {
1969        Box::pin(async move {
1970            self.state_machine
1971                .plan_next_cold_flush_batch(
1972                    request.min_hot_bytes,
1973                    request.max_flush_bytes,
1974                    request.max_batch_bytes,
1975                    max_candidates,
1976                )
1977                .map_err(stream_response_error)
1978        })
1979    }
1980
1981    fn cold_hot_backlog<'a>(
1982        &'a mut self,
1983        stream_id: BucketStreamId,
1984        _placement: ShardPlacement,
1985    ) -> GroupColdHotBacklogFuture<'a> {
1986        Box::pin(async move { self.cold_hot_backlog_for(stream_id) })
1987    }
1988
1989    fn snapshot<'a>(&'a mut self, placement: ShardPlacement) -> GroupSnapshotFuture<'a> {
1990        Box::pin(async move { Ok(self.build_snapshot(placement)) })
1991    }
1992
1993    fn install_snapshot<'a>(
1994        &'a mut self,
1995        snapshot: GroupSnapshot,
1996    ) -> GroupInstallSnapshotFuture<'a> {
1997        Box::pin(async move { self.install_snapshot_inner(snapshot) })
1998    }
1999}
2000
2001#[derive(Debug, Clone, Default)]
2002pub struct InMemoryGroupEngineFactory {
2003    cold_store: Option<ColdStoreHandle>,
2004}
2005
2006impl InMemoryGroupEngineFactory {
2007    pub fn new() -> Self {
2008        Self::default()
2009    }
2010
2011    pub fn with_cold_store(cold_store: Option<ColdStoreHandle>) -> Self {
2012        Self { cold_store }
2013    }
2014}
2015
2016impl GroupEngineFactory for InMemoryGroupEngineFactory {
2017    fn create<'a>(
2018        &'a self,
2019        _placement: ShardPlacement,
2020        _metrics: GroupEngineMetrics,
2021    ) -> GroupEngineCreateFuture<'a> {
2022        Box::pin(async move {
2023            let mut engine = InMemoryGroupEngine::default();
2024            engine.set_cold_store(self.cold_store.clone());
2025            let engine: Box<dyn GroupEngine> = Box::new(engine);
2026            Ok(engine)
2027        })
2028    }
2029}
2030
2031pub(crate) fn compare_stream_ids(
2032    left: &BucketStreamId,
2033    right: &BucketStreamId,
2034) -> std::cmp::Ordering {
2035    left.bucket_id
2036        .cmp(&right.bucket_id)
2037        .then_with(|| left.stream_id.cmp(&right.stream_id))
2038}
2039pub(crate) fn ensure_bucket_exists(
2040    state_machine: &mut StreamStateMachine,
2041    stream_id: &BucketStreamId,
2042) -> Result<(), GroupEngineError> {
2043    if state_machine.bucket_exists(&stream_id.bucket_id) {
2044        return Ok(());
2045    }
2046
2047    match state_machine.apply(StreamCommand::CreateBucket {
2048        bucket_id: stream_id.bucket_id.clone(),
2049    }) {
2050        StreamResponse::BucketCreated { .. } | StreamResponse::BucketAlreadyExists { .. } => Ok(()),
2051        StreamResponse::Error {
2052            code,
2053            message,
2054            next_offset,
2055            context,
2056        } => Err(GroupEngineError::stream_with_context(
2057            code,
2058            message,
2059            next_offset,
2060            context,
2061        )),
2062        other => Err(GroupEngineError::new(format!(
2063            "unexpected create bucket response: {other:?}"
2064        ))),
2065    }
2066}
2067
2068/// Stream id a command targets, if any (bucket and GC commands have none).
2069fn command_stream_id(command: &StreamCommand) -> Option<BucketStreamId> {
2070    match command {
2071        StreamCommand::CreateBucket { .. }
2072        | StreamCommand::DeleteBucket { .. }
2073        | StreamCommand::PurgeBucket { .. }
2074        | StreamCommand::AckColdGc { .. }
2075        | StreamCommand::ImportSnapshot { .. }
2076        | StreamCommand::SetBucketQuota { .. } => None,
2077        StreamCommand::CreateStream { stream_id, .. }
2078        | StreamCommand::CreateExternal { stream_id, .. }
2079        | StreamCommand::Append { stream_id, .. }
2080        | StreamCommand::AppendExternal { stream_id, .. }
2081        | StreamCommand::AppendBatch { stream_id, .. }
2082        | StreamCommand::PublishSnapshot { stream_id, .. }
2083        | StreamCommand::AdvanceRetention { stream_id, .. }
2084        | StreamCommand::TouchStreamAccess { stream_id, .. }
2085        | StreamCommand::UpdateStreamAttrs { stream_id, .. }
2086        | StreamCommand::FlushCold { stream_id, .. }
2087        | StreamCommand::CompactCold { stream_id, .. }
2088        | StreamCommand::Close { stream_id, .. }
2089        | StreamCommand::DeleteStream { stream_id } => Some(stream_id.clone()),
2090    }
2091}
2092
2093fn command_producer(command: &StreamCommand) -> Option<ProducerRequest> {
2094    match command {
2095        StreamCommand::Close { producer, .. } => producer.clone(),
2096        _ => None,
2097    }
2098}
2099
2100fn require_response_stream_id(
2101    stream_id: Option<BucketStreamId>,
2102    response: &str,
2103) -> Result<BucketStreamId, GroupEngineError> {
2104    stream_id.ok_or_else(|| {
2105        GroupEngineError::new(format!(
2106            "{response} response for a command without a stream id"
2107        ))
2108    })
2109}
2110
2111pub(crate) fn stream_response_error(response: StreamResponse) -> GroupEngineError {
2112    match response {
2113        StreamResponse::Error {
2114            code,
2115            message,
2116            next_offset,
2117            context,
2118        } => GroupEngineError::stream_with_context(code, message, next_offset, context),
2119        other => GroupEngineError::new(format!("unexpected stream response error: {other:?}")),
2120    }
2121}
2122
2123pub(crate) fn restore_stream_append_counts(
2124    counts: Vec<StreamAppendCount>,
2125    snapshot_stream_ids: &HashSet<BucketStreamId>,
2126) -> Result<HashMap<BucketStreamId, u64>, GroupEngineError> {
2127    let mut restored = HashMap::with_capacity(counts.len());
2128    for count in counts {
2129        if !snapshot_stream_ids.contains(&count.stream_id) {
2130            return Err(GroupEngineError::new(format!(
2131                "append count references missing snapshot stream '{}'",
2132                count.stream_id
2133            )));
2134        }
2135        if restored
2136            .insert(count.stream_id.clone(), count.append_count)
2137            .is_some()
2138        {
2139            return Err(GroupEngineError::new(format!(
2140                "snapshot contains duplicate append count for stream '{}'",
2141                count.stream_id
2142            )));
2143        }
2144    }
2145    Ok(restored)
2146}