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 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 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 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}