Skip to main content

alopex_core/lsm/
wal.rs

1//! Write-Ahead Log (WAL) primitives for the LSM file mode.
2//!
3//! This module defines the on-disk WAL segment header and entry layout used by
4//! the circular buffer inside a single `.alopex` file. Writer/reader
5//! implementations build on top of these primitives.
6
7#[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
43/// WAL file magic ("AWAL").
44pub const WAL_MAGIC: [u8; 4] = *b"AWAL";
45/// WAL format version (v0.4.x).
46pub const WAL_FORMAT_VERSION_V04: u16 = 1;
47/// WAL format version (v0.5.0).
48pub const WAL_FORMAT_VERSION: u16 = 2;
49/// WAL format version used by this binary (uint16).
50pub const WAL_VERSION: u16 = WAL_FORMAT_VERSION;
51/// Fixed segment header size (bytes).
52pub const WAL_SEGMENT_HEADER_SIZE: usize = 28;
53/// WAL section header size (bytes) for circular buffer start/end offsets.
54pub const WAL_SECTION_HEADER_SIZE: usize = 20;
55/// Fixed overhead for an entry before payload bytes (LSN + length).
56pub const WAL_ENTRY_FIXED_HEADER: usize = 8 + 4;
57
58/// Top-level WAL entry type.
59#[repr(u8)]
60#[derive(Debug, Clone, Copy, PartialEq, Eq)]
61pub enum WalEntryType {
62    /// Put operation (single).
63    Put = 0,
64    /// Delete (tombstone) operation (single).
65    Delete = 1,
66    /// Batched operations encoded in a single entry.
67    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/// Operation type inside a batch payload.
86#[repr(u8)]
87#[derive(Debug, Clone, Copy, PartialEq, Eq)]
88pub enum WalOpType {
89    /// Put operation.
90    Put = 0,
91    /// Delete (tombstone) operation.
92    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/// Configuration for WAL circular buffer segments.
110#[derive(Debug, Clone, PartialEq, Eq)]
111pub struct WalConfig {
112    /// Segment size in bytes (default: 64MB).
113    pub segment_size: usize,
114    /// Maximum number of segments (default: 8).
115    pub max_segments: usize,
116    /// Fsync strategy.
117    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/// Sync policy for WAL writes.
131#[derive(Debug, Clone, PartialEq, Eq)]
132pub enum SyncMode {
133    /// fsync on every append (safest).
134    EveryWrite,
135    /// fsync periodically based on batch size or timeout.
136    ///
137    /// Note: this mode can leave the WAL's section/segment headers temporarily stale after a crash.
138    /// Recovery must validate headers and fall back safely on inconsistencies.
139    BatchSync {
140        /// Max bytes to buffer before fsync.
141        max_batch_size: usize,
142        /// Max wait time (ms) before forcing fsync.
143        max_wait_ms: u64,
144    },
145    /// Rely on OS buffering (fastest, least safe).
146    ///
147    /// Note: after a crash, recent appends and/or header updates may be missing.
148    NoSync,
149}
150
151/// Fixed-size WAL segment header (28 bytes).
152#[derive(Debug, Clone, PartialEq, Eq)]
153pub struct WalSegmentHeader {
154    /// Format version.
155    pub version: u16,
156    /// Segment identifier.
157    pub segment_id: u64,
158    /// First LSN stored in this segment (including fragments).
159    ///
160    /// Entries may cross segment boundaries, so the data region may begin with a fragment of an
161    /// entry with this LSN. Do not assume decoding can start at the segment's first byte.
162    pub first_lsn: u64,
163    /// CRC32 for bytes [0..22).
164    pub crc32: u32,
165    /// Reserved field (kept for alignment/forward-compat).
166    pub reserved: u16,
167}
168
169impl WalSegmentHeader {
170    /// Create a new header with computed CRC.
171    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    /// Serialize the header to fixed-size bytes (28B).
183    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    /// Deserialize and validate a header.
195    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    /// Deserialize and validate a header, allowing legacy format versions.
229    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    /// Returns true if the header version is older than current WAL format.
267    pub fn is_legacy_format(&self) -> bool {
268        self.version < WAL_FORMAT_VERSION
269    }
270}
271
272/// A single batch operation.
273#[derive(Debug, Clone, PartialEq, Eq)]
274pub struct WalBatchOp {
275    /// Operation type.
276    pub op_type: WalOpType,
277    /// Key bytes.
278    pub key: Vec<u8>,
279    /// Value bytes for Put (may be empty). None for Delete.
280    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/// WAL entry payload variants.
294#[derive(Debug, Clone, PartialEq, Eq)]
295pub enum WalEntryPayload {
296    /// Single Put entry.
297    Put {
298        /// Key bytes.
299        key: Vec<u8>,
300        /// Value bytes (may be empty).
301        value: Vec<u8>,
302    },
303    /// Single Delete entry.
304    Delete {
305        /// Key bytes.
306        key: Vec<u8>,
307    },
308    /// Batched Put/Delete operations.
309    Batch(Vec<WalBatchOp>),
310}
311
312/// WAL entry.
313#[derive(Debug, Clone, PartialEq, Eq)]
314pub struct WalEntry {
315    /// Monotonic log sequence number.
316    pub lsn: u64,
317    /// Entry payload.
318    pub payload: WalEntryPayload,
319}
320
321impl WalEntry {
322    /// Construct a Put entry (value may be empty).
323    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    /// Construct a Delete entry.
331    pub fn delete(lsn: u64, key: Vec<u8>) -> Self {
332        Self {
333            lsn,
334            payload: WalEntryPayload::Delete { key },
335        }
336    }
337
338    /// Construct a Batch entry.
339    pub fn batch(lsn: u64, operations: Vec<WalBatchOp>) -> Self {
340        Self {
341            lsn,
342            payload: WalEntryPayload::Batch(operations),
343        }
344    }
345
346    /// Total encoded length including header and CRC.
347    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 // + CRC32
364    }
365
366    /// Encode the entry to bytes.
367    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        // CRC over the entire entry except the CRC field itself (header + body).
413        let crc = crc32fast::hash(&out);
414        out.extend_from_slice(&crc.to_le_bytes());
415        Ok(out)
416    }
417
418    /// Decode a single entry from the provided buffer, returning the entry and bytes consumed.
419    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]; // up to CRC field (22 bytes)
604    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
677/// Detect the WAL format version from the first segment header.
678pub 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; // break magic
799        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        // Corrupt only the version field (bytes 4..6); magic and crc32 remain as originally
807        // computed, so the version check must fail before the CRC check is reached.
808        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        // Corrupt only the crc32 field (bytes 22..26); magic/version/segment_id/first_lsn stay
818        // valid, isolating the failure to the checksum comparison.
819        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; // corrupt LSN byte
877        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        // Corrupt only the crc32 field (bytes 16..20); start_offset/end_offset payload data
926        // stays untouched, isolating the failure to the checksum comparison itself.
927        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        // Corrupt the CRC-covered payload (start_offset within bytes 0..8) while leaving the
936        // stored crc32 field untouched, so the mismatch is detected on the data side.
937        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        // Simulate checkpoint freeing the first entry.
1017        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; // force boundary crossing
1123        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        // Re-open and validate segment headers + ring contents.
1137        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        // First two segment headers should have been initialized/updated with the base id.
1168        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/// WAL section header for circular buffer bookkeeping.
1190#[derive(Debug, Clone, PartialEq, Eq)]
1191pub struct WalSectionHeader {
1192    /// Read pointer offset (inclusive), as a logical ring offset.
1193    ///
1194    /// The serialized form stores a validity marker in the MSB, so the logical offset is masked
1195    /// with [`WalSectionHeader::OFFSET_MASK`] when persisted and decoded.
1196    pub start_offset: u64,
1197    /// Write pointer offset (exclusive), as a logical ring offset.
1198    pub end_offset: u64,
1199    /// Whether the buffer is full (disambiguates `start_offset == end_offset`).
1200    ///
1201    /// This is persisted by storing a flag in the MSB of the serialized `end_offset`. Therefore,
1202    /// the ring length must stay within [`WalSectionHeader::OFFSET_MASK`].
1203    pub is_full: bool,
1204    /// CRC32 for bytes [0..16).
1205    pub crc32: u32,
1206}
1207
1208impl WalSectionHeader {
1209    const FULL_FLAG: u64 = 1u64 << 63;
1210    const OFFSET_MASK: u64 = !Self::FULL_FLAG;
1211
1212    /// Create a new section header with computed CRC.
1213    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    /// Serialize to 20 bytes (start/end offsets + CRC32).
1224    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    /// Deserialize from 20 bytes and validate CRC32.
1239    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/// WAL writer for circular buffer sections.
1282#[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    /// Bytes currently used in the buffer.
1292    used_bytes: u64,
1293    /// Bytes written since last fsync (BatchSync mode).
1294    pending_sync: usize,
1295    /// Timestamp of last fsync (BatchSync mode).
1296    last_sync: Instant,
1297}
1298
1299/// WAL 追記の統計(メトリクス用)。
1300#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1301pub struct WalAppendStats {
1302    /// 実際に書き込んだファイルオフセット(物理)。
1303    pub file_offset: u64,
1304    /// 書き込んだバイト数(エントリ全体、CRC 含む)。
1305    pub bytes_written: u64,
1306    /// append 内で fsync が発生した場合の所要時間(ms)。発生しなければ 0。
1307    pub sync_duration_ms: u64,
1308}
1309
1310impl WalWriter {
1311    /// Create a new WAL writer and initialize section + segment headers.
1312    ///
1313    /// The WAL section size is `segment_size * max_segments` and each segment starts with a
1314    /// [`WalSegmentHeader`], followed by the segment's data region.
1315    ///
1316    /// `segment_id_base` is used to assign stable per-segment IDs as
1317    /// `segment_id_base + segment_index`.
1318    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, &section_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    /// Open an existing WAL writer, resuming from the stored end offset.
1367    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        // Validate all segment headers (magic/version/crc).
1398        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    /// Advance the read pointer (start_offset) after a checkpoint.
1432    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            // Prefer reusing the reclaimed region immediately by resetting pointers.
1451            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    /// Truncate the logical WAL tail after recovery stopped at the last valid boundary.
1462    ///
1463    /// CORE-5.2 recovery replays the valid prefix and stops on a corrupt/incomplete
1464    /// tail entry. Persisting the repaired end offset keeps future appends from
1465    /// landing behind that corrupt tail, which would otherwise hide newly-written
1466    /// records on the next replay.
1467    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    /// Append a WAL entry into the circular buffer, returning the file offset written.
1484    pub fn append(&mut self, entry: &WalEntry) -> Result<u64> {
1485        Ok(self.append_with_stats(entry)?.file_offset)
1486    }
1487
1488    /// WAL へ追記しつつ、メトリクス用の統計も返す。
1489    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 the buffer is empty, always start from the beginning of the data region.
1506        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                // no-op
1561                Ok(0)
1562            }
1563        }
1564    }
1565
1566    /// Force fsync regardless of configured SyncMode.
1567    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    /// Total logical ring length in bytes.
1577    pub fn ring_len(&self) -> u64 {
1578        self.ring_len
1579    }
1580
1581    /// Bytes currently used in the WAL ring buffer.
1582    pub fn used_bytes(&self) -> u64 {
1583        self.used_bytes
1584    }
1585
1586    /// Current end offset in the logical ring.
1587    pub fn end_offset(&self) -> u64 {
1588        self.section_header.end_offset
1589    }
1590}
1591
1592/// Result of WAL replay.
1593#[derive(Debug, Clone, PartialEq, Eq)]
1594pub struct WalReplay {
1595    /// Successfully decoded entries, in order.
1596    pub entries: Vec<WalEntry>,
1597    /// Non-fatal warnings recorded during replay (e.g. resynchronization occurred).
1598    pub warnings: Vec<String>,
1599    /// If replay stopped early, the logical ring offset where it stopped.
1600    pub stopped_at: Option<u64>,
1601    /// If replay stopped early, the reason (e.g. checksum mismatch / truncated entry).
1602    pub stop_reason: Option<String>,
1603}
1604
1605/// WAL reader for crash recovery.
1606#[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    /// Open an existing WAL section for replay.
1620    ///
1621    /// The reader replays entries from `start_offset` (inclusive) up to `end_offset` (exclusive)
1622    /// within the logical ring.
1623    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    /// Open an existing WAL section for replay, allowing legacy format versions.
1703    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    /// Replay entries in the WAL section.
1783    ///
1784    /// If a corrupted or incomplete entry is detected, replay stops and returns the entries
1785    /// decoded so far with `stop_reason` populated. This supports crash recovery where the tail
1786    /// of the WAL may be partially written.
1787    ///
1788    /// Note: this method assumes `start_offset` is aligned to an entry boundary (i.e. checkpoints
1789    /// advance in entry-sized increments). If you need best-effort recovery from a potentially
1790    /// misaligned `start_offset`, use [`WalReader::replay_with_resync`].
1791    pub fn replay(&mut self) -> Result<WalReplay> {
1792        self.replay_with_resync(0)
1793    }
1794
1795    /// Replay entries with optional resynchronization.
1796    ///
1797    /// If `max_resync_scan_bytes > 0`, upon a decode failure at the current cursor, the reader
1798    /// scans forward up to that many bytes (bounded by remaining bytes) to find the next valid
1799    /// entry boundary by validating entry framing and CRC.
1800    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            // remaining covers the logical region [start_offset, end_offset) with wrap support.
1970            // `total_len` is guaranteed to fit within `remaining` above.
1971            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    /// Current section header (start/end/is_full).
1984    pub fn section_header(&self) -> &WalSectionHeader {
1985        &self.section_header
1986    }
1987
1988    /// Total logical ring length in bytes.
1989    pub fn ring_len(&self) -> u64 {
1990        self.ring_len
1991    }
1992
1993    /// Segment ID base (segment_id = base + index).
1994    pub fn segment_id_base(&self) -> u64 {
1995        self.segment_id_base
1996    }
1997
1998    /// Configuration used by this reader.
1999    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        // Simulate a crash where the section header advanced for e2, but its bytes are not valid.
2075        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; // force boundary crossing
2093        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}