Skip to main content

journal_core/file/
file_payload.rs

1use super::file::{JournalFile, PayloadParts, validate_offset_alignment};
2use super::mmap::{MemoryMap, WindowManager};
3use super::object::*;
4use crate::error::{JournalError, Result};
5use crate::file::value_guard::ValueGuard;
6use std::num::NonZeroU64;
7use zerocopy::FromBytes;
8
9#[doc(hidden)]
10#[derive(Debug, Clone, Copy)]
11pub struct DataPayloadReadContext {
12    is_compact: bool,
13    header_size: u64,
14    arena_end: u64,
15    payload_prefix_size: u64,
16}
17
18#[doc(hidden)]
19#[derive(Debug, Clone, Copy)]
20pub struct DataPayloadObjectInfo {
21    size_needed: u64,
22    is_compressed: bool,
23}
24
25#[doc(hidden)]
26pub enum RowPinnedPayload<'a> {
27    Borrowed { ptr: *const u8, len: usize },
28    Decompressed(&'a [u8]),
29}
30
31#[doc(hidden)]
32struct DataLookupResult<T> {
33    next_hash_offset: Option<NonZeroU64>,
34    match_value: Option<T>,
35}
36
37#[derive(Debug, Clone, Copy)]
38struct DataLookupHeader {
39    flags: u8,
40    size_needed: u64,
41    stored_hash: u64,
42    next_hash_offset: Option<NonZeroU64>,
43    entry_array_offset: Option<NonZeroU64>,
44    n_entries: Option<NonZeroU64>,
45}
46
47#[derive(Debug, Clone, Copy, PartialEq, Eq)]
48pub(super) struct ResolvedDataLinkState {
49    pub(super) n_entries: Option<NonZeroU64>,
50    pub(super) entry_array_offset: Option<NonZeroU64>,
51    pub(super) compact_tail: Option<(NonZeroU64, u64)>,
52}
53
54impl ResolvedDataLinkState {
55    pub(super) fn empty() -> Self {
56        Self {
57            n_entries: None,
58            entry_array_offset: None,
59            compact_tail: None,
60        }
61    }
62}
63
64impl DataLookupHeader {
65    fn is_compressed(self) -> bool {
66        (self.flags
67            & (ObjectFlags::CompressedZstd as u8
68                | ObjectFlags::CompressedLz4 as u8
69                | ObjectFlags::CompressedXz as u8))
70            != 0
71    }
72}
73
74impl DataPayloadObjectInfo {
75    pub fn is_compressed(self) -> bool {
76        self.is_compressed
77    }
78}
79
80fn parse_data_payload_object_header(header_slice: &[u8]) -> Result<DataPayloadObjectInfo> {
81    let object_header =
82        ObjectHeader::ref_from_bytes(header_slice).map_err(|_| JournalError::ZerocopyFailure)?;
83
84    if object_header.type_ != ObjectType::Data as u8 {
85        return Err(JournalError::InvalidObjectType);
86    }
87
88    Ok(DataPayloadObjectInfo {
89        size_needed: object_header.validated_size()?,
90        is_compressed: object_header.is_compressed(),
91    })
92}
93
94impl<M: MemoryMap> JournalFile<M> {
95    #[doc(hidden)]
96    pub fn data_payload_read_context(&self) -> DataPayloadReadContext {
97        let journal_header = self.journal_header_ref();
98        let is_compact = journal_header.has_incompatible_flag(HeaderIncompatibleFlags::Compact);
99        let payload_prefix_size = std::mem::size_of::<DataObjectHeader>() as u64
100            + if is_compact {
101                std::mem::size_of::<CompactDataFields>() as u64
102            } else {
103                0
104            };
105        DataPayloadReadContext {
106            is_compact,
107            header_size: journal_header.header_size,
108            arena_end: journal_header.header_size + journal_header.arena_size,
109            payload_prefix_size,
110        }
111    }
112
113    #[doc(hidden)]
114    pub fn visit_data_payload_at<F>(
115        &self,
116        offset: NonZeroU64,
117        decompressed: &mut Vec<u8>,
118        visitor: F,
119    ) -> Result<()>
120    where
121        F: FnOnce(&[u8]) -> Result<()>,
122    {
123        let context = self.data_payload_read_context();
124        self.visit_data_payload_at_with_context(context, offset, decompressed, visitor)
125    }
126
127    #[doc(hidden)]
128    pub fn visit_data_payload_at_with_context<F>(
129        &self,
130        context: DataPayloadReadContext,
131        offset: NonZeroU64,
132        decompressed: &mut Vec<u8>,
133        visitor: F,
134    ) -> Result<()>
135    where
136        F: FnOnce(&[u8]) -> Result<()>,
137    {
138        Self::validate_data_payload_offset(context, offset)?;
139        self.window_manager.with_mut(|wm| {
140            let info = Self::data_payload_info_from_window(wm, context, offset)?;
141            let data = Self::data_slice_from_window(wm, offset, info.size_needed)?;
142            if !info.is_compressed {
143                return visitor(&data[context.payload_prefix_size as usize..]);
144            }
145            let object = DataObject::from_data(data, context.is_compact)
146                .ok_or(JournalError::ZerocopyFailure)?;
147            decompressed.clear();
148            let len = object.decompress(decompressed)?;
149            visitor(&decompressed[..len])
150        })
151    }
152
153    #[doc(hidden)]
154    pub fn data_payload_object_info_at(
155        &self,
156        context: DataPayloadReadContext,
157        offset: NonZeroU64,
158    ) -> Result<DataPayloadObjectInfo> {
159        validate_offset_alignment(offset)?;
160        if offset.get() < context.header_size {
161            return Err(JournalError::ObjectExceedsFileBounds);
162        }
163
164        self.window_manager
165            .with_mut(|wm| Self::data_payload_info_from_window(wm, context, offset))
166    }
167
168    fn validate_data_payload_offset(
169        context: DataPayloadReadContext,
170        offset: NonZeroU64,
171    ) -> Result<()> {
172        validate_offset_alignment(offset)?;
173        if offset.get() < context.header_size {
174            return Err(JournalError::ObjectExceedsFileBounds);
175        }
176        Ok(())
177    }
178
179    fn data_payload_info_from_window(
180        wm: &mut WindowManager<M>,
181        context: DataPayloadReadContext,
182        offset: NonZeroU64,
183    ) -> Result<DataPayloadObjectInfo> {
184        let object_header_size = std::mem::size_of::<ObjectHeader>() as u64;
185        let header_slice = wm.get_slice(offset.get(), object_header_size)?;
186        let info = parse_data_payload_object_header(header_slice)?;
187        Self::validate_data_payload_info(context, offset, info)?;
188        Ok(info)
189    }
190
191    fn validate_data_payload_info(
192        context: DataPayloadReadContext,
193        offset: NonZeroU64,
194        info: DataPayloadObjectInfo,
195    ) -> Result<()> {
196        let end_offset = offset
197            .get()
198            .checked_add(info.size_needed)
199            .ok_or(JournalError::ObjectExceedsFileBounds)?;
200        if end_offset > context.arena_end {
201            return Err(JournalError::ObjectExceedsFileBounds);
202        }
203        if info.size_needed < context.payload_prefix_size {
204            return Err(JournalError::InvalidObjectSize(info.size_needed));
205        }
206        Ok(())
207    }
208
209    fn data_slice_from_window<'w>(
210        wm: &'w mut WindowManager<M>,
211        offset: NonZeroU64,
212        size_needed: u64,
213    ) -> Result<&'w [u8]> {
214        if wm.active_window_contains(offset.get(), size_needed) {
215            return Ok(wm.active_slice(offset.get(), size_needed));
216        }
217        wm.get_slice(offset.get(), size_needed)
218    }
219
220    #[doc(hidden)]
221    pub fn raw_data_payload_ref_with_info(
222        &self,
223        context: DataPayloadReadContext,
224        offset: NonZeroU64,
225        info: DataPayloadObjectInfo,
226    ) -> Result<ValueGuard<'_, &[u8]>> {
227        validate_offset_alignment(offset)?;
228        if offset.get() < context.header_size {
229            return Err(JournalError::ObjectExceedsFileBounds);
230        }
231        if info.is_compressed {
232            return Err(JournalError::InvalidObjectType);
233        }
234        if info.size_needed < context.payload_prefix_size {
235            return Err(JournalError::InvalidObjectSize(info.size_needed));
236        }
237
238        self.window_manager.with_guarded(offset, |wm| {
239            if wm.active_window_contains(offset.get(), info.size_needed) {
240                let data = wm.active_slice(offset.get(), info.size_needed);
241                return Ok(&data[context.payload_prefix_size as usize..]);
242            }
243            let data = wm.get_slice(offset.get(), info.size_needed)?;
244            Ok(&data[context.payload_prefix_size as usize..])
245        })
246    }
247
248    #[doc(hidden)]
249    /// Returns an unguarded pointer to an uncompressed DATA payload.
250    ///
251    /// The caller must only expose the pointer while it can prove the backing
252    /// mmap window will not be remapped or evicted. This is intended for
253    /// whole-file mmap row-scoped facade enumeration. Do not call this for
254    /// windowed mmap; use `raw_data_payload_ref_with_info()` or copy the
255    /// payload instead.
256    pub fn raw_data_payload_ptr_with_info_unguarded(
257        &self,
258        context: DataPayloadReadContext,
259        offset: NonZeroU64,
260        info: DataPayloadObjectInfo,
261    ) -> Result<(*const u8, usize)> {
262        validate_offset_alignment(offset)?;
263        if offset.get() < context.header_size {
264            return Err(JournalError::ObjectExceedsFileBounds);
265        }
266        if info.is_compressed {
267            return Err(JournalError::InvalidObjectType);
268        }
269        if info.size_needed < context.payload_prefix_size {
270            return Err(JournalError::InvalidObjectSize(info.size_needed));
271        }
272
273        self.window_manager.with_mut(|wm| {
274            let data =
275                if let Some(data) = wm.active_slice_if_contains(offset.get(), info.size_needed) {
276                    data
277                } else {
278                    wm.get_slice(offset.get(), info.size_needed)?
279                };
280            let payload = &data[context.payload_prefix_size as usize..];
281            Ok((payload.as_ptr(), payload.len()))
282        })
283    }
284
285    #[doc(hidden)]
286    /// Returns a pointer to an uncompressed DATA payload and pins the backing
287    /// mmap window until row pins are explicitly cleared.
288    pub fn raw_data_payload_ptr_with_info_row_pinned(
289        &self,
290        context: DataPayloadReadContext,
291        offset: NonZeroU64,
292        info: DataPayloadObjectInfo,
293    ) -> Result<(*const u8, usize)> {
294        validate_offset_alignment(offset)?;
295        if offset.get() < context.header_size {
296            return Err(JournalError::ObjectExceedsFileBounds);
297        }
298        if info.is_compressed {
299            return Err(JournalError::InvalidObjectType);
300        }
301        if info.size_needed < context.payload_prefix_size {
302            return Err(JournalError::InvalidObjectSize(info.size_needed));
303        }
304
305        self.window_manager.with_mut(|wm| {
306            let data = wm.get_row_pinned_slice(offset.get(), info.size_needed)?;
307            let payload = &data[context.payload_prefix_size as usize..];
308            Ok((payload.as_ptr(), payload.len()))
309        })
310    }
311
312    #[doc(hidden)]
313    /// Returns a row-pinned pointer when the DATA object is uncompressed.
314    /// Compressed DATA returns `Ok(None)` so the caller can take the
315    /// decompression path.
316    pub fn raw_data_payload_ptr_row_pinned_if_uncompressed(
317        &self,
318        context: DataPayloadReadContext,
319        offset: NonZeroU64,
320    ) -> Result<Option<(*const u8, usize)>> {
321        Self::validate_data_payload_offset(context, offset)?;
322
323        self.window_manager.with_mut(|wm| {
324            let info = Self::data_payload_info_from_window(wm, context, offset)?;
325            if info.is_compressed {
326                return Ok(None);
327            }
328            let data = wm.get_row_pinned_slice(offset.get(), info.size_needed)?;
329            let payload = &data[context.payload_prefix_size as usize..];
330            Ok(Some((payload.as_ptr(), payload.len())))
331        })
332    }
333
334    #[doc(hidden)]
335    pub fn clear_row_payload_pins(&self) -> Result<()> {
336        self.window_manager.with_mut(|wm| {
337            wm.clear_row_pins();
338            Ok(())
339        })
340    }
341
342    #[doc(hidden)]
343    pub fn visit_data_payloads_row_pinned_with_context<F>(
344        &self,
345        context: DataPayloadReadContext,
346        offsets: &[NonZeroU64],
347        decompressed: &mut Vec<u8>,
348        mut visitor: F,
349    ) -> Result<()>
350    where
351        F: FnMut(RowPinnedPayload<'_>) -> Result<()>,
352    {
353        self.window_manager.with_mut(|wm| {
354            for offset in offsets.iter().copied() {
355                Self::visit_data_payload_row_pinned_from_window(
356                    wm,
357                    context,
358                    offset,
359                    decompressed,
360                    &mut visitor,
361                )?;
362            }
363            Ok(())
364        })
365    }
366
367    fn visit_data_payload_row_pinned_from_window<F>(
368        wm: &mut WindowManager<M>,
369        context: DataPayloadReadContext,
370        offset: NonZeroU64,
371        decompressed: &mut Vec<u8>,
372        visitor: &mut F,
373    ) -> Result<()>
374    where
375        F: FnMut(RowPinnedPayload<'_>) -> Result<()>,
376    {
377        Self::validate_data_payload_offset(context, offset)?;
378        let info = Self::data_payload_info_from_window(wm, context, offset)?;
379        if info.is_compressed {
380            return Self::visit_compressed_row_payload(
381                wm,
382                context,
383                offset,
384                info,
385                decompressed,
386                visitor,
387            );
388        }
389        Self::visit_borrowed_row_payload(wm, context, offset, info, visitor)
390    }
391
392    fn visit_borrowed_row_payload<F>(
393        wm: &mut WindowManager<M>,
394        context: DataPayloadReadContext,
395        offset: NonZeroU64,
396        info: DataPayloadObjectInfo,
397        visitor: &mut F,
398    ) -> Result<()>
399    where
400        F: FnMut(RowPinnedPayload<'_>) -> Result<()>,
401    {
402        let data = wm.get_row_pinned_slice(offset.get(), info.size_needed)?;
403        let payload = &data[context.payload_prefix_size as usize..];
404        visitor(RowPinnedPayload::Borrowed {
405            ptr: payload.as_ptr(),
406            len: payload.len(),
407        })
408    }
409
410    fn visit_compressed_row_payload<F>(
411        wm: &mut WindowManager<M>,
412        context: DataPayloadReadContext,
413        offset: NonZeroU64,
414        info: DataPayloadObjectInfo,
415        decompressed: &mut Vec<u8>,
416        visitor: &mut F,
417    ) -> Result<()>
418    where
419        F: FnMut(RowPinnedPayload<'_>) -> Result<()>,
420    {
421        let data = Self::data_slice_from_window(wm, offset, info.size_needed)?;
422        let object =
423            DataObject::from_data(data, context.is_compact).ok_or(JournalError::ZerocopyFailure)?;
424        decompressed.clear();
425        let len = object.decompress(decompressed)?;
426        visitor(RowPinnedPayload::Decompressed(&decompressed[..len]))
427    }
428
429    pub fn find_data_offset(&self, hash: u64, payload: &[u8]) -> Result<Option<NonZeroU64>> {
430        self.find_data_offset_parts(hash, PayloadParts::raw(payload))
431    }
432
433    pub fn find_data_offset_parts(
434        &self,
435        hash: u64,
436        payload: PayloadParts<'_>,
437    ) -> Result<Option<NonZeroU64>> {
438        Ok(self
439            .find_data_match_parts(hash, payload, |_, _, _| ())?
440            .map(|(offset, ())| offset))
441    }
442
443    pub(super) fn find_data_with_link_state_parts(
444        &self,
445        hash: u64,
446        payload: PayloadParts<'_>,
447    ) -> Result<Option<(NonZeroU64, ResolvedDataLinkState)>> {
448        self.find_data_match_parts(hash, payload, Self::resolved_data_link_state)
449    }
450
451    fn find_data_match_parts<T>(
452        &self,
453        hash: u64,
454        payload: PayloadParts<'_>,
455        matched: impl Fn(DataPayloadReadContext, DataLookupHeader, &[u8]) -> T,
456    ) -> Result<Option<(NonZeroU64, T)>> {
457        let hash_table = self
458            .data_hash_table_ref()
459            .ok_or(JournalError::MissingHashTable)?;
460        let context = self.data_payload_read_context();
461        let mut decompression_buffer = Vec::new();
462        let mut object_offset = hash_table.hash_item_ref(hash).head_hash_offset;
463
464        while let Some(offset) = object_offset {
465            let result = self.data_lookup_result_at(
466                context,
467                offset,
468                hash,
469                payload,
470                &mut decompression_buffer,
471                &matched,
472            )?;
473            if let Some(match_value) = result.match_value {
474                return Ok(Some((offset, match_value)));
475            }
476            object_offset = result.next_hash_offset;
477        }
478
479        Ok(None)
480    }
481
482    fn data_lookup_result_at<T>(
483        &self,
484        context: DataPayloadReadContext,
485        offset: NonZeroU64,
486        hash: u64,
487        payload: PayloadParts<'_>,
488        decompression_buffer: &mut Vec<u8>,
489        matched: &impl Fn(DataPayloadReadContext, DataLookupHeader, &[u8]) -> T,
490    ) -> Result<DataLookupResult<T>> {
491        Self::validate_data_payload_offset(context, offset)?;
492        self.window_manager.with_mut(|wm| {
493            let lookup = Self::data_lookup_header_from_window(wm, context, offset)?;
494            if lookup.stored_hash != hash {
495                return Ok(DataLookupResult {
496                    next_hash_offset: lookup.next_hash_offset,
497                    match_value: None,
498                });
499            }
500
501            let data = Self::data_slice_from_window(wm, offset, lookup.size_needed)?;
502            let matches = Self::data_lookup_payload_matches(
503                context,
504                lookup,
505                data,
506                payload,
507                decompression_buffer,
508            )?;
509            let match_value = matches.then(|| matched(context, lookup, data));
510            Ok(DataLookupResult {
511                next_hash_offset: lookup.next_hash_offset,
512                match_value,
513            })
514        })
515    }
516
517    fn data_lookup_header_from_window(
518        wm: &mut WindowManager<M>,
519        context: DataPayloadReadContext,
520        offset: NonZeroU64,
521    ) -> Result<DataLookupHeader> {
522        let header_slice =
523            wm.get_slice(offset.get(), std::mem::size_of::<DataObjectHeader>() as u64)?;
524        Self::parse_data_lookup_header(context, offset, header_slice)
525    }
526
527    fn parse_data_lookup_header(
528        context: DataPayloadReadContext,
529        offset: NonZeroU64,
530        header_slice: &[u8],
531    ) -> Result<DataLookupHeader> {
532        if header_slice[0] != ObjectType::Data as u8 {
533            return Err(JournalError::InvalidObjectType);
534        }
535        let size_needed = u64::from_le_bytes(header_slice[8..16].try_into().unwrap());
536        if size_needed < std::mem::size_of::<DataObjectHeader>() as u64 {
537            return Err(JournalError::InvalidObjectSize(size_needed));
538        }
539        let info = DataPayloadObjectInfo {
540            size_needed,
541            is_compressed: false,
542        };
543        Self::validate_data_payload_info(context, offset, info)?;
544        Ok(DataLookupHeader {
545            flags: header_slice[1],
546            size_needed,
547            stored_hash: u64::from_le_bytes(header_slice[16..24].try_into().unwrap()),
548            next_hash_offset: NonZeroU64::new(u64::from_le_bytes(
549                header_slice[24..32].try_into().unwrap(),
550            )),
551            entry_array_offset: Self::optional_nonzero_u64_field::<DataObjectHeader>(
552                header_slice,
553                std::mem::offset_of!(DataObjectHeader, entry_array_offset),
554            ),
555            n_entries: Self::optional_nonzero_u64_field::<DataObjectHeader>(
556                header_slice,
557                std::mem::offset_of!(DataObjectHeader, n_entries),
558            ),
559        })
560    }
561
562    fn optional_nonzero_u64_field<T>(data: &[u8], offset: usize) -> Option<NonZeroU64> {
563        debug_assert!(offset + std::mem::size_of::<u64>() <= std::mem::size_of::<T>());
564        NonZeroU64::new(u64::from_le_bytes(
565            data[offset..offset + std::mem::size_of::<u64>()]
566                .try_into()
567                .unwrap(),
568        ))
569    }
570
571    fn resolved_data_link_state(
572        context: DataPayloadReadContext,
573        lookup: DataLookupHeader,
574        data: &[u8],
575    ) -> ResolvedDataLinkState {
576        let compact_tail = context.is_compact.then(|| {
577            let fields_offset = std::mem::size_of::<DataObjectHeader>();
578            let tail_offset = NonZeroU64::new(u32::from_le_bytes(
579                data[fields_offset..fields_offset + 4].try_into().unwrap(),
580            ) as u64)?;
581            let tail_entries = u32::from_le_bytes(
582                data[fields_offset + 4..fields_offset + 8]
583                    .try_into()
584                    .unwrap(),
585            ) as u64;
586            (tail_entries != 0).then_some((tail_offset, tail_entries))
587        });
588        ResolvedDataLinkState {
589            n_entries: lookup.n_entries,
590            entry_array_offset: lookup.entry_array_offset,
591            compact_tail: compact_tail.flatten(),
592        }
593    }
594
595    fn data_lookup_payload_matches(
596        context: DataPayloadReadContext,
597        lookup: DataLookupHeader,
598        data: &[u8],
599        payload: PayloadParts<'_>,
600        decompression_buffer: &mut Vec<u8>,
601    ) -> Result<bool> {
602        if lookup.is_compressed() {
603            let object = DataObject::from_data(data, context.is_compact)
604                .ok_or(JournalError::ZerocopyFailure)?;
605            decompression_buffer.clear();
606            let len = object.decompress(decompression_buffer)?;
607            return Ok(payload.equals_slice(&decompression_buffer[..len]));
608        }
609        let payload_start = context.payload_prefix_size as usize;
610        Ok(payload.equals_slice(&data[payload_start..]))
611    }
612}