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::ColdChunkRef;
7use crate::model::ColdGcEntry;
8use crate::model::HotPayloadSegment;
9use crate::model::ObjectPayloadRef;
10use crate::model::ProducerSnapshot;
11use crate::model::StreamAttrs;
12use crate::model::StreamMessageRecord;
13use crate::model::StreamMetadata;
14use crate::model::StreamVisibleSnapshot;
15use crate::record_index::StreamRecordIndex;
16
17#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
18pub struct StreamSnapshot {
19    pub buckets: Vec<String>,
20    pub streams: Vec<StreamSnapshotEntry>,
21    #[serde(default)]
22    pub pending_cold_gc: Vec<ColdGcEntry>,
23    #[serde(default)]
24    pub next_cold_gc_seq: u64,
25}
26
27#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
28pub struct StreamSnapshotEntry {
29    pub metadata: StreamMetadata,
30    #[serde(default)]
31    pub attrs: Option<StreamAttrs>,
32    pub hot_start_offset: u64,
33    pub payload: Vec<u8>,
34    pub hot_segments: Vec<HotPayloadSegment>,
35    #[serde(default)]
36    pub cold_frontier_offset: u64,
37    #[serde(default)]
38    pub cold_index_generation: u64,
39    pub cold_chunks: Vec<ColdChunkRef>,
40    pub external_segments: Vec<ObjectPayloadRef>,
41    pub message_records: Vec<StreamMessageRecord>,
42    #[serde(default)]
43    pub record_index: Option<StreamRecordIndex>,
44    pub integrity: StreamIntegritySnapshot,
45    pub visible_snapshot: Option<StreamVisibleSnapshot>,
46    pub producer_states: Vec<ProducerSnapshot>,
47}
48
49#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
50pub enum StreamSnapshotError {
51    #[error("snapshot contains duplicate bucket '{0}'")]
52    DuplicateBucket(String),
53    #[error("snapshot contains duplicate stream '{0}'")]
54    DuplicateStream(BucketStreamId),
55    #[error("snapshot stream '{stream_id}' contains duplicate producer '{producer_id}'")]
56    DuplicateProducer {
57        stream_id: BucketStreamId,
58        producer_id: String,
59    },
60    #[error("snapshot stream '{0}' references a missing bucket")]
61    MissingBucket(BucketStreamId),
62    #[error(
63        "snapshot stream '{stream_id}' tail offset {tail_offset} does not match payload length {payload_len}"
64    )]
65    PayloadLengthMismatch {
66        stream_id: BucketStreamId,
67        tail_offset: u64,
68        payload_len: usize,
69    },
70    #[error("snapshot stream '{stream_id}' has inconsistent message boundaries")]
71    MessageBoundaryMismatch { stream_id: BucketStreamId },
72    #[error("snapshot stream '{stream_id}' has inconsistent record boundaries")]
73    RecordBoundaryMismatch { stream_id: BucketStreamId },
74    #[error("snapshot stream '{stream_id}' has inconsistent integrity setsums")]
75    IntegrityMismatch { stream_id: BucketStreamId },
76    #[error(
77        "snapshot stream '{stream_id}' visible snapshot offset {snapshot_offset} is beyond tail offset {tail_offset}"
78    )]
79    SnapshotOffsetOutOfRange {
80        stream_id: BucketStreamId,
81        snapshot_offset: u64,
82        tail_offset: u64,
83    },
84}