ursula_stream/
snapshot.rs1use 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}