Skip to main content

ursula_stream/
record_index.rs

1//! Exact retained record-ordinal to canonical-offset boundaries for JSON streams.
2
3use 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}