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 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 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 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 #[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 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(×tamps)?;
355
356 let entry_realtime = self.peek_entry_realtime(×tamps);
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(×tamps))?;
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 #[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 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 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(×tamps)?;
432
433 let entry_realtime = self.peek_entry_realtime(×tamps);
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(×tamps))?;
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 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 pub fn close(self) -> Result<()> {
591 self.close_impl(true)
592 }
593
594 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 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 let protected_file = self.protected_active_file();
806 self.apply_retention(protected_file.as_ref())?;
807
808 Ok(())
809 }
810
811 #[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 #[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 let json_value = serde_json::to_value(value).map_err(|e| {
886 WriterError::Serialization(format!("failed to serialize to JSON: {}", e))
887 })?;
888
889 let flattened = if let serde_json::Value::Object(map) = json_value {
891 flatten_json_map(&map)
892 } else {
893 return Err(WriterError::Serialization(
895 "value must be a JSON object, not a primitive or array".to_string(),
896 ));
897 };
898
899 let mut fields: Vec<Vec<u8>> = Vec::with_capacity(flattened.len());
901
902 for (key, value) in flattened.iter() {
903 let journal_key = key.to_uppercase().replace('.', "_");
906
907 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 _ => {
923 format!("{}={}", journal_key, value)
924 }
925 };
926
927 fields.push(field.into_bytes());
928 }
929
930 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 active_file.journal_file.journal_header_mut().state = JournalState::Archived as u8;
949
950 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}