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