use std::cmp::Ordering;
use std::cmp::Reverse;
use std::collections::BinaryHeap;
use std::collections::HashMap;
use std::collections::HashSet;
use std::collections::VecDeque;
use bytes::Bytes;
use slotmap::Key;
use slotmap::new_key_type;
use ursula_shard::BucketStreamId;
use self::cold_gc::ColdGcQueue;
use self::cold_state::StreamColdState;
use self::hot_buffer::HotBuffer;
use self::registry::StreamRegistry;
use self::ttl::TtlEntry;
use self::ttl::TtlIndex;
use crate::command::StreamCommand;
use crate::integrity::StreamIntegrity;
use crate::model::AppendExternalInput;
use crate::model::AppendStreamInput;
use crate::model::BucketQuota;
use crate::model::BucketQuotaSnapshot;
use crate::model::BucketUsage;
use crate::model::BucketUsageSnapshot;
use crate::model::COLD_INDEX_PAGE_SPAN_BYTES;
use crate::model::ColdChunkRef;
use crate::model::ColdFlushCandidate;
use crate::model::ColdGcEntry;
use crate::model::ColdGcTarget;
use crate::model::ExternalPayloadRef;
use crate::model::HotPayloadSegment;
use crate::model::MAX_STREAM_ATTRS_BYTES;
use crate::model::ObjectPayloadRef;
use crate::model::ProducerAppendRecord;
use crate::model::ProducerReceipt;
use crate::model::ProducerRequest;
use crate::model::ProducerSnapshot;
use crate::model::ProducerState;
use crate::model::StreamAttrs;
use crate::model::StreamBatchAppend;
use crate::model::StreamBatchAppendItem;
use crate::model::StreamBootstrapPlan;
use crate::model::StreamMessageRecord;
use crate::model::StreamMetadata;
use crate::model::StreamRead;
use crate::model::StreamReadColdIndexSegment;
use crate::model::StreamReadObjectSegment;
use crate::model::StreamReadPlan;
use crate::model::StreamReadSegment;
use crate::model::StreamStatus;
use crate::model::StreamVisibleSnapshot;
use crate::record_index::StreamRecordIndex;
use crate::record_index::canonical_json_record_ends;
use crate::record_index::is_json_record_content_type;
use crate::response::StreamErrorCode;
use crate::response::StreamErrorContext;
use crate::response::StreamResponse;
use crate::snapshot::StreamSnapshot;
use crate::snapshot::StreamSnapshotEntry;
use crate::snapshot::StreamSnapshotError;
use crate::validate::validate_bucket_id;
use crate::validate::validate_stream_id;
mod append;
mod cold;
mod cold_gc;
mod cold_state;
mod hot_buffer;
mod lifecycle;
mod persist;
mod query;
mod registry;
mod ttl;
const TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE: usize = 256;
new_key_type! {
struct StreamKey;
}
#[derive(Debug, Clone, Default)]
pub struct StreamStateMachine {
buckets: HashSet<String>,
registry: StreamRegistry,
hot_payload_bytes: u64,
cold_gc: ColdGcQueue,
shared_cold_object_refs: HashMap<String, u64>,
bucket_usage: HashMap<String, BucketUsage>,
bucket_quotas: HashMap<String, BucketQuota>,
}
#[derive(Debug, Clone)]
struct StreamSlot {
metadata: StreamMetadata,
attrs: Option<StreamAttrs>,
hot_buffer: HotBuffer,
cold: StreamColdState,
message_records: Vec<StreamMessageRecord>,
record_index: Option<StreamRecordIndex>,
integrity: StreamIntegrity,
retained_offset: u64,
visible_snapshot: Option<StreamVisibleSnapshot>,
producers: HashMap<String, ProducerState>,
}
impl StreamStateMachine {
pub fn new() -> Self {
Self::default()
}
fn stream_slot(&self, stream_id: &BucketStreamId) -> Option<&StreamSlot> {
self.registry.slot(stream_id)
}
fn stream_slot_mut(&mut self, stream_id: &BucketStreamId) -> Option<&mut StreamSlot> {
self.registry.slot_mut(stream_id)
}
fn stream_metadata(&self, stream_id: &BucketStreamId) -> Option<&StreamMetadata> {
self.registry.metadata(stream_id)
}
fn retain_shared_cold_object(&mut self, path: &str) {
let refs = self
.shared_cold_object_refs
.entry(path.to_owned())
.or_default();
*refs = refs.saturating_add(1);
}
fn release_shared_cold_objects(
&mut self,
paths: impl IntoIterator<Item = String>,
not_before_ms: u64,
) {
let mut reclaim = Vec::new();
for path in paths {
let Some(refs) = self.shared_cold_object_refs.get_mut(&path) else {
continue;
};
*refs = refs.saturating_sub(1);
if *refs == 0 {
self.shared_cold_object_refs.remove(&path);
reclaim.push(path);
}
}
if !reclaim.is_empty() {
self.cold_gc
.enqueue_after(ColdGcTarget::Paths(reclaim), not_before_ms);
}
}
fn stream_metadata_mut(&mut self, stream_id: &BucketStreamId) -> Option<&mut StreamMetadata> {
self.registry.metadata_mut(stream_id)
}
fn insert_stream_slot(&mut self, slot: StreamSlot) -> Option<StreamKey> {
let hot_payload_bytes = u64::try_from(slot.hot_buffer.len()).expect("payload len fits u64");
let key = self.registry.insert(slot)?;
self.hot_payload_bytes = self.hot_payload_bytes.saturating_add(hot_payload_bytes);
Some(key)
}
fn add_hot_payload_bytes(&mut self, bytes: u64) {
self.hot_payload_bytes = self.hot_payload_bytes.saturating_add(bytes);
}
fn remove_hot_payload_bytes(&mut self, bytes: u64) {
self.hot_payload_bytes = self.hot_payload_bytes.saturating_sub(bytes);
}
fn appended_record_count(record_ends: &[u64], payload_len: u64) -> u64 {
if !record_ends.is_empty() {
record_ends.len() as u64
} else if payload_len > 0 {
1
} else {
0
}
}
fn usage_mut(&mut self, bucket_id: &str) -> &mut BucketUsage {
self.bucket_usage.entry(bucket_id.to_owned()).or_default()
}
fn usage_on_append(&mut self, bucket_id: &str, payload_bytes: u64, records: u64) {
let usage = self.usage_mut(bucket_id);
usage.committed_append_bytes = usage.committed_append_bytes.saturating_add(payload_bytes);
usage.committed_records = usage.committed_records.saturating_add(records);
usage.retained_bytes = usage.retained_bytes.saturating_add(payload_bytes);
}
fn usage_on_stream_created(&mut self, bucket_id: &str, initial_bytes: u64, records: u64) {
let usage = self.usage_mut(bucket_id);
usage.stream_count = usage.stream_count.saturating_add(1);
usage.committed_append_bytes = usage.committed_append_bytes.saturating_add(initial_bytes);
usage.committed_records = usage.committed_records.saturating_add(records);
usage.retained_bytes = usage.retained_bytes.saturating_add(initial_bytes);
}
fn usage_on_retention(&mut self, bucket_id: &str, reclaimed_bytes: u64) {
let usage = self.usage_mut(bucket_id);
usage.retained_bytes = usage.retained_bytes.saturating_sub(reclaimed_bytes);
}
fn usage_on_stream_removed(&mut self, bucket_id: &str, retained_bytes: u64) {
let usage = self.usage_mut(bucket_id);
usage.stream_count = usage.stream_count.saturating_sub(1);
usage.retained_bytes = usage.retained_bytes.saturating_sub(retained_bytes);
}
fn set_bucket_quota(
&mut self,
bucket_id: String,
max_streams: Option<u64>,
max_retained_bytes: Option<u64>,
) -> StreamResponse {
if let Err(message) = validate_bucket_id(&bucket_id) {
return StreamResponse::error(StreamErrorCode::InvalidBucketId, message);
}
let quota = BucketQuota {
max_streams,
max_retained_bytes,
};
if quota.is_unlimited() {
self.bucket_quotas.remove(&bucket_id);
} else {
self.bucket_quotas.insert(bucket_id.clone(), quota);
}
StreamResponse::BucketQuotaSet { bucket_id }
}
fn check_create_quota(
&self,
bucket_id: &str,
initial_bytes: u64,
) -> Result<(), StreamResponse> {
let Some(quota) = self.bucket_quotas.get(bucket_id) else {
return Ok(());
};
let usage = self
.bucket_usage
.get(bucket_id)
.copied()
.unwrap_or_default();
if let Some(max_streams) = quota.max_streams
&& usage.stream_count >= max_streams
{
return Err(StreamResponse::error(
StreamErrorCode::QuotaExceeded,
format!(
"bucket '{bucket_id}' stream-count quota exceeded in this group ({max_streams} max)"
),
));
}
self.check_retained_quota_inner(bucket_id, quota, &usage, initial_bytes)
}
fn check_append_quota(
&self,
bucket_id: &str,
payload_bytes: u64,
) -> Result<(), StreamResponse> {
if payload_bytes == 0 {
return Ok(());
}
let Some(quota) = self.bucket_quotas.get(bucket_id) else {
return Ok(());
};
let usage = self
.bucket_usage
.get(bucket_id)
.copied()
.unwrap_or_default();
self.check_retained_quota_inner(bucket_id, quota, &usage, payload_bytes)
}
fn check_retained_quota_inner(
&self,
bucket_id: &str,
quota: &BucketQuota,
usage: &BucketUsage,
incoming_bytes: u64,
) -> Result<(), StreamResponse> {
if let Some(max_retained) = quota.max_retained_bytes
&& usage.retained_bytes.saturating_add(incoming_bytes) > max_retained
{
return Err(StreamResponse::error(
StreamErrorCode::QuotaExceeded,
format!(
"bucket '{bucket_id}' retained-bytes quota exceeded in this group ({max_retained} max)"
),
));
}
Ok(())
}
pub fn bucket_quota_report(&self) -> Vec<BucketQuotaSnapshot> {
let mut report = self
.bucket_quotas
.iter()
.map(|(bucket_id, quota)| BucketQuotaSnapshot {
bucket_id: bucket_id.clone(),
quota: *quota,
})
.collect::<Vec<_>>();
report.sort_by(|left, right| left.bucket_id.cmp(&right.bucket_id));
report
}
pub fn bucket_usage_report(&self) -> Vec<BucketUsageSnapshot> {
let mut report = self
.bucket_usage
.iter()
.map(|(bucket_id, usage)| BucketUsageSnapshot {
bucket_id: bucket_id.clone(),
usage: *usage,
})
.collect::<Vec<_>>();
report.sort_by(|left, right| left.bucket_id.cmp(&right.bucket_id));
report
}
fn refresh_ttl_entry(&mut self, stream_id: &BucketStreamId) {
self.registry.refresh_ttl(stream_id);
}
fn message_records_for_append(
start_offset: u64,
end_offset: u64,
record_ends: &[u64],
) -> Vec<StreamMessageRecord> {
if record_ends.is_empty() {
return (start_offset < end_offset)
.then_some(StreamMessageRecord {
start_offset,
end_offset,
})
.into_iter()
.collect();
}
let mut start = start_offset;
record_ends
.iter()
.map(|relative_end| {
let end = start_offset.saturating_add(*relative_end);
let record = StreamMessageRecord {
start_offset: start,
end_offset: end,
};
start = end;
record
})
.collect()
}
pub fn apply(&mut self, command: StreamCommand) -> StreamResponse {
match command {
StreamCommand::CreateBucket { bucket_id } => self.create_bucket(bucket_id),
StreamCommand::DeleteBucket { bucket_id } => self.delete_bucket(&bucket_id),
StreamCommand::CreateStream {
stream_id,
content_type,
initial_payload,
close_after,
stream_seq,
producer,
stream_ttl_seconds,
stream_expires_at_ms,
attrs,
now_ms,
} => {
let response = match canonical_json_record_ends(&content_type, &initial_payload) {
Ok(record_ends) => self.create_stream(CreateStreamInput {
stream_id,
content_type,
initial_payload: initial_payload.into(),
record_ends,
close_after,
stream_seq,
producer,
stream_ttl_seconds,
stream_expires_at_ms,
attrs,
now_ms,
}),
Err(_) => StreamResponse::error(
StreamErrorCode::InvalidRecordBoundaries,
"application/json initial payload must use canonical newline boundaries",
),
};
self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
response
}
StreamCommand::CreateExternal {
stream_id,
content_type,
initial_payload,
record_ends,
close_after,
stream_seq,
producer,
stream_ttl_seconds,
stream_expires_at_ms,
attrs,
now_ms,
} => {
let response = self.create_external_stream(CreateExternalStreamInput {
stream_id,
content_type,
initial_payload,
record_ends,
close_after,
stream_seq,
producer,
stream_ttl_seconds,
stream_expires_at_ms,
attrs,
now_ms,
});
self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
response
}
StreamCommand::Append {
stream_id,
content_type,
payload,
close_after,
stream_seq,
producer,
now_ms,
record_match,
} => {
let response = self.append_borrowed(AppendStreamInput {
stream_id,
content_type: content_type.as_deref(),
payload: &payload,
close_after,
stream_seq,
producer,
now_ms,
record_match,
});
self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
response
}
StreamCommand::AppendExternal {
stream_id,
content_type,
payload,
record_ends,
close_after,
stream_seq,
producer,
now_ms,
record_match,
} => {
let response = self.append_external(AppendExternalInput {
stream_id,
content_type: content_type.as_deref(),
payload,
record_ends,
close_after,
stream_seq,
producer,
now_ms,
record_match,
});
self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
response
}
StreamCommand::AppendBatch {
stream_id,
content_type,
payloads,
producer,
now_ms,
} => {
let response = match self.append_batch_borrowed(
stream_id,
content_type.as_deref(),
&payloads.iter().map(Bytes::as_ref).collect::<Vec<_>>(),
producer,
now_ms,
) {
Ok(batch) => batch
.items
.last()
.map(|item| StreamResponse::Appended {
offset: item.offset,
next_offset: item.next_offset,
closed: item.closed,
deduplicated: item.deduplicated,
producer: None,
})
.unwrap_or_else(|| {
StreamResponse::error(
StreamErrorCode::EmptyAppend,
"append batch must contain at least one payload",
)
}),
Err(response) => response,
};
self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
response
}
StreamCommand::PublishSnapshot {
stream_id,
snapshot_offset,
content_type,
payload,
expected_digest,
now_ms,
} => {
let response = self.publish_snapshot(
stream_id,
snapshot_offset,
content_type,
payload.into(),
expected_digest,
now_ms,
);
self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
response
}
StreamCommand::AdvanceRetention {
stream_id,
retained_offset,
now_ms,
} => {
let response = self.advance_retention(stream_id, retained_offset, now_ms);
self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
response
}
StreamCommand::TouchStreamAccess {
stream_id,
now_ms,
renew_ttl,
} => {
let response = self.touch_stream_access(&stream_id, now_ms, renew_ttl);
self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
response
}
StreamCommand::UpdateStreamAttrs {
stream_id,
attrs,
now_ms,
} => {
let response = self.update_stream_attrs(&stream_id, attrs, now_ms);
self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
response
}
StreamCommand::FlushCold { stream_id, chunk } => self.flush_cold(stream_id, chunk),
StreamCommand::CompactCold {
stream_id,
old_chunks,
replacement,
gc_not_before_ms,
} => self.compact_cold(stream_id, old_chunks, replacement, gc_not_before_ms),
StreamCommand::Close {
stream_id,
stream_seq,
producer,
now_ms,
} => {
let response = self.close(stream_id, stream_seq, producer, now_ms);
self.sweep_expired_streams(now_ms, TTL_EXPIRY_SWEEP_MAX_STREAMS_PER_WRITE);
response
}
StreamCommand::DeleteStream { stream_id } => self.delete_stream(&stream_id),
StreamCommand::PurgeBucket { bucket_id } => self.purge_bucket(&bucket_id),
StreamCommand::AckColdGc { up_to_seq } => self.ack_cold_gc(up_to_seq),
StreamCommand::ImportSnapshot { snapshot } => self.import_snapshot(*snapshot),
StreamCommand::SetBucketQuota {
bucket_id,
max_streams,
max_retained_bytes,
} => self.set_bucket_quota(bucket_id, max_streams, max_retained_bytes),
}
}
}
#[derive(Debug)]
struct CreateStreamInput {
stream_id: BucketStreamId,
content_type: String,
initial_payload: Vec<u8>,
record_ends: Vec<u64>,
close_after: bool,
stream_seq: Option<String>,
producer: Option<ProducerRequest>,
stream_ttl_seconds: Option<u64>,
stream_expires_at_ms: Option<u64>,
attrs: Option<StreamAttrs>,
now_ms: u64,
}
#[derive(Debug)]
struct CreateExternalStreamInput {
stream_id: BucketStreamId,
content_type: String,
initial_payload: ExternalPayloadRef,
record_ends: Vec<u64>,
close_after: bool,
stream_seq: Option<String>,
producer: Option<ProducerRequest>,
stream_ttl_seconds: Option<u64>,
stream_expires_at_ms: Option<u64>,
attrs: Option<StreamAttrs>,
now_ms: u64,
}
impl CreateStreamInput {
fn initial_len(&self) -> u64 {
u64::try_from(self.initial_payload.len()).expect("payload len fits u64")
}
}
fn normalize_stream_attrs(attrs: Option<StreamAttrs>) -> Option<StreamAttrs> {
attrs.filter(|attrs| !attrs.is_empty())
}
fn stream_expiry_at_ms(stream: &StreamMetadata) -> Option<u64> {
if let Some(expires_at_ms) = stream.stream_expires_at_ms {
return Some(expires_at_ms);
}
stream.stream_ttl_seconds.map(|ttl_seconds| {
stream
.last_ttl_touch_at_ms
.saturating_add(ttl_seconds.saturating_mul(1000))
})
}
fn stream_is_expired(stream: &StreamMetadata, now_ms: u64) -> bool {
stream_expiry_at_ms(stream).is_some_and(|expires_at_ms| now_ms >= expires_at_ms)
}
fn stream_ttl_renewal_due(stream: &StreamMetadata, now_ms: u64) -> bool {
let Some(ttl_seconds) = stream.stream_ttl_seconds else {
return false;
};
if stream.stream_expires_at_ms.is_some() {
return false;
}
let ttl_ms = ttl_seconds.saturating_mul(1000);
let renewal_interval_ms = ttl_ms.div_ceil(4).max(1);
now_ms.saturating_sub(stream.last_ttl_touch_at_ms) >= renewal_interval_ms
}
fn renew_stream_ttl(stream: &mut StreamMetadata, now_ms: u64) {
if stream.stream_ttl_seconds.is_some() && stream.stream_expires_at_ms.is_none() {
stream.last_ttl_touch_at_ms = now_ms;
}
}
fn validate_producer_request(producer: Option<&ProducerRequest>) -> Result<(), StreamResponse> {
let Some(producer) = producer else {
return Ok(());
};
if producer.producer_id.trim().is_empty() {
return Err(StreamResponse::error(
StreamErrorCode::InvalidProducer,
"producer id must not be empty",
));
}
const MAX_JS_SAFE_INTEGER: u64 = 9_007_199_254_740_991;
if producer.producer_epoch > MAX_JS_SAFE_INTEGER {
return Err(StreamResponse::error(
StreamErrorCode::InvalidProducer,
format!(
"producer epoch {} exceeds maximum {}",
producer.producer_epoch, MAX_JS_SAFE_INTEGER
),
));
}
if producer.producer_seq > MAX_JS_SAFE_INTEGER {
return Err(StreamResponse::error(
StreamErrorCode::InvalidProducer,
format!(
"producer sequence {} exceeds maximum {}",
producer.producer_seq, MAX_JS_SAFE_INTEGER
),
));
}
Ok(())
}
fn validate_external_payload_ref(payload: &ExternalPayloadRef) -> Result<(), StreamResponse> {
if payload.s3_path.trim().is_empty() {
return Err(StreamResponse::error(
StreamErrorCode::InvalidColdFlush,
"external payload S3 path must not be empty",
));
}
if payload.payload_len == 0 {
return Err(StreamResponse::error(
StreamErrorCode::EmptyAppend,
"external payload length must be greater than zero",
));
}
if payload.object_size < payload.payload_len {
return Err(StreamResponse::error(
StreamErrorCode::InvalidColdFlush,
"external payload object size must cover payload length",
));
}
Ok(())
}
fn build_record_index(
content_type: &str,
payload_len: u64,
record_ends: &[u64],
) -> Result<Option<StreamRecordIndex>, StreamResponse> {
if !is_json_record_content_type(content_type) {
return record_ends.is_empty().then_some(None).ok_or_else(|| {
StreamResponse::error(
StreamErrorCode::InvalidRecordBoundaries,
"record boundaries are only valid for application/json streams",
)
});
}
if payload_len > 0 && record_ends.is_empty() {
return Ok(None);
}
let mut index = StreamRecordIndex::new();
index
.append_relative_ends(0, payload_len, record_ends)
.map_err(|_| {
StreamResponse::error(
StreamErrorCode::InvalidRecordBoundaries,
"record boundaries do not match the canonical JSON payload",
)
})?;
Ok(Some(index))
}
fn prepare_record_append(
current: Option<&StreamRecordIndex>,
json_stream: bool,
base_offset: u64,
payload_len: u64,
record_ends: &[u64],
) -> Result<Option<crate::PreparedRecordAppend>, StreamResponse> {
let Some(current) = current else {
if json_stream {
return Ok(None);
}
return record_ends.is_empty().then_some(None).ok_or_else(|| {
StreamResponse::error(
StreamErrorCode::InvalidRecordBoundaries,
"binary streams cannot carry JSON record boundaries",
)
});
};
current
.prepare_append(base_offset, payload_len, record_ends)
.map(Some)
.map_err(|_| {
StreamResponse::error(
StreamErrorCode::InvalidRecordBoundaries,
"record boundaries do not match the canonical JSON payload",
)
})
}
fn compare_stream_ids(left: &BucketStreamId, right: &BucketStreamId) -> std::cmp::Ordering {
left.bucket_id
.cmp(&right.bucket_id)
.then_with(|| left.stream_id.cmp(&right.stream_id))
}
fn snapshot_digest(content_type: &str, payload: &[u8]) -> String {
let mut hasher = blake3::Hasher::new();
hasher.update(&(content_type.len() as u64).to_le_bytes());
hasher.update(content_type.as_bytes());
hasher.update(payload);
hasher.finalize().to_hex().to_string()
}
#[cfg(test)]
mod tests;