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