Skip to main content

journal/
lib.rs

1//! Pure-Rust systemd journal reader and writer SDK.
2//!
3//! This crate provides a public Rust layer over the imported Netdata journal
4//! reader/writer crates. It intentionally keeps the low-level file parsing in
5//! the imported implementation and adds byte-safe entries, directory reading,
6//! export/JSON formatting, and a libsystemd-style facade.
7
8mod indexed_snapshot;
9pub use indexed_snapshot::{
10    CapturedValue, IndexedSnapshot, IndexedSnapshotOptions, SnapshotControl, SnapshotEntry,
11    SnapshotMetadata,
12};
13
14mod directory;
15mod explorer;
16mod export;
17mod facade;
18pub mod netdata;
19mod parse;
20mod reader_helpers;
21mod sealed_verify;
22mod verify_graph;
23
24pub use directory::DirectoryReader;
25pub use explorer::{
26    ExplorerAnchor, ExplorerComparison, ExplorerControl, ExplorerFieldMode, ExplorerFilter,
27    ExplorerFtsPattern, ExplorerHistogram, ExplorerHistogramBucket, ExplorerProgress,
28    ExplorerQuery, ExplorerResult, ExplorerRow, ExplorerSampling, ExplorerStats,
29    ExplorerStopReason, ExplorerStrategy,
30};
31pub use export::{export_entry, export_entry_bytes, format_entry_text, json_entry};
32pub use parse::{ParseError, ParsedCursor, parse_cursor, parse_match_bytes, parse_match_string};
33pub use sealed_verify::{verify_file, verify_file_with_key, verify_index};
34
35use ouroboros::self_referencing;
36use std::collections::HashMap;
37use std::fmt;
38use std::num::NonZeroU64;
39use std::path::{Path, PathBuf};
40
41use directory::DirectoryEntryKey;
42#[cfg(test)]
43use directory::is_journal_file_name;
44use reader_helpers::*;
45#[cfg(test)]
46use sealed_verify::{
47    COMPACT_DATA_OBJECT_HEADER_SIZE, DATA_OBJECT_HEADER_SIZE, HEADER_MIN_SIZE,
48    INCOMPATIBLE_COMPACT, OBJECT_HEADER_SIZE, OBJECT_TYPE_DATA, OBJECT_TYPE_TAG, align8,
49};
50
51pub use facade::{
52    ERR_END_OF_ENTRIES, ERR_INVALID_CURSOR, ERR_NO_ENTRY, ERR_UNSUPPORTED, Error as FacadeError,
53    OutputMode, SdJournal, SdJournalAddConjunction, SdJournalAddDisjunction, SdJournalAddMatch,
54    SdJournalClose, SdJournalEnumerateAvailableData, SdJournalEnumerateAvailableUnique,
55    SdJournalEnumerateField, SdJournalEnumerateFields, SdJournalFlushMatches, SdJournalGetCursor,
56    SdJournalGetData, SdJournalGetEntry, SdJournalGetMonotonicUsec, SdJournalGetRealtimeUsec,
57    SdJournalGetSeqnum, SdJournalListBoots, SdJournalNext, SdJournalNextSkip, SdJournalOpen,
58    SdJournalOpenDirectory, SdJournalOpenDirectoryWithOptions, SdJournalOpenFile,
59    SdJournalOpenFileWithOptions, SdJournalOpenFiles, SdJournalOpenFilesWithOptions,
60    SdJournalPrevious, SdJournalPreviousSkip, SdJournalProcessOutput, SdJournalQueryUnique,
61    SdJournalQueryUniqueState, SdJournalRestartData, SdJournalRestartFields,
62    SdJournalRestartUnique, SdJournalSeekCursor, SdJournalSeekHead, SdJournalSeekRealtimeUsec,
63    SdJournalSeekTail, SdJournalSetOutputMode, SdJournalTestCursor, SdJournalVisitUniqueValues,
64};
65pub use journal_core::error::JournalError;
66use journal_core::file::ExperimentalMmapStrategy;
67pub use journal_core::file::{
68    BucketUtilization, Compression, Direction, EntryItemsType, FieldNamePolicy, HashableObject,
69    JournalFile, JournalReader, Location, Mmap, WindowManagerStats,
70};
71use journal_core::file::{CurrentRowMetadata, CurrentRowView};
72pub use journal_log_writer::{
73    Config, EntryTimestamps, Log, LogLifecycleEvent, LogLifecycleObserver, RetentionPolicy,
74    RotationPolicy, WriterError,
75};
76pub use journal_registry::{Origin, Source};
77
78pub type Result<T> = std::result::Result<T, SdkError>;
79
80#[derive(Debug)]
81pub enum SdkError {
82    Journal(JournalError),
83    InvalidPath(String),
84    InvalidCursor(String),
85    NoEntry,
86    Cancelled,
87    DecompressionFailed(String),
88    Unsupported(&'static str),
89    VerificationError(String),
90}
91
92impl fmt::Display for SdkError {
93    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
94        match self {
95            Self::Journal(err) => write!(f, "{err}"),
96            Self::InvalidPath(path) => write!(f, "invalid path: {path}"),
97            Self::InvalidCursor(cursor) => write!(f, "invalid cursor: {cursor}"),
98            Self::Cancelled => write!(f, "operation cancelled"),
99            Self::NoEntry => write!(f, "no entry at current position"),
100            Self::DecompressionFailed(err) => write!(f, "decompression failed: {err}"),
101            Self::Unsupported(op) => write!(f, "unsupported operation: {op}"),
102            Self::VerificationError(msg) => {
103                write!(f, "journal verification failed: corrupt file: {msg}")
104            }
105        }
106    }
107}
108
109impl std::error::Error for SdkError {}
110
111impl From<JournalError> for SdkError {
112    fn from(err: JournalError) -> Self {
113        Self::Journal(err)
114    }
115}
116
117impl From<std::io::Error> for SdkError {
118    fn from(err: std::io::Error) -> Self {
119        Self::Journal(JournalError::Io(err))
120    }
121}
122
123#[derive(Debug, Clone)]
124pub struct Field {
125    pub name: String,
126    pub value: Vec<u8>,
127}
128
129impl Field {
130    pub fn new(name: &str, value: &str) -> Self {
131        Self {
132            name: name.to_string(),
133            value: value.as_bytes().to_vec(),
134        }
135    }
136
137    pub fn with_bytes(name: &str, value: Vec<u8>) -> Self {
138        Self {
139            name: name.to_string(),
140            value,
141        }
142    }
143
144    pub fn payload(&self) -> Vec<u8> {
145        let mut payload = Vec::with_capacity(self.name.len() + 1 + self.value.len());
146        payload.extend_from_slice(self.name.as_bytes());
147        payload.push(b'=');
148        payload.extend_from_slice(&self.value);
149        payload
150    }
151}
152
153#[derive(Debug, Clone, Copy, PartialEq, Eq)]
154pub enum ReaderBounds {
155    /// Systemd-style mutable reader bounds.
156    ///
157    /// The reader keeps a cached file size and refreshes it only when a read
158    /// would go beyond the cached end of file, matching libsystemd's active
159    /// journal behavior without a metadata syscall on every object read.
160    Live,
161    /// Immutable reader bounds.
162    ///
163    /// The reader fixes the file size at open time, like
164    /// `SD_JOURNAL_ASSUME_IMMUTABLE`, for polling/query consumers that do not
165    /// need to observe appends during the current scan.
166    Snapshot,
167}
168
169pub const DEFAULT_READER_WINDOW_SIZE: u64 = 32 * 1024 * 1024;
170
171impl Default for ReaderBounds {
172    fn default() -> Self {
173        Self::Live
174    }
175}
176
177#[derive(Debug, Clone, Copy, PartialEq, Eq)]
178pub struct ReaderOptions {
179    pub window_size: u64,
180    pub bounds: ReaderBounds,
181    mmap_strategy: ExperimentalMmapStrategy,
182}
183
184impl Default for ReaderOptions {
185    fn default() -> Self {
186        Self {
187            window_size: DEFAULT_READER_WINDOW_SIZE,
188            bounds: ReaderBounds::Live,
189            mmap_strategy: ExperimentalMmapStrategy::Windowed,
190        }
191    }
192}
193
194impl ReaderOptions {
195    pub fn live() -> Self {
196        Self::default()
197    }
198
199    pub fn snapshot() -> Self {
200        Self {
201            bounds: ReaderBounds::Snapshot,
202            ..Self::default()
203        }
204    }
205
206    pub fn with_window_size(mut self, window_size: u64) -> Self {
207        self.window_size = window_size;
208        self
209    }
210
211    pub fn with_bounds(mut self, bounds: ReaderBounds) -> Self {
212        self.bounds = bounds;
213        self
214    }
215
216    #[doc(hidden)]
217    pub fn with_experimental_mmap_strategy(mut self, strategy: ExperimentalMmapStrategy) -> Self {
218        self.mmap_strategy = strategy;
219        self
220    }
221}
222
223#[derive(Debug, Clone, Copy, PartialEq, Eq)]
224pub struct RawField<'a> {
225    pub name: &'a [u8],
226    pub value: &'a [u8],
227}
228
229impl RawField<'_> {
230    pub fn payload(&self) -> Vec<u8> {
231        let mut payload = Vec::with_capacity(self.name.len() + 1 + self.value.len());
232        payload.extend_from_slice(self.name);
233        payload.push(b'=');
234        payload.extend_from_slice(self.value);
235        payload
236    }
237
238    pub fn name_str(&self) -> Option<&str> {
239        std::str::from_utf8(self.name).ok()
240    }
241}
242
243#[derive(Debug, Clone)]
244pub struct Entry {
245    /// Convenience map for UTF-8 field names. RAW-mode files may contain field
246    /// names that are not valid UTF-8; use `raw_fields()` or `get_raw_values()`
247    /// when byte-identical field-name identity matters.
248    pub fields: HashMap<String, Vec<u8>>,
249    /// Convenience repeated-value map for UTF-8 field names.
250    pub field_values: HashMap<String, Vec<Vec<u8>>>,
251    /// Full on-disk DATA payloads as `FIELD=value` bytes.
252    pub payloads: Vec<Vec<u8>>,
253    pub seqnum: u64,
254    pub realtime: u64,
255    pub monotonic: u64,
256    pub boot_id: [u8; 16],
257    pub cursor: String,
258}
259
260impl Entry {
261    pub fn get(&self, key: &str) -> Option<&[u8]> {
262        self.fields.get(key).map(Vec::as_slice)
263    }
264
265    pub fn get_str(&self, key: &str) -> Option<&str> {
266        self.get(key)
267            .and_then(|value| std::str::from_utf8(value).ok())
268    }
269
270    pub fn raw_fields(&self) -> impl Iterator<Item = RawField<'_>> {
271        self.payloads
272            .iter()
273            .filter_map(|payload| split_raw_payload(payload))
274    }
275
276    pub fn get_raw(&self, key: &[u8]) -> Option<&[u8]> {
277        self.raw_fields()
278            .find(|field| field.name == key)
279            .map(|field| field.value)
280    }
281
282    pub fn get_raw_values(&self, key: &[u8]) -> Vec<&[u8]> {
283        self.raw_fields()
284            .filter_map(|field| (field.name == key).then_some(field.value))
285            .collect()
286    }
287}
288
289fn split_raw_payload(payload: &[u8]) -> Option<RawField<'_>> {
290    let eq = payload.iter().position(|byte| *byte == b'=')?;
291    Some(RawField {
292        name: &payload[..eq],
293        value: &payload[eq + 1..],
294    })
295}
296
297#[derive(Debug, Clone)]
298pub struct BootInfo {
299    pub index: i64,
300    pub boot_id: String,
301    pub first_entry: i64,
302    pub last_entry: i64,
303}
304
305#[derive(Debug, Clone, Copy)]
306pub struct FileHeader {
307    pub signature: [u8; 8],
308    pub compatible_flags: u32,
309    pub incompatible_flags: u32,
310    pub state: u8,
311    pub file_id: [u8; 16],
312    pub machine_id: [u8; 16],
313    pub header_size: u64,
314    pub arena_size: u64,
315    pub data_hash_table_size: u64,
316    pub field_hash_table_size: u64,
317    pub n_objects: u64,
318    pub n_entries: u64,
319    pub head_entry_realtime: u64,
320    pub tail_entry_realtime: u64,
321    pub tail_entry_monotonic: u64,
322    pub head_entry_seqnum: u64,
323    pub tail_entry_seqnum: u64,
324    pub tail_entry_boot_id: [u8; 16],
325    pub seqnum_id: [u8; 16],
326    pub n_data: u64,
327    pub n_fields: u64,
328    pub n_tags: u64,
329    pub n_entry_arrays: u64,
330    pub data_hash_chain_depth: u64,
331    pub field_hash_chain_depth: u64,
332}
333
334#[derive(Debug, Clone, Copy)]
335pub(crate) struct FileHeaderSnapshot {
336    pub(crate) header: FileHeader,
337}
338
339impl FileHeaderSnapshot {
340    fn from_file(file: &JournalFile<Mmap>) -> Self {
341        let header = file.journal_header_ref();
342        Self {
343            header: FileHeader {
344                signature: header.signature,
345                compatible_flags: header.compatible_flags,
346                incompatible_flags: header.incompatible_flags,
347                state: header.state,
348                file_id: header.file_id,
349                machine_id: header.machine_id,
350                header_size: header.header_size,
351                arena_size: header.arena_size,
352                data_hash_table_size: header
353                    .data_hash_table_size
354                    .map(|value| value.get())
355                    .unwrap_or(0),
356                field_hash_table_size: header
357                    .field_hash_table_size
358                    .map(|value| value.get())
359                    .unwrap_or(0),
360                n_objects: header.n_objects,
361                n_entries: header.n_entries,
362                head_entry_realtime: header.head_entry_realtime,
363                tail_entry_realtime: header.tail_entry_realtime,
364                tail_entry_monotonic: header.tail_entry_monotonic,
365                head_entry_seqnum: header.head_entry_seqnum,
366                tail_entry_seqnum: header.tail_entry_seqnum,
367                tail_entry_boot_id: header.tail_entry_boot_id,
368                seqnum_id: header.seqnum_id,
369                n_data: header.n_data,
370                n_fields: header.n_fields,
371                n_tags: header.n_tags,
372                n_entry_arrays: header.n_entry_arrays,
373                data_hash_chain_depth: header.data_hash_chain_depth,
374                field_hash_chain_depth: header.field_hash_chain_depth,
375            },
376        }
377    }
378}
379
380#[self_referencing]
381struct ReaderCell {
382    file: JournalFile<Mmap>,
383    #[borrows(file)]
384    #[not_covariant]
385    reader: JournalReader<'this, Mmap>,
386}
387
388pub struct FileReader {
389    inner: ReaderCell,
390    temp_path: Option<PathBuf>,
391    row: CurrentRowView,
392    header_snapshot: FileHeaderSnapshot,
393    bounds: ReaderBounds,
394}
395
396fn key_from_metadata(metadata: CurrentRowMetadata) -> DirectoryEntryKey {
397    DirectoryEntryKey {
398        seqnum_id: metadata.seqnum_id,
399        seqnum: metadata.seqnum,
400        boot_id: metadata.boot_id,
401        monotonic: metadata.monotonic,
402        realtime: metadata.realtime,
403        xor_hash: metadata.xor_hash,
404    }
405}
406
407enum StepStatus {
408    Valid,
409    Skip,
410    End,
411}
412
413impl Drop for FileReader {
414    fn drop(&mut self) {
415        self.inner
416            .with_file(|file| self.row.clear_current_best_effort(file));
417        if let Some(path) = &self.temp_path {
418            let _ = std::fs::remove_file(path);
419        }
420    }
421}
422
423impl FileReader {
424    pub fn open(path: impl AsRef<Path>) -> Result<Self> {
425        Self::open_with_options(path, ReaderOptions::default())
426    }
427
428    pub fn open_with_options(path: impl AsRef<Path>, options: ReaderOptions) -> Result<Self> {
429        let path = path.as_ref();
430        if is_zst_file(path) {
431            return Self::open_zst(path, options);
432        }
433
434        let file = open_journal_file(path, options)?;
435        let header_snapshot = FileHeaderSnapshot::from_file(&file);
436        Ok(Self {
437            inner: ReaderCellBuilder {
438                file,
439                reader_builder: |_file| JournalReader::default(),
440            }
441            .build(),
442            temp_path: None,
443            row: CurrentRowView::default(),
444            header_snapshot,
445            bounds: options.bounds,
446        })
447    }
448
449    fn open_zst(path: &Path, options: ReaderOptions) -> Result<Self> {
450        let temp_path = decompress_zst_to_temp(path, "rust-sdk-journal")?;
451        let file = match open_journal_file(&temp_path, options) {
452            Ok(file) => file,
453            Err(err) => {
454                let _ = std::fs::remove_file(&temp_path);
455                return Err(err);
456            }
457        };
458        let header_snapshot = FileHeaderSnapshot::from_file(&file);
459        Ok(Self {
460            inner: ReaderCellBuilder {
461                file,
462                reader_builder: |_file| JournalReader::default(),
463            }
464            .build(),
465            temp_path: Some(temp_path),
466            row: CurrentRowView::default(),
467            header_snapshot,
468            bounds: options.bounds,
469        })
470    }
471
472    pub fn header(&self) -> FileHeader {
473        if self.bounds == ReaderBounds::Snapshot {
474            return self.header_snapshot.header;
475        }
476        self.live_header()
477    }
478
479    pub(crate) fn cached_header(&self) -> FileHeaderSnapshot {
480        self.header_snapshot
481    }
482
483    fn live_header(&self) -> FileHeader {
484        self.inner
485            .with_file(|file| FileHeaderSnapshot::from_file(file).header)
486    }
487
488    pub fn bucket_utilization(&self) -> Option<BucketUtilization> {
489        self.inner.with_file(JournalFile::bucket_utilization)
490    }
491
492    #[doc(hidden)]
493    pub fn mmap_stats(&self) -> Result<WindowManagerStats> {
494        self.inner
495            .with_file(|file| file.mmap_stats())
496            .map_err(Into::into)
497    }
498
499    pub fn seek_head(&mut self) {
500        self.inner
501            .with_file(|file| self.row.clear_current_best_effort(file));
502        self.inner.with_reader_mut(|reader| {
503            reader.set_location(Location::Head);
504        });
505    }
506
507    pub fn seek_tail(&mut self) {
508        self.inner
509            .with_file(|file| self.row.clear_current_best_effort(file));
510        self.inner.with_reader_mut(|reader| {
511            reader.set_location(Location::Tail);
512        });
513    }
514
515    pub fn seek_realtime(&mut self, usec: u64) {
516        self.inner
517            .with_file(|file| self.row.clear_current_best_effort(file));
518        self.inner.with_reader_mut(|reader| {
519            reader.set_location(Location::Realtime(usec));
520        });
521    }
522
523    pub fn seek_cursor(&mut self, cursor: &str) -> Result<()> {
524        let want = parse::parse_cursor_location(cursor, true)
525            .map_err(|err| SdkError::InvalidCursor(err.to_string()))?;
526        if want.realtime_set {
527            self.seek_realtime(want.realtime);
528        } else {
529            self.seek_head();
530        }
531        while self.next()? {
532            let current_cursor = self.get_cursor()?;
533            let got = parse::parse_cursor_location(&current_cursor, false)
534                .map_err(|err| SdkError::InvalidCursor(err.to_string()))?;
535            if parse::cursor_location_at_or_after(&got, &want) {
536                return Ok(());
537            }
538        }
539        self.seek_tail();
540        Ok(())
541    }
542
543    pub fn next(&mut self) -> Result<bool> {
544        self.step_valid(Direction::Forward)
545    }
546
547    pub fn previous(&mut self) -> Result<bool> {
548        self.step_valid(Direction::Backward)
549    }
550
551    fn step_valid(&mut self, direction: Direction) -> Result<bool> {
552        self.inner
553            .with_file(|file| self.row.clear_current(file))
554            .map_err(SdkError::from)?;
555        loop {
556            let row = &mut self.row;
557            let status = self.inner.with_mut(|fields| {
558                if !fields.reader.step(fields.file, direction)? {
559                    return Ok(StepStatus::End);
560                }
561
562                match fields
563                    .reader
564                    .get_entry_offset()
565                    .and_then(|offset| row.load_entry(fields.file, offset))
566                {
567                    Ok(_) => Ok(StepStatus::Valid),
568                    Err(err) if recoverable_entry_error(&err) => Ok(StepStatus::Skip),
569                    Err(err) => Err(err),
570                }
571            })?;
572
573            match status {
574                StepStatus::Valid => {
575                    return Ok(true);
576                }
577                StepStatus::Skip => continue,
578                StepStatus::End => {
579                    self.inner
580                        .with_file(|file| self.row.clear_current(file))
581                        .map_err(SdkError::from)?;
582                    return Ok(false);
583                }
584            }
585        }
586    }
587
588    pub fn get_entry(&mut self) -> Result<Entry> {
589        self.invalidate_entry_data_state();
590        let inner = &mut self.inner;
591        let row = &mut self.row;
592        inner.with_mut(|fields| {
593            if row.entry_offset().is_none() {
594                let offset = fields.reader.get_entry_offset()?;
595                row.load_entry(fields.file, offset)?;
596            }
597            read_current_row_entry(fields.file, row)
598        })
599    }
600
601    pub fn visit_entry_payloads<F>(&mut self, mut visitor: F) -> Result<()>
602    where
603        F: FnMut(&[u8]) -> Result<()>,
604    {
605        self.invalidate_entry_data_state();
606        let inner = &mut self.inner;
607        let row = &mut self.row;
608        inner.with_mut(|fields| {
609            fields.reader.release_object_guards();
610            if row.entry_offset().is_none() {
611                let offset = fields.reader.get_entry_offset()?;
612                row.load_entry(fields.file, offset)?;
613            }
614            row.restart_data()?;
615            loop {
616                let payload = match row.read_next_payload(fields.file) {
617                    Ok(Some(payload)) => payload,
618                    Ok(None) => break,
619                    Err(err) if recoverable_entry_data_error(&err) => continue,
620                    Err(err) => {
621                        let _ = row.reset_data_state(fields.file);
622                        return Err(err.into());
623                    }
624                };
625                let payload = row.payload_slice(payload);
626                if let Err(err) = visitor(payload) {
627                    let _ = row.reset_data_state(fields.file);
628                    return Err(err);
629                }
630            }
631            row.reset_data_state(fields.file)?;
632            Ok(())
633        })
634    }
635
636    pub fn clear_entry_data_state(&mut self) {
637        self.inner
638            .with_file(|file| self.row.reset_data_state_best_effort(file));
639        self.inner
640            .with_reader_mut(|reader| reader.entry_data_restart());
641    }
642
643    fn invalidate_entry_data_state(&mut self) {
644        if self.row.data_state_active() {
645            self.clear_entry_data_state();
646        }
647    }
648
649    pub fn entry_data_restart(&mut self) -> Result<()> {
650        self.inner
651            .with_file(|file| self.row.clear_pins(file))
652            .map_err(SdkError::from)?;
653        self.inner
654            .with_reader_mut(|reader| reader.entry_data_restart());
655        if self.row.entry_offset().is_none() {
656            let row = &mut self.row;
657            self.inner.with_mut(|fields| {
658                let offset = fields.reader.get_entry_offset()?;
659                row.load_entry(fields.file, offset).map(|_| ())
660            })?;
661        }
662        self.row.restart_data().map_err(Into::into)
663    }
664
665    pub fn enumerate_entry_payload(&mut self) -> Result<Option<&[u8]>> {
666        let row = &mut self.row;
667        let payload = self.inner.with_mut(|fields| {
668            fields.reader.release_object_guards();
669            row.read_next_payload(fields.file)
670        })?;
671        Ok(payload.map(|payload| self.row.payload_slice(payload)))
672    }
673
674    pub fn collect_entry_payloads(&mut self, payloads: &mut Vec<Vec<u8>>) -> Result<()> {
675        payloads.clear();
676        self.visit_entry_payloads(|payload| {
677            payloads.push(payload.to_vec());
678            Ok(())
679        })
680    }
681
682    pub fn get_entry_payload(&mut self, field: &[u8]) -> Result<Option<Vec<u8>>> {
683        let mut found = None;
684        self.visit_entry_payloads(|payload| {
685            if found.is_none()
686                && payload.len() > field.len()
687                && payload.starts_with(field)
688                && payload[field.len()] == b'='
689            {
690                found = Some(payload.to_vec());
691            }
692            Ok(())
693        })?;
694        Ok(found)
695    }
696
697    pub fn get_realtime_usec(&self) -> Result<u64> {
698        if let Some(metadata) = self.row.metadata() {
699            return Ok(metadata.realtime);
700        }
701        self.inner
702            .with(|fields| fields.reader.get_realtime_usec(fields.file))
703            .map_err(Into::into)
704    }
705
706    pub fn get_seqnum(&self) -> Result<(u64, [u8; 16])> {
707        let key = self.current_directory_entry_key()?;
708        Ok((key.seqnum, key.seqnum_id))
709    }
710
711    pub fn get_monotonic_usec(&self) -> Result<(u64, [u8; 16])> {
712        let key = self.current_directory_entry_key()?;
713        Ok((key.monotonic, key.boot_id))
714    }
715
716    pub fn get_cursor(&self) -> Result<String> {
717        if let Some(metadata) = self.row.metadata() {
718            return Ok(format_cursor_from_key(key_from_metadata(metadata)));
719        }
720        let seqnum_id = self.header_snapshot.header.seqnum_id;
721        self.inner
722            .with(|fields| build_cursor(fields.file, fields.reader, seqnum_id))
723    }
724
725    fn current_directory_entry_key(&self) -> Result<DirectoryEntryKey> {
726        if let Some(metadata) = self.row.metadata() {
727            return Ok(key_from_metadata(metadata));
728        }
729        self.inner.with(|fields| {
730            let offset = fields.reader.get_entry_offset()?;
731            let entry = fields.file.entry_ref(offset)?;
732            Ok(DirectoryEntryKey {
733                seqnum_id: self.header_snapshot.header.seqnum_id,
734                seqnum: entry.header.seqnum,
735                boot_id: entry.header.boot_id,
736                monotonic: entry.header.monotonic,
737                realtime: entry.header.realtime,
738                xor_hash: entry.header.xor_hash,
739            })
740        })
741    }
742
743    pub fn test_cursor(&self, cursor: &str) -> Result<bool> {
744        let current = match self.get_cursor() {
745            Ok(cursor) => cursor,
746            Err(SdkError::Journal(JournalError::UnsetCursor)) => return Err(SdkError::NoEntry),
747            Err(err) => return Err(err),
748        };
749        let current = parse::parse_cursor_location(&current, false)
750            .map_err(|err| SdkError::InvalidCursor(err.to_string()))?;
751        let Ok(want) = parse::parse_cursor_location(cursor, false) else {
752            return Ok(false);
753        };
754        Ok(parse::cursor_location_matches(&current, &want))
755    }
756
757    pub fn add_match(&mut self, data: &[u8]) {
758        self.inner.with_reader_mut(|reader| reader.add_match(data));
759    }
760
761    pub fn add_conjunction(&mut self) -> Result<()> {
762        self.inner
763            .with_mut(|fields| fields.reader.add_conjunction(fields.file))
764            .map_err(Into::into)
765    }
766
767    pub fn add_disjunction(&mut self) -> Result<()> {
768        self.inner
769            .with_mut(|fields| fields.reader.add_disjunction(fields.file))
770            .map_err(Into::into)
771    }
772
773    pub fn flush_matches(&mut self) {
774        self.inner.with_reader_mut(|reader| reader.flush_matches());
775    }
776}
777
778impl FileReader {
779    fn header_realtime_start(&self) -> u64 {
780        self.header_snapshot.header.head_entry_realtime
781    }
782
783    pub fn enumerate_fields(&mut self) -> Result<Vec<String>> {
784        self.invalidate_entry_data_state();
785        match self.enumerate_fields_indexed() {
786            Ok(fields) => Ok(fields),
787            Err(_) => enumerate_file_fields_by_scan(self),
788        }
789    }
790
791    pub(crate) fn enumerate_fields_indexed(&mut self) -> Result<Vec<String>> {
792        self.invalidate_entry_data_state();
793        self.inner.with_file(enumerate_file_fields_indexed)
794    }
795
796    pub fn query_unique(&mut self, field_name: &str) -> Result<Vec<Vec<u8>>> {
797        let mut out = Vec::new();
798        self.visit_unique_values(field_name, |value| {
799            out.push(value.to_vec());
800            Ok(())
801        })?;
802        Ok(out)
803    }
804
805    pub fn visit_unique_values<F>(&mut self, field_name: &str, visitor: F) -> Result<()>
806    where
807        F: FnMut(&[u8]) -> Result<()>,
808    {
809        self.invalidate_entry_data_state();
810        let decompressed = self.row.decompressed_mut();
811        self.inner.with_file(|file| {
812            visit_file_unique_values_indexed(file, field_name.as_bytes(), decompressed, visitor)
813        })
814    }
815
816    pub fn query_unique_state(&mut self, field_name: &str) -> Result<()> {
817        self.invalidate_entry_data_state();
818        self.inner.with_mut(|fields| {
819            fields
820                .reader
821                .field_data_query_unique(fields.file, field_name.as_bytes())
822        })?;
823        Ok(())
824    }
825
826    pub fn restart_unique_state(&mut self) {
827        self.inner
828            .with_reader_mut(|reader| reader.field_data_restart());
829    }
830
831    pub fn clear_unique_state(&mut self) {
832        self.inner
833            .with_reader_mut(|reader| reader.field_data_clear());
834    }
835
836    pub fn enumerate_unique_payload(&mut self, field_name: &str) -> Result<Option<Vec<u8>>> {
837        let field = field_name.as_bytes();
838        let decompressed = self.row.decompressed_mut();
839        self.inner.with_mut(|fields| {
840            let Some(data) = fields.reader.field_data_enumerate(fields.file)? else {
841                return Ok(None);
842            };
843            let payload = if data.is_compressed() {
844                decompressed.clear();
845                let len = data.decompress(decompressed)?;
846                &decompressed[..len]
847            } else {
848                data.raw_payload()
849            };
850            let Some(_) = payload
851                .strip_prefix(field)
852                .and_then(|rest| rest.strip_prefix(b"="))
853            else {
854                return Err(SdkError::VerificationError(
855                    "field DATA chain object does not match requested field".to_string(),
856                ));
857            };
858            Ok(Some(payload.to_vec()))
859        })
860    }
861}
862
863#[cfg(test)]
864mod tests;