1use std::collections::HashMap;
2use std::collections::HashSet;
3use std::collections::VecDeque;
4
5use ursula_shard::BucketStreamId;
6
7use crate::command::StreamCommand;
8use crate::integrity::StreamIntegrity;
9use crate::model::AppendExternalInput;
10use crate::model::AppendStreamInput;
11use crate::model::COLD_INDEX_PAGE_SPAN_BYTES;
12use crate::model::ColdChunkRef;
13use crate::model::ColdFlushCandidate;
14use crate::model::ColdGcEntry;
15use crate::model::ColdGcTarget;
16use crate::model::ExternalPayloadRef;
17use crate::model::HotPayloadSegment;
18use crate::model::ObjectPayloadRef;
19use crate::model::ProducerAppendRecord;
20use crate::model::ProducerRequest;
21use crate::model::ProducerSnapshot;
22use crate::model::ProducerState;
23use crate::model::StreamBatchAppend;
24use crate::model::StreamBatchAppendItem;
25use crate::model::StreamBootstrapPlan;
26use crate::model::StreamMessageRecord;
27use crate::model::StreamMetadata;
28use crate::model::StreamRead;
29use crate::model::StreamReadColdIndexSegment;
30use crate::model::StreamReadObjectSegment;
31use crate::model::StreamReadPlan;
32use crate::model::StreamReadSegment;
33use crate::model::StreamStatus;
34use crate::model::StreamVisibleSnapshot;
35use crate::response::StreamErrorCode;
36use crate::response::StreamErrorContext;
37use crate::response::StreamResponse;
38use crate::snapshot::StreamSnapshot;
39use crate::snapshot::StreamSnapshotEntry;
40use crate::snapshot::StreamSnapshotError;
41use crate::validate::validate_bucket_id;
42use crate::validate::validate_stream_id;
43
44#[derive(Debug, Clone, Default)]
45pub struct StreamStateMachine {
46 buckets: HashSet<String>,
47 streams: HashMap<BucketStreamId, StreamMetadata>,
48 hot_buffers: HashMap<BucketStreamId, HotBuffer>,
49 cold_index: InMemoryColdIndex,
50 message_records: HashMap<BucketStreamId, Vec<StreamMessageRecord>>,
51 integrities: HashMap<BucketStreamId, StreamIntegrity>,
52 visible_snapshots: HashMap<BucketStreamId, StreamVisibleSnapshot>,
53 producers: HashMap<BucketStreamId, HashMap<String, ProducerState>>,
54 pending_cold_gc: VecDeque<ColdGcEntry>,
55 next_cold_gc_seq: u64,
56}
57
58trait ColdIndexState {
59 fn cold_chunks(&self, stream_id: &BucketStreamId) -> &[ColdChunkRef];
60 fn external_segments(&self, stream_id: &BucketStreamId) -> &[ObjectPayloadRef];
61 fn cold_generation(&self, stream_id: &BucketStreamId) -> u64;
62 fn push_cold_chunk(&mut self, stream_id: BucketStreamId, chunk: ColdChunkRef);
63 fn push_external_segment(&mut self, stream_id: BucketStreamId, object: ObjectPayloadRef);
64 fn restore_stream(
65 &mut self,
66 stream_id: BucketStreamId,
67 cold_frontier_offset: u64,
68 cold_index_generation: u64,
69 cold_chunks: Vec<ColdChunkRef>,
70 external_segments: Vec<ObjectPayloadRef>,
71 );
72 fn remove_stream(&mut self, stream_id: &BucketStreamId) -> bool;
73 fn compact_before(&mut self, stream_id: &BucketStreamId, retained_offset: u64) -> Vec<String>;
74 fn cold_frontier_offset(&self, stream_id: &BucketStreamId, retained_offset: u64) -> u64;
75}
76
77#[derive(Debug, Clone, Default)]
78struct InMemoryColdIndex {
79 cold_chunks: HashMap<BucketStreamId, Vec<ColdChunkRef>>,
80 external_segments: HashMap<BucketStreamId, Vec<ObjectPayloadRef>>,
81 cold_frontiers: HashMap<BucketStreamId, u64>,
82}
83
84impl ColdIndexState for InMemoryColdIndex {
85 fn cold_chunks(&self, stream_id: &BucketStreamId) -> &[ColdChunkRef] {
86 self.cold_chunks
87 .get(stream_id)
88 .map(Vec::as_slice)
89 .unwrap_or(&[])
90 }
91
92 fn external_segments(&self, stream_id: &BucketStreamId) -> &[ObjectPayloadRef] {
93 self.external_segments
94 .get(stream_id)
95 .map(Vec::as_slice)
96 .unwrap_or(&[])
97 }
98
99 fn cold_generation(&self, _stream_id: &BucketStreamId) -> u64 {
100 0
101 }
102
103 fn push_cold_chunk(&mut self, stream_id: BucketStreamId, chunk: ColdChunkRef) {
104 self.cold_frontiers.insert(stream_id, chunk.end_offset);
105 }
106
107 fn push_external_segment(&mut self, stream_id: BucketStreamId, object: ObjectPayloadRef) {
108 let frontier = self
109 .cold_frontiers
110 .get(&stream_id)
111 .copied()
112 .unwrap_or(object.start_offset)
113 .max(object.end_offset);
114 self.cold_frontiers.insert(stream_id, frontier);
115 }
116
117 fn restore_stream(
118 &mut self,
119 stream_id: BucketStreamId,
120 cold_frontier_offset: u64,
121 _cold_index_generation: u64,
122 cold_chunks: Vec<ColdChunkRef>,
123 external_segments: Vec<ObjectPayloadRef>,
124 ) {
125 if cold_frontier_offset > 0 {
126 self.cold_frontiers
127 .insert(stream_id.clone(), cold_frontier_offset);
128 }
129 if !cold_chunks.is_empty() {
130 self.cold_chunks.insert(stream_id.clone(), cold_chunks);
131 }
132 if !external_segments.is_empty() {
133 self.external_segments.insert(stream_id, external_segments);
134 }
135 }
136
137 fn remove_stream(&mut self, stream_id: &BucketStreamId) -> bool {
138 self.cold_frontiers
139 .remove(stream_id)
140 .is_some_and(|frontier| frontier > 0)
141 || self
142 .cold_chunks
143 .remove(stream_id)
144 .is_some_and(|chunks| !chunks.is_empty())
145 || self
146 .external_segments
147 .remove(stream_id)
148 .is_some_and(|objects| !objects.is_empty())
149 }
150
151 fn compact_before(&mut self, stream_id: &BucketStreamId, retained_offset: u64) -> Vec<String> {
152 let mut dropped_cold_paths = Vec::new();
153 if let Some(chunks) = self.cold_chunks.get_mut(stream_id) {
154 chunks.retain(|chunk| {
155 let retain = chunk.end_offset > retained_offset;
156 if !retain {
157 dropped_cold_paths.push(chunk.s3_path.clone());
158 }
159 retain
160 });
161 if chunks.is_empty() {
162 self.cold_chunks.remove(stream_id);
163 }
164 }
165 if let Some(objects) = self.external_segments.get_mut(stream_id) {
166 objects.retain(|object| object.end_offset > retained_offset);
167 if objects.is_empty() {
168 self.external_segments.remove(stream_id);
169 }
170 }
171 dropped_cold_paths
172 }
173
174 fn cold_frontier_offset(&self, stream_id: &BucketStreamId, retained_offset: u64) -> u64 {
175 let external_segments = self.external_segments(stream_id);
176 let cold_frontier = self
177 .cold_frontiers
178 .get(stream_id)
179 .copied()
180 .unwrap_or(retained_offset);
181 let mut ranges = Vec::with_capacity(1 + external_segments.len());
182 if cold_frontier > retained_offset {
183 ranges.push((retained_offset, cold_frontier));
184 }
185 ranges.extend(
186 external_segments
187 .iter()
188 .map(|object| (object.start_offset, object.end_offset)),
189 );
190 ranges.sort_unstable();
191
192 let mut frontier = retained_offset;
193 for (start_offset, end_offset) in ranges {
194 if end_offset <= frontier {
195 continue;
196 }
197 if start_offset > frontier {
198 break;
199 }
200 frontier = end_offset;
201 }
202 frontier
203 }
204}
205
206#[derive(Debug, Clone, Default)]
207struct HotBuffer {
208 chunks: VecDeque<HotChunk>,
209}
210
211#[derive(Debug, Clone, PartialEq, Eq)]
212struct HotChunk {
213 start_offset: u64,
214 end_offset: u64,
215 bytes: Vec<u8>,
216}
217
218impl HotBuffer {
219 fn from_payload(start_offset: u64, payload: Vec<u8>) -> Self {
220 if payload.is_empty() {
221 return Self::default();
222 }
223 let end_offset = start_offset
224 .saturating_add(u64::try_from(payload.len()).expect("payload len fits u64"));
225 let mut chunks = VecDeque::new();
226 chunks.push_back(HotChunk {
227 start_offset,
228 end_offset,
229 bytes: payload,
230 });
231 Self { chunks }
232 }
233
234 fn from_snapshot(payload: Vec<u8>, segments: &[HotPayloadSegment]) -> Self {
235 let mut chunks = VecDeque::with_capacity(segments.len());
236 for segment in segments {
237 chunks.push_back(HotChunk {
238 start_offset: segment.start_offset,
239 end_offset: segment.end_offset,
240 bytes: payload[segment.payload_start..segment.payload_end].to_vec(),
241 });
242 }
243 Self { chunks }
244 }
245
246 fn len(&self) -> usize {
247 self.chunks.iter().map(|chunk| chunk.bytes.len()).sum()
248 }
249
250 fn hot_start_offset(&self) -> u64 {
251 self.chunks
252 .front()
253 .map(|chunk| chunk.start_offset)
254 .unwrap_or(0)
255 }
256
257 fn payload(&self) -> Vec<u8> {
258 let mut payload = Vec::with_capacity(self.len());
259 for chunk in &self.chunks {
260 payload.extend_from_slice(&chunk.bytes);
261 }
262 payload
263 }
264
265 fn hot_segments(&self) -> Vec<HotPayloadSegment> {
266 let mut payload_start = 0usize;
267 self.chunks
268 .iter()
269 .map(|chunk| {
270 let payload_end = payload_start + chunk.bytes.len();
271 let segment = HotPayloadSegment {
272 start_offset: chunk.start_offset,
273 end_offset: chunk.end_offset,
274 payload_start,
275 payload_end,
276 };
277 payload_start = payload_end;
278 segment
279 })
280 .collect()
281 }
282
283 fn push(&mut self, start_offset: u64, end_offset: u64, payload: &[u8]) {
284 if payload.is_empty() {
285 return;
286 }
287 self.chunks.push_back(HotChunk {
288 start_offset,
289 end_offset,
290 bytes: payload.to_vec(),
291 });
292 }
293
294 fn plan_cold_flush(
295 &self,
296 min_hot_bytes: usize,
297 max_flush_bytes: usize,
298 ) -> Option<(u64, u64, Vec<u8>)> {
299 let first = self.chunks.front()?;
300 let mut payload = Vec::new();
301 let mut end_offset = first.start_offset;
302 for chunk in &self.chunks {
303 if chunk.start_offset != end_offset || payload.len() >= max_flush_bytes {
304 break;
305 }
306 let remaining = max_flush_bytes - payload.len();
307 let take = chunk.bytes.len().min(remaining);
308 payload.extend_from_slice(&chunk.bytes[..take]);
309 end_offset = end_offset.saturating_add(u64::try_from(take).expect("take fits u64"));
310 if take < chunk.bytes.len() {
311 break;
312 }
313 }
314 if payload.len() < min_hot_bytes {
315 return None;
316 }
317 Some((first.start_offset, end_offset, payload))
318 }
319
320 fn read_segments(&self, offset: u64, next_offset: u64) -> Vec<(u64, StreamReadSegment)> {
321 let mut segments = Vec::new();
322 for chunk in &self.chunks {
323 let start = offset.max(chunk.start_offset);
324 let end = next_offset.min(chunk.end_offset);
325 if start < end {
326 let payload_start =
327 usize::try_from(start - chunk.start_offset).expect("hot start fits usize");
328 let payload_end =
329 usize::try_from(end - chunk.start_offset).expect("hot end fits usize");
330 segments.push((
331 start,
332 StreamReadSegment::Hot(chunk.bytes[payload_start..payload_end].to_vec()),
333 ));
334 }
335 }
336 segments
337 }
338
339 fn covers_prefix(&self, start_offset: u64, end_offset: u64) -> bool {
340 let Some(first) = self.chunks.front() else {
341 return false;
342 };
343 if first.start_offset != start_offset {
344 return false;
345 }
346 let mut covered_offset = start_offset;
347 for chunk in &self.chunks {
348 if chunk.start_offset != covered_offset {
349 return false;
350 }
351 if chunk.end_offset >= end_offset {
352 return true;
353 }
354 covered_offset = chunk.end_offset;
355 }
356 false
357 }
358
359 fn flush_prefix(&mut self, end_offset: u64) {
360 while self
361 .chunks
362 .front()
363 .is_some_and(|chunk| chunk.end_offset <= end_offset)
364 {
365 self.chunks.pop_front();
366 }
367 if let Some(front) = self.chunks.front_mut()
368 && front.start_offset < end_offset
369 {
370 let drain_len =
371 usize::try_from(end_offset - front.start_offset).expect("drain len fits usize");
372 front.bytes.drain(..drain_len);
373 front.start_offset = end_offset;
374 }
375 }
376
377 fn discard_before(&mut self, retained_offset: u64) {
378 self.flush_prefix(retained_offset);
379 }
380}
381
382impl StreamStateMachine {
383 pub fn new() -> Self {
384 Self::default()
385 }
386
387 pub fn apply(&mut self, command: StreamCommand) -> StreamResponse {
388 match command {
389 StreamCommand::CreateBucket { bucket_id } => self.create_bucket(bucket_id),
390 StreamCommand::DeleteBucket { bucket_id } => self.delete_bucket(&bucket_id),
391 StreamCommand::CreateStream {
392 stream_id,
393 content_type,
394 initial_payload,
395 close_after,
396 stream_seq,
397 producer,
398 stream_ttl_seconds,
399 stream_expires_at_ms,
400 forked_from,
401 fork_offset,
402 now_ms,
403 } => self.create_stream(CreateStreamInput {
404 stream_id,
405 content_type,
406 initial_payload,
407 close_after,
408 stream_seq,
409 producer,
410 stream_ttl_seconds,
411 stream_expires_at_ms,
412 forked_from,
413 fork_offset,
414 now_ms,
415 }),
416 StreamCommand::CreateExternal {
417 stream_id,
418 content_type,
419 initial_payload,
420 close_after,
421 stream_seq,
422 producer,
423 stream_ttl_seconds,
424 stream_expires_at_ms,
425 forked_from,
426 fork_offset,
427 now_ms,
428 } => self.create_external_stream(CreateExternalStreamInput {
429 stream_id,
430 content_type,
431 initial_payload,
432 close_after,
433 stream_seq,
434 producer,
435 stream_ttl_seconds,
436 stream_expires_at_ms,
437 forked_from,
438 fork_offset,
439 now_ms,
440 }),
441 StreamCommand::Append {
442 stream_id,
443 content_type,
444 payload,
445 close_after,
446 stream_seq,
447 producer,
448 now_ms,
449 } => self.append_borrowed(AppendStreamInput {
450 stream_id,
451 content_type: content_type.as_deref(),
452 payload: &payload,
453 close_after,
454 stream_seq,
455 producer,
456 now_ms,
457 }),
458 StreamCommand::AppendExternal {
459 stream_id,
460 content_type,
461 payload,
462 close_after,
463 stream_seq,
464 producer,
465 now_ms,
466 } => self.append_external(AppendExternalInput {
467 stream_id,
468 content_type: content_type.as_deref(),
469 payload,
470 close_after,
471 stream_seq,
472 producer,
473 now_ms,
474 }),
475 StreamCommand::AppendBatch {
476 stream_id,
477 content_type,
478 payloads,
479 producer,
480 now_ms,
481 } => match self.append_batch_borrowed(
482 stream_id,
483 content_type.as_deref(),
484 &payloads.iter().map(Vec::as_slice).collect::<Vec<_>>(),
485 producer,
486 now_ms,
487 ) {
488 Ok(batch) => batch
489 .items
490 .last()
491 .map(|item| StreamResponse::Appended {
492 offset: item.offset,
493 next_offset: item.next_offset,
494 closed: item.closed,
495 deduplicated: item.deduplicated,
496 producer: None,
497 })
498 .unwrap_or_else(|| {
499 StreamResponse::error(
500 StreamErrorCode::EmptyAppend,
501 "append batch must contain at least one payload",
502 )
503 }),
504 Err(response) => response,
505 },
506 StreamCommand::PublishSnapshot {
507 stream_id,
508 snapshot_offset,
509 content_type,
510 payload,
511 now_ms,
512 } => self.publish_snapshot(stream_id, snapshot_offset, content_type, payload, now_ms),
513 StreamCommand::TouchStreamAccess {
514 stream_id,
515 now_ms,
516 renew_ttl,
517 } => self.touch_stream_access(&stream_id, now_ms, renew_ttl),
518 StreamCommand::AddForkRef { stream_id, now_ms } => {
519 self.add_fork_ref(&stream_id, now_ms)
520 }
521 StreamCommand::ReleaseForkRef { stream_id } => self.release_fork_ref(&stream_id),
522 StreamCommand::FlushCold { stream_id, chunk } => self.flush_cold(stream_id, chunk),
523 StreamCommand::Close {
524 stream_id,
525 stream_seq,
526 producer,
527 now_ms,
528 } => self.close(stream_id, stream_seq, producer, now_ms),
529 StreamCommand::DeleteStream { stream_id } => self.delete_stream(&stream_id),
530 StreamCommand::AckColdGc { up_to_seq } => self.ack_cold_gc(up_to_seq),
531 }
532 }
533
534 pub fn head(&self, stream_id: &BucketStreamId) -> Option<&StreamMetadata> {
535 self.streams.get(stream_id)
536 }
537
538 pub fn head_at(&mut self, stream_id: &BucketStreamId, now_ms: u64) -> Option<&StreamMetadata> {
539 self.expire_stream_if_due(stream_id, now_ms);
540 self.streams.get(stream_id)
541 }
542
543 pub fn access_requires_write(
544 &self,
545 stream_id: &BucketStreamId,
546 now_ms: u64,
547 renew_ttl: bool,
548 ) -> Result<bool, StreamResponse> {
549 self.validate_stream_scope(stream_id)?;
550 let Some(stream) = self.streams.get(stream_id) else {
551 return Err(StreamResponse::error(
552 StreamErrorCode::StreamNotFound,
553 format!("stream '{stream_id}' does not exist"),
554 ));
555 };
556 if is_soft_deleted(stream) {
557 return Err(StreamResponse::error(
558 StreamErrorCode::StreamGone,
559 format!("stream '{stream_id}' is gone"),
560 ));
561 }
562 if stream_is_expired(stream, now_ms) {
563 return Ok(true);
564 }
565 Ok(renew_ttl
566 && stream.stream_ttl_seconds.is_some()
567 && stream.last_ttl_touch_at_ms != now_ms)
568 }
569
570 pub fn hot_start_offset(&self, stream_id: &BucketStreamId) -> u64 {
571 let Some(stream) = self.streams.get(stream_id) else {
572 return 0;
573 };
574 self.hot_buffers
575 .get(stream_id)
576 .and_then(|buffer| buffer.chunks.front().map(|chunk| chunk.start_offset))
577 .unwrap_or(stream.tail_offset)
578 }
579
580 pub fn cold_chunks(&self, stream_id: &BucketStreamId) -> &[ColdChunkRef] {
581 self.cold_index().cold_chunks(stream_id)
582 }
583
584 pub fn external_segments(&self, stream_id: &BucketStreamId) -> &[ObjectPayloadRef] {
585 self.cold_index().external_segments(stream_id)
586 }
587
588 pub fn hot_segments(&self, stream_id: &BucketStreamId) -> Vec<HotPayloadSegment> {
589 self.hot_buffers
590 .get(stream_id)
591 .map(HotBuffer::hot_segments)
592 .unwrap_or_default()
593 }
594
595 pub fn hot_payload_len(&self, stream_id: &BucketStreamId) -> Result<u64, StreamResponse> {
596 let Some(stream) = self.streams.get(stream_id) else {
597 return Err(StreamResponse::error(
598 StreamErrorCode::StreamNotFound,
599 format!("stream '{stream_id}' does not exist"),
600 ));
601 };
602 if is_soft_deleted(stream) {
603 return Err(StreamResponse::error(
604 StreamErrorCode::StreamGone,
605 format!("stream '{stream_id}' is gone"),
606 ));
607 }
608 let payload = self
609 .hot_buffers
610 .get(stream_id)
611 .expect("hot buffer exists for stream metadata");
612 Ok(u64::try_from(payload.len()).expect("payload len fits u64"))
613 }
614
615 pub fn total_hot_payload_bytes(&self) -> u64 {
616 self.hot_buffers
617 .values()
618 .map(|payload| u64::try_from(payload.len()).expect("payload len fits u64"))
619 .sum()
620 }
621
622 pub fn plan_cold_flush(
623 &self,
624 stream_id: &BucketStreamId,
625 min_hot_bytes: usize,
626 max_flush_bytes: usize,
627 ) -> Result<Option<ColdFlushCandidate>, StreamResponse> {
628 if max_flush_bytes == 0 {
629 return Ok(None);
630 }
631 let Some(stream) = self.streams.get(stream_id) else {
632 return Err(StreamResponse::error(
633 StreamErrorCode::StreamNotFound,
634 format!("stream '{stream_id}' does not exist"),
635 ));
636 };
637 if is_soft_deleted(stream) {
638 return Err(StreamResponse::error(
639 StreamErrorCode::StreamGone,
640 format!("stream '{stream_id}' is gone"),
641 ));
642 }
643 let Some(hot_buffer) = self.hot_buffers.get(stream_id) else {
644 return Ok(None);
645 };
646 let Some((start_offset, end_offset, payload)) =
647 hot_buffer.plan_cold_flush(min_hot_bytes, max_flush_bytes)
648 else {
649 return Ok(None);
650 };
651 Ok(Some(ColdFlushCandidate {
652 stream_id: stream_id.clone(),
653 start_offset,
654 end_offset,
655 payload,
656 }))
657 }
658
659 pub fn plan_next_cold_flush(
660 &self,
661 min_hot_bytes: usize,
662 max_flush_bytes: usize,
663 ) -> Result<Option<ColdFlushCandidate>, StreamResponse> {
664 if max_flush_bytes == 0 {
665 return Ok(None);
666 }
667 let mut stream_ids = self.streams.keys().cloned().collect::<Vec<_>>();
668 stream_ids.sort_by(compare_stream_ids);
669 for stream_id in &stream_ids {
670 match self.plan_cold_flush(stream_id, min_hot_bytes, max_flush_bytes) {
671 Ok(Some(candidate)) => return Ok(Some(candidate)),
672 Ok(None) => {}
673 Err(StreamResponse::Error {
674 code: StreamErrorCode::StreamGone | StreamErrorCode::StreamNotFound,
675 ..
676 }) => {}
677 Err(err) => return Err(err),
678 }
679 }
680 let group_min_hot_bytes = u64::try_from(min_hot_bytes).unwrap_or(u64::MAX);
681 let group_hot_bytes = self.total_hot_payload_bytes();
682 if group_hot_bytes < group_min_hot_bytes {
683 return Ok(None);
684 }
685 for stream_id in stream_ids {
686 match self.plan_cold_flush(&stream_id, 1, max_flush_bytes) {
687 Ok(Some(candidate)) => return Ok(Some(candidate)),
688 Ok(None) => {}
689 Err(StreamResponse::Error {
690 code: StreamErrorCode::StreamGone | StreamErrorCode::StreamNotFound,
691 ..
692 }) => {}
693 Err(err) => return Err(err),
694 }
695 }
696 Ok(None)
697 }
698
699 pub fn plan_next_cold_flush_batch(
700 &self,
701 min_hot_bytes: usize,
702 max_flush_bytes: usize,
703 max_candidates: usize,
704 ) -> Result<Vec<ColdFlushCandidate>, StreamResponse> {
705 if max_candidates == 0 || max_flush_bytes == 0 {
706 return Ok(Vec::new());
707 }
708 let mut preview = self.clone();
709 let mut candidates = Vec::with_capacity(max_candidates);
710 while candidates.len() < max_candidates {
711 let Some(candidate) = preview.plan_next_cold_flush(min_hot_bytes, max_flush_bytes)?
712 else {
713 break;
714 };
715 let chunk = ColdChunkRef {
716 start_offset: candidate.start_offset,
717 end_offset: candidate.end_offset,
718 s3_path: "planned-cold-flush-batch".to_owned(),
719 object_size: u64::try_from(candidate.payload.len()).expect("payload len fits u64"),
720 };
721 match preview.flush_cold(candidate.stream_id.clone(), chunk) {
722 StreamResponse::ColdFlushed { .. } => candidates.push(candidate),
723 StreamResponse::Error { .. } => break,
724 other => {
725 return Err(StreamResponse::error(
726 StreamErrorCode::InvalidColdFlush,
727 format!("unexpected cold flush planning response: {other:?}"),
728 ));
729 }
730 }
731 }
732 Ok(candidates)
733 }
734
735 pub fn bucket_exists(&self, bucket_id: &str) -> bool {
736 self.buckets.contains(bucket_id)
737 }
738
739 fn cold_index(&self) -> &impl ColdIndexState {
740 &self.cold_index
741 }
742
743 fn cold_index_mut(&mut self) -> &mut impl ColdIndexState {
744 &mut self.cold_index
745 }
746
747 pub fn integrity_snapshot(
748 &self,
749 stream_id: &BucketStreamId,
750 ) -> Result<crate::integrity::StreamIntegritySnapshot, StreamResponse> {
751 let Some(stream) = self.streams.get(stream_id) else {
752 return Err(StreamResponse::error(
753 StreamErrorCode::StreamNotFound,
754 format!("stream '{stream_id}' does not exist"),
755 ));
756 };
757 if is_soft_deleted(stream) {
758 return Err(StreamResponse::error(
759 StreamErrorCode::StreamGone,
760 format!("stream '{stream_id}' is gone"),
761 ));
762 }
763 Ok(self
764 .integrities
765 .get(stream_id)
766 .expect("integrity exists for stream metadata")
767 .snapshot(self.earliest_retained_offset(stream_id), stream.tail_offset))
768 }
769
770 pub fn snapshot(&self) -> StreamSnapshot {
771 let mut buckets = self.buckets.iter().cloned().collect::<Vec<_>>();
772 buckets.sort();
773
774 let mut streams = self
775 .streams
776 .values()
777 .cloned()
778 .map(|metadata| {
779 let stream_id = metadata.stream_id.clone();
780 let tail_offset = metadata.tail_offset;
781 let hot_buffer = self
782 .hot_buffers
783 .get(&stream_id)
784 .expect("hot buffer exists for stream metadata");
785 let payload = hot_buffer.payload();
786 let producer_states = self.producer_snapshot(&stream_id);
787 StreamSnapshotEntry {
788 metadata,
789 hot_start_offset: self.hot_start_offset(&stream_id),
790 payload,
791 hot_segments: hot_buffer.hot_segments(),
792 cold_frontier_offset: self.cold_frontier_offset(
793 &stream_id,
794 self.earliest_retained_offset(&stream_id),
795 ),
796 cold_index_generation: self.cold_index().cold_generation(&stream_id),
797 cold_chunks: self.cold_index().cold_chunks(&stream_id).to_vec(),
798 external_segments: self.cold_index().external_segments(&stream_id).to_vec(),
799 message_records: self
800 .message_records
801 .get(&stream_id)
802 .cloned()
803 .unwrap_or_default(),
804 integrity: self
805 .integrities
806 .get(&stream_id)
807 .expect("integrity exists for stream metadata")
808 .snapshot(self.earliest_retained_offset(&stream_id), tail_offset),
809 visible_snapshot: self.visible_snapshots.get(&stream_id).cloned(),
810 producer_states,
811 }
812 })
813 .collect::<Vec<_>>();
814 streams.sort_by(|left, right| {
815 compare_stream_ids(&left.metadata.stream_id, &right.metadata.stream_id)
816 });
817
818 StreamSnapshot {
819 buckets,
820 streams,
821 pending_cold_gc: self.pending_cold_gc.iter().cloned().collect(),
822 next_cold_gc_seq: self.next_cold_gc_seq,
823 }
824 }
825
826 pub fn restore(snapshot: StreamSnapshot) -> Result<Self, StreamSnapshotError> {
827 let mut machine = Self::default();
828 for bucket_id in snapshot.buckets {
829 if !machine.buckets.insert(bucket_id.clone()) {
830 return Err(StreamSnapshotError::DuplicateBucket(bucket_id));
831 }
832 }
833
834 for entry in snapshot.streams {
835 let stream_id = entry.metadata.stream_id.clone();
836 if !machine.buckets.contains(&stream_id.bucket_id) {
837 return Err(StreamSnapshotError::MissingBucket(stream_id));
838 }
839 if let Some(snapshot) = entry.visible_snapshot.as_ref()
840 && snapshot.offset > entry.metadata.tail_offset
841 {
842 return Err(StreamSnapshotError::SnapshotOffsetOutOfRange {
843 stream_id,
844 snapshot_offset: snapshot.offset,
845 tail_offset: entry.metadata.tail_offset,
846 });
847 }
848 let retained_offset = entry
849 .visible_snapshot
850 .as_ref()
851 .map(|snapshot| snapshot.offset)
852 .unwrap_or(0);
853 let hot_segments = if entry.hot_segments.is_empty() && !entry.payload.is_empty() {
854 vec![HotPayloadSegment {
855 start_offset: entry.hot_start_offset,
856 end_offset: entry.metadata.tail_offset,
857 payload_start: 0,
858 payload_end: entry.payload.len(),
859 }]
860 } else {
861 entry.hot_segments
862 };
863 if !hot_segments_match_payload(&hot_segments, entry.payload.len())
864 || !payload_sources_cover_retained_suffix(
865 entry.cold_frontier_offset,
866 &entry.cold_chunks,
867 &entry.external_segments,
868 &hot_segments,
869 retained_offset,
870 entry.metadata.tail_offset,
871 )
872 {
873 return Err(StreamSnapshotError::PayloadLengthMismatch {
874 stream_id,
875 tail_offset: entry.metadata.tail_offset,
876 payload_len: entry.payload.len(),
877 });
878 }
879 if !message_records_cover_retained_suffix(
880 &entry.message_records,
881 retained_offset,
882 entry.metadata.tail_offset,
883 ) {
884 return Err(StreamSnapshotError::MessageBoundaryMismatch { stream_id });
885 }
886 let integrity = StreamIntegrity::restore(entry.integrity).ok_or_else(|| {
887 StreamSnapshotError::IntegrityMismatch {
888 stream_id: stream_id.clone(),
889 }
890 })?;
891 if machine
892 .streams
893 .insert(entry.metadata.stream_id.clone(), entry.metadata)
894 .is_some()
895 {
896 return Err(StreamSnapshotError::DuplicateStream(stream_id));
897 }
898 let producer_states = restore_producer_states(&stream_id, entry.producer_states)?;
899 machine.hot_buffers.insert(
900 stream_id.clone(),
901 HotBuffer::from_snapshot(entry.payload, &hot_segments),
902 );
903 machine.cold_index_mut().restore_stream(
904 stream_id.clone(),
905 entry.cold_frontier_offset,
906 entry.cold_index_generation,
907 entry.cold_chunks,
908 entry.external_segments,
909 );
910 if !entry.message_records.is_empty() {
911 machine
912 .message_records
913 .insert(stream_id.clone(), entry.message_records);
914 }
915 machine.integrities.insert(stream_id.clone(), integrity);
916 if let Some(snapshot) = entry.visible_snapshot {
917 machine
918 .visible_snapshots
919 .insert(stream_id.clone(), snapshot);
920 }
921 machine.producers.insert(stream_id.clone(), producer_states);
922 }
923
924 machine.pending_cold_gc = snapshot.pending_cold_gc.into_iter().collect();
925 machine.next_cold_gc_seq = snapshot.next_cold_gc_seq;
926
927 Ok(machine)
928 }
929
930 pub fn read(
931 &self,
932 stream_id: &BucketStreamId,
933 offset: u64,
934 max_len: usize,
935 ) -> Result<StreamRead, StreamResponse> {
936 let plan = self.read_plan(stream_id, offset, max_len)?;
937 if plan.segments.iter().any(|segment| {
938 matches!(
939 segment,
940 StreamReadSegment::ColdIndex(_) | StreamReadSegment::Object(_)
941 )
942 }) {
943 return Err(StreamResponse::error_with_next_offset(
944 StreamErrorCode::InvalidColdFlush,
945 format!("stream '{stream_id}' read requires object payload store"),
946 plan.next_offset,
947 ));
948 }
949 let payload = plan
950 .segments
951 .iter()
952 .flat_map(|segment| match segment {
953 StreamReadSegment::Hot(payload) => payload.as_slice(),
954 StreamReadSegment::ColdIndex(_) | StreamReadSegment::Object(_) => {
955 unreachable!("object segments checked above")
956 }
957 })
958 .copied()
959 .collect();
960 Ok(StreamRead {
961 offset: plan.offset,
962 next_offset: plan.next_offset,
963 content_type: plan.content_type,
964 payload,
965 up_to_date: plan.up_to_date,
966 closed: plan.closed,
967 })
968 }
969
970 pub fn read_plan(
971 &self,
972 stream_id: &BucketStreamId,
973 offset: u64,
974 max_len: usize,
975 ) -> Result<StreamReadPlan, StreamResponse> {
976 self.read_plan_at(stream_id, offset, max_len, 0)
977 }
978
979 pub fn read_plan_at(
980 &self,
981 stream_id: &BucketStreamId,
982 offset: u64,
983 max_len: usize,
984 now_ms: u64,
985 ) -> Result<StreamReadPlan, StreamResponse> {
986 let Some(stream) = self.streams.get(stream_id) else {
987 return Err(StreamResponse::error(
988 StreamErrorCode::StreamNotFound,
989 format!("stream '{stream_id}' does not exist"),
990 ));
991 };
992 if is_soft_deleted(stream) {
993 return Err(StreamResponse::error(
994 StreamErrorCode::StreamGone,
995 format!("stream '{stream_id}' is gone"),
996 ));
997 }
998 if stream_is_expired(stream, now_ms) {
999 return Err(StreamResponse::error(
1000 StreamErrorCode::StreamNotFound,
1001 format!("stream '{stream_id}' does not exist"),
1002 ));
1003 }
1004 if offset > stream.tail_offset {
1005 return Err(StreamResponse::error_with_next_offset(
1006 StreamErrorCode::OffsetOutOfRange,
1007 format!(
1008 "offset {offset} is beyond stream '{}' tail {}",
1009 stream_id, stream.tail_offset
1010 ),
1011 stream.tail_offset,
1012 ));
1013 }
1014 let retained_offset = self.earliest_retained_offset(stream_id);
1015 if offset < retained_offset {
1016 return Err(StreamResponse::error_with_next_offset(
1017 StreamErrorCode::StreamGone,
1018 format!(
1019 "offset {offset} is older than stream '{}' retained offset {retained_offset}",
1020 stream_id
1021 ),
1022 retained_offset,
1023 ));
1024 }
1025
1026 let max_len_u64 = u64::try_from(max_len).unwrap_or(u64::MAX);
1027 let next_offset = stream.tail_offset.min(offset.saturating_add(max_len_u64));
1028 let mut segments = Vec::<(u64, StreamReadSegment)>::new();
1029 let hot_segments = self
1030 .hot_buffers
1031 .get(stream_id)
1032 .map(|hot_buffer| hot_buffer.read_segments(offset, next_offset))
1033 .unwrap_or_default();
1034 let cold_frontier = self.cold_frontier_offset(stream_id, retained_offset);
1035 let cold_index_end = next_offset.min(cold_frontier);
1036 let mut cursor = offset;
1037 for (hot_start, hot_segment) in &hot_segments {
1038 if cursor >= cold_index_end {
1039 break;
1040 }
1041 let Some(hot_end) = read_segment_end(*hot_start, hot_segment) else {
1042 continue;
1043 };
1044 if hot_end <= cursor {
1045 continue;
1046 }
1047 let gap_end = (*hot_start).min(cold_index_end);
1048 push_cold_index_segments(
1049 &mut segments,
1050 stream_id,
1051 self.cold_index().cold_generation(stream_id),
1052 cursor,
1053 gap_end,
1054 );
1055 cursor = cursor.max(hot_end);
1056 }
1057 push_cold_index_segments(
1058 &mut segments,
1059 stream_id,
1060 self.cold_index().cold_generation(stream_id),
1061 cursor,
1062 cold_index_end,
1063 );
1064 for chunk in self.cold_chunks(stream_id) {
1065 let start = offset.max(chunk.start_offset);
1066 let end = next_offset.min(chunk.end_offset);
1067 if start < end {
1068 segments.push((
1069 start,
1070 StreamReadSegment::Object(StreamReadObjectSegment {
1071 object: ObjectPayloadRef::from(chunk),
1072 read_start_offset: start,
1073 len: usize::try_from(end - start).expect("object read len fits usize"),
1074 }),
1075 ));
1076 }
1077 }
1078 for object in self.external_segments(stream_id) {
1079 let start = offset.max(object.start_offset);
1080 let end = next_offset.min(object.end_offset);
1081 if start < end {
1082 segments.push((
1083 start,
1084 StreamReadSegment::Object(StreamReadObjectSegment {
1085 object: object.clone(),
1086 read_start_offset: start,
1087 len: usize::try_from(end - start).expect("object read len fits usize"),
1088 }),
1089 ));
1090 }
1091 }
1092 segments.extend(hot_segments);
1093 segments.sort_by_key(|(start, _)| *start);
1094 if !segments_cover_range(&segments, offset, next_offset) {
1095 return Err(StreamResponse::error_with_next_offset(
1096 StreamErrorCode::InvalidColdFlush,
1097 format!("stream '{stream_id}' has missing payload segment metadata"),
1098 next_offset,
1099 ));
1100 }
1101 Ok(StreamReadPlan {
1102 offset,
1103 next_offset,
1104 content_type: stream.content_type.clone(),
1105 segments: segments.into_iter().map(|(_, segment)| segment).collect(),
1106 up_to_date: next_offset == stream.tail_offset,
1107 closed: stream.status == StreamStatus::Closed,
1108 })
1109 }
1110
1111 pub fn latest_snapshot(
1112 &self,
1113 stream_id: &BucketStreamId,
1114 ) -> Result<Option<StreamVisibleSnapshot>, StreamResponse> {
1115 let Some(stream) = self.streams.get(stream_id) else {
1116 return Err(StreamResponse::error(
1117 StreamErrorCode::StreamNotFound,
1118 format!("stream '{stream_id}' does not exist"),
1119 ));
1120 };
1121 if is_soft_deleted(stream) {
1122 return Err(StreamResponse::error(
1123 StreamErrorCode::StreamGone,
1124 format!("stream '{stream_id}' is gone"),
1125 ));
1126 }
1127 Ok(self.visible_snapshots.get(stream_id).cloned())
1128 }
1129
1130 pub fn read_snapshot(
1131 &self,
1132 stream_id: &BucketStreamId,
1133 snapshot_offset: u64,
1134 ) -> Result<StreamVisibleSnapshot, StreamResponse> {
1135 let snapshot = self.latest_snapshot(stream_id)?;
1136 match snapshot {
1137 Some(snapshot) if snapshot.offset == snapshot_offset => Ok(snapshot),
1138 _ => Err(StreamResponse::error(
1139 StreamErrorCode::SnapshotNotFound,
1140 format!("snapshot {snapshot_offset} for stream '{stream_id}' does not exist"),
1141 )),
1142 }
1143 }
1144
1145 pub fn delete_snapshot(
1146 &self,
1147 stream_id: &BucketStreamId,
1148 snapshot_offset: u64,
1149 ) -> StreamResponse {
1150 match self.latest_snapshot(stream_id) {
1151 Ok(Some(snapshot)) if snapshot.offset == snapshot_offset => StreamResponse::error(
1152 StreamErrorCode::SnapshotConflict,
1153 format!(
1154 "snapshot {snapshot_offset} for stream '{stream_id}' is the latest visible snapshot"
1155 ),
1156 ),
1157 Ok(_) => StreamResponse::error(
1158 StreamErrorCode::SnapshotNotFound,
1159 format!("snapshot {snapshot_offset} for stream '{stream_id}' does not exist"),
1160 ),
1161 Err(err) => err,
1162 }
1163 }
1164
1165 pub fn bootstrap_plan(
1166 &self,
1167 stream_id: &BucketStreamId,
1168 ) -> Result<StreamBootstrapPlan, StreamResponse> {
1169 let Some(stream) = self.streams.get(stream_id) else {
1170 return Err(StreamResponse::error(
1171 StreamErrorCode::StreamNotFound,
1172 format!("stream '{stream_id}' does not exist"),
1173 ));
1174 };
1175 if is_soft_deleted(stream) {
1176 return Err(StreamResponse::error(
1177 StreamErrorCode::StreamGone,
1178 format!("stream '{stream_id}' is gone"),
1179 ));
1180 }
1181 let snapshot = self.visible_snapshots.get(stream_id).cloned();
1182 let retained_offset = snapshot
1183 .as_ref()
1184 .map(|snapshot| snapshot.offset)
1185 .unwrap_or(0);
1186 let updates = self
1187 .message_records
1188 .get(stream_id)
1189 .map(|records| {
1190 records
1191 .iter()
1192 .filter(|record| record.start_offset >= retained_offset)
1193 .cloned()
1194 .collect::<Vec<_>>()
1195 })
1196 .unwrap_or_default();
1197 Ok(StreamBootstrapPlan {
1198 snapshot,
1199 updates,
1200 next_offset: stream.tail_offset,
1201 content_type: stream.content_type.clone(),
1202 up_to_date: true,
1203 closed: stream.status == StreamStatus::Closed,
1204 })
1205 }
1206
1207 fn publish_snapshot(
1208 &mut self,
1209 stream_id: BucketStreamId,
1210 snapshot_offset: u64,
1211 content_type: String,
1212 payload: Vec<u8>,
1213 now_ms: u64,
1214 ) -> StreamResponse {
1215 if let Err(response) = self.validate_stream_scope(&stream_id) {
1216 return response;
1217 }
1218 if content_type.trim().is_empty() {
1219 return StreamResponse::error(
1220 StreamErrorCode::InvalidSnapshot,
1221 "snapshot content type must not be empty",
1222 );
1223 }
1224 let Some(stream) = self.streams.get(&stream_id) else {
1225 return StreamResponse::error(
1226 StreamErrorCode::StreamNotFound,
1227 format!("stream '{stream_id}' does not exist"),
1228 );
1229 };
1230 if is_soft_deleted(stream) {
1231 return StreamResponse::error(
1232 StreamErrorCode::StreamGone,
1233 format!("stream '{stream_id}' is gone"),
1234 );
1235 }
1236 if stream_is_expired(stream, now_ms) {
1237 self.remove_stream_state(&stream_id);
1238 return StreamResponse::error(
1239 StreamErrorCode::StreamNotFound,
1240 format!("stream '{stream_id}' does not exist"),
1241 );
1242 }
1243 let tail_offset = stream.tail_offset;
1244 let retained_offset = self.earliest_retained_offset(&stream_id);
1245 if snapshot_offset < retained_offset {
1246 return StreamResponse::error_with_next_offset(
1247 StreamErrorCode::StreamGone,
1248 format!(
1249 "snapshot offset {snapshot_offset} is older than stream '{}' retained offset {retained_offset}",
1250 stream_id
1251 ),
1252 retained_offset,
1253 );
1254 }
1255 if snapshot_offset > tail_offset {
1256 return StreamResponse::error_with_next_offset(
1257 StreamErrorCode::SnapshotConflict,
1258 format!(
1259 "snapshot offset {snapshot_offset} is beyond stream '{}' tail {tail_offset}",
1260 stream_id
1261 ),
1262 tail_offset,
1263 );
1264 }
1265 if !self.snapshot_offset_aligned(&stream_id, snapshot_offset, retained_offset) {
1266 return StreamResponse::error_with_next_offset(
1267 StreamErrorCode::InvalidSnapshot,
1268 format!(
1269 "snapshot offset {snapshot_offset} is not aligned to a committed message boundary for stream '{stream_id}'"
1270 ),
1271 tail_offset,
1272 );
1273 }
1274
1275 self.visible_snapshots
1276 .insert(stream_id.clone(), StreamVisibleSnapshot {
1277 offset: snapshot_offset,
1278 content_type,
1279 payload,
1280 });
1281 self.compact_retained_prefix(&stream_id, snapshot_offset);
1282 StreamResponse::SnapshotPublished { snapshot_offset }
1283 }
1284
1285 fn flush_cold(&mut self, stream_id: BucketStreamId, chunk: ColdChunkRef) -> StreamResponse {
1286 if let Err(response) = self.validate_stream_scope(&stream_id) {
1287 return response;
1288 }
1289 if chunk.s3_path.trim().is_empty() {
1290 return StreamResponse::error(
1291 StreamErrorCode::InvalidColdFlush,
1292 "cold chunk S3 path must not be empty",
1293 );
1294 }
1295 if chunk.object_size == 0 {
1296 return StreamResponse::error(
1297 StreamErrorCode::InvalidColdFlush,
1298 "cold chunk object size must be greater than zero",
1299 );
1300 }
1301 let Some(stream) = self.streams.get(&stream_id) else {
1302 return StreamResponse::error(
1303 StreamErrorCode::StreamNotFound,
1304 format!("stream '{stream_id}' does not exist"),
1305 );
1306 };
1307 if is_soft_deleted(stream) {
1308 return StreamResponse::error(
1309 StreamErrorCode::StreamGone,
1310 format!("stream '{stream_id}' is gone"),
1311 );
1312 }
1313 if chunk.end_offset <= chunk.start_offset {
1314 return StreamResponse::error_with_next_offset(
1315 StreamErrorCode::InvalidColdFlush,
1316 "cold chunk must cover at least one byte",
1317 stream.tail_offset,
1318 );
1319 }
1320 if chunk.end_offset > stream.tail_offset {
1321 return StreamResponse::error_with_next_offset_and_context(
1322 StreamErrorCode::InvalidColdFlush,
1323 format!(
1324 "cold chunk end {} is beyond stream '{}' tail {}",
1325 chunk.end_offset, stream_id, stream.tail_offset
1326 ),
1327 stream.tail_offset,
1328 vec![StreamErrorContext::StaleColdFlushCandidate],
1329 );
1330 }
1331 let Some(hot_buffer) = self.hot_buffers.get(&stream_id) else {
1332 return StreamResponse::error_with_next_offset_and_context(
1333 StreamErrorCode::InvalidColdFlush,
1334 format!("cold chunk for stream '{stream_id}' does not match hot payload"),
1335 stream.tail_offset,
1336 vec![StreamErrorContext::StaleColdFlushCandidate],
1337 );
1338 };
1339 if hot_buffer.hot_start_offset() != chunk.start_offset {
1340 return StreamResponse::error_with_next_offset_and_context(
1341 StreamErrorCode::InvalidColdFlush,
1342 format!("cold chunk for stream '{stream_id}' must start at the hot prefix"),
1343 stream.tail_offset,
1344 vec![StreamErrorContext::StaleColdFlushCandidate],
1345 );
1346 }
1347 if !hot_buffer.covers_prefix(chunk.start_offset, chunk.end_offset) {
1348 return StreamResponse::error_with_next_offset_and_context(
1349 StreamErrorCode::InvalidColdFlush,
1350 format!(
1351 "cold chunk for stream '{stream_id}' does not cover contiguous hot payload"
1352 ),
1353 stream.tail_offset,
1354 vec![StreamErrorContext::StaleColdFlushCandidate],
1355 );
1356 }
1357 self.hot_buffers
1358 .get_mut(&stream_id)
1359 .expect("hot buffer exists for stream metadata")
1360 .flush_prefix(chunk.end_offset);
1361 self.cold_index_mut()
1362 .push_cold_chunk(stream_id.clone(), chunk.clone());
1363 self.compact_message_records_before(
1364 &stream_id,
1365 self.earliest_retained_offset(&stream_id),
1366 chunk.end_offset,
1367 );
1368 StreamResponse::ColdFlushed {
1369 hot_start_offset: self.hot_start_offset(&stream_id),
1370 }
1371 }
1372
1373 fn create_bucket(&mut self, bucket_id: String) -> StreamResponse {
1374 if let Err(message) = validate_bucket_id(&bucket_id) {
1375 return StreamResponse::error(StreamErrorCode::InvalidBucketId, message);
1376 }
1377 if !self.buckets.insert(bucket_id.clone()) {
1378 return StreamResponse::BucketAlreadyExists { bucket_id };
1379 }
1380 StreamResponse::BucketCreated { bucket_id }
1381 }
1382
1383 fn delete_bucket(&mut self, bucket_id: &str) -> StreamResponse {
1384 if let Err(message) = validate_bucket_id(bucket_id) {
1385 return StreamResponse::error(StreamErrorCode::InvalidBucketId, message);
1386 }
1387 if !self.buckets.contains(bucket_id) {
1388 return StreamResponse::error(
1389 StreamErrorCode::BucketNotFound,
1390 format!("bucket '{bucket_id}' does not exist"),
1391 );
1392 }
1393 if self
1394 .streams
1395 .keys()
1396 .any(|stream_id| stream_id.bucket_id == bucket_id)
1397 {
1398 return StreamResponse::error(
1399 StreamErrorCode::BucketNotEmpty,
1400 format!("bucket '{bucket_id}' is not empty"),
1401 );
1402 }
1403 self.buckets.remove(bucket_id);
1404 StreamResponse::BucketDeleted {
1405 bucket_id: bucket_id.to_owned(),
1406 }
1407 }
1408
1409 fn create_stream(&mut self, input: CreateStreamInput) -> StreamResponse {
1410 if let Err(response) = self.validate_stream_scope(&input.stream_id) {
1411 return response;
1412 }
1413 if let Err(response) =
1414 validate_retention(input.stream_ttl_seconds, input.stream_expires_at_ms)
1415 {
1416 return response;
1417 }
1418 if let Err(response) = validate_producer_request(input.producer.as_ref()) {
1419 return response;
1420 }
1421 if let Some(producer) = input.producer.as_ref()
1422 && producer.producer_seq != 0
1423 {
1424 return StreamResponse::error_with_context(
1425 StreamErrorCode::ProducerSeqConflict,
1426 format!(
1427 "producer '{}' expected sequence 0, received {}",
1428 producer.producer_id, producer.producer_seq
1429 ),
1430 vec![StreamErrorContext::ProducerSeqConflict {
1431 expected_seq: 0,
1432 received_seq: producer.producer_seq,
1433 }],
1434 );
1435 }
1436 if self
1437 .streams
1438 .get(&input.stream_id)
1439 .is_some_and(|existing| stream_is_expired(existing, input.now_ms))
1440 {
1441 self.remove_stream_state(&input.stream_id);
1442 }
1443
1444 if let Some(existing) = self.streams.get(&input.stream_id) {
1445 if is_soft_deleted(existing) {
1446 return StreamResponse::error(
1447 StreamErrorCode::StreamAlreadyExistsConflict,
1448 format!(
1449 "stream '{}' is gone and cannot be recreated yet",
1450 input.stream_id
1451 ),
1452 );
1453 }
1454 if existing.content_type == input.content_type
1455 && existing.status == status_from_closed(input.close_after)
1456 && existing.stream_ttl_seconds == input.stream_ttl_seconds
1457 && existing.stream_expires_at_ms == input.stream_expires_at_ms
1458 && existing.forked_from == input.forked_from
1459 && existing.fork_offset == input.fork_offset
1460 {
1461 return StreamResponse::AlreadyExists {
1462 next_offset: existing.tail_offset,
1463 closed: existing.status == StreamStatus::Closed,
1464 content_type: existing.content_type.clone(),
1465 stream_ttl_seconds: existing.stream_ttl_seconds,
1466 stream_expires_at_ms: existing.stream_expires_at_ms,
1467 };
1468 }
1469 return StreamResponse::error(
1470 StreamErrorCode::StreamAlreadyExistsConflict,
1471 format!(
1472 "stream '{}' already exists with different metadata",
1473 input.stream_id
1474 ),
1475 );
1476 }
1477
1478 let initial_len = input.initial_len();
1479 let metadata = StreamMetadata {
1480 stream_id: input.stream_id.clone(),
1481 content_type: input.content_type,
1482 status: status_from_closed(input.close_after),
1483 tail_offset: initial_len,
1484 last_stream_seq: input.stream_seq,
1485 stream_ttl_seconds: input.stream_ttl_seconds,
1486 stream_expires_at_ms: input.stream_expires_at_ms,
1487 created_at_ms: input.now_ms,
1488 last_ttl_touch_at_ms: input.now_ms,
1489 forked_from: input.forked_from,
1490 fork_offset: input.fork_offset,
1491 fork_ref_count: 0,
1492 };
1493 self.streams.insert(input.stream_id.clone(), metadata);
1494 self.hot_buffers.insert(
1495 input.stream_id.clone(),
1496 HotBuffer::from_payload(0, input.initial_payload),
1497 );
1498 let mut integrity = StreamIntegrity::default();
1499 if initial_len > 0 {
1500 let payload = self
1501 .hot_buffers
1502 .get(&input.stream_id)
1503 .expect("hot buffer exists for stream metadata")
1504 .payload();
1505 integrity.append_payload(&input.stream_id, 0, initial_len, &payload);
1506 }
1507 self.integrities.insert(input.stream_id.clone(), integrity);
1508 if initial_len > 0 {
1509 self.message_records
1510 .insert(input.stream_id.clone(), vec![StreamMessageRecord {
1511 start_offset: 0,
1512 end_offset: initial_len,
1513 }]);
1514 }
1515 let mut producer_states = HashMap::new();
1516 if let Some(producer) = input.producer {
1517 let last_item = ProducerAppendRecord {
1518 start_offset: 0,
1519 next_offset: initial_len,
1520 closed: input.close_after,
1521 };
1522 producer_states.insert(producer.producer_id, ProducerState {
1523 producer_epoch: producer.producer_epoch,
1524 producer_seq: producer.producer_seq,
1525 last_start_offset: last_item.start_offset,
1526 last_next_offset: last_item.next_offset,
1527 last_closed: last_item.closed,
1528 last_items: vec![last_item],
1529 });
1530 }
1531 self.producers
1532 .insert(input.stream_id.clone(), producer_states);
1533 StreamResponse::Created {
1534 stream_id: input.stream_id,
1535 next_offset: initial_len,
1536 closed: input.close_after,
1537 }
1538 }
1539
1540 fn create_external_stream(&mut self, input: CreateExternalStreamInput) -> StreamResponse {
1541 if let Err(response) = validate_external_payload_ref(&input.initial_payload) {
1542 return response;
1543 }
1544 if let Err(response) = self.validate_stream_scope(&input.stream_id) {
1545 return response;
1546 }
1547 if let Err(response) =
1548 validate_retention(input.stream_ttl_seconds, input.stream_expires_at_ms)
1549 {
1550 return response;
1551 }
1552 if let Err(response) = validate_producer_request(input.producer.as_ref()) {
1553 return response;
1554 }
1555 if let Some(producer) = input.producer.as_ref()
1556 && producer.producer_seq != 0
1557 {
1558 return StreamResponse::error_with_context(
1559 StreamErrorCode::ProducerSeqConflict,
1560 format!(
1561 "producer '{}' expected sequence 0, received {}",
1562 producer.producer_id, producer.producer_seq
1563 ),
1564 vec![StreamErrorContext::ProducerSeqConflict {
1565 expected_seq: 0,
1566 received_seq: producer.producer_seq,
1567 }],
1568 );
1569 }
1570 if self
1571 .streams
1572 .get(&input.stream_id)
1573 .is_some_and(|existing| stream_is_expired(existing, input.now_ms))
1574 {
1575 self.remove_stream_state(&input.stream_id);
1576 }
1577
1578 if let Some(existing) = self.streams.get(&input.stream_id) {
1579 if is_soft_deleted(existing) {
1580 return StreamResponse::error(
1581 StreamErrorCode::StreamAlreadyExistsConflict,
1582 format!(
1583 "stream '{}' is gone and cannot be recreated yet",
1584 input.stream_id
1585 ),
1586 );
1587 }
1588 if existing.content_type == input.content_type
1589 && existing.status == status_from_closed(input.close_after)
1590 && existing.stream_ttl_seconds == input.stream_ttl_seconds
1591 && existing.stream_expires_at_ms == input.stream_expires_at_ms
1592 && existing.forked_from == input.forked_from
1593 && existing.fork_offset == input.fork_offset
1594 {
1595 return StreamResponse::AlreadyExists {
1596 next_offset: existing.tail_offset,
1597 closed: existing.status == StreamStatus::Closed,
1598 content_type: existing.content_type.clone(),
1599 stream_ttl_seconds: existing.stream_ttl_seconds,
1600 stream_expires_at_ms: existing.stream_expires_at_ms,
1601 };
1602 }
1603 return StreamResponse::error(
1604 StreamErrorCode::StreamAlreadyExistsConflict,
1605 format!(
1606 "stream '{}' already exists with different metadata",
1607 input.stream_id
1608 ),
1609 );
1610 }
1611
1612 let initial_len = input.initial_payload.payload_len;
1613 let metadata = StreamMetadata {
1614 stream_id: input.stream_id.clone(),
1615 content_type: input.content_type,
1616 status: status_from_closed(input.close_after),
1617 tail_offset: initial_len,
1618 last_stream_seq: input.stream_seq,
1619 stream_ttl_seconds: input.stream_ttl_seconds,
1620 stream_expires_at_ms: input.stream_expires_at_ms,
1621 created_at_ms: input.now_ms,
1622 last_ttl_touch_at_ms: input.now_ms,
1623 forked_from: input.forked_from,
1624 fork_offset: input.fork_offset,
1625 fork_ref_count: 0,
1626 };
1627 self.streams.insert(input.stream_id.clone(), metadata);
1628 self.hot_buffers
1629 .insert(input.stream_id.clone(), HotBuffer::default());
1630 let object = ObjectPayloadRef {
1631 start_offset: 0,
1632 end_offset: initial_len,
1633 s3_path: input.initial_payload.s3_path,
1634 object_size: input.initial_payload.object_size,
1635 };
1636 self.cold_index_mut()
1637 .push_external_segment(input.stream_id.clone(), object.clone());
1638 let mut integrity = StreamIntegrity::default();
1639 if initial_len > 0 {
1640 integrity.append_external(
1641 &input.stream_id,
1642 object.start_offset,
1643 object.end_offset,
1644 &object.s3_path,
1645 object.object_size,
1646 );
1647 }
1648 self.integrities.insert(input.stream_id.clone(), integrity);
1649 self.message_records
1650 .insert(input.stream_id.clone(), vec![StreamMessageRecord {
1651 start_offset: 0,
1652 end_offset: initial_len,
1653 }]);
1654 let mut producer_states = HashMap::new();
1655 if let Some(producer) = input.producer {
1656 let last_item = ProducerAppendRecord {
1657 start_offset: 0,
1658 next_offset: initial_len,
1659 closed: input.close_after,
1660 };
1661 producer_states.insert(producer.producer_id, ProducerState {
1662 producer_epoch: producer.producer_epoch,
1663 producer_seq: producer.producer_seq,
1664 last_start_offset: last_item.start_offset,
1665 last_next_offset: last_item.next_offset,
1666 last_closed: last_item.closed,
1667 last_items: vec![last_item],
1668 });
1669 }
1670 self.producers
1671 .insert(input.stream_id.clone(), producer_states);
1672 StreamResponse::Created {
1673 stream_id: input.stream_id,
1674 next_offset: initial_len,
1675 closed: input.close_after,
1676 }
1677 }
1678
1679 pub fn append_borrowed(&mut self, input: AppendStreamInput<'_>) -> StreamResponse {
1680 let AppendStreamInput {
1681 stream_id,
1682 content_type,
1683 payload,
1684 close_after,
1685 stream_seq,
1686 producer,
1687 now_ms,
1688 } = input;
1689 if let Err(response) = self.validate_stream_scope(&stream_id) {
1690 return response;
1691 }
1692 if let Err(response) = validate_producer_request(producer.as_ref()) {
1693 return response;
1694 }
1695
1696 let Some(_) = self.streams.get(&stream_id) else {
1697 return StreamResponse::error(
1698 StreamErrorCode::StreamNotFound,
1699 format!("stream '{stream_id}' does not exist"),
1700 );
1701 };
1702 if self.expire_stream_if_due(&stream_id, now_ms) {
1703 return StreamResponse::error(
1704 StreamErrorCode::StreamNotFound,
1705 format!("stream '{stream_id}' does not exist"),
1706 );
1707 }
1708 if self.streams.get(&stream_id).is_some_and(is_soft_deleted) {
1709 return StreamResponse::error(
1710 StreamErrorCode::StreamGone,
1711 format!("stream '{stream_id}' is gone"),
1712 );
1713 }
1714 let producer_decision = match self.evaluate_producer(&stream_id, producer.as_ref()) {
1715 Ok(decision) => decision,
1716 Err(response) => return response,
1717 };
1718 if let ProducerDecision::Duplicate {
1719 offset,
1720 next_offset,
1721 closed,
1722 producer,
1723 ..
1724 } = producer_decision
1725 {
1726 if payload.is_empty() {
1727 return StreamResponse::Closed {
1728 next_offset,
1729 deduplicated: true,
1730 producer: Some(producer),
1731 };
1732 }
1733 return StreamResponse::Appended {
1734 offset,
1735 next_offset,
1736 closed,
1737 deduplicated: true,
1738 producer: Some(producer),
1739 };
1740 }
1741
1742 let Some(stream) = self.streams.get_mut(&stream_id) else {
1743 unreachable!("stream existence checked before producer evaluation");
1744 };
1745
1746 if stream.status == StreamStatus::Closed {
1747 if close_after && payload.is_empty() {
1748 return StreamResponse::Closed {
1749 next_offset: stream.tail_offset,
1750 deduplicated: false,
1751 producer: None,
1752 };
1753 }
1754 return StreamResponse::error_with_next_offset_and_context(
1755 StreamErrorCode::StreamClosed,
1756 format!("stream '{stream_id}' is closed"),
1757 stream.tail_offset,
1758 vec![StreamErrorContext::StreamClosed],
1759 );
1760 }
1761
1762 if payload.is_empty() && !close_after {
1763 return StreamResponse::error(
1764 StreamErrorCode::EmptyAppend,
1765 "append payload must be non-empty unless closing the stream",
1766 );
1767 }
1768
1769 if !payload.is_empty() {
1770 let Some(content_type) = content_type else {
1771 return StreamResponse::error(
1772 StreamErrorCode::MissingContentType,
1773 "append with a body must include content type",
1774 );
1775 };
1776 if content_type != stream.content_type {
1777 return StreamResponse::error_with_next_offset(
1778 StreamErrorCode::ContentTypeMismatch,
1779 format!(
1780 "append content type '{content_type}' does not match stream content type '{}'",
1781 stream.content_type
1782 ),
1783 stream.tail_offset,
1784 );
1785 }
1786 }
1787
1788 if let Err(response) = check_stream_seq(stream, stream_seq.as_deref()) {
1789 return response;
1790 }
1791
1792 let offset = stream.tail_offset;
1793 let payload_len = u64::try_from(payload.len()).expect("payload len fits u64");
1794 stream.tail_offset = stream.tail_offset.saturating_add(payload_len);
1795 if let Some(seq) = stream_seq {
1796 stream.last_stream_seq = Some(seq);
1797 }
1798 renew_stream_ttl(stream, now_ms);
1799 if close_after {
1800 stream.status = StreamStatus::Closed;
1801 }
1802 let closed = stream.status == StreamStatus::Closed;
1803 let next_offset = stream.tail_offset;
1804 let producer_ack = producer.clone();
1805 if let Some(producer) = producer {
1806 self.record_producer_success(
1807 stream_id.clone(),
1808 producer,
1809 ProducerAppendRecord {
1810 start_offset: offset,
1811 next_offset,
1812 closed,
1813 },
1814 vec![ProducerAppendRecord {
1815 start_offset: offset,
1816 next_offset,
1817 closed,
1818 }],
1819 );
1820 }
1821
1822 if payload.is_empty() {
1823 StreamResponse::Closed {
1824 next_offset,
1825 deduplicated: false,
1826 producer: producer_ack,
1827 }
1828 } else {
1829 self.hot_buffers
1830 .get_mut(&stream_id)
1831 .expect("hot buffer exists for stream metadata")
1832 .push(offset, next_offset, payload);
1833 self.integrities
1834 .get_mut(&stream_id)
1835 .expect("integrity exists for stream metadata")
1836 .append_payload(&stream_id, offset, next_offset, payload);
1837 self.message_records
1838 .entry(stream_id.clone())
1839 .or_default()
1840 .push(StreamMessageRecord {
1841 start_offset: offset,
1842 end_offset: next_offset,
1843 });
1844 StreamResponse::Appended {
1845 offset,
1846 next_offset,
1847 closed: close_after,
1848 deduplicated: false,
1849 producer: producer_ack,
1850 }
1851 }
1852 }
1853
1854 fn append_external(&mut self, input: AppendExternalInput<'_>) -> StreamResponse {
1855 let AppendExternalInput {
1856 stream_id,
1857 content_type,
1858 payload,
1859 close_after,
1860 stream_seq,
1861 producer,
1862 now_ms,
1863 } = input;
1864 if let Err(response) = validate_external_payload_ref(&payload) {
1865 return response;
1866 }
1867 if let Err(response) = self.validate_stream_scope(&stream_id) {
1868 return response;
1869 }
1870 if let Err(response) = validate_producer_request(producer.as_ref()) {
1871 return response;
1872 }
1873 let Some(_) = self.streams.get(&stream_id) else {
1874 return StreamResponse::error(
1875 StreamErrorCode::StreamNotFound,
1876 format!("stream '{stream_id}' does not exist"),
1877 );
1878 };
1879 if self.expire_stream_if_due(&stream_id, now_ms) {
1880 return StreamResponse::error(
1881 StreamErrorCode::StreamNotFound,
1882 format!("stream '{stream_id}' does not exist"),
1883 );
1884 }
1885 if self.streams.get(&stream_id).is_some_and(is_soft_deleted) {
1886 return StreamResponse::error(
1887 StreamErrorCode::StreamGone,
1888 format!("stream '{stream_id}' is gone"),
1889 );
1890 }
1891 let producer_decision = match self.evaluate_producer(&stream_id, producer.as_ref()) {
1892 Ok(decision) => decision,
1893 Err(response) => return response,
1894 };
1895 if let ProducerDecision::Duplicate {
1896 offset,
1897 next_offset,
1898 closed,
1899 producer,
1900 ..
1901 } = producer_decision
1902 {
1903 return StreamResponse::Appended {
1904 offset,
1905 next_offset,
1906 closed,
1907 deduplicated: true,
1908 producer: Some(producer),
1909 };
1910 }
1911
1912 let Some(stream) = self.streams.get(&stream_id) else {
1913 unreachable!("stream existence checked before producer evaluation");
1914 };
1915 if stream.status == StreamStatus::Closed {
1916 return StreamResponse::error_with_next_offset_and_context(
1917 StreamErrorCode::StreamClosed,
1918 format!("stream '{stream_id}' is closed"),
1919 stream.tail_offset,
1920 vec![StreamErrorContext::StreamClosed],
1921 );
1922 }
1923 let Some(content_type) = content_type else {
1924 return StreamResponse::error(
1925 StreamErrorCode::MissingContentType,
1926 "append with a body must include content type",
1927 );
1928 };
1929 if content_type != stream.content_type {
1930 return StreamResponse::error_with_next_offset(
1931 StreamErrorCode::ContentTypeMismatch,
1932 format!(
1933 "append content type '{content_type}' does not match stream content type '{}'",
1934 stream.content_type
1935 ),
1936 stream.tail_offset,
1937 );
1938 }
1939 if let Err(response) = check_stream_seq(stream, stream_seq.as_deref()) {
1940 return response;
1941 }
1942 let offset = stream.tail_offset;
1943 let next_offset = offset.saturating_add(payload.payload_len);
1944 let stream = self
1945 .streams
1946 .get_mut(&stream_id)
1947 .expect("stream existence checked before external append mutation");
1948 stream.tail_offset = next_offset;
1949 if let Some(seq) = stream_seq {
1950 stream.last_stream_seq = Some(seq);
1951 }
1952 renew_stream_ttl(stream, now_ms);
1953 if close_after {
1954 stream.status = StreamStatus::Closed;
1955 }
1956 let closed = stream.status == StreamStatus::Closed;
1957 let producer_ack = producer.clone();
1958 if let Some(producer) = producer {
1959 self.record_producer_success(
1960 stream_id.clone(),
1961 producer,
1962 ProducerAppendRecord {
1963 start_offset: offset,
1964 next_offset,
1965 closed,
1966 },
1967 vec![ProducerAppendRecord {
1968 start_offset: offset,
1969 next_offset,
1970 closed,
1971 }],
1972 );
1973 }
1974 let object = ObjectPayloadRef {
1975 start_offset: offset,
1976 end_offset: next_offset,
1977 s3_path: payload.s3_path,
1978 object_size: payload.object_size,
1979 };
1980 self.cold_index_mut()
1981 .push_external_segment(stream_id.clone(), object.clone());
1982 self.integrities
1983 .get_mut(&stream_id)
1984 .expect("integrity exists for stream metadata")
1985 .append_external(
1986 &stream_id,
1987 object.start_offset,
1988 object.end_offset,
1989 &object.s3_path,
1990 object.object_size,
1991 );
1992 self.message_records
1993 .entry(stream_id.clone())
1994 .or_default()
1995 .push(StreamMessageRecord {
1996 start_offset: offset,
1997 end_offset: next_offset,
1998 });
1999 StreamResponse::Appended {
2000 offset,
2001 next_offset,
2002 closed: close_after,
2003 deduplicated: false,
2004 producer: producer_ack,
2005 }
2006 }
2007
2008 pub fn append_batch_borrowed(
2009 &mut self,
2010 stream_id: BucketStreamId,
2011 content_type: Option<&str>,
2012 payloads: &[&[u8]],
2013 producer: Option<ProducerRequest>,
2014 now_ms: u64,
2015 ) -> Result<StreamBatchAppend, StreamResponse> {
2016 if payloads.is_empty() {
2017 return Err(StreamResponse::error(
2018 StreamErrorCode::EmptyAppend,
2019 "append batch must contain at least one payload",
2020 ));
2021 }
2022 self.validate_stream_scope(&stream_id)?;
2023 validate_producer_request(producer.as_ref())?;
2024 if self.expire_stream_if_due(&stream_id, now_ms) {
2025 return Err(StreamResponse::error(
2026 StreamErrorCode::StreamNotFound,
2027 format!("stream '{stream_id}' does not exist"),
2028 ));
2029 }
2030 if self.streams.get(&stream_id).is_some_and(is_soft_deleted) {
2031 return Err(StreamResponse::error(
2032 StreamErrorCode::StreamGone,
2033 format!("stream '{stream_id}' is gone"),
2034 ));
2035 }
2036 let producer_decision = self.evaluate_producer(&stream_id, producer.as_ref())?;
2037 if let ProducerDecision::Duplicate { items, .. } = producer_decision {
2038 return Ok(StreamBatchAppend {
2039 items: items
2040 .into_iter()
2041 .map(|item| StreamBatchAppendItem {
2042 offset: item.start_offset,
2043 next_offset: item.next_offset,
2044 closed: item.closed,
2045 deduplicated: true,
2046 })
2047 .collect(),
2048 deduplicated: true,
2049 });
2050 }
2051
2052 let Some(stream) = self.streams.get_mut(&stream_id) else {
2053 return Err(StreamResponse::error(
2054 StreamErrorCode::StreamNotFound,
2055 format!("stream '{stream_id}' does not exist"),
2056 ));
2057 };
2058 if stream.status == StreamStatus::Closed {
2059 return Err(StreamResponse::error_with_next_offset_and_context(
2060 StreamErrorCode::StreamClosed,
2061 format!("stream '{stream_id}' is closed"),
2062 stream.tail_offset,
2063 vec![StreamErrorContext::StreamClosed],
2064 ));
2065 }
2066 let Some(content_type) = content_type else {
2067 return Err(StreamResponse::error(
2068 StreamErrorCode::MissingContentType,
2069 "append batch must include content type",
2070 ));
2071 };
2072 if content_type != stream.content_type {
2073 return Err(StreamResponse::error_with_next_offset(
2074 StreamErrorCode::ContentTypeMismatch,
2075 format!(
2076 "append content type '{content_type}' does not match stream content type '{}'",
2077 stream.content_type
2078 ),
2079 stream.tail_offset,
2080 ));
2081 }
2082 if payloads.iter().any(|payload| payload.is_empty()) {
2083 return Err(StreamResponse::error(
2084 StreamErrorCode::EmptyAppend,
2085 "append batch payloads must be non-empty",
2086 ));
2087 }
2088
2089 let mut items = Vec::with_capacity(payloads.len());
2090 for payload in payloads {
2091 let offset = stream.tail_offset;
2092 let payload_len = u64::try_from(payload.len()).expect("payload len fits u64");
2093 stream.tail_offset = stream.tail_offset.saturating_add(payload_len);
2094 items.push(ProducerAppendRecord {
2095 start_offset: offset,
2096 next_offset: stream.tail_offset,
2097 closed: false,
2098 });
2099 }
2100 let last = items
2101 .last()
2102 .expect("payloads checked non-empty before append")
2103 .clone();
2104 renew_stream_ttl(stream, now_ms);
2105 if let Some(producer) = producer {
2106 self.record_producer_success(stream_id.clone(), producer, last.clone(), items.clone());
2107 }
2108 let hot_buffer = self
2109 .hot_buffers
2110 .get_mut(&stream_id)
2111 .expect("hot buffer exists for stream metadata");
2112 for (item, payload) in items.iter().zip(payloads.iter()) {
2113 hot_buffer.push(item.start_offset, item.next_offset, payload);
2114 }
2115 let integrity = self
2116 .integrities
2117 .get_mut(&stream_id)
2118 .expect("integrity exists for stream metadata");
2119 for (item, payload) in items.iter().zip(payloads.iter()) {
2120 integrity.append_payload(&stream_id, item.start_offset, item.next_offset, payload);
2121 }
2122 self.message_records
2123 .entry(stream_id.clone())
2124 .or_default()
2125 .extend(items.iter().map(|item| StreamMessageRecord {
2126 start_offset: item.start_offset,
2127 end_offset: item.next_offset,
2128 }));
2129 Ok(StreamBatchAppend {
2130 items: items
2131 .into_iter()
2132 .map(|item| StreamBatchAppendItem {
2133 offset: item.start_offset,
2134 next_offset: item.next_offset,
2135 closed: item.closed,
2136 deduplicated: false,
2137 })
2138 .collect(),
2139 deduplicated: false,
2140 })
2141 }
2142
2143 fn close(
2144 &mut self,
2145 stream_id: BucketStreamId,
2146 stream_seq: Option<String>,
2147 producer: Option<ProducerRequest>,
2148 now_ms: u64,
2149 ) -> StreamResponse {
2150 self.append_borrowed(AppendStreamInput {
2151 stream_id,
2152 content_type: None,
2153 payload: &[],
2154 close_after: true,
2155 stream_seq,
2156 producer,
2157 now_ms,
2158 })
2159 }
2160
2161 fn delete_stream(&mut self, stream_id: &BucketStreamId) -> StreamResponse {
2162 if let Err(response) = self.validate_stream_scope(stream_id) {
2163 return response;
2164 }
2165 let Some(stream) = self.streams.get_mut(stream_id) else {
2166 return StreamResponse::error(
2167 StreamErrorCode::StreamNotFound,
2168 format!("stream '{stream_id}' does not exist"),
2169 );
2170 };
2171 if is_soft_deleted(stream) {
2172 return StreamResponse::error(
2173 StreamErrorCode::StreamGone,
2174 format!("stream '{stream_id}' is gone"),
2175 );
2176 }
2177 if stream.fork_ref_count > 0 {
2178 stream.status = StreamStatus::SoftDeleted;
2179 return StreamResponse::Deleted {
2180 hard_deleted: false,
2181 parent_to_release: None,
2182 };
2183 }
2184 let parent_to_release = stream.forked_from.clone();
2185 self.remove_stream_state(stream_id);
2186 StreamResponse::Deleted {
2187 hard_deleted: true,
2188 parent_to_release,
2189 }
2190 }
2191
2192 fn add_fork_ref(&mut self, stream_id: &BucketStreamId, now_ms: u64) -> StreamResponse {
2193 if let Err(response) = self.validate_stream_scope(stream_id) {
2194 return response;
2195 }
2196 if self.expire_stream_if_due(stream_id, now_ms) {
2197 return StreamResponse::error(
2198 StreamErrorCode::StreamNotFound,
2199 format!("stream '{stream_id}' does not exist"),
2200 );
2201 }
2202 let Some(stream) = self.streams.get_mut(stream_id) else {
2203 return StreamResponse::error(
2204 StreamErrorCode::StreamNotFound,
2205 format!("stream '{stream_id}' does not exist"),
2206 );
2207 };
2208 if is_soft_deleted(stream) {
2209 return StreamResponse::error(
2210 StreamErrorCode::StreamGone,
2211 format!("stream '{stream_id}' is gone"),
2212 );
2213 }
2214 stream.fork_ref_count = stream.fork_ref_count.saturating_add(1);
2215 StreamResponse::ForkRefAdded {
2216 fork_ref_count: stream.fork_ref_count,
2217 }
2218 }
2219
2220 fn release_fork_ref(&mut self, stream_id: &BucketStreamId) -> StreamResponse {
2221 if let Err(response) = self.validate_stream_scope(stream_id) {
2222 return response;
2223 }
2224 let Some(stream) = self.streams.get_mut(stream_id) else {
2225 return StreamResponse::ForkRefReleased {
2226 hard_deleted: false,
2227 fork_ref_count: 0,
2228 parent_to_release: None,
2229 };
2230 };
2231 if stream.fork_ref_count == 0 {
2232 return StreamResponse::error(
2233 StreamErrorCode::InvalidFork,
2234 format!("stream '{stream_id}' has no fork reference to release"),
2235 );
2236 }
2237 stream.fork_ref_count -= 1;
2238 if stream.fork_ref_count == 0 && is_soft_deleted(stream) {
2239 let parent_to_release = stream.forked_from.clone();
2240 self.remove_stream_state(stream_id);
2241 return StreamResponse::ForkRefReleased {
2242 hard_deleted: true,
2243 fork_ref_count: 0,
2244 parent_to_release,
2245 };
2246 }
2247 StreamResponse::ForkRefReleased {
2248 hard_deleted: false,
2249 fork_ref_count: stream.fork_ref_count,
2250 parent_to_release: None,
2251 }
2252 }
2253
2254 fn touch_stream_access(
2255 &mut self,
2256 stream_id: &BucketStreamId,
2257 now_ms: u64,
2258 renew_ttl: bool,
2259 ) -> StreamResponse {
2260 if let Err(response) = self.validate_stream_scope(stream_id) {
2261 return response;
2262 }
2263 let Some(stream) = self.streams.get(stream_id) else {
2264 return StreamResponse::error(
2265 StreamErrorCode::StreamNotFound,
2266 format!("stream '{stream_id}' does not exist"),
2267 );
2268 };
2269 if is_soft_deleted(stream) {
2270 return StreamResponse::error(
2271 StreamErrorCode::StreamGone,
2272 format!("stream '{stream_id}' is gone"),
2273 );
2274 }
2275 if stream_is_expired(stream, now_ms) {
2276 self.remove_stream_state(stream_id);
2277 return StreamResponse::Accessed {
2278 changed: true,
2279 expired: true,
2280 };
2281 }
2282 let changed = if renew_ttl && stream.stream_ttl_seconds.is_some() {
2283 let stream = self
2284 .streams
2285 .get_mut(stream_id)
2286 .expect("stream existence checked before TTL renewal");
2287 let previous = stream.last_ttl_touch_at_ms;
2288 renew_stream_ttl(stream, now_ms);
2289 stream.last_ttl_touch_at_ms != previous
2290 } else {
2291 false
2292 };
2293 StreamResponse::Accessed {
2294 changed,
2295 expired: false,
2296 }
2297 }
2298
2299 fn expire_stream_if_due(&mut self, stream_id: &BucketStreamId, now_ms: u64) -> bool {
2300 if self
2301 .streams
2302 .get(stream_id)
2303 .is_some_and(|stream| stream_is_expired(stream, now_ms))
2304 {
2305 self.remove_stream_state(stream_id);
2306 return true;
2307 }
2308 false
2309 }
2310
2311 fn remove_stream_state(&mut self, stream_id: &BucketStreamId) -> bool {
2312 if self.streams.remove(stream_id).is_some() {
2313 self.hot_buffers.remove(stream_id);
2314 let had_cold = self.cold_index_mut().remove_stream(stream_id);
2315 self.message_records.remove(stream_id);
2316 self.integrities.remove(stream_id);
2317 self.visible_snapshots.remove(stream_id);
2318 self.producers.remove(stream_id);
2319 if had_cold {
2324 self.enqueue_cold_gc(ColdGcTarget::Stream(stream_id.clone()));
2325 }
2326 true
2327 } else {
2328 false
2329 }
2330 }
2331
2332 fn enqueue_cold_gc(&mut self, target: ColdGcTarget) {
2333 let seq = self.next_cold_gc_seq;
2334 self.next_cold_gc_seq = self.next_cold_gc_seq.saturating_add(1);
2335 self.pending_cold_gc.push_back(ColdGcEntry { seq, target });
2336 }
2337
2338 fn ack_cold_gc(&mut self, up_to_seq: u64) -> StreamResponse {
2339 let before = self.pending_cold_gc.len();
2340 while self
2341 .pending_cold_gc
2342 .front()
2343 .is_some_and(|entry| entry.seq <= up_to_seq)
2344 {
2345 self.pending_cold_gc.pop_front();
2346 }
2347 let removed = u64::try_from(before - self.pending_cold_gc.len()).expect("removed fits u64");
2348 StreamResponse::ColdGcAcked { removed }
2349 }
2350
2351 pub fn pending_cold_gc_batch(&self, max: usize) -> Vec<ColdGcEntry> {
2354 self.pending_cold_gc.iter().take(max).cloned().collect()
2355 }
2356
2357 pub fn pending_cold_gc_len(&self) -> usize {
2358 self.pending_cold_gc.len()
2359 }
2360
2361 fn validate_stream_scope(&self, stream_id: &BucketStreamId) -> Result<(), StreamResponse> {
2362 if let Err(message) = validate_bucket_id(&stream_id.bucket_id) {
2363 return Err(StreamResponse::error(
2364 StreamErrorCode::InvalidBucketId,
2365 message,
2366 ));
2367 }
2368 if let Err(message) = validate_stream_id(stream_id) {
2369 return Err(StreamResponse::error(
2370 StreamErrorCode::InvalidStreamId,
2371 message,
2372 ));
2373 }
2374 if !self.buckets.contains(&stream_id.bucket_id) {
2375 return Err(StreamResponse::error(
2376 StreamErrorCode::BucketNotFound,
2377 format!("bucket '{}' does not exist", stream_id.bucket_id),
2378 ));
2379 }
2380 Ok(())
2381 }
2382
2383 fn earliest_retained_offset(&self, stream_id: &BucketStreamId) -> u64 {
2384 self.visible_snapshots
2385 .get(stream_id)
2386 .map(|snapshot| snapshot.offset)
2387 .unwrap_or(0)
2388 }
2389
2390 fn snapshot_offset_aligned(
2391 &self,
2392 stream_id: &BucketStreamId,
2393 snapshot_offset: u64,
2394 retained_offset: u64,
2395 ) -> bool {
2396 snapshot_offset == retained_offset
2397 || snapshot_offset <= self.cold_frontier_offset(stream_id, retained_offset)
2398 || self
2399 .hot_buffers
2400 .get(stream_id)
2401 .is_some_and(|buffer| snapshot_offset <= buffer.hot_start_offset())
2402 || self.message_records.get(stream_id).is_some_and(|records| {
2403 records
2404 .iter()
2405 .any(|record| record.end_offset == snapshot_offset)
2406 })
2407 }
2408
2409 fn compact_retained_prefix(&mut self, stream_id: &BucketStreamId, retained_offset: u64) {
2410 let frontier = self.cold_frontier_offset(stream_id, retained_offset).max(
2411 self.hot_buffers
2412 .get(stream_id)
2413 .map(|buffer| buffer.hot_start_offset())
2414 .unwrap_or(retained_offset),
2415 );
2416 self.compact_message_records_before(stream_id, retained_offset, frontier);
2417 if let Some(integrity) = self.integrities.get_mut(stream_id) {
2418 integrity.evict_before(retained_offset);
2419 }
2420 let dropped_cold_paths = self
2421 .cold_index_mut()
2422 .compact_before(stream_id, retained_offset);
2423 if !dropped_cold_paths.is_empty() {
2424 self.enqueue_cold_gc(ColdGcTarget::Paths(dropped_cold_paths));
2425 }
2426
2427 if let Some(hot_buffer) = self.hot_buffers.get_mut(stream_id) {
2428 hot_buffer.discard_before(retained_offset);
2429 }
2430 }
2431
2432 fn compact_message_records_before(
2433 &mut self,
2434 stream_id: &BucketStreamId,
2435 retained_offset: u64,
2436 frontier: u64,
2437 ) {
2438 let Some(records) = self.message_records.remove(stream_id) else {
2439 return;
2440 };
2441 let frontier = frontier.max(retained_offset);
2442 let mut compacted = Vec::with_capacity(records.len());
2443 if frontier > retained_offset {
2444 compacted.push(StreamMessageRecord {
2445 start_offset: retained_offset,
2446 end_offset: frontier,
2447 });
2448 }
2449 compacted.extend(records.iter().filter_map(|record| {
2450 if record.end_offset <= frontier {
2451 return None;
2452 }
2453 let start_offset = record.start_offset.max(frontier).max(retained_offset);
2454 (record.end_offset > start_offset).then_some(StreamMessageRecord {
2455 start_offset,
2456 end_offset: record.end_offset,
2457 })
2458 }));
2459 if compacted.is_empty() {
2460 return;
2461 }
2462 self.message_records.insert(stream_id.clone(), compacted);
2463 }
2464
2465 fn cold_frontier_offset(&self, stream_id: &BucketStreamId, retained_offset: u64) -> u64 {
2466 self.cold_index()
2467 .cold_frontier_offset(stream_id, retained_offset)
2468 }
2469
2470 fn producer_snapshot(&self, stream_id: &BucketStreamId) -> Vec<ProducerSnapshot> {
2471 let mut producer_states = self
2472 .producers
2473 .get(stream_id)
2474 .into_iter()
2475 .flat_map(|states| states.iter())
2476 .map(|(producer_id, state)| ProducerSnapshot {
2477 producer_id: producer_id.clone(),
2478 producer_epoch: state.producer_epoch,
2479 producer_seq: state.producer_seq,
2480 last_start_offset: state.last_start_offset,
2481 last_next_offset: state.last_next_offset,
2482 last_closed: state.last_closed,
2483 last_items: state.last_items.clone(),
2484 })
2485 .collect::<Vec<_>>();
2486 producer_states.sort_by(|left, right| left.producer_id.cmp(&right.producer_id));
2487 producer_states
2488 }
2489
2490 fn evaluate_producer(
2491 &self,
2492 stream_id: &BucketStreamId,
2493 producer: Option<&ProducerRequest>,
2494 ) -> Result<ProducerDecision, StreamResponse> {
2495 let Some(producer) = producer else {
2496 return Ok(ProducerDecision::Accept);
2497 };
2498 let Some(states) = self.producers.get(stream_id) else {
2499 return Ok(ProducerDecision::Accept);
2500 };
2501 let Some(state) = states.get(&producer.producer_id) else {
2502 if producer.producer_seq == 0 {
2503 return Ok(ProducerDecision::Accept);
2504 }
2505 return Err(StreamResponse::error_with_context(
2506 StreamErrorCode::ProducerSeqConflict,
2507 format!(
2508 "producer '{}' expected sequence 0, received {}",
2509 producer.producer_id, producer.producer_seq
2510 ),
2511 vec![StreamErrorContext::ProducerSeqConflict {
2512 expected_seq: 0,
2513 received_seq: producer.producer_seq,
2514 }],
2515 ));
2516 };
2517
2518 if producer.producer_epoch < state.producer_epoch {
2519 return Err(StreamResponse::error_with_context(
2520 StreamErrorCode::ProducerEpochStale,
2521 format!(
2522 "producer '{}' epoch {} is stale; current epoch is {}",
2523 producer.producer_id, producer.producer_epoch, state.producer_epoch
2524 ),
2525 vec![StreamErrorContext::ProducerEpochStale {
2526 current_epoch: state.producer_epoch,
2527 }],
2528 ));
2529 }
2530 if producer.producer_epoch > state.producer_epoch {
2531 if producer.producer_seq == 0 {
2532 return Ok(ProducerDecision::Accept);
2533 }
2534 return Err(StreamResponse::error(
2535 StreamErrorCode::InvalidProducer,
2536 format!(
2537 "producer '{}' new epoch {} must start at sequence 0",
2538 producer.producer_id, producer.producer_epoch
2539 ),
2540 ));
2541 }
2542
2543 if producer.producer_seq <= state.producer_seq {
2544 return Ok(ProducerDecision::Duplicate {
2545 offset: state.last_start_offset,
2546 next_offset: state.last_next_offset,
2547 closed: state.last_closed,
2548 producer: ProducerRequest {
2549 producer_id: producer.producer_id.clone(),
2550 producer_epoch: state.producer_epoch,
2551 producer_seq: state.producer_seq,
2552 },
2553 items: state.last_items.clone(),
2554 });
2555 }
2556 if producer.producer_seq == state.producer_seq + 1 {
2557 return Ok(ProducerDecision::Accept);
2558 }
2559 Err(StreamResponse::error_with_context(
2560 StreamErrorCode::ProducerSeqConflict,
2561 format!(
2562 "producer '{}' expected sequence {}, received {}",
2563 producer.producer_id,
2564 state.producer_seq + 1,
2565 producer.producer_seq
2566 ),
2567 vec![StreamErrorContext::ProducerSeqConflict {
2568 expected_seq: state.producer_seq + 1,
2569 received_seq: producer.producer_seq,
2570 }],
2571 ))
2572 }
2573
2574 fn record_producer_success(
2575 &mut self,
2576 stream_id: BucketStreamId,
2577 producer: ProducerRequest,
2578 last: ProducerAppendRecord,
2579 last_items: Vec<ProducerAppendRecord>,
2580 ) {
2581 self.producers
2582 .entry(stream_id)
2583 .or_default()
2584 .insert(producer.producer_id, ProducerState {
2585 producer_epoch: producer.producer_epoch,
2586 producer_seq: producer.producer_seq,
2587 last_start_offset: last.start_offset,
2588 last_next_offset: last.next_offset,
2589 last_closed: last.closed,
2590 last_items,
2591 });
2592 }
2593}
2594
2595#[derive(Debug)]
2596struct CreateStreamInput {
2597 stream_id: BucketStreamId,
2598 content_type: String,
2599 initial_payload: Vec<u8>,
2600 close_after: bool,
2601 stream_seq: Option<String>,
2602 producer: Option<ProducerRequest>,
2603 stream_ttl_seconds: Option<u64>,
2604 stream_expires_at_ms: Option<u64>,
2605 forked_from: Option<BucketStreamId>,
2606 fork_offset: Option<u64>,
2607 now_ms: u64,
2608}
2609
2610#[derive(Debug)]
2611struct CreateExternalStreamInput {
2612 stream_id: BucketStreamId,
2613 content_type: String,
2614 initial_payload: ExternalPayloadRef,
2615 close_after: bool,
2616 stream_seq: Option<String>,
2617 producer: Option<ProducerRequest>,
2618 stream_ttl_seconds: Option<u64>,
2619 stream_expires_at_ms: Option<u64>,
2620 forked_from: Option<BucketStreamId>,
2621 fork_offset: Option<u64>,
2622 now_ms: u64,
2623}
2624
2625#[derive(Debug, Clone, PartialEq, Eq)]
2626enum ProducerDecision {
2627 Accept,
2628 Duplicate {
2629 offset: u64,
2630 next_offset: u64,
2631 closed: bool,
2632 producer: ProducerRequest,
2633 items: Vec<ProducerAppendRecord>,
2634 },
2635}
2636
2637impl CreateStreamInput {
2638 fn initial_len(&self) -> u64 {
2639 u64::try_from(self.initial_payload.len()).expect("payload len fits u64")
2640 }
2641}
2642
2643fn status_from_closed(closed: bool) -> StreamStatus {
2644 if closed {
2645 StreamStatus::Closed
2646 } else {
2647 StreamStatus::Open
2648 }
2649}
2650
2651fn is_soft_deleted(stream: &StreamMetadata) -> bool {
2652 stream.status == StreamStatus::SoftDeleted
2653}
2654
2655fn validate_retention(
2656 stream_ttl_seconds: Option<u64>,
2657 stream_expires_at_ms: Option<u64>,
2658) -> Result<(), StreamResponse> {
2659 if stream_ttl_seconds.is_some() && stream_expires_at_ms.is_some() {
2660 return Err(StreamResponse::error(
2661 StreamErrorCode::InvalidRetention,
2662 "stream ttl and expires-at cannot both be set",
2663 ));
2664 }
2665 if let Some(ttl_seconds) = stream_ttl_seconds
2666 && ttl_seconds.checked_mul(1000).is_none()
2667 {
2668 return Err(StreamResponse::error(
2669 StreamErrorCode::InvalidRetention,
2670 "stream ttl overflows millisecond range",
2671 ));
2672 }
2673 Ok(())
2674}
2675
2676fn stream_expiry_at_ms(stream: &StreamMetadata) -> Option<u64> {
2677 if let Some(expires_at_ms) = stream.stream_expires_at_ms {
2678 return Some(expires_at_ms);
2679 }
2680 stream.stream_ttl_seconds.map(|ttl_seconds| {
2681 stream
2682 .last_ttl_touch_at_ms
2683 .saturating_add(ttl_seconds.saturating_mul(1000))
2684 })
2685}
2686
2687fn stream_is_expired(stream: &StreamMetadata, now_ms: u64) -> bool {
2688 stream_expiry_at_ms(stream).is_some_and(|expires_at_ms| now_ms >= expires_at_ms)
2689}
2690
2691fn renew_stream_ttl(stream: &mut StreamMetadata, now_ms: u64) {
2692 if stream.stream_ttl_seconds.is_some() && stream.stream_expires_at_ms.is_none() {
2693 stream.last_ttl_touch_at_ms = now_ms;
2694 }
2695}
2696
2697fn check_stream_seq(stream: &StreamMetadata, incoming: Option<&str>) -> Result<(), StreamResponse> {
2698 let Some(incoming) = incoming else {
2699 return Ok(());
2700 };
2701 if let Some(last) = stream.last_stream_seq.as_deref()
2702 && incoming <= last
2703 {
2704 return Err(StreamResponse::error_with_next_offset(
2705 StreamErrorCode::StreamSeqConflict,
2706 format!("stream sequence '{incoming}' is not greater than last sequence '{last}'"),
2707 stream.tail_offset,
2708 ));
2709 }
2710 Ok(())
2711}
2712
2713fn validate_producer_request(producer: Option<&ProducerRequest>) -> Result<(), StreamResponse> {
2714 let Some(producer) = producer else {
2715 return Ok(());
2716 };
2717 if producer.producer_id.trim().is_empty() {
2718 return Err(StreamResponse::error(
2719 StreamErrorCode::InvalidProducer,
2720 "producer id must not be empty",
2721 ));
2722 }
2723 const MAX_JS_SAFE_INTEGER: u64 = 9_007_199_254_740_991;
2724 if producer.producer_epoch > MAX_JS_SAFE_INTEGER {
2725 return Err(StreamResponse::error(
2726 StreamErrorCode::InvalidProducer,
2727 format!(
2728 "producer epoch {} exceeds maximum {}",
2729 producer.producer_epoch, MAX_JS_SAFE_INTEGER
2730 ),
2731 ));
2732 }
2733 if producer.producer_seq > MAX_JS_SAFE_INTEGER {
2734 return Err(StreamResponse::error(
2735 StreamErrorCode::InvalidProducer,
2736 format!(
2737 "producer sequence {} exceeds maximum {}",
2738 producer.producer_seq, MAX_JS_SAFE_INTEGER
2739 ),
2740 ));
2741 }
2742 Ok(())
2743}
2744
2745fn validate_external_payload_ref(payload: &ExternalPayloadRef) -> Result<(), StreamResponse> {
2746 if payload.s3_path.trim().is_empty() {
2747 return Err(StreamResponse::error(
2748 StreamErrorCode::InvalidColdFlush,
2749 "external payload S3 path must not be empty",
2750 ));
2751 }
2752 if payload.payload_len == 0 {
2753 return Err(StreamResponse::error(
2754 StreamErrorCode::EmptyAppend,
2755 "external payload length must be greater than zero",
2756 ));
2757 }
2758 if payload.object_size < payload.payload_len {
2759 return Err(StreamResponse::error(
2760 StreamErrorCode::InvalidColdFlush,
2761 "external payload object size must cover payload length",
2762 ));
2763 }
2764 Ok(())
2765}
2766
2767fn restore_producer_states(
2768 stream_id: &BucketStreamId,
2769 snapshots: Vec<ProducerSnapshot>,
2770) -> Result<HashMap<String, ProducerState>, StreamSnapshotError> {
2771 let mut states = HashMap::with_capacity(snapshots.len());
2772 for snapshot in snapshots {
2773 if states
2774 .insert(snapshot.producer_id.clone(), ProducerState {
2775 producer_epoch: snapshot.producer_epoch,
2776 producer_seq: snapshot.producer_seq,
2777 last_start_offset: snapshot.last_start_offset,
2778 last_next_offset: snapshot.last_next_offset,
2779 last_closed: snapshot.last_closed,
2780 last_items: snapshot.last_items,
2781 })
2782 .is_some()
2783 {
2784 return Err(StreamSnapshotError::DuplicateProducer {
2785 stream_id: stream_id.clone(),
2786 producer_id: snapshot.producer_id,
2787 });
2788 }
2789 }
2790 Ok(states)
2791}
2792
2793fn valid_cold_chunk_ref(chunk: &ColdChunkRef) -> bool {
2794 chunk.end_offset > chunk.start_offset
2795 && !chunk.s3_path.trim().is_empty()
2796 && chunk.object_size >= chunk.end_offset - chunk.start_offset
2797}
2798
2799fn valid_object_payload_ref(object: &ObjectPayloadRef) -> bool {
2800 object.end_offset > object.start_offset
2801 && !object.s3_path.trim().is_empty()
2802 && object.object_size >= object.end_offset - object.start_offset
2803}
2804
2805fn hot_segments_match_payload(segments: &[HotPayloadSegment], payload_len: usize) -> bool {
2806 let mut expected_payload_start = 0;
2807 for segment in segments {
2808 if segment.end_offset <= segment.start_offset
2809 || segment.payload_start != expected_payload_start
2810 || segment.payload_end <= segment.payload_start
2811 || segment.payload_end > payload_len
2812 {
2813 return false;
2814 }
2815 let Ok(logical_len) = usize::try_from(segment.end_offset - segment.start_offset) else {
2816 return false;
2817 };
2818 if logical_len != segment.payload_end - segment.payload_start {
2819 return false;
2820 }
2821 expected_payload_start = segment.payload_end;
2822 }
2823 expected_payload_start == payload_len
2824}
2825
2826fn payload_sources_cover_retained_suffix(
2827 cold_frontier_offset: u64,
2828 cold_chunks: &[ColdChunkRef],
2829 external_segments: &[ObjectPayloadRef],
2830 hot_segments: &[HotPayloadSegment],
2831 retained_offset: u64,
2832 tail_offset: u64,
2833) -> bool {
2834 if tail_offset < retained_offset {
2835 return false;
2836 }
2837 let mut ranges =
2838 Vec::with_capacity(1 + cold_chunks.len() + external_segments.len() + hot_segments.len());
2839 if cold_frontier_offset > retained_offset {
2840 ranges.push((retained_offset, cold_frontier_offset));
2841 }
2842 for chunk in cold_chunks {
2843 if !valid_cold_chunk_ref(chunk) {
2844 return false;
2845 }
2846 ranges.push((chunk.start_offset, chunk.end_offset));
2847 }
2848 for object in external_segments {
2849 if !valid_object_payload_ref(object) {
2850 return false;
2851 }
2852 ranges.push((object.start_offset, object.end_offset));
2853 }
2854 for segment in hot_segments {
2855 if segment.end_offset <= segment.start_offset {
2856 return false;
2857 }
2858 ranges.push((segment.start_offset, segment.end_offset));
2859 }
2860 ranges.sort_unstable();
2861
2862 let mut expected_start = retained_offset;
2863 for (start_offset, end_offset) in ranges {
2864 if end_offset <= expected_start {
2865 continue;
2866 }
2867 if start_offset > expected_start {
2868 return false;
2869 }
2870 expected_start = end_offset;
2871 if expected_start >= tail_offset {
2872 return true;
2873 }
2874 }
2875 expected_start == tail_offset
2876}
2877
2878fn push_cold_index_segments(
2879 segments: &mut Vec<(u64, StreamReadSegment)>,
2880 _stream_id: &BucketStreamId,
2881 generation: u64,
2882 start_offset: u64,
2883 end_offset: u64,
2884) {
2885 let mut cursor = start_offset;
2886 while cursor < end_offset {
2887 let page_id = cursor / COLD_INDEX_PAGE_SPAN_BYTES;
2888 let page_end = page_id
2889 .saturating_add(1)
2890 .saturating_mul(COLD_INDEX_PAGE_SPAN_BYTES);
2891 let segment_end = end_offset.min(page_end);
2892 segments.push((
2893 cursor,
2894 StreamReadSegment::ColdIndex(StreamReadColdIndexSegment {
2895 generation,
2896 page_id,
2897 read_start_offset: cursor,
2898 len: usize::try_from(segment_end - cursor).expect("cold index read len fits usize"),
2899 }),
2900 ));
2901 cursor = segment_end;
2902 }
2903}
2904
2905fn segments_cover_range(
2906 segments: &[(u64, StreamReadSegment)],
2907 offset: u64,
2908 next_offset: u64,
2909) -> bool {
2910 if next_offset < offset {
2911 return false;
2912 }
2913 let mut expected_start = offset;
2914 for (segment_start, segment) in segments {
2915 let Some(segment_end) = read_segment_end(*segment_start, segment) else {
2916 return false;
2917 };
2918 if segment_end <= expected_start {
2919 continue;
2920 }
2921 if *segment_start > expected_start {
2922 return false;
2923 }
2924 expected_start = segment_end;
2925 if expected_start >= next_offset {
2926 return true;
2927 }
2928 }
2929 expected_start == next_offset
2930}
2931
2932fn read_segment_end(segment_start: u64, segment: &StreamReadSegment) -> Option<u64> {
2933 match segment {
2934 StreamReadSegment::Object(object) => {
2935 if object.len == 0
2936 || object.read_start_offset != segment_start
2937 || object.read_start_offset < object.object.start_offset
2938 {
2939 return None;
2940 }
2941 let len = u64::try_from(object.len).ok()?;
2942 let segment_end = object.read_start_offset.checked_add(len)?;
2943 if segment_end > object.object.end_offset {
2944 return None;
2945 }
2946 Some(segment_end)
2947 }
2948 StreamReadSegment::ColdIndex(index) => {
2949 if index.len == 0 || index.read_start_offset != segment_start {
2950 return None;
2951 }
2952 let len = u64::try_from(index.len).ok()?;
2953 segment_start.checked_add(len)
2954 }
2955 StreamReadSegment::Hot(payload) => {
2956 if payload.is_empty() {
2957 return None;
2958 }
2959 let len = u64::try_from(payload.len()).ok()?;
2960 segment_start.checked_add(len)
2961 }
2962 }
2963}
2964
2965fn message_records_cover_retained_suffix(
2966 records: &[StreamMessageRecord],
2967 retained_offset: u64,
2968 tail_offset: u64,
2969) -> bool {
2970 let mut expected_start = retained_offset;
2971 for record in records {
2972 if record.start_offset != expected_start || record.end_offset <= record.start_offset {
2973 return false;
2974 }
2975 expected_start = record.end_offset;
2976 }
2977 expected_start == tail_offset
2978}
2979
2980fn compare_stream_ids(left: &BucketStreamId, right: &BucketStreamId) -> std::cmp::Ordering {
2981 left.bucket_id
2982 .cmp(&right.bucket_id)
2983 .then_with(|| left.stream_id.cmp(&right.stream_id))
2984}
2985
2986#[cfg(test)]
2987mod tests;