Skip to main content

ursula_stream/
snapshot.rs

1use serde::Deserialize;
2use serde::Serialize;
3use ursula_shard::BucketStreamId;
4
5use crate::integrity::StreamIntegritySnapshot;
6use crate::model::BucketQuotaSnapshot;
7use crate::model::BucketUsageSnapshot;
8use crate::model::ColdChunkRef;
9use crate::model::ColdGcEntry;
10use crate::model::HotPayloadSegment;
11use crate::model::ObjectPayloadRef;
12use crate::model::ProducerSnapshot;
13use crate::model::StreamAttrs;
14use crate::model::StreamMessageRecord;
15use crate::model::StreamMetadata;
16use crate::model::StreamVisibleSnapshot;
17use crate::record_index::StreamRecordIndex;
18
19#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
20pub struct StreamSnapshot {
21    pub buckets: Vec<String>,
22    pub streams: Vec<StreamSnapshotEntry>,
23    #[serde(default)]
24    pub pending_cold_gc: Vec<ColdGcEntry>,
25    #[serde(default)]
26    pub next_cold_gc_seq: u64,
27    /// Monotonic per-bucket usage counters. Absent in legacy snapshots, in
28    /// which case the monotonic counters restart from the restored gauges.
29    #[serde(default)]
30    pub bucket_usage: Vec<BucketUsageSnapshot>,
31    /// Per-bucket data-plane quota records. Absent in legacy snapshots.
32    #[serde(default)]
33    pub bucket_quotas: Vec<BucketQuotaSnapshot>,
34}
35
36#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
37pub struct StreamSnapshotEntry {
38    pub metadata: StreamMetadata,
39    #[serde(default)]
40    pub attrs: Option<StreamAttrs>,
41    pub hot_start_offset: u64,
42    pub payload: Vec<u8>,
43    pub hot_segments: Vec<HotPayloadSegment>,
44    #[serde(default)]
45    pub cold_frontier_offset: u64,
46    #[serde(default)]
47    pub cold_index_generation: u64,
48    pub cold_chunks: Vec<ColdChunkRef>,
49    pub external_segments: Vec<ObjectPayloadRef>,
50    pub message_records: Vec<StreamMessageRecord>,
51    #[serde(default)]
52    pub record_index: Option<StreamRecordIndex>,
53    pub integrity: StreamIntegritySnapshot,
54    /// Independent destructive-retention floor. `None` denotes a legacy
55    /// snapshot where the visible snapshot offset also implied retention.
56    #[serde(default)]
57    pub retained_offset: Option<u64>,
58    pub visible_snapshot: Option<StreamVisibleSnapshot>,
59    pub producer_states: Vec<ProducerSnapshot>,
60}
61
62#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
63pub enum StreamSnapshotError {
64    #[error("snapshot contains duplicate bucket '{0}'")]
65    DuplicateBucket(String),
66    #[error("snapshot contains duplicate stream '{0}'")]
67    DuplicateStream(BucketStreamId),
68    #[error("snapshot stream '{stream_id}' contains duplicate producer '{producer_id}'")]
69    DuplicateProducer {
70        stream_id: BucketStreamId,
71        producer_id: String,
72    },
73    #[error("snapshot stream '{0}' references a missing bucket")]
74    MissingBucket(BucketStreamId),
75    #[error(
76        "snapshot stream '{stream_id}' tail offset {tail_offset} does not match payload length {payload_len}"
77    )]
78    PayloadLengthMismatch {
79        stream_id: BucketStreamId,
80        tail_offset: u64,
81        payload_len: usize,
82    },
83    #[error("snapshot stream '{stream_id}' has inconsistent message boundaries")]
84    MessageBoundaryMismatch { stream_id: BucketStreamId },
85    #[error("snapshot stream '{stream_id}' has inconsistent record boundaries")]
86    RecordBoundaryMismatch { stream_id: BucketStreamId },
87    #[error("snapshot stream '{stream_id}' has inconsistent integrity setsums")]
88    IntegrityMismatch { stream_id: BucketStreamId },
89    #[error(
90        "snapshot stream '{stream_id}' visible snapshot offset {snapshot_offset} is beyond tail offset {tail_offset}"
91    )]
92    SnapshotOffsetOutOfRange {
93        stream_id: BucketStreamId,
94        snapshot_offset: u64,
95        tail_offset: u64,
96    },
97}