1use std::cmp::Ordering;
17use std::cmp::Reverse;
18use std::collections::BinaryHeap;
19use std::collections::HashMap;
20use std::collections::HashSet;
21use std::collections::VecDeque;
22
23use bytes::Bytes;
24use slotmap::Key;
25use slotmap::new_key_type;
26use ursula_shard::BucketStreamId;
27
28use self::cold_gc::ColdGcQueue;
29use self::cold_state::StreamColdState;
30use self::hot_buffer::HotBuffer;
31use self::registry::StreamRegistry;
32use self::ttl::TtlEntry;
33use self::ttl::TtlIndex;
34use crate::command::StreamCommand;
35use crate::integrity::StreamIntegrity;
36use crate::model::AppendExternalInput;
37use crate::model::AppendStreamInput;
38use crate::model::BucketQuota;
39use crate::model::BucketQuotaSnapshot;
40use crate::model::BucketUsage;
41use crate::model::BucketUsageSnapshot;
42use crate::model::COLD_INDEX_PAGE_SPAN_BYTES;
43use crate::model::ColdChunkRef;
44use crate::model::ColdFlushCandidate;
45use crate::model::ColdGcEntry;
46use crate::model::ColdGcTarget;
47use crate::model::ExternalPayloadRef;
48use crate::model::HotPayloadSegment;
49use crate::model::MAX_STREAM_ATTRS_BYTES;
50use crate::model::ObjectPayloadRef;
51use crate::model::ProducerAppendRecord;
52use crate::model::ProducerReceipt;
53use crate::model::ProducerRequest;
54use crate::model::ProducerSnapshot;
55use crate::model::ProducerState;
56use crate::model::StreamAttrs;
57use crate::model::StreamBatchAppend;
58use crate::model::StreamBatchAppendItem;
59use crate::model::StreamBootstrapPlan;
60use crate::model::StreamMessageRecord;
61use crate::model::StreamMetadata;
62use crate::model::StreamRead;
63use crate::model::StreamReadColdIndexSegment;
64use crate::model::StreamReadObjectSegment;
65use crate::model::StreamReadPlan;
66use crate::model::StreamReadSegment;
67use crate::model::StreamStatus;
68use crate::model::StreamVisibleSnapshot;
69use crate::record_index::StreamRecordIndex;
70use crate::record_index::canonical_json_record_ends;
71use crate::record_index::is_json_record_content_type;
72use crate::response::StreamErrorCode;
73use crate::response::StreamErrorContext;
74use crate::response::StreamResponse;
75use crate::snapshot::StreamSnapshot;
76use crate::snapshot::StreamSnapshotEntry;
77use crate::snapshot::StreamSnapshotError;
78use crate::validate::validate_bucket_id;
79use crate::validate::validate_stream_id;
80
81mod append;
82mod cold;
83mod cold_gc;
84mod cold_state;
85mod hot_buffer;
86mod lifecycle;
87mod persist;
88mod query;
89mod registry;
90mod ttl;
91
92const TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE: usize = 256;
93pub const COMMITTED_WRITE_UNIT_BYTES: u64 = 10 * 1024;
99
100new_key_type! {
101 struct StreamKey;
102}
103
104#[derive(Debug, Clone, Default)]
105pub struct StreamStateMachine {
106 buckets: HashSet<String>,
107 registry: StreamRegistry,
108 hot_payload_bytes: u64,
111 cold_gc: ColdGcQueue,
112 shared_cold_object_refs: HashMap<String, u64>,
115 bucket_usage: HashMap<String, BucketUsage>,
119 bucket_quotas: HashMap<String, BucketQuota>,
122}
123
124#[derive(Debug, Clone)]
125struct StreamSlot {
126 metadata: StreamMetadata,
127 attrs: Option<StreamAttrs>,
128 hot_buffer: HotBuffer,
129 cold: StreamColdState,
130 message_records: Vec<StreamMessageRecord>,
131 record_index: Option<StreamRecordIndex>,
132 integrity: StreamIntegrity,
133 retained_offset: u64,
134 visible_snapshot: Option<StreamVisibleSnapshot>,
135 producers: HashMap<String, ProducerState>,
136}
137
138impl StreamStateMachine {
139 pub fn new() -> Self {
140 Self::default()
141 }
142
143 fn stream_slot(&self, stream_id: &BucketStreamId) -> Option<&StreamSlot> {
144 self.registry.slot(stream_id)
145 }
146
147 fn stream_slot_mut(&mut self, stream_id: &BucketStreamId) -> Option<&mut StreamSlot> {
148 self.registry.slot_mut(stream_id)
149 }
150
151 fn stream_metadata(&self, stream_id: &BucketStreamId) -> Option<&StreamMetadata> {
152 self.registry.metadata(stream_id)
153 }
154
155 fn retain_shared_cold_object(&mut self, path: &str) {
156 let refs = self
157 .shared_cold_object_refs
158 .entry(path.to_owned())
159 .or_default();
160 *refs = refs.saturating_add(1);
161 }
162
163 fn release_shared_cold_objects(
164 &mut self,
165 paths: impl IntoIterator<Item = String>,
166 not_before_ms: u64,
167 ) {
168 let mut reclaim = Vec::new();
169 for path in paths {
170 let Some(refs) = self.shared_cold_object_refs.get_mut(&path) else {
171 continue;
172 };
173 *refs = refs.saturating_sub(1);
174 if *refs == 0 {
175 self.shared_cold_object_refs.remove(&path);
176 reclaim.push(path);
177 }
178 }
179 if !reclaim.is_empty() {
180 self.cold_gc
181 .enqueue_after(ColdGcTarget::Paths(reclaim), not_before_ms);
182 }
183 }
184
185 fn stream_metadata_mut(&mut self, stream_id: &BucketStreamId) -> Option<&mut StreamMetadata> {
186 self.registry.metadata_mut(stream_id)
187 }
188
189 fn insert_stream_slot(&mut self, slot: StreamSlot) -> Option<StreamKey> {
190 let hot_payload_bytes = u64::try_from(slot.hot_buffer.len()).expect("payload len fits u64");
191 let key = self.registry.insert(slot)?;
192 self.hot_payload_bytes = self.hot_payload_bytes.saturating_add(hot_payload_bytes);
193 Some(key)
194 }
195
196 fn add_hot_payload_bytes(&mut self, bytes: u64) {
197 self.hot_payload_bytes = self.hot_payload_bytes.saturating_add(bytes);
198 }
199
200 fn remove_hot_payload_bytes(&mut self, bytes: u64) {
201 self.hot_payload_bytes = self.hot_payload_bytes.saturating_sub(bytes);
202 }
203
204 fn appended_record_count(record_ends: &[u64], payload_len: u64) -> u64 {
208 if !record_ends.is_empty() {
209 record_ends.len() as u64
210 } else if payload_len > 0 {
211 1
212 } else {
213 0
214 }
215 }
216
217 fn usage_mut(&mut self, bucket_id: &str) -> &mut BucketUsage {
218 self.bucket_usage.entry(bucket_id.to_owned()).or_default()
219 }
220
221 fn usage_on_append(&mut self, bucket_id: &str, payload_bytes: u64, records: u64) {
224 let usage = self.usage_mut(bucket_id);
225 usage.committed_append_bytes = usage.committed_append_bytes.saturating_add(payload_bytes);
226 usage.committed_records = usage.committed_records.saturating_add(records);
227 usage.committed_write_units = usage
228 .committed_write_units
229 .saturating_add(payload_bytes.div_ceil(COMMITTED_WRITE_UNIT_BYTES).max(1));
230 usage.retained_bytes = usage.retained_bytes.saturating_add(payload_bytes);
231 }
232
233 fn usage_on_stream_created(&mut self, bucket_id: &str, initial_bytes: u64, records: u64) {
236 let usage = self.usage_mut(bucket_id);
237 usage.stream_count = usage.stream_count.saturating_add(1);
238 usage.committed_append_bytes = usage.committed_append_bytes.saturating_add(initial_bytes);
239 usage.committed_records = usage.committed_records.saturating_add(records);
240 usage.committed_write_units = usage
241 .committed_write_units
242 .saturating_add(initial_bytes.div_ceil(COMMITTED_WRITE_UNIT_BYTES).max(1));
243 usage.retained_bytes = usage.retained_bytes.saturating_add(initial_bytes);
244 }
245
246 fn usage_on_retention(&mut self, bucket_id: &str, reclaimed_bytes: u64) {
248 let usage = self.usage_mut(bucket_id);
249 usage.retained_bytes = usage.retained_bytes.saturating_sub(reclaimed_bytes);
250 }
251
252 fn usage_on_stream_removed(&mut self, bucket_id: &str, retained_bytes: u64) {
255 let usage = self.usage_mut(bucket_id);
256 usage.stream_count = usage.stream_count.saturating_sub(1);
257 usage.retained_bytes = usage.retained_bytes.saturating_sub(retained_bytes);
258 }
259
260 fn set_bucket_quota(
263 &mut self,
264 bucket_id: String,
265 max_streams: Option<u64>,
266 max_retained_bytes: Option<u64>,
267 ) -> StreamResponse {
268 if let Err(message) = validate_bucket_id(&bucket_id) {
269 return StreamResponse::error(StreamErrorCode::InvalidBucketId, message);
270 }
271 let quota = BucketQuota {
272 max_streams,
273 max_retained_bytes,
274 };
275 if quota.is_unlimited() {
276 self.bucket_quotas.remove(&bucket_id);
277 } else {
278 self.bucket_quotas.insert(bucket_id.clone(), quota);
279 }
280 StreamResponse::BucketQuotaSet { bucket_id }
281 }
282
283 fn check_create_quota(
288 &self,
289 bucket_id: &str,
290 initial_bytes: u64,
291 ) -> Result<(), StreamResponse> {
292 let Some(quota) = self.bucket_quotas.get(bucket_id) else {
293 return Ok(());
294 };
295 let usage = self
296 .bucket_usage
297 .get(bucket_id)
298 .copied()
299 .unwrap_or_default();
300 if let Some(max_streams) = quota.max_streams
301 && usage.stream_count >= max_streams
302 {
303 return Err(StreamResponse::error(
304 StreamErrorCode::QuotaExceeded,
305 format!(
306 "bucket '{bucket_id}' stream-count quota exceeded in this group ({max_streams} max)"
307 ),
308 ));
309 }
310 self.check_retained_quota_inner(bucket_id, quota, &usage, initial_bytes)
311 }
312
313 fn check_append_quota(
318 &self,
319 bucket_id: &str,
320 payload_bytes: u64,
321 ) -> Result<(), StreamResponse> {
322 if payload_bytes == 0 {
323 return Ok(());
324 }
325 let Some(quota) = self.bucket_quotas.get(bucket_id) else {
326 return Ok(());
327 };
328 let usage = self
329 .bucket_usage
330 .get(bucket_id)
331 .copied()
332 .unwrap_or_default();
333 self.check_retained_quota_inner(bucket_id, quota, &usage, payload_bytes)
334 }
335
336 fn check_retained_quota_inner(
337 &self,
338 bucket_id: &str,
339 quota: &BucketQuota,
340 usage: &BucketUsage,
341 incoming_bytes: u64,
342 ) -> Result<(), StreamResponse> {
343 if let Some(max_retained) = quota.max_retained_bytes
344 && usage.retained_bytes.saturating_add(incoming_bytes) > max_retained
345 {
346 return Err(StreamResponse::error(
347 StreamErrorCode::QuotaExceeded,
348 format!(
349 "bucket '{bucket_id}' retained-bytes quota exceeded in this group ({max_retained} max)"
350 ),
351 ));
352 }
353 Ok(())
354 }
355
356 pub fn bucket_quota_report(&self) -> Vec<BucketQuotaSnapshot> {
359 let mut report = self
360 .bucket_quotas
361 .iter()
362 .map(|(bucket_id, quota)| BucketQuotaSnapshot {
363 bucket_id: bucket_id.clone(),
364 quota: *quota,
365 })
366 .collect::<Vec<_>>();
367 report.sort_by(|left, right| left.bucket_id.cmp(&right.bucket_id));
368 report
369 }
370
371 pub fn bucket_usage_report(&self) -> Vec<BucketUsageSnapshot> {
374 let mut report = self
375 .bucket_usage
376 .iter()
377 .map(|(bucket_id, usage)| BucketUsageSnapshot {
378 bucket_id: bucket_id.clone(),
379 usage: *usage,
380 })
381 .collect::<Vec<_>>();
382 report.sort_by(|left, right| left.bucket_id.cmp(&right.bucket_id));
383 report
384 }
385
386 fn refresh_ttl_entry(&mut self, stream_id: &BucketStreamId) {
387 self.registry.refresh_ttl(stream_id);
388 }
389
390 fn message_records_for_append(
391 start_offset: u64,
392 end_offset: u64,
393 record_ends: &[u64],
394 ) -> Vec<StreamMessageRecord> {
395 if record_ends.is_empty() {
396 return (start_offset < end_offset)
397 .then_some(StreamMessageRecord {
398 start_offset,
399 end_offset,
400 })
401 .into_iter()
402 .collect();
403 }
404 let mut start = start_offset;
405 record_ends
406 .iter()
407 .map(|relative_end| {
408 let end = start_offset.saturating_add(*relative_end);
409 let record = StreamMessageRecord {
410 start_offset: start,
411 end_offset: end,
412 };
413 start = end;
414 record
415 })
416 .collect()
417 }
418
419 pub fn apply(&mut self, command: StreamCommand) -> StreamResponse {
420 match command {
421 StreamCommand::CreateBucket { bucket_id } => self.create_bucket(bucket_id),
422 StreamCommand::DeleteBucket { bucket_id } => self.delete_bucket(&bucket_id),
423 StreamCommand::CreateStream {
424 stream_id,
425 content_type,
426 initial_payload,
427 close_after,
428 stream_seq,
429 producer,
430 stream_ttl_seconds,
431 stream_expires_at_ms,
432 attrs,
433 now_ms,
434 } => {
435 let response = match canonical_json_record_ends(&content_type, &initial_payload) {
436 Ok(record_ends) => self.create_stream(CreateStreamInput {
437 stream_id,
438 content_type,
439 initial_payload: initial_payload.into(),
440 record_ends,
441 close_after,
442 stream_seq,
443 producer,
444 stream_ttl_seconds,
445 stream_expires_at_ms,
446 attrs,
447 now_ms,
448 }),
449 Err(_) => StreamResponse::error(
450 StreamErrorCode::InvalidRecordBoundaries,
451 "application/json initial payload must use canonical newline boundaries",
452 ),
453 };
454 self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
455 response
456 }
457 StreamCommand::CreateExternal {
458 stream_id,
459 content_type,
460 initial_payload,
461 record_ends,
462 close_after,
463 stream_seq,
464 producer,
465 stream_ttl_seconds,
466 stream_expires_at_ms,
467 attrs,
468 now_ms,
469 } => {
470 let response = self.create_external_stream(CreateExternalStreamInput {
471 stream_id,
472 content_type,
473 initial_payload,
474 record_ends,
475 close_after,
476 stream_seq,
477 producer,
478 stream_ttl_seconds,
479 stream_expires_at_ms,
480 attrs,
481 now_ms,
482 });
483 self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
484 response
485 }
486 StreamCommand::Append {
487 stream_id,
488 content_type,
489 payload,
490 close_after,
491 stream_seq,
492 producer,
493 now_ms,
494 record_match,
495 } => {
496 let response = self.append_borrowed(AppendStreamInput {
497 stream_id,
498 content_type: content_type.as_deref(),
499 payload: &payload,
500 close_after,
501 stream_seq,
502 producer,
503 now_ms,
504 record_match,
505 });
506 self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
507 response
508 }
509 StreamCommand::AppendExternal {
510 stream_id,
511 content_type,
512 payload,
513 record_ends,
514 close_after,
515 stream_seq,
516 producer,
517 now_ms,
518 record_match,
519 } => {
520 let response = self.append_external(AppendExternalInput {
521 stream_id,
522 content_type: content_type.as_deref(),
523 payload,
524 record_ends,
525 close_after,
526 stream_seq,
527 producer,
528 now_ms,
529 record_match,
530 });
531 self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
532 response
533 }
534 StreamCommand::AppendBatch {
535 stream_id,
536 content_type,
537 payloads,
538 producer,
539 now_ms,
540 } => {
541 let response = match self.append_batch_borrowed(
542 stream_id,
543 content_type.as_deref(),
544 &payloads.iter().map(Bytes::as_ref).collect::<Vec<_>>(),
545 producer,
546 now_ms,
547 ) {
548 Ok(batch) => batch
549 .items
550 .last()
551 .map(|item| StreamResponse::Appended {
552 offset: item.offset,
553 next_offset: item.next_offset,
554 closed: item.closed,
555 deduplicated: item.deduplicated,
556 producer: None,
557 })
558 .unwrap_or_else(|| {
559 StreamResponse::error(
560 StreamErrorCode::EmptyAppend,
561 "append batch must contain at least one payload",
562 )
563 }),
564 Err(response) => response,
565 };
566 self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
567 response
568 }
569 StreamCommand::PublishSnapshot {
570 stream_id,
571 snapshot_offset,
572 content_type,
573 payload,
574 expected_digest,
575 now_ms,
576 } => {
577 let response = self.publish_snapshot(
578 stream_id,
579 snapshot_offset,
580 content_type,
581 payload.into(),
582 expected_digest,
583 now_ms,
584 );
585 self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
586 response
587 }
588 StreamCommand::AdvanceRetention {
589 stream_id,
590 retained_offset,
591 now_ms,
592 } => {
593 let response = self.advance_retention(stream_id, retained_offset, now_ms);
594 self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
595 response
596 }
597 StreamCommand::TouchStreamAccess {
598 stream_id,
599 now_ms,
600 renew_ttl,
601 } => {
602 let response = self.touch_stream_access(&stream_id, now_ms, renew_ttl);
603 self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
604 response
605 }
606 StreamCommand::UpdateStreamAttrs {
607 stream_id,
608 attrs,
609 now_ms,
610 } => {
611 let response = self.update_stream_attrs(&stream_id, attrs, now_ms);
612 self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
613 response
614 }
615 StreamCommand::FlushCold { stream_id, chunk } => self.flush_cold(stream_id, chunk),
616 StreamCommand::CompactCold {
617 stream_id,
618 old_chunks,
619 replacement,
620 gc_not_before_ms,
621 } => self.compact_cold(stream_id, old_chunks, replacement, gc_not_before_ms),
622 StreamCommand::Close {
623 stream_id,
624 stream_seq,
625 producer,
626 now_ms,
627 } => {
628 let response = self.close(stream_id, stream_seq, producer, now_ms);
629 self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
630 response
631 }
632 StreamCommand::DeleteStream { stream_id } => self.delete_stream(&stream_id),
633 StreamCommand::PurgeBucket { bucket_id } => self.purge_bucket(&bucket_id),
634 StreamCommand::AckColdGc { up_to_seq } => self.ack_cold_gc(up_to_seq),
635 StreamCommand::ImportSnapshot { snapshot } => self.import_snapshot(*snapshot),
636 StreamCommand::SetBucketQuota {
637 bucket_id,
638 max_streams,
639 max_retained_bytes,
640 } => self.set_bucket_quota(bucket_id, max_streams, max_retained_bytes),
641 }
642 }
643}
644
645#[derive(Debug)]
646struct CreateStreamInput {
647 stream_id: BucketStreamId,
648 content_type: String,
649 initial_payload: Vec<u8>,
650 record_ends: Vec<u64>,
651 close_after: bool,
652 stream_seq: Option<String>,
653 producer: Option<ProducerRequest>,
654 stream_ttl_seconds: Option<u64>,
655 stream_expires_at_ms: Option<u64>,
656 attrs: Option<StreamAttrs>,
657 now_ms: u64,
658}
659
660#[derive(Debug)]
661struct CreateExternalStreamInput {
662 stream_id: BucketStreamId,
663 content_type: String,
664 initial_payload: ExternalPayloadRef,
665 record_ends: Vec<u64>,
666 close_after: bool,
667 stream_seq: Option<String>,
668 producer: Option<ProducerRequest>,
669 stream_ttl_seconds: Option<u64>,
670 stream_expires_at_ms: Option<u64>,
671 attrs: Option<StreamAttrs>,
672 now_ms: u64,
673}
674
675impl CreateStreamInput {
676 fn initial_len(&self) -> u64 {
677 u64::try_from(self.initial_payload.len()).expect("payload len fits u64")
678 }
679}
680
681fn normalize_stream_attrs(attrs: Option<StreamAttrs>) -> Option<StreamAttrs> {
682 attrs.filter(|attrs| !attrs.is_empty())
683}
684
685fn stream_expiry_at_ms(stream: &StreamMetadata) -> Option<u64> {
686 if let Some(expires_at_ms) = stream.stream_expires_at_ms {
687 return Some(expires_at_ms);
688 }
689 stream.stream_ttl_seconds.map(|ttl_seconds| {
690 stream
691 .last_ttl_touch_at_ms
692 .saturating_add(ttl_seconds.saturating_mul(1000))
693 })
694}
695
696fn stream_is_expired(stream: &StreamMetadata, now_ms: u64) -> bool {
697 stream_expiry_at_ms(stream).is_some_and(|expires_at_ms| now_ms >= expires_at_ms)
698}
699
700fn stream_ttl_renewal_due(stream: &StreamMetadata, now_ms: u64) -> bool {
701 let Some(ttl_seconds) = stream.stream_ttl_seconds else {
702 return false;
703 };
704 if stream.stream_expires_at_ms.is_some() {
705 return false;
706 }
707 let ttl_ms = ttl_seconds.saturating_mul(1000);
708 let renewal_interval_ms = ttl_ms.div_ceil(4).max(1);
709 now_ms.saturating_sub(stream.last_ttl_touch_at_ms) >= renewal_interval_ms
710}
711
712fn renew_stream_ttl(stream: &mut StreamMetadata, now_ms: u64) {
713 if stream.stream_ttl_seconds.is_some() && stream.stream_expires_at_ms.is_none() {
714 stream.last_ttl_touch_at_ms = now_ms;
715 }
716}
717
718fn validate_producer_request(producer: Option<&ProducerRequest>) -> Result<(), StreamResponse> {
719 let Some(producer) = producer else {
720 return Ok(());
721 };
722 if producer.producer_id.trim().is_empty() {
723 return Err(StreamResponse::error(
724 StreamErrorCode::InvalidProducer,
725 "producer id must not be empty",
726 ));
727 }
728 const MAX_JS_SAFE_INTEGER: u64 = 9_007_199_254_740_991;
729 if producer.producer_epoch > MAX_JS_SAFE_INTEGER {
730 return Err(StreamResponse::error(
731 StreamErrorCode::InvalidProducer,
732 format!(
733 "producer epoch {} exceeds maximum {}",
734 producer.producer_epoch, MAX_JS_SAFE_INTEGER
735 ),
736 ));
737 }
738 if producer.producer_seq > MAX_JS_SAFE_INTEGER {
739 return Err(StreamResponse::error(
740 StreamErrorCode::InvalidProducer,
741 format!(
742 "producer sequence {} exceeds maximum {}",
743 producer.producer_seq, MAX_JS_SAFE_INTEGER
744 ),
745 ));
746 }
747 Ok(())
748}
749
750fn validate_external_payload_ref(payload: &ExternalPayloadRef) -> Result<(), StreamResponse> {
751 if payload.s3_path.trim().is_empty() {
752 return Err(StreamResponse::error(
753 StreamErrorCode::InvalidColdFlush,
754 "external payload S3 path must not be empty",
755 ));
756 }
757 if payload.payload_len == 0 {
758 return Err(StreamResponse::error(
759 StreamErrorCode::EmptyAppend,
760 "external payload length must be greater than zero",
761 ));
762 }
763 if payload.object_size < payload.payload_len {
764 return Err(StreamResponse::error(
765 StreamErrorCode::InvalidColdFlush,
766 "external payload object size must cover payload length",
767 ));
768 }
769 Ok(())
770}
771
772fn build_record_index(
773 content_type: &str,
774 payload_len: u64,
775 record_ends: &[u64],
776) -> Result<Option<StreamRecordIndex>, StreamResponse> {
777 if !is_json_record_content_type(content_type) {
778 return record_ends.is_empty().then_some(None).ok_or_else(|| {
779 StreamResponse::error(
780 StreamErrorCode::InvalidRecordBoundaries,
781 "record boundaries are only valid for application/json streams",
782 )
783 });
784 }
785 if payload_len > 0 && record_ends.is_empty() {
786 return Ok(None);
790 }
791 let mut index = StreamRecordIndex::new();
792 index
793 .append_relative_ends(0, payload_len, record_ends)
794 .map_err(|_| {
795 StreamResponse::error(
796 StreamErrorCode::InvalidRecordBoundaries,
797 "record boundaries do not match the canonical JSON payload",
798 )
799 })?;
800 Ok(Some(index))
801}
802
803fn prepare_record_append(
804 current: Option<&StreamRecordIndex>,
805 json_stream: bool,
806 base_offset: u64,
807 payload_len: u64,
808 record_ends: &[u64],
809) -> Result<Option<crate::PreparedRecordAppend>, StreamResponse> {
810 let Some(current) = current else {
811 if json_stream {
812 return Ok(None);
813 }
814 return record_ends.is_empty().then_some(None).ok_or_else(|| {
815 StreamResponse::error(
816 StreamErrorCode::InvalidRecordBoundaries,
817 "binary streams cannot carry JSON record boundaries",
818 )
819 });
820 };
821 current
822 .prepare_append(base_offset, payload_len, record_ends)
823 .map(Some)
824 .map_err(|_| {
825 StreamResponse::error(
826 StreamErrorCode::InvalidRecordBoundaries,
827 "record boundaries do not match the canonical JSON payload",
828 )
829 })
830}
831
832fn compare_stream_ids(left: &BucketStreamId, right: &BucketStreamId) -> std::cmp::Ordering {
833 left.bucket_id
834 .cmp(&right.bucket_id)
835 .then_with(|| left.stream_id.cmp(&right.stream_id))
836}
837
838fn snapshot_digest(content_type: &str, payload: &[u8]) -> String {
839 let mut hasher = blake3::Hasher::new();
840 hasher.update(&(content_type.len() as u64).to_le_bytes());
841 hasher.update(content_type.as_bytes());
842 hasher.update(payload);
843 hasher.finalize().to_hex().to_string()
844}
845
846#[cfg(test)]
847mod tests;