Skip to main content

ursula_stream/
state_machine.rs

1use std::collections::HashMap;
2use std::collections::HashSet;
3use std::collections::VecDeque;
4
5use ursula_shard::BucketStreamId;
6
7use crate::command::StreamCommand;
8use crate::integrity::StreamIntegrity;
9use crate::model::AppendExternalInput;
10use crate::model::AppendStreamInput;
11use crate::model::COLD_INDEX_PAGE_SPAN_BYTES;
12use crate::model::ColdChunkRef;
13use crate::model::ColdFlushCandidate;
14use crate::model::ColdGcEntry;
15use crate::model::ColdGcTarget;
16use crate::model::ExternalPayloadRef;
17use crate::model::HotPayloadSegment;
18use crate::model::ObjectPayloadRef;
19use crate::model::ProducerAppendRecord;
20use crate::model::ProducerRequest;
21use crate::model::ProducerSnapshot;
22use crate::model::ProducerState;
23use crate::model::StreamBatchAppend;
24use crate::model::StreamBatchAppendItem;
25use crate::model::StreamBootstrapPlan;
26use crate::model::StreamMessageRecord;
27use crate::model::StreamMetadata;
28use crate::model::StreamRead;
29use crate::model::StreamReadColdIndexSegment;
30use crate::model::StreamReadObjectSegment;
31use crate::model::StreamReadPlan;
32use crate::model::StreamReadSegment;
33use crate::model::StreamStatus;
34use crate::model::StreamVisibleSnapshot;
35use crate::response::StreamErrorCode;
36use crate::response::StreamErrorContext;
37use crate::response::StreamResponse;
38use crate::snapshot::StreamSnapshot;
39use crate::snapshot::StreamSnapshotEntry;
40use crate::snapshot::StreamSnapshotError;
41use crate::validate::validate_bucket_id;
42use crate::validate::validate_stream_id;
43
44#[derive(Debug, Clone, Default)]
45pub struct StreamStateMachine {
46    buckets: HashSet<String>,
47    streams: HashMap<BucketStreamId, StreamMetadata>,
48    hot_buffers: HashMap<BucketStreamId, HotBuffer>,
49    cold_index: InMemoryColdIndex,
50    message_records: HashMap<BucketStreamId, Vec<StreamMessageRecord>>,
51    integrities: HashMap<BucketStreamId, StreamIntegrity>,
52    visible_snapshots: HashMap<BucketStreamId, StreamVisibleSnapshot>,
53    producers: HashMap<BucketStreamId, HashMap<String, ProducerState>>,
54    pending_cold_gc: VecDeque<ColdGcEntry>,
55    next_cold_gc_seq: u64,
56}
57
58trait ColdIndexState {
59    fn cold_chunks(&self, stream_id: &BucketStreamId) -> &[ColdChunkRef];
60    fn external_segments(&self, stream_id: &BucketStreamId) -> &[ObjectPayloadRef];
61    fn cold_generation(&self, stream_id: &BucketStreamId) -> u64;
62    fn push_cold_chunk(&mut self, stream_id: BucketStreamId, chunk: ColdChunkRef);
63    fn push_external_segment(&mut self, stream_id: BucketStreamId, object: ObjectPayloadRef);
64    fn restore_stream(
65        &mut self,
66        stream_id: BucketStreamId,
67        cold_frontier_offset: u64,
68        cold_index_generation: u64,
69        cold_chunks: Vec<ColdChunkRef>,
70        external_segments: Vec<ObjectPayloadRef>,
71    );
72    fn remove_stream(&mut self, stream_id: &BucketStreamId) -> bool;
73    fn compact_before(&mut self, stream_id: &BucketStreamId, retained_offset: u64) -> Vec<String>;
74    fn cold_frontier_offset(&self, stream_id: &BucketStreamId, retained_offset: u64) -> u64;
75}
76
77#[derive(Debug, Clone, Default)]
78struct InMemoryColdIndex {
79    cold_chunks: HashMap<BucketStreamId, Vec<ColdChunkRef>>,
80    external_segments: HashMap<BucketStreamId, Vec<ObjectPayloadRef>>,
81    cold_frontiers: HashMap<BucketStreamId, u64>,
82}
83
84impl ColdIndexState for InMemoryColdIndex {
85    fn cold_chunks(&self, stream_id: &BucketStreamId) -> &[ColdChunkRef] {
86        self.cold_chunks
87            .get(stream_id)
88            .map(Vec::as_slice)
89            .unwrap_or(&[])
90    }
91
92    fn external_segments(&self, stream_id: &BucketStreamId) -> &[ObjectPayloadRef] {
93        self.external_segments
94            .get(stream_id)
95            .map(Vec::as_slice)
96            .unwrap_or(&[])
97    }
98
99    fn cold_generation(&self, _stream_id: &BucketStreamId) -> u64 {
100        0
101    }
102
103    fn push_cold_chunk(&mut self, stream_id: BucketStreamId, chunk: ColdChunkRef) {
104        self.cold_frontiers.insert(stream_id, chunk.end_offset);
105    }
106
107    fn push_external_segment(&mut self, stream_id: BucketStreamId, object: ObjectPayloadRef) {
108        let frontier = self
109            .cold_frontiers
110            .get(&stream_id)
111            .copied()
112            .unwrap_or(object.start_offset)
113            .max(object.end_offset);
114        self.cold_frontiers.insert(stream_id, frontier);
115    }
116
117    fn restore_stream(
118        &mut self,
119        stream_id: BucketStreamId,
120        cold_frontier_offset: u64,
121        _cold_index_generation: u64,
122        cold_chunks: Vec<ColdChunkRef>,
123        external_segments: Vec<ObjectPayloadRef>,
124    ) {
125        if cold_frontier_offset > 0 {
126            self.cold_frontiers
127                .insert(stream_id.clone(), cold_frontier_offset);
128        }
129        if !cold_chunks.is_empty() {
130            self.cold_chunks.insert(stream_id.clone(), cold_chunks);
131        }
132        if !external_segments.is_empty() {
133            self.external_segments.insert(stream_id, external_segments);
134        }
135    }
136
137    fn remove_stream(&mut self, stream_id: &BucketStreamId) -> bool {
138        self.cold_frontiers
139            .remove(stream_id)
140            .is_some_and(|frontier| frontier > 0)
141            || self
142                .cold_chunks
143                .remove(stream_id)
144                .is_some_and(|chunks| !chunks.is_empty())
145            || self
146                .external_segments
147                .remove(stream_id)
148                .is_some_and(|objects| !objects.is_empty())
149    }
150
151    fn compact_before(&mut self, stream_id: &BucketStreamId, retained_offset: u64) -> Vec<String> {
152        let mut dropped_cold_paths = Vec::new();
153        if let Some(chunks) = self.cold_chunks.get_mut(stream_id) {
154            chunks.retain(|chunk| {
155                let retain = chunk.end_offset > retained_offset;
156                if !retain {
157                    dropped_cold_paths.push(chunk.s3_path.clone());
158                }
159                retain
160            });
161            if chunks.is_empty() {
162                self.cold_chunks.remove(stream_id);
163            }
164        }
165        if let Some(objects) = self.external_segments.get_mut(stream_id) {
166            objects.retain(|object| object.end_offset > retained_offset);
167            if objects.is_empty() {
168                self.external_segments.remove(stream_id);
169            }
170        }
171        dropped_cold_paths
172    }
173
174    fn cold_frontier_offset(&self, stream_id: &BucketStreamId, retained_offset: u64) -> u64 {
175        let external_segments = self.external_segments(stream_id);
176        let cold_frontier = self
177            .cold_frontiers
178            .get(stream_id)
179            .copied()
180            .unwrap_or(retained_offset);
181        let mut ranges = Vec::with_capacity(1 + external_segments.len());
182        if cold_frontier > retained_offset {
183            ranges.push((retained_offset, cold_frontier));
184        }
185        ranges.extend(
186            external_segments
187                .iter()
188                .map(|object| (object.start_offset, object.end_offset)),
189        );
190        ranges.sort_unstable();
191
192        let mut frontier = retained_offset;
193        for (start_offset, end_offset) in ranges {
194            if end_offset <= frontier {
195                continue;
196            }
197            if start_offset > frontier {
198                break;
199            }
200            frontier = end_offset;
201        }
202        frontier
203    }
204}
205
206#[derive(Debug, Clone, Default)]
207struct HotBuffer {
208    chunks: VecDeque<HotChunk>,
209}
210
211#[derive(Debug, Clone, PartialEq, Eq)]
212struct HotChunk {
213    start_offset: u64,
214    end_offset: u64,
215    bytes: Vec<u8>,
216}
217
218impl HotBuffer {
219    fn from_payload(start_offset: u64, payload: Vec<u8>) -> Self {
220        if payload.is_empty() {
221            return Self::default();
222        }
223        let end_offset = start_offset
224            .saturating_add(u64::try_from(payload.len()).expect("payload len fits u64"));
225        let mut chunks = VecDeque::new();
226        chunks.push_back(HotChunk {
227            start_offset,
228            end_offset,
229            bytes: payload,
230        });
231        Self { chunks }
232    }
233
234    fn from_snapshot(payload: Vec<u8>, segments: &[HotPayloadSegment]) -> Self {
235        let mut chunks = VecDeque::with_capacity(segments.len());
236        for segment in segments {
237            chunks.push_back(HotChunk {
238                start_offset: segment.start_offset,
239                end_offset: segment.end_offset,
240                bytes: payload[segment.payload_start..segment.payload_end].to_vec(),
241            });
242        }
243        Self { chunks }
244    }
245
246    fn len(&self) -> usize {
247        self.chunks.iter().map(|chunk| chunk.bytes.len()).sum()
248    }
249
250    fn hot_start_offset(&self) -> u64 {
251        self.chunks
252            .front()
253            .map(|chunk| chunk.start_offset)
254            .unwrap_or(0)
255    }
256
257    fn payload(&self) -> Vec<u8> {
258        let mut payload = Vec::with_capacity(self.len());
259        for chunk in &self.chunks {
260            payload.extend_from_slice(&chunk.bytes);
261        }
262        payload
263    }
264
265    fn hot_segments(&self) -> Vec<HotPayloadSegment> {
266        let mut payload_start = 0usize;
267        self.chunks
268            .iter()
269            .map(|chunk| {
270                let payload_end = payload_start + chunk.bytes.len();
271                let segment = HotPayloadSegment {
272                    start_offset: chunk.start_offset,
273                    end_offset: chunk.end_offset,
274                    payload_start,
275                    payload_end,
276                };
277                payload_start = payload_end;
278                segment
279            })
280            .collect()
281    }
282
283    fn push(&mut self, start_offset: u64, end_offset: u64, payload: &[u8]) {
284        if payload.is_empty() {
285            return;
286        }
287        self.chunks.push_back(HotChunk {
288            start_offset,
289            end_offset,
290            bytes: payload.to_vec(),
291        });
292    }
293
294    fn plan_cold_flush(
295        &self,
296        min_hot_bytes: usize,
297        max_flush_bytes: usize,
298    ) -> Option<(u64, u64, Vec<u8>)> {
299        let first = self.chunks.front()?;
300        let mut payload = Vec::new();
301        let mut end_offset = first.start_offset;
302        for chunk in &self.chunks {
303            if chunk.start_offset != end_offset || payload.len() >= max_flush_bytes {
304                break;
305            }
306            let remaining = max_flush_bytes - payload.len();
307            let take = chunk.bytes.len().min(remaining);
308            payload.extend_from_slice(&chunk.bytes[..take]);
309            end_offset = end_offset.saturating_add(u64::try_from(take).expect("take fits u64"));
310            if take < chunk.bytes.len() {
311                break;
312            }
313        }
314        if payload.len() < min_hot_bytes {
315            return None;
316        }
317        Some((first.start_offset, end_offset, payload))
318    }
319
320    fn read_segments(&self, offset: u64, next_offset: u64) -> Vec<(u64, StreamReadSegment)> {
321        let mut segments = Vec::new();
322        for chunk in &self.chunks {
323            let start = offset.max(chunk.start_offset);
324            let end = next_offset.min(chunk.end_offset);
325            if start < end {
326                let payload_start =
327                    usize::try_from(start - chunk.start_offset).expect("hot start fits usize");
328                let payload_end =
329                    usize::try_from(end - chunk.start_offset).expect("hot end fits usize");
330                segments.push((
331                    start,
332                    StreamReadSegment::Hot(chunk.bytes[payload_start..payload_end].to_vec()),
333                ));
334            }
335        }
336        segments
337    }
338
339    fn covers_prefix(&self, start_offset: u64, end_offset: u64) -> bool {
340        let Some(first) = self.chunks.front() else {
341            return false;
342        };
343        if first.start_offset != start_offset {
344            return false;
345        }
346        let mut covered_offset = start_offset;
347        for chunk in &self.chunks {
348            if chunk.start_offset != covered_offset {
349                return false;
350            }
351            if chunk.end_offset >= end_offset {
352                return true;
353            }
354            covered_offset = chunk.end_offset;
355        }
356        false
357    }
358
359    fn flush_prefix(&mut self, end_offset: u64) {
360        while self
361            .chunks
362            .front()
363            .is_some_and(|chunk| chunk.end_offset <= end_offset)
364        {
365            self.chunks.pop_front();
366        }
367        if let Some(front) = self.chunks.front_mut()
368            && front.start_offset < end_offset
369        {
370            let drain_len =
371                usize::try_from(end_offset - front.start_offset).expect("drain len fits usize");
372            front.bytes.drain(..drain_len);
373            front.start_offset = end_offset;
374        }
375    }
376
377    fn discard_before(&mut self, retained_offset: u64) {
378        self.flush_prefix(retained_offset);
379    }
380}
381
382impl StreamStateMachine {
383    pub fn new() -> Self {
384        Self::default()
385    }
386
387    pub fn apply(&mut self, command: StreamCommand) -> StreamResponse {
388        match command {
389            StreamCommand::CreateBucket { bucket_id } => self.create_bucket(bucket_id),
390            StreamCommand::DeleteBucket { bucket_id } => self.delete_bucket(&bucket_id),
391            StreamCommand::CreateStream {
392                stream_id,
393                content_type,
394                initial_payload,
395                close_after,
396                stream_seq,
397                producer,
398                stream_ttl_seconds,
399                stream_expires_at_ms,
400                forked_from,
401                fork_offset,
402                now_ms,
403            } => self.create_stream(CreateStreamInput {
404                stream_id,
405                content_type,
406                initial_payload,
407                close_after,
408                stream_seq,
409                producer,
410                stream_ttl_seconds,
411                stream_expires_at_ms,
412                forked_from,
413                fork_offset,
414                now_ms,
415            }),
416            StreamCommand::CreateExternal {
417                stream_id,
418                content_type,
419                initial_payload,
420                close_after,
421                stream_seq,
422                producer,
423                stream_ttl_seconds,
424                stream_expires_at_ms,
425                forked_from,
426                fork_offset,
427                now_ms,
428            } => self.create_external_stream(CreateExternalStreamInput {
429                stream_id,
430                content_type,
431                initial_payload,
432                close_after,
433                stream_seq,
434                producer,
435                stream_ttl_seconds,
436                stream_expires_at_ms,
437                forked_from,
438                fork_offset,
439                now_ms,
440            }),
441            StreamCommand::Append {
442                stream_id,
443                content_type,
444                payload,
445                close_after,
446                stream_seq,
447                producer,
448                now_ms,
449            } => self.append_borrowed(AppendStreamInput {
450                stream_id,
451                content_type: content_type.as_deref(),
452                payload: &payload,
453                close_after,
454                stream_seq,
455                producer,
456                now_ms,
457            }),
458            StreamCommand::AppendExternal {
459                stream_id,
460                content_type,
461                payload,
462                close_after,
463                stream_seq,
464                producer,
465                now_ms,
466            } => self.append_external(AppendExternalInput {
467                stream_id,
468                content_type: content_type.as_deref(),
469                payload,
470                close_after,
471                stream_seq,
472                producer,
473                now_ms,
474            }),
475            StreamCommand::AppendBatch {
476                stream_id,
477                content_type,
478                payloads,
479                producer,
480                now_ms,
481            } => match self.append_batch_borrowed(
482                stream_id,
483                content_type.as_deref(),
484                &payloads.iter().map(Vec::as_slice).collect::<Vec<_>>(),
485                producer,
486                now_ms,
487            ) {
488                Ok(batch) => batch
489                    .items
490                    .last()
491                    .map(|item| StreamResponse::Appended {
492                        offset: item.offset,
493                        next_offset: item.next_offset,
494                        closed: item.closed,
495                        deduplicated: item.deduplicated,
496                        producer: None,
497                    })
498                    .unwrap_or_else(|| {
499                        StreamResponse::error(
500                            StreamErrorCode::EmptyAppend,
501                            "append batch must contain at least one payload",
502                        )
503                    }),
504                Err(response) => response,
505            },
506            StreamCommand::PublishSnapshot {
507                stream_id,
508                snapshot_offset,
509                content_type,
510                payload,
511                now_ms,
512            } => self.publish_snapshot(stream_id, snapshot_offset, content_type, payload, now_ms),
513            StreamCommand::TouchStreamAccess {
514                stream_id,
515                now_ms,
516                renew_ttl,
517            } => self.touch_stream_access(&stream_id, now_ms, renew_ttl),
518            StreamCommand::AddForkRef { stream_id, now_ms } => {
519                self.add_fork_ref(&stream_id, now_ms)
520            }
521            StreamCommand::ReleaseForkRef { stream_id } => self.release_fork_ref(&stream_id),
522            StreamCommand::FlushCold { stream_id, chunk } => self.flush_cold(stream_id, chunk),
523            StreamCommand::Close {
524                stream_id,
525                stream_seq,
526                producer,
527                now_ms,
528            } => self.close(stream_id, stream_seq, producer, now_ms),
529            StreamCommand::DeleteStream { stream_id } => self.delete_stream(&stream_id),
530            StreamCommand::AckColdGc { up_to_seq } => self.ack_cold_gc(up_to_seq),
531        }
532    }
533
534    pub fn head(&self, stream_id: &BucketStreamId) -> Option<&StreamMetadata> {
535        self.streams.get(stream_id)
536    }
537
538    pub fn head_at(&mut self, stream_id: &BucketStreamId, now_ms: u64) -> Option<&StreamMetadata> {
539        self.expire_stream_if_due(stream_id, now_ms);
540        self.streams.get(stream_id)
541    }
542
543    pub fn access_requires_write(
544        &self,
545        stream_id: &BucketStreamId,
546        now_ms: u64,
547        renew_ttl: bool,
548    ) -> Result<bool, StreamResponse> {
549        self.validate_stream_scope(stream_id)?;
550        let Some(stream) = self.streams.get(stream_id) else {
551            return Err(StreamResponse::error(
552                StreamErrorCode::StreamNotFound,
553                format!("stream '{stream_id}' does not exist"),
554            ));
555        };
556        if is_soft_deleted(stream) {
557            return Err(StreamResponse::error(
558                StreamErrorCode::StreamGone,
559                format!("stream '{stream_id}' is gone"),
560            ));
561        }
562        if stream_is_expired(stream, now_ms) {
563            return Ok(true);
564        }
565        Ok(renew_ttl
566            && stream.stream_ttl_seconds.is_some()
567            && stream.last_ttl_touch_at_ms != now_ms)
568    }
569
570    pub fn hot_start_offset(&self, stream_id: &BucketStreamId) -> u64 {
571        let Some(stream) = self.streams.get(stream_id) else {
572            return 0;
573        };
574        self.hot_buffers
575            .get(stream_id)
576            .and_then(|buffer| buffer.chunks.front().map(|chunk| chunk.start_offset))
577            .unwrap_or(stream.tail_offset)
578    }
579
580    pub fn cold_chunks(&self, stream_id: &BucketStreamId) -> &[ColdChunkRef] {
581        self.cold_index().cold_chunks(stream_id)
582    }
583
584    pub fn external_segments(&self, stream_id: &BucketStreamId) -> &[ObjectPayloadRef] {
585        self.cold_index().external_segments(stream_id)
586    }
587
588    pub fn hot_segments(&self, stream_id: &BucketStreamId) -> Vec<HotPayloadSegment> {
589        self.hot_buffers
590            .get(stream_id)
591            .map(HotBuffer::hot_segments)
592            .unwrap_or_default()
593    }
594
595    pub fn hot_payload_len(&self, stream_id: &BucketStreamId) -> Result<u64, StreamResponse> {
596        let Some(stream) = self.streams.get(stream_id) else {
597            return Err(StreamResponse::error(
598                StreamErrorCode::StreamNotFound,
599                format!("stream '{stream_id}' does not exist"),
600            ));
601        };
602        if is_soft_deleted(stream) {
603            return Err(StreamResponse::error(
604                StreamErrorCode::StreamGone,
605                format!("stream '{stream_id}' is gone"),
606            ));
607        }
608        let payload = self
609            .hot_buffers
610            .get(stream_id)
611            .expect("hot buffer exists for stream metadata");
612        Ok(u64::try_from(payload.len()).expect("payload len fits u64"))
613    }
614
615    pub fn total_hot_payload_bytes(&self) -> u64 {
616        self.hot_buffers
617            .values()
618            .map(|payload| u64::try_from(payload.len()).expect("payload len fits u64"))
619            .sum()
620    }
621
622    pub fn plan_cold_flush(
623        &self,
624        stream_id: &BucketStreamId,
625        min_hot_bytes: usize,
626        max_flush_bytes: usize,
627    ) -> Result<Option<ColdFlushCandidate>, StreamResponse> {
628        if max_flush_bytes == 0 {
629            return Ok(None);
630        }
631        let Some(stream) = self.streams.get(stream_id) else {
632            return Err(StreamResponse::error(
633                StreamErrorCode::StreamNotFound,
634                format!("stream '{stream_id}' does not exist"),
635            ));
636        };
637        if is_soft_deleted(stream) {
638            return Err(StreamResponse::error(
639                StreamErrorCode::StreamGone,
640                format!("stream '{stream_id}' is gone"),
641            ));
642        }
643        let Some(hot_buffer) = self.hot_buffers.get(stream_id) else {
644            return Ok(None);
645        };
646        let Some((start_offset, end_offset, payload)) =
647            hot_buffer.plan_cold_flush(min_hot_bytes, max_flush_bytes)
648        else {
649            return Ok(None);
650        };
651        Ok(Some(ColdFlushCandidate {
652            stream_id: stream_id.clone(),
653            start_offset,
654            end_offset,
655            payload,
656        }))
657    }
658
659    pub fn plan_next_cold_flush(
660        &self,
661        min_hot_bytes: usize,
662        max_flush_bytes: usize,
663    ) -> Result<Option<ColdFlushCandidate>, StreamResponse> {
664        if max_flush_bytes == 0 {
665            return Ok(None);
666        }
667        let mut stream_ids = self.streams.keys().cloned().collect::<Vec<_>>();
668        stream_ids.sort_by(compare_stream_ids);
669        for stream_id in &stream_ids {
670            match self.plan_cold_flush(stream_id, min_hot_bytes, max_flush_bytes) {
671                Ok(Some(candidate)) => return Ok(Some(candidate)),
672                Ok(None) => {}
673                Err(StreamResponse::Error {
674                    code: StreamErrorCode::StreamGone | StreamErrorCode::StreamNotFound,
675                    ..
676                }) => {}
677                Err(err) => return Err(err),
678            }
679        }
680        let group_min_hot_bytes = u64::try_from(min_hot_bytes).unwrap_or(u64::MAX);
681        let group_hot_bytes = self.total_hot_payload_bytes();
682        if group_hot_bytes < group_min_hot_bytes {
683            return Ok(None);
684        }
685        for stream_id in stream_ids {
686            match self.plan_cold_flush(&stream_id, 1, max_flush_bytes) {
687                Ok(Some(candidate)) => return Ok(Some(candidate)),
688                Ok(None) => {}
689                Err(StreamResponse::Error {
690                    code: StreamErrorCode::StreamGone | StreamErrorCode::StreamNotFound,
691                    ..
692                }) => {}
693                Err(err) => return Err(err),
694            }
695        }
696        Ok(None)
697    }
698
699    pub fn plan_next_cold_flush_batch(
700        &self,
701        min_hot_bytes: usize,
702        max_flush_bytes: usize,
703        max_candidates: usize,
704    ) -> Result<Vec<ColdFlushCandidate>, StreamResponse> {
705        if max_candidates == 0 || max_flush_bytes == 0 {
706            return Ok(Vec::new());
707        }
708        let mut preview = self.clone();
709        let mut candidates = Vec::with_capacity(max_candidates);
710        while candidates.len() < max_candidates {
711            let Some(candidate) = preview.plan_next_cold_flush(min_hot_bytes, max_flush_bytes)?
712            else {
713                break;
714            };
715            let chunk = ColdChunkRef {
716                start_offset: candidate.start_offset,
717                end_offset: candidate.end_offset,
718                s3_path: "planned-cold-flush-batch".to_owned(),
719                object_size: u64::try_from(candidate.payload.len()).expect("payload len fits u64"),
720            };
721            match preview.flush_cold(candidate.stream_id.clone(), chunk) {
722                StreamResponse::ColdFlushed { .. } => candidates.push(candidate),
723                StreamResponse::Error { .. } => break,
724                other => {
725                    return Err(StreamResponse::error(
726                        StreamErrorCode::InvalidColdFlush,
727                        format!("unexpected cold flush planning response: {other:?}"),
728                    ));
729                }
730            }
731        }
732        Ok(candidates)
733    }
734
735    pub fn bucket_exists(&self, bucket_id: &str) -> bool {
736        self.buckets.contains(bucket_id)
737    }
738
739    fn cold_index(&self) -> &impl ColdIndexState {
740        &self.cold_index
741    }
742
743    fn cold_index_mut(&mut self) -> &mut impl ColdIndexState {
744        &mut self.cold_index
745    }
746
747    pub fn integrity_snapshot(
748        &self,
749        stream_id: &BucketStreamId,
750    ) -> Result<crate::integrity::StreamIntegritySnapshot, StreamResponse> {
751        let Some(stream) = self.streams.get(stream_id) else {
752            return Err(StreamResponse::error(
753                StreamErrorCode::StreamNotFound,
754                format!("stream '{stream_id}' does not exist"),
755            ));
756        };
757        if is_soft_deleted(stream) {
758            return Err(StreamResponse::error(
759                StreamErrorCode::StreamGone,
760                format!("stream '{stream_id}' is gone"),
761            ));
762        }
763        Ok(self
764            .integrities
765            .get(stream_id)
766            .expect("integrity exists for stream metadata")
767            .snapshot(self.earliest_retained_offset(stream_id), stream.tail_offset))
768    }
769
770    pub fn snapshot(&self) -> StreamSnapshot {
771        let mut buckets = self.buckets.iter().cloned().collect::<Vec<_>>();
772        buckets.sort();
773
774        let mut streams = self
775            .streams
776            .values()
777            .cloned()
778            .map(|metadata| {
779                let stream_id = metadata.stream_id.clone();
780                let tail_offset = metadata.tail_offset;
781                let hot_buffer = self
782                    .hot_buffers
783                    .get(&stream_id)
784                    .expect("hot buffer exists for stream metadata");
785                let payload = hot_buffer.payload();
786                let producer_states = self.producer_snapshot(&stream_id);
787                StreamSnapshotEntry {
788                    metadata,
789                    hot_start_offset: self.hot_start_offset(&stream_id),
790                    payload,
791                    hot_segments: hot_buffer.hot_segments(),
792                    cold_frontier_offset: self.cold_frontier_offset(
793                        &stream_id,
794                        self.earliest_retained_offset(&stream_id),
795                    ),
796                    cold_index_generation: self.cold_index().cold_generation(&stream_id),
797                    cold_chunks: self.cold_index().cold_chunks(&stream_id).to_vec(),
798                    external_segments: self.cold_index().external_segments(&stream_id).to_vec(),
799                    message_records: self
800                        .message_records
801                        .get(&stream_id)
802                        .cloned()
803                        .unwrap_or_default(),
804                    integrity: self
805                        .integrities
806                        .get(&stream_id)
807                        .expect("integrity exists for stream metadata")
808                        .snapshot(self.earliest_retained_offset(&stream_id), tail_offset),
809                    visible_snapshot: self.visible_snapshots.get(&stream_id).cloned(),
810                    producer_states,
811                }
812            })
813            .collect::<Vec<_>>();
814        streams.sort_by(|left, right| {
815            compare_stream_ids(&left.metadata.stream_id, &right.metadata.stream_id)
816        });
817
818        StreamSnapshot {
819            buckets,
820            streams,
821            pending_cold_gc: self.pending_cold_gc.iter().cloned().collect(),
822            next_cold_gc_seq: self.next_cold_gc_seq,
823        }
824    }
825
826    pub fn restore(snapshot: StreamSnapshot) -> Result<Self, StreamSnapshotError> {
827        let mut machine = Self::default();
828        for bucket_id in snapshot.buckets {
829            if !machine.buckets.insert(bucket_id.clone()) {
830                return Err(StreamSnapshotError::DuplicateBucket(bucket_id));
831            }
832        }
833
834        for entry in snapshot.streams {
835            let stream_id = entry.metadata.stream_id.clone();
836            if !machine.buckets.contains(&stream_id.bucket_id) {
837                return Err(StreamSnapshotError::MissingBucket(stream_id));
838            }
839            if let Some(snapshot) = entry.visible_snapshot.as_ref()
840                && snapshot.offset > entry.metadata.tail_offset
841            {
842                return Err(StreamSnapshotError::SnapshotOffsetOutOfRange {
843                    stream_id,
844                    snapshot_offset: snapshot.offset,
845                    tail_offset: entry.metadata.tail_offset,
846                });
847            }
848            let retained_offset = entry
849                .visible_snapshot
850                .as_ref()
851                .map(|snapshot| snapshot.offset)
852                .unwrap_or(0);
853            let hot_segments = if entry.hot_segments.is_empty() && !entry.payload.is_empty() {
854                vec![HotPayloadSegment {
855                    start_offset: entry.hot_start_offset,
856                    end_offset: entry.metadata.tail_offset,
857                    payload_start: 0,
858                    payload_end: entry.payload.len(),
859                }]
860            } else {
861                entry.hot_segments
862            };
863            if !hot_segments_match_payload(&hot_segments, entry.payload.len())
864                || !payload_sources_cover_retained_suffix(
865                    entry.cold_frontier_offset,
866                    &entry.cold_chunks,
867                    &entry.external_segments,
868                    &hot_segments,
869                    retained_offset,
870                    entry.metadata.tail_offset,
871                )
872            {
873                return Err(StreamSnapshotError::PayloadLengthMismatch {
874                    stream_id,
875                    tail_offset: entry.metadata.tail_offset,
876                    payload_len: entry.payload.len(),
877                });
878            }
879            if !message_records_cover_retained_suffix(
880                &entry.message_records,
881                retained_offset,
882                entry.metadata.tail_offset,
883            ) {
884                return Err(StreamSnapshotError::MessageBoundaryMismatch { stream_id });
885            }
886            let integrity = StreamIntegrity::restore(entry.integrity).ok_or_else(|| {
887                StreamSnapshotError::IntegrityMismatch {
888                    stream_id: stream_id.clone(),
889                }
890            })?;
891            if machine
892                .streams
893                .insert(entry.metadata.stream_id.clone(), entry.metadata)
894                .is_some()
895            {
896                return Err(StreamSnapshotError::DuplicateStream(stream_id));
897            }
898            let producer_states = restore_producer_states(&stream_id, entry.producer_states)?;
899            machine.hot_buffers.insert(
900                stream_id.clone(),
901                HotBuffer::from_snapshot(entry.payload, &hot_segments),
902            );
903            machine.cold_index_mut().restore_stream(
904                stream_id.clone(),
905                entry.cold_frontier_offset,
906                entry.cold_index_generation,
907                entry.cold_chunks,
908                entry.external_segments,
909            );
910            if !entry.message_records.is_empty() {
911                machine
912                    .message_records
913                    .insert(stream_id.clone(), entry.message_records);
914            }
915            machine.integrities.insert(stream_id.clone(), integrity);
916            if let Some(snapshot) = entry.visible_snapshot {
917                machine
918                    .visible_snapshots
919                    .insert(stream_id.clone(), snapshot);
920            }
921            machine.producers.insert(stream_id.clone(), producer_states);
922        }
923
924        machine.pending_cold_gc = snapshot.pending_cold_gc.into_iter().collect();
925        machine.next_cold_gc_seq = snapshot.next_cold_gc_seq;
926
927        Ok(machine)
928    }
929
930    pub fn read(
931        &self,
932        stream_id: &BucketStreamId,
933        offset: u64,
934        max_len: usize,
935    ) -> Result<StreamRead, StreamResponse> {
936        let plan = self.read_plan(stream_id, offset, max_len)?;
937        if plan.segments.iter().any(|segment| {
938            matches!(
939                segment,
940                StreamReadSegment::ColdIndex(_) | StreamReadSegment::Object(_)
941            )
942        }) {
943            return Err(StreamResponse::error_with_next_offset(
944                StreamErrorCode::InvalidColdFlush,
945                format!("stream '{stream_id}' read requires object payload store"),
946                plan.next_offset,
947            ));
948        }
949        let payload = plan
950            .segments
951            .iter()
952            .flat_map(|segment| match segment {
953                StreamReadSegment::Hot(payload) => payload.as_slice(),
954                StreamReadSegment::ColdIndex(_) | StreamReadSegment::Object(_) => {
955                    unreachable!("object segments checked above")
956                }
957            })
958            .copied()
959            .collect();
960        Ok(StreamRead {
961            offset: plan.offset,
962            next_offset: plan.next_offset,
963            content_type: plan.content_type,
964            payload,
965            up_to_date: plan.up_to_date,
966            closed: plan.closed,
967        })
968    }
969
970    pub fn read_plan(
971        &self,
972        stream_id: &BucketStreamId,
973        offset: u64,
974        max_len: usize,
975    ) -> Result<StreamReadPlan, StreamResponse> {
976        self.read_plan_at(stream_id, offset, max_len, 0)
977    }
978
979    pub fn read_plan_at(
980        &self,
981        stream_id: &BucketStreamId,
982        offset: u64,
983        max_len: usize,
984        now_ms: u64,
985    ) -> Result<StreamReadPlan, StreamResponse> {
986        let Some(stream) = self.streams.get(stream_id) else {
987            return Err(StreamResponse::error(
988                StreamErrorCode::StreamNotFound,
989                format!("stream '{stream_id}' does not exist"),
990            ));
991        };
992        if is_soft_deleted(stream) {
993            return Err(StreamResponse::error(
994                StreamErrorCode::StreamGone,
995                format!("stream '{stream_id}' is gone"),
996            ));
997        }
998        if stream_is_expired(stream, now_ms) {
999            return Err(StreamResponse::error(
1000                StreamErrorCode::StreamNotFound,
1001                format!("stream '{stream_id}' does not exist"),
1002            ));
1003        }
1004        if offset > stream.tail_offset {
1005            return Err(StreamResponse::error_with_next_offset(
1006                StreamErrorCode::OffsetOutOfRange,
1007                format!(
1008                    "offset {offset} is beyond stream '{}' tail {}",
1009                    stream_id, stream.tail_offset
1010                ),
1011                stream.tail_offset,
1012            ));
1013        }
1014        let retained_offset = self.earliest_retained_offset(stream_id);
1015        if offset < retained_offset {
1016            return Err(StreamResponse::error_with_next_offset(
1017                StreamErrorCode::StreamGone,
1018                format!(
1019                    "offset {offset} is older than stream '{}' retained offset {retained_offset}",
1020                    stream_id
1021                ),
1022                retained_offset,
1023            ));
1024        }
1025
1026        let max_len_u64 = u64::try_from(max_len).unwrap_or(u64::MAX);
1027        let next_offset = stream.tail_offset.min(offset.saturating_add(max_len_u64));
1028        let mut segments = Vec::<(u64, StreamReadSegment)>::new();
1029        let hot_segments = self
1030            .hot_buffers
1031            .get(stream_id)
1032            .map(|hot_buffer| hot_buffer.read_segments(offset, next_offset))
1033            .unwrap_or_default();
1034        let cold_frontier = self.cold_frontier_offset(stream_id, retained_offset);
1035        let cold_index_end = next_offset.min(cold_frontier);
1036        let mut cursor = offset;
1037        for (hot_start, hot_segment) in &hot_segments {
1038            if cursor >= cold_index_end {
1039                break;
1040            }
1041            let Some(hot_end) = read_segment_end(*hot_start, hot_segment) else {
1042                continue;
1043            };
1044            if hot_end <= cursor {
1045                continue;
1046            }
1047            let gap_end = (*hot_start).min(cold_index_end);
1048            push_cold_index_segments(
1049                &mut segments,
1050                stream_id,
1051                self.cold_index().cold_generation(stream_id),
1052                cursor,
1053                gap_end,
1054            );
1055            cursor = cursor.max(hot_end);
1056        }
1057        push_cold_index_segments(
1058            &mut segments,
1059            stream_id,
1060            self.cold_index().cold_generation(stream_id),
1061            cursor,
1062            cold_index_end,
1063        );
1064        for chunk in self.cold_chunks(stream_id) {
1065            let start = offset.max(chunk.start_offset);
1066            let end = next_offset.min(chunk.end_offset);
1067            if start < end {
1068                segments.push((
1069                    start,
1070                    StreamReadSegment::Object(StreamReadObjectSegment {
1071                        object: ObjectPayloadRef::from(chunk),
1072                        read_start_offset: start,
1073                        len: usize::try_from(end - start).expect("object read len fits usize"),
1074                    }),
1075                ));
1076            }
1077        }
1078        for object in self.external_segments(stream_id) {
1079            let start = offset.max(object.start_offset);
1080            let end = next_offset.min(object.end_offset);
1081            if start < end {
1082                segments.push((
1083                    start,
1084                    StreamReadSegment::Object(StreamReadObjectSegment {
1085                        object: object.clone(),
1086                        read_start_offset: start,
1087                        len: usize::try_from(end - start).expect("object read len fits usize"),
1088                    }),
1089                ));
1090            }
1091        }
1092        segments.extend(hot_segments);
1093        segments.sort_by_key(|(start, _)| *start);
1094        if !segments_cover_range(&segments, offset, next_offset) {
1095            return Err(StreamResponse::error_with_next_offset(
1096                StreamErrorCode::InvalidColdFlush,
1097                format!("stream '{stream_id}' has missing payload segment metadata"),
1098                next_offset,
1099            ));
1100        }
1101        Ok(StreamReadPlan {
1102            offset,
1103            next_offset,
1104            content_type: stream.content_type.clone(),
1105            segments: segments.into_iter().map(|(_, segment)| segment).collect(),
1106            up_to_date: next_offset == stream.tail_offset,
1107            closed: stream.status == StreamStatus::Closed,
1108        })
1109    }
1110
1111    pub fn latest_snapshot(
1112        &self,
1113        stream_id: &BucketStreamId,
1114    ) -> Result<Option<StreamVisibleSnapshot>, StreamResponse> {
1115        let Some(stream) = self.streams.get(stream_id) else {
1116            return Err(StreamResponse::error(
1117                StreamErrorCode::StreamNotFound,
1118                format!("stream '{stream_id}' does not exist"),
1119            ));
1120        };
1121        if is_soft_deleted(stream) {
1122            return Err(StreamResponse::error(
1123                StreamErrorCode::StreamGone,
1124                format!("stream '{stream_id}' is gone"),
1125            ));
1126        }
1127        Ok(self.visible_snapshots.get(stream_id).cloned())
1128    }
1129
1130    pub fn read_snapshot(
1131        &self,
1132        stream_id: &BucketStreamId,
1133        snapshot_offset: u64,
1134    ) -> Result<StreamVisibleSnapshot, StreamResponse> {
1135        let snapshot = self.latest_snapshot(stream_id)?;
1136        match snapshot {
1137            Some(snapshot) if snapshot.offset == snapshot_offset => Ok(snapshot),
1138            _ => Err(StreamResponse::error(
1139                StreamErrorCode::SnapshotNotFound,
1140                format!("snapshot {snapshot_offset} for stream '{stream_id}' does not exist"),
1141            )),
1142        }
1143    }
1144
1145    pub fn delete_snapshot(
1146        &self,
1147        stream_id: &BucketStreamId,
1148        snapshot_offset: u64,
1149    ) -> StreamResponse {
1150        match self.latest_snapshot(stream_id) {
1151            Ok(Some(snapshot)) if snapshot.offset == snapshot_offset => StreamResponse::error(
1152                StreamErrorCode::SnapshotConflict,
1153                format!(
1154                    "snapshot {snapshot_offset} for stream '{stream_id}' is the latest visible snapshot"
1155                ),
1156            ),
1157            Ok(_) => StreamResponse::error(
1158                StreamErrorCode::SnapshotNotFound,
1159                format!("snapshot {snapshot_offset} for stream '{stream_id}' does not exist"),
1160            ),
1161            Err(err) => err,
1162        }
1163    }
1164
1165    pub fn bootstrap_plan(
1166        &self,
1167        stream_id: &BucketStreamId,
1168    ) -> Result<StreamBootstrapPlan, StreamResponse> {
1169        let Some(stream) = self.streams.get(stream_id) else {
1170            return Err(StreamResponse::error(
1171                StreamErrorCode::StreamNotFound,
1172                format!("stream '{stream_id}' does not exist"),
1173            ));
1174        };
1175        if is_soft_deleted(stream) {
1176            return Err(StreamResponse::error(
1177                StreamErrorCode::StreamGone,
1178                format!("stream '{stream_id}' is gone"),
1179            ));
1180        }
1181        let snapshot = self.visible_snapshots.get(stream_id).cloned();
1182        let retained_offset = snapshot
1183            .as_ref()
1184            .map(|snapshot| snapshot.offset)
1185            .unwrap_or(0);
1186        let updates = self
1187            .message_records
1188            .get(stream_id)
1189            .map(|records| {
1190                records
1191                    .iter()
1192                    .filter(|record| record.start_offset >= retained_offset)
1193                    .cloned()
1194                    .collect::<Vec<_>>()
1195            })
1196            .unwrap_or_default();
1197        Ok(StreamBootstrapPlan {
1198            snapshot,
1199            updates,
1200            next_offset: stream.tail_offset,
1201            content_type: stream.content_type.clone(),
1202            up_to_date: true,
1203            closed: stream.status == StreamStatus::Closed,
1204        })
1205    }
1206
1207    fn publish_snapshot(
1208        &mut self,
1209        stream_id: BucketStreamId,
1210        snapshot_offset: u64,
1211        content_type: String,
1212        payload: Vec<u8>,
1213        now_ms: u64,
1214    ) -> StreamResponse {
1215        if let Err(response) = self.validate_stream_scope(&stream_id) {
1216            return response;
1217        }
1218        if content_type.trim().is_empty() {
1219            return StreamResponse::error(
1220                StreamErrorCode::InvalidSnapshot,
1221                "snapshot content type must not be empty",
1222            );
1223        }
1224        let Some(stream) = self.streams.get(&stream_id) else {
1225            return StreamResponse::error(
1226                StreamErrorCode::StreamNotFound,
1227                format!("stream '{stream_id}' does not exist"),
1228            );
1229        };
1230        if is_soft_deleted(stream) {
1231            return StreamResponse::error(
1232                StreamErrorCode::StreamGone,
1233                format!("stream '{stream_id}' is gone"),
1234            );
1235        }
1236        if stream_is_expired(stream, now_ms) {
1237            self.remove_stream_state(&stream_id);
1238            return StreamResponse::error(
1239                StreamErrorCode::StreamNotFound,
1240                format!("stream '{stream_id}' does not exist"),
1241            );
1242        }
1243        let tail_offset = stream.tail_offset;
1244        let retained_offset = self.earliest_retained_offset(&stream_id);
1245        if snapshot_offset < retained_offset {
1246            return StreamResponse::error_with_next_offset(
1247                StreamErrorCode::StreamGone,
1248                format!(
1249                    "snapshot offset {snapshot_offset} is older than stream '{}' retained offset {retained_offset}",
1250                    stream_id
1251                ),
1252                retained_offset,
1253            );
1254        }
1255        if snapshot_offset > tail_offset {
1256            return StreamResponse::error_with_next_offset(
1257                StreamErrorCode::SnapshotConflict,
1258                format!(
1259                    "snapshot offset {snapshot_offset} is beyond stream '{}' tail {tail_offset}",
1260                    stream_id
1261                ),
1262                tail_offset,
1263            );
1264        }
1265        if !self.snapshot_offset_aligned(&stream_id, snapshot_offset, retained_offset) {
1266            return StreamResponse::error_with_next_offset(
1267                StreamErrorCode::InvalidSnapshot,
1268                format!(
1269                    "snapshot offset {snapshot_offset} is not aligned to a committed message boundary for stream '{stream_id}'"
1270                ),
1271                tail_offset,
1272            );
1273        }
1274
1275        self.visible_snapshots
1276            .insert(stream_id.clone(), StreamVisibleSnapshot {
1277                offset: snapshot_offset,
1278                content_type,
1279                payload,
1280            });
1281        self.compact_retained_prefix(&stream_id, snapshot_offset);
1282        StreamResponse::SnapshotPublished { snapshot_offset }
1283    }
1284
1285    fn flush_cold(&mut self, stream_id: BucketStreamId, chunk: ColdChunkRef) -> StreamResponse {
1286        if let Err(response) = self.validate_stream_scope(&stream_id) {
1287            return response;
1288        }
1289        if chunk.s3_path.trim().is_empty() {
1290            return StreamResponse::error(
1291                StreamErrorCode::InvalidColdFlush,
1292                "cold chunk S3 path must not be empty",
1293            );
1294        }
1295        if chunk.object_size == 0 {
1296            return StreamResponse::error(
1297                StreamErrorCode::InvalidColdFlush,
1298                "cold chunk object size must be greater than zero",
1299            );
1300        }
1301        let Some(stream) = self.streams.get(&stream_id) else {
1302            return StreamResponse::error(
1303                StreamErrorCode::StreamNotFound,
1304                format!("stream '{stream_id}' does not exist"),
1305            );
1306        };
1307        if is_soft_deleted(stream) {
1308            return StreamResponse::error(
1309                StreamErrorCode::StreamGone,
1310                format!("stream '{stream_id}' is gone"),
1311            );
1312        }
1313        if chunk.end_offset <= chunk.start_offset {
1314            return StreamResponse::error_with_next_offset(
1315                StreamErrorCode::InvalidColdFlush,
1316                "cold chunk must cover at least one byte",
1317                stream.tail_offset,
1318            );
1319        }
1320        if chunk.end_offset > stream.tail_offset {
1321            return StreamResponse::error_with_next_offset_and_context(
1322                StreamErrorCode::InvalidColdFlush,
1323                format!(
1324                    "cold chunk end {} is beyond stream '{}' tail {}",
1325                    chunk.end_offset, stream_id, stream.tail_offset
1326                ),
1327                stream.tail_offset,
1328                vec![StreamErrorContext::StaleColdFlushCandidate],
1329            );
1330        }
1331        let Some(hot_buffer) = self.hot_buffers.get(&stream_id) else {
1332            return StreamResponse::error_with_next_offset_and_context(
1333                StreamErrorCode::InvalidColdFlush,
1334                format!("cold chunk for stream '{stream_id}' does not match hot payload"),
1335                stream.tail_offset,
1336                vec![StreamErrorContext::StaleColdFlushCandidate],
1337            );
1338        };
1339        if hot_buffer.hot_start_offset() != chunk.start_offset {
1340            return StreamResponse::error_with_next_offset_and_context(
1341                StreamErrorCode::InvalidColdFlush,
1342                format!("cold chunk for stream '{stream_id}' must start at the hot prefix"),
1343                stream.tail_offset,
1344                vec![StreamErrorContext::StaleColdFlushCandidate],
1345            );
1346        }
1347        if !hot_buffer.covers_prefix(chunk.start_offset, chunk.end_offset) {
1348            return StreamResponse::error_with_next_offset_and_context(
1349                StreamErrorCode::InvalidColdFlush,
1350                format!(
1351                    "cold chunk for stream '{stream_id}' does not cover contiguous hot payload"
1352                ),
1353                stream.tail_offset,
1354                vec![StreamErrorContext::StaleColdFlushCandidate],
1355            );
1356        }
1357        self.hot_buffers
1358            .get_mut(&stream_id)
1359            .expect("hot buffer exists for stream metadata")
1360            .flush_prefix(chunk.end_offset);
1361        self.cold_index_mut()
1362            .push_cold_chunk(stream_id.clone(), chunk.clone());
1363        self.compact_message_records_before(
1364            &stream_id,
1365            self.earliest_retained_offset(&stream_id),
1366            chunk.end_offset,
1367        );
1368        StreamResponse::ColdFlushed {
1369            hot_start_offset: self.hot_start_offset(&stream_id),
1370        }
1371    }
1372
1373    fn create_bucket(&mut self, bucket_id: String) -> StreamResponse {
1374        if let Err(message) = validate_bucket_id(&bucket_id) {
1375            return StreamResponse::error(StreamErrorCode::InvalidBucketId, message);
1376        }
1377        if !self.buckets.insert(bucket_id.clone()) {
1378            return StreamResponse::BucketAlreadyExists { bucket_id };
1379        }
1380        StreamResponse::BucketCreated { bucket_id }
1381    }
1382
1383    fn delete_bucket(&mut self, bucket_id: &str) -> StreamResponse {
1384        if let Err(message) = validate_bucket_id(bucket_id) {
1385            return StreamResponse::error(StreamErrorCode::InvalidBucketId, message);
1386        }
1387        if !self.buckets.contains(bucket_id) {
1388            return StreamResponse::error(
1389                StreamErrorCode::BucketNotFound,
1390                format!("bucket '{bucket_id}' does not exist"),
1391            );
1392        }
1393        if self
1394            .streams
1395            .keys()
1396            .any(|stream_id| stream_id.bucket_id == bucket_id)
1397        {
1398            return StreamResponse::error(
1399                StreamErrorCode::BucketNotEmpty,
1400                format!("bucket '{bucket_id}' is not empty"),
1401            );
1402        }
1403        self.buckets.remove(bucket_id);
1404        StreamResponse::BucketDeleted {
1405            bucket_id: bucket_id.to_owned(),
1406        }
1407    }
1408
1409    fn create_stream(&mut self, input: CreateStreamInput) -> StreamResponse {
1410        if let Err(response) = self.validate_stream_scope(&input.stream_id) {
1411            return response;
1412        }
1413        if let Err(response) =
1414            validate_retention(input.stream_ttl_seconds, input.stream_expires_at_ms)
1415        {
1416            return response;
1417        }
1418        if let Err(response) = validate_producer_request(input.producer.as_ref()) {
1419            return response;
1420        }
1421        if let Some(producer) = input.producer.as_ref()
1422            && producer.producer_seq != 0
1423        {
1424            return StreamResponse::error_with_context(
1425                StreamErrorCode::ProducerSeqConflict,
1426                format!(
1427                    "producer '{}' expected sequence 0, received {}",
1428                    producer.producer_id, producer.producer_seq
1429                ),
1430                vec![StreamErrorContext::ProducerSeqConflict {
1431                    expected_seq: 0,
1432                    received_seq: producer.producer_seq,
1433                }],
1434            );
1435        }
1436        if self
1437            .streams
1438            .get(&input.stream_id)
1439            .is_some_and(|existing| stream_is_expired(existing, input.now_ms))
1440        {
1441            self.remove_stream_state(&input.stream_id);
1442        }
1443
1444        if let Some(existing) = self.streams.get(&input.stream_id) {
1445            if is_soft_deleted(existing) {
1446                return StreamResponse::error(
1447                    StreamErrorCode::StreamAlreadyExistsConflict,
1448                    format!(
1449                        "stream '{}' is gone and cannot be recreated yet",
1450                        input.stream_id
1451                    ),
1452                );
1453            }
1454            if existing.content_type == input.content_type
1455                && existing.status == status_from_closed(input.close_after)
1456                && existing.stream_ttl_seconds == input.stream_ttl_seconds
1457                && existing.stream_expires_at_ms == input.stream_expires_at_ms
1458                && existing.forked_from == input.forked_from
1459                && existing.fork_offset == input.fork_offset
1460            {
1461                return StreamResponse::AlreadyExists {
1462                    next_offset: existing.tail_offset,
1463                    closed: existing.status == StreamStatus::Closed,
1464                    content_type: existing.content_type.clone(),
1465                    stream_ttl_seconds: existing.stream_ttl_seconds,
1466                    stream_expires_at_ms: existing.stream_expires_at_ms,
1467                };
1468            }
1469            return StreamResponse::error(
1470                StreamErrorCode::StreamAlreadyExistsConflict,
1471                format!(
1472                    "stream '{}' already exists with different metadata",
1473                    input.stream_id
1474                ),
1475            );
1476        }
1477
1478        let initial_len = input.initial_len();
1479        let metadata = StreamMetadata {
1480            stream_id: input.stream_id.clone(),
1481            content_type: input.content_type,
1482            status: status_from_closed(input.close_after),
1483            tail_offset: initial_len,
1484            last_stream_seq: input.stream_seq,
1485            stream_ttl_seconds: input.stream_ttl_seconds,
1486            stream_expires_at_ms: input.stream_expires_at_ms,
1487            created_at_ms: input.now_ms,
1488            last_ttl_touch_at_ms: input.now_ms,
1489            forked_from: input.forked_from,
1490            fork_offset: input.fork_offset,
1491            fork_ref_count: 0,
1492        };
1493        self.streams.insert(input.stream_id.clone(), metadata);
1494        self.hot_buffers.insert(
1495            input.stream_id.clone(),
1496            HotBuffer::from_payload(0, input.initial_payload),
1497        );
1498        let mut integrity = StreamIntegrity::default();
1499        if initial_len > 0 {
1500            let payload = self
1501                .hot_buffers
1502                .get(&input.stream_id)
1503                .expect("hot buffer exists for stream metadata")
1504                .payload();
1505            integrity.append_payload(&input.stream_id, 0, initial_len, &payload);
1506        }
1507        self.integrities.insert(input.stream_id.clone(), integrity);
1508        if initial_len > 0 {
1509            self.message_records
1510                .insert(input.stream_id.clone(), vec![StreamMessageRecord {
1511                    start_offset: 0,
1512                    end_offset: initial_len,
1513                }]);
1514        }
1515        let mut producer_states = HashMap::new();
1516        if let Some(producer) = input.producer {
1517            let last_item = ProducerAppendRecord {
1518                start_offset: 0,
1519                next_offset: initial_len,
1520                closed: input.close_after,
1521            };
1522            producer_states.insert(producer.producer_id, ProducerState {
1523                producer_epoch: producer.producer_epoch,
1524                producer_seq: producer.producer_seq,
1525                last_start_offset: last_item.start_offset,
1526                last_next_offset: last_item.next_offset,
1527                last_closed: last_item.closed,
1528                last_items: vec![last_item],
1529            });
1530        }
1531        self.producers
1532            .insert(input.stream_id.clone(), producer_states);
1533        StreamResponse::Created {
1534            stream_id: input.stream_id,
1535            next_offset: initial_len,
1536            closed: input.close_after,
1537        }
1538    }
1539
1540    fn create_external_stream(&mut self, input: CreateExternalStreamInput) -> StreamResponse {
1541        if let Err(response) = validate_external_payload_ref(&input.initial_payload) {
1542            return response;
1543        }
1544        if let Err(response) = self.validate_stream_scope(&input.stream_id) {
1545            return response;
1546        }
1547        if let Err(response) =
1548            validate_retention(input.stream_ttl_seconds, input.stream_expires_at_ms)
1549        {
1550            return response;
1551        }
1552        if let Err(response) = validate_producer_request(input.producer.as_ref()) {
1553            return response;
1554        }
1555        if let Some(producer) = input.producer.as_ref()
1556            && producer.producer_seq != 0
1557        {
1558            return StreamResponse::error_with_context(
1559                StreamErrorCode::ProducerSeqConflict,
1560                format!(
1561                    "producer '{}' expected sequence 0, received {}",
1562                    producer.producer_id, producer.producer_seq
1563                ),
1564                vec![StreamErrorContext::ProducerSeqConflict {
1565                    expected_seq: 0,
1566                    received_seq: producer.producer_seq,
1567                }],
1568            );
1569        }
1570        if self
1571            .streams
1572            .get(&input.stream_id)
1573            .is_some_and(|existing| stream_is_expired(existing, input.now_ms))
1574        {
1575            self.remove_stream_state(&input.stream_id);
1576        }
1577
1578        if let Some(existing) = self.streams.get(&input.stream_id) {
1579            if is_soft_deleted(existing) {
1580                return StreamResponse::error(
1581                    StreamErrorCode::StreamAlreadyExistsConflict,
1582                    format!(
1583                        "stream '{}' is gone and cannot be recreated yet",
1584                        input.stream_id
1585                    ),
1586                );
1587            }
1588            if existing.content_type == input.content_type
1589                && existing.status == status_from_closed(input.close_after)
1590                && existing.stream_ttl_seconds == input.stream_ttl_seconds
1591                && existing.stream_expires_at_ms == input.stream_expires_at_ms
1592                && existing.forked_from == input.forked_from
1593                && existing.fork_offset == input.fork_offset
1594            {
1595                return StreamResponse::AlreadyExists {
1596                    next_offset: existing.tail_offset,
1597                    closed: existing.status == StreamStatus::Closed,
1598                    content_type: existing.content_type.clone(),
1599                    stream_ttl_seconds: existing.stream_ttl_seconds,
1600                    stream_expires_at_ms: existing.stream_expires_at_ms,
1601                };
1602            }
1603            return StreamResponse::error(
1604                StreamErrorCode::StreamAlreadyExistsConflict,
1605                format!(
1606                    "stream '{}' already exists with different metadata",
1607                    input.stream_id
1608                ),
1609            );
1610        }
1611
1612        let initial_len = input.initial_payload.payload_len;
1613        let metadata = StreamMetadata {
1614            stream_id: input.stream_id.clone(),
1615            content_type: input.content_type,
1616            status: status_from_closed(input.close_after),
1617            tail_offset: initial_len,
1618            last_stream_seq: input.stream_seq,
1619            stream_ttl_seconds: input.stream_ttl_seconds,
1620            stream_expires_at_ms: input.stream_expires_at_ms,
1621            created_at_ms: input.now_ms,
1622            last_ttl_touch_at_ms: input.now_ms,
1623            forked_from: input.forked_from,
1624            fork_offset: input.fork_offset,
1625            fork_ref_count: 0,
1626        };
1627        self.streams.insert(input.stream_id.clone(), metadata);
1628        self.hot_buffers
1629            .insert(input.stream_id.clone(), HotBuffer::default());
1630        let object = ObjectPayloadRef {
1631            start_offset: 0,
1632            end_offset: initial_len,
1633            s3_path: input.initial_payload.s3_path,
1634            object_size: input.initial_payload.object_size,
1635        };
1636        self.cold_index_mut()
1637            .push_external_segment(input.stream_id.clone(), object.clone());
1638        let mut integrity = StreamIntegrity::default();
1639        if initial_len > 0 {
1640            integrity.append_external(
1641                &input.stream_id,
1642                object.start_offset,
1643                object.end_offset,
1644                &object.s3_path,
1645                object.object_size,
1646            );
1647        }
1648        self.integrities.insert(input.stream_id.clone(), integrity);
1649        self.message_records
1650            .insert(input.stream_id.clone(), vec![StreamMessageRecord {
1651                start_offset: 0,
1652                end_offset: initial_len,
1653            }]);
1654        let mut producer_states = HashMap::new();
1655        if let Some(producer) = input.producer {
1656            let last_item = ProducerAppendRecord {
1657                start_offset: 0,
1658                next_offset: initial_len,
1659                closed: input.close_after,
1660            };
1661            producer_states.insert(producer.producer_id, ProducerState {
1662                producer_epoch: producer.producer_epoch,
1663                producer_seq: producer.producer_seq,
1664                last_start_offset: last_item.start_offset,
1665                last_next_offset: last_item.next_offset,
1666                last_closed: last_item.closed,
1667                last_items: vec![last_item],
1668            });
1669        }
1670        self.producers
1671            .insert(input.stream_id.clone(), producer_states);
1672        StreamResponse::Created {
1673            stream_id: input.stream_id,
1674            next_offset: initial_len,
1675            closed: input.close_after,
1676        }
1677    }
1678
1679    pub fn append_borrowed(&mut self, input: AppendStreamInput<'_>) -> StreamResponse {
1680        let AppendStreamInput {
1681            stream_id,
1682            content_type,
1683            payload,
1684            close_after,
1685            stream_seq,
1686            producer,
1687            now_ms,
1688        } = input;
1689        if let Err(response) = self.validate_stream_scope(&stream_id) {
1690            return response;
1691        }
1692        if let Err(response) = validate_producer_request(producer.as_ref()) {
1693            return response;
1694        }
1695
1696        let Some(_) = self.streams.get(&stream_id) else {
1697            return StreamResponse::error(
1698                StreamErrorCode::StreamNotFound,
1699                format!("stream '{stream_id}' does not exist"),
1700            );
1701        };
1702        if self.expire_stream_if_due(&stream_id, now_ms) {
1703            return StreamResponse::error(
1704                StreamErrorCode::StreamNotFound,
1705                format!("stream '{stream_id}' does not exist"),
1706            );
1707        }
1708        if self.streams.get(&stream_id).is_some_and(is_soft_deleted) {
1709            return StreamResponse::error(
1710                StreamErrorCode::StreamGone,
1711                format!("stream '{stream_id}' is gone"),
1712            );
1713        }
1714        let producer_decision = match self.evaluate_producer(&stream_id, producer.as_ref()) {
1715            Ok(decision) => decision,
1716            Err(response) => return response,
1717        };
1718        if let ProducerDecision::Duplicate {
1719            offset,
1720            next_offset,
1721            closed,
1722            producer,
1723            ..
1724        } = producer_decision
1725        {
1726            if payload.is_empty() {
1727                return StreamResponse::Closed {
1728                    next_offset,
1729                    deduplicated: true,
1730                    producer: Some(producer),
1731                };
1732            }
1733            return StreamResponse::Appended {
1734                offset,
1735                next_offset,
1736                closed,
1737                deduplicated: true,
1738                producer: Some(producer),
1739            };
1740        }
1741
1742        let Some(stream) = self.streams.get_mut(&stream_id) else {
1743            unreachable!("stream existence checked before producer evaluation");
1744        };
1745
1746        if stream.status == StreamStatus::Closed {
1747            if close_after && payload.is_empty() {
1748                return StreamResponse::Closed {
1749                    next_offset: stream.tail_offset,
1750                    deduplicated: false,
1751                    producer: None,
1752                };
1753            }
1754            return StreamResponse::error_with_next_offset_and_context(
1755                StreamErrorCode::StreamClosed,
1756                format!("stream '{stream_id}' is closed"),
1757                stream.tail_offset,
1758                vec![StreamErrorContext::StreamClosed],
1759            );
1760        }
1761
1762        if payload.is_empty() && !close_after {
1763            return StreamResponse::error(
1764                StreamErrorCode::EmptyAppend,
1765                "append payload must be non-empty unless closing the stream",
1766            );
1767        }
1768
1769        if !payload.is_empty() {
1770            let Some(content_type) = content_type else {
1771                return StreamResponse::error(
1772                    StreamErrorCode::MissingContentType,
1773                    "append with a body must include content type",
1774                );
1775            };
1776            if content_type != stream.content_type {
1777                return StreamResponse::error_with_next_offset(
1778                    StreamErrorCode::ContentTypeMismatch,
1779                    format!(
1780                        "append content type '{content_type}' does not match stream content type '{}'",
1781                        stream.content_type
1782                    ),
1783                    stream.tail_offset,
1784                );
1785            }
1786        }
1787
1788        if let Err(response) = check_stream_seq(stream, stream_seq.as_deref()) {
1789            return response;
1790        }
1791
1792        let offset = stream.tail_offset;
1793        let payload_len = u64::try_from(payload.len()).expect("payload len fits u64");
1794        stream.tail_offset = stream.tail_offset.saturating_add(payload_len);
1795        if let Some(seq) = stream_seq {
1796            stream.last_stream_seq = Some(seq);
1797        }
1798        renew_stream_ttl(stream, now_ms);
1799        if close_after {
1800            stream.status = StreamStatus::Closed;
1801        }
1802        let closed = stream.status == StreamStatus::Closed;
1803        let next_offset = stream.tail_offset;
1804        let producer_ack = producer.clone();
1805        if let Some(producer) = producer {
1806            self.record_producer_success(
1807                stream_id.clone(),
1808                producer,
1809                ProducerAppendRecord {
1810                    start_offset: offset,
1811                    next_offset,
1812                    closed,
1813                },
1814                vec![ProducerAppendRecord {
1815                    start_offset: offset,
1816                    next_offset,
1817                    closed,
1818                }],
1819            );
1820        }
1821
1822        if payload.is_empty() {
1823            StreamResponse::Closed {
1824                next_offset,
1825                deduplicated: false,
1826                producer: producer_ack,
1827            }
1828        } else {
1829            self.hot_buffers
1830                .get_mut(&stream_id)
1831                .expect("hot buffer exists for stream metadata")
1832                .push(offset, next_offset, payload);
1833            self.integrities
1834                .get_mut(&stream_id)
1835                .expect("integrity exists for stream metadata")
1836                .append_payload(&stream_id, offset, next_offset, payload);
1837            self.message_records
1838                .entry(stream_id.clone())
1839                .or_default()
1840                .push(StreamMessageRecord {
1841                    start_offset: offset,
1842                    end_offset: next_offset,
1843                });
1844            StreamResponse::Appended {
1845                offset,
1846                next_offset,
1847                closed: close_after,
1848                deduplicated: false,
1849                producer: producer_ack,
1850            }
1851        }
1852    }
1853
1854    fn append_external(&mut self, input: AppendExternalInput<'_>) -> StreamResponse {
1855        let AppendExternalInput {
1856            stream_id,
1857            content_type,
1858            payload,
1859            close_after,
1860            stream_seq,
1861            producer,
1862            now_ms,
1863        } = input;
1864        if let Err(response) = validate_external_payload_ref(&payload) {
1865            return response;
1866        }
1867        if let Err(response) = self.validate_stream_scope(&stream_id) {
1868            return response;
1869        }
1870        if let Err(response) = validate_producer_request(producer.as_ref()) {
1871            return response;
1872        }
1873        let Some(_) = self.streams.get(&stream_id) else {
1874            return StreamResponse::error(
1875                StreamErrorCode::StreamNotFound,
1876                format!("stream '{stream_id}' does not exist"),
1877            );
1878        };
1879        if self.expire_stream_if_due(&stream_id, now_ms) {
1880            return StreamResponse::error(
1881                StreamErrorCode::StreamNotFound,
1882                format!("stream '{stream_id}' does not exist"),
1883            );
1884        }
1885        if self.streams.get(&stream_id).is_some_and(is_soft_deleted) {
1886            return StreamResponse::error(
1887                StreamErrorCode::StreamGone,
1888                format!("stream '{stream_id}' is gone"),
1889            );
1890        }
1891        let producer_decision = match self.evaluate_producer(&stream_id, producer.as_ref()) {
1892            Ok(decision) => decision,
1893            Err(response) => return response,
1894        };
1895        if let ProducerDecision::Duplicate {
1896            offset,
1897            next_offset,
1898            closed,
1899            producer,
1900            ..
1901        } = producer_decision
1902        {
1903            return StreamResponse::Appended {
1904                offset,
1905                next_offset,
1906                closed,
1907                deduplicated: true,
1908                producer: Some(producer),
1909            };
1910        }
1911
1912        let Some(stream) = self.streams.get(&stream_id) else {
1913            unreachable!("stream existence checked before producer evaluation");
1914        };
1915        if stream.status == StreamStatus::Closed {
1916            return StreamResponse::error_with_next_offset_and_context(
1917                StreamErrorCode::StreamClosed,
1918                format!("stream '{stream_id}' is closed"),
1919                stream.tail_offset,
1920                vec![StreamErrorContext::StreamClosed],
1921            );
1922        }
1923        let Some(content_type) = content_type else {
1924            return StreamResponse::error(
1925                StreamErrorCode::MissingContentType,
1926                "append with a body must include content type",
1927            );
1928        };
1929        if content_type != stream.content_type {
1930            return StreamResponse::error_with_next_offset(
1931                StreamErrorCode::ContentTypeMismatch,
1932                format!(
1933                    "append content type '{content_type}' does not match stream content type '{}'",
1934                    stream.content_type
1935                ),
1936                stream.tail_offset,
1937            );
1938        }
1939        if let Err(response) = check_stream_seq(stream, stream_seq.as_deref()) {
1940            return response;
1941        }
1942        let offset = stream.tail_offset;
1943        let next_offset = offset.saturating_add(payload.payload_len);
1944        let stream = self
1945            .streams
1946            .get_mut(&stream_id)
1947            .expect("stream existence checked before external append mutation");
1948        stream.tail_offset = next_offset;
1949        if let Some(seq) = stream_seq {
1950            stream.last_stream_seq = Some(seq);
1951        }
1952        renew_stream_ttl(stream, now_ms);
1953        if close_after {
1954            stream.status = StreamStatus::Closed;
1955        }
1956        let closed = stream.status == StreamStatus::Closed;
1957        let producer_ack = producer.clone();
1958        if let Some(producer) = producer {
1959            self.record_producer_success(
1960                stream_id.clone(),
1961                producer,
1962                ProducerAppendRecord {
1963                    start_offset: offset,
1964                    next_offset,
1965                    closed,
1966                },
1967                vec![ProducerAppendRecord {
1968                    start_offset: offset,
1969                    next_offset,
1970                    closed,
1971                }],
1972            );
1973        }
1974        let object = ObjectPayloadRef {
1975            start_offset: offset,
1976            end_offset: next_offset,
1977            s3_path: payload.s3_path,
1978            object_size: payload.object_size,
1979        };
1980        self.cold_index_mut()
1981            .push_external_segment(stream_id.clone(), object.clone());
1982        self.integrities
1983            .get_mut(&stream_id)
1984            .expect("integrity exists for stream metadata")
1985            .append_external(
1986                &stream_id,
1987                object.start_offset,
1988                object.end_offset,
1989                &object.s3_path,
1990                object.object_size,
1991            );
1992        self.message_records
1993            .entry(stream_id.clone())
1994            .or_default()
1995            .push(StreamMessageRecord {
1996                start_offset: offset,
1997                end_offset: next_offset,
1998            });
1999        StreamResponse::Appended {
2000            offset,
2001            next_offset,
2002            closed: close_after,
2003            deduplicated: false,
2004            producer: producer_ack,
2005        }
2006    }
2007
2008    pub fn append_batch_borrowed(
2009        &mut self,
2010        stream_id: BucketStreamId,
2011        content_type: Option<&str>,
2012        payloads: &[&[u8]],
2013        producer: Option<ProducerRequest>,
2014        now_ms: u64,
2015    ) -> Result<StreamBatchAppend, StreamResponse> {
2016        if payloads.is_empty() {
2017            return Err(StreamResponse::error(
2018                StreamErrorCode::EmptyAppend,
2019                "append batch must contain at least one payload",
2020            ));
2021        }
2022        self.validate_stream_scope(&stream_id)?;
2023        validate_producer_request(producer.as_ref())?;
2024        if self.expire_stream_if_due(&stream_id, now_ms) {
2025            return Err(StreamResponse::error(
2026                StreamErrorCode::StreamNotFound,
2027                format!("stream '{stream_id}' does not exist"),
2028            ));
2029        }
2030        if self.streams.get(&stream_id).is_some_and(is_soft_deleted) {
2031            return Err(StreamResponse::error(
2032                StreamErrorCode::StreamGone,
2033                format!("stream '{stream_id}' is gone"),
2034            ));
2035        }
2036        let producer_decision = self.evaluate_producer(&stream_id, producer.as_ref())?;
2037        if let ProducerDecision::Duplicate { items, .. } = producer_decision {
2038            return Ok(StreamBatchAppend {
2039                items: items
2040                    .into_iter()
2041                    .map(|item| StreamBatchAppendItem {
2042                        offset: item.start_offset,
2043                        next_offset: item.next_offset,
2044                        closed: item.closed,
2045                        deduplicated: true,
2046                    })
2047                    .collect(),
2048                deduplicated: true,
2049            });
2050        }
2051
2052        let Some(stream) = self.streams.get_mut(&stream_id) else {
2053            return Err(StreamResponse::error(
2054                StreamErrorCode::StreamNotFound,
2055                format!("stream '{stream_id}' does not exist"),
2056            ));
2057        };
2058        if stream.status == StreamStatus::Closed {
2059            return Err(StreamResponse::error_with_next_offset_and_context(
2060                StreamErrorCode::StreamClosed,
2061                format!("stream '{stream_id}' is closed"),
2062                stream.tail_offset,
2063                vec![StreamErrorContext::StreamClosed],
2064            ));
2065        }
2066        let Some(content_type) = content_type else {
2067            return Err(StreamResponse::error(
2068                StreamErrorCode::MissingContentType,
2069                "append batch must include content type",
2070            ));
2071        };
2072        if content_type != stream.content_type {
2073            return Err(StreamResponse::error_with_next_offset(
2074                StreamErrorCode::ContentTypeMismatch,
2075                format!(
2076                    "append content type '{content_type}' does not match stream content type '{}'",
2077                    stream.content_type
2078                ),
2079                stream.tail_offset,
2080            ));
2081        }
2082        if payloads.iter().any(|payload| payload.is_empty()) {
2083            return Err(StreamResponse::error(
2084                StreamErrorCode::EmptyAppend,
2085                "append batch payloads must be non-empty",
2086            ));
2087        }
2088
2089        let mut items = Vec::with_capacity(payloads.len());
2090        for payload in payloads {
2091            let offset = stream.tail_offset;
2092            let payload_len = u64::try_from(payload.len()).expect("payload len fits u64");
2093            stream.tail_offset = stream.tail_offset.saturating_add(payload_len);
2094            items.push(ProducerAppendRecord {
2095                start_offset: offset,
2096                next_offset: stream.tail_offset,
2097                closed: false,
2098            });
2099        }
2100        let last = items
2101            .last()
2102            .expect("payloads checked non-empty before append")
2103            .clone();
2104        renew_stream_ttl(stream, now_ms);
2105        if let Some(producer) = producer {
2106            self.record_producer_success(stream_id.clone(), producer, last.clone(), items.clone());
2107        }
2108        let hot_buffer = self
2109            .hot_buffers
2110            .get_mut(&stream_id)
2111            .expect("hot buffer exists for stream metadata");
2112        for (item, payload) in items.iter().zip(payloads.iter()) {
2113            hot_buffer.push(item.start_offset, item.next_offset, payload);
2114        }
2115        let integrity = self
2116            .integrities
2117            .get_mut(&stream_id)
2118            .expect("integrity exists for stream metadata");
2119        for (item, payload) in items.iter().zip(payloads.iter()) {
2120            integrity.append_payload(&stream_id, item.start_offset, item.next_offset, payload);
2121        }
2122        self.message_records
2123            .entry(stream_id.clone())
2124            .or_default()
2125            .extend(items.iter().map(|item| StreamMessageRecord {
2126                start_offset: item.start_offset,
2127                end_offset: item.next_offset,
2128            }));
2129        Ok(StreamBatchAppend {
2130            items: items
2131                .into_iter()
2132                .map(|item| StreamBatchAppendItem {
2133                    offset: item.start_offset,
2134                    next_offset: item.next_offset,
2135                    closed: item.closed,
2136                    deduplicated: false,
2137                })
2138                .collect(),
2139            deduplicated: false,
2140        })
2141    }
2142
2143    fn close(
2144        &mut self,
2145        stream_id: BucketStreamId,
2146        stream_seq: Option<String>,
2147        producer: Option<ProducerRequest>,
2148        now_ms: u64,
2149    ) -> StreamResponse {
2150        self.append_borrowed(AppendStreamInput {
2151            stream_id,
2152            content_type: None,
2153            payload: &[],
2154            close_after: true,
2155            stream_seq,
2156            producer,
2157            now_ms,
2158        })
2159    }
2160
2161    fn delete_stream(&mut self, stream_id: &BucketStreamId) -> StreamResponse {
2162        if let Err(response) = self.validate_stream_scope(stream_id) {
2163            return response;
2164        }
2165        let Some(stream) = self.streams.get_mut(stream_id) else {
2166            return StreamResponse::error(
2167                StreamErrorCode::StreamNotFound,
2168                format!("stream '{stream_id}' does not exist"),
2169            );
2170        };
2171        if is_soft_deleted(stream) {
2172            return StreamResponse::error(
2173                StreamErrorCode::StreamGone,
2174                format!("stream '{stream_id}' is gone"),
2175            );
2176        }
2177        if stream.fork_ref_count > 0 {
2178            stream.status = StreamStatus::SoftDeleted;
2179            return StreamResponse::Deleted {
2180                hard_deleted: false,
2181                parent_to_release: None,
2182            };
2183        }
2184        let parent_to_release = stream.forked_from.clone();
2185        self.remove_stream_state(stream_id);
2186        StreamResponse::Deleted {
2187            hard_deleted: true,
2188            parent_to_release,
2189        }
2190    }
2191
2192    fn add_fork_ref(&mut self, stream_id: &BucketStreamId, now_ms: u64) -> StreamResponse {
2193        if let Err(response) = self.validate_stream_scope(stream_id) {
2194            return response;
2195        }
2196        if self.expire_stream_if_due(stream_id, now_ms) {
2197            return StreamResponse::error(
2198                StreamErrorCode::StreamNotFound,
2199                format!("stream '{stream_id}' does not exist"),
2200            );
2201        }
2202        let Some(stream) = self.streams.get_mut(stream_id) else {
2203            return StreamResponse::error(
2204                StreamErrorCode::StreamNotFound,
2205                format!("stream '{stream_id}' does not exist"),
2206            );
2207        };
2208        if is_soft_deleted(stream) {
2209            return StreamResponse::error(
2210                StreamErrorCode::StreamGone,
2211                format!("stream '{stream_id}' is gone"),
2212            );
2213        }
2214        stream.fork_ref_count = stream.fork_ref_count.saturating_add(1);
2215        StreamResponse::ForkRefAdded {
2216            fork_ref_count: stream.fork_ref_count,
2217        }
2218    }
2219
2220    fn release_fork_ref(&mut self, stream_id: &BucketStreamId) -> StreamResponse {
2221        if let Err(response) = self.validate_stream_scope(stream_id) {
2222            return response;
2223        }
2224        let Some(stream) = self.streams.get_mut(stream_id) else {
2225            return StreamResponse::ForkRefReleased {
2226                hard_deleted: false,
2227                fork_ref_count: 0,
2228                parent_to_release: None,
2229            };
2230        };
2231        if stream.fork_ref_count == 0 {
2232            return StreamResponse::error(
2233                StreamErrorCode::InvalidFork,
2234                format!("stream '{stream_id}' has no fork reference to release"),
2235            );
2236        }
2237        stream.fork_ref_count -= 1;
2238        if stream.fork_ref_count == 0 && is_soft_deleted(stream) {
2239            let parent_to_release = stream.forked_from.clone();
2240            self.remove_stream_state(stream_id);
2241            return StreamResponse::ForkRefReleased {
2242                hard_deleted: true,
2243                fork_ref_count: 0,
2244                parent_to_release,
2245            };
2246        }
2247        StreamResponse::ForkRefReleased {
2248            hard_deleted: false,
2249            fork_ref_count: stream.fork_ref_count,
2250            parent_to_release: None,
2251        }
2252    }
2253
2254    fn touch_stream_access(
2255        &mut self,
2256        stream_id: &BucketStreamId,
2257        now_ms: u64,
2258        renew_ttl: bool,
2259    ) -> StreamResponse {
2260        if let Err(response) = self.validate_stream_scope(stream_id) {
2261            return response;
2262        }
2263        let Some(stream) = self.streams.get(stream_id) else {
2264            return StreamResponse::error(
2265                StreamErrorCode::StreamNotFound,
2266                format!("stream '{stream_id}' does not exist"),
2267            );
2268        };
2269        if is_soft_deleted(stream) {
2270            return StreamResponse::error(
2271                StreamErrorCode::StreamGone,
2272                format!("stream '{stream_id}' is gone"),
2273            );
2274        }
2275        if stream_is_expired(stream, now_ms) {
2276            self.remove_stream_state(stream_id);
2277            return StreamResponse::Accessed {
2278                changed: true,
2279                expired: true,
2280            };
2281        }
2282        let changed = if renew_ttl && stream.stream_ttl_seconds.is_some() {
2283            let stream = self
2284                .streams
2285                .get_mut(stream_id)
2286                .expect("stream existence checked before TTL renewal");
2287            let previous = stream.last_ttl_touch_at_ms;
2288            renew_stream_ttl(stream, now_ms);
2289            stream.last_ttl_touch_at_ms != previous
2290        } else {
2291            false
2292        };
2293        StreamResponse::Accessed {
2294            changed,
2295            expired: false,
2296        }
2297    }
2298
2299    fn expire_stream_if_due(&mut self, stream_id: &BucketStreamId, now_ms: u64) -> bool {
2300        if self
2301            .streams
2302            .get(stream_id)
2303            .is_some_and(|stream| stream_is_expired(stream, now_ms))
2304        {
2305            self.remove_stream_state(stream_id);
2306            return true;
2307        }
2308        false
2309    }
2310
2311    fn remove_stream_state(&mut self, stream_id: &BucketStreamId) -> bool {
2312        if self.streams.remove(stream_id).is_some() {
2313            self.hot_buffers.remove(stream_id);
2314            let had_cold = self.cold_index_mut().remove_stream(stream_id);
2315            self.message_records.remove(stream_id);
2316            self.integrities.remove(stream_id);
2317            self.visible_snapshots.remove(stream_id);
2318            self.producers.remove(stream_id);
2319            // The cold objects we wrote for this stream are now unreferenced.
2320            // Enqueue the whole prefix for the background GC worker to reclaim;
2321            // cold objects are stream-exclusive (forks copy, never share), so a
2322            // prefix sweep is safe and keeps the queue O(streams) not O(chunks).
2323            if had_cold {
2324                self.enqueue_cold_gc(ColdGcTarget::Stream(stream_id.clone()));
2325            }
2326            true
2327        } else {
2328            false
2329        }
2330    }
2331
2332    fn enqueue_cold_gc(&mut self, target: ColdGcTarget) {
2333        let seq = self.next_cold_gc_seq;
2334        self.next_cold_gc_seq = self.next_cold_gc_seq.saturating_add(1);
2335        self.pending_cold_gc.push_back(ColdGcEntry { seq, target });
2336    }
2337
2338    fn ack_cold_gc(&mut self, up_to_seq: u64) -> StreamResponse {
2339        let before = self.pending_cold_gc.len();
2340        while self
2341            .pending_cold_gc
2342            .front()
2343            .is_some_and(|entry| entry.seq <= up_to_seq)
2344        {
2345            self.pending_cold_gc.pop_front();
2346        }
2347        let removed = u64::try_from(before - self.pending_cold_gc.len()).expect("removed fits u64");
2348        StreamResponse::ColdGcAcked { removed }
2349    }
2350
2351    /// A bounded snapshot of the front of the GC queue for the leader's worker
2352    /// to reclaim. Read-only; draining is confirmed by a replicated `AckColdGc`.
2353    pub fn pending_cold_gc_batch(&self, max: usize) -> Vec<ColdGcEntry> {
2354        self.pending_cold_gc.iter().take(max).cloned().collect()
2355    }
2356
2357    pub fn pending_cold_gc_len(&self) -> usize {
2358        self.pending_cold_gc.len()
2359    }
2360
2361    fn validate_stream_scope(&self, stream_id: &BucketStreamId) -> Result<(), StreamResponse> {
2362        if let Err(message) = validate_bucket_id(&stream_id.bucket_id) {
2363            return Err(StreamResponse::error(
2364                StreamErrorCode::InvalidBucketId,
2365                message,
2366            ));
2367        }
2368        if let Err(message) = validate_stream_id(stream_id) {
2369            return Err(StreamResponse::error(
2370                StreamErrorCode::InvalidStreamId,
2371                message,
2372            ));
2373        }
2374        if !self.buckets.contains(&stream_id.bucket_id) {
2375            return Err(StreamResponse::error(
2376                StreamErrorCode::BucketNotFound,
2377                format!("bucket '{}' does not exist", stream_id.bucket_id),
2378            ));
2379        }
2380        Ok(())
2381    }
2382
2383    fn earliest_retained_offset(&self, stream_id: &BucketStreamId) -> u64 {
2384        self.visible_snapshots
2385            .get(stream_id)
2386            .map(|snapshot| snapshot.offset)
2387            .unwrap_or(0)
2388    }
2389
2390    fn snapshot_offset_aligned(
2391        &self,
2392        stream_id: &BucketStreamId,
2393        snapshot_offset: u64,
2394        retained_offset: u64,
2395    ) -> bool {
2396        snapshot_offset == retained_offset
2397            || snapshot_offset <= self.cold_frontier_offset(stream_id, retained_offset)
2398            || self
2399                .hot_buffers
2400                .get(stream_id)
2401                .is_some_and(|buffer| snapshot_offset <= buffer.hot_start_offset())
2402            || self.message_records.get(stream_id).is_some_and(|records| {
2403                records
2404                    .iter()
2405                    .any(|record| record.end_offset == snapshot_offset)
2406            })
2407    }
2408
2409    fn compact_retained_prefix(&mut self, stream_id: &BucketStreamId, retained_offset: u64) {
2410        let frontier = self.cold_frontier_offset(stream_id, retained_offset).max(
2411            self.hot_buffers
2412                .get(stream_id)
2413                .map(|buffer| buffer.hot_start_offset())
2414                .unwrap_or(retained_offset),
2415        );
2416        self.compact_message_records_before(stream_id, retained_offset, frontier);
2417        if let Some(integrity) = self.integrities.get_mut(stream_id) {
2418            integrity.evict_before(retained_offset);
2419        }
2420        let dropped_cold_paths = self
2421            .cold_index_mut()
2422            .compact_before(stream_id, retained_offset);
2423        if !dropped_cold_paths.is_empty() {
2424            self.enqueue_cold_gc(ColdGcTarget::Paths(dropped_cold_paths));
2425        }
2426
2427        if let Some(hot_buffer) = self.hot_buffers.get_mut(stream_id) {
2428            hot_buffer.discard_before(retained_offset);
2429        }
2430    }
2431
2432    fn compact_message_records_before(
2433        &mut self,
2434        stream_id: &BucketStreamId,
2435        retained_offset: u64,
2436        frontier: u64,
2437    ) {
2438        let Some(records) = self.message_records.remove(stream_id) else {
2439            return;
2440        };
2441        let frontier = frontier.max(retained_offset);
2442        let mut compacted = Vec::with_capacity(records.len());
2443        if frontier > retained_offset {
2444            compacted.push(StreamMessageRecord {
2445                start_offset: retained_offset,
2446                end_offset: frontier,
2447            });
2448        }
2449        compacted.extend(records.iter().filter_map(|record| {
2450            if record.end_offset <= frontier {
2451                return None;
2452            }
2453            let start_offset = record.start_offset.max(frontier).max(retained_offset);
2454            (record.end_offset > start_offset).then_some(StreamMessageRecord {
2455                start_offset,
2456                end_offset: record.end_offset,
2457            })
2458        }));
2459        if compacted.is_empty() {
2460            return;
2461        }
2462        self.message_records.insert(stream_id.clone(), compacted);
2463    }
2464
2465    fn cold_frontier_offset(&self, stream_id: &BucketStreamId, retained_offset: u64) -> u64 {
2466        self.cold_index()
2467            .cold_frontier_offset(stream_id, retained_offset)
2468    }
2469
2470    fn producer_snapshot(&self, stream_id: &BucketStreamId) -> Vec<ProducerSnapshot> {
2471        let mut producer_states = self
2472            .producers
2473            .get(stream_id)
2474            .into_iter()
2475            .flat_map(|states| states.iter())
2476            .map(|(producer_id, state)| ProducerSnapshot {
2477                producer_id: producer_id.clone(),
2478                producer_epoch: state.producer_epoch,
2479                producer_seq: state.producer_seq,
2480                last_start_offset: state.last_start_offset,
2481                last_next_offset: state.last_next_offset,
2482                last_closed: state.last_closed,
2483                last_items: state.last_items.clone(),
2484            })
2485            .collect::<Vec<_>>();
2486        producer_states.sort_by(|left, right| left.producer_id.cmp(&right.producer_id));
2487        producer_states
2488    }
2489
2490    fn evaluate_producer(
2491        &self,
2492        stream_id: &BucketStreamId,
2493        producer: Option<&ProducerRequest>,
2494    ) -> Result<ProducerDecision, StreamResponse> {
2495        let Some(producer) = producer else {
2496            return Ok(ProducerDecision::Accept);
2497        };
2498        let Some(states) = self.producers.get(stream_id) else {
2499            return Ok(ProducerDecision::Accept);
2500        };
2501        let Some(state) = states.get(&producer.producer_id) else {
2502            if producer.producer_seq == 0 {
2503                return Ok(ProducerDecision::Accept);
2504            }
2505            return Err(StreamResponse::error_with_context(
2506                StreamErrorCode::ProducerSeqConflict,
2507                format!(
2508                    "producer '{}' expected sequence 0, received {}",
2509                    producer.producer_id, producer.producer_seq
2510                ),
2511                vec![StreamErrorContext::ProducerSeqConflict {
2512                    expected_seq: 0,
2513                    received_seq: producer.producer_seq,
2514                }],
2515            ));
2516        };
2517
2518        if producer.producer_epoch < state.producer_epoch {
2519            return Err(StreamResponse::error_with_context(
2520                StreamErrorCode::ProducerEpochStale,
2521                format!(
2522                    "producer '{}' epoch {} is stale; current epoch is {}",
2523                    producer.producer_id, producer.producer_epoch, state.producer_epoch
2524                ),
2525                vec![StreamErrorContext::ProducerEpochStale {
2526                    current_epoch: state.producer_epoch,
2527                }],
2528            ));
2529        }
2530        if producer.producer_epoch > state.producer_epoch {
2531            if producer.producer_seq == 0 {
2532                return Ok(ProducerDecision::Accept);
2533            }
2534            return Err(StreamResponse::error(
2535                StreamErrorCode::InvalidProducer,
2536                format!(
2537                    "producer '{}' new epoch {} must start at sequence 0",
2538                    producer.producer_id, producer.producer_epoch
2539                ),
2540            ));
2541        }
2542
2543        if producer.producer_seq <= state.producer_seq {
2544            return Ok(ProducerDecision::Duplicate {
2545                offset: state.last_start_offset,
2546                next_offset: state.last_next_offset,
2547                closed: state.last_closed,
2548                producer: ProducerRequest {
2549                    producer_id: producer.producer_id.clone(),
2550                    producer_epoch: state.producer_epoch,
2551                    producer_seq: state.producer_seq,
2552                },
2553                items: state.last_items.clone(),
2554            });
2555        }
2556        if producer.producer_seq == state.producer_seq + 1 {
2557            return Ok(ProducerDecision::Accept);
2558        }
2559        Err(StreamResponse::error_with_context(
2560            StreamErrorCode::ProducerSeqConflict,
2561            format!(
2562                "producer '{}' expected sequence {}, received {}",
2563                producer.producer_id,
2564                state.producer_seq + 1,
2565                producer.producer_seq
2566            ),
2567            vec![StreamErrorContext::ProducerSeqConflict {
2568                expected_seq: state.producer_seq + 1,
2569                received_seq: producer.producer_seq,
2570            }],
2571        ))
2572    }
2573
2574    fn record_producer_success(
2575        &mut self,
2576        stream_id: BucketStreamId,
2577        producer: ProducerRequest,
2578        last: ProducerAppendRecord,
2579        last_items: Vec<ProducerAppendRecord>,
2580    ) {
2581        self.producers
2582            .entry(stream_id)
2583            .or_default()
2584            .insert(producer.producer_id, ProducerState {
2585                producer_epoch: producer.producer_epoch,
2586                producer_seq: producer.producer_seq,
2587                last_start_offset: last.start_offset,
2588                last_next_offset: last.next_offset,
2589                last_closed: last.closed,
2590                last_items,
2591            });
2592    }
2593}
2594
2595#[derive(Debug)]
2596struct CreateStreamInput {
2597    stream_id: BucketStreamId,
2598    content_type: String,
2599    initial_payload: Vec<u8>,
2600    close_after: bool,
2601    stream_seq: Option<String>,
2602    producer: Option<ProducerRequest>,
2603    stream_ttl_seconds: Option<u64>,
2604    stream_expires_at_ms: Option<u64>,
2605    forked_from: Option<BucketStreamId>,
2606    fork_offset: Option<u64>,
2607    now_ms: u64,
2608}
2609
2610#[derive(Debug)]
2611struct CreateExternalStreamInput {
2612    stream_id: BucketStreamId,
2613    content_type: String,
2614    initial_payload: ExternalPayloadRef,
2615    close_after: bool,
2616    stream_seq: Option<String>,
2617    producer: Option<ProducerRequest>,
2618    stream_ttl_seconds: Option<u64>,
2619    stream_expires_at_ms: Option<u64>,
2620    forked_from: Option<BucketStreamId>,
2621    fork_offset: Option<u64>,
2622    now_ms: u64,
2623}
2624
2625#[derive(Debug, Clone, PartialEq, Eq)]
2626enum ProducerDecision {
2627    Accept,
2628    Duplicate {
2629        offset: u64,
2630        next_offset: u64,
2631        closed: bool,
2632        producer: ProducerRequest,
2633        items: Vec<ProducerAppendRecord>,
2634    },
2635}
2636
2637impl CreateStreamInput {
2638    fn initial_len(&self) -> u64 {
2639        u64::try_from(self.initial_payload.len()).expect("payload len fits u64")
2640    }
2641}
2642
2643fn status_from_closed(closed: bool) -> StreamStatus {
2644    if closed {
2645        StreamStatus::Closed
2646    } else {
2647        StreamStatus::Open
2648    }
2649}
2650
2651fn is_soft_deleted(stream: &StreamMetadata) -> bool {
2652    stream.status == StreamStatus::SoftDeleted
2653}
2654
2655fn validate_retention(
2656    stream_ttl_seconds: Option<u64>,
2657    stream_expires_at_ms: Option<u64>,
2658) -> Result<(), StreamResponse> {
2659    if stream_ttl_seconds.is_some() && stream_expires_at_ms.is_some() {
2660        return Err(StreamResponse::error(
2661            StreamErrorCode::InvalidRetention,
2662            "stream ttl and expires-at cannot both be set",
2663        ));
2664    }
2665    if let Some(ttl_seconds) = stream_ttl_seconds
2666        && ttl_seconds.checked_mul(1000).is_none()
2667    {
2668        return Err(StreamResponse::error(
2669            StreamErrorCode::InvalidRetention,
2670            "stream ttl overflows millisecond range",
2671        ));
2672    }
2673    Ok(())
2674}
2675
2676fn stream_expiry_at_ms(stream: &StreamMetadata) -> Option<u64> {
2677    if let Some(expires_at_ms) = stream.stream_expires_at_ms {
2678        return Some(expires_at_ms);
2679    }
2680    stream.stream_ttl_seconds.map(|ttl_seconds| {
2681        stream
2682            .last_ttl_touch_at_ms
2683            .saturating_add(ttl_seconds.saturating_mul(1000))
2684    })
2685}
2686
2687fn stream_is_expired(stream: &StreamMetadata, now_ms: u64) -> bool {
2688    stream_expiry_at_ms(stream).is_some_and(|expires_at_ms| now_ms >= expires_at_ms)
2689}
2690
2691fn renew_stream_ttl(stream: &mut StreamMetadata, now_ms: u64) {
2692    if stream.stream_ttl_seconds.is_some() && stream.stream_expires_at_ms.is_none() {
2693        stream.last_ttl_touch_at_ms = now_ms;
2694    }
2695}
2696
2697fn check_stream_seq(stream: &StreamMetadata, incoming: Option<&str>) -> Result<(), StreamResponse> {
2698    let Some(incoming) = incoming else {
2699        return Ok(());
2700    };
2701    if let Some(last) = stream.last_stream_seq.as_deref()
2702        && incoming <= last
2703    {
2704        return Err(StreamResponse::error_with_next_offset(
2705            StreamErrorCode::StreamSeqConflict,
2706            format!("stream sequence '{incoming}' is not greater than last sequence '{last}'"),
2707            stream.tail_offset,
2708        ));
2709    }
2710    Ok(())
2711}
2712
2713fn validate_producer_request(producer: Option<&ProducerRequest>) -> Result<(), StreamResponse> {
2714    let Some(producer) = producer else {
2715        return Ok(());
2716    };
2717    if producer.producer_id.trim().is_empty() {
2718        return Err(StreamResponse::error(
2719            StreamErrorCode::InvalidProducer,
2720            "producer id must not be empty",
2721        ));
2722    }
2723    const MAX_JS_SAFE_INTEGER: u64 = 9_007_199_254_740_991;
2724    if producer.producer_epoch > MAX_JS_SAFE_INTEGER {
2725        return Err(StreamResponse::error(
2726            StreamErrorCode::InvalidProducer,
2727            format!(
2728                "producer epoch {} exceeds maximum {}",
2729                producer.producer_epoch, MAX_JS_SAFE_INTEGER
2730            ),
2731        ));
2732    }
2733    if producer.producer_seq > MAX_JS_SAFE_INTEGER {
2734        return Err(StreamResponse::error(
2735            StreamErrorCode::InvalidProducer,
2736            format!(
2737                "producer sequence {} exceeds maximum {}",
2738                producer.producer_seq, MAX_JS_SAFE_INTEGER
2739            ),
2740        ));
2741    }
2742    Ok(())
2743}
2744
2745fn validate_external_payload_ref(payload: &ExternalPayloadRef) -> Result<(), StreamResponse> {
2746    if payload.s3_path.trim().is_empty() {
2747        return Err(StreamResponse::error(
2748            StreamErrorCode::InvalidColdFlush,
2749            "external payload S3 path must not be empty",
2750        ));
2751    }
2752    if payload.payload_len == 0 {
2753        return Err(StreamResponse::error(
2754            StreamErrorCode::EmptyAppend,
2755            "external payload length must be greater than zero",
2756        ));
2757    }
2758    if payload.object_size < payload.payload_len {
2759        return Err(StreamResponse::error(
2760            StreamErrorCode::InvalidColdFlush,
2761            "external payload object size must cover payload length",
2762        ));
2763    }
2764    Ok(())
2765}
2766
2767fn restore_producer_states(
2768    stream_id: &BucketStreamId,
2769    snapshots: Vec<ProducerSnapshot>,
2770) -> Result<HashMap<String, ProducerState>, StreamSnapshotError> {
2771    let mut states = HashMap::with_capacity(snapshots.len());
2772    for snapshot in snapshots {
2773        if states
2774            .insert(snapshot.producer_id.clone(), ProducerState {
2775                producer_epoch: snapshot.producer_epoch,
2776                producer_seq: snapshot.producer_seq,
2777                last_start_offset: snapshot.last_start_offset,
2778                last_next_offset: snapshot.last_next_offset,
2779                last_closed: snapshot.last_closed,
2780                last_items: snapshot.last_items,
2781            })
2782            .is_some()
2783        {
2784            return Err(StreamSnapshotError::DuplicateProducer {
2785                stream_id: stream_id.clone(),
2786                producer_id: snapshot.producer_id,
2787            });
2788        }
2789    }
2790    Ok(states)
2791}
2792
2793fn valid_cold_chunk_ref(chunk: &ColdChunkRef) -> bool {
2794    chunk.end_offset > chunk.start_offset
2795        && !chunk.s3_path.trim().is_empty()
2796        && chunk.object_size >= chunk.end_offset - chunk.start_offset
2797}
2798
2799fn valid_object_payload_ref(object: &ObjectPayloadRef) -> bool {
2800    object.end_offset > object.start_offset
2801        && !object.s3_path.trim().is_empty()
2802        && object.object_size >= object.end_offset - object.start_offset
2803}
2804
2805fn hot_segments_match_payload(segments: &[HotPayloadSegment], payload_len: usize) -> bool {
2806    let mut expected_payload_start = 0;
2807    for segment in segments {
2808        if segment.end_offset <= segment.start_offset
2809            || segment.payload_start != expected_payload_start
2810            || segment.payload_end <= segment.payload_start
2811            || segment.payload_end > payload_len
2812        {
2813            return false;
2814        }
2815        let Ok(logical_len) = usize::try_from(segment.end_offset - segment.start_offset) else {
2816            return false;
2817        };
2818        if logical_len != segment.payload_end - segment.payload_start {
2819            return false;
2820        }
2821        expected_payload_start = segment.payload_end;
2822    }
2823    expected_payload_start == payload_len
2824}
2825
2826fn payload_sources_cover_retained_suffix(
2827    cold_frontier_offset: u64,
2828    cold_chunks: &[ColdChunkRef],
2829    external_segments: &[ObjectPayloadRef],
2830    hot_segments: &[HotPayloadSegment],
2831    retained_offset: u64,
2832    tail_offset: u64,
2833) -> bool {
2834    if tail_offset < retained_offset {
2835        return false;
2836    }
2837    let mut ranges =
2838        Vec::with_capacity(1 + cold_chunks.len() + external_segments.len() + hot_segments.len());
2839    if cold_frontier_offset > retained_offset {
2840        ranges.push((retained_offset, cold_frontier_offset));
2841    }
2842    for chunk in cold_chunks {
2843        if !valid_cold_chunk_ref(chunk) {
2844            return false;
2845        }
2846        ranges.push((chunk.start_offset, chunk.end_offset));
2847    }
2848    for object in external_segments {
2849        if !valid_object_payload_ref(object) {
2850            return false;
2851        }
2852        ranges.push((object.start_offset, object.end_offset));
2853    }
2854    for segment in hot_segments {
2855        if segment.end_offset <= segment.start_offset {
2856            return false;
2857        }
2858        ranges.push((segment.start_offset, segment.end_offset));
2859    }
2860    ranges.sort_unstable();
2861
2862    let mut expected_start = retained_offset;
2863    for (start_offset, end_offset) in ranges {
2864        if end_offset <= expected_start {
2865            continue;
2866        }
2867        if start_offset > expected_start {
2868            return false;
2869        }
2870        expected_start = end_offset;
2871        if expected_start >= tail_offset {
2872            return true;
2873        }
2874    }
2875    expected_start == tail_offset
2876}
2877
2878fn push_cold_index_segments(
2879    segments: &mut Vec<(u64, StreamReadSegment)>,
2880    _stream_id: &BucketStreamId,
2881    generation: u64,
2882    start_offset: u64,
2883    end_offset: u64,
2884) {
2885    let mut cursor = start_offset;
2886    while cursor < end_offset {
2887        let page_id = cursor / COLD_INDEX_PAGE_SPAN_BYTES;
2888        let page_end = page_id
2889            .saturating_add(1)
2890            .saturating_mul(COLD_INDEX_PAGE_SPAN_BYTES);
2891        let segment_end = end_offset.min(page_end);
2892        segments.push((
2893            cursor,
2894            StreamReadSegment::ColdIndex(StreamReadColdIndexSegment {
2895                generation,
2896                page_id,
2897                read_start_offset: cursor,
2898                len: usize::try_from(segment_end - cursor).expect("cold index read len fits usize"),
2899            }),
2900        ));
2901        cursor = segment_end;
2902    }
2903}
2904
2905fn segments_cover_range(
2906    segments: &[(u64, StreamReadSegment)],
2907    offset: u64,
2908    next_offset: u64,
2909) -> bool {
2910    if next_offset < offset {
2911        return false;
2912    }
2913    let mut expected_start = offset;
2914    for (segment_start, segment) in segments {
2915        let Some(segment_end) = read_segment_end(*segment_start, segment) else {
2916            return false;
2917        };
2918        if segment_end <= expected_start {
2919            continue;
2920        }
2921        if *segment_start > expected_start {
2922            return false;
2923        }
2924        expected_start = segment_end;
2925        if expected_start >= next_offset {
2926            return true;
2927        }
2928    }
2929    expected_start == next_offset
2930}
2931
2932fn read_segment_end(segment_start: u64, segment: &StreamReadSegment) -> Option<u64> {
2933    match segment {
2934        StreamReadSegment::Object(object) => {
2935            if object.len == 0
2936                || object.read_start_offset != segment_start
2937                || object.read_start_offset < object.object.start_offset
2938            {
2939                return None;
2940            }
2941            let len = u64::try_from(object.len).ok()?;
2942            let segment_end = object.read_start_offset.checked_add(len)?;
2943            if segment_end > object.object.end_offset {
2944                return None;
2945            }
2946            Some(segment_end)
2947        }
2948        StreamReadSegment::ColdIndex(index) => {
2949            if index.len == 0 || index.read_start_offset != segment_start {
2950                return None;
2951            }
2952            let len = u64::try_from(index.len).ok()?;
2953            segment_start.checked_add(len)
2954        }
2955        StreamReadSegment::Hot(payload) => {
2956            if payload.is_empty() {
2957                return None;
2958            }
2959            let len = u64::try_from(payload.len()).ok()?;
2960            segment_start.checked_add(len)
2961        }
2962    }
2963}
2964
2965fn message_records_cover_retained_suffix(
2966    records: &[StreamMessageRecord],
2967    retained_offset: u64,
2968    tail_offset: u64,
2969) -> bool {
2970    let mut expected_start = retained_offset;
2971    for record in records {
2972        if record.start_offset != expected_start || record.end_offset <= record.start_offset {
2973            return false;
2974        }
2975        expected_start = record.end_offset;
2976    }
2977    expected_start == tail_offset
2978}
2979
2980fn compare_stream_ids(left: &BucketStreamId, right: &BucketStreamId) -> std::cmp::Ordering {
2981    left.bucket_id
2982        .cmp(&right.bucket_id)
2983        .then_with(|| left.stream_id.cmp(&right.stream_id))
2984}
2985
2986#[cfg(test)]
2987mod tests;