1use 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 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}