1mod 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 Live,
161 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 pub fields: HashMap<String, Vec<u8>>,
249 pub field_values: HashMap<String, Vec<Vec<u8>>>,
251 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(¤t_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(¤t, 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(¤t, &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;