ursula_stream/state_machine/
cold.rs1use 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 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}