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        Ok(Some(ColdFlushCandidate {
52            stream_id: stream_id.clone(),
53            start_offset,
54            end_offset,
55            payload,
56        }))
57    }
58
59    pub(super) fn plan_next_cold_flush_from_start(
60        &self,
61        mut start_fn: impl FnMut(&BucketStreamId) -> u64,
62        min_hot_bytes: usize,
63        max_flush_bytes: usize,
64        group_hot_bytes: u64,
65    ) -> Result<Option<ColdFlushCandidate>, StreamResponse> {
66        if max_flush_bytes == 0 {
67            return Ok(None);
68        }
69        let mut stream_ids = self.registry.stream_ids().cloned().collect::<Vec<_>>();
70        stream_ids.sort_by(compare_stream_ids);
71        for stream_id in &stream_ids {
72            let start = start_fn(stream_id);
73            match self.plan_cold_flush_with_start(stream_id, start, min_hot_bytes, max_flush_bytes)
74            {
75                Ok(Some(candidate)) => return Ok(Some(candidate)),
76                Ok(None) => {}
77                Err(StreamResponse::Error {
78                    code: StreamErrorCode::StreamGone | StreamErrorCode::StreamNotFound,
79                    ..
80                }) => {}
81                Err(err) => return Err(err),
82            }
83        }
84        let group_min_hot_bytes = u64::try_from(min_hot_bytes).unwrap_or(u64::MAX);
85        if group_hot_bytes < group_min_hot_bytes {
86            return Ok(None);
87        }
88        for stream_id in stream_ids {
89            let start = start_fn(&stream_id);
90            match self.plan_cold_flush_with_start(&stream_id, start, 1, max_flush_bytes) {
91                Ok(Some(candidate)) => return Ok(Some(candidate)),
92                Ok(None) => {}
93                Err(StreamResponse::Error {
94                    code: StreamErrorCode::StreamGone | StreamErrorCode::StreamNotFound,
95                    ..
96                }) => {}
97                Err(err) => return Err(err),
98            }
99        }
100        Ok(None)
101    }
102
103    pub fn plan_next_cold_flush_batch(
104        &self,
105        min_hot_bytes: usize,
106        max_flush_bytes: usize,
107        max_candidates: usize,
108    ) -> Result<Vec<ColdFlushCandidate>, StreamResponse> {
109        if max_candidates == 0 || max_flush_bytes == 0 {
110            return Ok(Vec::new());
111        }
112        let mut planned_flush_offsets: HashMap<BucketStreamId, u64> = HashMap::new();
113        let mut candidates = Vec::with_capacity(max_candidates);
114        while candidates.len() < max_candidates {
115            let start_for = |stream_id: &BucketStreamId| -> u64 {
116                planned_flush_offsets
117                    .get(stream_id)
118                    .copied()
119                    .unwrap_or_else(|| self.hot_start_offset(stream_id))
120            };
121            let group_hot_bytes: u64 = self
122                .registry
123                .stream_ids()
124                .map(|stream_id| {
125                    let start = start_for(stream_id);
126                    self.stream_slot(stream_id)
127                        .map(|slot| {
128                            u64::try_from(slot.hot_buffer.remaining_len_from(start))
129                                .expect("len fits u64")
130                        })
131                        .unwrap_or(0)
132                })
133                .sum();
134            let candidate = self.plan_next_cold_flush_from_start(
135                start_for,
136                min_hot_bytes,
137                max_flush_bytes,
138                group_hot_bytes,
139            )?;
140            let Some(candidate) = candidate else {
141                break;
142            };
143            planned_flush_offsets.insert(candidate.stream_id.clone(), candidate.end_offset);
144            candidates.push(candidate);
145        }
146        Ok(candidates)
147    }
148
149    pub(super) fn publish_snapshot(
150        &mut self,
151        stream_id: BucketStreamId,
152        snapshot_offset: u64,
153        content_type: String,
154        payload: Vec<u8>,
155        now_ms: u64,
156    ) -> StreamResponse {
157        if let Err(response) = self.validate_stream_scope(&stream_id) {
158            return response;
159        }
160        if content_type.trim().is_empty() {
161            return StreamResponse::error(
162                StreamErrorCode::InvalidSnapshot,
163                "snapshot content type must not be empty",
164            );
165        }
166        let Some(stream) = self.stream_metadata(&stream_id) else {
167            return StreamResponse::error(
168                StreamErrorCode::StreamNotFound,
169                format!("stream '{stream_id}' does not exist"),
170            );
171        };
172        if stream_is_expired(stream, now_ms) {
173            self.remove_stream_state(&stream_id);
174            return StreamResponse::error(
175                StreamErrorCode::StreamNotFound,
176                format!("stream '{stream_id}' does not exist"),
177            );
178        }
179        let tail_offset = stream.tail_offset;
180        let retained_offset = self.earliest_retained_offset(&stream_id);
181        if snapshot_offset < retained_offset {
182            return StreamResponse::error_with_next_offset(
183                StreamErrorCode::StreamGone,
184                format!(
185                    "snapshot offset {snapshot_offset} is older than stream '{}' retained offset {retained_offset}",
186                    stream_id
187                ),
188                retained_offset,
189            );
190        }
191        if snapshot_offset > tail_offset {
192            return StreamResponse::error_with_next_offset(
193                StreamErrorCode::SnapshotConflict,
194                format!(
195                    "snapshot offset {snapshot_offset} is beyond stream '{}' tail {tail_offset}",
196                    stream_id
197                ),
198                tail_offset,
199            );
200        }
201        if !self.snapshot_offset_aligned(&stream_id, snapshot_offset, retained_offset) {
202            return StreamResponse::error_with_next_offset(
203                StreamErrorCode::InvalidSnapshot,
204                format!(
205                    "snapshot offset {snapshot_offset} is not aligned to a committed message boundary for stream '{stream_id}'"
206                ),
207                tail_offset,
208            );
209        }
210
211        let mut retained_record_index = self
212            .stream_slot(&stream_id)
213            .expect("stream existence checked before snapshot publish")
214            .record_index
215            .clone();
216        let record_range = if let Some(record_index) = retained_record_index.as_mut() {
217            if record_index
218                .retain_from_offset(snapshot_offset, tail_offset)
219                .is_err()
220            {
221                return StreamResponse::error_with_next_offset(
222                    StreamErrorCode::InvalidRecordBoundaries,
223                    format!(
224                        "snapshot offset {snapshot_offset} is not a retained record boundary for stream '{stream_id}'"
225                    ),
226                    tail_offset,
227                );
228            }
229            match record_index.range() {
230                Ok(range) => Some(range),
231                Err(_) => {
232                    return StreamResponse::error_with_next_offset(
233                        StreamErrorCode::InvalidRecordBoundaries,
234                        format!("stream '{stream_id}' has an invalid retained record index"),
235                        tail_offset,
236                    );
237                }
238            }
239        } else {
240            None
241        };
242
243        self.stream_slot_mut(&stream_id)
244            .expect("stream existence checked before snapshot publish")
245            .visible_snapshot = Some(StreamVisibleSnapshot {
246            offset: snapshot_offset,
247            content_type,
248            payload,
249        });
250        self.compact_retained_prefix(&stream_id, snapshot_offset, retained_record_index);
251        StreamResponse::SnapshotPublished {
252            snapshot_offset,
253            record_range,
254        }
255    }
256
257    pub(super) fn flush_cold(
258        &mut self,
259        stream_id: BucketStreamId,
260        chunk: ColdChunkRef,
261    ) -> StreamResponse {
262        if let Err(response) = self.validate_stream_scope(&stream_id) {
263            return response;
264        }
265        if chunk.s3_path.trim().is_empty() {
266            return StreamResponse::error(
267                StreamErrorCode::InvalidColdFlush,
268                "cold chunk S3 path must not be empty",
269            );
270        }
271        if chunk.object_size == 0 {
272            return StreamResponse::error(
273                StreamErrorCode::InvalidColdFlush,
274                "cold chunk object size must be greater than zero",
275            );
276        }
277        let Some(slot) = self.stream_slot(&stream_id) else {
278            return StreamResponse::error(
279                StreamErrorCode::StreamNotFound,
280                format!("stream '{stream_id}' does not exist"),
281            );
282        };
283        let stream = &slot.metadata;
284        if chunk.end_offset <= chunk.start_offset {
285            return StreamResponse::error_with_next_offset(
286                StreamErrorCode::InvalidColdFlush,
287                "cold chunk must cover at least one byte",
288                stream.tail_offset,
289            );
290        }
291        if chunk.end_offset > stream.tail_offset {
292            return StreamResponse::error_with_next_offset_and_context(
293                StreamErrorCode::InvalidColdFlush,
294                format!(
295                    "cold chunk end {} is beyond stream '{}' tail {}",
296                    chunk.end_offset, stream_id, stream.tail_offset
297                ),
298                stream.tail_offset,
299                vec![StreamErrorContext::StaleColdFlushCandidate],
300            );
301        }
302        let hot_buffer = &slot.hot_buffer;
303        if hot_buffer.hot_start_offset() != chunk.start_offset {
304            return StreamResponse::error_with_next_offset_and_context(
305                StreamErrorCode::InvalidColdFlush,
306                format!("cold chunk for stream '{stream_id}' must start at the hot prefix"),
307                stream.tail_offset,
308                vec![StreamErrorContext::StaleColdFlushCandidate],
309            );
310        }
311        if !hot_buffer.covers_prefix(chunk.start_offset, chunk.end_offset) {
312            return StreamResponse::error_with_next_offset_and_context(
313                StreamErrorCode::InvalidColdFlush,
314                format!(
315                    "cold chunk for stream '{stream_id}' does not cover contiguous hot payload"
316                ),
317                stream.tail_offset,
318                vec![StreamErrorContext::StaleColdFlushCandidate],
319            );
320        }
321        let slot = self
322            .stream_slot_mut(&stream_id)
323            .expect("stream existence checked before cold flush mutation");
324        slot.hot_buffer.flush_prefix(chunk.end_offset);
325        slot.cold.push_cold_chunk(chunk.clone());
326        self.compact_message_records_before(
327            &stream_id,
328            self.earliest_retained_offset(&stream_id),
329            chunk.end_offset,
330        );
331        StreamResponse::ColdFlushed {
332            hot_start_offset: self.hot_start_offset(&stream_id),
333        }
334    }
335
336    pub fn delete_snapshot(
337        &self,
338        stream_id: &BucketStreamId,
339        snapshot_offset: u64,
340    ) -> StreamResponse {
341        match self.latest_snapshot(stream_id) {
342            Ok(Some(snapshot)) if snapshot.offset == snapshot_offset => StreamResponse::error(
343                StreamErrorCode::SnapshotConflict,
344                format!(
345                    "snapshot {snapshot_offset} for stream '{stream_id}' is the latest visible snapshot"
346                ),
347            ),
348            Ok(_) => StreamResponse::error(
349                StreamErrorCode::SnapshotNotFound,
350                format!("snapshot {snapshot_offset} for stream '{stream_id}' does not exist"),
351            ),
352            Err(err) => err,
353        }
354    }
355
356    pub(super) fn ack_cold_gc(&mut self, up_to_seq: u64) -> StreamResponse {
357        let removed = self.cold_gc.ack(up_to_seq);
358        StreamResponse::ColdGcAcked { removed }
359    }
360
361    /// A bounded snapshot of the front of the GC queue for the leader's worker
362    /// to reclaim. Read-only; draining is confirmed by a replicated `AckColdGc`.
363    pub fn pending_cold_gc_batch(&self, max: usize) -> Vec<ColdGcEntry> {
364        self.cold_gc.batch(max)
365    }
366
367    pub fn pending_cold_gc_len(&self) -> usize {
368        self.cold_gc.len()
369    }
370
371    pub(super) fn earliest_retained_offset(&self, stream_id: &BucketStreamId) -> u64 {
372        self.stream_slot(stream_id)
373            .and_then(|slot| slot.visible_snapshot.as_ref())
374            .map(|snapshot| snapshot.offset)
375            .unwrap_or(0)
376    }
377
378    pub(super) fn snapshot_offset_aligned(
379        &self,
380        stream_id: &BucketStreamId,
381        snapshot_offset: u64,
382        retained_offset: u64,
383    ) -> bool {
384        snapshot_offset == retained_offset
385            || snapshot_offset <= self.cold_frontier_offset(stream_id, retained_offset)
386            || self
387                .stream_slot(stream_id)
388                .is_some_and(|slot| snapshot_offset <= slot.hot_buffer.hot_start_offset())
389            || self.stream_slot(stream_id).is_some_and(|slot| {
390                slot.message_records
391                    .iter()
392                    .any(|record| record.end_offset == snapshot_offset)
393            })
394    }
395
396    pub(super) fn compact_retained_prefix(
397        &mut self,
398        stream_id: &BucketStreamId,
399        retained_offset: u64,
400        retained_record_index: Option<crate::StreamRecordIndex>,
401    ) {
402        let frontier = self.cold_frontier_offset(stream_id, retained_offset).max(
403            self.stream_slot(stream_id)
404                .map(|slot| slot.hot_buffer.hot_start_offset())
405                .unwrap_or(retained_offset),
406        );
407        self.compact_message_records_before(stream_id, retained_offset, frontier);
408        let slot = self
409            .stream_slot_mut(stream_id)
410            .expect("stream existence checked before retained-prefix compaction");
411        slot.record_index = retained_record_index;
412        slot.integrity.evict_before(retained_offset);
413        let dropped_cold_paths = slot.cold.compact_before(retained_offset);
414        if !dropped_cold_paths.is_empty() {
415            self.cold_gc
416                .enqueue(ColdGcTarget::Paths(dropped_cold_paths));
417        }
418
419        self.stream_slot_mut(stream_id)
420            .expect("stream existence checked before hot compact")
421            .hot_buffer
422            .discard_before(retained_offset);
423    }
424
425    pub(super) fn compact_message_records_before(
426        &mut self,
427        stream_id: &BucketStreamId,
428        retained_offset: u64,
429        frontier: u64,
430    ) {
431        let slot = self
432            .stream_slot_mut(stream_id)
433            .expect("stream existence checked before message-record compaction");
434        let records = std::mem::take(&mut slot.message_records);
435        let frontier = frontier.max(retained_offset);
436        let mut compacted = Vec::with_capacity(records.len());
437        if frontier > retained_offset {
438            compacted.push(StreamMessageRecord {
439                start_offset: retained_offset,
440                end_offset: frontier,
441            });
442        }
443        compacted.extend(records.iter().filter_map(|record| {
444            if record.end_offset <= frontier {
445                return None;
446            }
447            let start_offset = record.start_offset.max(frontier).max(retained_offset);
448            (record.end_offset > start_offset).then_some(StreamMessageRecord {
449                start_offset,
450                end_offset: record.end_offset,
451            })
452        }));
453        if compacted.is_empty() {
454            return;
455        }
456        self.stream_slot_mut(stream_id)
457            .expect("stream existence checked before message record compact")
458            .message_records = compacted;
459    }
460
461    pub(super) fn cold_frontier_offset(
462        &self,
463        stream_id: &BucketStreamId,
464        retained_offset: u64,
465    ) -> u64 {
466        self.stream_slot(stream_id)
467            .map(|slot| slot.cold.cold_frontier_offset(retained_offset))
468            .unwrap_or(retained_offset)
469    }
470}