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    /// Permanent tenant-erasure fences. A purged bucket name can never be
108    /// reused, including after snapshot restore, so bytes cannot reappear
109    /// behind an already-issued physical absence proof.
110    erased_buckets: HashSet<String>,
111    registry: StreamRegistry,
112    /// Group-wide hot payload gauge. Kept incrementally so append admission
113    /// and responses do not scan every stream in the group.
114    hot_payload_bytes: u64,
115    cold_gc: ColdGcQueue,
116    /// Live logical references to group-scoped shared cold objects. This is
117    /// derived from per-stream cold refs when snapshots are restored.
118    shared_cold_object_refs: HashMap<String, u64>,
119    /// Every bucket whose bytes have ever occupied a still-live shared object.
120    /// Keep owners after an individual bucket releases its references so the
121    /// eventual physical-delete work is attributed to every erasure proof.
122    shared_cold_object_owners: HashMap<String, HashSet<String>>,
123    /// Per-bucket committed usage for this group; see [`BucketUsage`] for the
124    /// monotonic-versus-gauge split. Mutated only by the accounting helpers
125    /// below so every counter change stays deterministic and auditable.
126    bucket_usage: HashMap<String, BucketUsage>,
127    /// Per-bucket data-plane quota backstops enforced against this group's
128    /// local counters; see [`BucketQuota`] for the enforcement semantics.
129    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    /// Records committed by one accepted append. JSON streams provide exact
233    /// canonical boundaries; a byte stream counts one message record per
234    /// non-empty append.
235    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    /// One accepted (non-deduplicated) append: monotonic counters grow and
250    /// the retained gauge grows by the same bytes.
251    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    /// A newly created stream, including any initial payload it was created
262    /// with.
263    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    /// Destructive retention reclaimed `reclaimed_bytes` of logical prefix.
275    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    /// A stream left the registry (delete or TTL expiry); its remaining
281    /// retained bytes leave the gauge with it.
282    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    /// Sets or clears the quota record for one bucket. Both limits `None`
289    /// removes the record so cleared quotas leave no residue in snapshots.
290    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    /// Data-plane backstop for stream creation: this group's local stream
318    /// count and the incoming initial payload must fit under the bucket's
319    /// quota. Runs after the idempotent already-exists paths so replays of
320    /// accepted creates never fail retroactively.
321    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    /// Data-plane backstop for appends: the payload must fit under the
348    /// bucket's retained-bytes quota against this group's local gauge.
349    /// Producer-deduplicated retries return before this check, so an
350    /// accepted append replay can never fail retroactively.
351    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    /// Current per-bucket quotas for this group, sorted for deterministic
391    /// output.
392    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    /// Current per-bucket usage for this group, sorted for deterministic
406    /// output.
407    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        // Pre-extension WAL and snapshot entries have no boundary metadata.
821        // Keep those JSON streams readable without activating coordinates
822        // part-way through their history.
823        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;