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
51pub 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 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 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 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 #[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 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(×tamps)?;
387
388 let entry_realtime = self.peek_entry_realtime(×tamps);
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(×tamps))?;
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 #[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 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 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(×tamps)?;
465
466 let entry_realtime = self.peek_entry_realtime(×tamps);
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(×tamps))?;
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 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 pub fn close(self) -> Result<()> {
628 self.close_impl(true)
629 }
630
631 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 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 let protected_file = self.protected_active_file();
852 self.apply_retention(protected_file.as_ref())?;
853
854 Ok(())
855 }
856
857 #[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 #[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 let json_value = serde_json::to_value(value).map_err(|e| {
932 WriterError::Serialization(format!("failed to serialize to JSON: {}", e))
933 })?;
934
935 let flattened = if let serde_json::Value::Object(map) = json_value {
937 flatten_json_map(&map)
938 } else {
939 return Err(WriterError::Serialization(
941 "value must be a JSON object, not a primitive or array".to_string(),
942 ));
943 };
944
945 let mut fields: Vec<Vec<u8>> = Vec::with_capacity(flattened.len());
947
948 for (key, value) in flattened.iter() {
949 let journal_key = key.to_uppercase().replace('.', "_");
952
953 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 _ => {
969 format!("{}={}", journal_key, value)
970 }
971 };
972
973 fields.push(field.into_bytes());
974 }
975
976 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 active_file.journal_file.journal_header_mut().state = JournalState::Archived as u8;
998
999 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 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}