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}
104
105impl Default for MmapUnifiedLoggerBuilder {
106 fn default() -> Self {
107 Self::new()
108 }
109}
110
111impl MmapUnifiedLoggerBuilder {
112 pub fn new() -> Self {
113 Self {
114 file_base_name: None,
115 preallocated_size: None,
116 write: false,
117 create: false, }
119 }
120
121 pub fn file_base_name(mut self, file_path: &Path) -> Self {
123 self.file_base_name = Some(file_path.to_path_buf());
124 self
125 }
126
127 pub fn preallocated_size(mut self, preallocated_size: usize) -> Self {
128 self.preallocated_size = Some(preallocated_size);
129 self
130 }
131
132 pub fn write(mut self, write: bool) -> Self {
133 self.write = write;
134 self
135 }
136
137 pub fn create(mut self, create: bool) -> Self {
138 self.create = create;
139 self
140 }
141
142 pub fn build(self) -> io::Result<MmapUnifiedLogger> {
143 let page_size = page_size::get();
144
145 if self.write && self.create {
146 let file_path = self.file_base_name.ok_or_else(|| {
147 io::Error::new(
148 io::ErrorKind::InvalidInput,
149 "File path is required for write mode",
150 )
151 })?;
152 let preallocated_size = self.preallocated_size.ok_or_else(|| {
153 io::Error::new(
154 io::ErrorKind::InvalidInput,
155 "Preallocated size is required for write mode",
156 )
157 })?;
158 let ulw = MmapUnifiedLoggerWrite::new(&file_path, preallocated_size, page_size)?;
159 Ok(MmapUnifiedLogger::Write(ulw))
160 } else {
161 let file_path = self.file_base_name.ok_or_else(|| {
162 io::Error::new(io::ErrorKind::InvalidInput, "File path is required")
163 })?;
164 let ulr = MmapUnifiedLoggerRead::new(&file_path)?;
165 Ok(MmapUnifiedLogger::Read(ulr))
166 }
167 }
168}
169
170struct SlabEntry {
171 file: File,
172 mmap_buffer: ManuallyDrop<MmapMut>,
173 current_global_position: usize,
174 sections_offsets_in_flight: Vec<usize>,
175 flushed_until_offset: usize,
176 page_size: usize,
177 temporary_end_marker: Option<usize>,
178 #[cfg(test)]
179 closed_sections: Vec<(usize, usize)>,
180 #[cfg(test)]
181 flushed_ranges: Vec<(usize, usize)>,
182 #[cfg(all(test, feature = "mmap-fsync"))]
183 sync_call_count: usize,
184}
185
186impl Drop for SlabEntry {
187 fn drop(&mut self) {
188 self.flush_until(self.current_global_position);
189 unsafe { ManuallyDrop::drop(&mut self.mmap_buffer) };
191 if let Err(error) = self.file.set_len(self.current_global_position as u64) {
192 eprintln!("Failed to trim datalogger file: {}", error);
193 }
194 self.sync_file();
195
196 if !self.sections_offsets_in_flight.is_empty() {
197 eprintln!("Error: Slab not full flushed.");
198 }
199 }
200}
201
202impl SlabEntry {
203 fn new(file: File, page_size: usize) -> io::Result<Self> {
204 let mmap_buffer = ManuallyDrop::new(
205 unsafe { MmapMut::map_mut(&file) }
207 .map_err(|e| io::Error::new(e.kind(), format!("Failed to map file: {e}")))?,
208 );
209 Ok(Self {
210 file,
211 mmap_buffer,
212 current_global_position: 0,
213 sections_offsets_in_flight: Vec::with_capacity(16),
214 flushed_until_offset: 0,
215 page_size,
216 temporary_end_marker: None,
217 #[cfg(test)]
218 closed_sections: Vec::new(),
219 #[cfg(test)]
220 flushed_ranges: Vec::new(),
221 #[cfg(all(test, feature = "mmap-fsync"))]
222 sync_call_count: 0,
223 })
224 }
225
226 fn flush_range(&mut self, start: usize, len: usize) {
227 if len == 0 {
228 return;
229 }
230 self.mmap_buffer
231 .flush_async_range(start, len)
232 .expect("Failed to flush memory map");
233 self.sync_file();
234 #[cfg(test)]
235 self.record_flushed_range(start, len);
236 }
237
238 fn sync_file(&mut self) {
239 #[cfg(feature = "mmap-fsync")]
240 {
241 self.file.sync_all().expect("Failed to fsync log file");
242 #[cfg(test)]
243 {
244 self.sync_call_count += 1;
245 }
246 }
247 }
248 fn flush_until(&mut self, until_position: usize) {
250 if (self.flushed_until_offset == until_position) || (until_position == 0) {
252 return;
253 }
254 self.flush_range(
255 self.flushed_until_offset,
256 until_position - self.flushed_until_offset,
257 );
258 self.flushed_until_offset = until_position;
259 }
260
261 fn clear_temporary_end_marker(&mut self) {
262 if let Some(marker_start) = self.temporary_end_marker.take() {
263 self.current_global_position = marker_start;
264 if self.flushed_until_offset > marker_start {
265 self.flushed_until_offset = marker_start;
266 }
267 }
268 }
269
270 fn write_end_marker(&mut self, temporary: bool) -> CuResult<()> {
271 let block_size = SECTION_HEADER_COMPACT_SIZE as usize;
272 let marker_start = self.align_to_next_page(self.current_global_position);
273 let total_marker_size = block_size; let marker_end = marker_start + total_marker_size;
275 if marker_end > self.mmap_buffer.len() {
276 return Err("Not enough space to write end-of-log marker".into());
277 }
278
279 let header = SectionHeader {
280 magic: SECTION_MAGIC,
281 block_size: SECTION_HEADER_COMPACT_SIZE,
282 entry_type: UnifiedLogType::LastEntry,
283 offset_to_next_section: total_marker_size as u32,
284 used: 0,
285 is_open: temporary,
286 };
287
288 encode_into_slice(
289 &header,
290 &mut self.mmap_buffer
291 [marker_start..marker_start + SECTION_HEADER_COMPACT_SIZE as usize],
292 standard(),
293 )
294 .map_err(|e| CuError::new_with_cause("Failed to encode end-of-log header", e))?;
295
296 self.temporary_end_marker = Some(marker_start);
297 self.current_global_position = marker_end;
298 Ok(())
299 }
300
301 fn is_it_my_section(&self, section: &SectionHandle<MmapSectionStorage>) -> bool {
302 let storage = section.get_storage();
303 let ptr = storage.buffer_ptr();
304 (ptr >= self.mmap_buffer.as_ptr())
305 && (ptr as usize)
306 < (self.mmap_buffer.as_ref().as_ptr() as usize + self.mmap_buffer.as_ref().len())
307 }
308
309 fn flush_section(&mut self, section: &mut SectionHandle<MmapSectionStorage>) {
312 section
313 .post_update_header()
314 .expect("Failed to update section header");
315
316 let storage = section.get_storage();
317 let ptr = storage.buffer_ptr();
318
319 if ptr < self.mmap_buffer.as_ptr()
320 || ptr as usize > self.mmap_buffer.as_ptr() as usize + self.mmap_buffer.len()
321 {
322 panic!("Invalid section buffer, not in the slab");
323 }
324
325 let base = self.mmap_buffer.as_ptr() as usize;
326 let section_start = ptr as usize - base;
327 let section_len = section.header.offset_to_next_section as usize;
328 #[cfg(test)]
329 self.record_closed_section(section_start, section_len);
330 self.sections_offsets_in_flight
331 .retain(|&x| x != section_start);
332
333 if self.sections_offsets_in_flight.is_empty() {
334 self.flush_until(self.current_global_position);
335 return;
336 }
337 let next_open_offset = self.sections_offsets_in_flight[0];
338 if self.flushed_until_offset < next_open_offset {
339 self.flush_until(next_open_offset);
340 }
341 if section_start + section_len > self.flushed_until_offset {
342 self.flush_range(section_start, section_len);
345 }
346 }
347
348 #[cfg(test)]
349 fn record_closed_section(&mut self, start: usize, len: usize) {
350 self.closed_sections.push((start, len));
351 }
352
353 #[cfg(test)]
354 fn record_flushed_range(&mut self, start: usize, len: usize) {
355 let mut merged_start = start;
356 let mut merged_end = start + len;
357 let mut merged_ranges = Vec::with_capacity(self.flushed_ranges.len() + 1);
358 let mut inserted = false;
359
360 for (range_start, range_len) in self.flushed_ranges.drain(..) {
361 let range_end = range_start + range_len;
362 if range_end < merged_start {
363 merged_ranges.push((range_start, range_len));
364 continue;
365 }
366 if merged_end < range_start {
367 if !inserted {
368 merged_ranges.push((merged_start, merged_end - merged_start));
369 inserted = true;
370 }
371 merged_ranges.push((range_start, range_len));
372 continue;
373 }
374
375 merged_start = merged_start.min(range_start);
376 merged_end = merged_end.max(range_end);
377 }
378
379 if !inserted {
380 merged_ranges.push((merged_start, merged_end - merged_start));
381 }
382
383 self.flushed_ranges = merged_ranges;
384 }
385
386 #[cfg(test)]
387 fn pending_closed_bytes(&self) -> usize {
388 let mut pending = 0;
389
390 for (section_start, section_len) in &self.closed_sections {
391 let section_end = section_start + section_len;
392 let mut cursor = *section_start;
393
394 for (range_start, range_len) in &self.flushed_ranges {
395 let range_end = range_start + range_len;
396 if range_end <= cursor {
397 continue;
398 }
399 if *range_start >= section_end {
400 break;
401 }
402 if *range_start > cursor {
403 pending += *range_start - cursor;
404 }
405 cursor = cursor.max(range_end);
406 if cursor >= section_end {
407 break;
408 }
409 }
410
411 if cursor < section_end {
412 pending += section_end - cursor;
413 }
414 }
415
416 pending
417 }
418
419 #[inline]
420 fn align_to_next_page(&self, ptr: usize) -> usize {
421 (ptr + self.page_size - 1) & !(self.page_size - 1)
422 }
423
424 fn add_section(
426 &mut self,
427 entry_type: UnifiedLogType,
428 requested_section_size: usize,
429 ) -> AllocatedSection<MmapSectionStorage> {
430 self.current_global_position = self.align_to_next_page(self.current_global_position);
432 let section_size = self.align_to_next_page(requested_section_size) as u32;
433
434 if self.current_global_position + section_size as usize > self.mmap_buffer.len() {
436 return AllocatedSection::NoMoreSpace;
437 }
438
439 #[cfg(feature = "compact")]
440 let block_size = SECTION_HEADER_COMPACT_SIZE;
441
442 #[cfg(not(feature = "compact"))]
443 let block_size = self.page_size as u16;
444
445 let section_header = SectionHeader {
446 magic: SECTION_MAGIC,
447 block_size,
448 entry_type,
449 offset_to_next_section: section_size,
450 used: 0u32,
451 is_open: true,
452 };
453
454 self.sections_offsets_in_flight
456 .push(self.current_global_position);
457 let end_of_section = self.current_global_position + requested_section_size;
458 let user_buffer = &mut self.mmap_buffer[self.current_global_position..end_of_section];
459
460 let handle_buffer =
462 unsafe { from_raw_parts_mut(user_buffer.as_mut_ptr(), user_buffer.len()) };
463 let storage = MmapSectionStorage::new(handle_buffer, block_size as usize);
464
465 self.current_global_position = end_of_section;
466
467 Section(SectionHandle::create(section_header, storage).expect("Failed to create section"))
468 }
469
470 #[cfg(test)]
471 fn used(&self) -> usize {
472 self.current_global_position
473 }
474}
475
476pub struct MmapUnifiedLoggerWrite {
478 front_slab: SlabEntry,
480 back_slabs: Vec<SlabEntry>,
482 base_file_path: PathBuf,
484 slab_size: usize,
486 front_slab_suffix: usize,
488}
489
490fn build_slab_path(base_file_path: &Path, slab_index: usize) -> io::Result<PathBuf> {
491 let mut file_path = base_file_path.to_path_buf();
492 let stem = file_path.file_stem().ok_or_else(|| {
493 io::Error::new(
494 io::ErrorKind::InvalidInput,
495 "Base file path has no file name",
496 )
497 })?;
498 let stem = stem.to_str().ok_or_else(|| {
499 io::Error::new(
500 io::ErrorKind::InvalidInput,
501 "Base file name is not valid UTF-8",
502 )
503 })?;
504 let extension = file_path.extension().ok_or_else(|| {
505 io::Error::new(
506 io::ErrorKind::InvalidInput,
507 "Base file path has no extension",
508 )
509 })?;
510 let extension = extension.to_str().ok_or_else(|| {
511 io::Error::new(
512 io::ErrorKind::InvalidInput,
513 "Base file extension is not valid UTF-8",
514 )
515 })?;
516 if stem.is_empty() {
517 return Err(io::Error::new(
518 io::ErrorKind::InvalidInput,
519 "Base file name is empty",
520 ));
521 }
522 let file_name = format!("{stem}_{slab_index}.{extension}");
523 file_path.set_file_name(file_name);
524 Ok(file_path)
525}
526
527fn make_slab_file(base_file_path: &Path, slab_size: usize, slab_suffix: usize) -> io::Result<File> {
528 let file_path = build_slab_path(base_file_path, slab_suffix)?;
529 let file = OpenOptions::new()
530 .read(true)
531 .write(true)
532 .create(true)
533 .truncate(true)
534 .open(&file_path)
535 .map_err(|e| {
536 io::Error::new(
537 e.kind(),
538 format!("Failed to open file {}: {e}", file_path.display()),
539 )
540 })?;
541 file.set_len(slab_size as u64).map_err(|e| {
542 io::Error::new(
543 e.kind(),
544 format!("Failed to set file length for {}: {e}", file_path.display()),
545 )
546 })?;
547 Ok(file)
548}
549
550fn remove_existing_alias(base_file_path: &Path) -> io::Result<()> {
551 match std::fs::symlink_metadata(base_file_path) {
552 Ok(meta) => {
553 if meta.is_dir() {
554 return Err(io::Error::new(
555 io::ErrorKind::AlreadyExists,
556 format!(
557 "Cannot create base log alias at {} because a directory already exists there",
558 base_file_path.display()
559 ),
560 ));
561 }
562 std::fs::remove_file(base_file_path).map_err(|e| {
563 io::Error::new(
564 e.kind(),
565 format!(
566 "Failed to remove existing base log alias {}: {e}",
567 base_file_path.display()
568 ),
569 )
570 })
571 }
572 Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(()),
573 Err(e) => Err(io::Error::new(
574 e.kind(),
575 format!(
576 "Failed to inspect existing base log alias {}: {e}",
577 base_file_path.display()
578 ),
579 )),
580 }
581}
582
583fn create_base_alias_link(base_file_path: &Path) -> io::Result<()> {
584 let first_slab_path = build_slab_path(base_file_path, 0)?;
585 remove_existing_alias(base_file_path)?;
586
587 #[cfg(unix)]
588 {
589 use std::os::unix::fs::symlink;
590 let relative_target = Path::new(first_slab_path.file_name().ok_or_else(|| {
591 io::Error::new(
592 io::ErrorKind::InvalidInput,
593 "First slab file has no name component",
594 )
595 })?);
596 symlink(relative_target, base_file_path).map_err(|e| {
597 io::Error::new(
598 e.kind(),
599 format!(
600 "Failed to create base log alias {} -> {}: {e}",
601 base_file_path.display(),
602 first_slab_path.display()
603 ),
604 )
605 })
606 }
607
608 #[cfg(windows)]
609 {
610 use std::os::windows::fs::symlink_file;
611 let relative_target = Path::new(first_slab_path.file_name().ok_or_else(|| {
612 io::Error::new(
613 io::ErrorKind::InvalidInput,
614 "First slab file has no name component",
615 )
616 })?);
617 match symlink_file(relative_target, base_file_path) {
618 Ok(()) => Ok(()),
619 Err(symlink_err) => std::fs::hard_link(&first_slab_path, base_file_path).map_err(
620 |hard_link_err| {
621 io::Error::other(format!(
622 "Failed to create base log alias {}. Symlink error: {symlink_err}. Hard-link fallback error: {hard_link_err}",
623 base_file_path.display()
624 ))
625 },
626 ),
627 }?;
628 Ok(())
629 }
630
631 #[cfg(not(any(unix, windows)))]
632 {
633 std::fs::hard_link(&first_slab_path, base_file_path).map_err(|e| {
634 io::Error::new(
635 e.kind(),
636 format!(
637 "Failed to create base log alias {} -> {}: {e}",
638 base_file_path.display(),
639 first_slab_path.display()
640 ),
641 )
642 })
643 }
644}
645
646impl UnifiedLogWrite<MmapSectionStorage> for MmapUnifiedLoggerWrite {
647 fn add_section(
649 &mut self,
650 entry_type: UnifiedLogType,
651 requested_section_size: usize,
652 ) -> CuResult<SectionHandle<MmapSectionStorage>> {
653 self.garbage_collect_backslabs(); self.front_slab.clear_temporary_end_marker();
655 let maybe_section = self
656 .front_slab
657 .add_section(entry_type, requested_section_size);
658
659 match maybe_section {
660 AllocatedSection::NoMoreSpace => {
661 let new_slab = self.create_slab()?;
663 self.back_slabs
665 .push(mem::replace(&mut self.front_slab, new_slab));
666 match self
667 .front_slab
668 .add_section(entry_type, requested_section_size)
669 {
670 AllocatedSection::NoMoreSpace => Err(CuError::from("out of space")),
671 Section(section) => {
672 self.place_end_marker(true)?;
673 Ok(section)
674 }
675 }
676 }
677 Section(section) => {
678 self.place_end_marker(true)?;
679 Ok(section)
680 }
681 }
682 }
683
684 fn flush_section(&mut self, section: &mut SectionHandle<MmapSectionStorage>) {
685 section.mark_closed();
686 for slab in self.back_slabs.iter_mut() {
687 if slab.is_it_my_section(section) {
688 slab.flush_section(section);
689 return;
690 }
691 }
692 self.front_slab.flush_section(section);
693 }
694
695 fn status(&self) -> UnifiedLogStatus {
696 UnifiedLogStatus {
697 total_used_space: self.front_slab.current_global_position,
698 total_allocated_space: self.slab_size * self.front_slab_suffix,
699 }
700 }
701}
702
703impl MmapUnifiedLoggerWrite {
704 fn next_slab(&mut self) -> io::Result<File> {
705 let next_suffix = self.front_slab_suffix + 1;
706 let file = make_slab_file(&self.base_file_path, self.slab_size, next_suffix)?;
707 self.front_slab_suffix = next_suffix;
708 Ok(file)
709 }
710
711 fn new(base_file_path: &Path, slab_size: usize, page_size: usize) -> io::Result<Self> {
712 let file = make_slab_file(base_file_path, slab_size, 0)?;
713 create_base_alias_link(base_file_path)?;
714 let mut front_slab = SlabEntry::new(file, page_size)?;
715
716 let main_header = MainHeader {
718 magic: MAIN_MAGIC,
719 format_version: UNIFIED_LOG_FORMAT_VERSION,
720 first_section_offset: page_size as u16,
721 page_size: page_size as u16,
722 };
723 let nb_bytes = encode_into_slice(&main_header, &mut front_slab.mmap_buffer[..], standard())
724 .map_err(|e| io::Error::other(format!("Failed to encode main header: {e}")))?;
725 assert!(nb_bytes < page_size);
726 front_slab.current_global_position = page_size; Ok(Self {
729 front_slab,
730 back_slabs: Vec::new(),
731 base_file_path: base_file_path.to_path_buf(),
732 slab_size,
733 front_slab_suffix: 0,
734 })
735 }
736
737 fn garbage_collect_backslabs(&mut self) {
738 self.back_slabs
739 .retain_mut(|slab| !slab.sections_offsets_in_flight.is_empty());
740 }
741
742 fn place_end_marker(&mut self, temporary: bool) -> CuResult<()> {
743 match self.front_slab.write_end_marker(temporary) {
744 Ok(_) => Ok(()),
745 Err(_) => {
746 let new_slab = self.create_slab()?;
748 self.back_slabs
749 .push(mem::replace(&mut self.front_slab, new_slab));
750 self.front_slab.write_end_marker(temporary)
751 }
752 }
753 }
754
755 pub fn stats(&self) -> (usize, Vec<usize>, usize) {
756 (
757 self.front_slab.current_global_position,
758 self.front_slab.sections_offsets_in_flight.clone(),
759 self.back_slabs.len(),
760 )
761 }
762
763 fn create_slab(&mut self) -> CuResult<SlabEntry> {
764 let file = self
765 .next_slab()
766 .map_err(|e| CuError::new_with_cause("Failed to create slab file", e))?;
767 SlabEntry::new(file, self.front_slab.page_size)
768 .map_err(|e| CuError::new_with_cause("Failed to create slab memory map", e))
769 }
770}
771
772impl Drop for MmapUnifiedLoggerWrite {
773 fn drop(&mut self) {
774 #[cfg(debug_assertions)]
775 eprintln!("Flushing the unified Logger ... "); self.front_slab.clear_temporary_end_marker();
778 if let Err(e) = self.place_end_marker(false) {
779 panic!("Failed to flush the unified logger: {}", e);
780 }
781 self.front_slab
782 .flush_until(self.front_slab.current_global_position);
783 self.garbage_collect_backslabs();
784 #[cfg(debug_assertions)]
785 eprintln!("Unified Logger flushed."); }
787}
788
789fn open_slab_index(
790 base_file_path: &Path,
791 slab_index: usize,
792) -> io::Result<(File, Mmap, u16, Option<MainHeader>)> {
793 let mut options = OpenOptions::new();
794 let options = options.read(true);
795
796 let file_path = build_slab_path(base_file_path, slab_index)?;
797 let file = options.open(&file_path).map_err(|e| {
798 io::Error::new(
799 e.kind(),
800 format!("Failed to open slab file {}: {e}", file_path.display()),
801 )
802 })?;
803 let mmap = unsafe { Mmap::map(&file) }
805 .map_err(|e| io::Error::new(e.kind(), format!("Failed to map slab file: {e}")))?;
806 let mut prolog = 0u16;
807 let mut maybe_main_header: Option<MainHeader> = None;
808 if slab_index == 0 {
809 let main_header: MainHeader;
810 let _read: usize;
811 (main_header, _read) = decode_from_slice(&mmap[..], standard()).map_err(|e| {
812 io::Error::new(
813 io::ErrorKind::InvalidData,
814 format!("Failed to decode main header: {e}"),
815 )
816 })?;
817 if main_header.magic != MAIN_MAGIC {
818 return Err(io::Error::new(
819 io::ErrorKind::InvalidData,
820 "Invalid magic number in main header",
821 ));
822 }
823 if main_header.format_version != UNIFIED_LOG_FORMAT_VERSION {
824 return Err(io::Error::new(
825 io::ErrorKind::InvalidData,
826 format!(
827 "Unsupported unified log format version {} in main header; this reader supports version {}",
828 main_header.format_version, UNIFIED_LOG_FORMAT_VERSION
829 ),
830 ));
831 }
832 prolog = main_header.first_section_offset;
833 maybe_main_header = Some(main_header);
834 }
835 Ok((file, mmap, prolog, maybe_main_header))
836}
837
838pub struct MmapUnifiedLoggerRead {
840 base_file_path: PathBuf,
841 main_header: MainHeader,
842 current_mmap_buffer: Mmap,
843 current_file: File,
844 current_slab_index: usize,
845 current_reading_position: usize,
846}
847
848#[derive(Clone, Copy, Debug, PartialEq, Eq)]
850pub struct LogPosition {
851 pub slab_index: usize,
852 pub offset: usize,
853}
854
855impl UnifiedLogRead for MmapUnifiedLoggerRead {
856 fn read_next_section_type(&mut self, datalogtype: UnifiedLogType) -> CuResult<Option<Vec<u8>>> {
857 loop {
859 if self.current_reading_position >= self.current_mmap_buffer.len() {
860 self.next_slab().map_err(|e| {
861 CuError::new_with_cause("Failed to read next slab, is the log complete?", e)
862 })?;
863 }
864
865 let header_result = self.read_section_header();
866 let header = header_result.map_err(|error| {
867 CuError::new_with_cause(
868 &format!(
869 "Could not read a sections header: {}/{}:{}",
870 self.base_file_path.as_os_str().to_string_lossy(),
871 self.current_slab_index,
872 self.current_reading_position,
873 ),
874 error,
875 )
876 })?;
877
878 if header.entry_type == UnifiedLogType::LastEntry {
880 return Ok(None);
881 }
882
883 if header.entry_type == datalogtype {
885 let result = Some(self.read_section_content(&header)?);
886 self.current_reading_position += header.offset_to_next_section as usize;
887 return Ok(result);
888 }
889
890 self.current_reading_position += header.offset_to_next_section as usize;
892 }
893 }
894
895 fn raw_read_section(&mut self) -> CuResult<(SectionHeader, Vec<u8>)> {
897 if self.current_reading_position >= self.current_mmap_buffer.len() {
898 self.next_slab().map_err(|e| {
899 CuError::new_with_cause("Failed to read next slab, is the log complete?", e)
900 })?;
901 }
902
903 let read_result = self.read_section_header();
904
905 match read_result {
906 Err(error) => Err(CuError::new_with_cause(
907 &format!(
908 "Could not read a sections header: {}/{}:{}",
909 self.base_file_path.as_os_str().to_string_lossy(),
910 self.current_slab_index,
911 self.current_reading_position,
912 ),
913 error,
914 )),
915 Ok(header) => {
916 let data = self.read_section_content(&header)?;
917 self.current_reading_position += header.offset_to_next_section as usize;
918 Ok((header, data))
919 }
920 }
921 }
922}
923
924impl MmapUnifiedLoggerRead {
925 pub fn raw_skip_section(&mut self) -> CuResult<SectionHeader> {
927 if self.current_reading_position >= self.current_mmap_buffer.len() {
928 self.next_slab().map_err(|e| {
929 CuError::new_with_cause("Failed to read next slab, is the log complete?", e)
930 })?;
931 }
932
933 let header = self.read_section_header().map_err(|error| {
934 CuError::new_with_cause(
935 &format!(
936 "Could not read a sections header: {}/{}:{}",
937 self.base_file_path.as_os_str().to_string_lossy(),
938 self.current_slab_index,
939 self.current_reading_position,
940 ),
941 error,
942 )
943 })?;
944 self.current_reading_position += header.offset_to_next_section as usize;
945 Ok(header)
946 }
947
948 pub fn new(base_file_path: &Path) -> io::Result<Self> {
949 let (file, mmap, prolog, header) = open_slab_index(base_file_path, 0)?;
950 let main_header = header.ok_or_else(|| {
951 io::Error::new(io::ErrorKind::InvalidData, "Missing main header in slab 0")
952 })?;
953
954 Ok(Self {
955 base_file_path: base_file_path.to_path_buf(),
956 main_header,
957 current_file: file,
958 current_mmap_buffer: mmap,
959 current_slab_index: 0,
960 current_reading_position: prolog as usize,
961 })
962 }
963
964 pub fn position(&self) -> LogPosition {
966 LogPosition {
967 slab_index: self.current_slab_index,
968 offset: self.current_reading_position,
969 }
970 }
971
972 pub fn seek(&mut self, pos: LogPosition) -> CuResult<()> {
974 if pos.slab_index != self.current_slab_index {
975 let (file, mmap, _prolog, _header) =
976 open_slab_index(&self.base_file_path, pos.slab_index).map_err(|e| {
977 CuError::new_with_cause(
978 &format!("Failed to open slab {} for seek", pos.slab_index),
979 e,
980 )
981 })?;
982 self.current_file = file;
983 self.current_mmap_buffer = mmap;
984 self.current_slab_index = pos.slab_index;
985 }
986 self.current_reading_position = pos.offset;
987 Ok(())
988 }
989
990 fn next_slab(&mut self) -> io::Result<()> {
991 self.current_slab_index += 1;
992 let (file, mmap, prolog, _) =
993 open_slab_index(&self.base_file_path, self.current_slab_index)?;
994 self.current_file = file;
995 self.current_mmap_buffer = mmap;
996 self.current_reading_position = prolog as usize;
997 Ok(())
998 }
999
1000 pub fn raw_main_header(&self) -> &MainHeader {
1001 &self.main_header
1002 }
1003
1004 pub fn scan_section_bytes(&mut self, datalogtype: UnifiedLogType) -> CuResult<u64> {
1005 let mut total = 0u64;
1006
1007 loop {
1008 if self.current_reading_position >= self.current_mmap_buffer.len() {
1009 self.next_slab().map_err(|e| {
1010 CuError::new_with_cause("Failed to read next slab, is the log complete?", e)
1011 })?;
1012 }
1013
1014 let header = self.read_section_header()?;
1015
1016 if header.entry_type == UnifiedLogType::LastEntry {
1017 return Ok(total);
1018 }
1019
1020 if header.entry_type == datalogtype {
1021 total = total.saturating_add(header.used as u64);
1022 }
1023
1024 self.current_reading_position += header.offset_to_next_section as usize;
1025 }
1026 }
1027
1028 fn read_section_content(&mut self, header: &SectionHeader) -> CuResult<Vec<u8>> {
1030 let mut section_data = vec![0; header.used as usize];
1032 let start_of_data = self.current_reading_position + header.block_size as usize;
1033 section_data.copy_from_slice(
1034 &self.current_mmap_buffer[start_of_data..start_of_data + header.used as usize],
1035 );
1036
1037 Ok(section_data)
1038 }
1039
1040 fn read_section_header(&mut self) -> CuResult<SectionHeader> {
1041 let section_header: SectionHeader;
1042 (section_header, _) = decode_from_slice(
1043 &self.current_mmap_buffer[self.current_reading_position..],
1044 standard(),
1045 )
1046 .map_err(|e| {
1047 CuError::new_with_cause(
1048 &format!(
1049 "Could not read a sections header: {}/{}:{}",
1050 self.base_file_path.as_os_str().to_string_lossy(),
1051 self.current_slab_index,
1052 self.current_reading_position,
1053 ),
1054 e,
1055 )
1056 })?;
1057 if section_header.magic != SECTION_MAGIC {
1058 return Err("Invalid magic number in section header".into());
1059 }
1060
1061 Ok(section_header)
1062 }
1063}
1064
1065pub struct UnifiedLoggerIOReader {
1067 logger: MmapUnifiedLoggerRead,
1068 log_type: UnifiedLogType,
1069 buffer: Vec<u8>,
1070 buffer_pos: usize,
1071}
1072
1073impl UnifiedLoggerIOReader {
1074 pub fn new(logger: MmapUnifiedLoggerRead, log_type: UnifiedLogType) -> Self {
1075 Self {
1076 logger,
1077 log_type,
1078 buffer: Vec::new(),
1079 buffer_pos: 0,
1080 }
1081 }
1082
1083 fn fill_buffer(&mut self) -> io::Result<bool> {
1085 match self.logger.read_next_section_type(self.log_type) {
1086 Ok(Some(section)) => {
1087 self.buffer = section;
1088 self.buffer_pos = 0;
1089 Ok(true)
1090 }
1091 Ok(None) => Ok(false), Err(e) => Err(io::Error::other(e.to_string())),
1093 }
1094 }
1095}
1096
1097impl Read for UnifiedLoggerIOReader {
1098 fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
1099 if self.buffer_pos >= self.buffer.len() && !self.fill_buffer()? {
1100 return Ok(0);
1102 }
1103
1104 if self.buffer_pos >= self.buffer.len() {
1106 return Ok(0);
1107 }
1108
1109 let len = std::cmp::min(buf.len(), self.buffer.len() - self.buffer_pos);
1111 buf[..len].copy_from_slice(&self.buffer[self.buffer_pos..self.buffer_pos + len]);
1112 self.buffer_pos += len;
1113 Ok(len)
1114 }
1115}
1116
1117#[cfg(feature = "std")]
1118#[cfg(test)]
1119mod tests {
1120 use super::*;
1121 use crate::stream_write;
1122 use bincode::de::read::SliceReader;
1123 use bincode::{Decode, Encode, decode_from_reader, decode_from_slice};
1124 use cu29_traits::WriteStream;
1125 use std::io::{Seek, SeekFrom, Write};
1126 use std::path::PathBuf;
1127 use std::sync::{Arc, Mutex};
1128 use tempfile::TempDir;
1129
1130 const LARGE_SLAB: usize = 100 * 1024; const SMALL_SLAB: usize = 16 * 2 * 1024; fn make_a_logger(
1134 tmp_dir: &TempDir,
1135 slab_size: usize,
1136 ) -> (Arc<Mutex<MmapUnifiedLoggerWrite>>, PathBuf) {
1137 let file_path = tmp_dir.path().join("test.bin");
1138 let MmapUnifiedLogger::Write(data_logger) = MmapUnifiedLoggerBuilder::new()
1139 .write(true)
1140 .create(true)
1141 .file_base_name(&file_path)
1142 .preallocated_size(slab_size)
1143 .build()
1144 .expect("Failed to create logger")
1145 else {
1146 panic!("Failed to create logger")
1147 };
1148
1149 (Arc::new(Mutex::new(data_logger)), file_path)
1150 }
1151
1152 #[test]
1153 fn test_truncation_and_sections_creations() {
1154 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1155 let file_path = tmp_dir.path().join("test.bin");
1156 let _used = {
1157 let MmapUnifiedLogger::Write(mut logger) = MmapUnifiedLoggerBuilder::new()
1158 .write(true)
1159 .create(true)
1160 .file_base_name(&file_path)
1161 .preallocated_size(100000)
1162 .build()
1163 .expect("Failed to create logger")
1164 else {
1165 panic!("Failed to create logger")
1166 };
1167 logger
1168 .add_section(UnifiedLogType::StructuredLogLine, 1024)
1169 .unwrap();
1170 logger
1171 .add_section(UnifiedLogType::CopperList, 2048)
1172 .unwrap();
1173 let used = logger.front_slab.used();
1174 assert!(used < 4 * page_size::get()); used
1178 };
1179
1180 let _file = OpenOptions::new()
1181 .read(true)
1182 .open(tmp_dir.path().join("test_0.bin"))
1183 .expect("Could not reopen the file");
1184 }
1191
1192 #[test]
1193 fn test_unsupported_main_header_format_version_is_rejected() {
1194 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1195 let file_path = tmp_dir.path().join("test.bin");
1196 {
1197 let MmapUnifiedLogger::Write(_logger) = MmapUnifiedLoggerBuilder::new()
1198 .write(true)
1199 .create(true)
1200 .file_base_name(&file_path)
1201 .preallocated_size(100000)
1202 .build()
1203 .expect("Failed to create logger")
1204 else {
1205 panic!("Failed to create logger")
1206 };
1207 }
1208
1209 let mut file = OpenOptions::new()
1210 .read(true)
1211 .write(true)
1212 .open(tmp_dir.path().join("test_0.bin"))
1213 .expect("Could not reopen the slab");
1214 let unsupported_version = UNIFIED_LOG_FORMAT_VERSION + 1;
1215 file.seek(SeekFrom::Start(MAIN_MAGIC.len() as u64))
1216 .expect("Could not seek to format version");
1217 file.write_all(&[unsupported_version])
1218 .expect("Could not write unsupported format version");
1219 drop(file);
1220
1221 let err = match MmapUnifiedLoggerBuilder::new()
1222 .file_base_name(&file_path)
1223 .build()
1224 {
1225 Ok(_) => panic!("Reader accepted unsupported unified log format version"),
1226 Err(err) => err,
1227 };
1228
1229 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1230 assert_eq!(
1231 err.to_string(),
1232 format!(
1233 "Unsupported unified log format version {unsupported_version} in main header; this reader supports version {UNIFIED_LOG_FORMAT_VERSION}"
1234 )
1235 );
1236 }
1237
1238 #[test]
1239 fn test_base_alias_exists_and_matches_first_slab() {
1240 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1241 let file_path = tmp_dir.path().join("test.bin");
1242 let _logger = MmapUnifiedLoggerBuilder::new()
1243 .write(true)
1244 .create(true)
1245 .file_base_name(&file_path)
1246 .preallocated_size(LARGE_SLAB)
1247 .build()
1248 .expect("Failed to create logger");
1249
1250 let first_slab = build_slab_path(&file_path, 0).expect("Failed to build first slab path");
1251 assert!(file_path.exists(), "base alias does not exist");
1252 assert!(first_slab.exists(), "first slab does not exist");
1253
1254 let alias_bytes = std::fs::read(&file_path).expect("Failed to read base alias");
1255 let slab_bytes = std::fs::read(&first_slab).expect("Failed to read first slab");
1256 assert_eq!(alias_bytes, slab_bytes);
1257 }
1258
1259 #[test]
1260 fn test_one_section_self_cleaning() {
1261 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1262 let (logger, _) = make_a_logger(&tmp_dir, LARGE_SLAB);
1263 {
1264 let _stream = stream_write::<(), MmapSectionStorage>(
1265 logger.clone(),
1266 UnifiedLogType::StructuredLogLine,
1267 1024,
1268 );
1269 assert_eq!(
1270 logger
1271 .lock()
1272 .unwrap()
1273 .front_slab
1274 .sections_offsets_in_flight
1275 .len(),
1276 1
1277 );
1278 }
1279 assert_eq!(
1280 logger
1281 .lock()
1282 .unwrap()
1283 .front_slab
1284 .sections_offsets_in_flight
1285 .len(),
1286 0
1287 );
1288 let logger = logger.lock().unwrap();
1289 assert_eq!(
1290 logger.front_slab.flushed_until_offset,
1291 logger.front_slab.current_global_position
1292 );
1293 }
1294
1295 #[test]
1296 fn test_temporary_end_marker_is_created() {
1297 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1298 let (logger, _) = make_a_logger(&tmp_dir, LARGE_SLAB);
1299 {
1300 let mut stream = stream_write::<u32, MmapSectionStorage>(
1301 logger.clone(),
1302 UnifiedLogType::StructuredLogLine,
1303 1024,
1304 )
1305 .unwrap();
1306 stream.log(&42u32).unwrap();
1307 }
1308
1309 let logger_guard = logger.lock().unwrap();
1310 let slab = &logger_guard.front_slab;
1311 let marker_start = slab
1312 .temporary_end_marker
1313 .expect("temporary end-of-log marker missing");
1314 let (eof_header, _) =
1315 decode_from_slice::<SectionHeader, _>(&slab.mmap_buffer[marker_start..], standard())
1316 .expect("Could not decode end-of-log marker header");
1317 assert_eq!(eof_header.entry_type, UnifiedLogType::LastEntry);
1318 assert!(eof_header.is_open);
1319 assert_eq!(eof_header.used, 0);
1320 }
1321
1322 #[test]
1323 fn test_final_end_marker_is_not_temporary() {
1324 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1325 let (logger, f) = make_a_logger(&tmp_dir, LARGE_SLAB);
1326 {
1327 let mut stream = stream_write::<u32, MmapSectionStorage>(
1328 logger.clone(),
1329 UnifiedLogType::CopperList,
1330 1024,
1331 )
1332 .unwrap();
1333 stream.log(&1u32).unwrap();
1334 }
1335 drop(logger);
1336
1337 let MmapUnifiedLogger::Read(mut reader) = MmapUnifiedLoggerBuilder::new()
1338 .file_base_name(&f)
1339 .build()
1340 .expect("Failed to build reader")
1341 else {
1342 panic!("Failed to create reader");
1343 };
1344
1345 loop {
1346 let (header, _data) = reader
1347 .raw_read_section()
1348 .expect("Failed to read section while searching for EOF");
1349 if header.entry_type == UnifiedLogType::LastEntry {
1350 assert!(!header.is_open);
1351 break;
1352 }
1353 }
1354 }
1355
1356 #[test]
1357 fn test_two_sections_self_cleaning_in_order() {
1358 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1359 let (logger, _) = make_a_logger(&tmp_dir, LARGE_SLAB);
1360 let s1 = stream_write::<(), MmapSectionStorage>(
1361 logger.clone(),
1362 UnifiedLogType::StructuredLogLine,
1363 1024,
1364 );
1365 assert_eq!(
1366 logger
1367 .lock()
1368 .unwrap()
1369 .front_slab
1370 .sections_offsets_in_flight
1371 .len(),
1372 1
1373 );
1374 let s2 = stream_write::<(), MmapSectionStorage>(
1375 logger.clone(),
1376 UnifiedLogType::StructuredLogLine,
1377 1024,
1378 );
1379 assert_eq!(
1380 logger
1381 .lock()
1382 .unwrap()
1383 .front_slab
1384 .sections_offsets_in_flight
1385 .len(),
1386 2
1387 );
1388 drop(s2);
1389 assert_eq!(
1390 logger
1391 .lock()
1392 .unwrap()
1393 .front_slab
1394 .sections_offsets_in_flight
1395 .len(),
1396 1
1397 );
1398 drop(s1);
1399 let lg = logger.lock().unwrap();
1400 assert_eq!(lg.front_slab.sections_offsets_in_flight.len(), 0);
1401 assert_eq!(
1402 lg.front_slab.flushed_until_offset,
1403 lg.front_slab.current_global_position
1404 );
1405 }
1406
1407 #[test]
1408 fn test_two_sections_self_cleaning_out_of_order() {
1409 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1410 let (logger, _) = make_a_logger(&tmp_dir, LARGE_SLAB);
1411 let s1 = stream_write::<(), MmapSectionStorage>(
1412 logger.clone(),
1413 UnifiedLogType::StructuredLogLine,
1414 1024,
1415 );
1416 assert_eq!(
1417 logger
1418 .lock()
1419 .unwrap()
1420 .front_slab
1421 .sections_offsets_in_flight
1422 .len(),
1423 1
1424 );
1425 let s2 = stream_write::<(), MmapSectionStorage>(
1426 logger.clone(),
1427 UnifiedLogType::StructuredLogLine,
1428 1024,
1429 );
1430 assert_eq!(
1431 logger
1432 .lock()
1433 .unwrap()
1434 .front_slab
1435 .sections_offsets_in_flight
1436 .len(),
1437 2
1438 );
1439 drop(s1);
1440 assert_eq!(
1441 logger
1442 .lock()
1443 .unwrap()
1444 .front_slab
1445 .sections_offsets_in_flight
1446 .len(),
1447 1
1448 );
1449 drop(s2);
1450 let lg = logger.lock().unwrap();
1451 assert_eq!(lg.front_slab.sections_offsets_in_flight.len(), 0);
1452 assert_eq!(
1453 lg.front_slab.flushed_until_offset,
1454 lg.front_slab.current_global_position
1455 );
1456 }
1457
1458 #[test]
1459 fn test_closed_section_flushes_behind_open_earlier_section() {
1460 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1461 let (logger, _) = make_a_logger(&tmp_dir, LARGE_SLAB);
1462 let s1 = stream_write::<(), MmapSectionStorage>(
1463 logger.clone(),
1464 UnifiedLogType::StructuredLogLine,
1465 1024,
1466 )
1467 .unwrap();
1468 {
1469 let mut s2 = stream_write::<u32, MmapSectionStorage>(
1470 logger.clone(),
1471 UnifiedLogType::CopperList,
1472 1024,
1473 )
1474 .unwrap();
1475 s2.log(&42u32).unwrap();
1476 }
1477
1478 let logger_guard = logger.lock().unwrap();
1479 assert_eq!(logger_guard.front_slab.sections_offsets_in_flight.len(), 1);
1480 assert!(
1481 logger_guard.front_slab.flushed_until_offset
1482 < logger_guard.front_slab.current_global_position
1483 );
1484 assert_eq!(logger_guard.front_slab.pending_closed_bytes(), 0);
1485 drop(logger_guard);
1486 drop(s1);
1487 }
1488
1489 #[test]
1490 fn test_write_then_read_one_section() {
1491 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1492 let (logger, f) = make_a_logger(&tmp_dir, LARGE_SLAB);
1493 {
1494 let mut stream =
1495 stream_write(logger.clone(), UnifiedLogType::StructuredLogLine, 1024).unwrap();
1496 stream.log(&1u32).unwrap();
1497 stream.log(&2u32).unwrap();
1498 stream.log(&3u32).unwrap();
1499 }
1500 drop(logger);
1501 let MmapUnifiedLogger::Read(mut dl) = MmapUnifiedLoggerBuilder::new()
1502 .file_base_name(&f)
1503 .build()
1504 .expect("Failed to build logger")
1505 else {
1506 panic!("Failed to build logger");
1507 };
1508 let section = dl
1509 .read_next_section_type(UnifiedLogType::StructuredLogLine)
1510 .expect("Failed to read section");
1511 assert!(section.is_some());
1512 let section = section.unwrap();
1513 let mut reader = SliceReader::new(§ion[..]);
1514 let v1: u32 = decode_from_reader(&mut reader, standard()).unwrap();
1515 let v2: u32 = decode_from_reader(&mut reader, standard()).unwrap();
1516 let v3: u32 = decode_from_reader(&mut reader, standard()).unwrap();
1517 assert_eq!(v1, 1);
1518 assert_eq!(v2, 2);
1519 assert_eq!(v3, 3);
1520 }
1521
1522 #[cfg(feature = "mmap-fsync")]
1523 #[test]
1524 fn test_fsync_feature_syncs_on_section_flush() {
1525 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1526 let (logger, _) = make_a_logger(&tmp_dir, LARGE_SLAB);
1527 {
1528 let mut stream =
1529 stream_write(logger.clone(), UnifiedLogType::StructuredLogLine, 1024).unwrap();
1530 stream.log(&1u32).unwrap();
1531 }
1532
1533 let logger = logger.lock().unwrap();
1534 assert!(
1535 logger.front_slab.sync_call_count > 0,
1536 "expected mmap-fsync to issue at least one sync_all call"
1537 );
1538 }
1539
1540 #[derive(Debug, Encode, Decode)]
1543 enum CopperListStateMock {
1544 Free,
1545 ProcessingTasks,
1546 BeingSerialized,
1547 }
1548
1549 #[derive(Encode, Decode)]
1550 struct CopperList<P: bincode::enc::Encode> {
1551 state: CopperListStateMock,
1552 payload: P, }
1554
1555 #[test]
1556 fn test_copperlist_list_like_logging() {
1557 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1558 let (logger, f) = make_a_logger(&tmp_dir, LARGE_SLAB);
1559 {
1560 let mut stream =
1561 stream_write(logger.clone(), UnifiedLogType::CopperList, 1024).unwrap();
1562 let cl0 = CopperList {
1563 state: CopperListStateMock::Free,
1564 payload: (1u32, 2u32, 3u32),
1565 };
1566 let cl1 = CopperList {
1567 state: CopperListStateMock::ProcessingTasks,
1568 payload: (4u32, 5u32, 6u32),
1569 };
1570 stream.log(&cl0).unwrap();
1571 stream.log(&cl1).unwrap();
1572 }
1573 drop(logger);
1574
1575 let MmapUnifiedLogger::Read(mut dl) = MmapUnifiedLoggerBuilder::new()
1576 .file_base_name(&f)
1577 .build()
1578 .expect("Failed to build logger")
1579 else {
1580 panic!("Failed to build logger");
1581 };
1582 let section = dl
1583 .read_next_section_type(UnifiedLogType::CopperList)
1584 .expect("Failed to read section");
1585 assert!(section.is_some());
1586 let section = section.unwrap();
1587
1588 let mut reader = SliceReader::new(§ion[..]);
1589 let cl0: CopperList<(u32, u32, u32)> = decode_from_reader(&mut reader, standard()).unwrap();
1590 let cl1: CopperList<(u32, u32, u32)> = decode_from_reader(&mut reader, standard()).unwrap();
1591 assert_eq!(cl0.payload.1, 2);
1592 assert_eq!(cl1.payload.2, 6);
1593 }
1594
1595 #[test]
1596 fn test_multi_slab_end2end() {
1597 let tmp_dir = TempDir::new().expect("could not create a tmp dir");
1598 let (logger, f) = make_a_logger(&tmp_dir, SMALL_SLAB);
1599 {
1600 let mut stream =
1601 stream_write(logger.clone(), UnifiedLogType::CopperList, 1024).unwrap();
1602 let cl0 = CopperList {
1603 state: CopperListStateMock::Free,
1604 payload: (1u32, 2u32, 3u32),
1605 };
1606 for _ in 0..10000 {
1608 stream.log(&cl0).unwrap();
1609 }
1610 }
1611 drop(logger);
1612
1613 let MmapUnifiedLogger::Read(mut dl) = MmapUnifiedLoggerBuilder::new()
1614 .file_base_name(&f)
1615 .build()
1616 .expect("Failed to build logger")
1617 else {
1618 panic!("Failed to build logger");
1619 };
1620 let mut total_readback = 0;
1621 loop {
1622 let section = dl.read_next_section_type(UnifiedLogType::CopperList);
1623 if section.is_err() {
1624 break;
1625 }
1626 let section = section.unwrap();
1627 if section.is_none() {
1628 break;
1629 }
1630 let section = section.unwrap();
1631
1632 let mut reader = SliceReader::new(§ion[..]);
1633 loop {
1634 let maybe_cl: Result<CopperList<(u32, u32, u32)>, _> =
1635 decode_from_reader(&mut reader, standard());
1636 if maybe_cl.is_ok() {
1637 total_readback += 1;
1638 } else {
1639 break;
1640 }
1641 }
1642 }
1643 assert_eq!(total_readback, 10000);
1644 }
1645}