Skip to main content

ursula_stream/
state_machine.rs

1//! Deterministic stream state machine driving a single Raft group.
2//!
3//! The state machine lives in this module root; its behavior is split across
4//! cohesive submodules to keep each surface readable:
5//!
6//! - [`query`]: read paths — heads, accessors, read plans, snapshots, bootstrap.
7//! - [`append`]: append paths and idempotent producer bookkeeping.
8//! - [`lifecycle`]: bucket/stream create, close, delete, attrs, and TTL expiry.
9//! - [`cold`]: cold-tier flush planning, GC, retention compaction, snapshot publishing.
10//! - [`persist`]: snapshot / restore / integrity serialization.
11//! - [`hot_buffer`], [`cold_state`], [`ttl`]: internal per-stream data structures.
12//!
13//! The root keeps the [`StreamStateMachine`] type, its core slot/TTL accessors,
14//! the [`StreamStateMachine::apply`] command dispatcher, and cross-cutting helpers.
15
16use 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;
93/// Size of the self-described derived write unit exported by Ursula.
94///
95/// Raw committed bytes and records remain available beside it. Keeping the
96/// unit in the usage contract, rather than its field name, lets consumers
97/// validate the interpretation before using the derived counter.
98pub 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    /// Group-wide hot payload gauge. Kept incrementally so append admission
109    /// and responses do not scan every stream in the group.
110    hot_payload_bytes: u64,
111    cold_gc: ColdGcQueue,
112    /// Live logical references to group-scoped shared cold objects. This is
113    /// derived from per-stream cold refs when snapshots are restored.
114    shared_cold_object_refs: HashMap<String, u64>,
115    /// Per-bucket committed usage for this group; see [`BucketUsage`] for the
116    /// monotonic-versus-gauge split. Mutated only by the accounting helpers
117    /// below so every counter change stays deterministic and auditable.
118    bucket_usage: HashMap<String, BucketUsage>,
119    /// Per-bucket data-plane quota backstops enforced against this group's
120    /// local counters; see [`BucketQuota`] for the enforcement semantics.
121    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    /// Records committed by one accepted append. JSON streams provide exact
205    /// canonical boundaries; a byte stream counts one message record per
206    /// non-empty append.
207    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    /// One accepted (non-deduplicated) append: monotonic counters grow and
222    /// the retained gauge grows by the same bytes.
223    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    /// A newly created stream, including any initial payload it was created
234    /// with.
235    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    /// Destructive retention reclaimed `reclaimed_bytes` of logical prefix.
247    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    /// A stream left the registry (delete or TTL expiry); its remaining
253    /// retained bytes leave the gauge with it.
254    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    /// Sets or clears the quota record for one bucket. Both limits `None`
261    /// removes the record so cleared quotas leave no residue in snapshots.
262    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    /// Data-plane backstop for stream creation: this group's local stream
284    /// count and the incoming initial payload must fit under the bucket's
285    /// quota. Runs after the idempotent already-exists paths so replays of
286    /// accepted creates never fail retroactively.
287    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    /// Data-plane backstop for appends: the payload must fit under the
314    /// bucket's retained-bytes quota against this group's local gauge.
315    /// Producer-deduplicated retries return before this check, so an
316    /// accepted append replay can never fail retroactively.
317    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    /// Current per-bucket quotas for this group, sorted for deterministic
357    /// output.
358    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    /// Current per-bucket usage for this group, sorted for deterministic
372    /// output.
373    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        // Pre-extension WAL and snapshot entries have no boundary metadata.
787        // Keep those JSON streams readable without activating coordinates
788        // part-way through their history.
789        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;