1#[cfg(test)]
8use std::cell::Cell;
9use std::convert::TryFrom;
10use std::fs::{File, OpenOptions};
11use std::io::{Read, Seek, SeekFrom, Write};
12use std::path::Path;
13use std::time::{Duration, Instant};
14
15use crate::error::{Error, Result};
16
17#[cfg(test)]
18thread_local! {
19 static WAL_SYNC_CALLS: Cell<usize> = const { Cell::new(0) };
20}
21
22#[cfg(test)]
23#[allow(dead_code)]
24pub(crate) fn reset_sync_calls() {
25 WAL_SYNC_CALLS.with(|counter| counter.set(0));
26}
27
28#[cfg(test)]
29#[allow(dead_code)]
30pub(crate) fn sync_calls() -> usize {
31 WAL_SYNC_CALLS.with(|counter| counter.get())
32}
33
34fn sync_file(file: &File) -> Result<()> {
35 #[cfg(test)]
36 {
37 WAL_SYNC_CALLS.with(|counter| counter.set(counter.get() + 1));
38 }
39 file.sync_data()?;
40 Ok(())
41}
42
43pub const WAL_MAGIC: [u8; 4] = *b"AWAL";
45pub const WAL_FORMAT_VERSION_V04: u16 = 1;
47pub const WAL_FORMAT_VERSION: u16 = 2;
49pub const WAL_VERSION: u16 = WAL_FORMAT_VERSION;
51pub const WAL_SEGMENT_HEADER_SIZE: usize = 28;
53pub const WAL_SECTION_HEADER_SIZE: usize = 20;
55pub const WAL_ENTRY_FIXED_HEADER: usize = 8 + 4;
57
58#[repr(u8)]
60#[derive(Debug, Clone, Copy, PartialEq, Eq)]
61pub enum WalEntryType {
62 Put = 0,
64 Delete = 1,
66 Batch = 2,
68}
69
70impl TryFrom<u8> for WalEntryType {
71 type Error = Error;
72
73 fn try_from(value: u8) -> Result<Self> {
74 match value {
75 0 => Ok(Self::Put),
76 1 => Ok(Self::Delete),
77 2 => Ok(Self::Batch),
78 other => Err(Error::InvalidFormat(format!(
79 "unknown WAL entry type: {other}"
80 ))),
81 }
82 }
83}
84
85#[repr(u8)]
87#[derive(Debug, Clone, Copy, PartialEq, Eq)]
88pub enum WalOpType {
89 Put = 0,
91 Delete = 1,
93}
94
95impl TryFrom<u8> for WalOpType {
96 type Error = Error;
97
98 fn try_from(value: u8) -> Result<Self> {
99 match value {
100 0 => Ok(Self::Put),
101 1 => Ok(Self::Delete),
102 other => Err(Error::InvalidFormat(format!(
103 "unknown WAL batch op type: {other}"
104 ))),
105 }
106 }
107}
108
109#[derive(Debug, Clone, PartialEq, Eq)]
111pub struct WalConfig {
112 pub segment_size: usize,
114 pub max_segments: usize,
116 pub sync_mode: SyncMode,
118}
119
120impl Default for WalConfig {
121 fn default() -> Self {
122 Self {
123 segment_size: 64 * 1024 * 1024,
124 max_segments: 8,
125 sync_mode: SyncMode::EveryWrite,
126 }
127 }
128}
129
130#[derive(Debug, Clone, PartialEq, Eq)]
132pub enum SyncMode {
133 EveryWrite,
135 BatchSync {
140 max_batch_size: usize,
142 max_wait_ms: u64,
144 },
145 NoSync,
149}
150
151#[derive(Debug, Clone, PartialEq, Eq)]
153pub struct WalSegmentHeader {
154 pub version: u16,
156 pub segment_id: u64,
158 pub first_lsn: u64,
163 pub crc32: u32,
165 pub reserved: u16,
167}
168
169impl WalSegmentHeader {
170 pub fn new(segment_id: u64, first_lsn: u64) -> Self {
172 let crc32 = compute_crc(WAL_VERSION, segment_id, first_lsn);
173 Self {
174 version: WAL_VERSION,
175 segment_id,
176 first_lsn,
177 crc32,
178 reserved: 0,
179 }
180 }
181
182 pub fn to_bytes(&self) -> [u8; WAL_SEGMENT_HEADER_SIZE] {
184 let mut buf = [0u8; WAL_SEGMENT_HEADER_SIZE];
185 buf[0..4].copy_from_slice(&WAL_MAGIC);
186 buf[4..6].copy_from_slice(&self.version.to_le_bytes());
187 buf[6..14].copy_from_slice(&self.segment_id.to_le_bytes());
188 buf[14..22].copy_from_slice(&self.first_lsn.to_le_bytes());
189 buf[22..26].copy_from_slice(&self.crc32.to_le_bytes());
190 buf[26..28].copy_from_slice(&self.reserved.to_le_bytes());
191 buf
192 }
193
194 pub fn from_bytes(bytes: &[u8; WAL_SEGMENT_HEADER_SIZE]) -> Result<Self> {
196 if bytes[0..4] != WAL_MAGIC {
197 return Err(Error::InvalidFormat("WAL magic mismatch".into()));
198 }
199
200 let version = u16::from_le_bytes([bytes[4], bytes[5]]);
201 if version != WAL_VERSION {
202 return Err(Error::InvalidFormat(format!(
203 "unsupported WAL version: {version}"
204 )));
205 }
206
207 let segment_id = u64::from_le_bytes(bytes[6..14].try_into().expect("fixed slice length"));
208 let first_lsn = u64::from_le_bytes(bytes[14..22].try_into().expect("fixed slice length"));
209 let stored_crc = u32::from_le_bytes(bytes[22..26].try_into().expect("fixed slice length"));
210 let reserved = u16::from_le_bytes(bytes[26..28].try_into().expect("fixed slice length"));
211
212 let header = Self {
213 version,
214 segment_id,
215 first_lsn,
216 crc32: stored_crc,
217 reserved,
218 };
219
220 let computed = header.compute_crc();
221 if computed != stored_crc {
222 return Err(Error::ChecksumMismatch);
223 }
224
225 Ok(header)
226 }
227
228 pub fn from_bytes_allow_legacy(bytes: &[u8; WAL_SEGMENT_HEADER_SIZE]) -> Result<Self> {
230 if bytes[0..4] != WAL_MAGIC {
231 return Err(Error::InvalidFormat("WAL magic mismatch".into()));
232 }
233
234 let version = u16::from_le_bytes([bytes[4], bytes[5]]);
235 if version > WAL_FORMAT_VERSION {
236 return Err(Error::InvalidFormat(format!(
237 "unsupported WAL version: {version}"
238 )));
239 }
240
241 let segment_id = u64::from_le_bytes(bytes[6..14].try_into().expect("fixed slice length"));
242 let first_lsn = u64::from_le_bytes(bytes[14..22].try_into().expect("fixed slice length"));
243 let stored_crc = u32::from_le_bytes(bytes[22..26].try_into().expect("fixed slice length"));
244 let reserved = u16::from_le_bytes(bytes[26..28].try_into().expect("fixed slice length"));
245
246 let header = Self {
247 version,
248 segment_id,
249 first_lsn,
250 crc32: stored_crc,
251 reserved,
252 };
253
254 let computed = header.compute_crc();
255 if computed != stored_crc {
256 return Err(Error::ChecksumMismatch);
257 }
258
259 Ok(header)
260 }
261
262 fn compute_crc(&self) -> u32 {
263 compute_crc(self.version, self.segment_id, self.first_lsn)
264 }
265
266 pub fn is_legacy_format(&self) -> bool {
268 self.version < WAL_FORMAT_VERSION
269 }
270}
271
272#[derive(Debug, Clone, PartialEq, Eq)]
274pub struct WalBatchOp {
275 pub op_type: WalOpType,
277 pub key: Vec<u8>,
279 pub value: Option<Vec<u8>>,
281}
282
283impl WalBatchOp {
284 fn encoded_len(&self) -> usize {
285 let val_len = self.value.as_ref().map(|v| v.len()).unwrap_or(0);
286 1 + varint_len(self.key.len() as u64)
287 + self.key.len()
288 + varint_len(val_len as u64)
289 + val_len
290 }
291}
292
293#[derive(Debug, Clone, PartialEq, Eq)]
295pub enum WalEntryPayload {
296 Put {
298 key: Vec<u8>,
300 value: Vec<u8>,
302 },
303 Delete {
305 key: Vec<u8>,
307 },
308 Batch(Vec<WalBatchOp>),
310}
311
312#[derive(Debug, Clone, PartialEq, Eq)]
314pub struct WalEntry {
315 pub lsn: u64,
317 pub payload: WalEntryPayload,
319}
320
321impl WalEntry {
322 pub fn put(lsn: u64, key: Vec<u8>, value: Vec<u8>) -> Self {
324 Self {
325 lsn,
326 payload: WalEntryPayload::Put { key, value },
327 }
328 }
329
330 pub fn delete(lsn: u64, key: Vec<u8>) -> Self {
332 Self {
333 lsn,
334 payload: WalEntryPayload::Delete { key },
335 }
336 }
337
338 pub fn batch(lsn: u64, operations: Vec<WalBatchOp>) -> Self {
340 Self {
341 lsn,
342 payload: WalEntryPayload::Batch(operations),
343 }
344 }
345
346 pub fn encoded_len(&self) -> usize {
348 let body_len = match &self.payload {
349 WalEntryPayload::Put { key, value } => {
350 1 + varint_len(key.len() as u64)
351 + key.len()
352 + varint_len(value.len() as u64)
353 + value.len()
354 }
355 WalEntryPayload::Delete { key } => {
356 1 + varint_len(key.len() as u64) + key.len() + varint_len(0)
357 }
358 WalEntryPayload::Batch(ops) => {
359 1 + varint_len(ops.len() as u64)
360 + ops.iter().map(WalBatchOp::encoded_len).sum::<usize>()
361 }
362 };
363 WAL_ENTRY_FIXED_HEADER + body_len + 4 }
365
366 pub fn encode(&self) -> Result<Vec<u8>> {
368 let mut body = Vec::with_capacity(self.encoded_len() - WAL_ENTRY_FIXED_HEADER - 4);
369 match &self.payload {
370 WalEntryPayload::Put { key, value } => {
371 body.push(WalEntryType::Put as u8);
372 encode_varint(key.len() as u64, &mut body);
373 body.extend_from_slice(key);
374 encode_varint(value.len() as u64, &mut body);
375 body.extend_from_slice(value);
376 }
377 WalEntryPayload::Delete { key } => {
378 body.push(WalEntryType::Delete as u8);
379 encode_varint(key.len() as u64, &mut body);
380 body.extend_from_slice(key);
381 encode_varint(0, &mut body);
382 }
383 WalEntryPayload::Batch(ops) => {
384 body.push(WalEntryType::Batch as u8);
385 encode_varint(ops.len() as u64, &mut body);
386 for op in ops {
387 body.push(op.op_type as u8);
388 encode_varint(op.key.len() as u64, &mut body);
389 body.extend_from_slice(&op.key);
390 let val_len = op.value.as_ref().map(|v| v.len()).unwrap_or(0);
391 encode_varint(val_len as u64, &mut body);
392 if let Some(value) = &op.value {
393 body.extend_from_slice(value);
394 }
395 }
396 }
397 }
398
399 let payload_len = body.len();
400 let total_len_field = payload_len
401 .checked_add(4)
402 .ok_or_else(|| Error::InvalidFormat("WAL entry too large".into()))?;
403 if total_len_field > u32::MAX as usize {
404 return Err(Error::InvalidFormat("WAL entry too large".into()));
405 }
406
407 let mut out = Vec::with_capacity(WAL_ENTRY_FIXED_HEADER + payload_len + 4);
408 out.extend_from_slice(&self.lsn.to_le_bytes());
409 out.extend_from_slice(&(total_len_field as u32).to_le_bytes());
410 out.extend_from_slice(&body);
411
412 let crc = crc32fast::hash(&out);
414 out.extend_from_slice(&crc.to_le_bytes());
415 Ok(out)
416 }
417
418 pub fn decode(buf: &[u8]) -> Result<(Self, usize)> {
420 if buf.len() < WAL_ENTRY_FIXED_HEADER {
421 return Err(Error::InvalidFormat(
422 "buffer too small for WAL entry header".into(),
423 ));
424 }
425
426 let lsn = u64::from_le_bytes(buf[0..8].try_into().expect("fixed slice length"));
427 let payload_and_crc_len =
428 u32::from_le_bytes(buf[8..12].try_into().expect("fixed slice length")) as usize;
429 let total_len = WAL_ENTRY_FIXED_HEADER + payload_and_crc_len;
430 if buf.len() < total_len {
431 return Err(Error::InvalidFormat(
432 "buffer truncated for WAL entry payload".into(),
433 ));
434 }
435 if payload_and_crc_len < 1 + 4 {
436 return Err(Error::InvalidFormat("WAL entry payload too small".into()));
437 }
438
439 let body_len = payload_and_crc_len - 4;
440 let body = &buf[WAL_ENTRY_FIXED_HEADER..WAL_ENTRY_FIXED_HEADER + body_len];
441 let stored_crc = u32::from_le_bytes(
442 buf[WAL_ENTRY_FIXED_HEADER + body_len..total_len]
443 .try_into()
444 .expect("fixed slice length"),
445 );
446 let computed_crc = crc32fast::hash(&buf[..WAL_ENTRY_FIXED_HEADER + body_len]);
447 if stored_crc != computed_crc {
448 return Err(Error::ChecksumMismatch);
449 }
450
451 let entry_type = WalEntryType::try_from(body[0])?;
452 let mut cursor = 1;
453
454 let payload = match entry_type {
455 WalEntryType::Put => {
456 let (key_len, key_len_bytes) = decode_varint(&body[cursor..])?;
457 cursor += key_len_bytes;
458 let key_len = key_len as usize;
459 if body_len < cursor + key_len {
460 return Err(Error::InvalidFormat("WAL entry truncated (key)".into()));
461 }
462 let key = body[cursor..cursor + key_len].to_vec();
463 cursor += key_len;
464
465 let (val_len, val_len_bytes) = decode_varint(&body[cursor..])?;
466 cursor += val_len_bytes;
467 let val_len = val_len as usize;
468 if body_len < cursor + val_len {
469 return Err(Error::InvalidFormat("WAL entry truncated (value)".into()));
470 }
471 let value = body[cursor..cursor + val_len].to_vec();
472 cursor += val_len;
473
474 if cursor != body_len {
475 return Err(Error::InvalidFormat(
476 "WAL entry has trailing bytes after Put".into(),
477 ));
478 }
479
480 WalEntryPayload::Put { key, value }
481 }
482 WalEntryType::Delete => {
483 let (key_len, key_len_bytes) = decode_varint(&body[cursor..])?;
484 cursor += key_len_bytes;
485 let key_len = key_len as usize;
486 if body_len < cursor + key_len {
487 return Err(Error::InvalidFormat("WAL entry truncated (key)".into()));
488 }
489 let key = body[cursor..cursor + key_len].to_vec();
490 cursor += key_len;
491
492 let (val_len, val_len_bytes) = decode_varint(&body[cursor..])?;
493 cursor += val_len_bytes;
494 if val_len != 0 {
495 return Err(Error::InvalidFormat(
496 "delete entry must have zero-length value".into(),
497 ));
498 }
499 if cursor != body_len {
500 return Err(Error::InvalidFormat(
501 "WAL entry has trailing bytes after Delete".into(),
502 ));
503 }
504 WalEntryPayload::Delete { key }
505 }
506 WalEntryType::Batch => {
507 let (op_count, op_count_bytes) = decode_varint(&body[cursor..])?;
508 cursor += op_count_bytes;
509 let op_count = op_count as usize;
510 let mut ops = Vec::with_capacity(op_count);
511 for _ in 0..op_count {
512 if cursor >= body_len {
513 return Err(Error::InvalidFormat(
514 "WAL batch truncated before op type".into(),
515 ));
516 }
517 let op_type = WalOpType::try_from(body[cursor])?;
518 cursor += 1;
519
520 let (key_len, key_len_bytes) = decode_varint(&body[cursor..])?;
521 cursor += key_len_bytes;
522 let key_len = key_len as usize;
523 if body_len < cursor + key_len {
524 return Err(Error::InvalidFormat("WAL batch truncated (key)".into()));
525 }
526 let key = body[cursor..cursor + key_len].to_vec();
527 cursor += key_len;
528
529 let (val_len, val_len_bytes) = decode_varint(&body[cursor..])?;
530 cursor += val_len_bytes;
531 let val_len = val_len as usize;
532 if body_len < cursor + val_len {
533 return Err(Error::InvalidFormat("WAL batch truncated (value)".into()));
534 }
535 let value = if op_type == WalOpType::Delete {
536 if val_len != 0 {
537 return Err(Error::InvalidFormat(
538 "batch delete must have zero-length value".into(),
539 ));
540 }
541 None
542 } else {
543 Some(body[cursor..cursor + val_len].to_vec())
544 };
545 cursor += val_len;
546
547 ops.push(WalBatchOp {
548 op_type,
549 key,
550 value,
551 });
552 }
553
554 if cursor != body_len {
555 return Err(Error::InvalidFormat(
556 "WAL batch has trailing unparsed bytes".into(),
557 ));
558 }
559
560 WalEntryPayload::Batch(ops)
561 }
562 };
563
564 Ok((Self { lsn, payload }, total_len))
565 }
566}
567
568fn encode_varint(mut n: u64, buf: &mut Vec<u8>) {
569 while n >= 0x80 {
570 buf.push((n as u8) | 0x80);
571 n >>= 7;
572 }
573 buf.push(n as u8);
574}
575
576fn decode_varint(data: &[u8]) -> Result<(u64, usize)> {
577 let mut result = 0u64;
578 let mut shift = 0;
579 for (i, &byte) in data.iter().enumerate() {
580 let bits = (byte & 0x7F) as u64;
581 result |= bits << shift;
582 if byte & 0x80 == 0 {
583 return Ok((result, i + 1));
584 }
585 shift += 7;
586 if shift >= 64 {
587 return Err(Error::InvalidFormat("varint overflow".into()));
588 }
589 }
590 Err(Error::InvalidFormat("varint truncated".into()))
591}
592
593fn varint_len(mut n: u64) -> usize {
594 let mut len = 1;
595 while n >= 0x80 {
596 n >>= 7;
597 len += 1;
598 }
599 len
600}
601
602fn compute_crc(version: u16, segment_id: u64, first_lsn: u64) -> u32 {
603 let mut buf = [0u8; WAL_SEGMENT_HEADER_SIZE - 6]; buf[0..4].copy_from_slice(&WAL_MAGIC);
605 buf[4..6].copy_from_slice(&version.to_le_bytes());
606 buf[6..14].copy_from_slice(&segment_id.to_le_bytes());
607 buf[14..22].copy_from_slice(&first_lsn.to_le_bytes());
608 crc32fast::hash(&buf)
609}
610
611fn ring_distance(start: u64, end: u64, len: u64) -> u64 {
612 if start <= end {
613 end - start
614 } else {
615 len - (start - end)
616 }
617}
618
619fn compute_ring_layout(config: &WalConfig) -> Result<(u64, u64, u64, u64)> {
620 if config.max_segments == 0 {
621 return Err(Error::InvalidFormat("max_segments must be >= 1".into()));
622 }
623 let segment_size = config.segment_size as u64;
624 let segment_header_bytes = WAL_SEGMENT_HEADER_SIZE as u64;
625 if segment_size <= segment_header_bytes {
626 return Err(Error::InvalidFormat(
627 "WAL segment size too small for segment header".into(),
628 ));
629 }
630 let segment_data_len = segment_size - segment_header_bytes;
631 let max_segments = config.max_segments as u64;
632 let ring_len = segment_data_len
633 .checked_mul(max_segments)
634 .ok_or_else(|| Error::InvalidFormat("WAL ring length overflow".into()))?;
635 if ring_len > WalSectionHeader::OFFSET_MASK {
636 return Err(Error::InvalidFormat(
637 "WAL ring length exceeds offset encoding capacity".into(),
638 ));
639 }
640 Ok((segment_size, segment_data_len, max_segments, ring_len))
641}
642
643fn segment_header_offset(segment_size: u64, segment_index: u64) -> u64 {
644 (WAL_SECTION_HEADER_SIZE as u64) + (segment_index * segment_size)
645}
646
647fn read_segment_header(
648 file: &mut File,
649 segment_size: u64,
650 segment_index: u64,
651) -> Result<WalSegmentHeader> {
652 let off = segment_header_offset(segment_size, segment_index);
653 file.seek(SeekFrom::Start(off))?;
654 let mut bytes = [0u8; WAL_SEGMENT_HEADER_SIZE];
655 file.read_exact(&mut bytes)?;
656 WalSegmentHeader::from_bytes(&bytes).map_err(|err| Error::CorruptedSegment {
657 segment_id: segment_index,
658 reason: format!("WAL segment header invalid: {err}"),
659 })
660}
661
662fn read_segment_header_allow_legacy(
663 file: &mut File,
664 segment_size: u64,
665 segment_index: u64,
666) -> Result<WalSegmentHeader> {
667 let off = segment_header_offset(segment_size, segment_index);
668 file.seek(SeekFrom::Start(off))?;
669 let mut bytes = [0u8; WAL_SEGMENT_HEADER_SIZE];
670 file.read_exact(&mut bytes)?;
671 WalSegmentHeader::from_bytes_allow_legacy(&bytes).map_err(|err| Error::CorruptedSegment {
672 segment_id: segment_index,
673 reason: format!("WAL segment header invalid: {err}"),
674 })
675}
676
677pub fn detect_wal_format_version(path: &Path, config: &WalConfig) -> Result<u16> {
679 let mut file = OpenOptions::new().read(true).open(path)?;
680 let (segment_size, _segment_data_len, _max_segments, _ring_len) = compute_ring_layout(config)?;
681 let header = read_segment_header_allow_legacy(&mut file, segment_size, 0)?;
682 Ok(header.version)
683}
684
685fn write_segment_header(
686 file: &mut File,
687 segment_size: u64,
688 segment_index: u64,
689 header: &WalSegmentHeader,
690) -> Result<()> {
691 let off = segment_header_offset(segment_size, segment_index);
692 file.seek(SeekFrom::Start(off))?;
693 file.write_all(&header.to_bytes())?;
694 Ok(())
695}
696
697fn ring_logical_to_physical(
698 logical_offset: u64,
699 segment_size: u64,
700 segment_data_len: u64,
701) -> Result<u64> {
702 if segment_data_len == 0 {
703 return Err(Error::InvalidFormat("segment data length is zero".into()));
704 }
705 let segment_index = logical_offset / segment_data_len;
706 let offset_in_segment = logical_offset % segment_data_len;
707 Ok(segment_header_offset(segment_size, segment_index)
708 + (WAL_SEGMENT_HEADER_SIZE as u64)
709 + offset_in_segment)
710}
711
712fn read_ring_bytes(
713 file: &mut File,
714 mut logical_offset: u64,
715 len: usize,
716 segment_size: u64,
717 segment_data_len: u64,
718 ring_len: u64,
719) -> Result<Vec<u8>> {
720 let mut out = Vec::with_capacity(len);
721 while out.len() < len {
722 let offset_in_segment = logical_offset % segment_data_len;
723 let remaining_in_segment = (segment_data_len - offset_in_segment) as usize;
724 let chunk_len = remaining_in_segment.min(len - out.len());
725 let phys = ring_logical_to_physical(logical_offset, segment_size, segment_data_len)?;
726 file.seek(SeekFrom::Start(phys))?;
727 let mut buf = vec![0u8; chunk_len];
728 file.read_exact(&mut buf)?;
729 out.extend_from_slice(&buf);
730 logical_offset = (logical_offset + (chunk_len as u64)) % ring_len;
731 }
732 Ok(out)
733}
734
735impl WalWriter {
736 fn write_ring(
737 &mut self,
738 mut logical_offset: u64,
739 mut data: &[u8],
740 entry_lsn: u64,
741 ) -> Result<()> {
742 while !data.is_empty() {
743 let offset_in_segment = logical_offset % self.segment_data_len;
744 if offset_in_segment == 0 {
745 let segment_index = logical_offset / self.segment_data_len;
746 let header = WalSegmentHeader::new(self.segment_id_base + segment_index, entry_lsn);
747 write_segment_header(&mut self.file, self.segment_size, segment_index, &header)?;
748 }
749 let remaining_in_segment = (self.segment_data_len - offset_in_segment) as usize;
750 let chunk_len = remaining_in_segment.min(data.len());
751
752 let phys =
753 ring_logical_to_physical(logical_offset, self.segment_size, self.segment_data_len)?;
754 self.file.seek(SeekFrom::Start(phys))?;
755 self.file.write_all(&data[..chunk_len])?;
756
757 logical_offset = (logical_offset + (chunk_len as u64)) % self.ring_len;
758 data = &data[chunk_len..];
759 }
760 Ok(())
761 }
762}
763
764fn persist_section_header(file: &mut File, offset: u64, header: &WalSectionHeader) -> Result<()> {
765 file.seek(SeekFrom::Start(offset))?;
766 file.write_all(&header.to_bytes())?;
767 Ok(())
768}
769
770fn load_section_header(file: &mut File, offset: u64) -> Result<WalSectionHeader> {
771 file.seek(SeekFrom::Start(offset))?;
772 let mut bytes = [0u8; WAL_SECTION_HEADER_SIZE];
773 file.read_exact(&mut bytes)?;
774 WalSectionHeader::from_bytes(&bytes)
775}
776
777#[cfg(all(test, not(target_arch = "wasm32")))]
778mod tests {
779 use super::*;
780 use std::fs::File;
781 use std::io::{Read, Seek, SeekFrom};
782 use tempfile::tempdir;
783
784 #[test]
785 fn segment_header_roundtrip() {
786 let header = WalSegmentHeader::new(42, 100);
787 let bytes = header.to_bytes();
788 assert_eq!(bytes.len(), WAL_SEGMENT_HEADER_SIZE);
789 let decoded = WalSegmentHeader::from_bytes(&bytes).unwrap();
790 assert_eq!(header.segment_id, decoded.segment_id);
791 assert_eq!(header.first_lsn, decoded.first_lsn);
792 assert_eq!(header.version, decoded.version);
793 }
794
795 #[test]
796 fn segment_header_crc_mismatch() {
797 let mut header = WalSegmentHeader::new(1, 1).to_bytes();
798 header[0] ^= 0xFF; let err = WalSegmentHeader::from_bytes(&header).unwrap_err();
800 assert!(matches!(err, Error::InvalidFormat(_)));
801 }
802
803 #[test]
804 fn segment_header_version_mismatch() {
805 let mut header = WalSegmentHeader::new(1, 1).to_bytes();
806 let corrupted_version = WAL_VERSION.wrapping_add(1);
809 header[4..6].copy_from_slice(&corrupted_version.to_le_bytes());
810 let err = WalSegmentHeader::from_bytes(&header).unwrap_err();
811 assert!(matches!(err, Error::InvalidFormat(_)));
812 }
813
814 #[test]
815 fn segment_header_crc_field_only_corruption() {
816 let mut header = WalSegmentHeader::new(1, 1).to_bytes();
817 header[22] ^= 0xFF;
820 let err = WalSegmentHeader::from_bytes(&header).unwrap_err();
821 assert!(matches!(err, Error::ChecksumMismatch));
822 }
823
824 #[test]
825 fn wal_entry_encode_decode_put() {
826 let entry = WalEntry::put(10, b"key".to_vec(), b"value".to_vec());
827 let encoded = entry.encode().unwrap();
828 let (decoded, consumed) = WalEntry::decode(&encoded).unwrap();
829 assert_eq!(consumed, encoded.len());
830 assert_eq!(decoded, entry);
831 }
832
833 #[test]
834 fn wal_entry_encode_decode_delete() {
835 let entry = WalEntry::delete(11, b"gone".to_vec());
836 let encoded = entry.encode().unwrap();
837 let (decoded, consumed) = WalEntry::decode(&encoded).unwrap();
838 assert_eq!(consumed, encoded.len());
839 assert_eq!(decoded, entry);
840 }
841
842 #[test]
843 fn wal_entry_crc_detects_corruption() {
844 let entry = WalEntry::put(12, b"k".to_vec(), b"v".to_vec());
845 let mut encoded = entry.encode().unwrap();
846 *encoded.last_mut().unwrap() ^= 0x10;
847 let err = WalEntry::decode(&encoded).unwrap_err();
848 assert!(matches!(err, Error::ChecksumMismatch));
849 }
850
851 #[test]
852 fn varint_helpers() {
853 let values = [
854 0u64,
855 1,
856 127,
857 128,
858 16384,
859 u32::MAX as u64,
860 u64::from(u32::MAX) + 1,
861 ];
862 for &v in &values {
863 let mut buf = Vec::new();
864 encode_varint(v, &mut buf);
865 let (decoded, read) = decode_varint(&buf).unwrap();
866 assert_eq!(decoded, v);
867 assert_eq!(read, buf.len());
868 assert_eq!(buf.len(), varint_len(v));
869 }
870 }
871
872 #[test]
873 fn wal_entry_crc_covers_header() {
874 let entry = WalEntry::put(20, b"key".to_vec(), b"value".to_vec());
875 let mut encoded = entry.encode().unwrap();
876 encoded[0] ^= 0xFF; let err = WalEntry::decode(&encoded).unwrap_err();
878 assert!(matches!(err, Error::ChecksumMismatch));
879 }
880
881 #[test]
882 fn wal_entry_retains_empty_value_put() {
883 let entry = WalEntry::put(30, b"key".to_vec(), Vec::new());
884 let encoded = entry.encode().unwrap();
885 let (decoded, _) = WalEntry::decode(&encoded).unwrap();
886 if let WalEntryPayload::Put { value, .. } = decoded.payload {
887 assert_eq!(value, Vec::<u8>::new());
888 } else {
889 panic!("expected Put payload");
890 }
891 }
892
893 #[test]
894 fn wal_batch_roundtrip() {
895 let ops = vec![
896 WalBatchOp {
897 op_type: WalOpType::Put,
898 key: b"a".to_vec(),
899 value: Some(b"1".to_vec()),
900 },
901 WalBatchOp {
902 op_type: WalOpType::Delete,
903 key: b"b".to_vec(),
904 value: None,
905 },
906 ];
907 let entry = WalEntry::batch(40, ops.clone());
908 let encoded = entry.encode().unwrap();
909 let (decoded, consumed) = WalEntry::decode(&encoded).unwrap();
910 assert_eq!(consumed, encoded.len());
911 assert_eq!(decoded.payload, WalEntryPayload::Batch(ops));
912 }
913
914 #[test]
915 fn wal_section_header_roundtrip() {
916 let header = WalSectionHeader::new(128, 4096, false);
917 let bytes = header.to_bytes();
918 let decoded = WalSectionHeader::from_bytes(&bytes).unwrap();
919 assert_eq!(decoded, header);
920 }
921
922 #[test]
923 fn wal_section_header_crc_field_only_corruption() {
924 let mut bytes = WalSectionHeader::new(128, 4096, false).to_bytes();
925 bytes[16] ^= 0xFF;
928 let err = WalSectionHeader::from_bytes(&bytes).unwrap_err();
929 assert!(matches!(err, Error::ChecksumMismatch));
930 }
931
932 #[test]
933 fn wal_section_header_data_corruption_detected_by_crc() {
934 let mut bytes = WalSectionHeader::new(128, 4096, false).to_bytes();
935 bytes[0] ^= 0xFF;
938 let err = WalSectionHeader::from_bytes(&bytes).unwrap_err();
939 assert!(matches!(err, Error::ChecksumMismatch));
940 }
941
942 #[test]
943 fn wal_writer_appends_and_updates_header() {
944 let dir = tempdir().unwrap();
945 let path = dir.path().join("wal");
946 let config = WalConfig {
947 segment_size: 4096,
948 max_segments: 1,
949 ..Default::default()
950 };
951 let mut writer = WalWriter::create(&path, config, 1, 100).unwrap();
952 let entry = WalEntry::put(100, b"key".to_vec(), b"value".to_vec());
953 let encoded_len = entry.encode().unwrap().len() as u64;
954
955 let offset = writer.append(&entry).unwrap();
956 assert_eq!(
957 offset,
958 (WAL_SECTION_HEADER_SIZE + WAL_SEGMENT_HEADER_SIZE) as u64
959 );
960
961 let mut file = File::open(&path).unwrap();
962 let mut hdr = [0u8; WAL_SECTION_HEADER_SIZE];
963 file.read_exact(&mut hdr).unwrap();
964 let section = WalSectionHeader::from_bytes(&hdr).unwrap();
965 assert_eq!(section.start_offset, 0);
966 assert_eq!(section.end_offset, encoded_len);
967 assert!(!section.is_full);
968
969 file.seek(SeekFrom::Start(offset)).unwrap();
970 let mut buf = vec![0u8; encoded_len as usize];
971 file.read_exact(&mut buf).unwrap();
972 let (decoded, consumed) = WalEntry::decode(&buf).unwrap();
973 assert_eq!(consumed, encoded_len as usize);
974 assert_eq!(decoded, entry);
975 }
976
977 #[test]
978 fn wal_writer_force_sync_calls_fsync() {
979 let dir = tempdir().unwrap();
980 let path = dir.path().join("wal_force_sync");
981 let config = WalConfig {
982 segment_size: 4096,
983 max_segments: 1,
984 sync_mode: SyncMode::NoSync,
985 };
986 let mut writer = WalWriter::create(&path, config, 1, 1).unwrap();
987 reset_sync_calls();
988
989 let entry = WalEntry::put(1, b"key".to_vec(), b"value".to_vec());
990 writer.append(&entry).unwrap();
991 assert_eq!(sync_calls(), 0);
992
993 writer.force_sync().unwrap();
994 assert_eq!(sync_calls(), 1);
995 }
996
997 #[test]
998 fn wal_writer_wraps_when_full() {
999 let dir = tempdir().unwrap();
1000 let path = dir.path().join("wal_wrap");
1001 let entry = WalEntry::put(1, b"a".to_vec(), b"1".to_vec());
1002 let entry_len = entry.encode().unwrap().len() as u64;
1003 let header_bytes = (WAL_SECTION_HEADER_SIZE + WAL_SEGMENT_HEADER_SIZE) as u64;
1004 let ring_len = (entry_len + (entry_len / 2)).max(entry_len + 1);
1005 let segment_size = (WAL_SEGMENT_HEADER_SIZE as u64) + ring_len;
1006 let config = WalConfig {
1007 segment_size: segment_size as usize,
1008 max_segments: 1,
1009 ..Default::default()
1010 };
1011
1012 let mut writer = WalWriter::create(&path, config, 2, 10).unwrap();
1013 let first_offset = writer.append(&entry).unwrap();
1014 assert_eq!(first_offset, header_bytes);
1015
1016 writer.advance_start(entry_len).unwrap();
1018
1019 let second_offset = writer.append(&entry).unwrap();
1020 assert_eq!(second_offset, header_bytes);
1021
1022 let mut file = File::open(&path).unwrap();
1023 let mut hdr = [0u8; WAL_SECTION_HEADER_SIZE];
1024 file.read_exact(&mut hdr).unwrap();
1025 let section = WalSectionHeader::from_bytes(&hdr).unwrap();
1026 assert_eq!(section.end_offset, entry_len);
1027 assert!(!section.is_full);
1028
1029 file.seek(SeekFrom::Start(second_offset)).unwrap();
1030 let mut buf = vec![0u8; entry_len as usize];
1031 file.read_exact(&mut buf).unwrap();
1032 let (decoded, consumed) = WalEntry::decode(&buf).unwrap();
1033 assert_eq!(consumed as u64, entry_len);
1034 assert_eq!(decoded, entry);
1035 }
1036
1037 #[test]
1038 fn wal_writer_refuses_overwrite_without_checkpoint() {
1039 let dir = tempdir().unwrap();
1040 let path = dir.path().join("wal_no_overwrite");
1041 let entry = WalEntry::put(1, b"a".to_vec(), b"1".to_vec());
1042 let entry_len = entry.encode().unwrap().len() as u64;
1043 let ring_len = entry_len + (entry_len / 2);
1044 let segment_size = (WAL_SEGMENT_HEADER_SIZE as u64) + ring_len;
1045 let config = WalConfig {
1046 segment_size: segment_size as usize,
1047 max_segments: 1,
1048 ..Default::default()
1049 };
1050 let mut writer = WalWriter::create(&path, config, 2, 10).unwrap();
1051 writer.append(&entry).unwrap();
1052 assert!(writer.append(&entry).is_err());
1053 }
1054
1055 #[test]
1056 fn wal_writer_advances_start_and_reclaims_space() {
1057 let dir = tempdir().unwrap();
1058 let path = dir.path().join("wal_reclaim");
1059 let entry = WalEntry::put(1, b"k".to_vec(), b"v".to_vec());
1060 let entry_len = entry.encode().unwrap().len() as u64;
1061 let header_bytes = (WAL_SECTION_HEADER_SIZE + WAL_SEGMENT_HEADER_SIZE) as u64;
1062 let ring_len = entry_len * 2;
1063 let segment_size = (WAL_SEGMENT_HEADER_SIZE as u64) + ring_len;
1064 let config = WalConfig {
1065 segment_size: segment_size as usize,
1066 max_segments: 1,
1067 ..Default::default()
1068 };
1069 let mut writer = WalWriter::create(&path, config, 5, 50).unwrap();
1070
1071 writer.append(&entry).unwrap();
1072 writer.advance_start(entry_len).unwrap();
1073
1074 let second = writer.append(&entry).unwrap();
1075 assert_eq!(second, header_bytes);
1076 }
1077
1078 #[test]
1079 fn wal_section_header_can_represent_full_buffer() {
1080 let header = WalSectionHeader::new(0, 0, true);
1081 let bytes = header.to_bytes();
1082 let decoded = WalSectionHeader::from_bytes(&bytes).unwrap();
1083 assert_eq!(decoded, header);
1084 }
1085
1086 #[test]
1087 fn wal_writer_persists_full_state_and_open_respects_it() {
1088 let dir = tempdir().unwrap();
1089 let path = dir.path().join("wal_full");
1090 let entry = WalEntry::put(1, b"k".to_vec(), vec![0; 32]);
1091 let entry_len = entry.encode().unwrap().len() as u64;
1092 let segment_size = (WAL_SEGMENT_HEADER_SIZE as u64) + entry_len;
1093 let config = WalConfig {
1094 segment_size: segment_size as usize,
1095 max_segments: 1,
1096 ..Default::default()
1097 };
1098
1099 let mut writer = WalWriter::create(&path, config.clone(), 7, 1).unwrap();
1100 writer.append(&entry).unwrap();
1101
1102 let mut file = File::open(&path).unwrap();
1103 let mut hdr = [0u8; WAL_SECTION_HEADER_SIZE];
1104 file.read_exact(&mut hdr).unwrap();
1105 let section = WalSectionHeader::from_bytes(&hdr).unwrap();
1106 assert!(section.is_full);
1107 assert_eq!(section.start_offset, 0);
1108 assert_eq!(section.end_offset, 0);
1109
1110 let mut reopened = WalWriter::open(&path, config).unwrap();
1111 assert!(reopened
1112 .append(&WalEntry::put(2, b"x".to_vec(), b"y".to_vec()))
1113 .is_err());
1114 }
1115
1116 #[test]
1117 fn wal_writer_multi_segment_entry_crosses_boundary_and_is_readable() {
1118 let dir = tempdir().unwrap();
1119 let path = dir.path().join("wal_multi");
1120 let entry = WalEntry::put(10, b"k".to_vec(), vec![0xAB; 64]);
1121 let encoded = entry.encode().unwrap();
1122 let segment_data_len = (encoded.len() - 1) as u64; let config = WalConfig {
1124 segment_size: (WAL_SEGMENT_HEADER_SIZE as u64 + segment_data_len) as usize,
1125 max_segments: 2,
1126 ..Default::default()
1127 };
1128
1129 let mut writer = WalWriter::create(&path, config.clone(), 1000, entry.lsn).unwrap();
1130 let start_physical = writer.append(&entry).unwrap();
1131 assert_eq!(
1132 start_physical,
1133 (WAL_SECTION_HEADER_SIZE + WAL_SEGMENT_HEADER_SIZE) as u64
1134 );
1135
1136 let _reopened = WalWriter::open(&path, config.clone()).unwrap();
1138
1139 fn read_ring_bytes(
1140 file: &mut File,
1141 mut logical_offset: u64,
1142 len: usize,
1143 segment_size: u64,
1144 segment_data_len: u64,
1145 ring_len: u64,
1146 ) -> Vec<u8> {
1147 let mut out = Vec::with_capacity(len);
1148 while out.len() < len {
1149 let offset_in_segment = logical_offset % segment_data_len;
1150 let remaining_in_segment = (segment_data_len - offset_in_segment) as usize;
1151 let chunk_len = remaining_in_segment.min(len - out.len());
1152 let phys = ring_logical_to_physical(logical_offset, segment_size, segment_data_len)
1153 .unwrap();
1154 file.seek(SeekFrom::Start(phys)).unwrap();
1155 let mut buf = vec![0u8; chunk_len];
1156 file.read_exact(&mut buf).unwrap();
1157 out.extend_from_slice(&buf);
1158 logical_offset = (logical_offset + chunk_len as u64) % ring_len;
1159 }
1160 out
1161 }
1162
1163 let mut file = File::open(&path).unwrap();
1164 let (segment_size, segment_data_len, _max_segments, ring_len) =
1165 compute_ring_layout(&config).unwrap();
1166
1167 let h0 = read_segment_header(&mut file, segment_size, 0).unwrap();
1169 let h1 = read_segment_header(&mut file, segment_size, 1).unwrap();
1170 assert_eq!(h0.segment_id, 1000);
1171 assert_eq!(h1.segment_id, 1001);
1172 assert_eq!(h0.first_lsn, entry.lsn);
1173 assert_eq!(h1.first_lsn, entry.lsn);
1174
1175 let bytes = read_ring_bytes(
1176 &mut file,
1177 0,
1178 encoded.len(),
1179 segment_size,
1180 segment_data_len,
1181 ring_len,
1182 );
1183 let (decoded, consumed) = WalEntry::decode(&bytes).unwrap();
1184 assert_eq!(consumed, encoded.len());
1185 assert_eq!(decoded, entry);
1186 }
1187}
1188
1189#[derive(Debug, Clone, PartialEq, Eq)]
1191pub struct WalSectionHeader {
1192 pub start_offset: u64,
1197 pub end_offset: u64,
1199 pub is_full: bool,
1204 pub crc32: u32,
1206}
1207
1208impl WalSectionHeader {
1209 const FULL_FLAG: u64 = 1u64 << 63;
1210 const OFFSET_MASK: u64 = !Self::FULL_FLAG;
1211
1212 pub fn new(start_offset: u64, end_offset: u64, is_full: bool) -> Self {
1214 let crc32 = Self::compute_crc(start_offset, end_offset, is_full);
1215 Self {
1216 start_offset,
1217 end_offset,
1218 is_full,
1219 crc32,
1220 }
1221 }
1222
1223 pub fn to_bytes(&self) -> [u8; WAL_SECTION_HEADER_SIZE] {
1225 let mut buf = [0u8; WAL_SECTION_HEADER_SIZE];
1226 let start = (self.start_offset & Self::OFFSET_MASK) | Self::FULL_FLAG;
1227 buf[0..8].copy_from_slice(&start.to_le_bytes());
1228 let mut end = self.end_offset & Self::OFFSET_MASK;
1229 if self.is_full {
1230 end |= Self::FULL_FLAG;
1231 }
1232 buf[8..16].copy_from_slice(&end.to_le_bytes());
1233 let crc32 = crc32fast::hash(&buf[0..16]);
1234 buf[16..20].copy_from_slice(&crc32.to_le_bytes());
1235 buf
1236 }
1237
1238 pub fn from_bytes(bytes: &[u8; WAL_SECTION_HEADER_SIZE]) -> Result<Self> {
1240 let raw_start = u64::from_le_bytes(bytes[0..8].try_into().expect("fixed slice length"));
1241 if (raw_start & Self::FULL_FLAG) == 0 {
1242 return Err(Error::InvalidFormat(
1243 "WAL section header valid marker missing".into(),
1244 ));
1245 }
1246 let start_offset = raw_start & Self::OFFSET_MASK;
1247 let raw_end = u64::from_le_bytes(bytes[8..16].try_into().expect("fixed slice length"));
1248 let stored_crc = u32::from_le_bytes(bytes[16..20].try_into().expect("fixed slice length"));
1249 let computed_crc = crc32fast::hash(&bytes[0..16]);
1250 if stored_crc != computed_crc {
1251 return Err(Error::ChecksumMismatch);
1252 }
1253
1254 let is_full = (raw_end & Self::FULL_FLAG) != 0;
1255 let end_offset = raw_end & Self::OFFSET_MASK;
1256 Ok(Self {
1257 start_offset,
1258 end_offset,
1259 is_full,
1260 crc32: stored_crc,
1261 })
1262 }
1263
1264 fn compute_crc(start_offset: u64, end_offset: u64, is_full: bool) -> u32 {
1265 let mut buf = [0u8; 16];
1266 let start = (start_offset & Self::OFFSET_MASK) | Self::FULL_FLAG;
1267 buf[0..8].copy_from_slice(&start.to_le_bytes());
1268 let mut end = end_offset & Self::OFFSET_MASK;
1269 if is_full {
1270 end |= Self::FULL_FLAG;
1271 }
1272 buf[8..16].copy_from_slice(&end.to_le_bytes());
1273 crc32fast::hash(&buf)
1274 }
1275
1276 pub(crate) fn refresh_crc(&mut self) {
1277 self.crc32 = Self::compute_crc(self.start_offset, self.end_offset, self.is_full);
1278 }
1279}
1280
1281#[derive(Debug)]
1283pub struct WalWriter {
1284 file: File,
1285 config: WalConfig,
1286 section_header: WalSectionHeader,
1287 segment_id_base: u64,
1288 segment_size: u64,
1289 segment_data_len: u64,
1290 ring_len: u64,
1291 used_bytes: u64,
1293 pending_sync: usize,
1295 last_sync: Instant,
1297}
1298
1299#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1301pub struct WalAppendStats {
1302 pub file_offset: u64,
1304 pub bytes_written: u64,
1306 pub sync_duration_ms: u64,
1308}
1309
1310impl WalWriter {
1311 pub fn create(
1319 path: &Path,
1320 config: WalConfig,
1321 segment_id_base: u64,
1322 first_lsn: u64,
1323 ) -> Result<Self> {
1324 let mut file = OpenOptions::new()
1325 .create(true)
1326 .read(true)
1327 .write(true)
1328 .truncate(true)
1329 .open(path)?;
1330
1331 let (segment_size, segment_data_len, max_segments, ring_len) =
1332 compute_ring_layout(&config)?;
1333 let wal_section_size = (WAL_SECTION_HEADER_SIZE as u64)
1334 .checked_add(
1335 segment_size
1336 .checked_mul(max_segments)
1337 .ok_or_else(|| Error::InvalidFormat("WAL section size overflow".into()))?,
1338 )
1339 .ok_or_else(|| Error::InvalidFormat("WAL section size overflow".into()))?;
1340 file.set_len(wal_section_size)?;
1341
1342 let section_header = WalSectionHeader::new(0, 0, false);
1343 persist_section_header(&mut file, 0, §ion_header)?;
1344
1345 for segment_index in 0..max_segments {
1346 let header_first_lsn = if segment_index == 0 { first_lsn } else { 0 };
1347 let header = WalSegmentHeader::new(segment_id_base + segment_index, header_first_lsn);
1348 write_segment_header(&mut file, segment_size, segment_index, &header)?;
1349 }
1350 sync_file(&file)?;
1351
1352 Ok(Self {
1353 file,
1354 config,
1355 section_header,
1356 segment_id_base,
1357 segment_size,
1358 segment_data_len,
1359 ring_len,
1360 used_bytes: 0,
1361 pending_sync: 0,
1362 last_sync: Instant::now(),
1363 })
1364 }
1365
1366 pub fn open(path: &Path, config: WalConfig) -> Result<Self> {
1368 let mut file = OpenOptions::new().read(true).write(true).open(path)?;
1369 let (segment_size, segment_data_len, max_segments, ring_len) =
1370 compute_ring_layout(&config)?;
1371
1372 let mut header_bytes = [0u8; WAL_SECTION_HEADER_SIZE];
1373 file.read_exact(&mut header_bytes)?;
1374 let section_header = WalSectionHeader::from_bytes(&header_bytes)?;
1375
1376 if section_header.start_offset >= ring_len {
1377 return Err(Error::InvalidFormat(
1378 "WAL start offset exceeds ring length".into(),
1379 ));
1380 }
1381 if section_header.end_offset >= ring_len {
1382 return Err(Error::InvalidFormat(
1383 "WAL end offset exceeds ring length".into(),
1384 ));
1385 }
1386
1387 let used_bytes = if section_header.is_full {
1388 ring_len
1389 } else {
1390 ring_distance(
1391 section_header.start_offset,
1392 section_header.end_offset,
1393 ring_len,
1394 )
1395 };
1396
1397 let mut segment_id_base: Option<u64> = None;
1399 for segment_index in 0..max_segments {
1400 let header = read_segment_header(&mut file, segment_size, segment_index)?;
1401 if segment_index == 0 {
1402 segment_id_base = Some(header.segment_id);
1403 } else if let Some(base) = segment_id_base {
1404 if header.segment_id != base + segment_index {
1405 return Err(Error::CorruptedSegment {
1406 segment_id: header.segment_id,
1407 reason: format!(
1408 "WAL segment_id sequence mismatch: expected {}, got {}",
1409 base + segment_index,
1410 header.segment_id
1411 ),
1412 });
1413 }
1414 }
1415 }
1416
1417 Ok(Self {
1418 file,
1419 config,
1420 section_header,
1421 segment_id_base: segment_id_base.unwrap_or(0),
1422 segment_size,
1423 segment_data_len,
1424 ring_len,
1425 used_bytes,
1426 pending_sync: 0,
1427 last_sync: Instant::now(),
1428 })
1429 }
1430
1431 pub fn advance_start(&mut self, new_start: u64) -> Result<()> {
1433 if new_start >= self.ring_len {
1434 return Err(Error::InvalidFormat(
1435 "WAL start offset exceeds ring length".into(),
1436 ));
1437 }
1438 let current = self.section_header.start_offset;
1439 let distance = ring_distance(current, new_start, self.ring_len);
1440 if distance > self.used_bytes {
1441 return Err(Error::InvalidFormat(
1442 "WAL start offset advances beyond written data".into(),
1443 ));
1444 }
1445 self.used_bytes -= distance;
1446 self.section_header.start_offset = new_start;
1447 self.section_header.is_full = false;
1448
1449 if self.used_bytes == 0 {
1450 self.section_header.start_offset = 0;
1452 self.section_header.end_offset = 0;
1453 self.section_header.is_full = false;
1454 }
1455
1456 self.section_header.refresh_crc();
1457 persist_section_header(&mut self.file, 0, &self.section_header)?;
1458 Ok(())
1459 }
1460
1461 pub fn truncate_tail_to(&mut self, new_end: u64) -> Result<()> {
1468 if new_end >= self.ring_len {
1469 return Err(Error::InvalidFormat(
1470 "WAL end offset exceeds ring length".into(),
1471 ));
1472 }
1473
1474 self.section_header.end_offset = new_end;
1475 self.section_header.is_full = false;
1476 self.used_bytes = ring_distance(self.section_header.start_offset, new_end, self.ring_len);
1477 self.section_header.refresh_crc();
1478 persist_section_header(&mut self.file, 0, &self.section_header)?;
1479 sync_file(&self.file)?;
1480 Ok(())
1481 }
1482
1483 pub fn append(&mut self, entry: &WalEntry) -> Result<u64> {
1485 Ok(self.append_with_stats(entry)?.file_offset)
1486 }
1487
1488 pub fn append_with_stats(&mut self, entry: &WalEntry) -> Result<WalAppendStats> {
1490 let encoded = entry.encode()?;
1491 let entry_len = encoded.len() as u64;
1492 if entry_len > self.ring_len {
1493 return Err(Error::InvalidFormat(
1494 "WAL entry exceeds ring capacity".into(),
1495 ));
1496 }
1497
1498 let free_space = self.ring_len - self.used_bytes;
1499 if entry_len > free_space {
1500 return Err(Error::InvalidFormat(
1501 "WAL buffer is full; cannot append entry".into(),
1502 ));
1503 }
1504
1505 if self.used_bytes == 0 && !self.section_header.is_full {
1507 self.section_header.start_offset = 0;
1508 self.section_header.end_offset = 0;
1509 }
1510
1511 let write_offset = self.section_header.end_offset;
1512 let file_offset =
1513 ring_logical_to_physical(write_offset, self.segment_size, self.segment_data_len)?;
1514 self.write_ring(write_offset, &encoded, entry.lsn)?;
1515
1516 let new_end = write_offset + entry_len;
1517 self.section_header.end_offset = new_end % self.ring_len;
1518 self.used_bytes += entry_len;
1519 self.section_header.is_full = self.used_bytes == self.ring_len;
1520 self.section_header.refresh_crc();
1521 persist_section_header(&mut self.file, 0, &self.section_header)?;
1522
1523 let sync_duration_ms = self.maybe_sync_with_stats(encoded.len())?;
1524 Ok(WalAppendStats {
1525 file_offset,
1526 bytes_written: entry_len,
1527 sync_duration_ms,
1528 })
1529 }
1530
1531 fn maybe_sync_with_stats(&mut self, bytes_written: usize) -> Result<u64> {
1532 match self.config.sync_mode {
1533 SyncMode::EveryWrite => {
1534 let start = Instant::now();
1535 sync_file(&self.file)?;
1536 let ms = start.elapsed().as_millis() as u64;
1537 self.pending_sync = 0;
1538 self.last_sync = Instant::now();
1539 Ok(ms)
1540 }
1541 SyncMode::BatchSync {
1542 max_batch_size,
1543 max_wait_ms,
1544 } => {
1545 self.pending_sync += bytes_written;
1546 let elapsed = self.last_sync.elapsed();
1547 let should_sync = self.pending_sync >= max_batch_size
1548 || elapsed >= Duration::from_millis(max_wait_ms);
1549 if should_sync {
1550 let start = Instant::now();
1551 sync_file(&self.file)?;
1552 let ms = start.elapsed().as_millis() as u64;
1553 self.pending_sync = 0;
1554 self.last_sync = Instant::now();
1555 return Ok(ms);
1556 }
1557 Ok(0)
1558 }
1559 SyncMode::NoSync => {
1560 Ok(0)
1562 }
1563 }
1564 }
1565
1566 pub fn force_sync(&mut self) -> Result<u64> {
1568 let start = Instant::now();
1569 sync_file(&self.file)?;
1570 let ms = start.elapsed().as_millis() as u64;
1571 self.pending_sync = 0;
1572 self.last_sync = Instant::now();
1573 Ok(ms)
1574 }
1575
1576 pub fn ring_len(&self) -> u64 {
1578 self.ring_len
1579 }
1580
1581 pub fn used_bytes(&self) -> u64 {
1583 self.used_bytes
1584 }
1585
1586 pub fn end_offset(&self) -> u64 {
1588 self.section_header.end_offset
1589 }
1590}
1591
1592#[derive(Debug, Clone, PartialEq, Eq)]
1594pub struct WalReplay {
1595 pub entries: Vec<WalEntry>,
1597 pub warnings: Vec<String>,
1599 pub stopped_at: Option<u64>,
1601 pub stop_reason: Option<String>,
1603}
1604
1605#[derive(Debug)]
1607pub struct WalReader {
1608 file: File,
1609 config: WalConfig,
1610 section_header: WalSectionHeader,
1611 segment_id_base: u64,
1612 segment_size: u64,
1613 segment_data_len: u64,
1614 ring_len: u64,
1615 used_bytes: u64,
1616}
1617
1618impl WalReader {
1619 pub fn open(path: &Path, config: WalConfig) -> Result<Self> {
1624 let mut file = OpenOptions::new().read(true).open(path)?;
1625
1626 let (segment_size, segment_data_len, max_segments, ring_len) =
1627 compute_ring_layout(&config)?;
1628 let wal_section_size = (WAL_SECTION_HEADER_SIZE as u64)
1629 .checked_add(
1630 segment_size
1631 .checked_mul(max_segments)
1632 .ok_or_else(|| Error::InvalidFormat("WAL section size overflow".into()))?,
1633 )
1634 .ok_or_else(|| Error::InvalidFormat("WAL section size overflow".into()))?;
1635
1636 let file_len = file.metadata()?.len();
1637 if file_len < wal_section_size {
1638 return Err(Error::InvalidFormat(
1639 "WAL file is smaller than configured section size".into(),
1640 ));
1641 }
1642
1643 let section_header = load_section_header(&mut file, 0)?;
1644 if section_header.start_offset >= ring_len {
1645 return Err(Error::InvalidFormat(
1646 "WAL start offset exceeds ring length".into(),
1647 ));
1648 }
1649 if section_header.end_offset >= ring_len {
1650 return Err(Error::InvalidFormat(
1651 "WAL end offset exceeds ring length".into(),
1652 ));
1653 }
1654 if section_header.is_full && section_header.start_offset != section_header.end_offset {
1655 return Err(Error::InvalidFormat(
1656 "WAL section header inconsistent: is_full=true but start_offset != end_offset"
1657 .into(),
1658 ));
1659 }
1660
1661 let used_bytes = if section_header.is_full {
1662 ring_len
1663 } else {
1664 ring_distance(
1665 section_header.start_offset,
1666 section_header.end_offset,
1667 ring_len,
1668 )
1669 };
1670
1671 let mut segment_id_base: Option<u64> = None;
1672 for segment_index in 0..max_segments {
1673 let header = read_segment_header(&mut file, segment_size, segment_index)?;
1674 if segment_index == 0 {
1675 segment_id_base = Some(header.segment_id);
1676 } else if let Some(base) = segment_id_base {
1677 if header.segment_id != base + segment_index {
1678 return Err(Error::CorruptedSegment {
1679 segment_id: header.segment_id,
1680 reason: format!(
1681 "WAL segment_id sequence mismatch: expected {}, got {}",
1682 base + segment_index,
1683 header.segment_id
1684 ),
1685 });
1686 }
1687 }
1688 }
1689
1690 Ok(Self {
1691 file,
1692 config,
1693 section_header,
1694 segment_id_base: segment_id_base.unwrap_or(0),
1695 segment_size,
1696 segment_data_len,
1697 ring_len,
1698 used_bytes,
1699 })
1700 }
1701
1702 pub(crate) fn open_allow_legacy(path: &Path, config: WalConfig) -> Result<Self> {
1704 let mut file = OpenOptions::new().read(true).open(path)?;
1705
1706 let (segment_size, segment_data_len, max_segments, ring_len) =
1707 compute_ring_layout(&config)?;
1708 let wal_section_size = (WAL_SECTION_HEADER_SIZE as u64)
1709 .checked_add(
1710 segment_size
1711 .checked_mul(max_segments)
1712 .ok_or_else(|| Error::InvalidFormat("WAL section size overflow".into()))?,
1713 )
1714 .ok_or_else(|| Error::InvalidFormat("WAL section size overflow".into()))?;
1715
1716 let file_len = file.metadata()?.len();
1717 if file_len < wal_section_size {
1718 return Err(Error::InvalidFormat(
1719 "WAL file is smaller than configured section size".into(),
1720 ));
1721 }
1722
1723 let section_header = load_section_header(&mut file, 0)?;
1724 if section_header.start_offset >= ring_len {
1725 return Err(Error::InvalidFormat(
1726 "WAL start offset exceeds ring length".into(),
1727 ));
1728 }
1729 if section_header.end_offset >= ring_len {
1730 return Err(Error::InvalidFormat(
1731 "WAL end offset exceeds ring length".into(),
1732 ));
1733 }
1734 if section_header.is_full && section_header.start_offset != section_header.end_offset {
1735 return Err(Error::InvalidFormat(
1736 "WAL section header inconsistent: is_full=true but start_offset != end_offset"
1737 .into(),
1738 ));
1739 }
1740
1741 let used_bytes = if section_header.is_full {
1742 ring_len
1743 } else {
1744 ring_distance(
1745 section_header.start_offset,
1746 section_header.end_offset,
1747 ring_len,
1748 )
1749 };
1750
1751 let mut segment_id_base: Option<u64> = None;
1752 for segment_index in 0..max_segments {
1753 let header = read_segment_header_allow_legacy(&mut file, segment_size, segment_index)?;
1754 if segment_index == 0 {
1755 segment_id_base = Some(header.segment_id);
1756 } else if let Some(base) = segment_id_base {
1757 if header.segment_id != base + segment_index {
1758 return Err(Error::CorruptedSegment {
1759 segment_id: header.segment_id,
1760 reason: format!(
1761 "WAL segment_id sequence mismatch: expected {}, got {}",
1762 base + segment_index,
1763 header.segment_id
1764 ),
1765 });
1766 }
1767 }
1768 }
1769
1770 Ok(Self {
1771 file,
1772 config,
1773 section_header,
1774 segment_id_base: segment_id_base.unwrap_or(0),
1775 segment_size,
1776 segment_data_len,
1777 ring_len,
1778 used_bytes,
1779 })
1780 }
1781
1782 pub fn replay(&mut self) -> Result<WalReplay> {
1792 self.replay_with_resync(0)
1793 }
1794
1795 pub fn replay_with_resync(&mut self, max_resync_scan_bytes: usize) -> Result<WalReplay> {
1801 let mut entries = Vec::new();
1802 let mut warnings = Vec::new();
1803 let mut cursor = self.section_header.start_offset;
1804 let mut remaining = self.used_bytes;
1805 let mut last_lsn: Option<u64> = None;
1806
1807 while remaining > 0 {
1808 if remaining < WAL_ENTRY_FIXED_HEADER as u64 {
1809 return Ok(WalReplay {
1810 entries,
1811 warnings,
1812 stopped_at: Some(cursor),
1813 stop_reason: Some("WAL entry header truncated".into()),
1814 });
1815 }
1816
1817 let header_bytes = read_ring_bytes(
1818 &mut self.file,
1819 cursor,
1820 WAL_ENTRY_FIXED_HEADER,
1821 self.segment_size,
1822 self.segment_data_len,
1823 self.ring_len,
1824 )?;
1825 let payload_and_crc_len =
1826 u32::from_le_bytes(header_bytes[8..12].try_into().expect("fixed slice length"))
1827 as u64;
1828 let total_len = (WAL_ENTRY_FIXED_HEADER as u64)
1829 .checked_add(payload_and_crc_len)
1830 .ok_or_else(|| Error::InvalidFormat("WAL entry length overflow".into()))?;
1831
1832 if total_len == 0 || total_len > self.ring_len {
1833 return Ok(WalReplay {
1834 entries,
1835 warnings,
1836 stopped_at: Some(cursor),
1837 stop_reason: Some("WAL entry length is invalid".into()),
1838 });
1839 }
1840
1841 if total_len > remaining {
1842 return Ok(WalReplay {
1843 entries,
1844 warnings,
1845 stopped_at: Some(cursor),
1846 stop_reason: Some("WAL entry truncated at tail".into()),
1847 });
1848 }
1849
1850 let entry_bytes = read_ring_bytes(
1851 &mut self.file,
1852 cursor,
1853 total_len as usize,
1854 self.segment_size,
1855 self.segment_data_len,
1856 self.ring_len,
1857 )?;
1858
1859 let decoded = match WalEntry::decode(&entry_bytes) {
1860 Ok((entry, consumed)) => {
1861 if consumed as u64 != total_len {
1862 return Ok(WalReplay {
1863 entries,
1864 warnings,
1865 stopped_at: Some(cursor),
1866 stop_reason: Some("WAL entry decode consumed unexpected length".into()),
1867 });
1868 }
1869 entry
1870 }
1871 Err(err) => {
1872 if max_resync_scan_bytes == 0 {
1873 return Ok(WalReplay {
1874 entries,
1875 warnings,
1876 stopped_at: Some(cursor),
1877 stop_reason: Some(format!("WAL entry decode failed: {err}")),
1878 });
1879 }
1880
1881 let max_scan = max_resync_scan_bytes.min(remaining as usize);
1882 let mut resynced: Option<(u64, WalEntry, u64, u64)> = None;
1883 for delta in 1..=max_scan {
1884 if remaining < (delta as u64) + (WAL_ENTRY_FIXED_HEADER as u64) {
1885 break;
1886 }
1887 let candidate = (cursor + (delta as u64)) % self.ring_len;
1888 let header = read_ring_bytes(
1889 &mut self.file,
1890 candidate,
1891 WAL_ENTRY_FIXED_HEADER,
1892 self.segment_size,
1893 self.segment_data_len,
1894 self.ring_len,
1895 )?;
1896 let payload_and_crc_len = u32::from_le_bytes(
1897 header[8..12].try_into().expect("fixed slice length"),
1898 ) as u64;
1899 let cand_total_len = (WAL_ENTRY_FIXED_HEADER as u64)
1900 .checked_add(payload_and_crc_len)
1901 .ok_or_else(|| {
1902 Error::InvalidFormat("WAL entry length overflow".into())
1903 })?;
1904 if cand_total_len == 0 || cand_total_len > self.ring_len {
1905 continue;
1906 }
1907 let remaining_after_skip = remaining - (delta as u64);
1908 if cand_total_len > remaining_after_skip {
1909 continue;
1910 }
1911
1912 let bytes = read_ring_bytes(
1913 &mut self.file,
1914 candidate,
1915 cand_total_len as usize,
1916 self.segment_size,
1917 self.segment_data_len,
1918 self.ring_len,
1919 )?;
1920 let Ok((entry, consumed)) = WalEntry::decode(&bytes) else {
1921 continue;
1922 };
1923 if consumed as u64 != cand_total_len {
1924 continue;
1925 }
1926 if let Some(prev) = last_lsn {
1927 if entry.lsn <= prev {
1928 continue;
1929 }
1930 }
1931 resynced = Some((candidate, entry, cand_total_len, delta as u64));
1932 break;
1933 }
1934
1935 if let Some((candidate, entry, cand_total_len, skipped)) = resynced {
1936 warnings.push(format!(
1937 "WAL replay resynchronized: skipped {skipped} bytes at offset {cursor} -> {candidate}"
1938 ));
1939 entries.push(entry);
1940 cursor = (candidate + cand_total_len) % self.ring_len;
1941 remaining -= skipped + cand_total_len;
1942 continue;
1943 }
1944
1945 return Ok(WalReplay {
1946 entries,
1947 warnings,
1948 stopped_at: Some(cursor),
1949 stop_reason: Some(format!(
1950 "WAL entry decode failed and resync could not find next boundary: {err}"
1951 )),
1952 });
1953 }
1954 };
1955
1956 if let Some(prev) = last_lsn {
1957 if decoded.lsn <= prev {
1958 return Ok(WalReplay {
1959 entries,
1960 warnings,
1961 stopped_at: Some(cursor),
1962 stop_reason: Some("WAL LSN is not strictly increasing".into()),
1963 });
1964 }
1965 }
1966 last_lsn = Some(decoded.lsn);
1967
1968 entries.push(decoded);
1969 cursor = (cursor + total_len) % self.ring_len;
1972 remaining -= total_len;
1973 }
1974
1975 Ok(WalReplay {
1976 entries,
1977 warnings,
1978 stopped_at: None,
1979 stop_reason: None,
1980 })
1981 }
1982
1983 pub fn section_header(&self) -> &WalSectionHeader {
1985 &self.section_header
1986 }
1987
1988 pub fn ring_len(&self) -> u64 {
1990 self.ring_len
1991 }
1992
1993 pub fn segment_id_base(&self) -> u64 {
1995 self.segment_id_base
1996 }
1997
1998 pub fn config(&self) -> &WalConfig {
2000 &self.config
2001 }
2002}
2003
2004#[cfg(all(test, not(target_arch = "wasm32")))]
2005mod reader {
2006 use super::*;
2007 use tempfile::tempdir;
2008
2009 #[test]
2010 fn wal_reader_replays_entries_in_order() {
2011 let dir = tempdir().unwrap();
2012 let path = dir.path().join("wal_reader_basic");
2013 let config = WalConfig {
2014 segment_size: 4096,
2015 max_segments: 1,
2016 ..Default::default()
2017 };
2018
2019 let mut writer = WalWriter::create(&path, config.clone(), 10, 1).unwrap();
2020 let e1 = WalEntry::put(1, b"a".to_vec(), b"1".to_vec());
2021 let e2 = WalEntry::delete(2, b"b".to_vec());
2022 writer.append(&e1).unwrap();
2023 writer.append(&e2).unwrap();
2024
2025 let mut reader = WalReader::open(&path, config).unwrap();
2026 let replay = reader.replay().unwrap();
2027 assert_eq!(replay.stop_reason, None);
2028 assert!(replay.warnings.is_empty());
2029 assert_eq!(replay.entries, vec![e1, e2]);
2030 }
2031
2032 #[test]
2033 fn wal_reader_skips_entries_before_start_offset() {
2034 let dir = tempdir().unwrap();
2035 let path = dir.path().join("wal_reader_start");
2036 let config = WalConfig {
2037 segment_size: 4096,
2038 max_segments: 1,
2039 ..Default::default()
2040 };
2041
2042 let mut writer = WalWriter::create(&path, config.clone(), 10, 1).unwrap();
2043 let e1 = WalEntry::put(1, b"a".to_vec(), b"1".to_vec());
2044 let e2 = WalEntry::put(2, b"b".to_vec(), b"2".to_vec());
2045 let e1_len = e1.encode().unwrap().len() as u64;
2046 writer.append(&e1).unwrap();
2047 writer.append(&e2).unwrap();
2048
2049 writer.advance_start(e1_len).unwrap();
2050
2051 let mut reader = WalReader::open(&path, config).unwrap();
2052 let replay = reader.replay().unwrap();
2053 assert_eq!(replay.stop_reason, None);
2054 assert!(replay.warnings.is_empty());
2055 assert_eq!(replay.entries, vec![e2]);
2056 }
2057
2058 #[test]
2059 fn wal_reader_stops_on_corrupt_tail_and_returns_prefix() {
2060 let dir = tempdir().unwrap();
2061 let path = dir.path().join("wal_reader_corrupt_tail");
2062 let config = WalConfig {
2063 segment_size: 4096,
2064 max_segments: 1,
2065 ..Default::default()
2066 };
2067
2068 let mut writer = WalWriter::create(&path, config.clone(), 10, 1).unwrap();
2069 let e1 = WalEntry::put(1, b"a".to_vec(), b"1".to_vec());
2070 let e2 = WalEntry::put(2, b"b".to_vec(), b"2".to_vec());
2071 let e2_len = e2.encode().unwrap().len() as u64;
2072 writer.append(&e1).unwrap();
2073
2074 let start_of_e2 = writer.section_header.end_offset;
2076 let mut corrupted_header = writer.section_header.clone();
2077 corrupted_header.end_offset = (start_of_e2 + e2_len) % writer.ring_len;
2078 persist_section_header(&mut writer.file, 0, &corrupted_header).unwrap();
2079
2080 let mut reader = WalReader::open(&path, config).unwrap();
2081 let replay = reader.replay().unwrap();
2082 assert_eq!(replay.entries, vec![e1]);
2083 assert!(replay.stop_reason.is_some());
2084 }
2085
2086 #[test]
2087 fn wal_reader_replays_entry_crossing_segment_boundary() {
2088 let dir = tempdir().unwrap();
2089 let path = dir.path().join("wal_reader_multi");
2090 let entry = WalEntry::put(10, b"k".to_vec(), vec![0xCD; 64]);
2091 let encoded = entry.encode().unwrap();
2092 let segment_data_len = (encoded.len() - 1) as u64; let config = WalConfig {
2094 segment_size: (WAL_SEGMENT_HEADER_SIZE as u64 + segment_data_len) as usize,
2095 max_segments: 2,
2096 ..Default::default()
2097 };
2098
2099 let mut writer = WalWriter::create(&path, config.clone(), 2000, entry.lsn).unwrap();
2100 writer.append(&entry).unwrap();
2101
2102 let mut reader = WalReader::open(&path, config).unwrap();
2103 let replay = reader.replay().unwrap();
2104 assert_eq!(replay.stop_reason, None);
2105 assert!(replay.warnings.is_empty());
2106 assert_eq!(replay.entries, vec![entry]);
2107 }
2108
2109 #[test]
2110 fn wal_reader_open_rejects_inconsistent_full_flag() {
2111 let dir = tempdir().unwrap();
2112 let path = dir.path().join("wal_reader_inconsistent_full");
2113 let entry = WalEntry::put(1, b"k".to_vec(), vec![0; 32]);
2114 let entry_len = entry.encode().unwrap().len() as u64;
2115 let segment_size = (WAL_SEGMENT_HEADER_SIZE as u64) + entry_len;
2116 let config = WalConfig {
2117 segment_size: segment_size as usize,
2118 max_segments: 1,
2119 ..Default::default()
2120 };
2121
2122 let mut writer = WalWriter::create(&path, config.clone(), 1, 1).unwrap();
2123 writer.append(&entry).unwrap();
2124 assert!(writer.section_header.is_full);
2125
2126 let mut bad = writer.section_header.clone();
2127 bad.start_offset = 1;
2128 persist_section_header(&mut writer.file, 0, &bad).unwrap();
2129
2130 let err = WalReader::open(&path, config).unwrap_err();
2131 matches!(err, Error::InvalidFormat(_));
2132 }
2133
2134 #[test]
2135 fn wal_reader_can_resync_when_start_offset_is_misaligned() {
2136 let dir = tempdir().unwrap();
2137 let path = dir.path().join("wal_reader_resync");
2138 let config = WalConfig {
2139 segment_size: 4096,
2140 max_segments: 1,
2141 ..Default::default()
2142 };
2143
2144 let mut writer = WalWriter::create(&path, config.clone(), 10, 1).unwrap();
2145 let e1 = WalEntry::put(1, b"a".to_vec(), vec![0xAA; 128]);
2146 let e2 = WalEntry::put(2, b"b".to_vec(), vec![0xBB; 128]);
2147 writer.append(&e1).unwrap();
2148 writer.append(&e2).unwrap();
2149
2150 let mut misaligned = writer.section_header.clone();
2151 misaligned.start_offset = (misaligned.start_offset + 1) % writer.ring_len;
2152 persist_section_header(&mut writer.file, 0, &misaligned).unwrap();
2153
2154 let mut reader = WalReader::open(&path, config).unwrap();
2155 let replay = reader.replay_with_resync(4096).unwrap();
2156 assert_eq!(replay.entries, vec![e2]);
2157 assert!(replay.stop_reason.is_none());
2158 assert!(!replay.warnings.is_empty());
2159 }
2160}