1use super::file_payload::ResolvedDataLinkState;
2use super::mmap::MmapMut;
3use crate::error::{JournalError, Result};
4use crate::file::{
5 Compression, DEFAULT_COMPRESS_THRESHOLD, DataObject, DataPayloadType, EntryObjectHeader,
6 FieldObjectHeader, HashableObjectMut, HeaderIncompatibleFlags, JournalFile, JournalHeader,
7 ObjectFlags, ObjectType, PayloadParts, hash::jenkins_hash64_parts,
8 normalize_compress_threshold,
9};
10use rustc_hash::FxHashMap;
11use std::io::Cursor;
12use std::num::NonZeroU64;
13
14pub(super) const OBJECT_ALIGNMENT: u64 = 8;
15pub(super) const JOURNAL_COMPACT_SIZE_MAX: u64 = u32::MAX as u64;
16const FILE_SIZE_INCREASE: u64 = 8 * 1024 * 1024;
17const FIELD_CACHE_MAX_ENTRIES: usize = 1024;
18const FIELD_CACHE_MAX_PAYLOAD_LEN: usize = 128;
19pub(super) fn round_up_to_file_size_increment(value: u64) -> Result<u64> {
20 value
21 .checked_add(FILE_SIZE_INCREASE - 1)
22 .map(|v| v & !(FILE_SIZE_INCREASE - 1))
23 .ok_or(JournalError::ObjectExceedsFileBounds)
24}
25
26#[derive(Debug, Clone, Copy)]
27pub(super) struct EntryItem {
28 pub(super) offset: NonZeroU64,
29 pub(super) hash: u64,
30 pub(super) link_state: Option<ResolvedDataLinkState>,
31}
32
33#[derive(Debug, Clone, Copy)]
34pub struct StructuredField<'a> {
35 pub name: &'a [u8],
37 pub value: &'a [u8],
39}
40
41impl<'a> StructuredField<'a> {
42 pub fn new(name: &'a [u8], value: &'a [u8]) -> Self {
44 Self { name, value }
45 }
46}
47
48#[derive(Debug, Clone, Copy)]
49pub enum EntryField<'a> {
50 Raw(&'a [u8]),
52 Structured(StructuredField<'a>),
55}
56
57impl<'a> EntryField<'a> {
58 pub fn raw(payload: &'a [u8]) -> Self {
60 Self::Raw(payload)
61 }
62
63 pub fn structured(name: &'a [u8], value: &'a [u8]) -> Self {
65 Self::Structured(StructuredField::new(name, value))
66 }
67
68 fn payload_parts(self) -> PayloadParts<'a> {
69 match self {
70 Self::Raw(payload) => PayloadParts::raw(payload),
71 Self::Structured(field) => PayloadParts::structured(field.name, field.value),
72 }
73 }
74
75 fn field_name(self) -> Option<&'a [u8]> {
76 match self {
77 Self::Raw(payload) => payload
78 .iter()
79 .position(|&b| b == b'=')
80 .map(|pos| &payload[..pos]),
81 Self::Structured(field) => Some(field.name),
82 }
83 }
84}
85
86#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
87pub enum FieldNamePolicy {
88 #[default]
91 Journald,
92 Raw,
95 JournalApp,
98}
99
100#[derive(Debug, Clone, Copy, Default)]
101pub struct EntryWriteOptions {
102 pub trusted_unique_payloads: bool,
111 pub field_name_policy: FieldNamePolicy,
113 pub seqnum: Option<u64>,
119 pub boot_id: Option<uuid::Uuid>,
124}
125
126impl EntryWriteOptions {
127 pub fn trusted_unique_payloads(mut self, enabled: bool) -> Self {
132 self.trusted_unique_payloads = enabled;
133 self
134 }
135
136 pub fn field_name_policy(mut self, policy: FieldNamePolicy) -> Self {
138 self.field_name_policy = policy;
139 self
140 }
141
142 pub fn seqnum(mut self, seqnum: u64) -> Self {
144 self.seqnum = Some(seqnum);
145 self
146 }
147
148 pub fn boot_id(mut self, boot_id: uuid::Uuid) -> Self {
150 self.boot_id = Some(boot_id);
151 self
152 }
153}
154
155fn is_journal_field_name_valid(field_name: &[u8], allow_protected: bool) -> bool {
156 if field_name.is_empty() || field_name.len() > 64 {
157 return false;
158 }
159 if field_name[0] == b'_' && !allow_protected {
160 return false;
161 }
162 if field_name[0].is_ascii_digit() {
163 return false;
164 }
165 field_name
166 .iter()
167 .all(|&b| b.is_ascii_uppercase() || b.is_ascii_digit() || b == b'_')
168}
169
170fn is_raw_field_name_valid(field_name: &[u8]) -> bool {
171 !field_name.is_empty() && !field_name.contains(&b'=')
172}
173
174fn accept_entry_field(field: EntryField<'_>, policy: FieldNamePolicy) -> Result<bool> {
175 let Some(field_name) = field.field_name() else {
176 return Err(JournalError::InvalidField);
177 };
178 let valid = match policy {
179 FieldNamePolicy::Raw => is_raw_field_name_valid(field_name),
180 FieldNamePolicy::Journald => is_journal_field_name_valid(field_name, true),
181 FieldNamePolicy::JournalApp => is_journal_field_name_valid(field_name, false),
182 };
183 if valid {
184 return Ok(true);
185 }
186 if matches!(policy, FieldNamePolicy::JournalApp) {
187 return Ok(false);
188 }
189 Err(JournalError::InvalidField)
190}
191
192#[derive(Debug)]
193struct FieldCache {
194 entries: FxHashMap<Box<[u8]>, NonZeroU64>,
195}
196
197impl FieldCache {
198 fn new() -> Self {
199 Self {
200 entries: FxHashMap::default(),
201 }
202 }
203
204 fn get(&self, payload: &[u8]) -> Option<NonZeroU64> {
205 self.entries.get(payload).copied()
206 }
207
208 fn insert(&mut self, payload: &[u8], offset: NonZeroU64) {
209 if payload.len() > FIELD_CACHE_MAX_PAYLOAD_LEN {
210 return;
211 }
212
213 if self.entries.len() >= FIELD_CACHE_MAX_ENTRIES && self.entries.get(payload).is_none() {
214 self.entries.clear();
215 }
216
217 self.entries
218 .insert(payload.to_vec().into_boxed_slice(), offset);
219 }
220
221 #[cfg(test)]
222 fn len(&self) -> usize {
223 self.entries.len()
224 }
225}
226
227enum StoredDataPayload<'a> {
228 Uncompressed(PayloadParts<'a>),
229 Compressed(Vec<u8>, u8),
230}
231
232impl StoredDataPayload<'_> {
233 fn len(&self) -> usize {
234 match self {
235 Self::Uncompressed(payload) => payload.len(),
236 Self::Compressed(payload, _) => payload.len(),
237 }
238 }
239
240 fn object_flags(&self) -> u8 {
241 match self {
242 Self::Uncompressed(_) => 0,
243 Self::Compressed(_, flags) => *flags,
244 }
245 }
246
247 fn copy_to_data_object(&self, data: &mut DataObject<&mut [u8]>) {
248 match self {
249 Self::Uncompressed(payload) => match &mut data.payload {
250 DataPayloadType::Regular(dst) => payload.copy_to_slice(dst),
251 DataPayloadType::Compact { payload: dst, .. } => payload.copy_to_slice(dst),
252 },
253 Self::Compressed(payload, _) => data.set_payload(payload),
254 }
255 }
256}
257
258pub struct JournalWriter {
259 pub(super) tail_object_offset: NonZeroU64,
260 pub(super) append_offset: NonZeroU64,
261 next_seqnum: u64,
262 num_written_objects: u64,
263 pub(super) first_tag_written: bool,
264 pub(super) entry_items: Vec<EntryItem>,
265 field_cache: FieldCache,
266 first_entry_monotonic: Option<u64>,
267 boot_id: uuid::Uuid,
268 compression: Compression,
269 compress_threshold: usize,
270 live_publish_every_entries: u64,
271 entries_since_live_publication: u64,
272 pub(super) seal: Option<crate::seal::SealState>,
273}
274
275impl JournalWriter {
276 pub fn current_file_size(&self) -> u64 {
278 self.append_offset.get()
279 }
280
281 pub fn first_entry_monotonic(&self) -> Option<u64> {
283 self.first_entry_monotonic
284 }
285
286 pub fn next_seqnum(&self) -> u64 {
288 self.next_seqnum
289 }
290
291 pub fn boot_id(&self) -> uuid::Uuid {
293 self.boot_id
294 }
295
296 pub fn set_live_publish_every_entries(&mut self, entries: u64) {
304 self.live_publish_every_entries = entries;
305 self.entries_since_live_publication = 0;
306 }
307
308 pub fn live_publish_every_entries(&self) -> u64 {
313 self.live_publish_every_entries
314 }
315
316 pub fn new(
317 journal_file: &mut JournalFile<MmapMut>,
318 next_seqnum: u64,
319 boot_id: uuid::Uuid,
320 ) -> Result<Self> {
321 let compression = match journal_file.journal_header_ref() {
322 header if header.has_incompatible_flag(HeaderIncompatibleFlags::CompressedZstd) => {
323 Compression::Zstd
324 }
325 header if header.has_incompatible_flag(HeaderIncompatibleFlags::CompressedXz) => {
326 Compression::Xz
327 }
328 header if header.has_incompatible_flag(HeaderIncompatibleFlags::CompressedLz4) => {
329 Compression::Lz4
330 }
331 _ => Compression::None,
332 };
333
334 Self::new_with_compression(
335 journal_file,
336 next_seqnum,
337 boot_id,
338 compression,
339 DEFAULT_COMPRESS_THRESHOLD,
340 )
341 }
342
343 pub fn new_with_compression(
344 journal_file: &mut JournalFile<MmapMut>,
345 next_seqnum: u64,
346 boot_id: uuid::Uuid,
347 compression: Compression,
348 compress_threshold: usize,
349 ) -> Result<Self> {
350 let current_header_size = std::mem::size_of::<JournalHeader>() as u64;
351 let header = journal_file.journal_header_ref();
352 if !header.has_incompatible_flag(HeaderIncompatibleFlags::KeyedHash)
353 || header.header_size < current_header_size
354 {
355 return Err(JournalError::UnsupportedJournalFile);
356 }
357
358 let append_offset = {
359 let header = journal_file.journal_header_ref();
360
361 let Some(tail_object_offset) = header.tail_object_offset else {
362 return Err(JournalError::InvalidMagicNumber);
363 };
364
365 let tail_object = journal_file.object_header_ref(tail_object_offset)?;
366
367 tail_object_offset.saturating_add(tail_object.size)
368 };
369
370 let seal = journal_file
371 .seal_options
372 .as_ref()
373 .map(|opts| crate::seal::SealState::new(opts))
374 .transpose()?;
375
376 let mut writer = Self {
377 tail_object_offset: journal_file
378 .journal_header_ref()
379 .tail_object_offset
380 .unwrap(),
381 append_offset,
382 next_seqnum,
383 num_written_objects: 0,
384 first_tag_written: false,
385 entry_items: Vec::with_capacity(128),
386 field_cache: FieldCache::new(),
387 first_entry_monotonic: None,
388 boot_id,
389 compression,
390 compress_threshold: normalize_compress_threshold(compress_threshold),
391 live_publish_every_entries: 1,
392 entries_since_live_publication: 0,
393 seal,
394 };
395
396 if writer.seal.is_some() && journal_file.journal_header_ref().n_tags == 0 {
397 writer.ensure_first_tag(journal_file)?;
398 {
399 let header = journal_file.journal_header_mut();
400 header.n_objects += writer.num_written_objects;
401 header.tail_object_offset = Some(writer.tail_object_offset);
402 }
403 writer.num_written_objects = 0;
404 }
405
406 Ok(writer)
407 }
408
409 pub fn create_successor(&self, journal_file: &mut JournalFile<MmapMut>) -> Result<Self> {
411 Self::new_with_compression(
412 journal_file,
413 self.next_seqnum,
414 self.boot_id,
415 self.compression,
416 self.compress_threshold,
417 )
418 }
419
420 pub fn add_entry(
421 &mut self,
422 journal_file: &mut JournalFile<MmapMut>,
423 items: &[&[u8]],
424 realtime: u64,
425 monotonic: u64,
426 ) -> Result<()> {
427 self.add_entry_fields_with_options(
428 journal_file,
429 items.iter().copied().map(EntryField::raw),
430 realtime,
431 monotonic,
432 EntryWriteOptions::default(),
433 )
434 }
435
436 pub fn add_entry_structured(
437 &mut self,
438 journal_file: &mut JournalFile<MmapMut>,
439 fields: &[StructuredField<'_>],
440 realtime: u64,
441 monotonic: u64,
442 ) -> Result<()> {
443 self.add_entry_fields_with_options(
444 journal_file,
445 fields.iter().copied().map(EntryField::Structured),
446 realtime,
447 monotonic,
448 EntryWriteOptions::default(),
449 )
450 }
451
452 pub fn add_entry_structured_with_options(
453 &mut self,
454 journal_file: &mut JournalFile<MmapMut>,
455 fields: &[StructuredField<'_>],
456 realtime: u64,
457 monotonic: u64,
458 options: EntryWriteOptions,
459 ) -> Result<()> {
460 self.add_entry_fields_with_options(
461 journal_file,
462 fields.iter().copied().map(EntryField::Structured),
463 realtime,
464 monotonic,
465 options,
466 )
467 }
468
469 pub fn add_entry_fields<'a>(
470 &mut self,
471 journal_file: &mut JournalFile<MmapMut>,
472 fields: impl IntoIterator<Item = EntryField<'a>>,
473 realtime: u64,
474 monotonic: u64,
475 ) -> Result<()> {
476 self.add_entry_fields_with_options(
477 journal_file,
478 fields,
479 realtime,
480 monotonic,
481 EntryWriteOptions::default(),
482 )
483 }
484
485 pub fn add_entry_fields_with_options<'a>(
486 &mut self,
487 journal_file: &mut JournalFile<MmapMut>,
488 fields: impl IntoIterator<Item = EntryField<'a>>,
489 realtime: u64,
490 monotonic: u64,
491 options: EntryWriteOptions,
492 ) -> Result<()> {
493 self.ensure_keyed_append(journal_file)?;
494 let entry_seqnum = self.entry_seqnum_for_options(options)?;
495 let entry_boot_id = options.boot_id.unwrap_or(self.boot_id);
496 let monotonic = self.clamp_same_boot_monotonic(journal_file, entry_boot_id, monotonic)?;
497 let xor_hash = self.prepare_entry_items(journal_file, fields, realtime, options)?;
498 let entry_offset = self.write_entry_object(
499 journal_file,
500 entry_seqnum,
501 entry_boot_id,
502 realtime,
503 monotonic,
504 xor_hash,
505 )?;
506 self.publish_entry_links(journal_file, entry_offset)?;
507 self.entry_added(
508 journal_file.journal_header_mut(),
509 entry_offset,
510 entry_seqnum,
511 entry_boot_id,
512 realtime,
513 monotonic,
514 );
515 self.publish_after_entry(journal_file)
516 }
517
518 fn ensure_keyed_append(&self, journal_file: &JournalFile<MmapMut>) -> Result<()> {
519 let header = journal_file.journal_header_ref();
520 if header.has_incompatible_flag(HeaderIncompatibleFlags::KeyedHash) {
521 return Ok(());
522 }
523 Err(JournalError::UnsupportedJournalFile)
524 }
525
526 fn entry_seqnum_for_options(&self, options: EntryWriteOptions) -> Result<u64> {
527 let entry_seqnum = options.seqnum.unwrap_or(self.next_seqnum);
528 if entry_seqnum == 0 || entry_seqnum == u64::MAX || entry_seqnum < self.next_seqnum {
529 return Err(JournalError::InvalidField);
530 }
531 Ok(entry_seqnum)
532 }
533
534 fn clamp_same_boot_monotonic(
535 &self,
536 journal_file: &JournalFile<MmapMut>,
537 entry_boot_id: uuid::Uuid,
538 monotonic: u64,
539 ) -> Result<u64> {
540 let header = journal_file.journal_header_ref();
541 if header.n_entries == 0
542 || header.tail_entry_boot_id != *entry_boot_id.as_bytes()
543 || monotonic > header.tail_entry_monotonic
544 {
545 return Ok(monotonic);
546 }
547 header
548 .tail_entry_monotonic
549 .checked_add(1)
550 .ok_or(JournalError::InvalidField)
551 }
552
553 fn prepare_entry_items<'a>(
554 &mut self,
555 journal_file: &mut JournalFile<MmapMut>,
556 fields: impl IntoIterator<Item = EntryField<'a>>,
557 realtime: u64,
558 options: EntryWriteOptions,
559 ) -> Result<u64> {
560 let mut xor_hash = 0;
561 self.entry_items.clear();
562 let mut publication_ready = false;
563 for field in fields {
564 if !accept_entry_field(field, options.field_name_policy)? {
565 continue;
566 }
567 self.ensure_entry_publication_ready(journal_file, realtime, &mut publication_ready)?;
568 xor_hash ^= self.add_entry_field_item(journal_file, field)?;
569 }
570 self.finish_entry_items(options.trusted_unique_payloads)?;
571 Ok(xor_hash)
572 }
573
574 fn ensure_entry_publication_ready(
575 &mut self,
576 journal_file: &mut JournalFile<MmapMut>,
577 realtime: u64,
578 publication_ready: &mut bool,
579 ) -> Result<()> {
580 if *publication_ready {
581 return Ok(());
582 }
583 self.ensure_first_tag(journal_file)?;
584 self.maybe_append_tag(journal_file, realtime)?;
585 *publication_ready = true;
586 Ok(())
587 }
588
589 fn add_entry_field_item(
590 &mut self,
591 journal_file: &mut JournalFile<MmapMut>,
592 field: EntryField<'_>,
593 ) -> Result<u64> {
594 let entry_item = self.add_data(journal_file, field)?;
595 self.entry_items.push(entry_item);
596 Ok(jenkins_hash64_parts(field.payload_parts().iter()))
597 }
598
599 fn finish_entry_items(&mut self, trusted_unique_payloads: bool) -> Result<()> {
600 if self.entry_items.is_empty() {
601 return Err(JournalError::InvalidField);
602 }
603 if !self.entry_items_are_sorted() {
604 self.entry_items
605 .sort_unstable_by(|a, b| a.offset.cmp(&b.offset));
606 }
607 if !trusted_unique_payloads {
608 self.entry_items.dedup_by(|a, b| a.offset == b.offset);
609 } else {
610 for index in 1..self.entry_items.len() {
611 if self.entry_items[index - 1].offset == self.entry_items[index].offset {
612 self.entry_items[index].link_state = None;
613 }
614 }
615 }
616 Ok(())
617 }
618
619 fn entry_items_are_sorted(&self) -> bool {
620 self.entry_items
621 .windows(2)
622 .all(|items| items[0].offset <= items[1].offset)
623 }
624
625 fn write_entry_object(
626 &mut self,
627 journal_file: &mut JournalFile<MmapMut>,
628 entry_seqnum: u64,
629 entry_boot_id: uuid::Uuid,
630 realtime: u64,
631 monotonic: u64,
632 xor_hash: u64,
633 ) -> Result<NonZeroU64> {
634 let entry_offset = self.append_offset;
635 let is_compact = Self::is_compact(journal_file);
636 let entry_payload_size = self.entry_items.len() as u64 * Self::entry_item_size(is_compact);
637 Self::ensure_compact_object_fits(
638 is_compact,
639 entry_offset,
640 std::mem::size_of::<EntryObjectHeader>() as u64 + entry_payload_size,
641 )?;
642 let entry_size = {
643 let size = Some(entry_payload_size);
644 let mut entry_guard = journal_file.entry_mut(entry_offset, size)?;
645
646 entry_guard.header.seqnum = entry_seqnum;
647 entry_guard.header.xor_hash = xor_hash;
648 entry_guard.header.boot_id = *entry_boot_id.as_bytes();
649 entry_guard.header.monotonic = monotonic;
650 entry_guard.header.realtime = realtime;
651
652 for (index, entry_item) in self.entry_items.iter().enumerate() {
654 Self::ensure_compact_offset(is_compact, entry_item.offset)?;
655 let item_hash = (!is_compact).then_some(entry_item.hash);
656 entry_guard.items.set(index, entry_item.offset, item_hash);
657 }
658
659 entry_guard.header.object_header.aligned_size()
660 };
661 self.hmac_put_object(journal_file, entry_offset.get(), ObjectType::Entry)?;
662 self.object_added(journal_file, entry_offset, entry_size)?;
663 Ok(entry_offset)
664 }
665
666 fn publish_entry_links(
667 &mut self,
668 journal_file: &mut JournalFile<MmapMut>,
669 entry_offset: NonZeroU64,
670 ) -> Result<()> {
671 self.append_to_entry_array(journal_file, entry_offset)?;
672 for entry_item_index in 0..self.entry_items.len() {
673 self.link_data_to_entry(journal_file, entry_offset, entry_item_index)?;
674 }
675 Ok(())
676 }
677
678 fn publish_after_entry(&mut self, journal_file: &mut JournalFile<MmapMut>) -> Result<()> {
679 match self.live_publish_every_entries {
680 0 => Ok(()),
681 1 => journal_file.post_change(),
682 interval => {
683 self.entries_since_live_publication += 1;
684 if self.entries_since_live_publication >= interval {
685 self.entries_since_live_publication = 0;
686 journal_file.post_change()
687 } else {
688 Ok(())
689 }
690 }
691 }
692 }
693
694 pub(super) fn object_added(
695 &mut self,
696 journal_file: &mut JournalFile<MmapMut>,
697 object_offset: NonZeroU64,
698 object_size: u64,
699 ) -> Result<()> {
700 self.tail_object_offset = object_offset;
701 self.append_offset = object_offset
702 .checked_add(object_size)
703 .ok_or(JournalError::ObjectExceedsFileBounds)?;
704 self.num_written_objects += 1;
705
706 let header = journal_file.journal_header_mut();
707 let old_size = header
708 .header_size
709 .checked_add(header.arena_size)
710 .ok_or(JournalError::ObjectExceedsFileBounds)?;
711 if self.append_offset.get() > old_size {
712 let new_size = round_up_to_file_size_increment(self.append_offset.get())?;
713 header.arena_size = new_size
714 .checked_sub(header.header_size)
715 .ok_or(JournalError::ObjectExceedsFileBounds)?;
716 }
717
718 Ok(())
719 }
720
721 fn entry_added(
722 &mut self,
723 header: &mut JournalHeader,
724 entry_offset: NonZeroU64,
725 entry_seqnum: u64,
726 entry_boot_id: uuid::Uuid,
727 realtime: u64,
728 monotonic: u64,
729 ) {
730 header.n_objects += self.num_written_objects;
731 header.tail_object_offset = Some(self.tail_object_offset);
732
733 if header.head_entry_seqnum == 0 {
734 header.head_entry_seqnum = entry_seqnum;
735 }
736 if header.head_entry_realtime == 0 {
737 header.head_entry_realtime = realtime;
738 }
739 if self.first_entry_monotonic.is_none() {
740 self.first_entry_monotonic = Some(monotonic);
741 }
742
743 header.tail_entry_seqnum = entry_seqnum;
744 header.tail_entry_realtime = realtime;
745 header.tail_entry_monotonic = monotonic;
746 header.tail_entry_boot_id = *entry_boot_id.as_bytes();
747 header.tail_entry_offset = entry_offset.get();
748 header.n_entries += 1;
749
750 self.next_seqnum = entry_seqnum + 1;
751 self.num_written_objects = 0;
752 }
753
754 fn add_data(
755 &mut self,
756 journal_file: &mut JournalFile<MmapMut>,
757 field: EntryField<'_>,
758 ) -> Result<EntryItem> {
759 let payload = field.payload_parts();
760 let field_name = field.field_name().ok_or(JournalError::InvalidField)?;
761 let hash = journal_file.hash_parts(payload);
762 if let Some((data_offset, link_state)) =
763 journal_file.find_data_with_link_state_parts(hash, payload)?
764 {
765 return Ok(Self::entry_item(data_offset, hash, link_state));
766 }
767 self.add_new_data(journal_file, payload, field_name, hash)
768 }
769
770 fn entry_item(offset: NonZeroU64, hash: u64, link_state: ResolvedDataLinkState) -> EntryItem {
771 EntryItem {
772 offset,
773 hash,
774 link_state: Some(link_state),
775 }
776 }
777
778 fn add_new_data<'a>(
779 &mut self,
780 journal_file: &mut JournalFile<MmapMut>,
781 payload: PayloadParts<'a>,
782 field_name: &'a [u8],
783 hash: u64,
784 ) -> Result<EntryItem> {
785 let data_offset = self.write_new_data_object(journal_file, payload, hash)?;
786 self.publish_new_data_object(journal_file, data_offset, hash)?;
787 self.link_data_to_field(journal_file, data_offset, field_name)?;
788 Ok(Self::entry_item(
789 data_offset,
790 hash,
791 ResolvedDataLinkState::empty(),
792 ))
793 }
794
795 fn write_new_data_object<'a>(
796 &mut self,
797 journal_file: &mut JournalFile<MmapMut>,
798 payload: PayloadParts<'a>,
799 hash: u64,
800 ) -> Result<NonZeroU64> {
801 let data_offset = self.append_offset;
802 let stored_payload = self.stored_data_payload(payload);
803 self.ensure_data_object_fits(journal_file, data_offset, stored_payload.len() as u64)?;
804 let data_size = {
805 let mut data_guard =
806 journal_file.data_mut(data_offset, Some(stored_payload.len() as u64))?;
807 data_guard.header.hash = hash;
808 stored_payload.copy_to_data_object(&mut data_guard);
809 data_guard.header.object_header.flags = stored_payload.object_flags();
810 data_guard.header.object_header.aligned_size()
811 };
812 self.hmac_put_object(journal_file, data_offset.get(), ObjectType::Data)?;
813 self.object_added(journal_file, data_offset, data_size)?;
814 Ok(data_offset)
815 }
816
817 fn ensure_data_object_fits(
818 &self,
819 journal_file: &JournalFile<MmapMut>,
820 data_offset: NonZeroU64,
821 payload_size: u64,
822 ) -> Result<()> {
823 let is_compact = Self::is_compact(journal_file);
824 Self::ensure_compact_object_fits(
825 is_compact,
826 data_offset,
827 Self::data_object_size(is_compact, payload_size),
828 )
829 }
830
831 fn publish_new_data_object(
832 &mut self,
833 journal_file: &mut JournalFile<MmapMut>,
834 data_offset: NonZeroU64,
835 hash: u64,
836 ) -> Result<()> {
837 journal_file.data_hash_table_set_tail_offset(hash, data_offset)?;
838 Self::update_data_hash_chain_depth(journal_file, hash)?;
839 journal_file.journal_header_mut().n_data += 1;
840 Ok(())
841 }
842
843 fn link_data_to_field(
844 &mut self,
845 journal_file: &mut JournalFile<MmapMut>,
846 data_offset: NonZeroU64,
847 field_name: &[u8],
848 ) -> Result<()> {
849 let field_offset = self.add_field(journal_file, field_name)?;
850 let head_data_offset = {
851 let field_guard = journal_file.field_ref(field_offset)?;
852 field_guard.header.head_data_offset
853 };
854 {
855 let mut data_guard = journal_file.data_mut(data_offset, None)?;
856 data_guard.header.next_field_offset = head_data_offset;
857 }
858 let mut field_guard = journal_file.field_mut(field_offset, None)?;
859 field_guard.header.head_data_offset = Some(data_offset);
860 Ok(())
861 }
862
863 fn stored_data_payload<'a>(&self, payload: PayloadParts<'a>) -> StoredDataPayload<'a> {
864 if payload.len() >= self.compress_threshold {
865 let full_payload;
866 let payload_bytes = if let Some(raw) = payload.as_single_slice() {
867 raw
868 } else {
869 full_payload = payload.to_vec();
872 full_payload.as_slice()
873 };
874 match self.compression {
875 Compression::Zstd => {
876 let compressed = ruzstd::encoding::compress_to_vec(
877 Cursor::new(payload_bytes),
878 ruzstd::encoding::CompressionLevel::Fastest,
879 );
880 let compressed = zstd_frame_with_content_size(compressed, payload_bytes.len());
881 if compressed.len() < payload_bytes.len() {
882 return StoredDataPayload::Compressed(
883 compressed,
884 ObjectFlags::CompressedZstd as u8,
885 );
886 }
887 }
888 Compression::Xz => {
889 if payload_bytes.len() >= 80 {
890 if let Ok(compressed) = xz_compress(payload_bytes) {
891 if compressed.len() < payload_bytes.len() {
892 return StoredDataPayload::Compressed(
893 compressed,
894 ObjectFlags::CompressedXz as u8,
895 );
896 }
897 }
898 }
899 }
900 Compression::Lz4 => {
901 if payload_bytes.len() >= 9 {
902 let compressed = lz4_compress(payload_bytes);
903 if compressed.len() < payload_bytes.len() {
904 return StoredDataPayload::Compressed(
905 compressed,
906 ObjectFlags::CompressedLz4 as u8,
907 );
908 }
909 }
910 }
911 Compression::None => {}
912 }
913 }
914
915 StoredDataPayload::Uncompressed(payload)
916 }
917
918 fn add_field(
919 &mut self,
920 journal_file: &mut JournalFile<MmapMut>,
921 payload: &[u8],
922 ) -> Result<NonZeroU64> {
923 self.ensure_first_tag(journal_file)?;
924
925 if let Some(field_offset) = self.field_cache.get(payload) {
926 return Ok(field_offset);
927 }
928
929 let hash = journal_file.hash(payload);
930
931 match journal_file.find_field_offset(hash, payload)? {
932 Some(field_offset) => {
933 self.field_cache.insert(payload, field_offset);
934 Ok(field_offset)
935 }
936 None => {
937 let field_offset = self.append_offset;
940 let is_compact = Self::is_compact(journal_file);
941 Self::ensure_compact_object_fits(
942 is_compact,
943 field_offset,
944 std::mem::size_of::<FieldObjectHeader>() as u64 + payload.len() as u64,
945 )?;
946 let field_size = {
947 let mut field_guard =
948 journal_file.field_mut(field_offset, Some(payload.len() as u64))?;
949
950 field_guard.header.hash = hash;
951 field_guard.set_payload(payload);
952 field_guard.header.object_header.aligned_size()
953 };
954 self.hmac_put_object(journal_file, field_offset.get(), ObjectType::Field)?;
955 self.object_added(journal_file, field_offset, field_size)?;
956
957 journal_file.field_hash_table_set_tail_offset(hash, field_offset)?;
959 let depth = Self::current_field_hash_chain_depth(journal_file, hash)?;
960 let max_depth = journal_file
961 .journal_header_ref()
962 .field_hash_chain_depth
963 .max(depth);
964 journal_file.journal_header_mut().field_hash_chain_depth = max_depth;
965 journal_file.journal_header_mut().n_fields += 1;
966
967 self.field_cache.insert(payload, field_offset);
968
969 Ok(field_offset)
971 }
972 }
973 }
974}
975
976fn zstd_frame_with_content_size(frame: Vec<u8>, content_size: usize) -> Vec<u8> {
977 const ZSTD_MAGIC: [u8; 4] = [0x28, 0xb5, 0x2f, 0xfd];
978 const SINGLE_SEGMENT_FLAG: u8 = 1 << 5;
979 const CONTENT_CHECKSUM_FLAG: u8 = 1 << 2;
980
981 if frame.len() < 6 || frame[0..4] != ZSTD_MAGIC {
982 return frame;
983 }
984
985 let descriptor = frame[4];
986 let dictionary_id_flag = descriptor & 0x03;
987 let frame_content_size_flag = descriptor >> 6;
988 if dictionary_id_flag != 0
989 || frame_content_size_flag != 0
990 || (descriptor & SINGLE_SEGMENT_FLAG) != 0
991 {
992 return frame;
993 }
994
995 let (new_frame_content_size_flag, frame_content_size) = if content_size <= 255 {
996 (0u8, vec![content_size as u8])
997 } else if content_size <= 65_791 {
998 (1u8, ((content_size - 256) as u16).to_le_bytes().to_vec())
999 } else if u32::try_from(content_size).is_ok() {
1000 (2u8, (content_size as u32).to_le_bytes().to_vec())
1001 } else {
1002 (3u8, (content_size as u64).to_le_bytes().to_vec())
1003 };
1004
1005 let mut patched = Vec::with_capacity(frame.len() + frame_content_size.len() - 1);
1006 patched.extend_from_slice(&frame[..4]);
1007 patched.push(
1008 (new_frame_content_size_flag << 6)
1009 | SINGLE_SEGMENT_FLAG
1010 | (descriptor & CONTENT_CHECKSUM_FLAG),
1011 );
1012 patched.extend_from_slice(&frame_content_size);
1013 patched.extend_from_slice(&frame[6..]);
1014 patched
1015}
1016
1017fn xz_compress(payload: &[u8]) -> std::io::Result<Vec<u8>> {
1018 use lzma_rust2::{XzOptions, XzWriter};
1019 use std::io::Write;
1020
1021 let mut options = XzOptions::with_preset(0);
1022 options.set_check_sum_type(lzma_rust2::CheckType::None);
1023 let mut writer = XzWriter::new(Vec::new(), options)?;
1024 writer.write_all(payload)?;
1025 writer.finish()
1026}
1027
1028fn lz4_compress(payload: &[u8]) -> Vec<u8> {
1029 let compressed = lz4_flex::block::compress(payload);
1030 let mut out = Vec::with_capacity(8 + compressed.len());
1031 out.extend_from_slice(&(payload.len() as u64).to_le_bytes());
1032 out.extend_from_slice(&compressed);
1033 out
1034}
1035
1036#[cfg(test)]
1037#[path = "writer_tests.rs"]
1038mod tests;