Skip to main content

journal_log_writer/log/
mod.rs

1mod chain;
2use chain::OwnedChain;
3
4mod config;
5pub use config::{Config, LogIdentityMode, LogOpenMode, RetentionPolicy, RotationPolicy};
6
7mod helpers;
8mod startup;
9use helpers::*;
10use startup::{ActiveFile, RotationState, build_startup_state};
11
12use crate::{Result, WriterError};
13use itoa::Buffer as ItoaBuffer;
14pub use journal_common::EntryTimestamps;
15use journal_common::{Microseconds, RealtimeClock};
16use journal_core::error::JournalError;
17use journal_core::file::mmap::MmapMut;
18use journal_core::file::{
19    Compression, EntryField, EntryWriteOptions, FieldNamePolicy, JournalFile, JournalFileOptions,
20    JournalWriter, StructuredField,
21};
22use journal_registry::repository;
23use std::path::{Path, PathBuf};
24use std::sync::Arc;
25#[cfg(test)]
26use std::sync::atomic::{AtomicUsize, Ordering};
27
28const STACK_ENTRY_REF_LIMIT: usize = 128;
29const SOURCE_REALTIME_PREFIX: &[u8] = b"_SOURCE_REALTIME_TIMESTAMP=";
30const DERIVED_ROTATION_FRACTION: u64 = 20;
31const JOURNAL_FILE_SIZE_MIN: u64 = 512 * 1024;
32const PAGE_SIZE: u64 = 4096;
33const JOURNAL_COMPACT_SIZE_MAX: u64 = u32::MAX as u64;
34
35#[cfg(test)]
36static ARCHIVE_SYNC_CALLS: AtomicUsize = AtomicUsize::new(0);
37
38fn sync_archive_journal_file(
39    sync_on_archive: bool,
40    journal_file: &mut JournalFile<MmapMut>,
41) -> Result<()> {
42    if !sync_on_archive {
43        return Ok(());
44    }
45    #[cfg(test)]
46    ARCHIVE_SYNC_CALLS.fetch_add(1, Ordering::Relaxed);
47    journal_file.sync()?;
48    Ok(())
49}
50
51/// Tracks rotation state for size and count limits.
52pub struct Log {
53    poisoned: bool,
54    configured_dir: PathBuf,
55    chain: OwnedChain,
56    config: Config,
57    active_file: Option<ActiveFile>,
58    rotation_state: RotationState,
59    boot_id: uuid::Uuid,
60    seqnum_id: uuid::Uuid,
61    current_seqnum: u64,
62    clock: RealtimeClock,
63    last_monotonic_usec: u64,
64    lifecycle_observer: Option<Arc<dyn LogLifecycleObserver>>,
65    artifact_sizer: Option<Arc<dyn LogArtifactSizer>>,
66    retention_on_open_applied: bool,
67    boot_id_field: Vec<u8>,
68    source_realtime_field: Vec<u8>,
69}
70
71#[derive(Debug, Clone, Copy, PartialEq, Eq)]
72pub enum LogLifecycleReason {
73    Append,
74    EagerOpen,
75    Rotation,
76    Retention,
77}
78
79#[derive(Debug, Clone)]
80pub enum LogLifecycleEvent {
81    Created {
82        active: repository::File,
83        reason: LogLifecycleReason,
84    },
85    Rotated {
86        archived: repository::File,
87        active: repository::File,
88    },
89    RetainedDeleted {
90        files: Vec<repository::File>,
91    },
92}
93
94pub trait LogLifecycleObserver: Send + Sync {
95    fn on_event(&self, event: &LogLifecycleEvent);
96}
97
98pub trait LogArtifactSizer: Send + Sync {
99    fn journal_artifact_size(&self, journal_path: &Path) -> Result<u64>;
100}
101
102impl Log {
103    /// An uncertain writer must only be dropped/closed to release resources.
104    pub fn is_poisoned(&self) -> bool {
105        self.poisoned
106            || self
107                .active_file
108                .as_ref()
109                .is_some_and(|file| file.writer.is_poisoned())
110    }
111
112    fn ensure_healthy(&self) -> Result<()> {
113        if self.is_poisoned() {
114            Err(JournalError::WriterPoisoned.into())
115        } else {
116            Ok(())
117        }
118    }
119
120    fn duration_to_micros(duration: std::time::Duration) -> u64 {
121        duration.as_micros().try_into().unwrap_or(u64::MAX)
122    }
123
124    fn peek_entry_realtime(&self, timestamps: &EntryTimestamps) -> u64 {
125        let candidate = timestamps
126            .entry_realtime_usec
127            .unwrap_or_else(|| Microseconds::now().get());
128        let last_seen = self.clock.last_seen().get();
129        if candidate > last_seen {
130            candidate
131        } else {
132            last_seen.saturating_add(1)
133        }
134    }
135
136    fn should_rotate_for_realtime(&self, realtime: u64) -> bool {
137        let Some(active_file) = &self.active_file else {
138            return true;
139        };
140        if self.rotation_state.should_rotate() {
141            return true;
142        }
143        let Some(max_duration) = self.config.rotation_policy.duration_of_journal_file else {
144            return false;
145        };
146        let header = active_file.journal_file.journal_header_ref();
147        header.n_entries > 0
148            && header.head_entry_realtime > 0
149            && realtime.saturating_sub(header.head_entry_realtime)
150                >= Self::duration_to_micros(max_duration)
151    }
152
153    fn append_rotation_reason(&self) -> LogLifecycleReason {
154        if self.active_file.is_none() {
155            LogLifecycleReason::Append
156        } else {
157            LogLifecycleReason::Rotation
158        }
159    }
160
161    fn prepare_append_for_realtime(&mut self, entry_realtime: u64) -> Result<()> {
162        self.ensure_healthy()?;
163        self.apply_retention_on_open()?;
164        let opened_first_active = self.active_file.is_none();
165        if self.should_rotate_for_realtime(entry_realtime) {
166            self.rotate(entry_realtime, self.append_rotation_reason())?;
167            if opened_first_active {
168                self.retention_on_open_applied = true;
169            }
170        }
171        self.apply_retention_on_open()
172    }
173
174    fn raw_items_for_policy<'a>(&self, items: &'a [&'a [u8]]) -> Result<Option<Vec<&'a [u8]>>> {
175        if self.config.field_name_policy != FieldNamePolicy::JournalApp {
176            for item in items {
177                journal_core::file::writer::accept_entry_field(
178                    EntryField::raw(item),
179                    self.config.field_name_policy,
180                )?;
181            }
182            return Ok(None);
183        }
184        let filtered_items = filter_raw_items_for_journal_app(items)?;
185        if filtered_items.is_empty() {
186            return Err(WriterError::EmptyEntry);
187        }
188        Ok(Some(filtered_items))
189    }
190
191    fn structured_fields_for_policy<'a>(
192        &self,
193        fields: &'a [StructuredField<'a>],
194    ) -> Result<Option<Vec<StructuredField<'a>>>> {
195        if self.config.field_name_policy != FieldNamePolicy::JournalApp {
196            for field in fields {
197                journal_core::file::writer::accept_entry_field(
198                    EntryField::Structured(*field),
199                    self.config.field_name_policy,
200                )?;
201            }
202            return Ok(None);
203        }
204        let filtered_fields = filter_structured_fields_for_journal_app(fields);
205        if filtered_fields.is_empty() {
206            return Err(WriterError::EmptyEntry);
207        }
208        Ok(Some(filtered_fields))
209    }
210
211    fn apply_retention(&mut self, protected_file: Option<&repository::File>) -> Result<()> {
212        if let Some(sizer) = &self.artifact_sizer {
213            self.chain.refresh_retained_sizes(|file| {
214                sizer.journal_artifact_size(Path::new(file.path()))
215            })?;
216        }
217        let retention = self
218            .chain
219            .retain(&self.config.retention_policy, protected_file);
220        let deleted_files = retention.deleted_files;
221        if !deleted_files.is_empty()
222            && let Some(observer) = &self.lifecycle_observer
223        {
224            observer.on_event(&LogLifecycleEvent::RetainedDeleted {
225                files: deleted_files,
226            });
227        }
228        if let Some(error) = retention.error {
229            return Err(error);
230        }
231
232        Ok(())
233    }
234
235    fn apply_retention_on_open(&mut self) -> Result<()> {
236        if self.retention_on_open_applied || self.active_file.is_none() {
237            return Ok(());
238        }
239        self.enforce_retention()?;
240        self.retention_on_open_applied = true;
241        Ok(())
242    }
243
244    /// Captures both realtime and monotonic timestamps, similar to systemd's dual_timestamp_now().
245    ///
246    /// Returns (realtime_usec, monotonic_usec) where:
247    /// - realtime: microseconds since Unix epoch (CLOCK_REALTIME), monotonically increasing
248    /// - monotonic: microseconds since boot (CLOCK_MONOTONIC)
249    fn capture_dual_timestamp(
250        &mut self,
251        timestamp_override: Option<&EntryTimestamps>,
252    ) -> Result<(u64, u64)> {
253        let realtime = match timestamp_override.and_then(|ts| ts.entry_realtime_usec) {
254            Some(ts) => self.clock.observe(Microseconds::new(ts)).get(),
255            None => self.clock.now().get(),
256        };
257
258        let desired_monotonic = timestamp_override
259            .and_then(|ts| ts.entry_monotonic_usec)
260            .ok_or_else(|| {
261                WriterError::InvalidConfig("entry monotonic timestamp is required".to_string())
262            })?;
263
264        let monotonic = if desired_monotonic > self.last_monotonic_usec {
265            desired_monotonic
266        } else {
267            self.last_monotonic_usec.saturating_add(1)
268        };
269        self.last_monotonic_usec = monotonic;
270
271        Ok((realtime, monotonic))
272    }
273
274    fn require_entry_monotonic(timestamps: &EntryTimestamps) -> Result<()> {
275        if timestamps.entry_monotonic_usec.is_none() {
276            return Err(WriterError::InvalidConfig(
277                "entry monotonic timestamp is required".to_string(),
278            ));
279        }
280        Ok(())
281    }
282
283    /// Creates a new journal log.
284    pub fn new(path: &Path, config: Config) -> Result<Self> {
285        Self::new_inner(path, config, None, None)
286    }
287
288    pub fn new_with_lifecycle_observer(
289        path: &Path,
290        config: Config,
291        observer: Arc<dyn LogLifecycleObserver>,
292    ) -> Result<Self> {
293        Self::new_inner(path, config, Some(observer), None)
294    }
295
296    pub fn new_with_hooks(
297        path: &Path,
298        config: Config,
299        observer: Option<Arc<dyn LogLifecycleObserver>>,
300        artifact_sizer: Option<Arc<dyn LogArtifactSizer>>,
301    ) -> Result<Self> {
302        Self::new_inner(path, config, observer, artifact_sizer)
303    }
304
305    fn new_inner(
306        path: &Path,
307        config: Config,
308        lifecycle_observer: Option<Arc<dyn LogLifecycleObserver>>,
309        artifact_sizer: Option<Arc<dyn LogArtifactSizer>>,
310    ) -> Result<Self> {
311        let startup = build_startup_state(path, config)?;
312
313        let mut log = Log {
314            poisoned: false,
315            configured_dir: path.to_path_buf(),
316            chain: startup.chain,
317            config: startup.config,
318            active_file: startup.active_file,
319            rotation_state: startup.rotation_state,
320            boot_id: startup.boot_id,
321            seqnum_id: startup.seqnum_id,
322            current_seqnum: startup.current_seqnum,
323            clock: startup.clock,
324            last_monotonic_usec: startup.last_monotonic_usec,
325            lifecycle_observer,
326            artifact_sizer,
327            retention_on_open_applied: false,
328            boot_id_field: format!("_BOOT_ID={}", startup.boot_id.as_simple()).into_bytes(),
329            source_realtime_field: Vec::with_capacity(SOURCE_REALTIME_PREFIX.len() + 20),
330        };
331        if log.config.open_mode == LogOpenMode::Eager && log.active_file.is_none() {
332            let realtime = log.peek_entry_realtime(&EntryTimestamps::default());
333            log.rotate(realtime, LogLifecycleReason::EagerOpen)?;
334            log.retention_on_open_applied = true;
335        }
336        log.apply_retention_on_open()?;
337        Ok(log)
338    }
339
340    pub fn with_lifecycle_observer(mut self, observer: Arc<dyn LogLifecycleObserver>) -> Self {
341        self.lifecycle_observer = Some(observer);
342        self
343    }
344
345    pub fn with_artifact_sizer(mut self, sizer: Arc<dyn LogArtifactSizer>) -> Self {
346        self.artifact_sizer = Some(sizer);
347        self
348    }
349
350    /// Writes a journal entry.
351    ///
352    /// This compatibility method always returns an error under the strict
353    /// writer contract. Use [`Log::write_entry_with_timestamps`] and provide
354    /// an explicit entry monotonic timestamp.
355    #[deprecated(
356        since = "0.7.2",
357        note = "use write_entry_with_timestamps and provide an explicit entry monotonic timestamp"
358    )]
359    pub fn write_entry(
360        &mut self,
361        items: &[&[u8]],
362        source_realtime_usec: Option<u64>,
363    ) -> Result<()> {
364        self.write_entry_with_timestamps(
365            items,
366            EntryTimestamps {
367                source_realtime_usec,
368                ..EntryTimestamps::default()
369            },
370        )
371    }
372
373    /// Writes a journal entry with optional source and entry timestamp overrides.
374    ///
375    /// Overrides are safe by construction:
376    /// - entry realtime is clamped to strict monotonic progression (`last + 1us` floor)
377    /// - entry monotonic is also clamped to strict monotonic progression (`last + 1us` floor)
378    pub fn write_entry_with_timestamps(
379        &mut self,
380        items: &[&[u8]],
381        timestamps: EntryTimestamps,
382    ) -> Result<()> {
383        if items.is_empty() {
384            return Err(WriterError::EmptyEntry);
385        }
386        Self::require_entry_monotonic(&timestamps)?;
387
388        let entry_realtime = self.peek_entry_realtime(&timestamps);
389        self.ensure_healthy()?;
390        let filtered_items = self.raw_items_for_policy(items)?;
391        self.prepare_append_for_realtime(entry_realtime)?;
392        let write_items = filtered_items.as_deref().unwrap_or(items);
393
394        let (realtime, monotonic) = self.capture_dual_timestamp(Some(&timestamps))?;
395        self.write_raw_entry_fields(
396            write_items,
397            timestamps.source_realtime_usec,
398            realtime,
399            monotonic,
400            self.low_level_entry_options(EntryWriteOptions::default()),
401        )?;
402
403        let active_file = self.active_file.as_ref().unwrap();
404        self.rotation_state.update(&active_file.writer);
405        self.current_seqnum += 1;
406
407        Ok(())
408    }
409
410    /// Writes a journal entry from structured field names and binary-safe values.
411    ///
412    /// This is the preferred path when the producer already has split field
413    /// names and values. If `source_realtime_usec` is provided, a
414    /// `_SOURCE_REALTIME_TIMESTAMP` field is added.
415    ///
416    /// This compatibility method always returns an error under the strict
417    /// writer contract. Use [`Log::write_fields_with_timestamps`] and provide
418    /// an explicit entry monotonic timestamp.
419    #[deprecated(
420        since = "0.7.2",
421        note = "use write_fields_with_timestamps and provide an explicit entry monotonic timestamp"
422    )]
423    pub fn write_fields(
424        &mut self,
425        fields: &[StructuredField<'_>],
426        source_realtime_usec: Option<u64>,
427    ) -> Result<()> {
428        self.write_fields_with_timestamps(
429            fields,
430            EntryTimestamps {
431                source_realtime_usec,
432                ..EntryTimestamps::default()
433            },
434        )
435    }
436
437    /// Writes structured fields with optional source and entry timestamp overrides.
438    ///
439    /// Entry monotonic timestamp is required. Entry realtime and monotonic
440    /// overrides use the same clamping rules as
441    /// [`Log::write_entry_with_timestamps`].
442    pub fn write_fields_with_timestamps(
443        &mut self,
444        fields: &[StructuredField<'_>],
445        timestamps: EntryTimestamps,
446    ) -> Result<()> {
447        self.write_fields_with_options(fields, timestamps, EntryWriteOptions::default())
448    }
449
450    /// Writes structured fields with explicit low-level entry write options.
451    ///
452    /// Use this only when the caller can satisfy any invariants required by the
453    /// selected [`EntryWriteOptions`], especially no duplicate full `KEY=value`
454    /// payloads when `trusted_unique_payloads` is enabled.
455    pub fn write_fields_with_options(
456        &mut self,
457        fields: &[StructuredField<'_>],
458        timestamps: EntryTimestamps,
459        options: EntryWriteOptions,
460    ) -> Result<()> {
461        if fields.is_empty() {
462            return Err(WriterError::EmptyEntry);
463        }
464        Self::require_entry_monotonic(&timestamps)?;
465
466        let entry_realtime = self.peek_entry_realtime(&timestamps);
467        self.ensure_healthy()?;
468        let filtered_fields = self.structured_fields_for_policy(fields)?;
469        self.prepare_append_for_realtime(entry_realtime)?;
470        let write_fields = filtered_fields.as_deref().unwrap_or(fields);
471
472        let (realtime, monotonic) = self.capture_dual_timestamp(Some(&timestamps))?;
473        self.write_structured_entry_fields(
474            write_fields,
475            timestamps.source_realtime_usec,
476            realtime,
477            monotonic,
478            self.low_level_entry_options(options),
479        )?;
480
481        let active_file = self.active_file.as_ref().unwrap();
482        self.rotation_state.update(&active_file.writer);
483        self.current_seqnum += 1;
484
485        Ok(())
486    }
487
488    fn write_raw_entry_fields(
489        &mut self,
490        items: &[&[u8]],
491        source_realtime_usec: Option<u64>,
492        realtime: u64,
493        monotonic: u64,
494        options: EntryWriteOptions,
495    ) -> Result<()> {
496        let source_field = if let Some(timestamp_usec) = source_realtime_usec {
497            self.prepare_source_realtime_field(timestamp_usec);
498            Some(self.source_realtime_field.as_slice())
499        } else {
500            None
501        };
502
503        let total_items = items.len() + 1 + usize::from(source_field.is_some());
504        if total_items <= STACK_ENTRY_REF_LIMIT {
505            let mut refs = [EntryField::raw(&[]); STACK_ENTRY_REF_LIMIT];
506            let mut len = 0usize;
507            refs[len] = EntryField::raw(self.boot_id_field.as_slice());
508            len += 1;
509            if let Some(source_field) = source_field {
510                refs[len] = EntryField::raw(source_field);
511                len += 1;
512            }
513            for item in items {
514                refs[len] = EntryField::raw(item);
515                len += 1;
516            }
517            self.active_file.as_mut().unwrap().write_entry_fields(
518                refs[..len].iter().copied(),
519                realtime,
520                monotonic,
521                options,
522            )?;
523        } else {
524            let mut refs = Vec::with_capacity(total_items);
525            refs.push(EntryField::raw(self.boot_id_field.as_slice()));
526            if let Some(source_field) = source_field {
527                refs.push(EntryField::raw(source_field));
528            }
529            refs.extend(items.iter().copied().map(EntryField::raw));
530            self.active_file.as_mut().unwrap().write_entry_fields(
531                refs.iter().copied(),
532                realtime,
533                monotonic,
534                options,
535            )?;
536        }
537
538        Ok(())
539    }
540
541    fn write_structured_entry_fields(
542        &mut self,
543        fields: &[StructuredField<'_>],
544        source_realtime_usec: Option<u64>,
545        realtime: u64,
546        monotonic: u64,
547        options: EntryWriteOptions,
548    ) -> Result<()> {
549        let source_field = if let Some(timestamp_usec) = source_realtime_usec {
550            self.prepare_source_realtime_field(timestamp_usec);
551            Some(self.source_realtime_field.as_slice())
552        } else {
553            None
554        };
555
556        let total_items = fields.len() + 1 + usize::from(source_field.is_some());
557        if total_items <= STACK_ENTRY_REF_LIMIT {
558            let mut refs = [EntryField::raw(&[]); STACK_ENTRY_REF_LIMIT];
559            let mut len = 0usize;
560            refs[len] = EntryField::raw(self.boot_id_field.as_slice());
561            len += 1;
562            if let Some(source_field) = source_field {
563                refs[len] = EntryField::raw(source_field);
564                len += 1;
565            }
566            for field in fields {
567                refs[len] = EntryField::Structured(*field);
568                len += 1;
569            }
570            self.active_file.as_mut().unwrap().write_entry_fields(
571                refs[..len].iter().copied(),
572                realtime,
573                monotonic,
574                options,
575            )?;
576        } else {
577            let mut refs = Vec::with_capacity(total_items);
578            refs.push(EntryField::raw(self.boot_id_field.as_slice()));
579            if let Some(source_field) = source_field {
580                refs.push(EntryField::raw(source_field));
581            }
582            refs.extend(fields.iter().copied().map(EntryField::Structured));
583            self.active_file.as_mut().unwrap().write_entry_fields(
584                refs.iter().copied(),
585                realtime,
586                monotonic,
587                options,
588            )?;
589        }
590
591        Ok(())
592    }
593
594    fn low_level_entry_options(&self, options: EntryWriteOptions) -> EntryWriteOptions {
595        options.field_name_policy(log_writer_field_name_policy(self.config.field_name_policy))
596    }
597
598    fn prepare_source_realtime_field(&mut self, timestamp_usec: u64) {
599        self.source_realtime_field.clear();
600        self.source_realtime_field
601            .extend_from_slice(SOURCE_REALTIME_PREFIX);
602        let mut buffer = ItoaBuffer::new();
603        self.source_realtime_field
604            .extend_from_slice(buffer.format(timestamp_usec).as_bytes());
605    }
606
607    /// Syncs all written data to disk, ensuring durability.
608    ///
609    /// This should be called after writing a batch of log entries to ensure
610    /// they are persisted to disk before acknowledging the request.
611    pub fn sync(&mut self) -> Result<()> {
612        self.ensure_healthy()?;
613        if let Some(active_file) = &mut self.active_file {
614            self.poisoned = true;
615            active_file.journal_file.sync()?;
616            self.poisoned = false;
617        }
618        Ok(())
619    }
620
621    /// Archives and closes the active file.
622    ///
623    /// In strict systemd naming mode this renames `<source>.journal` to the
624    /// chain filename before retention, matching the explicit close behavior of
625    /// the other SDK implementations. `Drop` remains best-effort for callers
626    /// that do not explicitly close.
627    pub fn close(self) -> Result<()> {
628        self.close_impl(true)
629    }
630
631    /// Archives and closes the active file without applying retention.
632    ///
633    /// Use before reopening with a changed policy. This consumes the writer
634    /// and preserves [`Self::close`]'s archive and durability behavior; the
635    /// caller owns subsequent retention enforcement.
636    pub fn close_without_retention(self) -> Result<()> {
637        self.close_impl(false)
638    }
639
640    fn close_impl(mut self, enforce_retention: bool) -> Result<()> {
641        use journal_core::file::JournalState;
642
643        if self.is_poisoned() {
644            self.active_file.take();
645            return Err(JournalError::WriterPoisoned.into());
646        }
647        let Some(mut active_file) = self.active_file.take() else {
648            return Ok(());
649        };
650
651        let n_entries = active_file.journal_file.journal_header_ref().n_entries;
652        if self.config.strict_systemd_naming && n_entries == 0 {
653            match std::fs::remove_file(active_file.repository_file.path()) {
654                Ok(()) => {}
655                Err(err) if err.kind() == std::io::ErrorKind::NotFound => {}
656                Err(err) => return Err(err.into()),
657            }
658            self.chain.remove_tracked_file(&active_file.repository_file);
659            return Ok(());
660        }
661
662        self.chain.update_file_size(
663            &active_file.repository_file,
664            active_file.current_file_size(),
665        );
666        active_file.journal_file.journal_header_mut().state = JournalState::Archived as u8;
667        sync_archive_journal_file(self.config.sync_on_archive, &mut active_file.journal_file)?;
668
669        let protected_file = if self.config.strict_systemd_naming {
670            let header = active_file.journal_file.journal_header_ref();
671            self.chain.archive_file(
672                &active_file.repository_file,
673                uuid::Uuid::from_bytes(header.seqnum_id),
674                header.head_entry_seqnum,
675                header.head_entry_realtime,
676            )?
677        } else {
678            active_file.repository_file.clone()
679        };
680
681        if enforce_retention {
682            self.apply_retention(Some(&protected_file))?;
683        }
684
685        Ok(())
686    }
687
688    pub fn active_file(&self) -> Option<&repository::File> {
689        self.active_file
690            .as_ref()
691            .map(|active_file| &active_file.repository_file)
692    }
693
694    pub fn active_path(&self) -> Option<&Path> {
695        self.active_file
696            .as_ref()
697            .map(|active_file| Path::new(active_file.repository_file.path()))
698    }
699
700    pub fn configured_directory(&self) -> &Path {
701        &self.configured_dir
702    }
703
704    pub fn journal_directory(&self) -> &Path {
705        &self.chain.path
706    }
707
708    pub fn machine_id(&self) -> uuid::Uuid {
709        self.chain.machine_id
710    }
711
712    pub fn boot_id(&self) -> uuid::Uuid {
713        self.boot_id
714    }
715
716    pub fn source(&self) -> &journal_registry::Source {
717        &self.chain.source
718    }
719
720    /// Applies the configured retention policy without requiring a rotation or
721    /// close. The current active file is counted in retention envelopes and is
722    /// protected from deletion.
723    pub fn enforce_retention(&mut self) -> Result<()> {
724        self.ensure_healthy()?;
725        let protected_file = if let Some(active_file) = &self.active_file {
726            self.chain.update_file_size(
727                &active_file.repository_file,
728                active_file.current_file_size(),
729            );
730            Some(active_file.repository_file.clone())
731        } else {
732            None
733        };
734        self.apply_retention(protected_file.as_ref())
735    }
736
737    fn update_active_file_size(&mut self) {
738        if let Some(active_file) = &self.active_file {
739            self.chain.update_file_size(
740                &active_file.repository_file,
741                active_file.current_file_size(),
742            );
743        }
744    }
745
746    fn prepare_initial_rotation(&mut self) -> Result<()> {
747        self.update_active_file_size();
748        if self.active_file.is_none() && self.config.strict_systemd_naming {
749            self.chain
750                .archive_existing_active_file(self.config.sync_on_archive)?;
751        }
752        Ok(())
753    }
754
755    fn archive_rotated_file(&mut self, old_file: &ActiveFile) -> Result<repository::File> {
756        if !self.config.strict_systemd_naming {
757            return Ok(old_file.repository_file.clone());
758        }
759        let old_header = old_file.journal_file.journal_header_ref();
760        self.chain.archive_file(
761            &old_file.repository_file,
762            uuid::Uuid::from_bytes(old_header.seqnum_id),
763            old_header.head_entry_seqnum,
764            old_header.head_entry_realtime,
765        )
766    }
767
768    fn rotate_existing_active_file(
769        &mut self,
770        mut old_file: ActiveFile,
771        max_file_size: Option<u64>,
772        head_realtime: u64,
773    ) -> Result<(ActiveFile, LogLifecycleEvent)> {
774        use journal_core::file::JournalState;
775
776        old_file.journal_file.journal_header_mut().state = JournalState::Archived as u8;
777        sync_archive_journal_file(self.config.sync_on_archive, &mut old_file.journal_file)?;
778        let archived = self.archive_rotated_file(&old_file)?;
779        let new_file = old_file.rotate(
780            &mut self.chain,
781            max_file_size,
782            head_realtime,
783            self.config.compression,
784            self.config.compression_threshold,
785            self.config.strict_systemd_naming,
786            self.config.live_publish_every_entries,
787            self.config.file_mode,
788        )?;
789        let active = new_file.repository_file.clone();
790        Ok((new_file, LogLifecycleEvent::Rotated { archived, active }))
791    }
792
793    fn create_initial_active_file(
794        &mut self,
795        max_file_size: Option<u64>,
796        head_realtime: u64,
797        reason: LogLifecycleReason,
798    ) -> Result<(ActiveFile, LogLifecycleEvent)> {
799        let new_file = ActiveFile::create(
800            &mut self.chain,
801            self.seqnum_id,
802            self.boot_id,
803            self.current_seqnum + 1,
804            max_file_size,
805            head_realtime,
806            self.config.compression,
807            self.config.compression_threshold,
808            self.config.compact,
809            self.config.strict_systemd_naming,
810            self.config.live_publish_every_entries,
811            self.config.file_mode,
812        )?;
813        let active = new_file.repository_file.clone();
814        Ok((new_file, LogLifecycleEvent::Created { active, reason }))
815    }
816
817    fn emit_lifecycle_event(&self, event: &LogLifecycleEvent) {
818        if let Some(observer) = &self.lifecycle_observer {
819            observer.on_event(event);
820        }
821    }
822
823    fn protected_active_file(&self) -> Option<repository::File> {
824        self.active_file
825            .as_ref()
826            .map(|active_file| active_file.repository_file.clone())
827    }
828
829    #[tracing::instrument(skip_all, fields(active_file))]
830    fn rotate(&mut self, head_realtime: u64, reason: LogLifecycleReason) -> Result<()> {
831        self.ensure_healthy()?;
832        self.poisoned = true;
833        self.prepare_initial_rotation()?;
834        let max_file_size = self.config.rotation_policy.size_of_journal_file;
835        let (new_file, lifecycle_event) = if let Some(old_file) = self.active_file.take() {
836            self.rotate_existing_active_file(old_file, max_file_size, head_realtime)?
837        } else {
838            self.create_initial_active_file(max_file_size, head_realtime, reason)?
839        };
840
841        tracing::Span::current().record("new_file", new_file.repository_file.path());
842
843        self.active_file = Some(new_file);
844        self.poisoned = false;
845        self.rotation_state.reset();
846        self.update_active_file_size();
847        self.emit_lifecycle_event(&lifecycle_event);
848
849        // Retention runs after the post-rotation current file is known, so the
850        // tracked current file counts in the envelope and is never deleted.
851        let protected_file = self.protected_active_file();
852        self.apply_retention(protected_file.as_ref())?;
853
854        Ok(())
855    }
856
857    /// Writes a journal entry from a serializable value.
858    ///
859    /// This method serializes the value to JSON, flattens it, and writes it to the journal.
860    /// The flattened structure converts nested JSON into KEY=VALUE pairs suitable for journal entries.
861    ///
862    /// # Example
863    ///
864    /// ```no_run
865    /// use serde::Serialize;
866    /// use journal_log_writer::{Config, EntryTimestamps, Log, RotationPolicy, RetentionPolicy};
867    /// use journal_registry::Origin;
868    /// use std::path::Path;
869    ///
870    /// #[derive(Serialize)]
871    /// struct LogEntry {
872    ///     message: String,
873    ///     level: String,
874    ///     user: User,
875    /// }
876    ///
877    /// #[derive(Serialize)]
878    /// struct User {
879    ///     id: u64,
880    ///     name: String,
881    /// }
882    ///
883    /// # fn main() -> Result<(), Box<dyn std::error::Error>> {
884    /// let origin = Origin {
885    ///     machine_id: Some("00112233445566778899aabbccddeeff".parse()?),
886    ///     namespace: None,
887    ///     source: journal_registry::Source::System,
888    /// };
889    /// let config = Config::new(origin, RotationPolicy::default(), RetentionPolicy::default())
890    ///     .with_boot_id("ffeeddccbbaa99887766554433221100".parse()?);
891    /// let mut log = Log::new(Path::new("/tmp/test-journal"), config)?;
892    ///
893    /// let entry = LogEntry {
894    ///     message: "User logged in".to_string(),
895    ///     level: "INFO".to_string(),
896    ///     user: User {
897    ///         id: 42,
898    ///         name: "alice".to_string(),
899    ///     },
900    /// };
901    ///
902    /// // This will write fields like:
903    /// // MESSAGE=User logged in
904    /// // LEVEL=INFO
905    /// // USER_ID=42
906    /// // USER_NAME=alice
907    /// let timestamps = EntryTimestamps::default()
908    ///     .with_entry_realtime_usec(1_700_000_000_000_000)
909    ///     .with_entry_monotonic_usec(1);
910    /// log.write_structured_with_timestamps(&entry, timestamps)?;
911    /// # Ok(())
912    /// # }
913    /// ```
914    #[cfg(feature = "serde-api")]
915    #[deprecated(
916        since = "0.7.2",
917        note = "use write_structured_with_timestamps and provide an explicit entry monotonic timestamp"
918    )]
919    pub fn write_structured<T: serde::Serialize>(&mut self, value: &T) -> Result<()> {
920        self.write_structured_with_timestamps(value, EntryTimestamps::default())
921    }
922
923    /// Writes a journal entry from a serializable value with explicit timestamps.
924    #[cfg(feature = "serde-api")]
925    pub fn write_structured_with_timestamps<T: serde::Serialize>(
926        &mut self,
927        value: &T,
928        timestamps: EntryTimestamps,
929    ) -> Result<()> {
930        // Serialize to JSON value
931        let json_value = serde_json::to_value(value).map_err(|e| {
932            WriterError::Serialization(format!("failed to serialize to JSON: {}", e))
933        })?;
934
935        // Flatten the JSON structure - requires a JSON object (Map)
936        let flattened = if let serde_json::Value::Object(map) = json_value {
937            flatten_json_map(&map)
938        } else {
939            // If not an object, return error
940            return Err(WriterError::Serialization(
941                "value must be a JSON object, not a primitive or array".to_string(),
942            ));
943        };
944
945        // Convert to journal field format (KEY=VALUE)
946        let mut fields: Vec<Vec<u8>> = Vec::with_capacity(flattened.len());
947
948        for (key, value) in flattened.iter() {
949            // Convert key to uppercase and replace dots with underscores
950            // (journal convention)
951            let journal_key = key.to_uppercase().replace('.', "_");
952
953            // Format as KEY=VALUE
954            let field = match value {
955                serde_json::Value::String(s) => {
956                    format!("{}={}", journal_key, s)
957                }
958                serde_json::Value::Number(n) => {
959                    format!("{}={}", journal_key, n)
960                }
961                serde_json::Value::Bool(b) => {
962                    format!("{}={}", journal_key, if *b { "true" } else { "false" })
963                }
964                serde_json::Value::Null => {
965                    format!("{}=", journal_key)
966                }
967                // Arrays and objects should be flattened already, but just in case
968                _ => {
969                    format!("{}={}", journal_key, value)
970                }
971            };
972
973            fields.push(field.into_bytes());
974        }
975
976        // Convert Vec<Vec<u8>> to Vec<&[u8]> for write_entry
977        let field_refs: Vec<&[u8]> = fields.iter().map(|f| f.as_slice()).collect();
978
979        self.write_entry_with_timestamps(&field_refs, timestamps)
980    }
981}
982
983#[cfg(all(test, feature = "serde-api"))]
984mod serde_api_tests;
985
986impl Drop for Log {
987    fn drop(&mut self) {
988        use journal_core::file::JournalState;
989
990        if self.is_poisoned() {
991            return;
992        }
993        if let Some(ref mut active_file) = self.active_file {
994            // Keep the active path stable on close so file-backed readers that
995            // already follow system.journal can finish. The next writer startup
996            // archives this stale active file before creating a fresh one.
997            active_file.journal_file.journal_header_mut().state = JournalState::Archived as u8;
998
999            // Best/Last-effort sync just to be on the cautious side.
1000            let _ = sync_archive_journal_file(
1001                self.config.sync_on_archive,
1002                &mut active_file.journal_file,
1003            );
1004        }
1005    }
1006}
1007
1008#[cfg(test)]
1009mod tests {
1010    use super::*;
1011    use journal_registry::Origin;
1012    use std::sync::Mutex;
1013    use tempfile::TempDir;
1014
1015    pub(super) static ARCHIVE_SYNC_TEST_LOCK: Mutex<()> = Mutex::new(());
1016
1017    fn test_uuid(seed: u8) -> uuid::Uuid {
1018        uuid::Uuid::from_bytes([seed; 16])
1019    }
1020
1021    fn test_config() -> Config {
1022        Config::new(
1023            Origin {
1024                machine_id: Some(test_uuid(1)),
1025                namespace: None,
1026                source: journal_registry::Source::System,
1027            },
1028            RotationPolicy::default().with_number_of_entries(1),
1029            RetentionPolicy::default(),
1030        )
1031        .with_boot_id(test_uuid(2))
1032    }
1033
1034    fn write_test_entry(log: &mut Log, message: &[u8], realtime: u64) {
1035        log.write_entry_with_timestamps(
1036            &[message],
1037            EntryTimestamps::default()
1038                .with_entry_realtime_usec(realtime)
1039                .with_entry_monotonic_usec(realtime),
1040        )
1041        .expect("write entry");
1042    }
1043
1044    fn journal_paths(dir: &TempDir) -> Vec<PathBuf> {
1045        let journal_dir = dir.path().join(test_uuid(1).as_simple().to_string());
1046        let mut paths: Vec<_> = std::fs::read_dir(journal_dir)
1047            .expect("read journal dir")
1048            .map(|entry| entry.expect("read entry").path())
1049            .filter(|path| path.extension().is_some_and(|ext| ext == "journal"))
1050            .collect();
1051        paths.sort();
1052        paths
1053    }
1054
1055    fn create_online_chain_active(dir: &TempDir) -> PathBuf {
1056        let mut log = Log::new(
1057            dir.path(),
1058            test_config()
1059                .with_rotation_policy(RotationPolicy::default().with_number_of_entries(10)),
1060        )
1061        .expect("create log");
1062        write_test_entry(&mut log, b"MESSAGE=stale-online", 1);
1063        log.sync().expect("sync stale online active");
1064        let active_path = log.active_path().expect("active path").to_path_buf();
1065        let active_file = log.active_file.take().expect("active file");
1066        drop(active_file);
1067        drop(log);
1068        active_path
1069    }
1070
1071    #[test]
1072    fn default_rotation_syncs_archived_file_on_caller_path() {
1073        let _guard = ARCHIVE_SYNC_TEST_LOCK.lock().expect("lock sync test");
1074        ARCHIVE_SYNC_CALLS.store(0, Ordering::Relaxed);
1075        let dir = tempfile::tempdir().expect("create temp dir");
1076        let mut log = Log::new(dir.path(), test_config()).expect("create log");
1077
1078        write_test_entry(&mut log, b"MESSAGE=first", 1);
1079        write_test_entry(&mut log, b"MESSAGE=second", 2);
1080
1081        assert_eq!(
1082            ARCHIVE_SYNC_CALLS.load(Ordering::Relaxed),
1083            1,
1084            "default rotation should sync the outgoing archived file"
1085        );
1086    }
1087
1088    #[test]
1089    fn sync_on_archive_false_skips_archive_sync_and_keeps_files_readable() {
1090        let _guard = ARCHIVE_SYNC_TEST_LOCK.lock().expect("lock sync test");
1091        ARCHIVE_SYNC_CALLS.store(0, Ordering::Relaxed);
1092        let dir = tempfile::tempdir().expect("create temp dir");
1093        let mut log =
1094            Log::new(dir.path(), test_config().with_sync_on_archive(false)).expect("create log");
1095
1096        write_test_entry(&mut log, b"MESSAGE=first", 1);
1097        write_test_entry(&mut log, b"MESSAGE=second", 2);
1098        log.close().expect("close log");
1099
1100        assert_eq!(
1101            ARCHIVE_SYNC_CALLS.load(Ordering::Relaxed),
1102            0,
1103            "opt-out must not sync archived files on the caller path"
1104        );
1105
1106        let paths = journal_paths(&dir);
1107        assert_eq!(
1108            paths.len(),
1109            2,
1110            "rotation plus close should leave two journals"
1111        );
1112        for path in paths {
1113            let file = JournalFile::<journal_core::file::Mmap>::open_path(&path, 32 * 1024 * 1024)
1114                .expect("open archived journal");
1115            assert_eq!(file.journal_header_ref().n_entries, 1);
1116        }
1117    }
1118
1119    #[test]
1120    fn sync_on_archive_false_skips_drop_archive_sync() {
1121        let _guard = ARCHIVE_SYNC_TEST_LOCK.lock().expect("lock sync test");
1122        ARCHIVE_SYNC_CALLS.store(0, Ordering::Relaxed);
1123        let dir = tempfile::tempdir().expect("create temp dir");
1124
1125        {
1126            let mut log = Log::new(dir.path(), test_config().with_sync_on_archive(false))
1127                .expect("create log");
1128            write_test_entry(&mut log, b"MESSAGE=drop-opt-out", 1);
1129        }
1130
1131        assert_eq!(
1132            ARCHIVE_SYNC_CALLS.load(Ordering::Relaxed),
1133            0,
1134            "opt-out must skip best-effort Drop archive sync"
1135        );
1136    }
1137
1138    #[test]
1139    fn close_without_retention_preserves_archive_sync_policy() {
1140        let _guard = ARCHIVE_SYNC_TEST_LOCK.lock().expect("lock sync test");
1141        for strict in [false, true] {
1142            for sync in [false, true] {
1143                ARCHIVE_SYNC_CALLS.store(0, Ordering::Relaxed);
1144                let dir = tempfile::tempdir().expect("create temp dir");
1145                let mut log = Log::new(
1146                    dir.path(),
1147                    test_config()
1148                        .with_strict_systemd_naming(strict)
1149                        .with_sync_on_archive(sync),
1150                )
1151                .expect("create log");
1152                write_test_entry(&mut log, b"MESSAGE=close-without-retention", 1);
1153                log.close_without_retention().expect("close log");
1154                assert_eq!(
1155                    ARCHIVE_SYNC_CALLS.load(Ordering::Relaxed),
1156                    usize::from(sync)
1157                );
1158                let paths = journal_paths(&dir);
1159                assert_eq!(paths.len(), 1);
1160                let file =
1161                    JournalFile::<journal_core::file::Mmap>::open_path(&paths[0], 32 * 1024 * 1024)
1162                        .expect("open archive");
1163                assert_eq!(file.journal_header_ref().n_entries, 1);
1164                assert_eq!(
1165                    file.journal_header_ref().state,
1166                    journal_core::file::JournalState::Archived as u8
1167                );
1168            }
1169        }
1170    }
1171
1172    #[test]
1173    fn strict_startup_sync_on_archive_policy_applies_to_online_chain_active() {
1174        let _guard = ARCHIVE_SYNC_TEST_LOCK.lock().expect("lock sync test");
1175
1176        let default_dir = tempfile::tempdir().expect("create default temp dir");
1177        let default_active = create_online_chain_active(&default_dir);
1178        ARCHIVE_SYNC_CALLS.store(0, Ordering::Relaxed);
1179        let _default_log = Log::new(
1180            default_dir.path(),
1181            test_config()
1182                .with_rotation_policy(RotationPolicy::default().with_number_of_entries(10))
1183                .with_strict_systemd_naming(true),
1184        )
1185        .expect("open strict log");
1186        assert_eq!(
1187            ARCHIVE_SYNC_CALLS.load(Ordering::Relaxed),
1188            1,
1189            "default strict startup should sync the stale online chain file"
1190        );
1191        let file =
1192            JournalFile::<journal_core::file::Mmap>::open_path(&default_active, 32 * 1024 * 1024)
1193                .expect("open default startup archived journal");
1194        assert_eq!(
1195            file.journal_header_ref().state,
1196            journal_core::file::JournalState::Archived as u8
1197        );
1198
1199        let opt_out_dir = tempfile::tempdir().expect("create opt-out temp dir");
1200        let opt_out_active = create_online_chain_active(&opt_out_dir);
1201        ARCHIVE_SYNC_CALLS.store(0, Ordering::Relaxed);
1202        let _opt_out_log = Log::new(
1203            opt_out_dir.path(),
1204            test_config()
1205                .with_rotation_policy(RotationPolicy::default().with_number_of_entries(10))
1206                .with_strict_systemd_naming(true)
1207                .with_sync_on_archive(false),
1208        )
1209        .expect("open strict opt-out log");
1210        assert_eq!(
1211            ARCHIVE_SYNC_CALLS.load(Ordering::Relaxed),
1212            0,
1213            "opt-out strict startup must not sync the stale online chain file"
1214        );
1215        let file =
1216            JournalFile::<journal_core::file::Mmap>::open_path(&opt_out_active, 32 * 1024 * 1024)
1217                .expect("open opt-out startup archived journal");
1218        assert_eq!(
1219            file.journal_header_ref().state,
1220            journal_core::file::JournalState::Archived as u8
1221        );
1222    }
1223}
1224
1225#[cfg(test)]
1226mod poison_tests {
1227    use super::*;
1228
1229    #[test]
1230    fn poisoned_first_append_close_and_drop_preserve_active_file() {
1231        let _guard = super::tests::ARCHIVE_SYNC_TEST_LOCK.lock().unwrap();
1232        for compact in [false, true] {
1233            for explicit_close in [false, true] {
1234                let dir = tempfile::tempdir().unwrap();
1235                let config = Config::new(
1236                    journal_registry::Origin {
1237                        machine_id: Some(uuid::Uuid::from_bytes([1; 16])),
1238                        namespace: None,
1239                        source: journal_registry::Source::System,
1240                    },
1241                    RotationPolicy::default(),
1242                    RetentionPolicy::default(),
1243                )
1244                .with_boot_id(uuid::Uuid::from_bytes([2; 16]))
1245                .with_strict_systemd_naming(true)
1246                .with_compact(compact);
1247                let mut log = Log::new(dir.path(), config).unwrap();
1248                log.prepare_append_for_realtime(1_000_000).unwrap();
1249                let path = log.active_path().unwrap().to_path_buf();
1250                let active = log.active_file.as_mut().unwrap();
1251                // A real partial mutation: first DATA is linked, then the next
1252                // streamed field fails validation before ENTRY publication.
1253                assert!(
1254                    active
1255                        .write_entry_fields(
1256                            [
1257                                EntryField::raw(b"MESSAGE=partial"),
1258                                EntryField::raw(b"invalid")
1259                            ],
1260                            1_000_000,
1261                            1,
1262                            EntryWriteOptions::default()
1263                        )
1264                        .is_err()
1265                );
1266                assert_eq!(active.journal_file.journal_header_ref().n_entries, 0);
1267                assert!(log.is_poisoned());
1268                let before = std::fs::read(&path).unwrap();
1269                assert!(log.sync().is_err());
1270                assert!(
1271                    log.write_entry_with_timestamps(
1272                        &[b"MESSAGE=retry"],
1273                        EntryTimestamps::default()
1274                            .with_entry_realtime_usec(2_000_000)
1275                            .with_entry_monotonic_usec(2)
1276                    )
1277                    .is_err()
1278                );
1279                if explicit_close {
1280                    assert!(log.close().is_err());
1281                } else {
1282                    drop(log);
1283                }
1284                assert_eq!(std::fs::read(&path).unwrap(), before);
1285                assert_eq!(
1286                    std::fs::read_dir(path.parent().unwrap()).unwrap().count(),
1287                    1
1288                );
1289            }
1290        }
1291    }
1292
1293    #[test]
1294    fn invalid_input_keeps_log_reusable_before_rotation_or_mutation() {
1295        let _guard = super::tests::ARCHIVE_SYNC_TEST_LOCK.lock().unwrap();
1296        let dir = tempfile::tempdir().unwrap();
1297        let config = Config::new(
1298            journal_registry::Origin {
1299                machine_id: Some(uuid::Uuid::from_bytes([1; 16])),
1300                namespace: None,
1301                source: journal_registry::Source::System,
1302            },
1303            RotationPolicy::default(),
1304            RetentionPolicy::default(),
1305        )
1306        .with_boot_id(uuid::Uuid::from_bytes([2; 16]));
1307        let mut log = Log::new(dir.path(), config).unwrap();
1308        let timestamps = EntryTimestamps::default()
1309            .with_entry_realtime_usec(1_000_000)
1310            .with_entry_monotonic_usec(1);
1311        assert!(
1312            log.write_entry_with_timestamps(&[b"MESSAGE=valid", b"invalid"], timestamps)
1313                .is_err()
1314        );
1315        assert!(!log.is_poisoned());
1316        log.write_entry_with_timestamps(&[b"MESSAGE=valid"], timestamps)
1317            .unwrap();
1318        log.close().unwrap();
1319    }
1320}