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}