1use serde::Deserialize;
4use serde::Serialize;
5
6#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
7pub struct StreamRecordIndex {
8 first_record: u64,
9 record_offsets: Vec<u64>,
10}
11
12#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
13pub struct StreamRecordRange {
14 pub first_record: u64,
15 pub next_record: u64,
16}
17
18#[derive(Debug)]
19pub(crate) struct PreparedRecordAppend {
20 range: StreamRecordRange,
21 record_offsets: Vec<u64>,
22}
23
24impl PreparedRecordAppend {
25 pub(crate) fn range(&self) -> StreamRecordRange {
26 self.range
27 }
28}
29
30#[derive(Debug, Clone, Copy, PartialEq, Eq)]
31pub enum RecordIndexError {
32 InvalidBoundaries,
33 ArithmeticOverflow,
34 RecordGone { first_record: u64, next_record: u64 },
35 RecordBeyondTail { next_record: u64 },
36 OffsetNotRecordBoundary,
37}
38
39pub fn is_json_record_content_type(content_type: &str) -> bool {
40 content_type
41 .split(';')
42 .next()
43 .is_some_and(|value| value.trim().eq_ignore_ascii_case("application/json"))
44}
45
46pub fn canonical_json_record_ends(
47 content_type: &str,
48 payload: &[u8],
49) -> Result<Vec<u64>, RecordIndexError> {
50 if !is_json_record_content_type(content_type) {
51 return Ok(Vec::new());
52 }
53 if payload.is_empty() {
54 return Ok(Vec::new());
55 }
56 if payload.last() != Some(&b'\n') {
57 return Err(RecordIndexError::InvalidBoundaries);
58 }
59 payload
60 .iter()
61 .enumerate()
62 .filter_map(|(index, byte)| (*byte == b'\n').then_some(index + 1))
63 .map(|end| u64::try_from(end).map_err(|_| RecordIndexError::ArithmeticOverflow))
64 .collect()
65}
66
67impl StreamRecordIndex {
68 pub fn new() -> Self {
69 Self::default()
70 }
71
72 pub fn restore(
73 first_record: u64,
74 record_offsets: Vec<u64>,
75 retained_offset: u64,
76 tail_offset: u64,
77 ) -> Result<Self, RecordIndexError> {
78 let index = Self {
79 first_record,
80 record_offsets,
81 };
82 index.validate(retained_offset, tail_offset)?;
83 Ok(index)
84 }
85
86 pub fn range(&self) -> Result<StreamRecordRange, RecordIndexError> {
87 let retained = u64::try_from(self.record_offsets.len())
88 .map_err(|_| RecordIndexError::ArithmeticOverflow)?;
89 let next_record = self
90 .first_record
91 .checked_add(retained)
92 .ok_or(RecordIndexError::ArithmeticOverflow)?;
93 Ok(StreamRecordRange {
94 first_record: self.first_record,
95 next_record,
96 })
97 }
98
99 pub fn record_offsets(&self) -> &[u64] {
100 &self.record_offsets
101 }
102
103 pub fn append_relative_ends(
104 &mut self,
105 base_offset: u64,
106 payload_len: u64,
107 relative_ends: &[u64],
108 ) -> Result<StreamRecordRange, RecordIndexError> {
109 let prepared = self.prepare_append(base_offset, payload_len, relative_ends)?;
110 Ok(self.commit_append(prepared))
111 }
112
113 pub(crate) fn prepare_append(
114 &self,
115 base_offset: u64,
116 payload_len: u64,
117 relative_ends: &[u64],
118 ) -> Result<PreparedRecordAppend, RecordIndexError> {
119 validate_relative_ends(payload_len, relative_ends)?;
120 let record_start = self.range()?.next_record;
121 let appended =
122 u64::try_from(relative_ends.len()).map_err(|_| RecordIndexError::ArithmeticOverflow)?;
123 let record_next = record_start
124 .checked_add(appended)
125 .ok_or(RecordIndexError::ArithmeticOverflow)?;
126 if self
127 .record_offsets
128 .last()
129 .is_some_and(|last| *last >= base_offset)
130 {
131 return Err(RecordIndexError::InvalidBoundaries);
132 }
133 let mut starts = Vec::with_capacity(relative_ends.len());
134 let mut previous_end = 0;
135 for end in relative_ends {
136 let start_offset = base_offset
137 .checked_add(previous_end)
138 .ok_or(RecordIndexError::ArithmeticOverflow)?;
139 starts.push(start_offset);
140 previous_end = *end;
141 }
142 Ok(PreparedRecordAppend {
143 range: StreamRecordRange {
144 first_record: record_start,
145 next_record: record_next,
146 },
147 record_offsets: starts,
148 })
149 }
150
151 pub(crate) fn commit_append(&mut self, prepared: PreparedRecordAppend) -> StreamRecordRange {
152 self.record_offsets.extend(prepared.record_offsets);
153 prepared.range
154 }
155
156 pub(crate) fn append_checkpoint(&self) -> usize {
157 self.record_offsets.len()
158 }
159
160 pub(crate) fn rollback_appends(&mut self, checkpoint: usize) {
161 self.record_offsets.truncate(checkpoint);
162 }
163
164 pub fn offset_for(&self, record: u64, tail_offset: u64) -> Result<u64, RecordIndexError> {
165 let range = self.range()?;
166 if record < range.first_record {
167 return Err(RecordIndexError::RecordGone {
168 first_record: range.first_record,
169 next_record: range.next_record,
170 });
171 }
172 if record > range.next_record {
173 return Err(RecordIndexError::RecordBeyondTail {
174 next_record: range.next_record,
175 });
176 }
177 if record == range.next_record {
178 return Ok(tail_offset);
179 }
180 let relative = record
181 .checked_sub(range.first_record)
182 .ok_or(RecordIndexError::ArithmeticOverflow)?;
183 let index = usize::try_from(relative).map_err(|_| RecordIndexError::ArithmeticOverflow)?;
184 self.record_offsets
185 .get(index)
186 .copied()
187 .ok_or(RecordIndexError::InvalidBoundaries)
188 }
189
190 pub fn record_for_offset(
191 &self,
192 offset: u64,
193 tail_offset: u64,
194 ) -> Result<u64, RecordIndexError> {
195 let range = self.range()?;
196 if offset == tail_offset {
197 return Ok(range.next_record);
198 }
199 let relative = self
200 .record_offsets
201 .binary_search(&offset)
202 .map_err(|_| RecordIndexError::OffsetNotRecordBoundary)?;
203 range
204 .first_record
205 .checked_add(u64::try_from(relative).map_err(|_| RecordIndexError::ArithmeticOverflow)?)
206 .ok_or(RecordIndexError::ArithmeticOverflow)
207 }
208
209 pub fn retain_from_offset(
210 &mut self,
211 retained_offset: u64,
212 tail_offset: u64,
213 ) -> Result<u64, RecordIndexError> {
214 let removed = if retained_offset == tail_offset {
215 self.record_offsets.len()
216 } else {
217 self.record_offsets
218 .binary_search(&retained_offset)
219 .map_err(|_| RecordIndexError::OffsetNotRecordBoundary)?
220 };
221 let removed_u64 =
222 u64::try_from(removed).map_err(|_| RecordIndexError::ArithmeticOverflow)?;
223 self.first_record = self
224 .first_record
225 .checked_add(removed_u64)
226 .ok_or(RecordIndexError::ArithmeticOverflow)?;
227 self.record_offsets.drain(..removed);
228 Ok(self.first_record)
229 }
230
231 pub fn validate(&self, retained_offset: u64, tail_offset: u64) -> Result<(), RecordIndexError> {
232 let _ = self.range()?;
233 if retained_offset > tail_offset {
234 return Err(RecordIndexError::InvalidBoundaries);
235 }
236 if self.record_offsets.is_empty() {
237 return (retained_offset == tail_offset)
238 .then_some(())
239 .ok_or(RecordIndexError::InvalidBoundaries);
240 }
241 if self.record_offsets.first().copied() != Some(retained_offset)
242 || self
243 .record_offsets
244 .iter()
245 .any(|offset| *offset >= tail_offset)
246 || self.record_offsets.windows(2).any(|pair| {
247 let [left, right] = pair else {
248 return true;
249 };
250 left >= right
251 })
252 {
253 return Err(RecordIndexError::InvalidBoundaries);
254 }
255 Ok(())
256 }
257}
258
259fn validate_relative_ends(payload_len: u64, relative_ends: &[u64]) -> Result<(), RecordIndexError> {
260 if payload_len == 0 {
261 return relative_ends
262 .is_empty()
263 .then_some(())
264 .ok_or(RecordIndexError::InvalidBoundaries);
265 }
266 if relative_ends.last().copied() != Some(payload_len)
267 || relative_ends.first().copied() == Some(0)
268 || relative_ends.windows(2).any(|pair| {
269 let [left, right] = pair else {
270 return true;
271 };
272 left >= right
273 })
274 {
275 return Err(RecordIndexError::InvalidBoundaries);
276 }
277 Ok(())
278}
279
280#[cfg(test)]
281mod tests {
282 use super::RecordIndexError;
283 use super::StreamRecordIndex;
284 use super::StreamRecordRange;
285 use super::canonical_json_record_ends;
286
287 #[test]
288 fn append_maps_contiguous_ordinals_to_exact_offsets() {
289 let mut index = StreamRecordIndex::new();
290 assert_eq!(
291 index.append_relative_ends(0, 18, &[9, 18]),
292 Ok(StreamRecordRange {
293 first_record: 0,
294 next_record: 2,
295 })
296 );
297 assert_eq!(index.record_offsets(), &[0, 9]);
298 assert_eq!(index.offset_for(0, 18), Ok(0));
299 assert_eq!(index.offset_for(1, 18), Ok(9));
300 assert_eq!(index.offset_for(2, 18), Ok(18));
301
302 assert_eq!(
303 index.append_relative_ends(18, 9, &[9]),
304 Ok(StreamRecordRange {
305 first_record: 2,
306 next_record: 3,
307 })
308 );
309 assert_eq!(index.record_offsets(), &[0, 9, 18]);
310 assert_eq!(index.offset_for(3, 27), Ok(27));
311 }
312
313 #[test]
314 fn retention_drops_offsets_without_renumbering() {
315 let mut index = StreamRecordIndex::new();
316 index
317 .append_relative_ends(0, 27, &[9, 18, 27])
318 .expect("append boundaries");
319 assert_eq!(index.retain_from_offset(18, 27), Ok(2));
320 assert_eq!(index.record_offsets(), &[18]);
321 assert_eq!(
322 index.offset_for(1, 27),
323 Err(RecordIndexError::RecordGone {
324 first_record: 2,
325 next_record: 3,
326 })
327 );
328 assert_eq!(index.offset_for(2, 27), Ok(18));
329 }
330
331 #[test]
332 fn restore_rejects_misaligned_or_non_monotonic_offsets() {
333 assert_eq!(
334 StreamRecordIndex::restore(0, vec![1, 9], 0, 18),
335 Err(RecordIndexError::InvalidBoundaries)
336 );
337 assert_eq!(
338 StreamRecordIndex::restore(0, vec![0, 0], 0, 18),
339 Err(RecordIndexError::InvalidBoundaries)
340 );
341 assert!(StreamRecordIndex::restore(2, vec![18], 18, 27).is_ok());
342 }
343
344 #[test]
345 fn relative_ends_must_cover_the_payload_exactly() {
346 let mut index = StreamRecordIndex::new();
347 assert_eq!(
348 index.append_relative_ends(0, 18, &[9]),
349 Err(RecordIndexError::InvalidBoundaries)
350 );
351 assert_eq!(
352 index.append_relative_ends(0, 18, &[9, 9, 18]),
353 Err(RecordIndexError::InvalidBoundaries)
354 );
355 assert_eq!(
356 index.append_relative_ends(0, 0, &[0]),
357 Err(RecordIndexError::InvalidBoundaries)
358 );
359 }
360
361 #[test]
362 fn canonical_json_payload_exposes_each_ndjson_boundary() {
363 assert_eq!(
364 canonical_json_record_ends(
365 "application/json; charset=utf-8",
366 b"{\"a\":1}\n{\"b\":2}\n"
367 ),
368 Ok(vec![8, 16])
369 );
370 assert_eq!(
371 canonical_json_record_ends("application/octet-stream", b"x"),
372 Ok(vec![])
373 );
374 assert_eq!(
375 canonical_json_record_ends("application/json", b"{}"),
376 Err(RecordIndexError::InvalidBoundaries)
377 );
378 }
379}