Skip to main content

journal_core/file/
writer.rs

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    /// Field name without the `=` separator.
36    pub name: &'a [u8],
37    /// Field value bytes. Values may contain NUL bytes and `=` bytes.
38    pub value: &'a [u8],
39}
40
41impl<'a> StructuredField<'a> {
42    /// Creates a structured journal field from a name and binary-safe value.
43    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    /// Full `KEY=value` payload, matching systemd's low-level writer shape.
51    Raw(&'a [u8]),
52    /// Split field name and value, avoiding `KEY=value` reconstruction for
53    /// already-structured producers.
54    Structured(StructuredField<'a>),
55}
56
57impl<'a> EntryField<'a> {
58    /// Creates a raw full-field entry item from a `KEY=value` byte payload.
59    pub fn raw(payload: &'a [u8]) -> Self {
60        Self::Raw(payload)
61    }
62
63    /// Creates a structured entry item from a name and binary-safe value.
64    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    /// Trusted journald-compatible field names. Protected `_...` names are
89    /// allowed.
90    #[default]
91    Journald,
92    /// Journal DATA structure capability only. Stock systemd tooling
93    /// compatibility is not guaranteed for names outside JOURNALD.
94    Raw,
95    /// Untrusted application input accepted by journald. Invalid or protected
96    /// caller fields are dropped.
97    JournalApp,
98}
99
100#[derive(Debug, Clone, Copy, Default)]
101pub struct EntryWriteOptions {
102    /// Skips duplicate DATA reference elimination for this ENTRY.
103    ///
104    /// Set this only when the caller guarantees that the entry contains no
105    /// duplicate full `KEY=value` payloads after field-name policy filtering.
106    /// Offset sorting by DATA object offset is always performed regardless of
107    /// this flag.
108    /// Misuse can write duplicate DATA offsets into one ENTRY object. Keep the
109    /// default `false` unless the producer owns and enforces that invariant.
110    pub trusted_unique_payloads: bool,
111    /// Field-name validation policy for caller-provided fields.
112    pub field_name_policy: FieldNamePolicy,
113    /// Optional low-level ENTRY seqnum override.
114    ///
115    /// This is for exact journal regeneration and must be monotonically
116    /// increasing relative to previously written entries. Leave unset for the
117    /// normal systemd-style auto-incrementing sequence.
118    pub seqnum: Option<u64>,
119    /// Optional low-level ENTRY boot ID override.
120    ///
121    /// This is for exact journal regeneration of multi-boot files. Leave unset
122    /// for the normal writer-wide boot ID.
123    pub boot_id: Option<uuid::Uuid>,
124}
125
126impl EntryWriteOptions {
127    /// Enables or disables the trusted unique-payload fast path.
128    ///
129    /// See [`EntryWriteOptions::trusted_unique_payloads`] for the caller
130    /// invariant required before enabling this option.
131    pub fn trusted_unique_payloads(mut self, enabled: bool) -> Self {
132        self.trusted_unique_payloads = enabled;
133        self
134    }
135
136    /// Selects the field-name validation policy for caller-provided fields.
137    pub fn field_name_policy(mut self, policy: FieldNamePolicy) -> Self {
138        self.field_name_policy = policy;
139        self
140    }
141
142    /// Uses a caller-provided ENTRY seqnum for this entry.
143    pub fn seqnum(mut self, seqnum: u64) -> Self {
144        self.seqnum = Some(seqnum);
145        self
146    }
147
148    /// Uses a caller-provided ENTRY boot ID for this entry.
149    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    /// Get current file size in bytes
277    pub fn current_file_size(&self) -> u64 {
278        self.append_offset.get()
279    }
280
281    /// Get the monotonic timestamp of the first entry written to this file
282    pub fn first_entry_monotonic(&self) -> Option<u64> {
283        self.first_entry_monotonic
284    }
285
286    /// Get the next sequence number that will be written
287    pub fn next_seqnum(&self) -> u64 {
288        self.next_seqnum
289    }
290
291    /// Get the boot ID for this writer
292    pub fn boot_id(&self) -> uuid::Uuid {
293        self.boot_id
294    }
295
296    /// Sets how often the writer explicitly publishes live-reader visibility.
297    ///
298    /// `1` is the default and matches systemd-style publication after every
299    /// appended entry. `0` disables this explicit publication; closed-file
300    /// verification and reads after sync/close are unchanged, but stock
301    /// follow-reader visibility while the writer is active is not guaranteed.
302    /// Values greater than `1` publish after every N appended entries.
303    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    /// Returns the configured live-reader publication cadence.
309    ///
310    /// See [`JournalWriter::set_live_publish_every_entries`] for the meaning of
311    /// `0`, `1`, and larger values.
312    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    /// Creates a successor writer for a new journal file
410    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            // set each entry item
653            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                // Structured payloads need a contiguous buffer only when compression is
870                // enabled and the payload is large enough to attempt compression.
871                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                // We will have to write the new field object at the current
938                // tail offset
939                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                // Update hash table
958                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                // Return the offset where we wrote the newly added data object
970                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;