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