1use crate::{
5 AllocatedSection, MAIN_MAGIC, MainHeader, SECTION_MAGIC, SectionHandle, SectionHeader,
6 SectionStorage, UNIFIED_LOG_FORMAT_VERSION, UnifiedLogRead, UnifiedLogStatus, UnifiedLogWrite,
7};
8
9use crate::SECTION_HEADER_COMPACT_SIZE;
10
11use AllocatedSection::Section;
12use bincode::config::standard;
13use bincode::enc::EncoderImpl;
14use bincode::enc::write::SliceWriter;
15use bincode::error::EncodeError;
16use bincode::{Encode, decode_from_slice, encode_into_slice};
17use core::slice::from_raw_parts_mut;
18use cu29_traits::{
19 CuError, CuResult, ObservedWriter, UnifiedLogType, abort_observed_encode,
20 begin_observed_encode, finish_observed_encode,
21};
22use memmap2::{Mmap, MmapMut};
23use std::fs::{File, OpenOptions};
24use std::io::Read;
25use std::mem::ManuallyDrop;
26use std::path::{Path, PathBuf};
27use std::{io, mem};
28
29pub struct MmapSectionStorage {
30 buffer: &'static mut [u8],
31 offset: usize,
32 block_size: usize,
33}
34
35impl MmapSectionStorage {
36 pub fn new(buffer: &'static mut [u8], block_size: usize) -> Self {
37 Self {
38 buffer,
39 offset: 0,
40 block_size,
41 }
42 }
43
44 pub fn buffer_ptr(&self) -> *const u8 {
45 &self.buffer[0] as *const u8
46 }
47}
48
49impl SectionStorage for MmapSectionStorage {
50 fn initialize<E: Encode>(&mut self, header: &E) -> Result<usize, EncodeError> {
51 self.post_update_header(header)?;
52 self.offset = self.block_size;
53 Ok(self.offset)
54 }
55
56 fn post_update_header<E: Encode>(&mut self, header: &E) -> Result<usize, EncodeError> {
57 encode_into_slice(header, &mut self.buffer[0..], standard())
58 }
59
60 fn append<E: Encode>(&mut self, entry: &E) -> Result<usize, EncodeError> {
61 begin_observed_encode();
62 let result = (|| {
63 let mut encoder = EncoderImpl::new(
64 ObservedWriter::new(SliceWriter::new(&mut self.buffer[self.offset..])),
65 standard(),
66 );
67 entry.encode(&mut encoder)?;
68 Ok(encoder.into_writer().into_inner().bytes_written())
69 })();
70 let size = match result {
71 Ok(size) => {
72 debug_assert_eq!(size, finish_observed_encode());
73 size
74 }
75 Err(err) => {
76 abort_observed_encode();
77 return Err(err);
78 }
79 };
80 self.offset += size;
81 Ok(size)
82 }
83
84 fn flush(&mut self) -> CuResult<usize> {
85 Ok(self.offset)
87 }
88}
89
90pub enum MmapUnifiedLogger {
93 Read(MmapUnifiedLoggerRead),
94 Write(MmapUnifiedLoggerWrite),
95}
96
97pub struct MmapUnifiedLoggerBuilder {
99 file_base_name: Option<PathBuf>,
100 preallocated_size: Option<usize>,
101 write: bool,
102 create: bool,
103 append: bool,
104}
105
106impl Default for MmapUnifiedLoggerBuilder {
107 fn default() -> Self {
108 Self::new()
109 }
110}
111
112impl MmapUnifiedLoggerBuilder {
113 pub fn new() -> Self {
114 Self {
115 file_base_name: None,
116 preallocated_size: None,
117 write: false,
118 create: false, append: false,
120 }
121 }
122
123 pub fn file_base_name(mut self, file_path: &Path) -> Self {
125 self.file_base_name = Some(file_path.to_path_buf());
126 self
127 }
128
129 pub fn preallocated_size(mut self, preallocated_size: usize) -> Self {
130 self.preallocated_size = Some(preallocated_size);
131 self
132 }
133
134 pub fn write(mut self, write: bool) -> Self {
135 self.write = write;
136 self
137 }
138
139 pub fn create(mut self, create: bool) -> Self {
140 self.create = create;
141 self
142 }
143
144 pub fn append(mut self, append: bool) -> Self {
150 self.append = append;
151 self
152 }
153
154 pub fn build(self) -> io::Result<MmapUnifiedLogger> {
155 let page_size = page_size::get();
156
157 if self.write && self.create {
158 let file_path = self.file_base_name.ok_or_else(|| {
159 io::Error::new(
160 io::ErrorKind::InvalidInput,
161 "File path is required for write mode",
162 )
163 })?;
164 let preallocated_size = self.preallocated_size.ok_or_else(|| {
165 io::Error::new(
166 io::ErrorKind::InvalidInput,
167 "Preallocated size is required for write mode",
168 )
169 })?;
170 let ulw = if self.append {
171 MmapUnifiedLoggerWrite::append(&file_path, preallocated_size)?
172 } else {
173 MmapUnifiedLoggerWrite::new(&file_path, preallocated_size, page_size)?
174 };
175 Ok(MmapUnifiedLogger::Write(ulw))
176 } else {
177 let file_path = self.file_base_name.ok_or_else(|| {
178 io::Error::new(io::ErrorKind::InvalidInput, "File path is required")
179 })?;
180 let ulr = MmapUnifiedLoggerRead::new(&file_path)?;
181 Ok(MmapUnifiedLogger::Read(ulr))
182 }
183 }
184}
185
186struct SlabEntry {
187 file: File,
188 mmap_buffer: ManuallyDrop<MmapMut>,
189 current_global_position: usize,
190 sections_offsets_in_flight: Vec<usize>,
191 flushed_until_offset: usize,
192 page_size: usize,
193 temporary_end_marker: Option<usize>,
194 #[cfg(test)]
195 closed_sections: Vec<(usize, usize)>,
196 #[cfg(test)]
197 flushed_ranges: Vec<(usize, usize)>,
198 #[cfg(all(test, feature = "mmap-fsync"))]
199 sync_call_count: usize,
200}
201
202impl Drop for SlabEntry {
203 fn drop(&mut self) {
204 self.flush_until(self.current_global_position);
205 unsafe { ManuallyDrop::drop(&mut self.mmap_buffer) };
207 if let Err(error) = self.file.set_len(self.current_global_position as u64) {
208 eprintln!("Failed to trim datalogger file: {}", error);
209 }
210 self.sync_file();
211
212 if !self.sections_offsets_in_flight.is_empty() {
213 eprintln!("Error: Slab not full flushed.");
214 }
215 }
216}
217
218impl SlabEntry {
219 fn new(file: File, page_size: usize) -> io::Result<Self> {
220 let mmap_buffer = ManuallyDrop::new(
221 unsafe { MmapMut::map_mut(&file) }
223 .map_err(|e| io::Error::new(e.kind(), format!("Failed to map file: {e}")))?,
224 );
225 Ok(Self {
226 file,
227 mmap_buffer,
228 current_global_position: 0,
229 sections_offsets_in_flight: Vec::with_capacity(16),
230 flushed_until_offset: 0,
231 page_size,
232 temporary_end_marker: None,
233 #[cfg(test)]
234 closed_sections: Vec::new(),
235 #[cfg(test)]
236 flushed_ranges: Vec::new(),
237 #[cfg(all(test, feature = "mmap-fsync"))]
238 sync_call_count: 0,
239 })
240 }
241
242 fn flush_range(&mut self, start: usize, len: usize) {
243 if len == 0 {
244 return;
245 }
246 self.mmap_buffer
247 .flush_async_range(start, len)
248 .expect("Failed to flush memory map");
249 self.sync_file();
250 #[cfg(test)]
251 self.record_flushed_range(start, len);
252 }
253
254 fn sync_file(&mut self) {
255 #[cfg(feature = "mmap-fsync")]
256 {
257 self.file.sync_all().expect("Failed to fsync log file");
258 #[cfg(test)]
259 {
260 self.sync_call_count += 1;
261 }
262 }
263 }
264 fn flush_until(&mut self, until_position: usize) {
266 if (self.flushed_until_offset == until_position) || (until_position == 0) {
268 return;
269 }
270 self.flush_range(
271 self.flushed_until_offset,
272 until_position - self.flushed_until_offset,
273 );
274 self.flushed_until_offset = until_position;
275 }
276
277 fn clear_temporary_end_marker(&mut self) {
278 if let Some(marker_start) = self.temporary_end_marker.take() {
279 self.current_global_position = marker_start;
280 if self.flushed_until_offset > marker_start {
281 self.flushed_until_offset = marker_start;
282 }
283 }
284 }
285
286 fn write_end_marker(&mut self, temporary: bool) -> CuResult<()> {
287 let block_size = SECTION_HEADER_COMPACT_SIZE as usize;
288 let marker_start = self.align_to_next_page(self.current_global_position);
289 let total_marker_size = block_size; let marker_end = marker_start + total_marker_size;
291 if marker_end > self.mmap_buffer.len() {
292 return Err("Not enough space to write end-of-log marker".into());
293 }
294
295 let header = SectionHeader {
296 magic: SECTION_MAGIC,
297 block_size: SECTION_HEADER_COMPACT_SIZE,
298 entry_type: UnifiedLogType::LastEntry,
299 offset_to_next_section: total_marker_size as u32,
300 used: 0,
301 is_open: temporary,
302 };
303
304 encode_into_slice(
305 &header,
306 &mut self.mmap_buffer
307 [marker_start..marker_start + SECTION_HEADER_COMPACT_SIZE as usize],
308 standard(),
309 )
310 .map_err(|e| CuError::new_with_cause("Failed to encode end-of-log header", e))?;
311
312 self.temporary_end_marker = Some(marker_start);
313 self.current_global_position = marker_end;
314 Ok(())
315 }
316
317 fn is_it_my_section(&self, section: &SectionHandle<MmapSectionStorage>) -> bool {
318 let storage = section.get_storage();
319 let ptr = storage.buffer_ptr();
320 (ptr >= self.mmap_buffer.as_ptr())
321 && (ptr as usize)
322 < (self.mmap_buffer.as_ref().as_ptr() as usize + self.mmap_buffer.as_ref().len())
323 }
324
325 fn flush_section(&mut self, section: &mut SectionHandle<MmapSectionStorage>) {
328 section
329 .post_update_header()
330 .expect("Failed to update section header");
331
332 let storage = section.get_storage();
333 let ptr = storage.buffer_ptr();
334
335 if ptr < self.mmap_buffer.as_ptr()
336 || ptr as usize > self.mmap_buffer.as_ptr() as usize + self.mmap_buffer.len()
337 {
338 panic!("Invalid section buffer, not in the slab");
339 }
340
341 let base = self.mmap_buffer.as_ptr() as usize;
342 let section_start = ptr as usize - base;
343 let section_len = section.header.offset_to_next_section as usize;
344 #[cfg(test)]
345 self.record_closed_section(section_start, section_len);
346 self.sections_offsets_in_flight
347 .retain(|&x| x != section_start);
348
349 if self.sections_offsets_in_flight.is_empty() {
350 self.flush_until(self.current_global_position);
351 return;
352 }
353 let next_open_offset = self.sections_offsets_in_flight[0];
354 if self.flushed_until_offset < next_open_offset {
355 self.flush_until(next_open_offset);
356 }
357 if section_start + section_len > self.flushed_until_offset {
358 self.flush_range(section_start, section_len);
361 }
362 }
363
364 #[cfg(test)]
365 fn record_closed_section(&mut self, start: usize, len: usize) {
366 self.closed_sections.push((start, len));
367 }
368
369 #[cfg(test)]
370 fn record_flushed_range(&mut self, start: usize, len: usize) {
371 let mut merged_start = start;
372 let mut merged_end = start + len;
373 let mut merged_ranges = Vec::with_capacity(self.flushed_ranges.len() + 1);
374 let mut inserted = false;
375
376 for (range_start, range_len) in self.flushed_ranges.drain(..) {
377 let range_end = range_start + range_len;
378 if range_end < merged_start {
379 merged_ranges.push((range_start, range_len));
380 continue;
381 }
382 if merged_end < range_start {
383 if !inserted {
384 merged_ranges.push((merged_start, merged_end - merged_start));
385 inserted = true;
386 }
387 merged_ranges.push((range_start, range_len));
388 continue;
389 }
390
391 merged_start = merged_start.min(range_start);
392 merged_end = merged_end.max(range_end);
393 }
394
395 if !inserted {
396 merged_ranges.push((merged_start, merged_end - merged_start));
397 }
398
399 self.flushed_ranges = merged_ranges;
400 }
401
402 #[cfg(test)]
403 fn pending_closed_bytes(&self) -> usize {
404 let mut pending = 0;
405
406 for (section_start, section_len) in &self.closed_sections {
407 let section_end = section_start + section_len;
408 let mut cursor = *section_start;
409
410 for (range_start, range_len) in &self.flushed_ranges {
411 let range_end = range_start + range_len;
412 if range_end <= cursor {
413 continue;
414 }
415 if *range_start >= section_end {
416 break;
417 }
418 if *range_start > cursor {
419 pending += *range_start - cursor;
420 }
421 cursor = cursor.max(range_end);
422 if cursor >= section_end {
423 break;
424 }
425 }
426
427 if cursor < section_end {
428 pending += section_end - cursor;
429 }
430 }
431
432 pending
433 }
434
435 #[inline]
436 fn align_to_next_page(&self, ptr: usize) -> usize {
437 (ptr + self.page_size - 1) & !(self.page_size - 1)
438 }
439
440 fn add_section(
442 &mut self,
443 entry_type: UnifiedLogType,
444 requested_section_size: usize,
445 ) -> AllocatedSection<MmapSectionStorage> {
446 self.current_global_position = self.align_to_next_page(self.current_global_position);
448 let section_size = self.align_to_next_page(requested_section_size) as u32;
449
450 if self.current_global_position + section_size as usize > self.mmap_buffer.len() {
452 return AllocatedSection::NoMoreSpace;
453 }
454
455 #[cfg(feature = "compact")]
456 let block_size = SECTION_HEADER_COMPACT_SIZE;
457
458 #[cfg(not(feature = "compact"))]
459 let block_size = self.page_size as u16;
460
461 let section_header = SectionHeader {
462 magic: SECTION_MAGIC,
463 block_size,
464 entry_type,
465 offset_to_next_section: section_size,
466 used: 0u32,
467 is_open: true,
468 };
469
470 self.sections_offsets_in_flight
472 .push(self.current_global_position);
473 let end_of_section = self.current_global_position + requested_section_size;
474 let user_buffer = &mut self.mmap_buffer[self.current_global_position..end_of_section];
475
476 let handle_buffer =
478 unsafe { from_raw_parts_mut(user_buffer.as_mut_ptr(), user_buffer.len()) };
479 let storage = MmapSectionStorage::new(handle_buffer, block_size as usize);
480
481 self.current_global_position = end_of_section;
482
483 Section(SectionHandle::create(section_header, storage).expect("Failed to create section"))
484 }
485
486 #[cfg(test)]
487 fn used(&self) -> usize {
488 self.current_global_position
489 }
490}
491
492pub struct MmapUnifiedLoggerWrite {
494 front_slab: SlabEntry,
496 back_slabs: Vec<SlabEntry>,
498 base_file_path: PathBuf,
500 slab_size: usize,
502 front_slab_suffix: usize,
504}
505
506fn build_slab_path(base_file_path: &Path, slab_index: usize) -> io::Result<PathBuf> {
507 let mut file_path = base_file_path.to_path_buf();
508 let stem = file_path.file_stem().ok_or_else(|| {
509 io::Error::new(
510 io::ErrorKind::InvalidInput,
511 "Base file path has no file name",
512 )
513 })?;
514 let stem = stem.to_str().ok_or_else(|| {
515 io::Error::new(
516 io::ErrorKind::InvalidInput,
517 "Base file name is not valid UTF-8",
518 )
519 })?;
520 let extension = file_path.extension().ok_or_else(|| {
521 io::Error::new(
522 io::ErrorKind::InvalidInput,
523 "Base file path has no extension",
524 )
525 })?;
526 let extension = extension.to_str().ok_or_else(|| {
527 io::Error::new(
528 io::ErrorKind::InvalidInput,
529 "Base file extension is not valid UTF-8",
530 )
531 })?;
532 if stem.is_empty() {
533 return Err(io::Error::new(
534 io::ErrorKind::InvalidInput,
535 "Base file name is empty",
536 ));
537 }
538 let file_name = format!("{stem}_{slab_index}.{extension}");
539 file_path.set_file_name(file_name);
540 Ok(file_path)
541}
542
543fn make_slab_file(base_file_path: &Path, slab_size: usize, slab_suffix: usize) -> io::Result<File> {
544 let file_path = build_slab_path(base_file_path, slab_suffix)?;
545 let file = OpenOptions::new()
546 .read(true)
547 .write(true)
548 .create(true)
549 .truncate(true)
550 .open(&file_path)
551 .map_err(|e| {
552 io::Error::new(
553 e.kind(),
554 format!("Failed to open file {}: {e}", file_path.display()),
555 )
556 })?;
557 file.set_len(slab_size as u64).map_err(|e| {
558 io::Error::new(
559 e.kind(),
560 format!("Failed to set file length for {}: {e}", file_path.display()),
561 )
562 })?;
563 Ok(file)
564}
565
566fn remove_existing_alias(base_file_path: &Path) -> io::Result<()> {
567 match std::fs::symlink_metadata(base_file_path) {
568 Ok(meta) => {
569 if meta.is_dir() {
570 return Err(io::Error::new(
571 io::ErrorKind::AlreadyExists,
572 format!(
573 "Cannot create base log alias at {} because a directory already exists there",
574 base_file_path.display()
575 ),
576 ));
577 }
578 std::fs::remove_file(base_file_path).map_err(|e| {
579 io::Error::new(
580 e.kind(),
581 format!(
582 "Failed to remove existing base log alias {}: {e}",
583 base_file_path.display()
584 ),
585 )
586 })
587 }
588 Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(()),
589 Err(e) => Err(io::Error::new(
590 e.kind(),
591 format!(
592 "Failed to inspect existing base log alias {}: {e}",
593 base_file_path.display()
594 ),
595 )),
596 }
597}
598
599fn create_base_alias_link(base_file_path: &Path) -> io::Result<()> {
600 let first_slab_path = build_slab_path(base_file_path, 0)?;
601 remove_existing_alias(base_file_path)?;
602
603 #[cfg(unix)]
604 {
605 use std::os::unix::fs::symlink;
606 let relative_target = Path::new(first_slab_path.file_name().ok_or_else(|| {
607 io::Error::new(
608 io::ErrorKind::InvalidInput,
609 "First slab file has no name component",
610 )
611 })?);
612 symlink(relative_target, base_file_path).map_err(|e| {
613 io::Error::new(
614 e.kind(),
615 format!(
616 "Failed to create base log alias {} -> {}: {e}",
617 base_file_path.display(),
618 first_slab_path.display()
619 ),
620 )
621 })
622 }
623
624 #[cfg(windows)]
625 {
626 use std::os::windows::fs::symlink_file;
627 let relative_target = Path::new(first_slab_path.file_name().ok_or_else(|| {
628 io::Error::new(
629 io::ErrorKind::InvalidInput,
630 "First slab file has no name component",
631 )
632 })?);
633 match symlink_file(relative_target, base_file_path) {
634 Ok(()) => Ok(()),
635 Err(symlink_err) => std::fs::hard_link(&first_slab_path, base_file_path).map_err(
636 |hard_link_err| {
637 io::Error::other(format!(
638 "Failed to create base log alias {}. Symlink error: {symlink_err}. Hard-link fallback error: {hard_link_err}",
639 base_file_path.display()
640 ))
641 },
642 ),
643 }?;
644 Ok(())
645 }
646
647 #[cfg(not(any(unix, windows)))]
648 {
649 std::fs::hard_link(&first_slab_path, base_file_path).map_err(|e| {
650 io::Error::new(
651 e.kind(),
652 format!(
653 "Failed to create base log alias {} -> {}: {e}",
654 base_file_path.display(),
655 first_slab_path.display()
656 ),
657 )
658 })
659 }
660}
661
662impl UnifiedLogWrite<MmapSectionStorage> for MmapUnifiedLoggerWrite {
663 fn add_section(
665 &mut self,
666 entry_type: UnifiedLogType,
667 requested_section_size: usize,
668 ) -> CuResult<SectionHandle<MmapSectionStorage>> {
669 self.garbage_collect_backslabs(); self.front_slab.clear_temporary_end_marker();
671 let maybe_section = self
672 .front_slab
673 .add_section(entry_type, requested_section_size);
674
675 match maybe_section {
676 AllocatedSection::NoMoreSpace => {
677 let new_slab = self.create_slab()?;
679 self.back_slabs
681 .push(mem::replace(&mut self.front_slab, new_slab));
682 match self
683 .front_slab
684 .add_section(entry_type, requested_section_size)
685 {
686 AllocatedSection::NoMoreSpace => Err(CuError::from("out of space")),
687 Section(section) => {
688 self.place_end_marker(true)?;
689 Ok(section)
690 }
691 }
692 }
693 Section(section) => {
694 self.place_end_marker(true)?;
695 Ok(section)
696 }
697 }
698 }
699
700 fn flush_section(&mut self, section: &mut SectionHandle<MmapSectionStorage>) {
701 section.mark_closed();
702 for slab in self.back_slabs.iter_mut() {
703 if slab.is_it_my_section(section) {
704 slab.flush_section(section);
705 return;
706 }
707 }
708 self.front_slab.flush_section(section);
709 }
710
711 fn status(&self) -> UnifiedLogStatus {
712 UnifiedLogStatus {
713 total_used_space: self.front_slab.current_global_position,
714 total_allocated_space: self.slab_size * self.front_slab_suffix,
715 }
716 }
717}
718
719impl MmapUnifiedLoggerWrite {
720 fn next_slab(&mut self) -> io::Result<File> {
721 let next_suffix = self.front_slab_suffix + 1;
722 let file = make_slab_file(&self.base_file_path, self.slab_size, next_suffix)?;
723 self.front_slab_suffix = next_suffix;
724 Ok(file)
725 }
726
727 fn new(base_file_path: &Path, slab_size: usize, page_size: usize) -> io::Result<Self> {
728 let file = make_slab_file(base_file_path, slab_size, 0)?;
729 create_base_alias_link(base_file_path)?;
730 let mut front_slab = SlabEntry::new(file, page_size)?;
731
732 let main_header = MainHeader {
734 magic: MAIN_MAGIC,
735 format_version: UNIFIED_LOG_FORMAT_VERSION,
736 first_section_offset: page_size as u16,
737 page_size: page_size as u16,
738 };
739 let nb_bytes = encode_into_slice(&main_header, &mut front_slab.mmap_buffer[..], standard())
740 .map_err(|e| io::Error::other(format!("Failed to encode main header: {e}")))?;
741 assert!(nb_bytes < page_size);
742 front_slab.current_global_position = page_size; Ok(Self {
745 front_slab,
746 back_slabs: Vec::new(),
747 base_file_path: base_file_path.to_path_buf(),
748 slab_size,
749 front_slab_suffix: 0,
750 })
751 }
752
753 fn append(base_file_path: &Path, slab_size: usize) -> io::Result<Self> {
760 let mut reader = MmapUnifiedLoggerRead::new(base_file_path)?;
762 let end = reader.end_of_log().map_err(io::Error::other)?;
763 let last_slab_index = end.slab_index;
764 let resume_offset = end.offset;
765 let original_page_size = reader.raw_main_header().page_size as usize;
767 drop(reader);
768
769 let slab_path = build_slab_path(base_file_path, last_slab_index)?;
770 let file = OpenOptions::new()
771 .read(true)
772 .write(true)
773 .open(&slab_path)
774 .map_err(|e| {
775 io::Error::new(
776 e.kind(),
777 format!(
778 "Failed to open slab {} for append: {e}",
779 slab_path.display()
780 ),
781 )
782 })?;
783
784 let current_len = file
787 .metadata()
788 .map_err(|e| {
789 io::Error::new(
790 e.kind(),
791 format!("Failed to read metadata for {}", slab_path.display()),
792 )
793 })?
794 .len();
795 if (current_len as usize) < slab_size {
796 file.set_len(slab_size as u64).map_err(|e| {
797 io::Error::new(
798 e.kind(),
799 format!(
800 "Failed to extend slab {} for append: {e}",
801 slab_path.display()
802 ),
803 )
804 })?;
805 }
806
807 let mut front_slab = SlabEntry::new(file, original_page_size)?;
808 front_slab.current_global_position = resume_offset;
809 front_slab.flushed_until_offset = resume_offset;
810
811 Ok(Self {
812 front_slab,
813 back_slabs: Vec::new(),
814 base_file_path: base_file_path.to_path_buf(),
815 slab_size,
816 front_slab_suffix: last_slab_index,
817 })
818 }
819
820 fn garbage_collect_backslabs(&mut self) {
821 self.back_slabs
822 .retain_mut(|slab| !slab.sections_offsets_in_flight.is_empty());
823 }
824
825 fn place_end_marker(&mut self, temporary: bool) -> CuResult<()> {
826 match self.front_slab.write_end_marker(temporary) {
827 Ok(_) => Ok(()),
828 Err(_) => {
829 let new_slab = self.create_slab()?;
831 self.back_slabs
832 .push(mem::replace(&mut self.front_slab, new_slab));
833 self.front_slab.write_end_marker(temporary)
834 }
835 }
836 }
837
838 pub fn stats(&self) -> (usize, Vec<usize>, usize) {
839 (
840 self.front_slab.current_global_position,
841 self.front_slab.sections_offsets_in_flight.clone(),
842 self.back_slabs.len(),
843 )
844 }
845
846 fn create_slab(&mut self) -> CuResult<SlabEntry> {
847 let file = self
848 .next_slab()
849 .map_err(|e| CuError::new_with_cause("Failed to create slab file", e))?;
850 SlabEntry::new(file, self.front_slab.page_size)
851 .map_err(|e| CuError::new_with_cause("Failed to create slab memory map", e))
852 }
853}
854
855impl Drop for MmapUnifiedLoggerWrite {
856 fn drop(&mut self) {
857 #[cfg(debug_assertions)]
858 eprintln!("Flushing the unified Logger ... "); self.front_slab.clear_temporary_end_marker();
861 if let Err(e) = self.place_end_marker(false) {
862 panic!("Failed to flush the unified logger: {}", e);
863 }
864 self.front_slab
865 .flush_until(self.front_slab.current_global_position);
866 self.garbage_collect_backslabs();
867 #[cfg(debug_assertions)]
868 eprintln!("Unified Logger flushed."); }
870}
871
872fn open_slab_index(
873 base_file_path: &Path,
874 slab_index: usize,
875) -> io::Result<(File, Mmap, u16, Option<MainHeader>)> {
876 let mut options = OpenOptions::new();
877 let options = options.read(true);
878
879 let file_path = build_slab_path(base_file_path, slab_index)?;
880 let file = options.open(&file_path).map_err(|e| {
881 io::Error::new(
882 e.kind(),
883 format!("Failed to open slab file {}: {e}", file_path.display()),
884 )
885 })?;
886 let mmap = unsafe { Mmap::map(&file) }
888 .map_err(|e| io::Error::new(e.kind(), format!("Failed to map slab file: {e}")))?;
889 let mut prolog = 0u16;
890 let mut maybe_main_header: Option<MainHeader> = None;
891 if slab_index == 0 {
892 let main_header: MainHeader;
893 let _read: usize;
894 (main_header, _read) = decode_from_slice(&mmap[..], standard()).map_err(|e| {
895 io::Error::new(
896 io::ErrorKind::InvalidData,
897 format!("Failed to decode main header: {e}"),
898 )
899 })?;
900 if main_header.magic != MAIN_MAGIC {
901 return Err(io::Error::new(
902 io::ErrorKind::InvalidData,
903 "Invalid magic number in main header",
904 ));
905 }
906 if main_header.format_version != UNIFIED_LOG_FORMAT_VERSION {
907 return Err(io::Error::new(
908 io::ErrorKind::InvalidData,
909 format!(
910 "Unsupported unified log format version {} in main header; this reader supports version {}",
911 main_header.format_version, UNIFIED_LOG_FORMAT_VERSION
912 ),
913 ));
914 }
915 prolog = main_header.first_section_offset;
916 maybe_main_header = Some(main_header);
917 }
918 Ok((file, mmap, prolog, maybe_main_header))
919}
920
921pub struct MmapUnifiedLoggerRead {
923 base_file_path: PathBuf,
924 main_header: MainHeader,
925 current_mmap_buffer: Mmap,
926 current_file: File,
927 current_slab_index: usize,
928 current_reading_position: usize,
929}
930
931#[derive(Clone, Copy, Debug, PartialEq, Eq)]
933pub struct LogPosition {
934 pub slab_index: usize,
935 pub offset: usize,
936}
937
938impl UnifiedLogRead for MmapUnifiedLoggerRead {
939 fn read_next_section_type(&mut self, datalogtype: UnifiedLogType) -> CuResult<Option<Vec<u8>>> {
940 loop {
942 if self.current_reading_position >= self.current_mmap_buffer.len() {
943 self.next_slab().map_err(|e| {
944 CuError::new_with_cause("Failed to read next slab, is the log complete?", e)
945 })?;
946 }
947
948 let header_result = self.read_section_header();
949 let header = header_result.map_err(|error| {
950 CuError::new_with_cause(
951 &format!(
952 "Could not read a sections header: {}/{}:{}",
953 self.base_file_path.as_os_str().to_string_lossy(),
954 self.current_slab_index,
955 self.current_reading_position,
956 ),
957 error,
958 )
959 })?;
960
961 if header.entry_type == UnifiedLogType::LastEntry {
963 return Ok(None);
964 }
965
966 if header.entry_type == datalogtype {
968 let result = Some(self.read_section_content(&header)?);
969 self.current_reading_position += header.offset_to_next_section as usize;
970 return Ok(result);
971 }
972
973 self.current_reading_position += header.offset_to_next_section as usize;
975 }
976 }
977
978 fn raw_read_section(&mut self) -> CuResult<(SectionHeader, Vec<u8>)> {
980 if self.current_reading_position >= self.current_mmap_buffer.len() {
981 self.next_slab().map_err(|e| {
982 CuError::new_with_cause("Failed to read next slab, is the log complete?", e)
983 })?;
984 }
985
986 let read_result = self.read_section_header();
987
988 match read_result {
989 Err(error) => Err(CuError::new_with_cause(
990 &format!(
991 "Could not read a sections header: {}/{}:{}",
992 self.base_file_path.as_os_str().to_string_lossy(),
993 self.current_slab_index,
994 self.current_reading_position,
995 ),
996 error,
997 )),
998 Ok(header) => {
999 let data = self.read_section_content(&header)?;
1000 self.current_reading_position += header.offset_to_next_section as usize;
1001 Ok((header, data))
1002 }
1003 }
1004 }
1005}
1006
1007impl MmapUnifiedLoggerRead {
1008 pub fn raw_skip_section(&mut self) -> CuResult<SectionHeader> {
1010 if self.current_reading_position >= self.current_mmap_buffer.len() {
1011 self.next_slab().map_err(|e| {
1012 CuError::new_with_cause("Failed to read next slab, is the log complete?", e)
1013 })?;
1014 }
1015
1016 let header = self.read_section_header().map_err(|error| {
1017 CuError::new_with_cause(
1018 &format!(
1019 "Could not read a sections header: {}/{}:{}",
1020 self.base_file_path.as_os_str().to_string_lossy(),
1021 self.current_slab_index,
1022 self.current_reading_position,
1023 ),
1024 error,
1025 )
1026 })?;
1027 self.current_reading_position += header.offset_to_next_section as usize;
1028 Ok(header)
1029 }
1030
1031 pub fn new(base_file_path: &Path) -> io::Result<Self> {
1032 let (file, mmap, prolog, header) = open_slab_index(base_file_path, 0)?;
1033 let main_header = header.ok_or_else(|| {
1034 io::Error::new(io::ErrorKind::InvalidData, "Missing main header in slab 0")
1035 })?;
1036
1037 Ok(Self {
1038 base_file_path: base_file_path.to_path_buf(),
1039 main_header,
1040 current_file: file,
1041 current_mmap_buffer: mmap,
1042 current_slab_index: 0,
1043 current_reading_position: prolog as usize,
1044 })
1045 }
1046
1047 pub fn position(&self) -> LogPosition {
1049 LogPosition {
1050 slab_index: self.current_slab_index,
1051 offset: self.current_reading_position,
1052 }
1053 }
1054
1055 pub fn seek(&mut self, pos: LogPosition) -> CuResult<()> {
1057 if pos.slab_index != self.current_slab_index {
1058 let (file, mmap, _prolog, _header) =
1059 open_slab_index(&self.base_file_path, pos.slab_index).map_err(|e| {
1060 CuError::new_with_cause(
1061 &format!("Failed to open slab {} for seek", pos.slab_index),
1062 e,
1063 )
1064 })?;
1065 self.current_file = file;
1066 self.current_mmap_buffer = mmap;
1067 self.current_slab_index = pos.slab_index;
1068 }
1069 self.current_reading_position = pos.offset;
1070 Ok(())
1071 }
1072
1073 fn next_slab(&mut self) -> io::Result<()> {
1074 self.current_slab_index += 1;
1075 let (file, mmap, prolog, _) =
1076 open_slab_index(&self.base_file_path, self.current_slab_index)?;
1077 self.current_file = file;
1078 self.current_mmap_buffer = mmap;
1079 self.current_reading_position = prolog as usize;
1080 Ok(())
1081 }
1082
1083 pub fn raw_main_header(&self) -> &MainHeader {
1084 &self.main_header
1085 }
1086
1087 pub fn end_of_log(&mut self) -> CuResult<LogPosition> {
1090 loop {
1091 if self.current_reading_position >= self.current_mmap_buffer.len() {
1092 self.next_slab().map_err(|e| {
1093 CuError::new_with_cause(
1094 "Failed to advance to the next slab while locating the end of the log",
1095 e,
1096 )
1097 })?;
1098 }
1099
1100 let header = self.read_section_header()?;
1101 if header.entry_type == UnifiedLogType::LastEntry {
1102 if header.is_open {
1103 return Err(CuError::from(format!(
1104 "Log {} was not cleanly closed: temporary end-of-log marker in slab {}",
1105 self.base_file_path.display(),
1106 self.current_slab_index,
1107 )));
1108 }
1109 return Ok(self.position());
1110 }
1111 self.current_reading_position += header.offset_to_next_section as usize;
1112 }
1113 }
1114
1115 pub fn scan_section_bytes(&mut self, datalogtype: UnifiedLogType) -> CuResult<u64> {
1116 let mut total = 0u64;
1117
1118 loop {
1119 if self.current_reading_position >= self.current_mmap_buffer.len() {
1120 self.next_slab().map_err(|e| {
1121 CuError::new_with_cause("Failed to read next slab, is the log complete?", e)
1122 })?;
1123 }
1124
1125 let header = self.read_section_header()?;
1126
1127 if header.entry_type == UnifiedLogType::LastEntry {
1128 return Ok(total);
1129 }
1130
1131 if header.entry_type == datalogtype {
1132 total = total.saturating_add(header.used as u64);
1133 }
1134
1135 self.current_reading_position += header.offset_to_next_section as usize;
1136 }
1137 }
1138
1139 fn read_section_content(&mut self, header: &SectionHeader) -> CuResult<Vec<u8>> {
1141 let mut section_data = vec![0; header.used as usize];
1143 let start_of_data = self.current_reading_position + header.block_size as usize;
1144 section_data.copy_from_slice(
1145 &self.current_mmap_buffer[start_of_data..start_of_data + header.used as usize],
1146 );
1147
1148 Ok(section_data)
1149 }
1150
1151 fn read_section_header(&mut self) -> CuResult<SectionHeader> {
1152 let section_header: SectionHeader;
1153 (section_header, _) = decode_from_slice(
1154 &self.current_mmap_buffer[self.current_reading_position..],
1155 standard(),
1156 )
1157 .map_err(|e| {
1158 CuError::new_with_cause(
1159 &format!(
1160 "Could not read a sections header: {}/{}:{}",
1161 self.base_file_path.as_os_str().to_string_lossy(),
1162 self.current_slab_index,
1163 self.current_reading_position,
1164 ),
1165 e,
1166 )
1167 })?;
1168 if section_header.magic != SECTION_MAGIC {
1169 return Err("Invalid magic number in section header".into());
1170 }
1171
1172 Ok(section_header)
1173 }
1174}
1175
1176pub struct UnifiedLoggerIOReader {
1178 logger: MmapUnifiedLoggerRead,
1179 log_type: UnifiedLogType,
1180 buffer: Vec<u8>,
1181 buffer_pos: usize,
1182}
1183
1184impl UnifiedLoggerIOReader {
1185 pub fn new(logger: MmapUnifiedLoggerRead, log_type: UnifiedLogType) -> Self {
1186 Self {
1187 logger,
1188 log_type,
1189 buffer: Vec::new(),
1190 buffer_pos: 0,
1191 }
1192 }
1193
1194 fn fill_buffer(&mut self) -> io::Result<bool> {
1196 match self.logger.read_next_section_type(self.log_type) {
1197 Ok(Some(section)) => {
1198 self.buffer = section;
1199 self.buffer_pos = 0;
1200 Ok(true)
1201 }
1202 Ok(None) => Ok(false), Err(e) => Err(io::Error::other(e.to_string())),
1204 }
1205 }
1206}
1207
1208impl Read for UnifiedLoggerIOReader {
1209 fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
1210 if self.buffer_pos >= self.buffer.len() && !self.fill_buffer()? {
1211 return Ok(0);
1213 }
1214
1215 if self.buffer_pos >= self.buffer.len() {
1217 return Ok(0);
1218 }
1219
1220 let len = std::cmp::min(buf.len(), self.buffer.len() - self.buffer_pos);
1222 buf[..len].copy_from_slice(&self.buffer[self.buffer_pos..self.buffer_pos + len]);
1223 self.buffer_pos += len;
1224 Ok(len)
1225 }
1226}
1227
1228#[cfg(feature = "std")]
1229#[cfg(test)]
1230mod tests {
1231 use super::*;
1232 use crate::stream_write;
1233 use bincode::de::read::SliceReader;
1234 use bincode::{Decode, Encode, decode_from_reader, decode_from_slice};
1235 use cu29_traits::WriteStream;
1236 use std::io::{Seek, SeekFrom, Write};
1237 use std::path::PathBuf;
1238 use std::sync::{Arc, Mutex};
1239 use tempfile::TempDir;
1240
1241 const LARGE_SLAB: usize = 100 * 1024; const SMALL_SLAB: usize = 16 * 2 * 1024; fn make_a_logger(
1245 tmp_dir: &TempDir,
1246 slab_size: usize,
1247 ) -> (Arc<Mutex<MmapUnifiedLoggerWrite>>, PathBuf) {
1248 let file_path = tmp_dir.path().join("test.bin");
1249 let MmapUnifiedLogger::Write(data_logger) = MmapUnifiedLoggerBuilder::new()
1250 .write(true)
1251 .create(true)
1252 .file_base_name(&file_path)
1253 .preallocated_size(slab_size)
1254 .build()
1255 .expect("Failed to create logger")
1256 else {
1257 panic!("Failed to create logger")
1258 };
1259
1260 (Arc::new(Mutex::new(data_logger)), file_path)
1261 }
1262
1263 #[test]
1264 fn test_truncation_and_sections_creations() {
1265 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1266 let file_path = tmp_dir.path().join("test.bin");
1267 let _used = {
1268 let MmapUnifiedLogger::Write(mut logger) = MmapUnifiedLoggerBuilder::new()
1269 .write(true)
1270 .create(true)
1271 .file_base_name(&file_path)
1272 .preallocated_size(100000)
1273 .build()
1274 .expect("Failed to create logger")
1275 else {
1276 panic!("Failed to create logger")
1277 };
1278 logger
1279 .add_section(UnifiedLogType::StructuredLogLine, 1024)
1280 .unwrap();
1281 logger
1282 .add_section(UnifiedLogType::CopperList, 2048)
1283 .unwrap();
1284 let used = logger.front_slab.used();
1285 assert!(used < 4 * page_size::get()); used
1289 };
1290
1291 let _file = OpenOptions::new()
1292 .read(true)
1293 .open(tmp_dir.path().join("test_0.bin"))
1294 .expect("Could not reopen the file");
1295 }
1302
1303 #[test]
1304 fn test_unsupported_main_header_format_version_is_rejected() {
1305 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1306 let file_path = tmp_dir.path().join("test.bin");
1307 {
1308 let MmapUnifiedLogger::Write(_logger) = MmapUnifiedLoggerBuilder::new()
1309 .write(true)
1310 .create(true)
1311 .file_base_name(&file_path)
1312 .preallocated_size(100000)
1313 .build()
1314 .expect("Failed to create logger")
1315 else {
1316 panic!("Failed to create logger")
1317 };
1318 }
1319
1320 let mut file = OpenOptions::new()
1321 .read(true)
1322 .write(true)
1323 .open(tmp_dir.path().join("test_0.bin"))
1324 .expect("Could not reopen the slab");
1325 let unsupported_version = UNIFIED_LOG_FORMAT_VERSION + 1;
1326 file.seek(SeekFrom::Start(MAIN_MAGIC.len() as u64))
1327 .expect("Could not seek to format version");
1328 file.write_all(&[unsupported_version])
1329 .expect("Could not write unsupported format version");
1330 drop(file);
1331
1332 let err = match MmapUnifiedLoggerBuilder::new()
1333 .file_base_name(&file_path)
1334 .build()
1335 {
1336 Ok(_) => panic!("Reader accepted unsupported unified log format version"),
1337 Err(err) => err,
1338 };
1339
1340 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1341 assert_eq!(
1342 err.to_string(),
1343 format!(
1344 "Unsupported unified log format version {unsupported_version} in main header; this reader supports version {UNIFIED_LOG_FORMAT_VERSION}"
1345 )
1346 );
1347 }
1348
1349 #[test]
1350 fn test_base_alias_exists_and_matches_first_slab() {
1351 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1352 let file_path = tmp_dir.path().join("test.bin");
1353 let _logger = MmapUnifiedLoggerBuilder::new()
1354 .write(true)
1355 .create(true)
1356 .file_base_name(&file_path)
1357 .preallocated_size(LARGE_SLAB)
1358 .build()
1359 .expect("Failed to create logger");
1360
1361 let first_slab = build_slab_path(&file_path, 0).expect("Failed to build first slab path");
1362 assert!(file_path.exists(), "base alias does not exist");
1363 assert!(first_slab.exists(), "first slab does not exist");
1364
1365 let alias_bytes = std::fs::read(&file_path).expect("Failed to read base alias");
1366 let slab_bytes = std::fs::read(&first_slab).expect("Failed to read first slab");
1367 assert_eq!(alias_bytes, slab_bytes);
1368 }
1369
1370 #[test]
1371 fn test_one_section_self_cleaning() {
1372 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1373 let (logger, _) = make_a_logger(&tmp_dir, LARGE_SLAB);
1374 {
1375 let _stream = stream_write::<(), MmapSectionStorage>(
1376 logger.clone(),
1377 UnifiedLogType::StructuredLogLine,
1378 1024,
1379 );
1380 assert_eq!(
1381 logger
1382 .lock()
1383 .unwrap()
1384 .front_slab
1385 .sections_offsets_in_flight
1386 .len(),
1387 1
1388 );
1389 }
1390 assert_eq!(
1391 logger
1392 .lock()
1393 .unwrap()
1394 .front_slab
1395 .sections_offsets_in_flight
1396 .len(),
1397 0
1398 );
1399 let logger = logger.lock().unwrap();
1400 assert_eq!(
1401 logger.front_slab.flushed_until_offset,
1402 logger.front_slab.current_global_position
1403 );
1404 }
1405
1406 #[test]
1407 fn test_temporary_end_marker_is_created() {
1408 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1409 let (logger, _) = make_a_logger(&tmp_dir, LARGE_SLAB);
1410 {
1411 let mut stream = stream_write::<u32, MmapSectionStorage>(
1412 logger.clone(),
1413 UnifiedLogType::StructuredLogLine,
1414 1024,
1415 )
1416 .unwrap();
1417 stream.log(&42u32).unwrap();
1418 }
1419
1420 let logger_guard = logger.lock().unwrap();
1421 let slab = &logger_guard.front_slab;
1422 let marker_start = slab
1423 .temporary_end_marker
1424 .expect("temporary end-of-log marker missing");
1425 let (eof_header, _) =
1426 decode_from_slice::<SectionHeader, _>(&slab.mmap_buffer[marker_start..], standard())
1427 .expect("Could not decode end-of-log marker header");
1428 assert_eq!(eof_header.entry_type, UnifiedLogType::LastEntry);
1429 assert!(eof_header.is_open);
1430 assert_eq!(eof_header.used, 0);
1431 }
1432
1433 #[test]
1434 fn test_final_end_marker_is_not_temporary() {
1435 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1436 let (logger, f) = make_a_logger(&tmp_dir, LARGE_SLAB);
1437 {
1438 let mut stream = stream_write::<u32, MmapSectionStorage>(
1439 logger.clone(),
1440 UnifiedLogType::CopperList,
1441 1024,
1442 )
1443 .unwrap();
1444 stream.log(&1u32).unwrap();
1445 }
1446 drop(logger);
1447
1448 let MmapUnifiedLogger::Read(mut reader) = MmapUnifiedLoggerBuilder::new()
1449 .file_base_name(&f)
1450 .build()
1451 .expect("Failed to build reader")
1452 else {
1453 panic!("Failed to create reader");
1454 };
1455
1456 loop {
1457 let (header, _data) = reader
1458 .raw_read_section()
1459 .expect("Failed to read section while searching for EOF");
1460 if header.entry_type == UnifiedLogType::LastEntry {
1461 assert!(!header.is_open);
1462 break;
1463 }
1464 }
1465 }
1466
1467 #[test]
1468 fn test_two_sections_self_cleaning_in_order() {
1469 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1470 let (logger, _) = make_a_logger(&tmp_dir, LARGE_SLAB);
1471 let s1 = stream_write::<(), MmapSectionStorage>(
1472 logger.clone(),
1473 UnifiedLogType::StructuredLogLine,
1474 1024,
1475 );
1476 assert_eq!(
1477 logger
1478 .lock()
1479 .unwrap()
1480 .front_slab
1481 .sections_offsets_in_flight
1482 .len(),
1483 1
1484 );
1485 let s2 = stream_write::<(), MmapSectionStorage>(
1486 logger.clone(),
1487 UnifiedLogType::StructuredLogLine,
1488 1024,
1489 );
1490 assert_eq!(
1491 logger
1492 .lock()
1493 .unwrap()
1494 .front_slab
1495 .sections_offsets_in_flight
1496 .len(),
1497 2
1498 );
1499 drop(s2);
1500 assert_eq!(
1501 logger
1502 .lock()
1503 .unwrap()
1504 .front_slab
1505 .sections_offsets_in_flight
1506 .len(),
1507 1
1508 );
1509 drop(s1);
1510 let lg = logger.lock().unwrap();
1511 assert_eq!(lg.front_slab.sections_offsets_in_flight.len(), 0);
1512 assert_eq!(
1513 lg.front_slab.flushed_until_offset,
1514 lg.front_slab.current_global_position
1515 );
1516 }
1517
1518 #[test]
1519 fn test_two_sections_self_cleaning_out_of_order() {
1520 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1521 let (logger, _) = make_a_logger(&tmp_dir, LARGE_SLAB);
1522 let s1 = stream_write::<(), MmapSectionStorage>(
1523 logger.clone(),
1524 UnifiedLogType::StructuredLogLine,
1525 1024,
1526 );
1527 assert_eq!(
1528 logger
1529 .lock()
1530 .unwrap()
1531 .front_slab
1532 .sections_offsets_in_flight
1533 .len(),
1534 1
1535 );
1536 let s2 = stream_write::<(), MmapSectionStorage>(
1537 logger.clone(),
1538 UnifiedLogType::StructuredLogLine,
1539 1024,
1540 );
1541 assert_eq!(
1542 logger
1543 .lock()
1544 .unwrap()
1545 .front_slab
1546 .sections_offsets_in_flight
1547 .len(),
1548 2
1549 );
1550 drop(s1);
1551 assert_eq!(
1552 logger
1553 .lock()
1554 .unwrap()
1555 .front_slab
1556 .sections_offsets_in_flight
1557 .len(),
1558 1
1559 );
1560 drop(s2);
1561 let lg = logger.lock().unwrap();
1562 assert_eq!(lg.front_slab.sections_offsets_in_flight.len(), 0);
1563 assert_eq!(
1564 lg.front_slab.flushed_until_offset,
1565 lg.front_slab.current_global_position
1566 );
1567 }
1568
1569 #[test]
1570 fn test_closed_section_flushes_behind_open_earlier_section() {
1571 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1572 let (logger, _) = make_a_logger(&tmp_dir, LARGE_SLAB);
1573 let s1 = stream_write::<(), MmapSectionStorage>(
1574 logger.clone(),
1575 UnifiedLogType::StructuredLogLine,
1576 1024,
1577 )
1578 .unwrap();
1579 {
1580 let mut s2 = stream_write::<u32, MmapSectionStorage>(
1581 logger.clone(),
1582 UnifiedLogType::CopperList,
1583 1024,
1584 )
1585 .unwrap();
1586 s2.log(&42u32).unwrap();
1587 }
1588
1589 let logger_guard = logger.lock().unwrap();
1590 assert_eq!(logger_guard.front_slab.sections_offsets_in_flight.len(), 1);
1591 assert!(
1592 logger_guard.front_slab.flushed_until_offset
1593 < logger_guard.front_slab.current_global_position
1594 );
1595 assert_eq!(logger_guard.front_slab.pending_closed_bytes(), 0);
1596 drop(logger_guard);
1597 drop(s1);
1598 }
1599
1600 #[test]
1601 fn test_append_rejects_log_without_clean_close() {
1602 let tmp_dir =
1603 TempDir::new_in(env!("CARGO_MANIFEST_DIR")).expect("could not create a tmp dir");
1604 let (logger, file_path) = make_a_logger(&tmp_dir, LARGE_SLAB);
1605 {
1606 let mut stream = stream_write::<u32, MmapSectionStorage>(
1607 logger.clone(),
1608 UnifiedLogType::StructuredLogLine,
1609 1024,
1610 )
1611 .unwrap();
1612 stream.log(&1u32).unwrap();
1613 }
1614
1615 std::mem::forget(logger);
1618
1619 let MmapUnifiedLogger::Read(mut reader) = MmapUnifiedLoggerBuilder::new()
1620 .file_base_name(&file_path)
1621 .build()
1622 .expect("Failed to open logger for reading")
1623 else {
1624 panic!("Failed to open logger for reading")
1625 };
1626 let section = reader
1627 .read_next_section_type(UnifiedLogType::StructuredLogLine)
1628 .unwrap()
1629 .expect("Missing logged section");
1630 assert_eq!(
1631 decode_from_slice::<u32, _>(§ion, standard()).unwrap().0,
1632 1
1633 );
1634 let header = reader.read_section_header().unwrap();
1635 assert_eq!(header.entry_type, UnifiedLogType::LastEntry);
1636 assert!(header.is_open, "Expected a temporary end-of-log marker");
1637 drop(reader);
1638
1639 let result = MmapUnifiedLoggerBuilder::new()
1640 .write(true)
1641 .create(true)
1642 .append(true)
1643 .file_base_name(&file_path)
1644 .preallocated_size(LARGE_SLAB)
1645 .build();
1646
1647 assert!(
1648 result.is_err(),
1649 "Append must reject a log with a temporary end-of-log marker"
1650 );
1651 }
1652
1653 #[test]
1654 fn test_append_preserves_existing_and_adds_new_sections() {
1655 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1656 let file_path = tmp_dir.path().join("test.bin");
1657
1658 {
1660 let MmapUnifiedLogger::Write(logger) = MmapUnifiedLoggerBuilder::new()
1661 .write(true)
1662 .create(true)
1663 .file_base_name(&file_path)
1664 .preallocated_size(LARGE_SLAB)
1665 .build()
1666 .expect("Failed to create logger")
1667 else {
1668 panic!("Failed to create logger")
1669 };
1670 let logger = Arc::new(Mutex::new(logger));
1671 {
1672 let mut stream = stream_write::<u32, MmapSectionStorage>(
1673 logger.clone(),
1674 UnifiedLogType::StructuredLogLine,
1675 1024,
1676 )
1677 .unwrap();
1678 stream.log(&1u32).unwrap();
1679 stream.log(&2u32).unwrap();
1680 stream.log(&3u32).unwrap();
1681 }
1682 } {
1686 let MmapUnifiedLogger::Write(logger) = MmapUnifiedLoggerBuilder::new()
1687 .write(true)
1688 .create(true)
1689 .append(true)
1690 .file_base_name(&file_path)
1691 .preallocated_size(LARGE_SLAB)
1692 .build()
1693 .expect("Failed to append to logger")
1694 else {
1695 panic!("Failed to append to logger")
1696 };
1697 let logger = Arc::new(Mutex::new(logger));
1698 {
1699 let mut stream = stream_write::<u32, MmapSectionStorage>(
1700 logger.clone(),
1701 UnifiedLogType::StructuredLogLine,
1702 1024,
1703 )
1704 .unwrap();
1705 stream.log(&4u32).unwrap();
1706 stream.log(&5u32).unwrap();
1707 }
1708 }
1709
1710 let MmapUnifiedLogger::Read(mut dl) = MmapUnifiedLoggerBuilder::new()
1712 .file_base_name(&file_path)
1713 .build()
1714 .expect("Failed to build logger")
1715 else {
1716 panic!("Failed to build logger")
1717 };
1718
1719 let first = dl
1720 .read_next_section_type(UnifiedLogType::StructuredLogLine)
1721 .expect("Failed to read first section")
1722 .expect("Missing first section");
1723 let mut reader = SliceReader::new(&first[..]);
1724 assert_eq!(
1725 decode_from_reader::<u32, _, _>(&mut reader, standard()).unwrap(),
1726 1
1727 );
1728 assert_eq!(
1729 decode_from_reader::<u32, _, _>(&mut reader, standard()).unwrap(),
1730 2
1731 );
1732 assert_eq!(
1733 decode_from_reader::<u32, _, _>(&mut reader, standard()).unwrap(),
1734 3
1735 );
1736
1737 let second = dl
1738 .read_next_section_type(UnifiedLogType::StructuredLogLine)
1739 .expect("Failed to read second section")
1740 .expect("Missing second section");
1741 let mut reader = SliceReader::new(&second[..]);
1742 assert_eq!(
1743 decode_from_reader::<u32, _, _>(&mut reader, standard()).unwrap(),
1744 4
1745 );
1746 assert_eq!(
1747 decode_from_reader::<u32, _, _>(&mut reader, standard()).unwrap(),
1748 5
1749 );
1750
1751 assert!(
1752 dl.read_next_section_type(UnifiedLogType::StructuredLogLine)
1753 .expect("Failed to read past sections")
1754 .is_none()
1755 );
1756 }
1757
1758 #[test]
1759 fn test_write_then_read_one_section() {
1760 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1761 let (logger, f) = make_a_logger(&tmp_dir, LARGE_SLAB);
1762 {
1763 let mut stream =
1764 stream_write(logger.clone(), UnifiedLogType::StructuredLogLine, 1024).unwrap();
1765 stream.log(&1u32).unwrap();
1766 stream.log(&2u32).unwrap();
1767 stream.log(&3u32).unwrap();
1768 }
1769 drop(logger);
1770 let MmapUnifiedLogger::Read(mut dl) = MmapUnifiedLoggerBuilder::new()
1771 .file_base_name(&f)
1772 .build()
1773 .expect("Failed to build logger")
1774 else {
1775 panic!("Failed to build logger");
1776 };
1777 let section = dl
1778 .read_next_section_type(UnifiedLogType::StructuredLogLine)
1779 .expect("Failed to read section");
1780 assert!(section.is_some());
1781 let section = section.unwrap();
1782 let mut reader = SliceReader::new(§ion[..]);
1783 let v1: u32 = decode_from_reader(&mut reader, standard()).unwrap();
1784 let v2: u32 = decode_from_reader(&mut reader, standard()).unwrap();
1785 let v3: u32 = decode_from_reader(&mut reader, standard()).unwrap();
1786 assert_eq!(v1, 1);
1787 assert_eq!(v2, 2);
1788 assert_eq!(v3, 3);
1789 }
1790
1791 #[cfg(feature = "mmap-fsync")]
1792 #[test]
1793 fn test_fsync_feature_syncs_on_section_flush() {
1794 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1795 let (logger, _) = make_a_logger(&tmp_dir, LARGE_SLAB);
1796 {
1797 let mut stream =
1798 stream_write(logger.clone(), UnifiedLogType::StructuredLogLine, 1024).unwrap();
1799 stream.log(&1u32).unwrap();
1800 }
1801
1802 let logger = logger.lock().unwrap();
1803 assert!(
1804 logger.front_slab.sync_call_count > 0,
1805 "expected mmap-fsync to issue at least one sync_all call"
1806 );
1807 }
1808
1809 #[derive(Debug, Encode, Decode)]
1812 enum CopperListStateMock {
1813 Free,
1814 ProcessingTasks,
1815 BeingSerialized,
1816 }
1817
1818 #[derive(Encode, Decode)]
1819 struct CopperList<P: bincode::enc::Encode> {
1820 state: CopperListStateMock,
1821 payload: P, }
1823
1824 #[test]
1825 fn test_copperlist_list_like_logging() {
1826 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1827 let (logger, f) = make_a_logger(&tmp_dir, LARGE_SLAB);
1828 {
1829 let mut stream =
1830 stream_write(logger.clone(), UnifiedLogType::CopperList, 1024).unwrap();
1831 let cl0 = CopperList {
1832 state: CopperListStateMock::Free,
1833 payload: (1u32, 2u32, 3u32),
1834 };
1835 let cl1 = CopperList {
1836 state: CopperListStateMock::ProcessingTasks,
1837 payload: (4u32, 5u32, 6u32),
1838 };
1839 stream.log(&cl0).unwrap();
1840 stream.log(&cl1).unwrap();
1841 }
1842 drop(logger);
1843
1844 let MmapUnifiedLogger::Read(mut dl) = MmapUnifiedLoggerBuilder::new()
1845 .file_base_name(&f)
1846 .build()
1847 .expect("Failed to build logger")
1848 else {
1849 panic!("Failed to build logger");
1850 };
1851 let section = dl
1852 .read_next_section_type(UnifiedLogType::CopperList)
1853 .expect("Failed to read section");
1854 assert!(section.is_some());
1855 let section = section.unwrap();
1856
1857 let mut reader = SliceReader::new(§ion[..]);
1858 let cl0: CopperList<(u32, u32, u32)> = decode_from_reader(&mut reader, standard()).unwrap();
1859 let cl1: CopperList<(u32, u32, u32)> = decode_from_reader(&mut reader, standard()).unwrap();
1860 assert_eq!(cl0.payload.1, 2);
1861 assert_eq!(cl1.payload.2, 6);
1862 }
1863
1864 #[test]
1865 fn test_multi_slab_end2end() {
1866 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1867 let (logger, f) = make_a_logger(&tmp_dir, SMALL_SLAB);
1868 {
1869 let mut stream =
1870 stream_write(logger.clone(), UnifiedLogType::CopperList, 1024).unwrap();
1871 let cl0 = CopperList {
1872 state: CopperListStateMock::Free,
1873 payload: (1u32, 2u32, 3u32),
1874 };
1875 for _ in 0..10000 {
1877 stream.log(&cl0).unwrap();
1878 }
1879 }
1880 drop(logger);
1881
1882 let MmapUnifiedLogger::Read(mut dl) = MmapUnifiedLoggerBuilder::new()
1883 .file_base_name(&f)
1884 .build()
1885 .expect("Failed to build logger")
1886 else {
1887 panic!("Failed to build logger");
1888 };
1889 let mut total_readback = 0;
1890 loop {
1891 let section = dl.read_next_section_type(UnifiedLogType::CopperList);
1892 if section.is_err() {
1893 break;
1894 }
1895 let section = section.unwrap();
1896 if section.is_none() {
1897 break;
1898 }
1899 let section = section.unwrap();
1900
1901 let mut reader = SliceReader::new(§ion[..]);
1902 loop {
1903 let maybe_cl: Result<CopperList<(u32, u32, u32)>, _> =
1904 decode_from_reader(&mut reader, standard());
1905 if maybe_cl.is_ok() {
1906 total_readback += 1;
1907 } else {
1908 break;
1909 }
1910 }
1911 }
1912 assert_eq!(total_readback, 10000);
1913 }
1914}