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