ursula_stream/
snapshot.rs1use 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 #[serde(default)]
30 pub bucket_usage: Vec<BucketUsageSnapshot>,
31 #[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 #[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}