Skip to main content

ursula_stream/
integrity.rs

1use serde::Deserialize;
2use serde::Serialize;
3use setsum::Setsum;
4use ursula_shard::BucketStreamId;
5
6#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
7pub struct StreamIntegritySnapshot {
8    pub live_setsum: String,
9    pub evicted_setsum: String,
10    pub total_setsum: String,
11    pub live_start_offset: u64,
12    pub tail_offset: u64,
13    pub live_records: u64,
14    pub evicted_records: u64,
15    pub total_records: u64,
16}
17
18#[derive(Debug, Clone, Default, PartialEq, Eq)]
19pub(crate) struct StreamIntegrity {
20    live: Setsum,
21    total: Setsum,
22    total_records: u64,
23}
24
25impl StreamIntegrity {
26    pub(crate) fn append_payload(
27        &mut self,
28        stream_id: &BucketStreamId,
29        start_offset: u64,
30        end_offset: u64,
31        payload: &[u8],
32    ) {
33        let record = record_setsum(stream_id, start_offset, end_offset, b"inline", &[payload]);
34        self.append_record(start_offset, end_offset, record);
35    }
36
37    pub(crate) fn append_external(
38        &mut self,
39        stream_id: &BucketStreamId,
40        start_offset: u64,
41        end_offset: u64,
42        s3_path: &str,
43        object_size: u64,
44    ) {
45        let object_size = object_size.to_le_bytes();
46        let record = record_setsum(stream_id, start_offset, end_offset, b"external", &[
47            s3_path.as_bytes(),
48            &object_size,
49        ]);
50        self.append_record(start_offset, end_offset, record);
51    }
52
53    pub(crate) fn evict_before(&mut self, retained_offset: u64) {
54        let _ = retained_offset;
55    }
56
57    pub(crate) fn snapshot(
58        &self,
59        live_start_offset: u64,
60        tail_offset: u64,
61    ) -> StreamIntegritySnapshot {
62        StreamIntegritySnapshot {
63            live_setsum: self.live.hexdigest(),
64            evicted_setsum: Setsum::default().hexdigest(),
65            total_setsum: self.total.hexdigest(),
66            live_start_offset,
67            tail_offset,
68            live_records: self.total_records,
69            evicted_records: 0,
70            total_records: self.total_records,
71        }
72    }
73
74    pub(crate) fn restore(snapshot: StreamIntegritySnapshot) -> Option<Self> {
75        let live = Setsum::from_hexdigest(&snapshot.live_setsum)?;
76        let evicted = Setsum::from_hexdigest(&snapshot.evicted_setsum)?;
77        let total = Setsum::from_hexdigest(&snapshot.total_setsum)?;
78        if live + evicted != total {
79            return None;
80        }
81        Some(Self {
82            live: total,
83            total,
84            total_records: snapshot.total_records,
85        })
86    }
87
88    fn append_record(&mut self, start_offset: u64, end_offset: u64, record: Setsum) {
89        if start_offset == end_offset {
90            return;
91        }
92        self.live += record;
93        self.total += record;
94        self.total_records = self.total_records.saturating_add(1);
95    }
96}
97
98fn record_setsum(
99    stream_id: &BucketStreamId,
100    start_offset: u64,
101    end_offset: u64,
102    kind: &[u8],
103    pieces: &[&[u8]],
104) -> Setsum {
105    let mut setsum = Setsum::default();
106    let start = start_offset.to_le_bytes();
107    let end = end_offset.to_le_bytes();
108    let mut item = vec![
109        b"ursula-stream-record-v1".as_slice(),
110        stream_id.bucket_id.as_bytes(),
111        b"\0",
112        stream_id.stream_id.as_bytes(),
113        b"\0",
114        &start,
115        &end,
116        kind,
117    ];
118    item.extend_from_slice(pieces);
119    setsum.insert_vectored(&item);
120    setsum
121}
122
123#[cfg(test)]
124mod tests {
125    use super::*;
126
127    fn stream_id() -> BucketStreamId {
128        BucketStreamId::new("benchcmp", "integrity")
129    }
130
131    #[test]
132    fn snapshot_omits_per_append_records() {
133        let stream_id = stream_id();
134        let mut integrity = StreamIntegrity::default();
135        integrity.append_payload(&stream_id, 0, 3, b"abc");
136        integrity.append_payload(&stream_id, 3, 5, b"de");
137
138        let snapshot = integrity.snapshot(0, 5);
139
140        assert_eq!(snapshot.live_records, 2);
141        assert_eq!(snapshot.evicted_records, 0);
142        assert_eq!(snapshot.total_records, 2);
143        assert_eq!(snapshot.live_setsum, snapshot.total_setsum);
144    }
145
146    #[test]
147    fn snapshot_wire_format_has_no_records_field() {
148        let stream_id = stream_id();
149        let mut integrity = StreamIntegrity::default();
150        integrity.append_payload(&stream_id, 0, 3, b"abc");
151
152        let encoded =
153            serde_json::to_value(integrity.snapshot(0, 3)).expect("encode compacted integrity");
154
155        assert!(encoded.get("records").is_none());
156    }
157}