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