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 pub fn sync(&mut self) -> Result<()> {
102 self.header_map.flush()?;
104
105 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 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 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}