Skip to main content

journal/
indexed_snapshot.rs

1//! Bounded native-index traversal. Capture requires caller-owned writer exclusion.
2use super::*;
3use journal_core::file::{HashTable, JournalHeader, JournalState};
4
5/// Cooperative cancellation, checked between objects, chunks, and payloads.
6/// A syscall or one decompression is not interruptible.
7#[derive(Default)]
8pub struct SnapshotControl<'a> {
9    pub cancelled: Option<&'a dyn Fn() -> bool>,
10}
11impl SnapshotControl<'_> {
12    pub(crate) fn check(&self) -> Result<()> {
13        if self.cancelled.is_some_and(|f| f()) {
14            Err(SdkError::Cancelled)
15        } else {
16            Ok(())
17        }
18    }
19}
20
21#[derive(Clone, Default)]
22pub struct IndexedSnapshotOptions {
23    pub reader: ReaderOptions,
24    pub capture_fields: Vec<Vec<u8>>,
25    /// Exact `(name, value)` pairs. Names and values are arbitrary bytes.
26    pub capture_values: Vec<(Vec<u8>, Vec<u8>)>,
27}
28
29#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
30pub struct CapturedValue {
31    pub present: bool,
32    pub entry_count: u64,
33}
34
35#[derive(Clone, Copy, Debug)]
36pub struct SnapshotMetadata {
37    pub seqnum: u64,
38    pub realtime: u64,
39    pub monotonic: u64,
40    pub boot_id: [u8; 16],
41}
42
43/// One callback-scoped entry. Payload bytes cannot outlive their visitor call.
44pub struct SnapshotEntry<'a> {
45    snapshot: &'a IndexedSnapshot,
46    control: &'a SnapshotControl<'a>,
47    metadata: SnapshotMetadata,
48    offsets: Vec<NonZeroU64>,
49}
50impl SnapshotEntry<'_> {
51    pub fn metadata(&self) -> SnapshotMetadata {
52        self.metadata
53    }
54    pub fn visit_payloads(&mut self, mut visit: impl FnMut(&[u8]) -> Result<()>) -> Result<()> {
55        let mut buffer = Vec::new();
56        for &offset in &self.offsets {
57            self.control.check()?;
58            self.snapshot.check_object(offset, 1)?;
59            let data = self.snapshot.file.data_ref(offset)?;
60            if data.is_compressed() {
61                buffer.clear();
62                let len = data.decompress(&mut buffer)?;
63                visit(&buffer[..len])?;
64            } else {
65                visit(data.raw_payload())?;
66            }
67        }
68        Ok(())
69    }
70}
71
72/// A single-consumer bounded view of a valid native journal index graph.
73///
74/// Open while appends, rotation and truncation are excluded by the caller.
75/// After capture, supported append-only writers may resume. In-place mutation
76/// and truncation are unsupported. Uncertain files require `verify_index` first.
77/// This API never falls back to a row scan for an indexed operation. `.journal.zst`
78/// uses whole-file staging and therefore has full-file decompression opening cost.
79pub struct IndexedSnapshot {
80    file: JournalFile<Mmap>,
81    _staging: Option<SnapshotStaging>,
82    header: JournalHeader,
83    object_end: u64,
84    entry_tail: u64,
85    fields: HashMap<Vec<u8>, Option<NonZeroU64>>,
86    values: HashMap<(Vec<u8>, Vec<u8>), CapturedValue>,
87}
88// Fields drop in declaration order: release maps/handles before deleting staging.
89struct SnapshotStaging(PathBuf);
90impl Drop for SnapshotStaging {
91    fn drop(&mut self) {
92        let _ = std::fs::remove_file(&self.0);
93    }
94}
95fn corrupt(message: &str) -> SdkError {
96    SdkError::VerificationError(message.into())
97}
98fn raw_name(name: &[u8]) -> Result<()> {
99    if name.is_empty() || name.contains(&b'=') {
100        Err(JournalError::InvalidField.into())
101    } else {
102        Ok(())
103    }
104}
105fn number(offset: Option<NonZeroU64>) -> u64 {
106    offset.map_or(0, NonZeroU64::get)
107}
108
109impl IndexedSnapshot {
110    pub fn open(
111        path: impl AsRef<Path>,
112        mut options: IndexedSnapshotOptions,
113        control: &SnapshotControl<'_>,
114    ) -> Result<Self> {
115        control.check()?;
116        options.reader.bounds = ReaderBounds::Snapshot;
117        let path = path.as_ref();
118        let temp_path = if is_zst_file(path) {
119            Some(decompress_zst_to_temp(path, "rust-indexed-snapshot")?)
120        } else {
121            None
122        };
123        let file = match open_journal_file(temp_path.as_deref().unwrap_or(path), options.reader) {
124            Ok(file) => file,
125            Err(err) => {
126                if let Some(path) = temp_path {
127                    let _ = std::fs::remove_file(path);
128                }
129                return Err(err);
130            }
131        };
132        let header = *file.journal_header_ref();
133        let mut snapshot = Self {
134            file,
135            _staging: temp_path.map(SnapshotStaging),
136            header,
137            object_end: 0,
138            entry_tail: 0,
139            fields: HashMap::new(),
140            values: HashMap::new(),
141        };
142        control.check()?;
143        snapshot.capture_object_bounds()?;
144        snapshot.capture_entry_bounds(control)?;
145        for field in options.capture_fields {
146            raw_name(&field)?;
147            let head = snapshot.find_field(&field, control)?;
148            snapshot.fields.insert(field, head);
149        }
150        for (name, value) in options.capture_values {
151            raw_name(&name)?;
152            let captured = if let Some(offset) = snapshot.find_data(&name, &value, control, true)? {
153                snapshot.capture_postings(offset, control)?
154            } else {
155                CapturedValue::default()
156            };
157            snapshot.values.insert((name, value), captured);
158        }
159        control.check()?;
160        Ok(snapshot)
161    }
162    fn capture_object_bounds(&mut self) -> Result<()> {
163        let header = &self.header;
164        if header.tail_object_offset.is_none() != (header.n_objects == 0) {
165            return Err(corrupt("object count and tail disagree"));
166        }
167        self.object_end = if let Some(tail) = header.tail_object_offset {
168            if tail.get() < header.header_size || tail.get() % 8 != 0 {
169                return Err(corrupt("invalid object tail"));
170            }
171            let object = self.file.object_header_ref(tail)?;
172            tail.get()
173                .checked_add(object.validated_size()?)
174                .ok_or_else(|| corrupt("object end overflow"))?
175        } else {
176            header.header_size
177        };
178        if self.object_end > header.validated_arena_end(self.file.reader_file_size()?)? {
179            return Err(corrupt("object tail exceeds declared arena"));
180        }
181        Ok(())
182    }
183    fn capture_entry_bounds(&mut self, control: &SnapshotControl<'_>) -> Result<()> {
184        let header = &self.header;
185        let (tail, last_array, last_count) =
186            self.array_tail(header.entry_array_offset, header.n_entries, control)?;
187        self.entry_tail = number(tail);
188        if header.header_size >= 272 && header.tail_entry_offset != self.entry_tail {
189            return Err(corrupt("tail entry hint disagrees with committed count"));
190        }
191        if header.header_size >= 264
192            && (header.tail_entry_array_offset as u64 != number(last_array)
193                || header.tail_entry_array_n_entries as u64 != last_count)
194        {
195            return Err(corrupt("tail array hints disagree with committed count"));
196        }
197        if let Some(tail) = tail {
198            self.check_object(tail, 3)?;
199            let entry = self.file.entry_ref(tail)?;
200            if entry.header.seqnum != header.tail_entry_seqnum
201                || entry.header.realtime != header.tail_entry_realtime
202                || (header.compatible_flags & 2 != 0
203                    && (entry.header.monotonic != header.tail_entry_monotonic
204                        || entry.header.boot_id != header.tail_entry_boot_id))
205            {
206                return Err(corrupt("tail metadata disagrees with committed entry"));
207            }
208        } else {
209            header.validate_empty_entry_metadata()?;
210        }
211        Ok(())
212    }
213    fn capture_postings(
214        &self,
215        offset: NonZeroU64,
216        control: &SnapshotControl<'_>,
217    ) -> Result<CapturedValue> {
218        let (first, array, count) = self.postings(offset)?;
219        if count == 0 || count > self.header.n_entries {
220            return Err(corrupt(
221                "captured posting count exceeds committed population",
222            ));
223        }
224        self.check_entry(first.ok_or_else(|| corrupt("missing inline posting"))?)?;
225        let (last, last_array, used) = self.array_tail(array, count - 1, control)?;
226        let data = self.file.data_ref(offset)?;
227        if let Some((hint_offset, hint_count)) = data.tail_entry_array_hint() {
228            if hint_offset as u64 != number(last_array) || hint_count as u64 != used {
229                return Err(corrupt("DATA tail array hint disagrees with count"));
230            }
231        }
232        drop(data);
233        if let Some(last) = last {
234            self.check_entry(last)?;
235        }
236        Ok(CapturedValue {
237            present: true,
238            entry_count: count,
239        })
240    }
241    pub fn entry_count(&self) -> u64 {
242        self.header.n_entries
243    }
244    /// Whether the captured header was archived. This remains fixed if the
245    /// file is archived later; it does not certify index integrity.
246    pub fn is_archived(&self) -> bool {
247        self.header.state == JournalState::Archived as u8
248    }
249    pub fn captured_value(&self, name: &[u8], value: &[u8]) -> Result<CapturedValue> {
250        self.values
251            .get(&(name.to_vec(), value.to_vec()))
252            .copied()
253            .ok_or(SdkError::Unsupported("undeclared captured value"))
254    }
255    fn check_object(&self, offset: NonZeroU64, kind: u8) -> Result<()> {
256        let n = offset.get();
257        if n < self.header.header_size || n > number(self.header.tail_object_offset) || n % 8 != 0 {
258            return Err(corrupt("object outside captured bounds"));
259        }
260        let object = self.file.object_header_ref(offset)?;
261        if object.type_ != kind
262            || n.checked_add(object.validated_size()?)
263                .is_none_or(|end| end > self.object_end)
264        {
265            return Err(corrupt("invalid bounded object"));
266        }
267        Ok(())
268    }
269    fn check_entry(&self, offset: NonZeroU64) -> Result<()> {
270        if offset.get() > self.entry_tail {
271            return Err(corrupt("entry exceeds committed boundary"));
272        }
273        self.check_object(offset, 3)
274    }
275    // O(array nodes), not O(entries): only each node's terminal used slot is read.
276    fn array_tail(
277        &self,
278        mut current: Option<NonZeroU64>,
279        mut count: u64,
280        control: &SnapshotControl<'_>,
281    ) -> Result<(Option<NonZeroU64>, Option<NonZeroU64>, u64)> {
282        if count == 0 {
283            if current.is_some() {
284                return Err(corrupt("array present for empty population"));
285            }
286            return Ok((None, None, 0));
287        }
288        let mut previous = 0;
289        loop {
290            control.check()?;
291            let offset = current.ok_or_else(|| corrupt("array ended early"))?;
292            if offset.get() <= previous {
293                return Err(corrupt("array chain does not progress"));
294            }
295            self.check_object(offset, 6)?;
296            let array = self.file.offset_array_ref(offset)?;
297            let used = count.min(array.capacity() as u64);
298            if used == 0 {
299                return Err(corrupt("empty array"));
300            }
301            let last = array
302                .items
303                .get(used as usize - 1)
304                .ok_or_else(|| corrupt("zero used array slot"))?;
305            let next = array.header.next_offset_array;
306            drop(array);
307            self.check_object(last, 3)?;
308            count -= used;
309            if count == 0 {
310                return Ok((Some(last), Some(offset), used));
311            }
312            previous = offset.get();
313            current = next;
314        }
315    }
316    fn find_field(&self, name: &[u8], control: &SnapshotControl<'_>) -> Result<Option<NonZeroU64>> {
317        let hash = self.file.hash(name);
318        let table = self
319            .file
320            .field_hash_table_ref()
321            .ok_or(JournalError::MissingHashTable)?;
322        if table.is_empty() {
323            return Err(JournalError::MissingHashTable.into());
324        }
325        let mut current = table.hash_item_ref(hash).head_hash_offset;
326        let mut previous = 0;
327        while let Some(offset) = current {
328            control.check()?;
329            if offset.get() > number(self.header.tail_object_offset) {
330                return Err(corrupt("FIELD capture exceeds object bounds"));
331            }
332            if offset.get() <= previous {
333                return Err(corrupt("FIELD hash chain does not progress"));
334            }
335            self.check_object(offset, 2)?;
336            let field = self.file.field_ref(offset)?;
337            if field.header.next_hash_offset.is_some_and(|next| {
338                next.get() <= offset.get() || next.get() > number(self.header.tail_object_offset)
339            }) {
340                return Err(corrupt("invalid captured FIELD hash link"));
341            }
342            if field.header.hash == hash && field.payload == name {
343                let head = field.header.head_data_offset;
344                if head.is_none() {
345                    return Err(corrupt("FIELD has no DATA chain"));
346                }
347                drop(field);
348                if let Some(head) = head {
349                    self.check_object(head, 1)?;
350                }
351                return Ok(head);
352            }
353            previous = offset.get();
354            current = field.header.next_hash_offset;
355        }
356        Ok(None)
357    }
358    fn find_data(
359        &self,
360        name: &[u8],
361        value: &[u8],
362        control: &SnapshotControl<'_>,
363        capturing: bool,
364    ) -> Result<Option<NonZeroU64>> {
365        raw_name(name)?;
366        let payload = journal_core::file::PayloadParts::structured(name, value);
367        let hash = self.file.hash_parts(payload);
368        let table = self
369            .file
370            .data_hash_table_ref()
371            .ok_or(JournalError::MissingHashTable)?;
372        if table.is_empty() {
373            return Err(JournalError::MissingHashTable.into());
374        }
375        let mut current = table.hash_item_ref(hash).head_hash_offset;
376        let mut previous = 0;
377        let mut buffer = Vec::new();
378        while let Some(offset) = current {
379            control.check()?;
380            if offset.get() > number(self.header.tail_object_offset) {
381                if capturing {
382                    return Err(corrupt("DATA capture exceeds object bounds"));
383                }
384                break;
385            }
386            if offset.get() <= previous {
387                return Err(corrupt("DATA hash chain does not progress"));
388            }
389            self.check_object(offset, 1)?;
390            let data = self.file.data_ref(offset)?;
391            if capturing
392                && data.header.next_hash_offset.is_some_and(|next| {
393                    next.get() <= offset.get()
394                        || next.get() > number(self.header.tail_object_offset)
395                })
396            {
397                return Err(corrupt("invalid captured DATA hash link"));
398            }
399            if data.header.hash == hash {
400                let bytes = if data.is_compressed() {
401                    buffer.clear();
402                    let len = data.decompress(&mut buffer)?;
403                    &buffer[..len]
404                } else {
405                    data.raw_payload()
406                };
407                if payload.equals_slice(bytes) {
408                    return Ok(Some(offset));
409                }
410            }
411            previous = offset.get();
412            current = data.header.next_hash_offset;
413        }
414        Ok(None)
415    }
416    fn postings(
417        &self,
418        offset: NonZeroU64,
419    ) -> Result<(Option<NonZeroU64>, Option<NonZeroU64>, u64)> {
420        self.check_object(offset, 1)?;
421        let data = self.file.data_ref(offset)?;
422        let count = number(data.header.n_entries);
423        Ok((
424            data.header.entry_offset,
425            data.header.entry_array_offset,
426            count,
427        ))
428    }
429    fn emit(
430        &self,
431        offset: NonZeroU64,
432        control: &SnapshotControl<'_>,
433        visit: &mut impl FnMut(&mut SnapshotEntry<'_>) -> Result<()>,
434    ) -> Result<()> {
435        control.check()?;
436        self.check_entry(offset)?;
437        let (metadata, offsets) = {
438            let entry = self.file.entry_ref(offset)?;
439            let mut offsets = Vec::with_capacity(entry.items.len());
440            entry.collect_offsets(&mut offsets)?;
441            (
442                SnapshotMetadata {
443                    seqnum: entry.header.seqnum,
444                    realtime: entry.header.realtime,
445                    monotonic: entry.header.monotonic,
446                    boot_id: entry.header.boot_id,
447                },
448                offsets,
449            )
450        };
451        let mut entry = SnapshotEntry {
452            snapshot: self,
453            control,
454            metadata,
455            offsets,
456        };
457        visit(&mut entry)
458    }
459    // A live count can be newer than a copied pointer or slot. Writers publish
460    // links before counts: refresh required zeros before treating them as damage.
461    fn required_posting(
462        &self,
463        cached: Option<NonZeroU64>,
464        position: u64,
465        compact: bool,
466    ) -> Result<NonZeroU64> {
467        if let Some(offset) = cached {
468            return Ok(offset);
469        }
470        let mut bytes = [0; 8];
471        let size = if compact { 4 } else { 8 };
472        self.file
473            .read_fresh_bytes_at(position, &mut bytes[..size])?;
474        NonZeroU64::new(u64::from_le_bytes(bytes))
475            .ok_or_else(|| corrupt("missing required posting"))
476    }
477    fn posting_array_offset(
478        &self,
479        current: Option<NonZeroU64>,
480        previous_array: u64,
481        clip: bool,
482    ) -> Result<NonZeroU64> {
483        match current {
484            Some(offset) => Ok(offset),
485            None if clip && previous_array != 0 => {
486                self.required_posting(None, previous_array + 16, false)
487            }
488            None => Err(corrupt("missing posting array")),
489        }
490    }
491    fn posting_slot(
492        &self,
493        cached: Option<NonZeroU64>,
494        array: NonZeroU64,
495        index: usize,
496        clip: bool,
497    ) -> Result<NonZeroU64> {
498        match cached {
499            Some(entry) => Ok(entry),
500            None if clip => {
501                let compact = self.header.incompatible_flags & 16 != 0;
502                let width = if compact { 4 } else { 8 };
503                self.required_posting(None, array.get() + 24 + index as u64 * width, compact)
504            }
505            None => Err(corrupt("zero posting")),
506        }
507    }
508    fn visit_array(
509        &self,
510        mut current: Option<NonZeroU64>,
511        mut remaining: u64,
512        mut last: u64,
513        clip: bool,
514        control: &SnapshotControl<'_>,
515        visit: &mut impl FnMut(&mut SnapshotEntry<'_>) -> Result<()>,
516    ) -> Result<()> {
517        let mut previous_array = 0;
518        let mut chunk = Vec::with_capacity(256);
519        while remaining > 0 {
520            control.check()?;
521            let offset = self.posting_array_offset(current, previous_array, clip)?;
522            if clip && offset.get() > number(self.header.tail_object_offset) {
523                return Ok(());
524            }
525            if offset.get() <= previous_array {
526                return Err(corrupt("array chain does not progress"));
527            }
528            self.check_object(offset, 6)?;
529            let (capacity, next) = {
530                let array = self.file.offset_array_ref(offset)?;
531                (array.capacity(), array.header.next_offset_array)
532            };
533            if capacity == 0 {
534                return Err(corrupt("empty posting array"));
535            }
536            let used = remaining.min(capacity as u64) as usize;
537            for start in (0..used).step_by(256) {
538                control.check()?;
539                chunk.clear();
540                {
541                    let array = self.file.offset_array_ref(offset)?;
542                    for index in start..used.min(start + 256) {
543                        chunk.push(array.items.get(index));
544                    }
545                }
546                for (index, &cached) in chunk.iter().enumerate() {
547                    let entry = self.posting_slot(cached, offset, start + index, clip)?;
548                    if entry.get() <= last {
549                        return Err(corrupt("postings do not progress"));
550                    }
551                    if clip && entry.get() > self.entry_tail {
552                        return Ok(());
553                    }
554                    self.emit(entry, control, visit)?;
555                    last = entry.get();
556                }
557            }
558            remaining -= used as u64;
559            previous_array = offset.get();
560            current = next;
561        }
562        Ok(())
563    }
564    fn visit_data(
565        &self,
566        offset: NonZeroU64,
567        control: &SnapshotControl<'_>,
568        visit: &mut impl FnMut(&mut SnapshotEntry<'_>) -> Result<()>,
569    ) -> Result<()> {
570        let (first, array, count) = self.postings(offset)?;
571        if count == 0 {
572            return Err(corrupt("DATA has no postings"));
573        }
574        let first = first.ok_or_else(|| corrupt("missing inline posting"))?;
575        if first.get() > self.entry_tail {
576            return Ok(());
577        }
578        self.emit(first, control, visit)?;
579        let array = if count > 1 {
580            Some(self.required_posting(array, offset.get() + 48, false)?)
581        } else {
582            array
583        };
584        self.visit_array(array, count - 1, first.get(), true, control, visit)
585    }
586    pub fn visit_match(
587        &mut self,
588        name: &[u8],
589        value: &[u8],
590        control: &SnapshotControl<'_>,
591        mut visit: impl FnMut(&mut SnapshotEntry<'_>) -> Result<()>,
592    ) -> Result<()> {
593        control.check()?;
594        if let Some(offset) = self.find_data(name, value, control, false)? {
595            self.visit_data(offset, control, &mut visit)?;
596        }
597        Ok(())
598    }
599    /// Streams postings for every accepted value; multivalued rows may repeat.
600    /// No union, deduplication, or sorting is performed.
601    pub fn visit_field(
602        &mut self,
603        name: &[u8],
604        control: &SnapshotControl<'_>,
605        mut accept: impl FnMut(&[u8]) -> Result<bool>,
606        mut visit: impl FnMut(&mut SnapshotEntry<'_>) -> Result<()>,
607    ) -> Result<()> {
608        control.check()?;
609        let mut current = *self
610            .fields
611            .get(name)
612            .ok_or(SdkError::Unsupported("undeclared captured field"))?;
613        let mut previous = u64::MAX;
614        let mut buffer = Vec::new();
615        while let Some(offset) = current {
616            control.check()?;
617            if offset.get() >= previous {
618                return Err(corrupt("FIELD DATA chain does not progress"));
619            }
620            self.check_object(offset, 1)?;
621            let (next, accepted) = {
622                let data = self.file.data_ref(offset)?;
623                let bytes = if data.is_compressed() {
624                    buffer.clear();
625                    let len = data.decompress(&mut buffer)?;
626                    &buffer[..len]
627                } else {
628                    data.raw_payload()
629                };
630                let value = bytes
631                    .strip_prefix(name)
632                    .and_then(|rest| rest.strip_prefix(b"="))
633                    .ok_or_else(|| corrupt("FIELD chain payload mismatch"))?;
634                (data.header.next_field_offset, accept(value)?)
635            };
636            if accepted {
637                self.visit_data(offset, control, &mut visit)?;
638            }
639            previous = offset.get();
640            current = next;
641        }
642        Ok(())
643    }
644    /// Lazily streams all captured entries, with O(entries) work and bounded chunks.
645    pub fn visit_entries(
646        &mut self,
647        control: &SnapshotControl<'_>,
648        mut visit: impl FnMut(&mut SnapshotEntry<'_>) -> Result<()>,
649    ) -> Result<()> {
650        control.check()?;
651        self.visit_array(
652            self.header.entry_array_offset,
653            self.header.n_entries,
654            0,
655            false,
656            control,
657            &mut visit,
658        )
659    }
660}