Skip to main content

journal_core/file/
file_mut.rs

1use super::file::{
2    Compression, JOURNAL_COMPACT_SIZE_MAX, JournalFile, JournalFileOptions, OBJECT_ALIGNMENT,
3    map_hash_table, round_up_to_file_size_increment, validate_offset_alignment,
4};
5use super::mmap::{MemoryMap, MemoryMapMut, WindowManager, read_file_exact_at};
6use super::object::*;
7use crate::error::{JournalError, Result};
8use crate::file::guarded_cell::GuardedCell;
9use crate::file::value_guard::ValueGuard;
10use std::fs::{File, OpenOptions};
11use std::num::NonZeroU64;
12#[cfg(unix)]
13use std::os::unix::fs::OpenOptionsExt;
14use zerocopy::FromBytes;
15
16#[derive(Debug, Clone, Copy)]
17struct CreateLayout {
18    data_hash_table_size: usize,
19    field_hash_table_size: usize,
20    data_hash_table_offset: u64,
21    field_hash_table_offset: u64,
22    data_hash_table_object_offset: u64,
23    file_size: u64,
24}
25
26#[derive(Debug, Clone, Copy)]
27struct MutableObjectContext {
28    object_type: ObjectType,
29    is_compact: bool,
30    arena_end: u64,
31}
32
33impl JournalFile<super::mmap::MmapMut> {
34    pub fn open_for_append(file: &crate::repository::File, window_size: u64) -> Result<Self> {
35        debug_assert_eq!(window_size % OBJECT_ALIGNMENT, 0);
36
37        let fd = OpenOptions::new()
38            .read(true)
39            .write(true)
40            .open(file.path())?;
41
42        let header_size = std::mem::size_of::<JournalHeader>() as u64;
43        let file_size = fd.metadata()?.len();
44        if file_size < header_size {
45            return Err(JournalError::ObjectExceedsFileBounds);
46        }
47        let mut header_bytes = [0u8; std::mem::size_of::<JournalHeader>()];
48        read_file_exact_at(&fd, 0, &mut header_bytes)?;
49        let header = JournalHeader::read_from_prefix(&header_bytes).unwrap().0;
50        if header.signature != *b"LPKSHHRH" {
51            return Err(JournalError::InvalidMagicNumber);
52        }
53        if header.header_size < header_size {
54            return Err(JournalError::UnsupportedJournalFile);
55        }
56        if !header.has_incompatible_flag(HeaderIncompatibleFlags::KeyedHash) {
57            return Err(JournalError::UnsupportedJournalFile);
58        }
59
60        header.validate_empty_entry_metadata()?;
61        Self::validate_committed_arena_header(&header, file_size, |offset, bytes| {
62            read_file_exact_at(&fd, offset, bytes)
63        })?;
64        let header_map = super::mmap::MmapMut::create_checked(&fd, 0, header_size, file_size)?;
65
66        let data_hash_table_map = map_hash_table(
67            &fd,
68            header.header_size,
69            header.data_hash_table_offset,
70            header.data_hash_table_size,
71        )?;
72        let field_hash_table_map = map_hash_table(
73            &fd,
74            header.header_size,
75            header.field_hash_table_offset,
76            header.field_hash_table_size,
77        )?;
78
79        let window_manager =
80            GuardedCell::new(WindowManager::new_writer_owned(fd, window_size, 32)?);
81
82        let journal = JournalFile {
83            file: file.clone(),
84            header_map,
85            sanitized_header: None,
86            data_hash_table_map,
87            field_hash_table_map,
88            window_manager,
89            seal_options: None,
90        };
91        Ok(journal)
92    }
93}
94
95impl<M: MemoryMapMut> JournalFile<M> {
96    /// Syncs all file data to disk, ensuring all changes are persisted
97    ///
98    /// This performs a two-step sync process:
99    /// 1. Flushes memory-mapped regions to the file page cache (msync)
100    /// 2. Syncs the file page cache to physical disk (fdatasync)
101    pub fn sync(&mut self) -> Result<()> {
102        // Flush memory-mapped header to file page cache
103        self.header_map.flush()?;
104
105        // Sync file page cache to disk
106        let (logical_size, header_size) = {
107            let header = self.journal_header_ref();
108            (header.header_size + header.arena_size, header.header_size)
109        };
110        let header_bytes = self.header_map[..header_size as usize].to_vec();
111        let window_manager = self.window_manager.get_mut();
112        window_manager.sync(logical_size, &header_bytes)?;
113
114        Ok(())
115    }
116
117    /// Trigger a stock-reader-visible post-change notification after mmap append.
118    pub fn post_change(&mut self) -> Result<()> {
119        let logical_size = {
120            let header = self.journal_header_ref();
121            header.header_size + header.arena_size
122        };
123        self.window_manager.get_mut().post_change(logical_size)
124    }
125
126    /// Creates a successor journal file with optimized bucket sizes based on this file's utilization
127    pub fn create_successor(
128        &self,
129        file: &crate::repository::File,
130        max_file_size: Option<u64>,
131    ) -> Result<Self> {
132        self.create_successor_with_file_mode(file, max_file_size, self.current_file_mode())
133    }
134
135    pub fn create_successor_with_file_mode(
136        &self,
137        file: &crate::repository::File,
138        max_file_size: Option<u64>,
139        file_mode: u32,
140    ) -> Result<Self> {
141        let header = self.journal_header_ref();
142        let bucket_utilization = self.bucket_utilization();
143
144        let options = JournalFileOptions::new(
145            uuid::Uuid::from_bytes(header.machine_id),
146            uuid::Uuid::from_bytes(header.tail_entry_boot_id),
147            uuid::Uuid::from_bytes(header.seqnum_id),
148        )
149        .with_window_size(8 * 1024 * 1024)
150        .with_keyed_hash(header.has_incompatible_flag(HeaderIncompatibleFlags::KeyedHash))
151        .with_compact(header.has_incompatible_flag(HeaderIncompatibleFlags::Compact))
152        .with_file_mode(file_mode)
153        .with_optimized_buckets(bucket_utilization, max_file_size);
154
155        let options = if header.has_incompatible_flag(HeaderIncompatibleFlags::CompressedZstd) {
156            options.with_compression(Compression::Zstd)
157        } else if header.has_incompatible_flag(HeaderIncompatibleFlags::CompressedXz) {
158            options.with_compression(Compression::Xz)
159        } else if header.has_incompatible_flag(HeaderIncompatibleFlags::CompressedLz4) {
160            options.with_compression(Compression::Lz4)
161        } else {
162            options
163        };
164
165        Self::create(file, options)
166    }
167
168    pub fn create(file: &crate::repository::File, options: JournalFileOptions) -> Result<Self> {
169        let fd = Self::open_new_file(file, options.file_mode)?;
170        let layout = Self::create_layout(&options)?;
171        if options.compact && layout.file_size > JOURNAL_COMPACT_SIZE_MAX {
172            return Err(JournalError::ObjectExceedsFileBounds);
173        }
174        fd.set_len(layout.file_size)?;
175        let mut header = Self::create_header(&options, layout);
176        let data_hash_table_map = map_hash_table(
177            &fd,
178            header.header_size,
179            header.data_hash_table_offset,
180            header.data_hash_table_size,
181        )?;
182        let field_hash_table_map = map_hash_table(
183            &fd,
184            header.header_size,
185            header.field_hash_table_offset,
186            header.field_hash_table_size,
187        )?;
188        let header_map = Self::create_header_map(&fd, &mut header)?;
189        let window_manager = GuardedCell::new(WindowManager::new_writer_owned_with_strategy(
190            fd,
191            options.window_size,
192            32,
193            options.experimental_mmap_strategy,
194        )?);
195
196        let mut jf = JournalFile {
197            file: file.clone(),
198            header_map,
199            sanitized_header: None,
200            data_hash_table_map,
201            field_hash_table_map,
202            window_manager,
203            seal_options: options.seal.clone(),
204        };
205
206        jf.write_initial_hash_table_headers(header)?;
207        jf.sync()?;
208        Ok(jf)
209    }
210
211    fn current_file_mode(&self) -> u32 {
212        #[cfg(unix)]
213        {
214            use std::os::unix::fs::PermissionsExt;
215            if let Ok(metadata) = std::fs::metadata(self.file.path()) {
216                return metadata.permissions().mode() & 0o777;
217            }
218        }
219        super::file::DEFAULT_JOURNAL_FILE_MODE
220    }
221
222    fn open_new_file(file: &crate::repository::File, mode: u32) -> Result<File> {
223        let mut open_options = OpenOptions::new();
224        open_options
225            .create(true)
226            .truncate(true)
227            .read(true)
228            .write(true);
229        #[cfg(unix)]
230        open_options.mode(mode);
231        Ok(open_options.open(file.path())?)
232    }
233
234    fn create_layout(options: &JournalFileOptions) -> Result<CreateLayout> {
235        let data_hash_table_size =
236            options.data_hash_table_buckets * std::mem::size_of::<HashItem>();
237        let field_hash_table_size =
238            options.field_hash_table_buckets * std::mem::size_of::<HashItem>();
239        let field_hash_table_offset = std::mem::size_of::<JournalHeader>() as u64
240            + std::mem::size_of::<ObjectHeader>() as u64;
241        let data_hash_table_offset = field_hash_table_offset
242            + field_hash_table_size as u64
243            + std::mem::size_of::<ObjectHeader>() as u64;
244        let data_hash_table_object_offset =
245            data_hash_table_offset - std::mem::size_of::<ObjectHeader>() as u64;
246        let append_offset = data_hash_table_offset + data_hash_table_size as u64;
247        let file_size = round_up_to_file_size_increment(append_offset)?;
248        Ok(CreateLayout {
249            data_hash_table_size,
250            field_hash_table_size,
251            data_hash_table_offset,
252            field_hash_table_offset,
253            data_hash_table_object_offset,
254            file_size,
255        })
256    }
257
258    fn create_header(options: &JournalFileOptions, layout: CreateLayout) -> JournalHeader {
259        let mut header = JournalHeader::default();
260        header.signature = *b"LPKSHHRH";
261        header.compatible_flags = HeaderCompatibleFlags::TailEntryBootId as u32;
262        if options.enable_keyed_hash {
263            header.incompatible_flags |= HeaderIncompatibleFlags::KeyedHash as u32;
264        }
265        header.incompatible_flags |= options.compression.as_incompatible_flag();
266        if options.compact {
267            header.incompatible_flags |= HeaderIncompatibleFlags::Compact as u32;
268        }
269        if options.seal.is_some() {
270            header.compatible_flags |= HeaderCompatibleFlags::Sealed as u32;
271            header.compatible_flags |= HeaderCompatibleFlags::SealedContinuous as u32;
272        }
273        header.data_hash_table_offset = NonZeroU64::new(layout.data_hash_table_offset);
274        header.data_hash_table_size = NonZeroU64::new(layout.data_hash_table_size as u64);
275        header.field_hash_table_offset = NonZeroU64::new(layout.field_hash_table_offset);
276        header.field_hash_table_size = NonZeroU64::new(layout.field_hash_table_size as u64);
277        header.tail_object_offset = NonZeroU64::new(layout.data_hash_table_object_offset);
278        header.header_size = std::mem::size_of::<JournalHeader>() as u64;
279        header.n_objects = 2;
280        header.arena_size = layout.file_size - header.header_size;
281        header.machine_id = *options.machine_id.as_bytes();
282        header.file_id = *options.file_id.as_bytes();
283        header.seqnum_id = *options.seqnum_id.as_bytes();
284        header
285    }
286
287    fn create_header_map(fd: &File, header: &mut JournalHeader) -> Result<M> {
288        let header_size = std::mem::size_of::<JournalHeader>() as u64;
289        let mut header_map = M::create(fd, 0, header_size)?;
290        {
291            let header_mut = JournalHeader::mut_from_prefix(&mut header_map).unwrap().0;
292            *header_mut = *header;
293            header_mut.state = JournalState::Online as u8;
294            header.state = JournalState::Online as u8;
295        }
296        Ok(header_map)
297    }
298
299    fn write_initial_hash_table_headers(&mut self, header: JournalHeader) -> Result<()> {
300        self.write_hash_table_object_header(
301            header.data_hash_table_offset.unwrap(),
302            header.data_hash_table_size.unwrap(),
303            ObjectType::DataHashTable,
304        )?;
305        self.write_hash_table_object_header(
306            header.field_hash_table_offset.unwrap(),
307            header.field_hash_table_size.unwrap(),
308            ObjectType::FieldHashTable,
309        )
310    }
311
312    fn write_hash_table_object_header(
313        &self,
314        table_offset: NonZeroU64,
315        table_size: NonZeroU64,
316        object_type: ObjectType,
317    ) -> Result<()> {
318        let object_offset =
319            NonZeroU64::new(table_offset.get() - std::mem::size_of::<ObjectHeader>() as u64)
320                .unwrap();
321        let object_header = self.object_header_mut(object_offset)?;
322        object_header.type_ = object_type as u8;
323        object_header.size = table_size.get() + std::mem::size_of::<ObjectHeader>() as u64;
324        Ok(())
325    }
326
327    pub fn journal_header_mut(&mut self) -> &mut JournalHeader {
328        JournalHeader::mut_from_prefix(&mut self.header_map)
329            .unwrap()
330            .0
331    }
332
333    pub fn data_hash_table_mut(&mut self) -> Option<DataHashTable<&mut [u8]>> {
334        self.data_hash_table_map
335            .as_mut()
336            .and_then(|m| DataHashTable::<&mut [u8]>::from_data_mut(m, false))
337    }
338
339    pub fn field_hash_table_mut(&mut self) -> Option<FieldHashTable<&mut [u8]>> {
340        self.field_hash_table_map
341            .as_mut()
342            .and_then(|m| FieldHashTable::<&mut [u8]>::from_data_mut(m, false))
343    }
344
345    #[allow(clippy::mut_from_ref)]
346    fn object_header_mut(&self, offset: NonZeroU64) -> Result<&mut ObjectHeader> {
347        validate_offset_alignment(offset)?;
348        let size_needed = std::mem::size_of::<ObjectHeader>() as u64;
349        let window_manager = self.window_manager.borrow_mut_checked()?;
350        let header_slice = window_manager.get_slice_mut(offset.get(), size_needed)?;
351        ObjectHeader::mut_from_bytes(header_slice).map_err(|_| JournalError::ZerocopyFailure)
352    }
353
354    fn journal_object_mut<'a, T>(
355        &'a self,
356        type_: ObjectType,
357        offset: NonZeroU64,
358        size: Option<u64>,
359    ) -> Result<ValueGuard<'a, T>>
360    where
361        T: JournalObjectMut<&'a mut [u8]>,
362    {
363        let context = self.mutable_object_context(type_, offset)?;
364        self.window_manager.with_guarded(offset, |wm| {
365            let size_needed = Self::mutable_object_size(wm, context, offset, size)?;
366            let data = wm.get_slice_mut(offset.get(), size_needed)?;
367            let value =
368                T::from_data_mut(data, context.is_compact).ok_or(JournalError::ZerocopyFailure)?;
369            Ok(value)
370        })
371    }
372
373    fn mutable_object_context(
374        &self,
375        object_type: ObjectType,
376        offset: NonZeroU64,
377    ) -> Result<MutableObjectContext> {
378        validate_offset_alignment(offset)?;
379        let journal_header = self.journal_header_ref();
380        let header_size = journal_header.header_size;
381        if offset.get() < header_size {
382            return Err(JournalError::ObjectExceedsFileBounds);
383        }
384        Ok(MutableObjectContext {
385            object_type,
386            is_compact: journal_header.has_incompatible_flag(HeaderIncompatibleFlags::Compact),
387            arena_end: header_size + journal_header.arena_size,
388        })
389    }
390
391    fn mutable_object_size(
392        wm: &mut WindowManager<M>,
393        context: MutableObjectContext,
394        offset: NonZeroU64,
395        size: Option<u64>,
396    ) -> Result<u64> {
397        match size {
398            Some(size) => Self::initialize_mutable_object_header(wm, context, offset, size),
399            None => Self::existing_mutable_object_size(wm, context, offset),
400        }
401    }
402
403    fn initialize_mutable_object_header(
404        wm: &mut WindowManager<M>,
405        context: MutableObjectContext,
406        offset: NonZeroU64,
407        size: u64,
408    ) -> Result<u64> {
409        let header_slice =
410            wm.get_slice_mut(offset.get(), std::mem::size_of::<ObjectHeader>() as u64)?;
411        let header = ObjectHeader::mut_from_bytes(header_slice)
412            .map_err(|_| JournalError::ZerocopyFailure)?;
413        header.type_ = context.object_type as u8;
414        header.size = size;
415        Ok(size)
416    }
417
418    fn existing_mutable_object_size(
419        wm: &mut WindowManager<M>,
420        context: MutableObjectContext,
421        offset: NonZeroU64,
422    ) -> Result<u64> {
423        let header_slice =
424            wm.get_slice(offset.get(), std::mem::size_of::<ObjectHeader>() as u64)?;
425        let header = ObjectHeader::ref_from_bytes(header_slice)
426            .map_err(|_| JournalError::ZerocopyFailure)?;
427        if header.type_ != context.object_type as u8 {
428            return Err(JournalError::InvalidObjectType);
429        }
430        let size_needed = header.validated_size()?;
431        Self::validate_mutable_object_bounds(context, offset, size_needed)?;
432        Ok(size_needed)
433    }
434
435    fn validate_mutable_object_bounds(
436        context: MutableObjectContext,
437        offset: NonZeroU64,
438        size_needed: u64,
439    ) -> Result<()> {
440        let end_offset = offset
441            .get()
442            .checked_add(size_needed)
443            .ok_or(JournalError::ObjectExceedsFileBounds)?;
444        if end_offset > context.arena_end {
445            return Err(JournalError::ObjectExceedsFileBounds);
446        }
447        Ok(())
448    }
449
450    pub fn offset_array_mut(
451        &self,
452        offset: NonZeroU64,
453        capacity: Option<NonZeroU64>,
454    ) -> Result<ValueGuard<'_, OffsetArrayObject<&mut [u8]>>> {
455        let size = capacity.map(|c| {
456            let mut size = std::mem::size_of::<OffsetArrayObjectHeader>() as u64;
457
458            let is_compact = self
459                .journal_header_ref()
460                .has_incompatible_flag(HeaderIncompatibleFlags::Compact);
461            if is_compact {
462                size += c.get() * std::mem::size_of::<u32>() as u64;
463            } else {
464                size += c.get() * std::mem::size_of::<u64>() as u64;
465            }
466
467            size
468        });
469
470        self.journal_object_mut(ObjectType::EntryArray, offset, size)
471    }
472
473    pub fn field_mut(
474        &self,
475        offset: NonZeroU64,
476        size: Option<u64>,
477    ) -> Result<ValueGuard<'_, FieldObject<&mut [u8]>>> {
478        let size = size.map(|n| std::mem::size_of::<FieldObjectHeader>() as u64 + n);
479        self.journal_object_mut(ObjectType::Field, offset, size)
480    }
481
482    pub fn entry_mut(
483        &self,
484        offset: NonZeroU64,
485        size: Option<u64>,
486    ) -> Result<ValueGuard<'_, EntryObject<&mut [u8]>>> {
487        let size = size.map(|n| std::mem::size_of::<EntryObjectHeader>() as u64 + n);
488        self.journal_object_mut(ObjectType::Entry, offset, size)
489    }
490
491    pub fn data_mut(
492        &self,
493        offset: NonZeroU64,
494        size: Option<u64>,
495    ) -> Result<ValueGuard<'_, DataObject<&mut [u8]>>> {
496        let size = size.map(|n| {
497            let mut size = std::mem::size_of::<DataObjectHeader>() as u64 + n;
498            if self
499                .journal_header_ref()
500                .has_incompatible_flag(HeaderIncompatibleFlags::Compact)
501            {
502                size += std::mem::size_of::<CompactDataFields>() as u64;
503            }
504            size
505        });
506        self.journal_object_mut(ObjectType::Data, offset, size)
507    }
508
509    pub fn tag_mut(
510        &self,
511        offset: NonZeroU64,
512        new: bool,
513    ) -> Result<ValueGuard<'_, TagObject<&mut [u8]>>> {
514        let size = if new {
515            Some(std::mem::size_of::<TagObjectHeader>() as u64)
516        } else {
517            None
518        };
519        self.journal_object_mut(ObjectType::Tag, offset, size)
520    }
521}
522
523macro_rules! impl_hash_table_set_tail_offset {
524    (
525        $method_name:ident,
526        $hash_table_ref:ident,
527        $hash_table_mut:ident,
528        $object_mut:ident
529    ) => {
530        pub fn $method_name(&mut self, hash: u64, object_offset: NonZeroU64) -> Result<()> {
531            let hash_item = {
532                let Some(ht) = self.$hash_table_ref() else {
533                    return Err(JournalError::MissingHashTable);
534                };
535                *ht.hash_item_ref(hash)
536            };
537
538            if let Some(tail_hash_offset) = hash_item.tail_hash_offset {
539                let mut tail_object = self.$object_mut(tail_hash_offset, None)?;
540                tail_object.set_next_hash_offset(object_offset);
541            }
542
543            let Some(mut ht) = self.$hash_table_mut() else {
544                return Err(JournalError::MissingHashTable);
545            };
546
547            let hash_item = ht.hash_item_mut(hash);
548            if hash_item.head_hash_offset.is_none() {
549                hash_item.head_hash_offset = Some(object_offset);
550            }
551            hash_item.tail_hash_offset = Some(object_offset);
552
553            Ok(())
554        }
555    };
556}
557
558impl<M: MemoryMapMut> JournalFile<M> {
559    impl_hash_table_set_tail_offset!(
560        data_hash_table_set_tail_offset,
561        data_hash_table_ref,
562        data_hash_table_mut,
563        data_mut
564    );
565
566    impl_hash_table_set_tail_offset!(
567        field_hash_table_set_tail_offset,
568        field_hash_table_ref,
569        field_hash_table_mut,
570        field_mut
571    );
572}