Skip to main content

ursula_stream/state_machine/
cold.rs

1//! Cold-tier flush planning, GC queue, retention compaction, and snapshot publishing.
2
3use super::BucketStreamId;
4use super::ColdChunkRef;
5use super::ColdFlushCandidate;
6use super::ColdGcEntry;
7use super::ColdGcTarget;
8use super::HashMap;
9use super::StreamErrorCode;
10use super::StreamErrorContext;
11use super::StreamMessageRecord;
12use super::StreamResponse;
13use super::StreamStateMachine;
14use super::StreamVisibleSnapshot;
15use super::compare_stream_ids;
16use super::stream_is_expired;
17
18impl StreamStateMachine {
19    pub fn plan_cold_flush(
20        &self,
21        stream_id: &BucketStreamId,
22        min_hot_bytes: usize,
23        max_flush_bytes: usize,
24    ) -> Result<Option<ColdFlushCandidate>, StreamResponse> {
25        let start_offset = self.hot_start_offset(stream_id);
26        self.plan_cold_flush_with_start(stream_id, start_offset, min_hot_bytes, max_flush_bytes)
27    }
28
29    pub(super) fn plan_cold_flush_with_start(
30        &self,
31        stream_id: &BucketStreamId,
32        start_offset: u64,
33        min_hot_bytes: usize,
34        max_flush_bytes: usize,
35    ) -> Result<Option<ColdFlushCandidate>, StreamResponse> {
36        if max_flush_bytes == 0 {
37            return Ok(None);
38        }
39        let Some(slot) = self.stream_slot(stream_id) else {
40            return Err(StreamResponse::error(
41                StreamErrorCode::StreamNotFound,
42                format!("stream '{stream_id}' does not exist"),
43            ));
44        };
45        let Some((start_offset, end_offset, payload)) =
46            slot.hot_buffer
47                .plan_cold_flush_from(start_offset, min_hot_bytes, max_flush_bytes)
48        else {
49            return Ok(None);
50        };
51        let payload_digest = blake3::hash(&payload).to_hex().to_string();
52        Ok(Some(ColdFlushCandidate {
53            stream_id: stream_id.clone(),
54            start_offset,
55            end_offset,
56            payload,
57            payload_digest,
58        }))
59    }
60
61    pub(super) fn plan_next_cold_flush_from_start(
62        &self,
63        mut start_fn: impl FnMut(&BucketStreamId) -> u64,
64        min_hot_bytes: usize,
65        max_flush_bytes: usize,
66        group_hot_bytes: u64,
67    ) -> Result<Option<ColdFlushCandidate>, StreamResponse> {
68        if max_flush_bytes == 0 {
69            return Ok(None);
70        }
71        let mut stream_ids = self.registry.stream_ids().cloned().collect::<Vec<_>>();
72        stream_ids.sort_by(compare_stream_ids);
73        for stream_id in &stream_ids {
74            let start = start_fn(stream_id);
75            match self.plan_cold_flush_with_start(stream_id, start, min_hot_bytes, max_flush_bytes)
76            {
77                Ok(Some(candidate)) => return Ok(Some(candidate)),
78                Ok(None) => {}
79                Err(StreamResponse::Error {
80                    code: StreamErrorCode::StreamGone | StreamErrorCode::StreamNotFound,
81                    ..
82                }) => {}
83                Err(err) => return Err(err),
84            }
85        }
86        let group_min_hot_bytes = u64::try_from(min_hot_bytes).unwrap_or(u64::MAX);
87        if group_hot_bytes < group_min_hot_bytes {
88            return Ok(None);
89        }
90        for stream_id in stream_ids {
91            let start = start_fn(&stream_id);
92            match self.plan_cold_flush_with_start(&stream_id, start, 1, max_flush_bytes) {
93                Ok(Some(candidate)) => return Ok(Some(candidate)),
94                Ok(None) => {}
95                Err(StreamResponse::Error {
96                    code: StreamErrorCode::StreamGone | StreamErrorCode::StreamNotFound,
97                    ..
98                }) => {}
99                Err(err) => return Err(err),
100            }
101        }
102        Ok(None)
103    }
104
105    pub fn plan_next_cold_flush_batch(
106        &self,
107        min_hot_bytes: usize,
108        max_flush_bytes: usize,
109        max_batch_bytes: usize,
110        max_candidates: usize,
111    ) -> Result<Vec<ColdFlushCandidate>, StreamResponse> {
112        if max_candidates == 0 || max_flush_bytes == 0 || max_batch_bytes == 0 {
113            return Ok(Vec::new());
114        }
115        let mut planned_flush_offsets: HashMap<BucketStreamId, u64> = HashMap::new();
116        let mut candidates = Vec::with_capacity(max_candidates);
117        let mut batch_bytes = 0usize;
118        let initial_group_hot_bytes = self
119            .registry
120            .stream_ids()
121            .map(|stream_id| {
122                u64::try_from(
123                    self.stream_slot(stream_id)
124                        .map(|slot| {
125                            slot.hot_buffer
126                                .remaining_len_from(self.hot_start_offset(stream_id))
127                        })
128                        .unwrap_or(0),
129                )
130                .expect("len fits u64")
131            })
132            .sum::<u64>();
133        let drain_group =
134            initial_group_hot_bytes >= u64::try_from(min_hot_bytes).unwrap_or(u64::MAX);
135        while candidates.len() < max_candidates {
136            let start_for = |stream_id: &BucketStreamId| -> u64 {
137                planned_flush_offsets
138                    .get(stream_id)
139                    .copied()
140                    .unwrap_or_else(|| self.hot_start_offset(stream_id))
141            };
142            let group_hot_bytes: u64 = self
143                .registry
144                .stream_ids()
145                .map(|stream_id| {
146                    let start = start_for(stream_id);
147                    self.stream_slot(stream_id)
148                        .map(|slot| {
149                            u64::try_from(slot.hot_buffer.remaining_len_from(start))
150                                .expect("len fits u64")
151                        })
152                        .unwrap_or(0)
153                })
154                .sum();
155            let candidate = self.plan_next_cold_flush_from_start(
156                start_for,
157                if drain_group { 1 } else { min_hot_bytes },
158                max_flush_bytes,
159                group_hot_bytes,
160            )?;
161            let Some(candidate) = candidate else {
162                break;
163            };
164            let Some(next_batch_bytes) = batch_bytes.checked_add(candidate.payload.len()) else {
165                break;
166            };
167            if !candidates.is_empty() && next_batch_bytes > max_batch_bytes {
168                break;
169            }
170            planned_flush_offsets.insert(candidate.stream_id.clone(), candidate.end_offset);
171            batch_bytes = next_batch_bytes;
172            candidates.push(candidate);
173            if batch_bytes >= max_batch_bytes {
174                break;
175            }
176        }
177        Ok(candidates)
178    }
179
180    pub(super) fn publish_snapshot(
181        &mut self,
182        stream_id: BucketStreamId,
183        snapshot_offset: u64,
184        content_type: String,
185        payload: Vec<u8>,
186        expected_digest: Option<String>,
187        now_ms: u64,
188    ) -> StreamResponse {
189        if let Err(response) = self.validate_stream_scope(&stream_id) {
190            return response;
191        }
192        if content_type.trim().is_empty() {
193            return StreamResponse::error(
194                StreamErrorCode::InvalidSnapshot,
195                "snapshot content type must not be empty",
196            );
197        }
198        let Some(stream) = self.stream_metadata(&stream_id) else {
199            return StreamResponse::error(
200                StreamErrorCode::StreamNotFound,
201                format!("stream '{stream_id}' does not exist"),
202            );
203        };
204        if stream_is_expired(stream, now_ms) {
205            self.remove_stream_state(&stream_id);
206            return StreamResponse::error(
207                StreamErrorCode::StreamNotFound,
208                format!("stream '{stream_id}' does not exist"),
209            );
210        }
211        let tail_offset = stream.tail_offset;
212        let retained_offset = self.earliest_retained_offset(&stream_id);
213        if snapshot_offset < retained_offset {
214            return StreamResponse::error_with_next_offset(
215                StreamErrorCode::StreamGone,
216                format!(
217                    "snapshot offset {snapshot_offset} is older than stream '{}' retained offset {retained_offset}",
218                    stream_id
219                ),
220                retained_offset,
221            );
222        }
223        if snapshot_offset > tail_offset {
224            return StreamResponse::error_with_next_offset(
225                StreamErrorCode::SnapshotConflict,
226                format!(
227                    "snapshot offset {snapshot_offset} is beyond stream '{}' tail {tail_offset}",
228                    stream_id
229                ),
230                tail_offset,
231            );
232        }
233        let digest = super::snapshot_digest(&content_type, &payload);
234        let current_snapshot = self
235            .stream_slot(&stream_id)
236            .and_then(|slot| slot.visible_snapshot.as_ref());
237        if let Some(expected_digest) = expected_digest.as_deref()
238            && current_snapshot.map(|snapshot| snapshot.digest.as_str()) != Some(expected_digest)
239        {
240            return StreamResponse::error_with_next_offset(
241                StreamErrorCode::SnapshotConflict,
242                "current snapshot digest does not match Stream-Snapshot-Match",
243                tail_offset,
244            );
245        }
246        if let Some(current) = current_snapshot {
247            if snapshot_offset < current.offset {
248                return StreamResponse::error_with_next_offset(
249                    StreamErrorCode::SnapshotConflict,
250                    format!(
251                        "snapshot offset {snapshot_offset} is older than latest snapshot offset {}",
252                        current.offset
253                    ),
254                    tail_offset,
255                );
256            }
257            if snapshot_offset == current.offset {
258                if current.digest == digest {
259                    return StreamResponse::SnapshotPublished {
260                        snapshot_offset,
261                        snapshot_digest: digest,
262                        record_range: self.record_range(&stream_id).ok().flatten(),
263                    };
264                }
265                return StreamResponse::error_with_next_offset(
266                    StreamErrorCode::SnapshotConflict,
267                    format!(
268                        "snapshot offset {snapshot_offset} already has a different payload digest"
269                    ),
270                    tail_offset,
271                );
272            }
273        }
274        if !self.snapshot_offset_aligned(&stream_id, snapshot_offset, retained_offset) {
275            return StreamResponse::error_with_next_offset(
276                StreamErrorCode::InvalidSnapshot,
277                format!(
278                    "snapshot offset {snapshot_offset} is not aligned to a committed message boundary for stream '{stream_id}'"
279                ),
280                tail_offset,
281            );
282        }
283
284        let record_range = self.record_range(&stream_id).ok().flatten();
285
286        self.stream_slot_mut(&stream_id)
287            .expect("stream existence checked before snapshot publish")
288            .visible_snapshot = Some(StreamVisibleSnapshot {
289            offset: snapshot_offset,
290            content_type,
291            payload,
292            digest: digest.clone(),
293        });
294        StreamResponse::SnapshotPublished {
295            snapshot_offset,
296            snapshot_digest: digest,
297            record_range,
298        }
299    }
300
301    pub(super) fn advance_retention(
302        &mut self,
303        stream_id: BucketStreamId,
304        retained_offset: u64,
305        now_ms: u64,
306    ) -> StreamResponse {
307        if let Err(response) = self.validate_stream_scope(&stream_id) {
308            return response;
309        }
310        let Some(stream) = self.stream_metadata(&stream_id) else {
311            return StreamResponse::error(
312                StreamErrorCode::StreamNotFound,
313                format!("stream '{stream_id}' does not exist"),
314            );
315        };
316        if stream_is_expired(stream, now_ms) {
317            self.remove_stream_state(&stream_id);
318            return StreamResponse::error(
319                StreamErrorCode::StreamNotFound,
320                format!("stream '{stream_id}' does not exist"),
321            );
322        }
323        let current = self.earliest_retained_offset(&stream_id);
324        if retained_offset < current {
325            return StreamResponse::error_with_next_offset(
326                StreamErrorCode::SnapshotConflict,
327                format!(
328                    "retention offset {retained_offset} is older than current retained offset {current}"
329                ),
330                stream.tail_offset,
331            );
332        }
333        let Some(snapshot) = self
334            .stream_slot(&stream_id)
335            .and_then(|slot| slot.visible_snapshot.as_ref())
336        else {
337            return StreamResponse::error_with_next_offset(
338                StreamErrorCode::SnapshotConflict,
339                "retention requires a published checkpoint",
340                stream.tail_offset,
341            );
342        };
343        if retained_offset > snapshot.offset {
344            return StreamResponse::error_with_next_offset(
345                StreamErrorCode::SnapshotConflict,
346                format!(
347                    "retention offset {retained_offset} is beyond latest checkpoint offset {}",
348                    snapshot.offset
349                ),
350                stream.tail_offset,
351            );
352        }
353        if retained_offset == current {
354            return StreamResponse::RetentionAdvanced {
355                retained_offset,
356                record_range: self.record_range(&stream_id).ok().flatten(),
357            };
358        }
359        if !self.snapshot_offset_aligned(&stream_id, retained_offset, current) {
360            return StreamResponse::error_with_next_offset(
361                StreamErrorCode::InvalidSnapshot,
362                format!(
363                    "retention offset {retained_offset} is not aligned to a committed message boundary for stream '{stream_id}'"
364                ),
365                stream.tail_offset,
366            );
367        }
368        let mut retained_record_index = self
369            .stream_slot(&stream_id)
370            .expect("stream existence checked before retention")
371            .record_index
372            .clone();
373        if let Some(record_index) = retained_record_index.as_mut()
374            && record_index
375                .retain_from_offset(retained_offset, stream.tail_offset)
376                .is_err()
377        {
378            return StreamResponse::error_with_next_offset(
379                StreamErrorCode::InvalidRecordBoundaries,
380                format!(
381                    "retention offset {retained_offset} is not a retained record boundary for stream '{stream_id}'"
382                ),
383                stream.tail_offset,
384            );
385        }
386        let slot = self
387            .stream_slot_mut(&stream_id)
388            .expect("stream existence checked before retention mutation");
389        let previous_retained_offset = slot.retained_offset;
390        slot.retained_offset = retained_offset;
391        self.usage_on_retention(
392            &stream_id.bucket_id,
393            retained_offset.saturating_sub(previous_retained_offset),
394        );
395        self.compact_retained_prefix(&stream_id, retained_offset, retained_record_index);
396        StreamResponse::RetentionAdvanced {
397            retained_offset,
398            record_range: self.record_range(&stream_id).ok().flatten(),
399        }
400    }
401
402    pub(super) fn flush_cold(
403        &mut self,
404        stream_id: BucketStreamId,
405        chunk: ColdChunkRef,
406    ) -> StreamResponse {
407        if let Err(response) = self.validate_stream_scope(&stream_id) {
408            return response;
409        }
410        if chunk.s3_path.trim().is_empty() {
411            return StreamResponse::error(
412                StreamErrorCode::InvalidColdFlush,
413                "cold chunk S3 path must not be empty",
414            );
415        }
416        if chunk.object_size == 0 {
417            return StreamResponse::error(
418                StreamErrorCode::InvalidColdFlush,
419                "cold chunk object size must be greater than zero",
420            );
421        }
422        let Some(slot) = self.stream_slot(&stream_id) else {
423            return StreamResponse::error(
424                StreamErrorCode::StreamNotFound,
425                format!("stream '{stream_id}' does not exist"),
426            );
427        };
428        let stream = &slot.metadata;
429        if chunk.end_offset <= chunk.start_offset {
430            return StreamResponse::error_with_next_offset(
431                StreamErrorCode::InvalidColdFlush,
432                "cold chunk must cover at least one byte",
433                stream.tail_offset,
434            );
435        }
436        let logical_size = chunk.end_offset.saturating_sub(chunk.start_offset);
437        if chunk
438            .object_offset
439            .checked_add(logical_size)
440            .is_none_or(|end| end > chunk.object_size)
441        {
442            return StreamResponse::error_with_next_offset(
443                StreamErrorCode::InvalidColdFlush,
444                "cold chunk slice is outside the physical object",
445                stream.tail_offset,
446            );
447        }
448        if !chunk.shared_object && chunk.object_offset != 0 {
449            return StreamResponse::error_with_next_offset(
450                StreamErrorCode::InvalidColdFlush,
451                "exclusive cold chunks must start at physical object offset zero",
452                stream.tail_offset,
453            );
454        }
455        if chunk.end_offset > stream.tail_offset {
456            return StreamResponse::error_with_next_offset_and_context(
457                StreamErrorCode::InvalidColdFlush,
458                format!(
459                    "cold chunk end {} is beyond stream '{}' tail {}",
460                    chunk.end_offset, stream_id, stream.tail_offset
461                ),
462                stream.tail_offset,
463                vec![StreamErrorContext::StaleColdFlushCandidate],
464            );
465        }
466        let hot_buffer = &slot.hot_buffer;
467        if hot_buffer.hot_start_offset() != chunk.start_offset {
468            return StreamResponse::error_with_next_offset_and_context(
469                StreamErrorCode::InvalidColdFlush,
470                format!("cold chunk for stream '{stream_id}' must start at the hot prefix"),
471                stream.tail_offset,
472                vec![StreamErrorContext::StaleColdFlushCandidate],
473            );
474        }
475        if !hot_buffer.covers_prefix(chunk.start_offset, chunk.end_offset) {
476            return StreamResponse::error_with_next_offset_and_context(
477                StreamErrorCode::InvalidColdFlush,
478                format!(
479                    "cold chunk for stream '{stream_id}' does not cover contiguous hot payload"
480                ),
481                stream.tail_offset,
482                vec![StreamErrorContext::StaleColdFlushCandidate],
483            );
484        }
485        if !chunk.payload_digest.is_empty()
486            && hot_buffer
487                .digest_prefix(chunk.start_offset, chunk.end_offset)
488                .as_deref()
489                != Some(chunk.payload_digest.as_str())
490        {
491            return StreamResponse::error_with_next_offset_and_context(
492                StreamErrorCode::InvalidColdFlush,
493                format!("cold chunk payload for stream '{stream_id}' is stale"),
494                stream.tail_offset,
495                vec![StreamErrorContext::StaleColdFlushCandidate],
496            );
497        }
498        let shared_path = chunk.shared_object.then(|| chunk.s3_path.clone());
499        let slot = self
500            .stream_slot_mut(&stream_id)
501            .expect("stream existence checked before cold flush mutation");
502        let hot_bytes_before = u64::try_from(slot.hot_buffer.len()).expect("payload len fits u64");
503        slot.hot_buffer.flush_prefix(chunk.end_offset);
504        let hot_bytes_after = u64::try_from(slot.hot_buffer.len()).expect("payload len fits u64");
505        slot.cold.push_cold_chunk(chunk.clone());
506        self.remove_hot_payload_bytes(hot_bytes_before.saturating_sub(hot_bytes_after));
507        if let Some(path) = shared_path {
508            self.retain_shared_cold_object(&path, &stream_id.bucket_id);
509        }
510        self.compact_message_records_before(
511            &stream_id,
512            self.earliest_retained_offset(&stream_id),
513            chunk.end_offset,
514        );
515        StreamResponse::ColdFlushed {
516            hot_start_offset: self.hot_start_offset(&stream_id),
517        }
518    }
519
520    pub(super) fn compact_cold(
521        &mut self,
522        stream_id: BucketStreamId,
523        old_chunks: Vec<ColdChunkRef>,
524        replacement: ColdChunkRef,
525        gc_not_before_ms: u64,
526    ) -> StreamResponse {
527        if let Err(response) = self.validate_stream_scope(&stream_id) {
528            return response;
529        }
530        if self.stream_slot(&stream_id).is_none() {
531            return StreamResponse::error(
532                StreamErrorCode::StreamNotFound,
533                format!("stream '{stream_id}' does not exist"),
534            );
535        }
536        if old_chunks.is_empty()
537            || (old_chunks.len() < 2 && old_chunks.iter().all(|chunk| !chunk.shared_object))
538        {
539            return StreamResponse::error(
540                StreamErrorCode::InvalidColdFlush,
541                "cold compaction requires two raw chunks or one legacy shared chunk",
542            );
543        }
544        if replacement.s3_path.trim().is_empty() || replacement.object_size == 0 {
545            return StreamResponse::error(
546                StreamErrorCode::InvalidColdFlush,
547                "cold compaction replacement must name a non-empty object",
548            );
549        }
550        let rewriting_shared = old_chunks.iter().all(|chunk| chunk.shared_object);
551        if old_chunks.iter().any(|chunk| chunk.shared_object) != rewriting_shared
552            || replacement.shared_object
553            || replacement.object_offset != 0
554        {
555            return StreamResponse::error(
556                StreamErrorCode::InvalidColdFlush,
557                "cold compaction cannot mix shared and raw inputs or publish a shared replacement",
558            );
559        }
560        let mut expected_start = old_chunks
561            .first()
562            .map_or(replacement.start_offset, |chunk| chunk.start_offset);
563        let mut compacted_bytes = 0_u64;
564        for chunk in &old_chunks {
565            let logical_bytes = chunk.end_offset.saturating_sub(chunk.start_offset);
566            if chunk.start_offset != expected_start
567                || chunk.end_offset <= chunk.start_offset
568                || (!chunk.shared_object && chunk.object_size != logical_bytes)
569            {
570                return StreamResponse::error(
571                    StreamErrorCode::InvalidColdFlush,
572                    "cold compaction inputs must be contiguous raw chunks",
573                );
574            }
575            expected_start = chunk.end_offset;
576            compacted_bytes = compacted_bytes.saturating_add(logical_bytes);
577        }
578        let first = old_chunks
579            .first()
580            .expect("cold compaction input count validated");
581        if replacement.start_offset != first.start_offset
582            || replacement.end_offset != expected_start
583            || replacement.object_size != compacted_bytes
584        {
585            return StreamResponse::error(
586                StreamErrorCode::InvalidColdFlush,
587                "cold compaction replacement must cover the exact input range",
588            );
589        }
590        if rewriting_shared {
591            let slot = self
592                .stream_slot_mut(&stream_id)
593                .expect("stream existence checked before cold compaction");
594            if !slot.cold.remove_shared_chunks(&old_chunks) {
595                return StreamResponse::error(
596                    StreamErrorCode::InvalidColdFlush,
597                    "legacy shared compaction input no longer matches the stream state",
598                );
599            }
600        }
601        let compacted_chunks = u64::try_from(old_chunks.len()).expect("chunk count fits u64");
602        let mut exclusive_paths = Vec::new();
603        let mut shared_paths = Vec::new();
604        for chunk in old_chunks {
605            if chunk.shared_object {
606                shared_paths.push(chunk.s3_path);
607            } else {
608                exclusive_paths.push(chunk.s3_path);
609            }
610        }
611        if !exclusive_paths.is_empty() {
612            self.cold_gc.enqueue_after(
613                stream_id.bucket_id.clone(),
614                ColdGcTarget::Paths(exclusive_paths),
615                gc_not_before_ms,
616            );
617        }
618        self.release_shared_cold_objects(&stream_id.bucket_id, shared_paths, gc_not_before_ms);
619        StreamResponse::ColdCompacted {
620            compacted_chunks,
621            compacted_bytes,
622        }
623    }
624
625    pub fn delete_snapshot(
626        &self,
627        stream_id: &BucketStreamId,
628        snapshot_offset: u64,
629    ) -> StreamResponse {
630        match self.latest_snapshot(stream_id) {
631            Ok(Some(snapshot)) if snapshot.offset == snapshot_offset => StreamResponse::error(
632                StreamErrorCode::SnapshotConflict,
633                format!(
634                    "snapshot {snapshot_offset} for stream '{stream_id}' is the latest visible snapshot"
635                ),
636            ),
637            Ok(_) => StreamResponse::error(
638                StreamErrorCode::SnapshotNotFound,
639                format!("snapshot {snapshot_offset} for stream '{stream_id}' does not exist"),
640            ),
641            Err(err) => err,
642        }
643    }
644
645    pub(super) fn ack_cold_gc(&mut self, up_to_seq: u64) -> StreamResponse {
646        let removed = self.cold_gc.ack(up_to_seq);
647        StreamResponse::ColdGcAcked { removed }
648    }
649
650    /// A bounded snapshot of the front of the GC queue for the leader's worker
651    /// to reclaim. Read-only; draining is confirmed by a replicated `AckColdGc`.
652    pub fn pending_cold_gc_batch(&self, max: usize) -> Vec<ColdGcEntry> {
653        self.cold_gc.batch(max)
654    }
655
656    pub fn pending_cold_gc_len(&self) -> usize {
657        self.cold_gc.len()
658    }
659
660    pub fn pending_cold_gc_len_for_bucket(&self, bucket_id: &str) -> usize {
661        self.cold_gc.len_for_bucket(bucket_id)
662    }
663
664    pub(super) fn earliest_retained_offset(&self, stream_id: &BucketStreamId) -> u64 {
665        self.stream_slot(stream_id)
666            .map(|slot| slot.retained_offset)
667            .unwrap_or(0)
668    }
669
670    pub(super) fn snapshot_offset_aligned(
671        &self,
672        stream_id: &BucketStreamId,
673        snapshot_offset: u64,
674        retained_offset: u64,
675    ) -> bool {
676        snapshot_offset == retained_offset
677            || snapshot_offset <= self.cold_frontier_offset(stream_id, retained_offset)
678            || self
679                .stream_slot(stream_id)
680                .is_some_and(|slot| snapshot_offset <= slot.hot_buffer.hot_start_offset())
681            || self.stream_slot(stream_id).is_some_and(|slot| {
682                slot.message_records
683                    .iter()
684                    .any(|record| record.end_offset == snapshot_offset)
685            })
686    }
687
688    pub(super) fn compact_retained_prefix(
689        &mut self,
690        stream_id: &BucketStreamId,
691        retained_offset: u64,
692        retained_record_index: Option<crate::StreamRecordIndex>,
693    ) {
694        let frontier = self.cold_frontier_offset(stream_id, retained_offset).max(
695            self.stream_slot(stream_id)
696                .map(|slot| slot.hot_buffer.hot_start_offset())
697                .unwrap_or(retained_offset),
698        );
699        self.compact_message_records_before(stream_id, retained_offset, frontier);
700        let slot = self
701            .stream_slot_mut(stream_id)
702            .expect("stream existence checked before retained-prefix compaction");
703        slot.record_index = retained_record_index;
704        slot.integrity.evict_before(retained_offset);
705        let dropped_cold_paths = slot.cold.compact_before(retained_offset);
706        self.release_shared_cold_objects(&stream_id.bucket_id, dropped_cold_paths, 0);
707
708        let slot = self
709            .stream_slot_mut(stream_id)
710            .expect("stream existence checked before hot compact");
711        let hot_bytes_before = u64::try_from(slot.hot_buffer.len()).expect("payload len fits u64");
712        slot.hot_buffer.discard_before(retained_offset);
713        let hot_bytes_after = u64::try_from(slot.hot_buffer.len()).expect("payload len fits u64");
714        self.remove_hot_payload_bytes(hot_bytes_before.saturating_sub(hot_bytes_after));
715    }
716
717    pub(super) fn compact_message_records_before(
718        &mut self,
719        stream_id: &BucketStreamId,
720        retained_offset: u64,
721        frontier: u64,
722    ) {
723        let slot = self
724            .stream_slot_mut(stream_id)
725            .expect("stream existence checked before message-record compaction");
726        let records = std::mem::take(&mut slot.message_records);
727        let frontier = frontier.max(retained_offset);
728        let mut compacted = Vec::with_capacity(records.len());
729        if frontier > retained_offset {
730            compacted.push(StreamMessageRecord {
731                start_offset: retained_offset,
732                end_offset: frontier,
733            });
734        }
735        compacted.extend(records.iter().filter_map(|record| {
736            if record.end_offset <= frontier {
737                return None;
738            }
739            let start_offset = record.start_offset.max(frontier).max(retained_offset);
740            (record.end_offset > start_offset).then_some(StreamMessageRecord {
741                start_offset,
742                end_offset: record.end_offset,
743            })
744        }));
745        if compacted.is_empty() {
746            return;
747        }
748        self.stream_slot_mut(stream_id)
749            .expect("stream existence checked before message record compact")
750            .message_records = compacted;
751    }
752
753    pub(super) fn cold_frontier_offset(
754        &self,
755        stream_id: &BucketStreamId,
756        retained_offset: u64,
757    ) -> u64 {
758        self.stream_slot(stream_id)
759            .map(|slot| slot.cold.cold_frontier_offset(retained_offset))
760            .unwrap_or(retained_offset)
761    }
762}