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 = if let Some(affinity_key) = &stream_id.affinity_key {
109 vec![
110 b"ursula-stream-record-v2".as_slice(),
111 stream_id.bucket_id.as_bytes(),
112 b"\0",
113 affinity_key.as_bytes(),
114 b"\0",
115 stream_id.stream_id.as_bytes(),
116 b"\0",
117 &start,
118 &end,
119 kind,
120 ]
121 } else {
122 vec![
123 b"ursula-stream-record-v1".as_slice(),
124 stream_id.bucket_id.as_bytes(),
125 b"\0",
126 stream_id.stream_id.as_bytes(),
127 b"\0",
128 &start,
129 &end,
130 kind,
131 ]
132 };
133 item.extend_from_slice(pieces);
134 setsum.insert_vectored(&item);
135 setsum
136}
137
138#[cfg(test)]
139mod tests {
140 use super::*;
141
142 fn stream_id() -> BucketStreamId {
143 BucketStreamId::new("benchcmp", "integrity")
144 }
145
146 #[test]
147 fn snapshot_omits_per_append_records() {
148 let stream_id = stream_id();
149 let mut integrity = StreamIntegrity::default();
150 integrity.append_payload(&stream_id, 0, 3, b"abc");
151 integrity.append_payload(&stream_id, 3, 5, b"de");
152
153 let snapshot = integrity.snapshot(0, 5);
154
155 assert_eq!(snapshot.live_records, 2);
156 assert_eq!(snapshot.evicted_records, 0);
157 assert_eq!(snapshot.total_records, 2);
158 assert_eq!(snapshot.live_setsum, snapshot.total_setsum);
159 }
160
161 #[test]
162 fn snapshot_wire_format_has_no_records_field() {
163 let stream_id = stream_id();
164 let mut integrity = StreamIntegrity::default();
165 integrity.append_payload(&stream_id, 0, 3, b"abc");
166
167 let encoded =
168 serde_json::to_value(integrity.snapshot(0, 3)).expect("encode compacted integrity");
169
170 assert!(encoded.get("records").is_none());
171 }
172}