1use std::fs::{File, OpenOptions};
46use std::io::{Read, Seek, SeekFrom, Write};
47use std::path::Path;
48use std::sync::Arc;
49
50use crate::recovery::classify_backend_error;
51use crate::PipelineError;
52use vyre_driver::backend::BackendError;
53use vyre_foundation::diagnostics::RetryClass;
54
55const LOG_MAGIC: &[u8; 8] = b"VRRL0001";
56const LOG_VERSION: u32 = 1;
57const RECORD_MAGIC: u32 = 0xDEAD_BEEF;
58const RECORD_BYTES: u64 = 64;
59const HEADER_BYTES: u64 = 32;
60const MAX_REPLAY_RECORDS: u64 = 1_048_576;
61
62#[derive(Debug, Clone, Copy, PartialEq, Eq)]
64pub struct RecordedSlot {
65 pub timestamp_ns: u64,
67 pub slot_idx: u32,
69 pub tenant_id: u32,
71 pub opcode: u32,
73 pub args: [u32; 4],
76 pub epoch: u32,
80}
81
82#[derive(Debug, Clone, Copy, PartialEq, Eq)]
84pub struct ReplayRecord {
85 pub slot: RecordedSlot,
87 pub failure: Option<ReplayFailureEvidence>,
89}
90
91#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
93pub enum ReplayFailureClass {
94 #[default]
96 None,
97 DeviceLoss,
99 TransientQueue,
101 ProgramBug,
103 Unclassified,
105}
106
107impl ReplayFailureClass {
108 const NONE: u32 = 0;
109 const DEVICE_LOSS: u32 = 1;
110 const TRANSIENT_QUEUE: u32 = 2;
111 const PROGRAM_BUG: u32 = 3;
112 const UNCLASSIFIED: u32 = 4;
113
114 const fn encode(self) -> u32 {
115 match self {
116 Self::None => Self::NONE,
117 Self::DeviceLoss => Self::DEVICE_LOSS,
118 Self::TransientQueue => Self::TRANSIENT_QUEUE,
119 Self::ProgramBug => Self::PROGRAM_BUG,
120 Self::Unclassified => Self::UNCLASSIFIED,
121 }
122 }
123
124 const fn decode(raw: u32) -> Self {
125 match raw {
126 Self::NONE => Self::None,
127 Self::DEVICE_LOSS => Self::DeviceLoss,
128 Self::TRANSIENT_QUEUE => Self::TransientQueue,
129 Self::PROGRAM_BUG => Self::ProgramBug,
130 Self::UNCLASSIFIED => Self::Unclassified,
131 _ => Self::Unclassified,
132 }
133 }
134
135 const fn from_retry_class(class: RetryClass) -> Self {
136 match class {
137 RetryClass::NewDevice => Self::DeviceLoss,
138 RetryClass::SameDevice => Self::TransientQueue,
139 RetryClass::Never | RetryClass::RecompileSource => Self::ProgramBug,
140 _ => Self::Unclassified,
141 }
142 }
143}
144
145#[derive(Debug, Clone, Copy, PartialEq, Eq)]
147pub struct ReplayFailureEvidence {
148 pub slot_status: u32,
150 pub failure_class: ReplayFailureClass,
152 pub backend_error_code: u32,
154 pub output_digest: u64,
156}
157
158impl ReplayFailureEvidence {
159 #[must_use]
161 pub fn from_backend_error(slot_status: u32, error: &BackendError, output_bytes: &[u8]) -> Self {
162 Self {
163 slot_status,
164 failure_class: ReplayFailureClass::from_retry_class(classify_backend_error(error)),
165 backend_error_code: error.code().stable_id(),
166 output_digest: output_digest(output_bytes),
167 }
168 }
169
170 fn from_words(
171 slot_status: u32,
172 failure_class: u32,
173 backend_error_code: u32,
174 output_digest: u64,
175 ) -> Option<Self> {
176 if slot_status == 0 && failure_class == 0 && backend_error_code == 0 && output_digest == 0 {
177 return None;
178 }
179 Some(Self {
180 slot_status,
181 failure_class: ReplayFailureClass::decode(failure_class),
182 backend_error_code,
183 output_digest,
184 })
185 }
186}
187
188#[derive(Debug, thiserror::Error)]
191#[non_exhaustive]
192pub enum ReplayLogError {
193 #[error("replay log {op} on `{path}` failed: {source}. Fix: check disk space + permissions.")]
195 Io {
196 op: &'static str,
198 path: Arc<str>,
200 #[source]
202 source: std::io::Error,
203 },
204 #[error("replay log `{path}` header mismatch. Fix: regenerate the log; VRRL format may have changed.")]
206 HeaderMismatch {
207 path: Arc<str>,
209 },
210 #[error("replay log capacity must be > 0. Fix: construct with at least one slot.")]
212 ZeroCapacity,
213 #[error("replay log capacity {count} exceeds max {max}. Fix: shard replay into smaller logs.")]
217 CapacityOverflow {
218 count: u64,
220 max: u64,
222 },
223}
224
225fn io_err(op: &'static str, path: &Path, source: std::io::Error) -> ReplayLogError {
226 ReplayLogError::Io {
227 op,
228 path: Arc::from(path.to_string_lossy().as_ref()),
229 source,
230 }
231}
232
233#[derive(Debug)]
237pub struct RingLog {
238 file: File,
239 path_repr: Arc<str>,
240 capacity: u64,
241 next_slot: u64,
242}
243
244impl RingLog {
245 pub fn open(path: impl AsRef<Path>, capacity: u64) -> Result<Self, ReplayLogError> {
256 if capacity == 0 {
257 return Err(ReplayLogError::ZeroCapacity);
258 }
259 validate_capacity(capacity)?;
260
261 let path = path.as_ref();
262 let path_repr: Arc<str> = Arc::from(path.to_string_lossy().as_ref());
263 let existed = path.exists();
264 let mut file = OpenOptions::new()
265 .create(true)
266 .truncate(false)
267 .read(true)
268 .write(true)
269 .open(path)
270 .map_err(|e| io_err("open", path, e))?;
271
272 if existed {
273 let mut magic = [0u8; 8];
274 file.read_exact(&mut magic)
275 .map_err(|e| io_err("read", path, e))?;
276 if &magic != LOG_MAGIC {
277 return Err(ReplayLogError::HeaderMismatch {
278 path: Arc::clone(&path_repr),
279 });
280 }
281 let mut version_bytes = [0u8; 4];
282 file.read_exact(&mut version_bytes)
283 .map_err(|e| io_err("read", path, e))?;
284 if u32::from_le_bytes(version_bytes) != LOG_VERSION {
285 return Err(ReplayLogError::HeaderMismatch {
286 path: Arc::clone(&path_repr),
287 });
288 }
289 let mut _flags = [0u8; 4];
290 file.read_exact(&mut _flags)
291 .map_err(|e| io_err("read", path, e))?;
292 let mut cap_bytes = [0u8; 8];
293 file.read_exact(&mut cap_bytes)
294 .map_err(|e| io_err("read", path, e))?;
295 let mut cursor_bytes = [0u8; 8];
296 file.read_exact(&mut cursor_bytes)
297 .map_err(|e| io_err("read", path, e))?;
298 let existing_cap = u64::from_le_bytes(cap_bytes);
299 validate_capacity(existing_cap)?;
300 let cursor = u64::from_le_bytes(cursor_bytes);
301 return Ok(Self {
302 file,
303 path_repr,
304 capacity: existing_cap,
305 next_slot: cursor % existing_cap,
306 });
307 }
308
309 let total_bytes = log_file_len(capacity)?;
313 file.set_len(total_bytes)
314 .map_err(|e| io_err("set_len", path, e))?;
315 file.seek(SeekFrom::Start(0))
316 .map_err(|e| io_err("seek", path, e))?;
317 file.write_all(LOG_MAGIC)
318 .map_err(|e| io_err("write", path, e))?;
319 file.write_all(&LOG_VERSION.to_le_bytes())
320 .map_err(|e| io_err("write", path, e))?;
321 file.write_all(&0u32.to_le_bytes())
322 .map_err(|e| io_err("write", path, e))?; file.write_all(&capacity.to_le_bytes())
324 .map_err(|e| io_err("write", path, e))?;
325 file.write_all(&0u64.to_le_bytes())
326 .map_err(|e| io_err("write", path, e))?; Ok(Self {
329 file,
330 path_repr,
331 capacity,
332 next_slot: 0,
333 })
334 }
335
336 #[must_use]
339 pub fn capacity(&self) -> u64 {
340 self.capacity
341 }
342
343 #[must_use]
345 pub fn cursor(&self) -> u64 {
346 self.next_slot
347 }
348
349 #[must_use]
351 pub fn path(&self) -> &str {
352 self.path_repr.as_ref()
353 }
354
355 pub fn append(&mut self, slot: RecordedSlot) -> Result<(), ReplayLogError> {
363 self.append_record(ReplayRecord {
364 slot,
365 failure: None,
366 })
367 }
368
369 pub fn append_with_failure(
375 &mut self,
376 slot: RecordedSlot,
377 failure: ReplayFailureEvidence,
378 ) -> Result<(), ReplayLogError> {
379 self.append_record(ReplayRecord {
380 slot,
381 failure: Some(failure),
382 })
383 }
384
385 fn append_record(&mut self, record: ReplayRecord) -> Result<(), ReplayLogError> {
386 let record_offset = log_record_offset(self.next_slot)?;
387 self.file
388 .seek(SeekFrom::Start(record_offset))
389 .map_err(|e| self.io_err("seek", e))?;
390
391 let mut buf = [0u8; RECORD_BYTES as usize];
392 buf[0..4].copy_from_slice(&RECORD_MAGIC.to_le_bytes());
393 buf[4..12].copy_from_slice(&record.slot.timestamp_ns.to_le_bytes());
394 buf[12..16].copy_from_slice(&record.slot.slot_idx.to_le_bytes());
395 buf[16..20].copy_from_slice(&record.slot.tenant_id.to_le_bytes());
396 buf[20..24].copy_from_slice(&record.slot.opcode.to_le_bytes());
397 buf[24..28].copy_from_slice(&record.slot.args[0].to_le_bytes());
398 buf[28..32].copy_from_slice(&record.slot.args[1].to_le_bytes());
399 buf[32..36].copy_from_slice(&record.slot.args[2].to_le_bytes());
400 buf[36..40].copy_from_slice(&record.slot.args[3].to_le_bytes());
401 buf[40..44].copy_from_slice(&record.slot.epoch.to_le_bytes());
402 if let Some(failure) = record.failure {
403 buf[44..48].copy_from_slice(&failure.slot_status.to_le_bytes());
404 buf[48..52].copy_from_slice(&failure.failure_class.encode().to_le_bytes());
405 buf[52..56].copy_from_slice(&failure.backend_error_code.to_le_bytes());
406 buf[56..64].copy_from_slice(&failure.output_digest.to_le_bytes());
407 }
408 self.file
409 .write_all(&buf)
410 .map_err(|e| self.io_err("write", e))?;
411
412 self.next_slot = (self.next_slot + 1) % self.capacity;
415 self.file
416 .seek(SeekFrom::Start(24)) .map_err(|e| self.io_err("seek", e))?;
418 self.file
419 .write_all(&self.next_slot.to_le_bytes())
420 .map_err(|e| self.io_err("write", e))?;
421
422 Ok(())
423 }
424
425 pub fn replay_all(&mut self) -> Result<Vec<RecordedSlot>, ReplayLogError> {
436 Ok(self
437 .replay_records()?
438 .into_iter()
439 .map(|record| record.slot)
440 .collect())
441 }
442
443 pub fn replay_records(&mut self) -> Result<Vec<ReplayRecord>, ReplayLogError> {
450 let capacity =
451 usize::try_from(self.capacity).map_err(|_| ReplayLogError::CapacityOverflow {
452 count: self.capacity,
453 max: MAX_REPLAY_RECORDS,
454 })?;
455 let mut out = Vec::with_capacity(capacity);
456 for step in 0..self.capacity {
457 let slot_index = (self.next_slot + step) % self.capacity;
458 let offset = log_record_offset(slot_index)?;
459 self.file
460 .seek(SeekFrom::Start(offset))
461 .map_err(|e| self.io_err("seek", e))?;
462 let mut buf = [0u8; RECORD_BYTES as usize];
463 self.file
464 .read_exact(&mut buf)
465 .map_err(|e| self.io_err("read", e))?;
466 let magic = read_u32(&buf, 0);
467 if magic == 0 {
468 tracing::warn!(
480 slot_index,
481 next_slot = self.next_slot,
482 log_capacity = self.capacity,
483 step,
484 "replay_records: zero-magic record at slot_index {slot_index} (step {step}). \
485 If the log has wrapped this is a corruption gap, the replay will be shorter than expected. \
486 Fix: ensure the replay-log file is not subject to external zeroing or partial-write truncation."
487 );
488 continue;
489 }
490 if magic != RECORD_MAGIC {
491 return Err(ReplayLogError::HeaderMismatch {
492 path: self.path_repr.clone(),
493 });
494 }
495 let slot = RecordedSlot {
496 timestamp_ns: read_u64(&buf, 4),
497 slot_idx: read_u32(&buf, 12),
498 tenant_id: read_u32(&buf, 16),
499 opcode: read_u32(&buf, 20),
500 args: [
501 read_u32(&buf, 24),
502 read_u32(&buf, 28),
503 read_u32(&buf, 32),
504 read_u32(&buf, 36),
505 ],
506 epoch: read_u32(&buf, 40),
507 };
508 out.push(ReplayRecord {
509 slot,
510 failure: ReplayFailureEvidence::from_words(
511 read_u32(&buf, 44),
512 read_u32(&buf, 48),
513 read_u32(&buf, 52),
514 read_u64(&buf, 56),
515 ),
516 });
517 }
518 Ok(out)
519 }
520
521 pub fn sync(&mut self) -> Result<(), ReplayLogError> {
529 self.file.sync_all().map_err(|e| self.io_err("sync", e))?;
530 Ok(())
531 }
532
533 fn io_err(&self, op: &'static str, source: std::io::Error) -> ReplayLogError {
534 ReplayLogError::Io {
535 op,
536 path: self.path_repr.clone(),
537 source,
538 }
539 }
540}
541
542fn validate_capacity(capacity: u64) -> Result<(), ReplayLogError> {
543 if capacity == 0 {
544 return Err(ReplayLogError::ZeroCapacity);
545 }
546 if capacity > MAX_REPLAY_RECORDS {
547 return Err(ReplayLogError::CapacityOverflow {
548 count: capacity,
549 max: MAX_REPLAY_RECORDS,
550 });
551 }
552 Ok(())
553}
554
555fn log_file_len(capacity: u64) -> Result<u64, ReplayLogError> {
556 log_record_position(capacity)
557}
558
559fn log_record_offset(slot_index: u64) -> Result<u64, ReplayLogError> {
560 log_record_position(slot_index)
561}
562
563fn log_record_position(record_index: u64) -> Result<u64, ReplayLogError> {
564 let record_bytes =
565 vyre_driver::accounting::checked_mul_u64_lazy(record_index, RECORD_BYTES, || {
566 replay_capacity_overflow(record_index)
567 })?;
568 vyre_driver::accounting::checked_add_u64_lazy(HEADER_BYTES, record_bytes, || {
569 replay_capacity_overflow(record_index)
570 })
571}
572
573fn replay_capacity_overflow(count: u64) -> ReplayLogError {
574 ReplayLogError::CapacityOverflow {
575 count,
576 max: MAX_REPLAY_RECORDS,
577 }
578}
579
580fn read_u32(buf: &[u8], offset: usize) -> u32 {
581 let mut bytes = [0u8; 4];
582 bytes.copy_from_slice(&buf[offset..offset + 4]);
583 u32::from_le_bytes(bytes)
584}
585
586fn read_u64(buf: &[u8], offset: usize) -> u64 {
587 let mut bytes = [0u8; 8];
588 bytes.copy_from_slice(&buf[offset..offset + 8]);
589 u64::from_le_bytes(bytes)
590}
591
592fn output_digest(bytes: &[u8]) -> u64 {
593 let digest = blake3::hash(bytes);
594 let mut out = [0u8; 8];
595 out.copy_from_slice(&digest.as_bytes()[..8]);
596 u64::from_le_bytes(out)
597}
598
599impl From<ReplayLogError> for PipelineError {
602 fn from(err: ReplayLogError) -> Self {
603 PipelineError::Backend(err.to_string())
604 }
605}
606
607#[cfg(test)]
608mod tests {
609 use super::*;
610
611 fn rec(slot_idx: u32, epoch: u32) -> RecordedSlot {
612 RecordedSlot {
613 timestamp_ns: 1_000_000 + slot_idx as u64,
614 slot_idx,
615 tenant_id: 0,
616 opcode: 0x4000_0000 + slot_idx,
617 args: [slot_idx, slot_idx * 2, slot_idx * 3, slot_idx * 4],
618 epoch,
619 }
620 }
621
622 #[test]
623 fn open_rejects_zero_capacity() {
624 let dir = tempfile::tempdir().unwrap();
625 let path = dir.path().join("log.vrrl");
626 let err = RingLog::open(&path, 0).expect_err("zero capacity must reject");
627 assert!(matches!(err, ReplayLogError::ZeroCapacity));
628 }
629
630 #[test]
631 fn append_and_replay_round_trip() {
632 let dir = tempfile::tempdir().unwrap();
633 let path = dir.path().join("log.vrrl");
634 let mut log = RingLog::open(&path, 4)
635 .expect("Fix: open fresh log; restore this invariant before continuing.");
636 log.append(rec(1, 10)).unwrap();
637 log.append(rec(2, 11)).unwrap();
638 log.sync().unwrap();
639
640 let replay = log
641 .replay_all()
642 .expect("Fix: replay; restore this invariant before continuing.");
643 assert_eq!(replay.len(), 2);
644 assert_eq!(replay[0].slot_idx, 1);
645 assert_eq!(replay[0].epoch, 10);
646 assert_eq!(replay[1].slot_idx, 2);
647 assert_eq!(replay[1].epoch, 11);
648 }
649
650 #[test]
651 fn append_with_failure_round_trips_reproduction_evidence() {
652 let dir = tempfile::tempdir().unwrap();
653 let path = dir.path().join("log.vrrl");
654 let mut log = RingLog::open(&path, 4)
655 .expect("Fix: open fresh log; restore this invariant before continuing.");
656 let backend_error = BackendError::DeviceLost {
657 backend: "fixture".to_string(),
658 device: "fixture-0".to_string(),
659 generation: 7,
660 message: "device loss after queue submit".to_string(),
661 };
662 let failure =
663 ReplayFailureEvidence::from_backend_error(3, &backend_error, b"partial-output");
664
665 assert_eq!(failure.failure_class, ReplayFailureClass::DeviceLoss);
666 assert_eq!(failure.backend_error_code, backend_error.code().stable_id());
667 assert_ne!(failure.output_digest, 0);
668
669 log.append_with_failure(rec(7, 44), failure).unwrap();
670 log.sync().unwrap();
671
672 let replay = log
673 .replay_records()
674 .expect("Fix: replay records; restore this invariant before continuing.");
675 assert_eq!(replay.len(), 1);
676 assert_eq!(replay[0].slot.slot_idx, 7);
677 assert_eq!(replay[0].slot.epoch, 44);
678 assert_eq!(replay[0].failure, Some(failure));
679 }
680
681 #[test]
682 fn append_without_failure_has_no_failure_evidence() {
683 let dir = tempfile::tempdir().unwrap();
684 let path = dir.path().join("log.vrrl");
685 let mut log = RingLog::open(&path, 2)
686 .expect("Fix: open fresh log; restore this invariant before continuing.");
687
688 log.append(rec(1, 10)).unwrap();
689
690 let replay = log
691 .replay_records()
692 .expect("Fix: replay records; restore this invariant before continuing.");
693 assert_eq!(replay.len(), 1);
694 assert_eq!(replay[0].slot.slot_idx, 1);
695 assert_eq!(replay[0].failure, None);
696 }
697
698 #[test]
699 fn log_rollover_preserves_most_recent() {
700 let dir = tempfile::tempdir().unwrap();
701 let path = dir.path().join("log.vrrl");
702 let mut log =
703 RingLog::open(&path, 3).expect("Fix: open; restore this invariant before continuing.");
704 for i in 0..5 {
705 log.append(rec(i, 100 + i)).unwrap();
706 }
707 let replay = log
708 .replay_all()
709 .expect("Fix: replay; restore this invariant before continuing.");
710 assert_eq!(replay.len(), 3, "capacity=3 must retain exactly 3 records");
711 let slot_ids: Vec<u32> = replay.iter().map(|r| r.slot_idx).collect();
712 assert_eq!(slot_ids, vec![2, 3, 4]);
716 }
717
718 #[test]
719 fn reopen_restores_cursor() {
720 let dir = tempfile::tempdir().unwrap();
721 let path = dir.path().join("log.vrrl");
722 {
723 let mut log = RingLog::open(&path, 4)
724 .expect("Fix: open fresh; restore this invariant before continuing.");
725 log.append(rec(1, 10)).unwrap();
726 log.append(rec(2, 11)).unwrap();
727 log.sync().unwrap();
728 }
729 let mut reopened = RingLog::open(&path, 4)
730 .expect("Fix: reopen; restore this invariant before continuing.");
731 assert_eq!(reopened.cursor(), 2);
732 let replay = reopened.replay_all().unwrap();
733 assert_eq!(replay.len(), 2);
734 }
735
736 #[test]
737 fn corrupted_magic_rejected() {
738 use std::io::Write as _;
739
740 let dir = tempfile::tempdir().unwrap();
741 let path = dir.path().join("log.vrrl");
742 {
743 let mut f = std::fs::File::create(&path).unwrap();
745 f.write_all(b"XXXX0001").unwrap();
746 f.write_all(&1u32.to_le_bytes()).unwrap();
747 f.write_all(&0u32.to_le_bytes()).unwrap();
748 f.write_all(&4u64.to_le_bytes()).unwrap();
749 f.write_all(&0u64.to_le_bytes()).unwrap();
750 f.set_len(HEADER_BYTES + 4 * RECORD_BYTES).unwrap();
752 }
753 let err = RingLog::open(&path, 4).expect_err("wrong magic must reject");
754 assert!(matches!(err, ReplayLogError::HeaderMismatch { .. }));
755 }
756
757 fn write_header(path: &Path, capacity: u64, cursor: u64) {
758 use std::io::Write as _;
759
760 let mut f = std::fs::File::create(path).unwrap();
761 f.write_all(LOG_MAGIC).unwrap();
762 f.write_all(&LOG_VERSION.to_le_bytes()).unwrap();
763 f.write_all(&0u32.to_le_bytes()).unwrap();
764 f.write_all(&capacity.to_le_bytes()).unwrap();
765 f.write_all(&cursor.to_le_bytes()).unwrap();
766 }
767
768 #[test]
769 fn existing_log_zero_capacity_rejected_before_cursor_modulo() {
770 let dir = tempfile::tempdir().unwrap();
771 let path = dir.path().join("log.vrrl");
772 write_header(&path, 0, 0);
773
774 let err = RingLog::open(&path, 4).expect_err("header capacity=0 must reject");
775 assert!(matches!(err, ReplayLogError::ZeroCapacity));
776 }
777
778 #[test]
779 fn existing_log_huge_capacity_rejected_before_replay_allocation() {
780 let dir = tempfile::tempdir().unwrap();
781 let path = dir.path().join("log.vrrl");
782 write_header(&path, MAX_REPLAY_RECORDS + 1, 0);
783
784 let err = RingLog::open(&path, 4).expect_err("huge header capacity must reject");
785 assert!(matches!(
786 err,
787 ReplayLogError::CapacityOverflow {
788 count,
789 max: MAX_REPLAY_RECORDS
790 } if count == MAX_REPLAY_RECORDS + 1
791 ));
792 }
793
794 #[test]
795 fn capacity_overflow_rejected() {
796 let dir = tempfile::tempdir().unwrap();
797 let path = dir.path().join("log.vrrl");
798 let err = RingLog::open(&path, MAX_REPLAY_RECORDS + 1)
799 .expect_err("over-size capacity must reject");
800 assert!(matches!(
801 err,
802 ReplayLogError::CapacityOverflow {
803 count,
804 max: MAX_REPLAY_RECORDS
805 } if count == MAX_REPLAY_RECORDS + 1
806 ));
807 }
808
809 #[test]
821 fn replay_zero_magic_mid_sequence_skips_gracefully_and_produces_shorter_result() {
822 use std::io::{Seek, SeekFrom, Write as _};
823
824 let dir = tempfile::tempdir().unwrap();
825 let path = dir.path().join("log.vrrl");
826 let mut log = RingLog::open(&path, 4)
827 .expect("Fix: open fresh log; restore this invariant before continuing.");
828
829 log.append(rec(10, 100)).unwrap();
831 log.append(rec(20, 200)).unwrap();
832 log.append(rec(30, 300)).unwrap();
833 log.sync().unwrap();
834
835 {
838 let records = log
839 .replay_all()
840 .expect("Fix: replay of 3 records must succeed");
841 assert_eq!(records.len(), 3, "Fix: 3 appended records must all replay");
842 }
843
844 let slot1_offset = HEADER_BYTES + RECORD_BYTES; {
848 let mut f = std::fs::OpenOptions::new().write(true).open(&path).unwrap();
849 f.seek(SeekFrom::Start(slot1_offset)).unwrap();
850 f.write_all(&[0u8; RECORD_BYTES as usize]).unwrap();
851 f.sync_all().unwrap();
852 }
853
854 let mut log2 = RingLog::open(&path, 4).expect("Fix: reopen after zeroing must succeed");
856
857 let records = log2
859 .replay_all()
860 .expect("Fix: replay with a zeroed slot must not error");
861
862 assert_eq!(
865 records.len(),
866 2,
867 "Fix: zeroed slot must be skipped, yielding 2 out of 3 records; got: {:?}",
868 records.iter().map(|r| r.slot_idx).collect::<Vec<_>>()
869 );
870 assert_eq!(
872 records[0].slot_idx, 10,
873 "Fix: first replayed record must be slot_idx=10"
874 );
875 assert_eq!(
876 records[1].slot_idx, 30,
877 "Fix: second replayed record must be slot_idx=30"
878 );
879 }
880}