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
174#[doc(hidden)]
175pub fn accept_entry_field(field: EntryField<'_>, policy: FieldNamePolicy) -> Result<bool> {
176    let Some(field_name) = field.field_name() else {
177        return Err(JournalError::InvalidField);
178    };
179    let valid = match policy {
180        FieldNamePolicy::Raw => is_raw_field_name_valid(field_name),
181        FieldNamePolicy::Journald => is_journal_field_name_valid(field_name, true),
182        FieldNamePolicy::JournalApp => is_journal_field_name_valid(field_name, false),
183    };
184    if valid {
185        return Ok(true);
186    }
187    if matches!(policy, FieldNamePolicy::JournalApp) {
188        return Ok(false);
189    }
190    Err(JournalError::InvalidField)
191}
192
193#[derive(Debug)]
194struct FieldCache {
195    entries: FxHashMap<Box<[u8]>, NonZeroU64>,
196}
197
198impl FieldCache {
199    fn new() -> Self {
200        Self {
201            entries: FxHashMap::default(),
202        }
203    }
204
205    fn get(&self, payload: &[u8]) -> Option<NonZeroU64> {
206        self.entries.get(payload).copied()
207    }
208
209    fn insert(&mut self, payload: &[u8], offset: NonZeroU64) {
210        if payload.len() > FIELD_CACHE_MAX_PAYLOAD_LEN {
211            return;
212        }
213
214        if self.entries.len() >= FIELD_CACHE_MAX_ENTRIES && self.entries.get(payload).is_none() {
215            self.entries.clear();
216        }
217
218        self.entries
219            .insert(payload.to_vec().into_boxed_slice(), offset);
220    }
221
222    #[cfg(test)]
223    fn len(&self) -> usize {
224        self.entries.len()
225    }
226}
227
228enum StoredDataPayload<'a> {
229    Uncompressed(PayloadParts<'a>),
230    Compressed(Vec<u8>, u8),
231}
232
233impl StoredDataPayload<'_> {
234    fn len(&self) -> usize {
235        match self {
236            Self::Uncompressed(payload) => payload.len(),
237            Self::Compressed(payload, _) => payload.len(),
238        }
239    }
240
241    fn object_flags(&self) -> u8 {
242        match self {
243            Self::Uncompressed(_) => 0,
244            Self::Compressed(_, flags) => *flags,
245        }
246    }
247
248    fn copy_to_data_object(&self, data: &mut DataObject<&mut [u8]>) {
249        match self {
250            Self::Uncompressed(payload) => match &mut data.payload {
251                DataPayloadType::Regular(dst) => payload.copy_to_slice(dst),
252                DataPayloadType::Compact { payload: dst, .. } => payload.copy_to_slice(dst),
253            },
254            Self::Compressed(payload, _) => data.set_payload(payload),
255        }
256    }
257}
258
259pub struct JournalWriter {
260    poisoned: bool,
261    #[cfg(test)]
262    fail_append_stage: u8,
263    pub(super) tail_object_offset: NonZeroU64,
264    pub(super) append_offset: NonZeroU64,
265    next_seqnum: u64,
266    num_written_objects: u64,
267    pub(super) first_tag_written: bool,
268    pub(super) entry_items: Vec<EntryItem>,
269    field_cache: FieldCache,
270    first_entry_monotonic: Option<u64>,
271    boot_id: uuid::Uuid,
272    compression: Compression,
273    compress_threshold: usize,
274    live_publish_every_entries: u64,
275    entries_since_live_publication: u64,
276    pub(super) seal: Option<crate::seal::SealState>,
277}
278
279impl JournalWriter {
280    /// True after a mutating failure. Drop resources without publishing a clean state.
281    pub fn is_poisoned(&self) -> bool {
282        self.poisoned
283    }
284
285    #[cfg(test)]
286    fn fail_append_at(&self, stage: u8) -> Result<()> {
287        if self.fail_append_stage == stage {
288            return Err(JournalError::Io(std::io::Error::other(
289                "injected append failure",
290            )));
291        }
292        Ok(())
293    }
294
295    /// Get current file size in bytes
296    pub fn current_file_size(&self) -> u64 {
297        self.append_offset.get()
298    }
299
300    /// Get the monotonic timestamp of the first entry written to this file
301    pub fn first_entry_monotonic(&self) -> Option<u64> {
302        self.first_entry_monotonic
303    }
304
305    /// Get the next sequence number that will be written
306    pub fn next_seqnum(&self) -> u64 {
307        self.next_seqnum
308    }
309
310    /// Get the boot ID for this writer
311    pub fn boot_id(&self) -> uuid::Uuid {
312        self.boot_id
313    }
314
315    /// Sets how often the writer explicitly publishes live-reader visibility.
316    ///
317    /// `1` is the default and matches systemd-style publication after every
318    /// appended entry. `0` disables this explicit publication; closed-file
319    /// verification and reads after sync/close are unchanged, but stock
320    /// follow-reader visibility while the writer is active is not guaranteed.
321    /// Values greater than `1` publish after every N appended entries.
322    pub fn set_live_publish_every_entries(&mut self, entries: u64) {
323        self.live_publish_every_entries = entries;
324        self.entries_since_live_publication = 0;
325    }
326
327    /// Returns the configured live-reader publication cadence.
328    ///
329    /// See [`JournalWriter::set_live_publish_every_entries`] for the meaning of
330    /// `0`, `1`, and larger values.
331    pub fn live_publish_every_entries(&self) -> u64 {
332        self.live_publish_every_entries
333    }
334
335    pub fn new(
336        journal_file: &mut JournalFile<MmapMut>,
337        next_seqnum: u64,
338        boot_id: uuid::Uuid,
339    ) -> Result<Self> {
340        let compression = match journal_file.journal_header_ref() {
341            header if header.has_incompatible_flag(HeaderIncompatibleFlags::CompressedZstd) => {
342                Compression::Zstd
343            }
344            header if header.has_incompatible_flag(HeaderIncompatibleFlags::CompressedXz) => {
345                Compression::Xz
346            }
347            header if header.has_incompatible_flag(HeaderIncompatibleFlags::CompressedLz4) => {
348                Compression::Lz4
349            }
350            _ => Compression::None,
351        };
352
353        Self::new_with_compression(
354            journal_file,
355            next_seqnum,
356            boot_id,
357            compression,
358            DEFAULT_COMPRESS_THRESHOLD,
359        )
360    }
361
362    pub fn new_with_compression(
363        journal_file: &mut JournalFile<MmapMut>,
364        next_seqnum: u64,
365        boot_id: uuid::Uuid,
366        compression: Compression,
367        compress_threshold: usize,
368    ) -> Result<Self> {
369        let current_header_size = std::mem::size_of::<JournalHeader>() as u64;
370        let header = journal_file.journal_header_ref();
371        if !header.has_incompatible_flag(HeaderIncompatibleFlags::KeyedHash)
372            || header.header_size < current_header_size
373        {
374            return Err(JournalError::UnsupportedJournalFile);
375        }
376
377        journal_file.validate_committed_arena()?;
378        header.validate_empty_entry_metadata()?;
379
380        let append_offset = {
381            let header = journal_file.journal_header_ref();
382
383            let Some(tail_object_offset) = header.tail_object_offset else {
384                return Err(JournalError::InvalidMagicNumber);
385            };
386
387            let tail_object = journal_file.object_header_ref(tail_object_offset)?;
388
389            tail_object_offset
390                .checked_add(tail_object.aligned_size())
391                .ok_or(JournalError::ObjectExceedsFileBounds)?
392        };
393
394        let seal = journal_file
395            .seal_options
396            .as_ref()
397            .map(|opts| crate::seal::SealState::new(opts))
398            .transpose()?;
399
400        let mut writer = Self {
401            poisoned: false,
402            #[cfg(test)]
403            fail_append_stage: 0,
404            tail_object_offset: journal_file
405                .journal_header_ref()
406                .tail_object_offset
407                .unwrap(),
408            append_offset,
409            next_seqnum,
410            num_written_objects: 0,
411            first_tag_written: false,
412            entry_items: Vec::with_capacity(128),
413            field_cache: FieldCache::new(),
414            first_entry_monotonic: None,
415            boot_id,
416            compression,
417            compress_threshold: normalize_compress_threshold(compress_threshold),
418            live_publish_every_entries: 1,
419            entries_since_live_publication: 0,
420            seal,
421        };
422
423        if writer.seal.is_some() && journal_file.journal_header_ref().n_tags == 0 {
424            writer.ensure_first_tag(journal_file)?;
425            {
426                let header = journal_file.journal_header_mut();
427                header.n_objects += writer.num_written_objects;
428                header.tail_object_offset = Some(writer.tail_object_offset);
429            }
430            writer.num_written_objects = 0;
431        }
432
433        Ok(writer)
434    }
435
436    /// Creates a successor writer for a new journal file
437    pub fn create_successor(&self, journal_file: &mut JournalFile<MmapMut>) -> Result<Self> {
438        if self.poisoned {
439            return Err(JournalError::WriterPoisoned);
440        }
441        Self::new_with_compression(
442            journal_file,
443            self.next_seqnum,
444            self.boot_id,
445            self.compression,
446            self.compress_threshold,
447        )
448    }
449
450    pub fn add_entry(
451        &mut self,
452        journal_file: &mut JournalFile<MmapMut>,
453        items: &[&[u8]],
454        realtime: u64,
455        monotonic: u64,
456    ) -> Result<()> {
457        self.add_entry_fields_with_options(
458            journal_file,
459            items.iter().copied().map(EntryField::raw),
460            realtime,
461            monotonic,
462            EntryWriteOptions::default(),
463        )
464    }
465
466    pub fn add_entry_structured(
467        &mut self,
468        journal_file: &mut JournalFile<MmapMut>,
469        fields: &[StructuredField<'_>],
470        realtime: u64,
471        monotonic: u64,
472    ) -> Result<()> {
473        self.add_entry_fields_with_options(
474            journal_file,
475            fields.iter().copied().map(EntryField::Structured),
476            realtime,
477            monotonic,
478            EntryWriteOptions::default(),
479        )
480    }
481
482    pub fn add_entry_structured_with_options(
483        &mut self,
484        journal_file: &mut JournalFile<MmapMut>,
485        fields: &[StructuredField<'_>],
486        realtime: u64,
487        monotonic: u64,
488        options: EntryWriteOptions,
489    ) -> Result<()> {
490        self.add_entry_fields_with_options(
491            journal_file,
492            fields.iter().copied().map(EntryField::Structured),
493            realtime,
494            monotonic,
495            options,
496        )
497    }
498
499    pub fn add_entry_fields<'a>(
500        &mut self,
501        journal_file: &mut JournalFile<MmapMut>,
502        fields: impl IntoIterator<Item = EntryField<'a>>,
503        realtime: u64,
504        monotonic: u64,
505    ) -> Result<()> {
506        self.add_entry_fields_with_options(
507            journal_file,
508            fields,
509            realtime,
510            monotonic,
511            EntryWriteOptions::default(),
512        )
513    }
514
515    pub fn add_entry_fields_with_options<'a>(
516        &mut self,
517        journal_file: &mut JournalFile<MmapMut>,
518        fields: impl IntoIterator<Item = EntryField<'a>>,
519        realtime: u64,
520        monotonic: u64,
521        options: EntryWriteOptions,
522    ) -> Result<()> {
523        if self.poisoned {
524            return Err(JournalError::WriterPoisoned);
525        }
526        self.ensure_keyed_append(journal_file)?;
527        let entry_seqnum = self.entry_seqnum_for_options(options)?;
528        let entry_boot_id = options.boot_id.unwrap_or(self.boot_id);
529        let monotonic = self.clamp_same_boot_monotonic(journal_file, entry_boot_id, monotonic)?;
530        if let Some(seal) = &self.seal {
531            seal.need_evolve(realtime)?;
532        }
533        let xor_hash = self.prepare_entry_items(journal_file, fields, realtime, options)?;
534        #[cfg(test)]
535        self.fail_append_at(1)?;
536        let entry_offset = self.write_entry_object(
537            journal_file,
538            entry_seqnum,
539            entry_boot_id,
540            realtime,
541            monotonic,
542            xor_hash,
543        )?;
544        #[cfg(test)]
545        self.fail_append_at(2)?;
546        self.publish_entry_links(journal_file, entry_offset)?;
547        #[cfg(test)]
548        self.fail_append_at(3)?;
549        self.entry_added(
550            journal_file.journal_header_mut(),
551            entry_offset,
552            entry_seqnum,
553            entry_boot_id,
554            realtime,
555            monotonic,
556        );
557        #[cfg(test)]
558        self.fail_append_at(4)?;
559        self.publish_after_entry(journal_file)?;
560        self.poisoned = false;
561        Ok(())
562    }
563
564    fn ensure_keyed_append(&self, journal_file: &JournalFile<MmapMut>) -> Result<()> {
565        let header = journal_file.journal_header_ref();
566        if header.has_incompatible_flag(HeaderIncompatibleFlags::KeyedHash) {
567            return Ok(());
568        }
569        Err(JournalError::UnsupportedJournalFile)
570    }
571
572    fn entry_seqnum_for_options(&self, options: EntryWriteOptions) -> Result<u64> {
573        let entry_seqnum = options.seqnum.unwrap_or(self.next_seqnum);
574        if entry_seqnum == 0 || entry_seqnum == u64::MAX || entry_seqnum < self.next_seqnum {
575            return Err(JournalError::InvalidField);
576        }
577        Ok(entry_seqnum)
578    }
579
580    fn clamp_same_boot_monotonic(
581        &self,
582        journal_file: &JournalFile<MmapMut>,
583        entry_boot_id: uuid::Uuid,
584        monotonic: u64,
585    ) -> Result<u64> {
586        let header = journal_file.journal_header_ref();
587        if header.n_entries == 0
588            || header.tail_entry_boot_id != *entry_boot_id.as_bytes()
589            || monotonic > header.tail_entry_monotonic
590        {
591            return Ok(monotonic);
592        }
593        header
594            .tail_entry_monotonic
595            .checked_add(1)
596            .ok_or(JournalError::InvalidField)
597    }
598
599    fn prepare_entry_items<'a>(
600        &mut self,
601        journal_file: &mut JournalFile<MmapMut>,
602        fields: impl IntoIterator<Item = EntryField<'a>>,
603        realtime: u64,
604        options: EntryWriteOptions,
605    ) -> Result<u64> {
606        let mut xor_hash = 0;
607        self.entry_items.clear();
608        let mut publication_ready = false;
609        for field in fields {
610            if !accept_entry_field(field, options.field_name_policy)? {
611                continue;
612            }
613            self.ensure_entry_publication_ready(journal_file, realtime, &mut publication_ready)?;
614            xor_hash ^= self.add_entry_field_item(journal_file, field)?;
615        }
616        self.finish_entry_items(options.trusted_unique_payloads)?;
617        Ok(xor_hash)
618    }
619
620    fn ensure_entry_publication_ready(
621        &mut self,
622        journal_file: &mut JournalFile<MmapMut>,
623        realtime: u64,
624        publication_ready: &mut bool,
625    ) -> Result<()> {
626        if *publication_ready {
627            return Ok(());
628        }
629        // Set before any potentially partial I/O; only a fully published append clears it.
630        // Seal state is publication state too, even before the next file write.
631        if let Some(seal) = &self.seal {
632            if !self.first_tag_written || seal.need_evolve(realtime)? {
633                self.poisoned = true;
634            }
635        }
636        self.ensure_first_tag(journal_file)?;
637        self.maybe_append_tag(journal_file, realtime)?;
638        *publication_ready = true;
639        Ok(())
640    }
641
642    fn add_entry_field_item(
643        &mut self,
644        journal_file: &mut JournalFile<MmapMut>,
645        field: EntryField<'_>,
646    ) -> Result<u64> {
647        let entry_item = self.add_data(journal_file, field)?;
648        self.entry_items.push(entry_item);
649        Ok(jenkins_hash64_parts(field.payload_parts().iter()))
650    }
651
652    fn finish_entry_items(&mut self, trusted_unique_payloads: bool) -> Result<()> {
653        if self.entry_items.is_empty() {
654            return Err(JournalError::InvalidField);
655        }
656        if !self.entry_items_are_sorted() {
657            self.entry_items
658                .sort_unstable_by(|a, b| a.offset.cmp(&b.offset));
659        }
660        if !trusted_unique_payloads {
661            self.entry_items.dedup_by(|a, b| a.offset == b.offset);
662        } else {
663            for index in 1..self.entry_items.len() {
664                if self.entry_items[index - 1].offset == self.entry_items[index].offset {
665                    self.entry_items[index].link_state = None;
666                }
667            }
668        }
669        Ok(())
670    }
671
672    fn entry_items_are_sorted(&self) -> bool {
673        self.entry_items
674            .windows(2)
675            .all(|items| items[0].offset <= items[1].offset)
676    }
677
678    fn write_entry_object(
679        &mut self,
680        journal_file: &mut JournalFile<MmapMut>,
681        entry_seqnum: u64,
682        entry_boot_id: uuid::Uuid,
683        realtime: u64,
684        monotonic: u64,
685        xor_hash: u64,
686    ) -> Result<NonZeroU64> {
687        let entry_offset = self.append_offset;
688        let is_compact = Self::is_compact(journal_file);
689        let entry_payload_size = self.entry_items.len() as u64 * Self::entry_item_size(is_compact);
690        Self::ensure_compact_object_fits(
691            is_compact,
692            entry_offset,
693            std::mem::size_of::<EntryObjectHeader>() as u64 + entry_payload_size,
694        )?;
695        self.poisoned = true;
696        let entry_size = {
697            let size = Some(entry_payload_size);
698            let mut entry_guard = journal_file.entry_mut(entry_offset, size)?;
699
700            entry_guard.header.seqnum = entry_seqnum;
701            entry_guard.header.xor_hash = xor_hash;
702            entry_guard.header.boot_id = *entry_boot_id.as_bytes();
703            entry_guard.header.monotonic = monotonic;
704            entry_guard.header.realtime = realtime;
705
706            // set each entry item
707            for (index, entry_item) in self.entry_items.iter().enumerate() {
708                Self::ensure_compact_offset(is_compact, entry_item.offset)?;
709                let item_hash = (!is_compact).then_some(entry_item.hash);
710                entry_guard.items.set(index, entry_item.offset, item_hash);
711            }
712
713            entry_guard.header.object_header.aligned_size()
714        };
715        self.hmac_put_object(journal_file, entry_offset.get(), ObjectType::Entry)?;
716        self.object_added(journal_file, entry_offset, entry_size)?;
717        Ok(entry_offset)
718    }
719
720    fn publish_entry_links(
721        &mut self,
722        journal_file: &mut JournalFile<MmapMut>,
723        entry_offset: NonZeroU64,
724    ) -> Result<()> {
725        self.append_to_entry_array(journal_file, entry_offset)?;
726        for entry_item_index in 0..self.entry_items.len() {
727            self.link_data_to_entry(journal_file, entry_offset, entry_item_index)?;
728        }
729        Ok(())
730    }
731
732    fn publish_after_entry(&mut self, journal_file: &mut JournalFile<MmapMut>) -> Result<()> {
733        match self.live_publish_every_entries {
734            0 => Ok(()),
735            1 => journal_file.post_change(),
736            interval => {
737                self.entries_since_live_publication += 1;
738                if self.entries_since_live_publication >= interval {
739                    self.entries_since_live_publication = 0;
740                    journal_file.post_change()
741                } else {
742                    Ok(())
743                }
744            }
745        }
746    }
747
748    pub(super) fn object_added(
749        &mut self,
750        journal_file: &mut JournalFile<MmapMut>,
751        object_offset: NonZeroU64,
752        object_size: u64,
753    ) -> Result<()> {
754        self.tail_object_offset = object_offset;
755        self.append_offset = object_offset
756            .checked_add(object_size)
757            .ok_or(JournalError::ObjectExceedsFileBounds)?;
758        self.num_written_objects += 1;
759
760        let header = journal_file.journal_header_mut();
761        let old_size = header
762            .header_size
763            .checked_add(header.arena_size)
764            .ok_or(JournalError::ObjectExceedsFileBounds)?;
765        if self.append_offset.get() > old_size {
766            let new_size = round_up_to_file_size_increment(self.append_offset.get())?;
767            header.arena_size = new_size
768                .checked_sub(header.header_size)
769                .ok_or(JournalError::ObjectExceedsFileBounds)?;
770        }
771
772        Ok(())
773    }
774
775    fn entry_added(
776        &mut self,
777        header: &mut JournalHeader,
778        entry_offset: NonZeroU64,
779        entry_seqnum: u64,
780        entry_boot_id: uuid::Uuid,
781        realtime: u64,
782        monotonic: u64,
783    ) {
784        header.n_objects += self.num_written_objects;
785        header.tail_object_offset = Some(self.tail_object_offset);
786
787        if header.head_entry_seqnum == 0 {
788            header.head_entry_seqnum = entry_seqnum;
789        }
790        if header.head_entry_realtime == 0 {
791            header.head_entry_realtime = realtime;
792        }
793        if self.first_entry_monotonic.is_none() {
794            self.first_entry_monotonic = Some(monotonic);
795        }
796
797        header.tail_entry_seqnum = entry_seqnum;
798        header.tail_entry_realtime = realtime;
799        header.tail_entry_monotonic = monotonic;
800        header.tail_entry_boot_id = *entry_boot_id.as_bytes();
801        header.tail_entry_offset = entry_offset.get();
802        header.n_entries += 1;
803
804        self.next_seqnum = entry_seqnum + 1;
805        self.num_written_objects = 0;
806    }
807
808    fn add_data(
809        &mut self,
810        journal_file: &mut JournalFile<MmapMut>,
811        field: EntryField<'_>,
812    ) -> Result<EntryItem> {
813        let payload = field.payload_parts();
814        let field_name = field.field_name().ok_or(JournalError::InvalidField)?;
815        let hash = journal_file.hash_parts(payload);
816        if let Some((data_offset, link_state)) =
817            journal_file.find_data_with_link_state_parts(hash, payload)?
818        {
819            return Ok(Self::entry_item(data_offset, hash, link_state));
820        }
821        self.add_new_data(journal_file, payload, field_name, hash)
822    }
823
824    fn entry_item(offset: NonZeroU64, hash: u64, link_state: ResolvedDataLinkState) -> EntryItem {
825        EntryItem {
826            offset,
827            hash,
828            link_state: Some(link_state),
829        }
830    }
831
832    fn add_new_data<'a>(
833        &mut self,
834        journal_file: &mut JournalFile<MmapMut>,
835        payload: PayloadParts<'a>,
836        field_name: &'a [u8],
837        hash: u64,
838    ) -> Result<EntryItem> {
839        let data_offset = self.write_new_data_object(journal_file, payload, hash)?;
840        self.publish_new_data_object(journal_file, data_offset, hash)?;
841        self.link_data_to_field(journal_file, data_offset, field_name)?;
842        Ok(Self::entry_item(
843            data_offset,
844            hash,
845            ResolvedDataLinkState::empty(),
846        ))
847    }
848
849    fn write_new_data_object<'a>(
850        &mut self,
851        journal_file: &mut JournalFile<MmapMut>,
852        payload: PayloadParts<'a>,
853        hash: u64,
854    ) -> Result<NonZeroU64> {
855        let data_offset = self.append_offset;
856        let stored_payload = self.stored_data_payload(payload);
857        self.ensure_data_object_fits(journal_file, data_offset, stored_payload.len() as u64)?;
858        self.poisoned = true;
859        let data_size = {
860            let mut data_guard =
861                journal_file.data_mut(data_offset, Some(stored_payload.len() as u64))?;
862            data_guard.header.hash = hash;
863            stored_payload.copy_to_data_object(&mut data_guard);
864            data_guard.header.object_header.flags = stored_payload.object_flags();
865            data_guard.header.object_header.aligned_size()
866        };
867        self.hmac_put_object(journal_file, data_offset.get(), ObjectType::Data)?;
868        self.object_added(journal_file, data_offset, data_size)?;
869        Ok(data_offset)
870    }
871
872    fn ensure_data_object_fits(
873        &self,
874        journal_file: &JournalFile<MmapMut>,
875        data_offset: NonZeroU64,
876        payload_size: u64,
877    ) -> Result<()> {
878        let is_compact = Self::is_compact(journal_file);
879        Self::ensure_compact_object_fits(
880            is_compact,
881            data_offset,
882            Self::data_object_size(is_compact, payload_size),
883        )
884    }
885
886    fn publish_new_data_object(
887        &mut self,
888        journal_file: &mut JournalFile<MmapMut>,
889        data_offset: NonZeroU64,
890        hash: u64,
891    ) -> Result<()> {
892        journal_file.data_hash_table_set_tail_offset(hash, data_offset)?;
893        Self::update_data_hash_chain_depth(journal_file, hash)?;
894        journal_file.journal_header_mut().n_data += 1;
895        Ok(())
896    }
897
898    fn link_data_to_field(
899        &mut self,
900        journal_file: &mut JournalFile<MmapMut>,
901        data_offset: NonZeroU64,
902        field_name: &[u8],
903    ) -> Result<()> {
904        let field_offset = self.add_field(journal_file, field_name)?;
905        let head_data_offset = {
906            let field_guard = journal_file.field_ref(field_offset)?;
907            field_guard.header.head_data_offset
908        };
909        {
910            let mut data_guard = journal_file.data_mut(data_offset, None)?;
911            data_guard.header.next_field_offset = head_data_offset;
912        }
913        let mut field_guard = journal_file.field_mut(field_offset, None)?;
914        field_guard.header.head_data_offset = Some(data_offset);
915        Ok(())
916    }
917
918    fn stored_data_payload<'a>(&self, payload: PayloadParts<'a>) -> StoredDataPayload<'a> {
919        if payload.len() >= self.compress_threshold {
920            let full_payload;
921            let payload_bytes = if let Some(raw) = payload.as_single_slice() {
922                raw
923            } else {
924                // Structured payloads need a contiguous buffer only when compression is
925                // enabled and the payload is large enough to attempt compression.
926                full_payload = payload.to_vec();
927                full_payload.as_slice()
928            };
929            match self.compression {
930                Compression::Zstd => {
931                    let compressed = ruzstd::encoding::compress_to_vec(
932                        Cursor::new(payload_bytes),
933                        ruzstd::encoding::CompressionLevel::Fastest,
934                    );
935                    let compressed = zstd_frame_with_content_size(compressed, payload_bytes.len());
936                    if compressed.len() < payload_bytes.len() {
937                        return StoredDataPayload::Compressed(
938                            compressed,
939                            ObjectFlags::CompressedZstd as u8,
940                        );
941                    }
942                }
943                Compression::Xz => {
944                    if payload_bytes.len() >= 80 {
945                        if let Ok(compressed) = xz_compress(payload_bytes) {
946                            if compressed.len() < payload_bytes.len() {
947                                return StoredDataPayload::Compressed(
948                                    compressed,
949                                    ObjectFlags::CompressedXz as u8,
950                                );
951                            }
952                        }
953                    }
954                }
955                Compression::Lz4 => {
956                    if payload_bytes.len() >= 9 {
957                        let compressed = lz4_compress(payload_bytes);
958                        if compressed.len() < payload_bytes.len() {
959                            return StoredDataPayload::Compressed(
960                                compressed,
961                                ObjectFlags::CompressedLz4 as u8,
962                            );
963                        }
964                    }
965                }
966                Compression::None => {}
967            }
968        }
969
970        StoredDataPayload::Uncompressed(payload)
971    }
972
973    fn add_field(
974        &mut self,
975        journal_file: &mut JournalFile<MmapMut>,
976        payload: &[u8],
977    ) -> Result<NonZeroU64> {
978        self.ensure_first_tag(journal_file)?;
979
980        if let Some(field_offset) = self.field_cache.get(payload) {
981            return Ok(field_offset);
982        }
983
984        let hash = journal_file.hash(payload);
985
986        match journal_file.find_field_offset(hash, payload)? {
987            Some(field_offset) => {
988                self.field_cache.insert(payload, field_offset);
989                Ok(field_offset)
990            }
991            None => {
992                // We will have to write the new field object at the current
993                // tail offset
994                let field_offset = self.append_offset;
995                let is_compact = Self::is_compact(journal_file);
996                Self::ensure_compact_object_fits(
997                    is_compact,
998                    field_offset,
999                    std::mem::size_of::<FieldObjectHeader>() as u64 + payload.len() as u64,
1000                )?;
1001                let field_size = {
1002                    let mut field_guard =
1003                        journal_file.field_mut(field_offset, Some(payload.len() as u64))?;
1004
1005                    field_guard.header.hash = hash;
1006                    field_guard.set_payload(payload);
1007                    field_guard.header.object_header.aligned_size()
1008                };
1009                self.hmac_put_object(journal_file, field_offset.get(), ObjectType::Field)?;
1010                self.object_added(journal_file, field_offset, field_size)?;
1011
1012                // Update hash table
1013                journal_file.field_hash_table_set_tail_offset(hash, field_offset)?;
1014                let depth = Self::current_field_hash_chain_depth(journal_file, hash)?;
1015                let max_depth = journal_file
1016                    .journal_header_ref()
1017                    .field_hash_chain_depth
1018                    .max(depth);
1019                journal_file.journal_header_mut().field_hash_chain_depth = max_depth;
1020                journal_file.journal_header_mut().n_fields += 1;
1021
1022                self.field_cache.insert(payload, field_offset);
1023
1024                // Return the offset where we wrote the newly added data object
1025                Ok(field_offset)
1026            }
1027        }
1028    }
1029}
1030
1031fn zstd_frame_with_content_size(frame: Vec<u8>, content_size: usize) -> Vec<u8> {
1032    const ZSTD_MAGIC: [u8; 4] = [0x28, 0xb5, 0x2f, 0xfd];
1033    const SINGLE_SEGMENT_FLAG: u8 = 1 << 5;
1034    const CONTENT_CHECKSUM_FLAG: u8 = 1 << 2;
1035
1036    if frame.len() < 6 || frame[0..4] != ZSTD_MAGIC {
1037        return frame;
1038    }
1039
1040    let descriptor = frame[4];
1041    let dictionary_id_flag = descriptor & 0x03;
1042    let frame_content_size_flag = descriptor >> 6;
1043    if dictionary_id_flag != 0
1044        || frame_content_size_flag != 0
1045        || (descriptor & SINGLE_SEGMENT_FLAG) != 0
1046    {
1047        return frame;
1048    }
1049
1050    let (new_frame_content_size_flag, frame_content_size) = if content_size <= 255 {
1051        (0u8, vec![content_size as u8])
1052    } else if content_size <= 65_791 {
1053        (1u8, ((content_size - 256) as u16).to_le_bytes().to_vec())
1054    } else if u32::try_from(content_size).is_ok() {
1055        (2u8, (content_size as u32).to_le_bytes().to_vec())
1056    } else {
1057        (3u8, (content_size as u64).to_le_bytes().to_vec())
1058    };
1059
1060    let mut patched = Vec::with_capacity(frame.len() + frame_content_size.len() - 1);
1061    patched.extend_from_slice(&frame[..4]);
1062    patched.push(
1063        (new_frame_content_size_flag << 6)
1064            | SINGLE_SEGMENT_FLAG
1065            | (descriptor & CONTENT_CHECKSUM_FLAG),
1066    );
1067    patched.extend_from_slice(&frame_content_size);
1068    patched.extend_from_slice(&frame[6..]);
1069    patched
1070}
1071
1072fn xz_compress(payload: &[u8]) -> std::io::Result<Vec<u8>> {
1073    use lzma_rust2::{XzOptions, XzWriter};
1074    use std::io::Write;
1075
1076    let mut options = XzOptions::with_preset(0);
1077    options.set_check_sum_type(lzma_rust2::CheckType::None);
1078    let mut writer = XzWriter::new(Vec::new(), options)?;
1079    writer.write_all(payload)?;
1080    writer.finish()
1081}
1082
1083fn lz4_compress(payload: &[u8]) -> Vec<u8> {
1084    let compressed = lz4_flex::block::compress(payload);
1085    let mut out = Vec::with_capacity(8 + compressed.len());
1086    out.extend_from_slice(&(payload.len() as u64).to_le_bytes());
1087    out.extend_from_slice(&compressed);
1088    out
1089}
1090
1091#[cfg(test)]
1092#[path = "writer_tests.rs"]
1093mod tests;