1use crate::error::{JournalError, Result};
2use journal_common::compat::is_multiple_of;
3use std::fs::File;
4#[cfg(not(unix))]
5use std::io::{Read, Seek, SeekFrom, Write};
6use std::ops::{Deref, DerefMut};
7#[cfg(unix)]
8use std::os::unix::fs::FileExt;
9use std::sync::atomic::{Ordering, fence};
10use tracing::error;
11
12pub use memmap2::{Mmap, MmapMut, MmapOptions};
14
15const PAGE_SIZE: u64 = 4096;
16
17pub(super) fn read_file_exact_at(file: &File, position: u64, output: &mut [u8]) -> Result<()> {
18 #[cfg(unix)]
19 {
20 let mut read = 0usize;
21 while read < output.len() {
22 let bytes_read = file.read_at(&mut output[read..], position + read as u64)?;
23 if bytes_read == 0 {
24 return Err(JournalError::Io(std::io::Error::new(
25 std::io::ErrorKind::UnexpectedEof,
26 "short journal file read",
27 )));
28 }
29 read += bytes_read;
30 }
31 }
32
33 #[cfg(not(unix))]
34 {
35 let mut file = file;
36 file.seek(SeekFrom::Start(position))?;
37 file.read_exact(output)?;
38 }
39
40 Ok(())
41}
42
43#[doc(hidden)]
44#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
45pub enum ExperimentalMmapStrategy {
46 #[default]
47 Windowed,
48 WholeFile,
49}
50
51#[doc(hidden)]
52#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
53pub struct WindowManagerStats {
54 pub strategy: ExperimentalMmapStrategy,
55 pub file_size: u64,
56 pub window_count: usize,
57 pub row_pin_count: usize,
58 pub row_pin_limit: usize,
59 pub row_overflow_object_count: usize,
60 pub current_mapped_bytes: u64,
61 pub max_mapped_bytes: u64,
62 pub map_count: u64,
63 pub remap_count: u64,
64 pub eviction_count: u64,
65}
66
67pub trait MemoryMap: Deref<Target = [u8]> {
68 fn create(file: &File, offset: u64, size: u64) -> Result<Self>
69 where
70 Self: Sized;
71
72 fn create_checked(file: &File, offset: u64, size: u64, file_size: u64) -> Result<Self>
73 where
74 Self: Sized,
75 {
76 let _ = file_size;
77 Self::create(file, offset, size)
78 }
79}
80
81pub trait MemoryMapMut: MemoryMap + DerefMut {
82 fn flush(&self) -> Result<()>;
84}
85
86impl MemoryMap for Mmap {
87 fn create(file: &File, offset: u64, size: u64) -> Result<Self> {
88 let end = offset
89 .checked_add(size)
90 .ok_or(JournalError::ObjectExceedsFileBounds)?;
91 let file_size = file.metadata()?.len();
92 if end > file_size {
93 return Err(JournalError::ObjectExceedsFileBounds);
94 }
95 Self::create_checked(file, offset, size, file_size)
96 }
97
98 fn create_checked(file: &File, offset: u64, size: u64, file_size: u64) -> Result<Self> {
99 let end = offset
100 .checked_add(size)
101 .ok_or(JournalError::ObjectExceedsFileBounds)?;
102 if end > file_size {
103 return Err(JournalError::ObjectExceedsFileBounds);
104 }
105 let mmap = unsafe {
109 MmapOptions::new()
110 .offset(offset)
111 .len(size as usize)
112 .map(file)?
113 };
114
115 Ok(mmap)
116 }
117}
118
119impl MemoryMap for MmapMut {
120 fn create(file: &File, offset: u64, size: u64) -> Result<Self> {
121 let required_size = offset
122 .checked_add(size)
123 .ok_or(JournalError::ObjectExceedsFileBounds)?;
124
125 let mut file_size = file.metadata()?.len();
126 if required_size > file_size {
127 file.set_len(required_size)?;
128 file_size = required_size;
129 }
130 Self::create_checked(file, offset, size, file_size)
131 }
132
133 fn create_checked(file: &File, offset: u64, size: u64, file_size: u64) -> Result<Self> {
134 let required_size = offset
135 .checked_add(size)
136 .ok_or(JournalError::ObjectExceedsFileBounds)?;
137 if required_size > file_size {
138 return Err(JournalError::ObjectExceedsFileBounds);
139 }
140
141 let mmap = unsafe {
145 MmapOptions::new()
146 .offset(offset)
147 .len(size as usize)
148 .map_mut(file)?
149 };
150
151 Ok(mmap)
152 }
153}
154
155impl MemoryMapMut for MmapMut {
156 fn flush(&self) -> Result<()> {
157 MmapMut::flush(self)?;
158 Ok(())
159 }
160}
161
162struct Window<M: MemoryMap> {
163 offset: u64,
164 size: u64,
165 mmap: M,
166 row_pinned: bool,
167}
168
169impl<M: MemoryMap> std::fmt::Debug for Window<M> {
170 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
171 f.debug_struct("Window")
172 .field("offset", &self.offset)
173 .field("size", &self.size)
174 .finish()
175 }
176}
177
178impl<M: MemoryMap> Window<M> {
179 fn end_offset(&self) -> Option<u64> {
180 self.offset.checked_add(self.size)
181 }
182
183 fn contains(&self, position: u64) -> bool {
184 self.end_offset()
185 .is_some_and(|end_offset| position >= self.offset && position < end_offset)
186 }
187
188 fn contains_range(&self, position: u64, size: u64) -> bool {
189 let Some(end) = position.checked_add(size) else {
190 return false;
191 };
192 self.end_offset()
193 .is_some_and(|end_offset| position >= self.offset && end <= end_offset)
194 }
195
196 fn get_slice(&self, position: u64, size: u64) -> &[u8] {
197 debug_assert!(self.contains_range(position, size));
198
199 let offset = (position - self.offset) as usize;
200 &self.mmap[offset..offset + size as usize]
201 }
202}
203
204impl<M: MemoryMapMut> Window<M> {
205 pub fn get_mut_slice(&mut self, position: u64, size: u64) -> &mut [u8] {
206 debug_assert!(self.contains_range(position, size));
207
208 let offset = (position - self.offset) as usize;
209 &mut self.mmap[offset..offset + size as usize]
210 }
211}
212
213pub struct WindowManager<M: MemoryMap> {
214 file: File,
215 file_size: u64,
216 retained_size: u64,
218 bounds_mode: BoundsMode,
219 strategy: ExperimentalMmapStrategy,
220 chunk_size: u64,
221 active_window_idx: Option<usize>,
222 max_windows: usize,
223 windows: Vec<Window<M>>,
224 row_pin_count: usize,
225 map_count: u64,
226 remap_count: u64,
227 eviction_count: u64,
228 max_mapped_bytes: u64,
229 row_overflow_objects: Vec<Box<[u8]>>,
230}
231
232#[derive(Clone, Copy, Debug, Eq, PartialEq)]
233enum BoundsMode {
234 LiveFile,
235 Snapshot,
236 WriterOwned,
237}
238
239impl<M: MemoryMap> WindowManager<M> {
240 pub fn new(file: File, chunk_size: u64, max_windows: usize) -> Result<Self> {
241 Self::new_with_strategy(
242 file,
243 chunk_size,
244 max_windows,
245 ExperimentalMmapStrategy::Windowed,
246 )
247 }
248
249 pub fn new_with_strategy(
250 file: File,
251 chunk_size: u64,
252 max_windows: usize,
253 strategy: ExperimentalMmapStrategy,
254 ) -> Result<Self> {
255 Self::new_with_bounds_mode(
256 file,
257 chunk_size,
258 max_windows,
259 BoundsMode::LiveFile,
260 strategy,
261 )
262 }
263
264 pub fn new_snapshot(
265 file: File,
266 chunk_size: u64,
267 max_windows: usize,
268 strategy: ExperimentalMmapStrategy,
269 ) -> Result<Self> {
270 Self::new_with_bounds_mode(
271 file,
272 chunk_size,
273 max_windows,
274 BoundsMode::Snapshot,
275 strategy,
276 )
277 }
278
279 pub fn new_writer_owned(file: File, chunk_size: u64, max_windows: usize) -> Result<Self> {
280 Self::new_writer_owned_with_strategy(
281 file,
282 chunk_size,
283 max_windows,
284 ExperimentalMmapStrategy::Windowed,
285 )
286 }
287
288 pub fn new_writer_owned_with_strategy(
289 file: File,
290 chunk_size: u64,
291 max_windows: usize,
292 strategy: ExperimentalMmapStrategy,
293 ) -> Result<Self> {
294 Self::new_with_bounds_mode(
295 file,
296 chunk_size,
297 max_windows,
298 BoundsMode::WriterOwned,
299 strategy,
300 )
301 }
302
303 fn new_with_bounds_mode(
304 file: File,
305 chunk_size: u64,
306 max_windows: usize,
307 bounds_mode: BoundsMode,
308 strategy: ExperimentalMmapStrategy,
309 ) -> Result<Self> {
310 debug_assert!(chunk_size != 0 && is_multiple_of(chunk_size, PAGE_SIZE));
311 debug_assert!(max_windows != 0);
312
313 let file_size = file.metadata()?.len();
314
315 Ok(WindowManager {
316 file,
317 file_size,
318 retained_size: file_size,
319 bounds_mode,
320 strategy,
321 chunk_size,
322 max_windows,
323 windows: Vec::new(),
324 row_pin_count: 0,
325 active_window_idx: None,
326 map_count: 0,
327 remap_count: 0,
328 eviction_count: 0,
329 max_mapped_bytes: 0,
330 row_overflow_objects: Vec::new(),
331 })
332 }
333
334 pub fn stats(&self) -> WindowManagerStats {
335 let current_mapped_bytes = self.current_mapped_bytes();
336 WindowManagerStats {
337 strategy: self.strategy,
338 file_size: self.file_size,
339 window_count: self.windows.len(),
340 row_pin_count: self.row_pin_count,
341 row_pin_limit: self.max_windows,
342 row_overflow_object_count: self.row_overflow_objects.len(),
343 current_mapped_bytes,
344 max_mapped_bytes: self.max_mapped_bytes.max(current_mapped_bytes),
345 map_count: self.map_count,
346 remap_count: self.remap_count,
347 eviction_count: self.eviction_count,
348 }
349 }
350
351 fn current_mapped_bytes(&self) -> u64 {
352 self.windows.iter().map(|window| window.size).sum()
353 }
354
355 fn record_mapped_bytes(&mut self) {
356 self.max_mapped_bytes = self.max_mapped_bytes.max(self.current_mapped_bytes());
357 }
358
359 fn refresh_file_size(&mut self) -> Result<u64> {
360 self.file_size = self.file.metadata()?.len();
361 Ok(self.file_size)
362 }
363
364 fn ensure_cached_file_contains(&mut self, end: u64) -> Result<()> {
365 if end <= self.file_size {
366 return Ok(());
367 }
368 if self.bounds_mode == BoundsMode::LiveFile && end <= self.refresh_file_size()? {
369 return Ok(());
370 }
371 Err(JournalError::ObjectExceedsFileBounds)
372 }
373
374 pub(crate) fn read_exact_at(&mut self, position: u64, output: &mut [u8]) -> Result<()> {
375 let end = position
376 .checked_add(output.len() as u64)
377 .ok_or(JournalError::ObjectExceedsFileBounds)?;
378 self.ensure_cached_file_contains(end)?;
379
380 read_file_exact_at(&self.file, position, output)
381 }
382
383 fn get_chunk_aligned_start(&self, position: u64) -> u64 {
384 (position / self.chunk_size) * self.chunk_size
385 }
386
387 fn get_chunk_aligned_end(&self, position: u64) -> Result<u64> {
388 position
389 .div_ceil(self.chunk_size)
390 .checked_mul(self.chunk_size)
391 .ok_or(JournalError::ObjectExceedsFileBounds)
392 }
393
394 fn create_window(&mut self, window_start: u64, chunk_count: u64) -> Result<Window<M>> {
395 debug_assert_ne!(chunk_count, 0);
396
397 let requested_size = chunk_count
398 .checked_mul(self.chunk_size)
399 .ok_or(JournalError::ObjectExceedsFileBounds)?;
400 let requested_end = window_start
401 .checked_add(requested_size)
402 .ok_or(JournalError::ObjectExceedsFileBounds)?;
403 let size = match self.bounds_mode {
404 BoundsMode::LiveFile => {
405 if window_start >= self.file_size {
406 self.refresh_file_size()?;
407 }
408 if window_start >= self.file_size {
409 return Err(JournalError::ObjectExceedsFileBounds);
410 }
411 requested_size.min(self.file_size - window_start)
412 }
413 BoundsMode::Snapshot => {
414 if window_start >= self.file_size {
415 return Err(JournalError::ObjectExceedsFileBounds);
416 }
417 requested_size.min(self.file_size - window_start)
418 }
419 BoundsMode::WriterOwned => {
420 if requested_end > self.file_size {
421 self.file.set_len(requested_end)?;
422 self.file_size = requested_end;
423 }
424 requested_size
425 }
426 };
427 let mmap =
428 M::create_checked(&self.file, window_start, size, self.file_size).map_err(|e| {
429 error!(
430 window_start,
431 size,
432 chunk_count,
433 chunk_size = self.chunk_size,
434 "mmap failed: {e}"
435 );
436 e
437 })?;
438 self.map_count += 1;
439 Ok(Window {
440 offset: window_start,
441 size,
442 mmap,
443 row_pinned: false,
444 })
445 }
446
447 fn lookup_window_by_range(&self, position: u64, size_needed: u64) -> Option<usize> {
448 if let Some(idx) = self.active_window_idx {
449 if self.windows[idx].contains_range(position, size_needed) {
450 return Some(idx);
451 }
452 }
453
454 for (idx, window) in self.windows.iter().enumerate() {
455 if window.contains_range(position, size_needed) {
456 return Some(idx);
457 }
458 }
459
460 None
461 }
462
463 fn lookup_window_by_position(&self, position: u64) -> Option<usize> {
464 if let Some(idx) = self.active_window_idx {
465 if self.windows[idx].contains(position) {
466 return Some(idx);
467 }
468 }
469
470 for (idx, window) in self.windows.iter().enumerate() {
471 if window.contains(position) {
472 return Some(idx);
473 }
474 }
475
476 None
477 }
478
479 pub(crate) fn active_slice_if_contains(&self, position: u64, size: u64) -> Option<&[u8]> {
480 let idx = self.active_window_idx?;
481 let window = &self.windows[idx];
482 if window.contains_range(position, size) {
483 Some(window.get_slice(position, size))
484 } else {
485 None
486 }
487 }
488
489 pub(crate) fn active_window_contains(&self, position: u64, size: u64) -> bool {
490 self.active_window_idx
491 .and_then(|idx| self.windows.get(idx))
492 .is_some_and(|window| window.contains_range(position, size))
493 }
494
495 pub(crate) fn active_slice(&self, position: u64, size: u64) -> &[u8] {
496 let idx = self
497 .active_window_idx
498 .expect("active window should exist when active_window_contains returned true");
499 let window = &self.windows[idx];
500 debug_assert!(window.contains_range(position, size));
501 window.get_slice(position, size)
502 }
503
504 pub(crate) fn clear_row_pins(&mut self) {
505 if self.row_pin_count == 0 {
506 self.row_overflow_objects.clear();
507 return;
508 }
509 for window in &mut self.windows {
510 window.row_pinned = false;
511 }
512 self.row_pin_count = 0;
513 self.row_overflow_objects.clear();
514 }
515
516 #[inline(always)]
517 pub(crate) fn row_pin_limit_reached(&self) -> bool {
518 self.strategy != ExperimentalMmapStrategy::WholeFile
519 && self.row_pin_count >= self.max_windows
520 }
521
522 #[cold]
523 #[inline(never)]
524 fn get_row_overflow_slice(&mut self, position: u64, size: u64) -> Result<&[u8]> {
525 let len = usize::try_from(size).map_err(|_| JournalError::ObjectExceedsFileBounds)?;
526 let mut data = vec![0u8; len].into_boxed_slice();
527 self.read_exact_at(position, &mut data)?;
528 self.row_overflow_objects.push(data);
529 Ok(self
530 .row_overflow_objects
531 .last()
532 .expect("just pushed row overflow object")
533 .as_ref())
534 }
535
536 pub(crate) fn get_row_pinned_slice(&mut self, position: u64, size: u64) -> Result<&[u8]> {
537 let end = position
538 .checked_add(size)
539 .ok_or(JournalError::ObjectExceedsFileBounds)?;
540 self.ensure_cached_file_contains(end)?;
541 let Some(idx) = self.get_window_index_preserving_row_pins(position, size)? else {
542 return self.get_row_overflow_slice(position, size);
543 };
544 self.active_window_idx = Some(idx);
545 if !self.windows[idx].row_pinned {
546 if self.row_pin_limit_reached() {
547 return self.get_row_overflow_slice(position, size);
548 }
549 self.windows[idx].row_pinned = true;
550 self.row_pin_count += 1;
551 }
552 let window = &mut self.windows[idx];
553 Ok(window.get_slice(position, size))
554 }
555
556 fn push_window(&mut self, window: Window<M>) -> usize {
557 self.windows.push(window);
558 self.record_mapped_bytes();
559 self.windows.len() - 1
560 }
561
562 fn chunk_span_for_range(&self, position: u64, size_needed: u64) -> Result<(u64, u64)> {
563 let range_end = position
564 .checked_add(size_needed)
565 .ok_or(JournalError::ObjectExceedsFileBounds)?;
566 let window_start = self.get_chunk_aligned_start(position);
567 let window_end = self.get_chunk_aligned_end(range_end)?;
568 Ok((window_start, (window_end - window_start) / self.chunk_size))
569 }
570
571 fn push_new_window_for_range(&mut self, position: u64, size_needed: u64) -> Result<usize> {
572 let (window_start, num_chunks) = self.chunk_span_for_range(position, size_needed)?;
573 let new_window = self.create_window(window_start, num_chunks)?;
574 Ok(self.push_window(new_window))
575 }
576
577 fn replace_window_for_range(
578 &mut self,
579 idx: usize,
580 position: u64,
581 size_needed: u64,
582 ) -> Result<usize> {
583 let (window_start, num_chunks) = self.chunk_span_for_range(position, size_needed)?;
584 let _window = self.windows.remove(idx);
585 self.active_window_idx = None;
586 let new_window = self.create_window(window_start, num_chunks)?;
587 self.remap_count += 1;
588 Ok(self.push_window(new_window))
589 }
590
591 fn evict_unpinned_window_if_full(&mut self) -> bool {
592 if self.windows.len() < self.max_windows {
593 return true;
594 }
595 let Some(idx) = self.windows.iter().position(|window| !window.row_pinned) else {
596 return false;
597 };
598 self.windows.remove(idx);
599 self.eviction_count += 1;
600 self.active_window_idx = None;
601 true
602 }
603
604 fn make_room_for_new_window(&mut self) {
605 if self.windows.len() < self.max_windows {
606 return;
607 }
608 let idx = if self.row_pin_count == 0 {
609 Some(
610 if self.active_window_idx == Some(0) && self.windows.len() > 1 {
611 1
612 } else {
613 0
614 },
615 )
616 } else {
617 self.windows.iter().position(|window| !window.row_pinned)
618 };
619 if let Some(idx) = idx {
620 self.windows.remove(idx);
621 self.eviction_count += 1;
622 self.active_window_idx = None;
623 }
624 }
625
626 fn get_window_index_preserving_row_pins(
627 &mut self,
628 position: u64,
629 size_needed: u64,
630 ) -> Result<Option<usize>> {
631 if self.strategy == ExperimentalMmapStrategy::WholeFile {
632 return self.get_whole_file_window_index_preserving_row_pins(position, size_needed);
633 }
634
635 if let Some(idx) = self.lookup_window_by_range(position, size_needed) {
636 return Ok(Some(idx));
637 }
638
639 if let Some(idx) = self.lookup_window_by_position(position) {
640 if !self.windows[idx].row_pinned {
641 if self.row_pin_limit_reached() {
642 return Ok(None);
643 }
644 return Ok(Some(self.replace_window_for_range(
649 idx,
650 position,
651 size_needed,
652 )?));
653 }
654 }
659
660 if !self.evict_unpinned_window_if_full() {
661 return Ok(None);
662 }
663
664 Ok(Some(self.push_new_window_for_range(position, size_needed)?))
665 }
666
667 fn get_whole_file_window_index_preserving_row_pins(
668 &mut self,
669 position: u64,
670 size_needed: u64,
671 ) -> Result<Option<usize>> {
672 let idx = self.get_whole_file_window_index(position, size_needed)?;
673 if !self.windows[idx].row_pinned {
674 self.windows[idx].row_pinned = true;
675 self.row_pin_count += 1;
676 }
677 Ok(Some(idx))
678 }
679
680 fn get_window_index(&mut self, position: u64, size_needed: u64) -> Result<usize> {
681 if self.strategy == ExperimentalMmapStrategy::WholeFile {
682 return self.get_whole_file_window_index(position, size_needed);
683 }
684 if let Some(idx) = self.lookup_window_by_range(position, size_needed) {
685 self.active_window_idx = Some(idx);
686 return Ok(idx);
687 }
688 if let Some(idx) = self.lookup_window_by_position(position) {
689 return self.get_overlapping_window_index(idx, position, size_needed);
690 }
691 self.get_new_window_index(position, size_needed)
692 }
693
694 fn get_overlapping_window_index(
695 &mut self,
696 idx: usize,
697 position: u64,
698 size_needed: u64,
699 ) -> Result<usize> {
700 if self.row_pin_count > 0 && self.windows[idx].row_pinned {
701 self.evict_unpinned_window_if_full();
704 let idx = self.push_new_window_for_range(position, size_needed)?;
705 self.active_window_idx = Some(idx);
706 return Ok(idx);
707 }
708 let idx = self.replace_window_for_range(idx, position, size_needed)?;
709 self.active_window_idx = Some(idx);
710 Ok(idx)
711 }
712
713 fn get_new_window_index(&mut self, position: u64, size_needed: u64) -> Result<usize> {
714 self.make_room_for_new_window();
715 let idx = self.push_new_window_for_range(position, size_needed)?;
716 self.active_window_idx = Some(idx);
717 Ok(idx)
718 }
719
720 fn get_window(&mut self, position: u64, size_needed: u64) -> Result<&mut Window<M>> {
721 let idx = self.get_window_index(position, size_needed)?;
722 Ok(&mut self.windows[idx])
723 }
724
725 fn get_whole_file_window_index(&mut self, position: u64, size_needed: u64) -> Result<usize> {
726 if let Some(idx) = self.lookup_window_by_range(position, size_needed) {
727 self.active_window_idx = Some(idx);
728 return Ok(idx);
729 }
730
731 let requested_end = position
732 .checked_add(size_needed)
733 .ok_or(JournalError::ObjectExceedsFileBounds)?;
734 match self.bounds_mode {
735 BoundsMode::LiveFile | BoundsMode::Snapshot => {
736 self.ensure_cached_file_contains(requested_end)?
737 }
738 BoundsMode::WriterOwned => {}
739 }
740 let target_end = requested_end.max(self.file_size);
741 let window_end = self.get_chunk_aligned_end(target_end)?;
742 let chunk_count = (window_end / self.chunk_size).max(1);
743
744 let had_windows = !self.windows.is_empty();
745 if had_windows {
746 self.windows.clear();
747 self.active_window_idx = None;
748 self.row_pin_count = 0;
749 self.row_overflow_objects.clear();
750 }
751
752 let new_window = self.create_window(0, chunk_count)?;
753 if had_windows {
754 self.remap_count += 1;
755 }
756 let idx = self.push_window(new_window);
757 self.active_window_idx = Some(idx);
758 Ok(idx)
759 }
760
761 pub fn get_slice(&mut self, position: u64, size: u64) -> Result<&[u8]> {
762 let end = position
763 .checked_add(size)
764 .ok_or(JournalError::ObjectExceedsFileBounds)?;
765 self.ensure_cached_file_contains(end)?;
766 let window = self.get_window(position, size)?;
767 Ok(window.get_slice(position, size))
768 }
769}
770
771impl<M: MemoryMapMut> WindowManager<M> {
772 pub fn get_slice_mut(&mut self, position: u64, size: u64) -> Result<&mut [u8]> {
773 let _end = position
774 .checked_add(size)
775 .ok_or(JournalError::ObjectExceedsFileBounds)?;
776 let window = self.get_window(position, size)?;
777 Ok(window.get_mut_slice(position, size))
778 }
779
780 pub fn sync(&mut self, logical_size: u64, header_bytes: &[u8]) -> Result<()> {
782 let logical_size = logical_size.max(self.retained_size);
783 for window in &self.windows {
784 window.mmap.flush()?;
785 }
786 self.windows.clear();
787 self.active_window_idx = None;
788 self.row_pin_count = 0;
789 self.row_overflow_objects.clear();
790 self.file.set_len(logical_size)?;
791 #[cfg(unix)]
792 {
793 let mut written = 0usize;
794 while written < header_bytes.len() {
795 written += self
796 .file
797 .write_at(&header_bytes[written..], written as u64)?;
798 }
799 }
800 #[cfg(not(unix))]
801 {
802 self.file.seek(SeekFrom::Start(0))?;
803 self.file.write_all(header_bytes)?;
804 }
805 self.file.sync_data()?;
806 self.file_size = logical_size;
807 self.retained_size = logical_size;
808 Ok(())
809 }
810
811 pub fn post_change(&mut self, logical_size: u64) -> Result<()> {
814 let logical_size = logical_size.max(self.retained_size);
815 fence(Ordering::SeqCst);
816 if logical_size < self.file_size {
817 self.windows.clear();
818 self.active_window_idx = None;
819 self.row_pin_count = 0;
820 self.row_overflow_objects.clear();
821 }
822 self.file.set_len(logical_size)?;
823 self.file_size = logical_size;
824 self.retained_size = logical_size;
825 Ok(())
826 }
827}
828
829#[cfg(test)]
830mod tests {
831 use super::*;
832 use crate::error::JournalError;
833 use std::cell::Cell;
834 use std::io::Write;
835 use std::rc::Rc;
836 use tempfile::NamedTempFile;
837
838 const PAGE_SIZE_TEST: u64 = 4096;
839
840 struct FailingMmap {
843 data: Vec<u8>,
844 }
845
846 impl Deref for FailingMmap {
847 type Target = [u8];
848 fn deref(&self) -> &[u8] {
849 &self.data
850 }
851 }
852
853 struct MockController {
855 fail_next_create: Cell<bool>,
856 create_count: Cell<usize>,
857 }
858
859 impl MockController {
860 fn new() -> Self {
861 Self {
862 fail_next_create: Cell::new(false),
863 create_count: Cell::new(0),
864 }
865 }
866
867 fn set_fail_next(&self, fail: bool) {
868 self.fail_next_create.set(fail);
869 }
870
871 fn should_fail(&self) -> bool {
872 let count = self.create_count.get();
873 self.create_count.set(count + 1);
874 self.fail_next_create.get()
875 }
876 }
877
878 thread_local! {
880 static MOCK_CONTROLLER: Rc<MockController> = Rc::new(MockController::new());
881 }
882
883 impl MemoryMap for FailingMmap {
884 fn create(_file: &File, _offset: u64, size: u64) -> Result<Self> {
885 let mmap_size = size as usize;
886 MOCK_CONTROLLER.with(|ctrl| {
887 if ctrl.should_fail() {
888 return Err(JournalError::Io(std::io::Error::new(
889 std::io::ErrorKind::Other,
890 "simulated mmap failure",
891 )));
892 }
893 Ok(FailingMmap {
895 data: vec![0u8; mmap_size],
896 })
897 })
898 }
899 }
900
901 #[test]
912 fn test_consistent_state_after_failed_remap() {
913 let mut temp_file = NamedTempFile::new().unwrap();
915 temp_file.write_all(&[0u8; 8192]).unwrap();
916 temp_file.flush().unwrap();
917
918 let file = File::open(temp_file.path()).unwrap();
919
920 let mut wm: WindowManager<FailingMmap> =
922 WindowManager::new(file, PAGE_SIZE_TEST, 1).unwrap();
923
924 MOCK_CONTROLLER.with(|ctrl| {
926 ctrl.set_fail_next(false);
927 ctrl.create_count.set(0);
928 });
929
930 {
932 let slice = wm.get_slice(0, 100).unwrap();
933 assert_eq!(slice.len(), 100);
934 }
935 assert_eq!(wm.windows.len(), 1);
936 assert_eq!(wm.active_window_idx, Some(0));
937
938 MOCK_CONTROLLER.with(|ctrl| ctrl.set_fail_next(true));
940
941 let remap_result = wm.get_slice(100, 4000);
948 assert!(remap_result.is_err(), "Expected remap to fail");
949
950 assert_eq!(wm.windows.len(), 0);
954 assert_eq!(wm.active_window_idx, None);
955
956 MOCK_CONTROLLER.with(|ctrl| ctrl.set_fail_next(false));
958
959 let result = wm.get_slice(0, 100);
961 assert!(
962 result.is_ok(),
963 "Expected get_slice to succeed after recovery"
964 );
965 assert_eq!(wm.windows.len(), 1);
966 }
967
968 #[test]
979 fn test_consistent_state_after_failed_eviction() {
980 let mut temp_file = NamedTempFile::new().unwrap();
982 temp_file.write_all(&[0u8; 8192]).unwrap();
983 temp_file.flush().unwrap();
984
985 let file = File::open(temp_file.path()).unwrap();
986
987 let mut wm: WindowManager<FailingMmap> =
989 WindowManager::new(file, PAGE_SIZE_TEST, 1).unwrap();
990
991 MOCK_CONTROLLER.with(|ctrl| {
993 ctrl.set_fail_next(false);
994 ctrl.create_count.set(0);
995 });
996
997 {
999 let _slice = wm.get_slice(0, 100).unwrap();
1000 }
1001 assert_eq!(wm.windows.len(), 1);
1002 assert_eq!(wm.active_window_idx, Some(0));
1003
1004 MOCK_CONTROLLER.with(|ctrl| ctrl.set_fail_next(true));
1006
1007 let result = wm.get_slice(4096, 100);
1015 assert!(result.is_err(), "Expected mmap to fail");
1016
1017 assert_eq!(wm.windows.len(), 0);
1021 assert_eq!(wm.active_window_idx, None);
1022
1023 MOCK_CONTROLLER.with(|ctrl| ctrl.set_fail_next(false));
1025
1026 let result = wm.get_slice(0, 100);
1028 assert!(
1029 result.is_ok(),
1030 "Expected get_slice to succeed after recovery"
1031 );
1032 assert_eq!(wm.windows.len(), 1);
1033 }
1034
1035 #[test]
1036 fn row_pinned_slice_uses_overflow_storage_at_window_limit_one() {
1037 let mut temp_file = NamedTempFile::new().unwrap();
1038 temp_file
1039 .write_all(&vec![1u8; PAGE_SIZE_TEST as usize])
1040 .unwrap();
1041 temp_file
1042 .write_all(&vec![2u8; PAGE_SIZE_TEST as usize])
1043 .unwrap();
1044 temp_file.flush().unwrap();
1045
1046 let file = File::open(temp_file.path()).unwrap();
1047 let mut wm: WindowManager<Mmap> = WindowManager::new(file, PAGE_SIZE_TEST, 1).unwrap();
1048
1049 let first = wm.get_row_pinned_slice(0, 16).unwrap();
1050 let first_ptr = first.as_ptr();
1051 let first_len = first.len();
1052 assert_eq!(first, &[1u8; 16]);
1053
1054 let second = wm.get_row_pinned_slice(PAGE_SIZE_TEST, 16).unwrap();
1055 assert_eq!(second, &[2u8; 16]);
1056
1057 let stats = wm.stats();
1058 assert_eq!(stats.row_pin_limit, 1);
1059 assert_eq!(stats.row_pin_count, 1);
1060 assert_eq!(stats.window_count, 1);
1061 assert_eq!(stats.row_overflow_object_count, 1);
1062
1063 let first_after_overflow = unsafe { std::slice::from_raw_parts(first_ptr, first_len) };
1068 assert_eq!(first_after_overflow, &[1u8; 16]);
1069
1070 wm.clear_row_pins();
1071 let stats = wm.stats();
1072 assert_eq!(stats.row_pin_count, 0);
1073 assert_eq!(stats.row_overflow_object_count, 0);
1074 }
1075
1076 #[test]
1077 fn live_reader_refreshes_file_size_only_when_access_exceeds_cache() {
1078 let mut temp_file = NamedTempFile::new().unwrap();
1079 temp_file
1080 .write_all(&vec![1u8; PAGE_SIZE_TEST as usize])
1081 .unwrap();
1082 temp_file.flush().unwrap();
1083
1084 let file = File::open(temp_file.path()).unwrap();
1085 let mut wm: WindowManager<Mmap> = WindowManager::new(file, PAGE_SIZE_TEST, 2).unwrap();
1086 assert_eq!(wm.stats().file_size, PAGE_SIZE_TEST);
1087
1088 assert_eq!(wm.get_slice(0, 16).unwrap(), &[1u8; 16]);
1089
1090 temp_file
1091 .write_all(&vec![2u8; PAGE_SIZE_TEST as usize])
1092 .unwrap();
1093 temp_file.flush().unwrap();
1094
1095 assert_eq!(wm.get_slice(128, 16).unwrap(), &[1u8; 16]);
1096 assert_eq!(wm.stats().file_size, PAGE_SIZE_TEST);
1097
1098 assert_eq!(wm.get_slice(PAGE_SIZE_TEST + 128, 16).unwrap(), &[2u8; 16]);
1099 assert_eq!(wm.stats().file_size, PAGE_SIZE_TEST * 2);
1100 }
1101
1102 #[test]
1103 fn snapshot_reader_does_not_refresh_file_size_after_growth() {
1104 let mut temp_file = NamedTempFile::new().unwrap();
1105 temp_file
1106 .write_all(&vec![1u8; PAGE_SIZE_TEST as usize])
1107 .unwrap();
1108 temp_file.flush().unwrap();
1109
1110 let file = File::open(temp_file.path()).unwrap();
1111 let mut wm: WindowManager<Mmap> = WindowManager::new_snapshot(
1112 file,
1113 PAGE_SIZE_TEST,
1114 2,
1115 ExperimentalMmapStrategy::Windowed,
1116 )
1117 .unwrap();
1118 assert_eq!(wm.stats().file_size, PAGE_SIZE_TEST);
1119
1120 temp_file
1121 .write_all(&vec![2u8; PAGE_SIZE_TEST as usize])
1122 .unwrap();
1123 temp_file.flush().unwrap();
1124
1125 assert!(matches!(
1126 wm.get_slice(PAGE_SIZE_TEST + 128, 16).unwrap_err(),
1127 JournalError::ObjectExceedsFileBounds
1128 ));
1129 assert_eq!(wm.stats().file_size, PAGE_SIZE_TEST);
1130 }
1131
1132 #[test]
1133 fn snapshot_whole_file_maps_cached_file_once() {
1134 let temp_file = NamedTempFile::new().unwrap();
1135 temp_file.as_file().set_len(PAGE_SIZE_TEST * 2).unwrap();
1136 let file = std::fs::OpenOptions::new()
1137 .read(true)
1138 .open(temp_file.path())
1139 .unwrap();
1140 let mut wm: WindowManager<Mmap> = WindowManager::new_snapshot(
1141 file,
1142 PAGE_SIZE_TEST,
1143 32,
1144 ExperimentalMmapStrategy::WholeFile,
1145 )
1146 .unwrap();
1147
1148 assert_eq!(wm.get_slice(PAGE_SIZE_TEST + 128, 16).unwrap(), &[0; 16]);
1149 assert_eq!(wm.get_slice(128, 16).unwrap(), &[0; 16]);
1150
1151 let stats = wm.stats();
1152 assert_eq!(stats.strategy, ExperimentalMmapStrategy::WholeFile);
1153 assert_eq!(stats.file_size, PAGE_SIZE_TEST * 2);
1154 assert_eq!(stats.current_mapped_bytes, PAGE_SIZE_TEST * 2);
1155 assert_eq!(stats.map_count, 1);
1156 assert_eq!(stats.remap_count, 0);
1157 }
1158
1159 #[test]
1160 fn snapshot_whole_file_does_not_refresh_file_size_after_growth() {
1161 let mut temp_file = NamedTempFile::new().unwrap();
1162 temp_file
1163 .write_all(&vec![1u8; PAGE_SIZE_TEST as usize])
1164 .unwrap();
1165 temp_file.flush().unwrap();
1166
1167 let file = File::open(temp_file.path()).unwrap();
1168 let mut wm: WindowManager<Mmap> = WindowManager::new_snapshot(
1169 file,
1170 PAGE_SIZE_TEST,
1171 32,
1172 ExperimentalMmapStrategy::WholeFile,
1173 )
1174 .unwrap();
1175 assert_eq!(wm.get_slice(128, 16).unwrap(), &[1u8; 16]);
1176
1177 temp_file
1178 .write_all(&vec![2u8; PAGE_SIZE_TEST as usize])
1179 .unwrap();
1180 temp_file.flush().unwrap();
1181
1182 assert!(matches!(
1183 wm.get_slice(PAGE_SIZE_TEST + 128, 16).unwrap_err(),
1184 JournalError::ObjectExceedsFileBounds
1185 ));
1186 assert_eq!(wm.stats().file_size, PAGE_SIZE_TEST);
1187 }
1188
1189 #[test]
1190 fn live_whole_file_maps_cached_file_once_and_remaps_on_growth() {
1191 let mut temp_file = NamedTempFile::new().unwrap();
1192 temp_file
1193 .write_all(&vec![1u8; PAGE_SIZE_TEST as usize])
1194 .unwrap();
1195 temp_file.flush().unwrap();
1196
1197 let file = File::open(temp_file.path()).unwrap();
1198 let mut wm: WindowManager<Mmap> = WindowManager::new_with_strategy(
1199 file,
1200 PAGE_SIZE_TEST,
1201 32,
1202 ExperimentalMmapStrategy::WholeFile,
1203 )
1204 .unwrap();
1205
1206 assert_eq!(wm.get_slice(128, 16).unwrap(), &[1u8; 16]);
1207 let stats = wm.stats();
1208 assert_eq!(stats.strategy, ExperimentalMmapStrategy::WholeFile);
1209 assert_eq!(stats.file_size, PAGE_SIZE_TEST);
1210 assert_eq!(stats.current_mapped_bytes, PAGE_SIZE_TEST);
1211 assert_eq!(stats.map_count, 1);
1212 assert_eq!(stats.remap_count, 0);
1213
1214 temp_file
1215 .write_all(&vec![2u8; PAGE_SIZE_TEST as usize])
1216 .unwrap();
1217 temp_file.flush().unwrap();
1218
1219 assert_eq!(wm.get_slice(256, 16).unwrap(), &[1u8; 16]);
1220 let stats = wm.stats();
1221 assert_eq!(stats.file_size, PAGE_SIZE_TEST);
1222 assert_eq!(stats.map_count, 1);
1223 assert_eq!(stats.remap_count, 0);
1224
1225 assert_eq!(wm.get_slice(PAGE_SIZE_TEST + 128, 16).unwrap(), &[2u8; 16]);
1226 let stats = wm.stats();
1227 assert_eq!(stats.file_size, PAGE_SIZE_TEST * 2);
1228 assert_eq!(stats.current_mapped_bytes, PAGE_SIZE_TEST * 2);
1229 assert_eq!(stats.map_count, 2);
1230 assert_eq!(stats.remap_count, 1);
1231 }
1232
1233 #[test]
1234 fn whole_file_writer_owned_remaps_after_post_change_growth() {
1235 let temp_file = NamedTempFile::new().unwrap();
1236 temp_file.as_file().set_len(PAGE_SIZE_TEST).unwrap();
1237 let file = std::fs::OpenOptions::new()
1238 .read(true)
1239 .write(true)
1240 .open(temp_file.path())
1241 .unwrap();
1242 let mut wm: WindowManager<MmapMut> = WindowManager::new_writer_owned_with_strategy(
1243 file,
1244 PAGE_SIZE_TEST,
1245 32,
1246 ExperimentalMmapStrategy::WholeFile,
1247 )
1248 .unwrap();
1249
1250 wm.get_slice_mut(0, 16).unwrap().copy_from_slice(&[1; 16]);
1251 assert_eq!(wm.stats().current_mapped_bytes, PAGE_SIZE_TEST);
1252
1253 wm.post_change(PAGE_SIZE_TEST * 2).unwrap();
1254 assert_eq!(wm.get_slice(0, 16).unwrap(), &[1; 16]);
1255
1256 let new_offset = PAGE_SIZE_TEST + 128;
1257 wm.get_slice_mut(new_offset, 16)
1258 .unwrap()
1259 .copy_from_slice(&[2; 16]);
1260 assert_eq!(wm.get_slice(new_offset, 16).unwrap(), &[2; 16]);
1261
1262 let stats = wm.stats();
1263 assert_eq!(stats.current_mapped_bytes, PAGE_SIZE_TEST * 2);
1264 assert_eq!(stats.max_mapped_bytes, PAGE_SIZE_TEST * 2);
1265 assert_eq!(stats.remap_count, 1);
1266 }
1267
1268 #[test]
1269 fn publication_preserves_retained_allocation_and_trims_only_window_growth() {
1270 for strategy in [
1271 ExperimentalMmapStrategy::Windowed,
1272 ExperimentalMmapStrategy::WholeFile,
1273 ] {
1274 for sync in [false, true] {
1275 let temp_file = NamedTempFile::new().unwrap();
1276 temp_file.as_file().set_len(PAGE_SIZE_TEST * 2).unwrap();
1277 let file = temp_file.reopen().unwrap();
1278 let mut wm: WindowManager<MmapMut> = WindowManager::new_writer_owned_with_strategy(
1279 file,
1280 PAGE_SIZE_TEST * 4,
1281 32,
1282 strategy,
1283 )
1284 .unwrap();
1285 for (requested, expected) in [(1, 2), (3, 3), (1, 3)] {
1286 wm.get_slice_mut(0, 16).unwrap().copy_from_slice(&[7; 16]);
1287 assert_eq!(wm.stats().file_size, PAGE_SIZE_TEST * 4);
1288 if sync {
1289 wm.sync(PAGE_SIZE_TEST * requested, &[]).unwrap();
1290 } else {
1291 wm.post_change(PAGE_SIZE_TEST * requested).unwrap();
1292 }
1293 assert_eq!(
1294 temp_file.as_file().metadata().unwrap().len(),
1295 PAGE_SIZE_TEST * expected
1296 );
1297 assert_eq!(wm.stats().window_count, 0);
1298 }
1299 }
1300 }
1301 }
1302
1303 #[test]
1304 fn post_change_drops_mappings_before_truncating_oversized_windows() {
1305 let temp_file = NamedTempFile::new().unwrap();
1306 temp_file.as_file().set_len(PAGE_SIZE_TEST).unwrap();
1307 let file = std::fs::OpenOptions::new()
1308 .read(true)
1309 .write(true)
1310 .open(temp_file.path())
1311 .unwrap();
1312 let oversized_window = PAGE_SIZE_TEST * 4;
1313 let mut wm: WindowManager<MmapMut> = WindowManager::new_writer_owned_with_strategy(
1314 file,
1315 oversized_window,
1316 32,
1317 ExperimentalMmapStrategy::WholeFile,
1318 )
1319 .unwrap();
1320
1321 wm.get_slice_mut(0, 16).unwrap().copy_from_slice(&[1; 16]);
1322 assert_eq!(wm.stats().current_mapped_bytes, oversized_window);
1323
1324 wm.post_change(PAGE_SIZE_TEST * 2).unwrap();
1325 let stats_after_truncate = wm.stats();
1326 assert_eq!(stats_after_truncate.file_size, PAGE_SIZE_TEST * 2);
1327 assert_eq!(stats_after_truncate.current_mapped_bytes, 0);
1328 assert_eq!(stats_after_truncate.window_count, 0);
1329
1330 let crossing_offset = PAGE_SIZE_TEST * 2 - 1024;
1331 let crossing_payload = vec![2; 2048];
1332 wm.get_slice_mut(crossing_offset, 2048)
1333 .unwrap()
1334 .copy_from_slice(&crossing_payload);
1335 assert_eq!(
1336 wm.get_slice(crossing_offset, 2048).unwrap(),
1337 crossing_payload.as_slice()
1338 );
1339 assert_eq!(wm.stats().current_mapped_bytes, oversized_window);
1340 }
1341
1342 #[test]
1343 fn sequential_boundary_crossing_slides_window_instead_of_growing_from_start() {
1344 let mut temp_file = NamedTempFile::new().unwrap();
1345 temp_file.write_all(&[0u8; 64 * 1024]).unwrap();
1346 temp_file.flush().unwrap();
1347
1348 let file = File::open(temp_file.path()).unwrap();
1349 let mut wm: WindowManager<FailingMmap> =
1350 WindowManager::new(file, PAGE_SIZE_TEST, 1).unwrap();
1351
1352 MOCK_CONTROLLER.with(|ctrl| {
1353 ctrl.set_fail_next(false);
1354 ctrl.create_count.set(0);
1355 });
1356
1357 let _ = wm.get_slice(0, 100).unwrap();
1358 assert_eq!(wm.windows[0].offset, 0);
1359 assert_eq!(wm.windows[0].size, PAGE_SIZE_TEST);
1360
1361 let _ = wm.get_slice(PAGE_SIZE_TEST - 6, 32).unwrap();
1362 assert_eq!(wm.windows[0].offset, 0);
1363 assert_eq!(wm.windows[0].size, PAGE_SIZE_TEST * 2);
1364
1365 let _ = wm.get_slice((PAGE_SIZE_TEST * 2) - 12, 32).unwrap();
1366 assert_eq!(wm.windows[0].offset, PAGE_SIZE_TEST);
1367 assert_eq!(wm.windows[0].size, PAGE_SIZE_TEST * 2);
1368
1369 let _ = wm.get_slice((PAGE_SIZE_TEST * 3) - 20, 32).unwrap();
1370 assert_eq!(wm.windows[0].offset, PAGE_SIZE_TEST * 2);
1371 assert_eq!(wm.windows[0].size, PAGE_SIZE_TEST * 2);
1372 }
1373}