Skip to main content

rhiza_log/
lib.rs

1use std::{
2    fmt, fs,
3    io::Write,
4    path::{Path, PathBuf},
5    sync::{
6        atomic::{AtomicU64, Ordering},
7        Mutex, MutexGuard,
8    },
9};
10
11use rhiza_core::{
12    ConfigurationState, EntryType, LogAnchor, LogEntry, LogHash, LogIndex, RecoveryAnchor,
13    SnapshotIdentity, StopBinding, SuccessorDescriptor, RECOVERY_ANCHOR_FORMAT_VERSION,
14};
15
16pub const QLOG_MAGIC: [u8; 4] = *b"QLOG";
17pub const QLOG_FORMAT_VERSION: u16 = 1;
18pub const QLOG_HEADER_LEN: usize = 76;
19pub const QLOG_FRAME_MAGIC: [u8; 4] = *b"QFRM";
20pub const QLOG_FOOTER_MAGIC: [u8; 4] = *b"QEND";
21pub const OPEN_SEGMENT_MAX_BYTES: usize = 8 * 1024 * 1024;
22pub const OPEN_SEGMENT_MAX_ENTRIES: usize = 4096;
23
24const HEADER_WITHOUT_CRC_LEN: usize = 72;
25const FRAME_PREFIX_LEN: usize = 108;
26const FRAME_MIN_LEN: usize = 144;
27const FOOTER_LEN: usize = 88;
28const TRUNCATE_INTENT_FILE_NAME: &str = ".truncate-intent";
29const TRUNCATE_INTENT_MAGIC: [u8; 4] = *b"QTRN";
30const TRUNCATE_INTENT_VERSION: u16 = 1;
31const TRUNCATE_INTENT_REPLACEMENT: u16 = 1;
32const ANCHOR_FILE_NAME: &str = "recovery.anchor";
33const ANCHOR_MAGIC: [u8; 4] = *b"QANC";
34const ANCHOR_VERSION: u16 = 4;
35const COMPACT_INTENT_FILE_NAME: &str = ".compact-intent";
36const COMPACT_INTENT_MAGIC: [u8; 4] = *b"QCMP";
37const COMPACT_INTENT_VERSION: u16 = 1;
38const COMPACT_INTENT_PREVIOUS_ANCHOR: u16 = 1;
39const COMPACT_INTENT_REPLACEMENT: u16 = 2;
40static NEXT_TEMP_FILE_ID: AtomicU64 = AtomicU64::new(0);
41
42pub type Result<T> = std::result::Result<T, Error>;
43
44#[derive(Clone, Debug, Eq, PartialEq)]
45pub enum Error {
46    InvalidIndexRange {
47        start: LogIndex,
48        end: LogIndex,
49    },
50    CompactionUnsupported,
51    CompactionAboveTip {
52        target: LogIndex,
53        tip: Option<LogIndex>,
54    },
55    CompactionHashMismatch {
56        index: LogIndex,
57    },
58    CompactionRegression {
59        target: LogIndex,
60        anchor: LogIndex,
61    },
62    CompactionConflict {
63        index: LogIndex,
64    },
65    TruncateCompactedPrefix {
66        from: LogIndex,
67        anchor: LogIndex,
68    },
69    Decode(String),
70    Io(String),
71}
72
73impl fmt::Display for Error {
74    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
75        match self {
76            Self::InvalidIndexRange { start, end } => {
77                write!(f, "invalid index range: start {start} is after end {end}")
78            }
79            Self::CompactionUnsupported => {
80                write!(f, "prefix compaction is unsupported by qlog v1")
81            }
82            Self::CompactionAboveTip { target, tip } => {
83                write!(f, "compaction target {target} is above log tip {tip:?}")
84            }
85            Self::CompactionHashMismatch { index } => {
86                write!(f, "compaction hash does not match log entry {index}")
87            }
88            Self::CompactionRegression { target, anchor } => write!(
89                f,
90                "compaction target {target} regresses persisted anchor {anchor}"
91            ),
92            Self::CompactionConflict { index } => {
93                write!(f, "compaction replay conflicts at index {index}")
94            }
95            Self::TruncateCompactedPrefix { from, anchor } => write!(
96                f,
97                "cannot truncate from {from} at or below compacted anchor {anchor}"
98            ),
99            Self::Decode(message) => write!(f, "qlog decode failed: {message}"),
100            Self::Io(message) => write!(f, "qlog io failed: {message}"),
101        }
102    }
103}
104
105impl std::error::Error for Error {}
106
107#[derive(Clone, Copy, Debug, Eq, PartialEq)]
108pub struct IndexRange {
109    start: LogIndex,
110    end: LogIndex,
111}
112
113impl IndexRange {
114    pub fn new(start: LogIndex, end: LogIndex) -> Result<Self> {
115        if start > end {
116            return Err(Error::InvalidIndexRange { start, end });
117        }
118
119        Ok(Self { start, end })
120    }
121
122    pub const fn start(&self) -> LogIndex {
123        self.start
124    }
125
126    pub const fn end(&self) -> LogIndex {
127        self.end
128    }
129}
130
131pub fn segment_file_name(range: IndexRange) -> String {
132    format!("{:020}-{:020}.qlog", range.start(), range.end())
133}
134
135pub fn encode_segment(entries: &[LogEntry]) -> Vec<u8> {
136    encode_segment_inner(entries, true)
137}
138
139pub fn encode_open_segment(entries: &[LogEntry]) -> Vec<u8> {
140    encode_segment_inner(entries, false)
141}
142
143fn encode_segment_inner(entries: &[LogEntry], closed: bool) -> Vec<u8> {
144    let Some(first) = entries.first() else {
145        return Vec::new();
146    };
147    let mut out = encode_header(first);
148    let mut entry_hashes = Vec::with_capacity(entries.len() * 32);
149    for entry in entries {
150        out.extend_from_slice(&encode_frame(entry));
151        entry_hashes.extend_from_slice(entry.hash.as_bytes());
152    }
153    if closed {
154        let last = entries.last().expect("non-empty entries");
155        out.extend_from_slice(&encode_footer(
156            &out[..QLOG_HEADER_LEN],
157            &entry_hashes,
158            last.index,
159            entries.len() as u64,
160            last.hash,
161        ));
162    }
163    out
164}
165
166pub fn decode_segment(bytes: &[u8]) -> Result<Vec<LogEntry>> {
167    decode_segment_for_cluster(bytes, "")
168}
169
170pub fn decode_segment_for_cluster(bytes: &[u8], cluster_id: &str) -> Result<Vec<LogEntry>> {
171    let (header, mut offset) = decode_header(bytes, cluster_id)?;
172    let cluster_id = if cluster_id.is_empty() {
173        header.cluster_id_hash.to_hex()
174    } else {
175        cluster_id.to_string()
176    };
177    let mut entries = Vec::new();
178    let mut entry_hashes = Vec::new();
179
180    loop {
181        if offset >= bytes.len() {
182            return Err(Error::Decode("missing qlog footer".into()));
183        }
184        if bytes.len() - offset >= 4 && bytes[offset..offset + 4] == QLOG_FOOTER_MAGIC {
185            decode_footer(
186                bytes,
187                offset,
188                &bytes[..QLOG_HEADER_LEN],
189                &entry_hashes,
190                &entries,
191            )?;
192            validate_entries(&header, &entries)?;
193            return Ok(entries);
194        }
195        let (entry, next_offset) = decode_frame(bytes, offset, &cluster_id)?;
196        entry_hashes.extend_from_slice(entry.hash.as_bytes());
197        entries.push(entry);
198        offset = next_offset;
199    }
200}
201
202pub fn write_segment_file(dir: impl Into<PathBuf>, entries: &[LogEntry]) -> Result<PathBuf> {
203    let dir = dir.into();
204    fs::create_dir_all(&dir).map_err(|err| Error::Io(err.to_string()))?;
205    publish_closed_segment(&dir, entries)
206}
207
208pub fn read_segment_file(path: impl Into<PathBuf>) -> Result<Vec<LogEntry>> {
209    let bytes = fs::read(path.into()).map_err(|err| Error::Io(err.to_string()))?;
210    decode_segment(&bytes)
211}
212
213pub fn recover_open_segment_file(
214    path: impl AsRef<Path>,
215    cluster_id: &str,
216) -> Result<Vec<LogEntry>> {
217    let path = path.as_ref();
218    let name = path
219        .file_name()
220        .and_then(|name| name.to_str())
221        .unwrap_or_default();
222    if !name.ends_with("-open.qlog") {
223        return Err(Error::Decode(
224            "refusing to recover non-open qlog segment".into(),
225        ));
226    }
227    let bytes = fs::read(path).map_err(|err| Error::Io(err.to_string()))?;
228    let (entries, valid_len) = recover_open_segment_prefix(&bytes, cluster_id)?;
229    if valid_len < bytes.len() {
230        let file = fs::OpenOptions::new()
231            .write(true)
232            .open(path)
233            .map_err(|err| Error::Io(err.to_string()))?;
234        file.set_len(valid_len as u64)
235            .map_err(|err| Error::Io(err.to_string()))?;
236        file.sync_all().map_err(|err| Error::Io(err.to_string()))?;
237    }
238    Ok(entries)
239}
240
241fn recover_open_segment_prefix(bytes: &[u8], cluster_id: &str) -> Result<(Vec<LogEntry>, usize)> {
242    let (header, mut offset) = decode_header(bytes, cluster_id)?;
243    let mut entries = Vec::new();
244    let mut valid_len = offset;
245    while offset < bytes.len() {
246        let (entry, next_offset) = match decode_frame(bytes, offset, cluster_id) {
247            Ok(frame) => frame,
248            Err(_) if final_frame_is_incomplete(bytes, offset)? => break,
249            Err(err) => return Err(err),
250        };
251        validate_decoded_entry(&header, entries.len(), entries.last(), &entry)?;
252        entries.push(entry);
253        offset = next_offset;
254        valid_len = offset;
255    }
256    Ok((entries, valid_len))
257}
258
259fn final_frame_is_incomplete(bytes: &[u8], offset: usize) -> Result<bool> {
260    let tail = &bytes[offset..];
261    let magic_prefix_len = tail.len().min(QLOG_FRAME_MAGIC.len());
262    if tail[..magic_prefix_len] != QLOG_FRAME_MAGIC[..magic_prefix_len] {
263        return Ok(false);
264    }
265    if tail.len() < 8 {
266        return Ok(true);
267    }
268
269    let frame_len = read_u32(tail, 4)? as usize;
270    if frame_len < FRAME_MIN_LEN {
271        return Ok(false);
272    }
273    if tail.len() >= FRAME_PREFIX_LEN {
274        let payload_len = read_u32(tail, 104)? as usize;
275        if FRAME_MIN_LEN.checked_add(payload_len) != Some(frame_len) {
276            return Ok(false);
277        }
278    }
279    Ok(frame_len > tail.len())
280}
281
282fn encode_header(first: &LogEntry) -> Vec<u8> {
283    let mut out = Vec::with_capacity(QLOG_HEADER_LEN);
284    out.extend_from_slice(&QLOG_MAGIC);
285    put_u16(&mut out, QLOG_FORMAT_VERSION);
286    put_u16(&mut out, QLOG_HEADER_LEN as u16);
287    put_u64(&mut out, first.index);
288    put_u64(&mut out, first.epoch);
289    put_u64(&mut out, first.config_id);
290    out.extend_from_slice(LogHash::digest(&[first.cluster_id.as_bytes()]).as_bytes());
291    put_u64(&mut out, 0);
292    let crc = crc32c(&out[..HEADER_WITHOUT_CRC_LEN]);
293    put_u32(&mut out, crc);
294    out
295}
296
297fn decode_header(bytes: &[u8], cluster_id: &str) -> Result<(SegmentHeader, usize)> {
298    if bytes.len() < QLOG_HEADER_LEN {
299        return Err(Error::Decode("short qlog header".into()));
300    }
301    if bytes[0..4] != QLOG_MAGIC {
302        return Err(Error::Decode("wrong qlog magic".into()));
303    }
304    let version = read_u16(bytes, 4)?;
305    if version != QLOG_FORMAT_VERSION {
306        return Err(Error::Decode("unsupported qlog version".into()));
307    }
308    let header_len = read_u16(bytes, 6)? as usize;
309    if header_len != QLOG_HEADER_LEN {
310        return Err(Error::Decode("invalid qlog header_len".into()));
311    }
312    let expected_crc = read_u32(bytes, HEADER_WITHOUT_CRC_LEN)?;
313    if crc32c(&bytes[..HEADER_WITHOUT_CRC_LEN]) != expected_crc {
314        return Err(Error::Decode("qlog header crc mismatch".into()));
315    }
316    let cluster_id_hash = read_hash(bytes, 32)?;
317    if !cluster_id.is_empty() && LogHash::digest(&[cluster_id.as_bytes()]) != cluster_id_hash {
318        return Err(Error::Decode("qlog cluster_id hash mismatch".into()));
319    }
320    Ok((
321        SegmentHeader::new_with_config(
322            cluster_id_hash,
323            read_u64(bytes, 16)?,
324            read_u64(bytes, 24)?,
325            read_u64(bytes, 8)?,
326            read_u64(bytes, 64)?,
327        ),
328        QLOG_HEADER_LEN,
329    ))
330}
331
332fn encode_frame(entry: &LogEntry) -> Vec<u8> {
333    let frame_len = FRAME_MIN_LEN + entry.payload.len();
334    let mut out = Vec::with_capacity(frame_len);
335    out.extend_from_slice(&QLOG_FRAME_MAGIC);
336    put_u32(&mut out, frame_len as u32);
337    put_u64(&mut out, entry.index);
338    put_u64(&mut out, entry.epoch);
339    put_u64(&mut out, entry.config_id);
340    out.push(entry.entry_type.as_u8());
341    out.extend_from_slice(&[0; 7]);
342    out.extend_from_slice(entry.prev_hash.as_bytes());
343    out.extend_from_slice(LogHash::digest(&[&entry.payload]).as_bytes());
344    put_u32(&mut out, entry.payload.len() as u32);
345    out.extend_from_slice(&entry.payload);
346    out.extend_from_slice(entry.hash.as_bytes());
347    let crc = crc32c(&out);
348    put_u32(&mut out, crc);
349    out
350}
351
352fn decode_frame(bytes: &[u8], offset: usize, cluster_id: &str) -> Result<(LogEntry, usize)> {
353    if bytes.len().saturating_sub(offset) < FRAME_MIN_LEN {
354        return Err(Error::Decode("short qlog frame".into()));
355    }
356    if bytes[offset..offset + 4] != QLOG_FRAME_MAGIC {
357        return Err(Error::Decode("wrong qlog frame magic".into()));
358    }
359    let frame_len = read_u32(bytes, offset + 4)? as usize;
360    let Some(frame_end) = offset.checked_add(frame_len) else {
361        return Err(Error::Decode("invalid qlog frame_len".into()));
362    };
363    if frame_len < FRAME_MIN_LEN || frame_end > bytes.len() {
364        return Err(Error::Decode("invalid qlog frame_len".into()));
365    }
366    let crc_offset = frame_end - 4;
367    let expected_crc = read_u32(bytes, crc_offset)?;
368    if crc32c(&bytes[offset..crc_offset]) != expected_crc {
369        return Err(Error::Decode("qlog frame crc mismatch".into()));
370    }
371    let index = read_u64(bytes, offset + 8)?;
372    let epoch = read_u64(bytes, offset + 16)?;
373    let config_id = read_u64(bytes, offset + 24)?;
374    let entry_type = EntryType::from_u8(bytes[offset + 32])
375        .ok_or_else(|| Error::Decode("invalid qlog entry_type".into()))?;
376    let prev_hash = read_hash(bytes, offset + 40)?;
377    let payload_hash = read_hash(bytes, offset + 72)?;
378    let payload_len = read_u32(bytes, offset + 104)? as usize;
379    if FRAME_MIN_LEN.checked_add(payload_len) != Some(frame_len) {
380        return Err(Error::Decode(
381            "qlog frame_len does not match payload_len".into(),
382        ));
383    }
384    let payload_start = offset + FRAME_PREFIX_LEN;
385    let payload_end = payload_start
386        .checked_add(payload_len)
387        .ok_or_else(|| Error::Decode("invalid qlog payload_len".into()))?;
388    let payload = bytes[payload_start..payload_end].to_vec();
389    if LogHash::digest(&[&payload]) != payload_hash {
390        return Err(Error::Decode("qlog payload_hash mismatch".into()));
391    }
392    let hash = read_hash(bytes, payload_end)?;
393    let entry = LogEntry {
394        cluster_id: cluster_id.to_string(),
395        epoch,
396        config_id,
397        index,
398        entry_type,
399        payload,
400        prev_hash,
401        hash,
402    };
403    if entry.recompute_hash() != hash {
404        return Err(Error::Decode("qlog entry_hash mismatch".into()));
405    }
406    Ok((entry, frame_end))
407}
408
409fn encode_footer(
410    header: &[u8],
411    entry_hashes: &[u8],
412    end_index: LogIndex,
413    entry_count: u64,
414    last_entry_hash: LogHash,
415) -> Vec<u8> {
416    let mut prefix = Vec::with_capacity(52);
417    prefix.extend_from_slice(&QLOG_FOOTER_MAGIC);
418    put_u64(&mut prefix, end_index);
419    put_u64(&mut prefix, entry_count);
420    prefix.extend_from_slice(last_entry_hash.as_bytes());
421    let segment_hash = LogHash::digest(&[header, entry_hashes, &prefix]);
422    let mut out = prefix;
423    out.extend_from_slice(segment_hash.as_bytes());
424    let crc = crc32c(&out);
425    put_u32(&mut out, crc);
426    out
427}
428
429fn decode_footer(
430    bytes: &[u8],
431    offset: usize,
432    header: &[u8],
433    entry_hashes: &[u8],
434    entries: &[LogEntry],
435) -> Result<()> {
436    if bytes.len().saturating_sub(offset) != FOOTER_LEN {
437        return Err(Error::Decode("invalid qlog footer length".into()));
438    }
439    let footer = &bytes[offset..offset + FOOTER_LEN];
440    let expected_crc = read_u32(footer, FOOTER_LEN - 4)?;
441    if crc32c(&footer[..FOOTER_LEN - 4]) != expected_crc {
442        return Err(Error::Decode("qlog footer crc mismatch".into()));
443    }
444    let end_index = read_u64(footer, 4)?;
445    let entry_count = read_u64(footer, 12)?;
446    let last_hash = read_hash(footer, 20)?;
447    let segment_hash = read_hash(footer, 52)?;
448    let expected_segment_hash = LogHash::digest(&[header, entry_hashes, &footer[..52]]);
449    if segment_hash != expected_segment_hash {
450        return Err(Error::Decode("qlog segment_hash mismatch".into()));
451    }
452    if entries.len() as u64 != entry_count {
453        return Err(Error::Decode("qlog footer entry_count mismatch".into()));
454    }
455    if let Some(last) = entries.last() {
456        if last.index != end_index || last.hash != last_hash {
457            return Err(Error::Decode("qlog footer last entry mismatch".into()));
458        }
459    } else if entry_count != 0 {
460        return Err(Error::Decode("qlog empty footer mismatch".into()));
461    }
462    Ok(())
463}
464
465fn validate_entries(header: &SegmentHeader, entries: &[LogEntry]) -> Result<()> {
466    for (position, entry) in entries.iter().enumerate() {
467        validate_decoded_entry(
468            header,
469            position,
470            position.checked_sub(1).map(|i| &entries[i]),
471            entry,
472        )?;
473    }
474    Ok(())
475}
476
477fn validate_decoded_entry(
478    header: &SegmentHeader,
479    position: usize,
480    previous: Option<&LogEntry>,
481    entry: &LogEntry,
482) -> Result<()> {
483    let expected_index = u64::try_from(position)
484        .ok()
485        .and_then(|position| header.start_index.checked_add(position));
486    if expected_index != Some(entry.index) {
487        return Err(Error::Decode("qlog index gap".into()));
488    }
489    if entry.epoch != header.epoch || entry.config_id != header.config_id {
490        return Err(Error::Decode("qlog epoch/config mismatch".into()));
491    }
492    if previous.is_some_and(|previous| entry.prev_hash != previous.hash) {
493        return Err(Error::Decode("qlog hash chain mismatch".into()));
494    }
495    Ok(())
496}
497
498fn read_hash(bytes: &[u8], offset: usize) -> Result<LogHash> {
499    let slice = bytes
500        .get(offset..offset + 32)
501        .ok_or_else(|| Error::Decode("short qlog hash".into()))?;
502    let mut out = [0; 32];
503    out.copy_from_slice(slice);
504    Ok(LogHash::from_bytes(out))
505}
506
507fn read_u16(bytes: &[u8], offset: usize) -> Result<u16> {
508    let slice = bytes
509        .get(offset..offset + 2)
510        .ok_or_else(|| Error::Decode("short qlog u16".into()))?;
511    Ok(u16::from_be_bytes(
512        slice.try_into().expect("u16 slice length"),
513    ))
514}
515
516fn read_u32(bytes: &[u8], offset: usize) -> Result<u32> {
517    let slice = bytes
518        .get(offset..offset + 4)
519        .ok_or_else(|| Error::Decode("short qlog u32".into()))?;
520    Ok(u32::from_be_bytes(
521        slice.try_into().expect("u32 slice length"),
522    ))
523}
524
525fn read_u64(bytes: &[u8], offset: usize) -> Result<u64> {
526    let slice = bytes
527        .get(offset..offset + 8)
528        .ok_or_else(|| Error::Decode("short qlog u64".into()))?;
529    Ok(u64::from_be_bytes(
530        slice.try_into().expect("u64 slice length"),
531    ))
532}
533
534fn put_u16(out: &mut Vec<u8>, value: u16) {
535    out.extend_from_slice(&value.to_be_bytes());
536}
537
538fn put_u32(out: &mut Vec<u8>, value: u32) {
539    out.extend_from_slice(&value.to_be_bytes());
540}
541
542fn put_u64(out: &mut Vec<u8>, value: u64) {
543    out.extend_from_slice(&value.to_be_bytes());
544}
545
546fn crc32c(bytes: &[u8]) -> u32 {
547    let mut crc = !0u32;
548    for byte in bytes {
549        crc ^= u32::from(*byte);
550        for _ in 0..8 {
551            let mask = (crc & 1).wrapping_neg();
552            crc = (crc >> 1) ^ (0x82f6_3b78 & mask);
553        }
554    }
555    !crc
556}
557
558#[derive(Clone, Copy, Debug, Eq, PartialEq)]
559pub struct SegmentHeader {
560    magic: [u8; 4],
561    cluster_id_hash: LogHash,
562    epoch: u64,
563    config_id: u64,
564    start_index: LogIndex,
565    created_at_unix_ms: u64,
566}
567
568impl SegmentHeader {
569    pub const fn new(
570        cluster_id_hash: LogHash,
571        epoch: u64,
572        start_index: LogIndex,
573        created_at_unix_ms: u64,
574    ) -> Self {
575        Self {
576            magic: QLOG_MAGIC,
577            cluster_id_hash,
578            epoch,
579            config_id: 0,
580            start_index,
581            created_at_unix_ms,
582        }
583    }
584
585    pub const fn new_with_config(
586        cluster_id_hash: LogHash,
587        epoch: u64,
588        config_id: u64,
589        start_index: LogIndex,
590        created_at_unix_ms: u64,
591    ) -> Self {
592        Self {
593            magic: QLOG_MAGIC,
594            cluster_id_hash,
595            epoch,
596            config_id,
597            start_index,
598            created_at_unix_ms,
599        }
600    }
601
602    pub const fn magic(&self) -> [u8; 4] {
603        self.magic
604    }
605
606    pub const fn epoch(&self) -> u64 {
607        self.epoch
608    }
609
610    pub const fn config_id(&self) -> u64 {
611        self.config_id
612    }
613
614    pub const fn start_index(&self) -> LogIndex {
615        self.start_index
616    }
617}
618
619#[derive(Clone, Debug, Eq, PartialEq)]
620pub struct SegmentFile {
621    range: IndexRange,
622    bytes: Vec<u8>,
623}
624
625impl SegmentFile {
626    pub fn new(range: IndexRange, bytes: Vec<u8>) -> Self {
627        Self { range, bytes }
628    }
629
630    pub const fn range(&self) -> IndexRange {
631        self.range
632    }
633
634    pub fn bytes(&self) -> &[u8] {
635        &self.bytes
636    }
637}
638
639pub trait LogStore {
640    fn append(&self, entry: &LogEntry) -> Result<()>;
641    fn append_batch(&self, entries: &[LogEntry]) -> Result<()>;
642    fn read(&self, index: LogIndex) -> Result<Option<LogEntry>>;
643    fn read_range(&self, range: IndexRange) -> Result<Vec<LogEntry>>;
644    fn last_index(&self) -> Result<Option<LogIndex>>;
645    fn truncate_suffix(&self, from: LogIndex) -> Result<()>;
646    fn compact_prefix(&self, verified_snapshot_anchor: &RecoveryAnchor) -> Result<()>;
647}
648
649#[derive(Clone, Debug, Eq, PartialEq)]
650pub struct LogState {
651    pub anchor: Option<RecoveryAnchor>,
652    pub first_retained_index: LogIndex,
653    pub tip: Option<LogAnchor>,
654}
655
656#[derive(Debug)]
657pub struct FileLogStore {
658    inner: Mutex<FileLogStoreInner>,
659}
660
661impl FileLogStore {
662    pub fn open(
663        dir: impl Into<PathBuf>,
664        cluster_id: impl Into<String>,
665        epoch: u64,
666        config_id: u64,
667    ) -> Result<Self> {
668        Self::open_with_configuration(
669            dir,
670            cluster_id,
671            epoch,
672            ConfigurationState::active(config_id, LogHash::ZERO),
673        )
674    }
675
676    pub fn open_with_configuration(
677        dir: impl Into<PathBuf>,
678        cluster_id: impl Into<String>,
679        epoch: u64,
680        initial_configuration: ConfigurationState,
681    ) -> Result<Self> {
682        let dir = dir.into();
683        let cluster_id = cluster_id.into();
684        if cluster_id.is_empty() {
685            return Err(Error::Decode("cluster_id must not be empty".into()));
686        }
687
688        let existed = dir.exists();
689        fs::create_dir_all(&dir).map_err(|err| Error::Io(err.to_string()))?;
690        if !existed {
691            if let Some(parent) = dir.parent() {
692                sync_directory(parent)?;
693            }
694        }
695        recover_truncate_intent(&dir)?;
696        recover_compact_intent(&dir)?;
697        let anchor = read_anchor(&dir)?;
698        validate_anchor_identity(anchor.as_ref(), &cluster_id, epoch)?;
699        let (segments, configuration_state) = scan_closed_segments(
700            &dir,
701            &cluster_id,
702            epoch,
703            &initial_configuration,
704            anchor.as_ref(),
705        )?;
706        let (open_segment, configuration_state) = scan_open_segment(
707            &dir,
708            &cluster_id,
709            epoch,
710            anchor.as_ref(),
711            &segments,
712            configuration_state,
713        )?;
714
715        Ok(Self {
716            inner: Mutex::new(FileLogStoreInner {
717                dir,
718                cluster_id,
719                epoch,
720                initial_configuration,
721                configuration_state,
722                anchor,
723                segments,
724                open_segment,
725            }),
726        })
727    }
728
729    fn lock(&self) -> Result<MutexGuard<'_, FileLogStoreInner>> {
730        self.inner
731            .lock()
732            .map_err(|_| Error::Io("file log store lock poisoned".into()))
733    }
734
735    pub fn logical_state(&self) -> Result<LogState> {
736        let inner = self.lock()?;
737        let first_retained_index = match &inner.anchor {
738            Some(anchor) => anchor
739                .compacted()
740                .index()
741                .checked_add(1)
742                .ok_or_else(|| Error::Decode("qlog anchor index overflow".into()))?,
743            None => 1,
744        };
745        let tip = inner
746            .open_segment
747            .as_ref()
748            .and_then(|segment| segment.entries.last())
749            .map(|entry| LogAnchor::new(entry.index, entry.hash))
750            .or_else(|| {
751                inner.segments.last().map(|segment| {
752                    let entry = segment.entries.last().expect("non-empty segment");
753                    LogAnchor::new(entry.index, entry.hash)
754                })
755            })
756            .or_else(|| inner.anchor.as_ref().map(|anchor| *anchor.compacted()));
757        Ok(LogState {
758            anchor: inner.anchor.clone(),
759            first_retained_index,
760            tip,
761        })
762    }
763
764    pub fn configuration_state(&self) -> Result<ConfigurationState> {
765        Ok(self.lock()?.configuration_state.clone())
766    }
767
768    pub fn install_recovery_anchor(
769        &self,
770        verified_anchor: &RecoveryAnchor,
771        expected_recovery_generation: u64,
772        expected_configuration: &ConfigurationState,
773    ) -> Result<()> {
774        let mut inner = self.lock()?;
775        validate_anchor(verified_anchor)?;
776        validate_anchor_identity(Some(verified_anchor), &inner.cluster_id, inner.epoch)?;
777        if verified_anchor.recovery_generation() != expected_recovery_generation {
778            return Err(Error::Decode(
779                "recovery anchor generation does not match expected generation".into(),
780            ));
781        }
782        if verified_anchor.configuration_state() != expected_configuration {
783            return Err(Error::Decode(
784                "recovery anchor configuration state does not match expected state".into(),
785            ));
786        }
787        if inner.anchor.is_some() || !inner.segments.is_empty() || inner.open_segment.is_some() {
788            return Err(Error::Decode(
789                "recovery anchor installation requires an empty qlog store".into(),
790            ));
791        }
792
793        install_anchor(
794            &inner.dir,
795            &CompactIntent {
796                previous_anchor: None,
797                anchor: verified_anchor.clone(),
798                old_segment_names: Vec::new(),
799                replacement: None,
800            },
801        )?;
802        sync_directory(&inner.dir)?;
803        inner.anchor = Some(verified_anchor.clone());
804        inner.configuration_state = verified_anchor.configuration_state().clone();
805        Ok(())
806    }
807
808    /// Appends validated entries without issuing a data sync.
809    ///
810    /// Call [`Self::sync`] before making the appended entries externally durable.
811    pub fn append_batch_buffered(&self, entries: &[LogEntry]) -> Result<()> {
812        let mut inner = self.lock()?;
813        append_batch_to_open_segment(&mut inner, entries, false)
814    }
815
816    /// Syncs all buffered appends and returns the durable logical tip.
817    pub fn sync(&self) -> Result<Option<LogIndex>> {
818        let inner = self.lock()?;
819        if let Some(open) = &inner.open_segment {
820            open.file
821                .sync_data()
822                .map_err(|err| Error::Io(err.to_string()))?;
823        }
824        Ok(inner.last_index())
825    }
826}
827
828impl LogStore for FileLogStore {
829    fn append(&self, entry: &LogEntry) -> Result<()> {
830        self.append_batch(std::slice::from_ref(entry))
831    }
832
833    fn append_batch(&self, entries: &[LogEntry]) -> Result<()> {
834        let mut inner = self.lock()?;
835        append_batch_to_open_segment(&mut inner, entries, true)
836    }
837
838    fn read(&self, index: LogIndex) -> Result<Option<LogEntry>> {
839        let inner = self.lock()?;
840        Ok(inner
841            .segments
842            .iter()
843            .find(|segment| segment.start() <= index && index <= segment.end())
844            .map(|segment| segment.entries[(index - segment.start()) as usize].clone())
845            .or_else(|| {
846                inner.open_segment.as_ref().and_then(|segment| {
847                    segment
848                        .entries
849                        .first()
850                        .filter(|first| first.index <= index)
851                        .and_then(|first| segment.entries.get((index - first.index) as usize))
852                        .cloned()
853                })
854            }))
855    }
856
857    fn read_range(&self, range: IndexRange) -> Result<Vec<LogEntry>> {
858        let inner = self.lock()?;
859        Ok(inner
860            .segments
861            .iter()
862            .flat_map(|segment| segment.entries.iter())
863            .chain(
864                inner
865                    .open_segment
866                    .iter()
867                    .flat_map(|segment| segment.entries.iter()),
868            )
869            .filter(|entry| range.start() <= entry.index && entry.index <= range.end())
870            .cloned()
871            .collect())
872    }
873
874    fn last_index(&self) -> Result<Option<LogIndex>> {
875        let inner = self.lock()?;
876        Ok(inner.last_index())
877    }
878
879    fn truncate_suffix(&self, from: LogIndex) -> Result<()> {
880        let mut inner = self.lock()?;
881        if let Some(anchor) = &inner.anchor {
882            if from <= anchor.compacted().index() {
883                return Err(Error::TruncateCompactedPrefix {
884                    from,
885                    anchor: anchor.compacted().index(),
886                });
887            }
888        }
889        seal_open_segment(&mut inner)?;
890        truncate_suffix_with_hook(&mut inner, from, &mut |_| Ok(()))?;
891        let (segments, configuration_state) = scan_closed_segments(
892            &inner.dir,
893            &inner.cluster_id,
894            inner.epoch,
895            &inner.initial_configuration,
896            inner.anchor.as_ref(),
897        )?;
898        inner.segments = segments;
899        inner.configuration_state = configuration_state;
900        Ok(())
901    }
902
903    fn compact_prefix(&self, verified_snapshot_anchor: &RecoveryAnchor) -> Result<()> {
904        let mut inner = self.lock()?;
905        if !validate_compaction(&inner, verified_snapshot_anchor)? {
906            return Ok(());
907        }
908        seal_open_segment(&mut inner)?;
909        compact_prefix_with_hook(&mut inner, verified_snapshot_anchor, &mut |_| Ok(()))?;
910        inner.anchor = read_anchor(&inner.dir)?;
911        let (segments, configuration_state) = scan_closed_segments(
912            &inner.dir,
913            &inner.cluster_id,
914            inner.epoch,
915            &inner.initial_configuration,
916            inner.anchor.as_ref(),
917        )?;
918        inner.segments = segments;
919        inner.configuration_state = configuration_state;
920        Ok(())
921    }
922}
923
924#[derive(Debug)]
925struct FileLogStoreInner {
926    dir: PathBuf,
927    cluster_id: String,
928    epoch: u64,
929    initial_configuration: ConfigurationState,
930    configuration_state: ConfigurationState,
931    anchor: Option<RecoveryAnchor>,
932    segments: Vec<ClosedSegment>,
933    open_segment: Option<OpenSegment>,
934}
935
936#[derive(Clone, Debug)]
937struct ClosedSegment {
938    entries: Vec<LogEntry>,
939}
940
941#[derive(Debug)]
942struct OpenSegment {
943    config_id: u64,
944    path: PathBuf,
945    file: fs::File,
946    bytes_len: usize,
947    entries: Vec<LogEntry>,
948}
949
950#[derive(Clone, Debug, Eq, PartialEq)]
951struct TruncateIntent {
952    old_segment_names: Vec<String>,
953    replacement: Option<TruncateReplacement>,
954}
955
956#[derive(Clone, Debug, Eq, PartialEq)]
957struct TruncateReplacement {
958    temp_name: String,
959    final_name: String,
960}
961
962#[derive(Clone, Debug, Eq, PartialEq)]
963struct CompactIntent {
964    previous_anchor: Option<RecoveryAnchor>,
965    anchor: RecoveryAnchor,
966    old_segment_names: Vec<String>,
967    replacement: Option<TruncateReplacement>,
968}
969
970#[derive(Clone, Copy, Debug, Eq, PartialEq)]
971enum CompactPhase {
972    ReplacementPrepared,
973    IntentRenamed,
974    IntentDurable,
975    AnchorInstalled,
976    AnchorDurable,
977    OldSegmentRemoved(usize),
978    ReplacementInstalled,
979    AppliedDirectorySynced,
980    IntentRemoved,
981    CompleteDirectorySynced,
982}
983
984#[derive(Clone, Copy, Debug, Eq, PartialEq)]
985enum TruncatePhase {
986    ReplacementPrepared,
987    IntentRenamed,
988    IntentDurable,
989    OldSegmentRemoved(usize),
990    ReplacementInstalled,
991    AppliedDirectorySynced,
992    IntentRemoved,
993    CompleteDirectorySynced,
994}
995
996impl ClosedSegment {
997    fn start(&self) -> LogIndex {
998        self.entries.first().expect("non-empty segment").index
999    }
1000
1001    fn end(&self) -> LogIndex {
1002        self.entries.last().expect("non-empty segment").index
1003    }
1004}
1005
1006impl FileLogStoreInner {
1007    fn last_index(&self) -> Option<LogIndex> {
1008        self.open_segment
1009            .as_ref()
1010            .and_then(|segment| segment.entries.last().map(|entry| entry.index))
1011            .or_else(|| self.segments.last().map(ClosedSegment::end))
1012            .or_else(|| {
1013                self.anchor
1014                    .as_ref()
1015                    .map(|anchor| anchor.compacted().index())
1016            })
1017    }
1018}
1019
1020fn scan_closed_segments(
1021    dir: &Path,
1022    cluster_id: &str,
1023    epoch: u64,
1024    initial_configuration: &ConfigurationState,
1025    anchor: Option<&RecoveryAnchor>,
1026) -> Result<(Vec<ClosedSegment>, ConfigurationState)> {
1027    let mut paths = Vec::new();
1028    for entry in fs::read_dir(dir).map_err(|err| Error::Io(err.to_string()))? {
1029        let entry = entry.map_err(|err| Error::Io(err.to_string()))?;
1030        let name = entry.file_name();
1031        let name = name.to_string_lossy();
1032        if name.ends_with("-open.qlog") || !name.ends_with(".qlog") {
1033            continue;
1034        }
1035        let range = parse_closed_segment_name(&name)?;
1036        paths.push((range, entry.path()));
1037    }
1038    paths.sort_by_key(|(range, _)| (range.start(), range.end()));
1039
1040    let mut segments = Vec::with_capacity(paths.len());
1041    let mut configuration_state = anchor
1042        .map(|anchor| anchor.configuration_state().clone())
1043        .unwrap_or_else(|| initial_configuration.clone());
1044    for (range, path) in paths {
1045        let bytes = fs::read(&path).map_err(|err| Error::Io(err.to_string()))?;
1046        let entries = decode_segment_for_cluster(&bytes, cluster_id)?;
1047        let first = entries
1048            .first()
1049            .ok_or_else(|| Error::Decode("closed qlog segment is empty".into()))?;
1050        let last = entries.last().expect("non-empty segment");
1051        if first.index != range.start() || last.index != range.end() {
1052            return Err(Error::Decode(
1053                "qlog segment filename range does not match entries".into(),
1054            ));
1055        }
1056        if first.epoch != epoch {
1057            return Err(Error::Decode(
1058                "qlog segment does not match configured epoch".into(),
1059            ));
1060        }
1061        if segments.is_empty() {
1062            let (expected_index, expected_prev_hash) = match anchor {
1063                Some(anchor) => (
1064                    anchor
1065                        .compacted()
1066                        .index()
1067                        .checked_add(1)
1068                        .ok_or_else(|| Error::Decode("qlog anchor index overflow".into()))?,
1069                    anchor.compacted().hash(),
1070                ),
1071                None => (1, LogHash::ZERO),
1072            };
1073            if first.index != expected_index || first.prev_hash != expected_prev_hash {
1074                return Err(Error::Decode(
1075                    "qlog first retained entry does not match recovery anchor".into(),
1076                ));
1077            }
1078        }
1079        if let Some(previous) = segments.last() {
1080            let previous: &ClosedSegment = previous;
1081            if previous.end().checked_add(1) != Some(first.index) {
1082                return Err(Error::Decode("qlog index gap across segments".into()));
1083            }
1084            if first.prev_hash != previous.entries.last().expect("non-empty segment").hash {
1085                return Err(Error::Decode(
1086                    "qlog hash chain mismatch across segments".into(),
1087                ));
1088            }
1089        }
1090        for entry in &entries {
1091            configuration_state = configuration_state
1092                .validate_entry(entry)
1093                .map_err(|err| Error::Decode(err.to_string()))?;
1094        }
1095        segments.push(ClosedSegment { entries });
1096    }
1097    Ok((segments, configuration_state))
1098}
1099
1100fn scan_open_segment(
1101    dir: &Path,
1102    cluster_id: &str,
1103    epoch: u64,
1104    anchor: Option<&RecoveryAnchor>,
1105    segments: &[ClosedSegment],
1106    mut configuration_state: ConfigurationState,
1107) -> Result<(Option<OpenSegment>, ConfigurationState)> {
1108    let mut paths = fs::read_dir(dir)
1109        .map_err(|err| Error::Io(err.to_string()))?
1110        .filter_map(|entry| entry.ok())
1111        .filter_map(|entry| {
1112            let name = entry.file_name().into_string().ok()?;
1113            name.ends_with("-open.qlog").then_some((name, entry.path()))
1114        })
1115        .collect::<Vec<_>>();
1116    paths.sort_by(|left, right| left.0.cmp(&right.0));
1117    if paths.len() > 1 {
1118        return Err(Error::Decode("multiple open qlog segments".into()));
1119    }
1120    let Some((name, path)) = paths.pop() else {
1121        return Ok((None, configuration_state));
1122    };
1123    let start = parse_open_segment_name(&name)?;
1124    let bytes = fs::read(&path).map_err(|err| Error::Io(err.to_string()))?;
1125    let (header, _) = decode_header(&bytes, cluster_id)?;
1126    if header.start_index != start || header.epoch != epoch {
1127        return Err(Error::Decode(
1128            "open qlog filename or epoch does not match header".into(),
1129        ));
1130    }
1131    let entries = recover_open_segment_file(&path, cluster_id)?;
1132
1133    if let Some(segment) = segments
1134        .iter()
1135        .find(|segment| segment.start() == start && segment.entries == entries)
1136    {
1137        if segment.end() != entries.last().map_or(start, |entry| entry.index) {
1138            return Err(Error::Decode("open qlog duplicate range mismatch".into()));
1139        }
1140        fs::remove_file(&path).map_err(|err| Error::Io(err.to_string()))?;
1141        sync_directory(dir)?;
1142        return Ok((None, configuration_state));
1143    }
1144
1145    let (expected_index, expected_prev_hash) = match segments.last() {
1146        Some(segment) => (
1147            segment
1148                .end()
1149                .checked_add(1)
1150                .ok_or_else(|| Error::Decode("qlog index overflow".into()))?,
1151            segment.entries.last().expect("non-empty segment").hash,
1152        ),
1153        None => match anchor {
1154            Some(anchor) => (
1155                anchor
1156                    .compacted()
1157                    .index()
1158                    .checked_add(1)
1159                    .ok_or_else(|| Error::Decode("qlog anchor index overflow".into()))?,
1160                anchor.compacted().hash(),
1161            ),
1162            None => (1, LogHash::ZERO),
1163        },
1164    };
1165    if start != expected_index {
1166        return Err(Error::Decode(
1167            "open qlog does not start at the retained tip".into(),
1168        ));
1169    }
1170    if let Some(first) = entries.first() {
1171        if first.prev_hash != expected_prev_hash {
1172            return Err(Error::Decode(
1173                "qlog hash chain mismatch before open segment".into(),
1174            ));
1175        }
1176    }
1177    for entry in &entries {
1178        configuration_state = configuration_state
1179            .validate_entry(entry)
1180            .map_err(|err| Error::Decode(err.to_string()))?;
1181    }
1182    let file = fs::OpenOptions::new()
1183        .append(true)
1184        .open(&path)
1185        .map_err(|err| Error::Io(err.to_string()))?;
1186    let bytes_len = usize::try_from(
1187        file.metadata()
1188            .map_err(|err| Error::Io(err.to_string()))?
1189            .len(),
1190    )
1191    .map_err(|_| Error::Io("open qlog segment is too large".into()))?;
1192    Ok((
1193        Some(OpenSegment {
1194            config_id: header.config_id,
1195            path,
1196            file,
1197            bytes_len,
1198            entries,
1199        }),
1200        configuration_state,
1201    ))
1202}
1203
1204fn parse_closed_segment_name(name: &str) -> Result<IndexRange> {
1205    let stem = name
1206        .strip_suffix(".qlog")
1207        .ok_or_else(|| Error::Decode("invalid closed qlog segment filename".into()))?;
1208    let (start, end) = stem
1209        .split_once('-')
1210        .ok_or_else(|| Error::Decode("invalid closed qlog segment filename".into()))?;
1211    if start.len() != 20
1212        || end.len() != 20
1213        || !start.bytes().all(|byte| byte.is_ascii_digit())
1214        || !end.bytes().all(|byte| byte.is_ascii_digit())
1215    {
1216        return Err(Error::Decode("invalid closed qlog segment filename".into()));
1217    }
1218    let start = start
1219        .parse()
1220        .map_err(|_| Error::Decode("invalid closed qlog segment filename".into()))?;
1221    let end = end
1222        .parse()
1223        .map_err(|_| Error::Decode("invalid closed qlog segment filename".into()))?;
1224    IndexRange::new(start, end)
1225}
1226
1227fn parse_open_segment_name(name: &str) -> Result<LogIndex> {
1228    let start = name
1229        .strip_suffix("-open.qlog")
1230        .ok_or_else(|| Error::Decode("invalid open qlog segment filename".into()))?;
1231    if start.len() != 20 || !start.bytes().all(|byte| byte.is_ascii_digit()) {
1232        return Err(Error::Decode("invalid open qlog segment filename".into()));
1233    }
1234    start
1235        .parse()
1236        .map_err(|_| Error::Decode("invalid open qlog segment filename".into()))
1237}
1238
1239fn open_segment_file_name(start: LogIndex) -> String {
1240    format!("{start:020}-open.qlog")
1241}
1242
1243fn validate_append(inner: &FileLogStoreInner, entries: &[LogEntry]) -> Result<ConfigurationState> {
1244    let open_tip = inner
1245        .open_segment
1246        .as_ref()
1247        .and_then(|segment| segment.entries.last());
1248    let (mut expected_index, mut expected_prev_hash) = match open_tip {
1249        Some(entry) => (
1250            entry
1251                .index
1252                .checked_add(1)
1253                .ok_or_else(|| Error::Decode("qlog index overflow".into()))?,
1254            entry.hash,
1255        ),
1256        None => match inner.segments.last() {
1257            Some(segment) => (
1258                segment
1259                    .end()
1260                    .checked_add(1)
1261                    .ok_or_else(|| Error::Decode("qlog index overflow".into()))?,
1262                segment.entries.last().expect("non-empty segment").hash,
1263            ),
1264            None => match &inner.anchor {
1265                Some(anchor) => (
1266                    anchor
1267                        .compacted()
1268                        .index()
1269                        .checked_add(1)
1270                        .ok_or_else(|| Error::Decode("qlog anchor index overflow".into()))?,
1271                    anchor.compacted().hash(),
1272                ),
1273                None => (1, LogHash::ZERO),
1274            },
1275        },
1276    };
1277
1278    let mut configuration_state = inner.configuration_state.clone();
1279    for (position, entry) in entries.iter().enumerate() {
1280        if entry.cluster_id != inner.cluster_id || entry.epoch != inner.epoch {
1281            return Err(Error::Decode(
1282                "qlog entry does not match configured cluster/epoch".into(),
1283            ));
1284        }
1285        if entry.index != expected_index {
1286            return Err(Error::Decode("qlog append index is not contiguous".into()));
1287        }
1288        if entry.prev_hash != expected_prev_hash {
1289            return Err(Error::Decode("qlog append prev_hash mismatch".into()));
1290        }
1291        if entry.recompute_hash() != entry.hash {
1292            return Err(Error::Decode("qlog append entry_hash mismatch".into()));
1293        }
1294        configuration_state = configuration_state
1295            .validate_entry(entry)
1296            .map_err(|err| Error::Decode(err.to_string()))?;
1297        expected_prev_hash = entry.hash;
1298        if position + 1 < entries.len() {
1299            expected_index = expected_index
1300                .checked_add(1)
1301                .ok_or_else(|| Error::Decode("qlog index overflow".into()))?;
1302        }
1303    }
1304    Ok(configuration_state)
1305}
1306
1307fn append_batch_to_open_segment(
1308    inner: &mut FileLogStoreInner,
1309    entries: &[LogEntry],
1310    sync: bool,
1311) -> Result<()> {
1312    if entries.is_empty() {
1313        return Ok(());
1314    }
1315    validate_append(inner, entries)?;
1316    for homogeneous in entries.chunk_by(|left, right| left.config_id == right.config_id) {
1317        let mut offset = 0;
1318        while offset < homogeneous.len() {
1319            ensure_open_segment(inner, &homogeneous[offset])?;
1320            let chunk_len = open_segment_chunk_len(
1321                inner.open_segment.as_ref().expect("open segment exists"),
1322                &homogeneous[offset..],
1323            )?;
1324            if chunk_len == 0 {
1325                seal_open_segment(inner)?;
1326                continue;
1327            }
1328
1329            let chunk = &homogeneous[offset..offset + chunk_len];
1330            let bytes = chunk.iter().flat_map(encode_frame).collect::<Vec<_>>();
1331            let open = inner.open_segment.as_mut().expect("open segment exists");
1332            let old_len = open.bytes_len as u64;
1333            if let Err(write_err) = open.file.write_all(&bytes) {
1334                if let Err(rollback_err) = open.file.set_len(old_len) {
1335                    return Err(Error::Io(format!(
1336                        "{write_err}; failed to roll back partial qlog append: {rollback_err}"
1337                    )));
1338                }
1339                return Err(Error::Io(write_err.to_string()));
1340            }
1341            open.bytes_len += bytes.len();
1342            open.entries.extend_from_slice(chunk);
1343            for entry in chunk {
1344                inner.configuration_state = inner
1345                    .configuration_state
1346                    .validate_entry(entry)
1347                    .map_err(|err| Error::Decode(err.to_string()))?;
1348            }
1349            offset += chunk_len;
1350        }
1351    }
1352    if sync {
1353        inner
1354            .open_segment
1355            .as_ref()
1356            .expect("non-empty append has open segment")
1357            .file
1358            .sync_data()
1359            .map_err(|err| Error::Io(err.to_string()))?;
1360    }
1361    Ok(())
1362}
1363
1364fn open_segment_chunk_len(open: &OpenSegment, entries: &[LogEntry]) -> Result<usize> {
1365    let mut bytes_len = open.bytes_len;
1366    let mut entry_count = open.entries.len();
1367    let mut chunk_len = 0;
1368    for entry in entries {
1369        let frame_len = FRAME_MIN_LEN
1370            .checked_add(entry.payload.len())
1371            .ok_or_else(|| Error::Decode("qlog frame length overflow".into()))?;
1372        let next_bytes_len = bytes_len
1373            .checked_add(frame_len)
1374            .ok_or_else(|| Error::Decode("qlog segment length overflow".into()))?;
1375        let next_entry_count = entry_count
1376            .checked_add(1)
1377            .ok_or_else(|| Error::Decode("qlog segment entry count overflow".into()))?;
1378        let oversized_first_entry = open.entries.is_empty() && chunk_len == 0;
1379        if !oversized_first_entry
1380            && (next_bytes_len > OPEN_SEGMENT_MAX_BYTES
1381                || next_entry_count > OPEN_SEGMENT_MAX_ENTRIES)
1382        {
1383            break;
1384        }
1385        bytes_len = next_bytes_len;
1386        entry_count = next_entry_count;
1387        chunk_len += 1;
1388    }
1389    Ok(chunk_len)
1390}
1391
1392fn ensure_open_segment(inner: &mut FileLogStoreInner, first: &LogEntry) -> Result<()> {
1393    if inner
1394        .open_segment
1395        .as_ref()
1396        .is_some_and(|segment| segment.config_id == first.config_id)
1397    {
1398        return Ok(());
1399    }
1400    seal_open_segment(inner)?;
1401
1402    let final_path = inner.dir.join(open_segment_file_name(first.index));
1403    if final_path.exists() {
1404        return Err(Error::Io(format!(
1405            "open qlog segment already exists: {}",
1406            final_path.display()
1407        )));
1408    }
1409    let (temp_path, mut file) = create_unique_temp_file(&inner.dir, &final_path)?;
1410    file.write_all(&encode_header(first))
1411        .and_then(|_| file.sync_all())
1412        .map_err(|err| Error::Io(err.to_string()))?;
1413    fs::rename(&temp_path, &final_path).map_err(|err| Error::Io(err.to_string()))?;
1414    sync_directory(&inner.dir)?;
1415    let file = fs::OpenOptions::new()
1416        .append(true)
1417        .open(&final_path)
1418        .map_err(|err| Error::Io(err.to_string()))?;
1419    inner.open_segment = Some(OpenSegment {
1420        config_id: first.config_id,
1421        path: final_path,
1422        file,
1423        bytes_len: QLOG_HEADER_LEN,
1424        entries: Vec::new(),
1425    });
1426    Ok(())
1427}
1428
1429fn seal_open_segment(inner: &mut FileLogStoreInner) -> Result<()> {
1430    let Some(open) = &inner.open_segment else {
1431        return Ok(());
1432    };
1433    open.file
1434        .sync_data()
1435        .map_err(|err| Error::Io(err.to_string()))?;
1436    let entries = open.entries.clone();
1437    let open_path = open.path.clone();
1438
1439    if entries.is_empty() {
1440        fs::remove_file(&open_path).map_err(|err| Error::Io(err.to_string()))?;
1441        inner.open_segment = None;
1442        sync_directory(&inner.dir)?;
1443        return Ok(());
1444    }
1445
1446    let range = IndexRange::new(entries[0].index, entries.last().expect("non-empty").index)?;
1447    let final_path = inner.dir.join(segment_file_name(range));
1448    if final_path.exists() {
1449        let existing = fs::read(&final_path).map_err(|err| Error::Io(err.to_string()))?;
1450        if decode_segment_for_cluster(&existing, &inner.cluster_id)? != entries {
1451            return Err(Error::Decode(
1452                "open and closed qlog segments disagree".into(),
1453            ));
1454        }
1455    } else {
1456        publish_closed_segment(&inner.dir, &entries)?;
1457    }
1458    fs::remove_file(&open_path).map_err(|err| Error::Io(err.to_string()))?;
1459    inner.open_segment = None;
1460    inner.segments.push(ClosedSegment { entries });
1461    sync_directory(&inner.dir)
1462}
1463
1464fn compact_prefix_with_hook(
1465    inner: &mut FileLogStoreInner,
1466    anchor: &RecoveryAnchor,
1467    hook: &mut impl FnMut(CompactPhase) -> Result<()>,
1468) -> Result<()> {
1469    if !validate_compaction(inner, anchor)? {
1470        return Ok(());
1471    }
1472    let target = anchor.compacted().index();
1473
1474    let old_segments = inner
1475        .segments
1476        .iter()
1477        .take_while(|segment| segment.start() <= target)
1478        .collect::<Vec<_>>();
1479    let old_segment_names = old_segments
1480        .iter()
1481        .map(|segment| {
1482            segment_file_name(
1483                IndexRange::new(segment.start(), segment.end())
1484                    .expect("closed segment range is valid"),
1485            )
1486        })
1487        .collect::<Vec<_>>();
1488    let replacement_entries = old_segments
1489        .last()
1490        .filter(|segment| segment.end() > target)
1491        .map(|segment| {
1492            segment
1493                .entries
1494                .iter()
1495                .filter(|entry| entry.index > target)
1496                .cloned()
1497                .collect::<Vec<_>>()
1498        });
1499    let replacement = match replacement_entries {
1500        Some(entries) if !entries.is_empty() => {
1501            let first = entries.first().expect("replacement is non-empty");
1502            let last = entries.last().expect("replacement is non-empty");
1503            let final_name = segment_file_name(IndexRange::new(first.index, last.index)?);
1504            let final_path = inner.dir.join(&final_name);
1505            let (temp_path, mut file) = create_unique_temp_file(&inner.dir, &final_path)?;
1506            file.write_all(&encode_segment(&entries))
1507                .and_then(|_| file.sync_all())
1508                .map_err(|err| Error::Io(err.to_string()))?;
1509            drop(file);
1510            sync_directory(&inner.dir)?;
1511            hook(CompactPhase::ReplacementPrepared)?;
1512            Some(TruncateReplacement {
1513                temp_name: file_name(&temp_path)?,
1514                final_name,
1515            })
1516        }
1517        _ => None,
1518    };
1519    let intent = CompactIntent {
1520        previous_anchor: inner.anchor.clone(),
1521        anchor: anchor.clone(),
1522        old_segment_names,
1523        replacement,
1524    };
1525    publish_compact_intent(&inner.dir, &intent, hook)?;
1526    apply_compact_intent(&inner.dir, &intent, hook)
1527}
1528
1529fn validate_compaction(inner: &FileLogStoreInner, anchor: &RecoveryAnchor) -> Result<bool> {
1530    validate_anchor(anchor)?;
1531    validate_anchor_identity(Some(anchor), &inner.cluster_id, inner.epoch)?;
1532    let target = anchor.compacted().index();
1533    let tip = inner.last_index();
1534    if tip.is_none_or(|tip| target > tip) {
1535        return Err(Error::CompactionAboveTip { target, tip });
1536    }
1537    if configuration_state_at(inner, target)? != *anchor.configuration_state() {
1538        return Err(Error::CompactionConflict { index: target });
1539    }
1540
1541    if let Some(current) = &inner.anchor {
1542        let current_index = current.compacted().index();
1543        if target < current_index {
1544            return Err(Error::CompactionRegression {
1545                target,
1546                anchor: current_index,
1547            });
1548        }
1549        if target == current_index {
1550            return if current == anchor {
1551                Ok(false)
1552            } else {
1553                Err(Error::CompactionConflict { index: target })
1554            };
1555        }
1556        if anchor.recovery_generation() != current.recovery_generation() {
1557            return Err(Error::CompactionConflict { index: target });
1558        }
1559    }
1560
1561    let entry = inner
1562        .segments
1563        .iter()
1564        .flat_map(|segment| &segment.entries)
1565        .chain(
1566            inner
1567                .open_segment
1568                .iter()
1569                .flat_map(|segment| &segment.entries),
1570        )
1571        .find(|entry| entry.index == target)
1572        .ok_or(Error::CompactionAboveTip { target, tip })?;
1573    if entry.hash != anchor.compacted().hash() {
1574        return Err(Error::CompactionHashMismatch { index: target });
1575    }
1576    Ok(true)
1577}
1578
1579fn configuration_state_at(
1580    inner: &FileLogStoreInner,
1581    target: LogIndex,
1582) -> Result<ConfigurationState> {
1583    let (mut state, start) = match &inner.anchor {
1584        Some(anchor) => (
1585            anchor.configuration_state().clone(),
1586            anchor.compacted().index(),
1587        ),
1588        None => (inner.initial_configuration.clone(), 0),
1589    };
1590    if target == start {
1591        return Ok(state);
1592    }
1593    let mut found = false;
1594    for entry in inner
1595        .segments
1596        .iter()
1597        .flat_map(|segment| &segment.entries)
1598        .chain(
1599            inner
1600                .open_segment
1601                .iter()
1602                .flat_map(|segment| &segment.entries),
1603        )
1604    {
1605        if entry.index > target {
1606            break;
1607        }
1608        state = state
1609            .validate_entry(entry)
1610            .map_err(|err| Error::Decode(err.to_string()))?;
1611        found = entry.index == target;
1612    }
1613    if found {
1614        Ok(state)
1615    } else {
1616        Err(Error::CompactionAboveTip {
1617            target,
1618            tip: inner.last_index(),
1619        })
1620    }
1621}
1622
1623fn read_anchor(dir: &Path) -> Result<Option<RecoveryAnchor>> {
1624    let path = dir.join(ANCHOR_FILE_NAME);
1625    if !path.exists() {
1626        return Ok(None);
1627    }
1628    let bytes = fs::read(path).map_err(|err| Error::Io(err.to_string()))?;
1629    decode_anchor(&bytes).map(Some)
1630}
1631
1632fn validate_anchor_identity(
1633    anchor: Option<&RecoveryAnchor>,
1634    cluster_id: &str,
1635    epoch: u64,
1636) -> Result<()> {
1637    if let Some(anchor) = anchor {
1638        validate_anchor(anchor)?;
1639        if anchor.cluster_id() != cluster_id || anchor.epoch() != epoch {
1640            return Err(Error::Decode(
1641                "recovery anchor does not match configured cluster/epoch".into(),
1642            ));
1643        }
1644    }
1645    Ok(())
1646}
1647
1648fn validate_anchor(anchor: &RecoveryAnchor) -> Result<()> {
1649    if anchor.format_version() != RECOVERY_ANCHOR_FORMAT_VERSION {
1650        return Err(Error::Decode("unsupported recovery anchor version".into()));
1651    }
1652    if anchor.cluster_id().is_empty()
1653        || anchor.recovery_generation() == 0
1654        || anchor.compacted().index() == 0
1655        || anchor.snapshot().snapshot_id().is_empty()
1656        || anchor.snapshot().size_bytes() == 0
1657    {
1658        return Err(Error::Decode("invalid recovery anchor identity".into()));
1659    }
1660    if anchor.configuration_state().config_id() != anchor.config_id()
1661        || anchor
1662            .configuration_state()
1663            .stop()
1664            .is_some_and(|stop| stop != anchor.compacted())
1665    {
1666        return Err(Error::Decode(
1667            "invalid recovery anchor configuration state".into(),
1668        ));
1669    }
1670    Ok(())
1671}
1672
1673fn encode_anchor(anchor: &RecoveryAnchor) -> Result<Vec<u8>> {
1674    validate_anchor(anchor)?;
1675    let mut out = Vec::new();
1676    out.extend_from_slice(&ANCHOR_MAGIC);
1677    put_u16(&mut out, ANCHOR_VERSION);
1678    put_u16(&mut out, 0);
1679    put_u64(&mut out, anchor.epoch());
1680    put_u64(&mut out, anchor.config_id());
1681    put_u64(&mut out, anchor.recovery_generation());
1682    put_u64(&mut out, anchor.compacted().index());
1683    out.extend_from_slice(anchor.compacted().hash().as_bytes());
1684    out.extend_from_slice(anchor.snapshot().digest().as_bytes());
1685    put_u64(&mut out, anchor.snapshot().size_bytes());
1686    out.push(1);
1687    out.extend_from_slice(anchor.executor_fingerprint().as_bytes());
1688    encode_configuration_state(&mut out, anchor.configuration_state())?;
1689    put_string(&mut out, anchor.cluster_id(), "anchor cluster_id")?;
1690    put_string(
1691        &mut out,
1692        anchor.snapshot().snapshot_id(),
1693        "anchor snapshot_id",
1694    )?;
1695    let crc = crc32c(&out);
1696    put_u32(&mut out, crc);
1697    Ok(out)
1698}
1699
1700fn decode_anchor(bytes: &[u8]) -> Result<RecoveryAnchor> {
1701    if bytes.len() < 120 || bytes.get(..4) != Some(ANCHOR_MAGIC.as_slice()) {
1702        return Err(Error::Decode("invalid recovery anchor magic".into()));
1703    }
1704    let crc_offset = bytes.len() - 4;
1705    if crc32c(&bytes[..crc_offset]) != read_u32(bytes, crc_offset)? {
1706        return Err(Error::Decode("recovery anchor crc mismatch".into()));
1707    }
1708    let version = read_u16(bytes, 4)?;
1709    if version != ANCHOR_VERSION || read_u16(bytes, 6)? != 0 {
1710        return Err(Error::Decode("unsupported recovery anchor version".into()));
1711    }
1712    if bytes.get(112) != Some(&1) {
1713        return Err(Error::Decode(
1714            "invalid executor fingerprint encoding".into(),
1715        ));
1716    }
1717    let executor_fingerprint = read_hash(bytes, 113)?;
1718    let mut cursor = 145;
1719    let config_id = read_u64(bytes, 16)?;
1720    let configuration_state =
1721        decode_configuration_state(bytes, &mut cursor, crc_offset, config_id)?;
1722    let cluster_id = read_string(bytes, &mut cursor, crc_offset, "anchor cluster_id")?;
1723    let snapshot_id = read_string(bytes, &mut cursor, crc_offset, "anchor snapshot_id")?;
1724    if cursor != crc_offset {
1725        return Err(Error::Decode("trailing recovery anchor bytes".into()));
1726    }
1727    let compacted = LogAnchor::new(read_u64(bytes, 32)?, read_hash(bytes, 40)?);
1728    let snapshot = SnapshotIdentity::new(
1729        snapshot_id,
1730        read_hash(bytes, 72)?,
1731        read_u64(bytes, 104)?,
1732        executor_fingerprint,
1733    );
1734    let anchor = RecoveryAnchor::new(
1735        cluster_id,
1736        read_u64(bytes, 8)?,
1737        configuration_state,
1738        read_u64(bytes, 24)?,
1739        compacted,
1740        snapshot,
1741    );
1742    validate_anchor(&anchor)?;
1743    Ok(anchor)
1744}
1745
1746fn encode_configuration_state(out: &mut Vec<u8>, state: &ConfigurationState) -> Result<()> {
1747    match state {
1748        ConfigurationState::Active { digest, .. } => {
1749            out.push(1);
1750            out.extend_from_slice(digest.as_bytes());
1751        }
1752        ConfigurationState::Stopped {
1753            digest,
1754            stop,
1755            binding,
1756            ..
1757        } => {
1758            out.push(2);
1759            out.extend_from_slice(digest.as_bytes());
1760            put_u64(out, stop.index());
1761            out.extend_from_slice(stop.hash().as_bytes());
1762            match binding {
1763                StopBinding::Unbound => out.push(1),
1764                StopBinding::Bound {
1765                    successor,
1766                    stop_command_hash,
1767                } => {
1768                    out.push(2);
1769                    encode_successor_descriptor(out, successor)?;
1770                    out.extend_from_slice(stop_command_hash.as_bytes());
1771                }
1772            }
1773        }
1774    }
1775    Ok(())
1776}
1777
1778fn decode_configuration_state(
1779    bytes: &[u8],
1780    cursor: &mut usize,
1781    end: usize,
1782    config_id: u64,
1783) -> Result<ConfigurationState> {
1784    let kind = *bytes
1785        .get(*cursor)
1786        .filter(|_| *cursor < end)
1787        .ok_or_else(|| Error::Decode("short recovery anchor configuration state".into()))?;
1788    *cursor += 1;
1789    let digest = read_state_hash(bytes, cursor, end)?;
1790    match kind {
1791        1 => Ok(ConfigurationState::active(config_id, digest)),
1792        2 => {
1793            let stop_index = read_state_u64(bytes, cursor, end)?;
1794            let stop_hash = read_state_hash(bytes, cursor, end)?;
1795            let binding = decode_stop_binding(bytes, cursor, end)?;
1796            Ok(ConfigurationState::Stopped {
1797                config_id,
1798                digest,
1799                stop: LogAnchor::new(stop_index, stop_hash),
1800                binding,
1801            })
1802        }
1803        _ => Err(Error::Decode(
1804            "invalid recovery anchor configuration state".into(),
1805        )),
1806    }
1807}
1808
1809fn encode_successor_descriptor(out: &mut Vec<u8>, successor: &SuccessorDescriptor) -> Result<()> {
1810    put_string(out, successor.cluster_id(), "successor cluster_id")?;
1811    put_u64(out, successor.predecessor_config_id());
1812    out.extend_from_slice(successor.predecessor_config_digest().as_bytes());
1813    put_u64(out, successor.config_id());
1814    out.extend_from_slice(successor.digest().as_bytes());
1815    put_u16(out, successor.members().len() as u16);
1816    for member in successor.members() {
1817        put_string(out, member, "successor member")?;
1818    }
1819    Ok(())
1820}
1821
1822fn decode_stop_binding(bytes: &[u8], cursor: &mut usize, end: usize) -> Result<StopBinding> {
1823    let kind = read_state_u8(bytes, cursor, end)?;
1824    match kind {
1825        1 => Ok(StopBinding::Unbound),
1826        2 => {
1827            let cluster_id = read_string(bytes, cursor, end, "successor cluster_id")?;
1828            let predecessor_config_id = read_state_u64(bytes, cursor, end)?;
1829            let predecessor_config_digest = read_state_hash(bytes, cursor, end)?;
1830            let successor_config_id = read_state_u64(bytes, cursor, end)?;
1831            let encoded_digest = read_state_hash(bytes, cursor, end)?;
1832            let member_count = usize::from(read_state_u16(bytes, cursor, end)?);
1833            let mut members = Vec::with_capacity(member_count);
1834            for _ in 0..member_count {
1835                members.push(read_string(bytes, cursor, end, "successor member")?);
1836            }
1837            let successor = SuccessorDescriptor::new(
1838                cluster_id,
1839                predecessor_config_id,
1840                predecessor_config_digest,
1841                successor_config_id,
1842                members,
1843            )
1844            .map_err(|_| Error::Decode("invalid recovery anchor successor descriptor".into()))?;
1845            if successor.digest() != encoded_digest {
1846                return Err(Error::Decode(
1847                    "recovery anchor successor digest mismatch".into(),
1848                ));
1849            }
1850            let stop_command_hash = read_state_hash(bytes, cursor, end)?;
1851            Ok(StopBinding::Bound {
1852                successor,
1853                stop_command_hash,
1854            })
1855        }
1856        _ => Err(Error::Decode("invalid recovery anchor stop binding".into())),
1857    }
1858}
1859
1860fn read_state_u8(bytes: &[u8], cursor: &mut usize, end: usize) -> Result<u8> {
1861    let value = *bytes
1862        .get(*cursor)
1863        .filter(|_| *cursor < end)
1864        .ok_or_else(|| Error::Decode("short recovery anchor configuration state".into()))?;
1865    *cursor += 1;
1866    Ok(value)
1867}
1868
1869fn read_state_u16(bytes: &[u8], cursor: &mut usize, end: usize) -> Result<u16> {
1870    let next = cursor
1871        .checked_add(2)
1872        .filter(|next| *next <= end)
1873        .ok_or_else(|| Error::Decode("short recovery anchor configuration state".into()))?;
1874    let value = read_u16(bytes, *cursor)?;
1875    *cursor = next;
1876    Ok(value)
1877}
1878
1879fn read_state_u64(bytes: &[u8], cursor: &mut usize, end: usize) -> Result<u64> {
1880    let next = cursor
1881        .checked_add(8)
1882        .filter(|next| *next <= end)
1883        .ok_or_else(|| Error::Decode("short recovery anchor configuration state".into()))?;
1884    let value = read_u64(bytes, *cursor)?;
1885    *cursor = next;
1886    Ok(value)
1887}
1888
1889fn read_state_hash(bytes: &[u8], cursor: &mut usize, end: usize) -> Result<LogHash> {
1890    let next = cursor
1891        .checked_add(32)
1892        .filter(|next| *next <= end)
1893        .ok_or_else(|| Error::Decode("short recovery anchor configuration state".into()))?;
1894    let value = read_hash(bytes, *cursor)?;
1895    *cursor = next;
1896    Ok(value)
1897}
1898
1899fn recover_compact_intent(dir: &Path) -> Result<()> {
1900    let path = dir.join(COMPACT_INTENT_FILE_NAME);
1901    if !path.exists() {
1902        return Ok(());
1903    }
1904    let bytes = fs::read(path).map_err(|err| Error::Io(err.to_string()))?;
1905    let intent = decode_compact_intent(&bytes)?;
1906    apply_compact_intent(dir, &intent, &mut |_| Ok(()))
1907}
1908
1909fn publish_compact_intent(
1910    dir: &Path,
1911    intent: &CompactIntent,
1912    hook: &mut impl FnMut(CompactPhase) -> Result<()>,
1913) -> Result<()> {
1914    let final_path = dir.join(COMPACT_INTENT_FILE_NAME);
1915    if final_path.exists() {
1916        return Err(Error::Io("compact intent already exists".into()));
1917    }
1918    let (temp_path, mut file) = create_unique_temp_file(dir, &final_path)?;
1919    let bytes = encode_compact_intent(intent)?;
1920    if let Err(err) = file
1921        .write_all(&bytes)
1922        .and_then(|_| file.sync_all())
1923        .map_err(|err| Error::Io(err.to_string()))
1924    {
1925        drop(file);
1926        let _ = fs::remove_file(&temp_path);
1927        return Err(err);
1928    }
1929    drop(file);
1930    fs::rename(&temp_path, &final_path).map_err(|err| Error::Io(err.to_string()))?;
1931    hook(CompactPhase::IntentRenamed)?;
1932    sync_directory(dir)?;
1933    hook(CompactPhase::IntentDurable)
1934}
1935
1936fn apply_compact_intent(
1937    dir: &Path,
1938    intent: &CompactIntent,
1939    hook: &mut impl FnMut(CompactPhase) -> Result<()>,
1940) -> Result<()> {
1941    validate_compact_intent(intent)?;
1942    install_anchor(dir, intent)?;
1943    hook(CompactPhase::AnchorInstalled)?;
1944    sync_directory(dir)?;
1945    hook(CompactPhase::AnchorDurable)?;
1946
1947    for (position, name) in intent.old_segment_names.iter().enumerate() {
1948        let path = dir.join(name);
1949        if path.exists() {
1950            fs::remove_file(path).map_err(|err| Error::Io(err.to_string()))?;
1951        }
1952        hook(CompactPhase::OldSegmentRemoved(position))?;
1953    }
1954    install_replacement(dir, intent.replacement.as_ref(), "compact")?;
1955    hook(CompactPhase::ReplacementInstalled)?;
1956    sync_directory(dir)?;
1957    hook(CompactPhase::AppliedDirectorySynced)?;
1958
1959    let path = dir.join(COMPACT_INTENT_FILE_NAME);
1960    if path.exists() {
1961        fs::remove_file(path).map_err(|err| Error::Io(err.to_string()))?;
1962    }
1963    hook(CompactPhase::IntentRemoved)?;
1964    sync_directory(dir)?;
1965    hook(CompactPhase::CompleteDirectorySynced)
1966}
1967
1968fn install_anchor(dir: &Path, intent: &CompactIntent) -> Result<()> {
1969    let current = read_anchor(dir)?;
1970    if current.as_ref() == Some(&intent.anchor) {
1971        return Ok(());
1972    }
1973    if current != intent.previous_anchor {
1974        return Err(Error::CompactionConflict {
1975            index: intent.anchor.compacted().index(),
1976        });
1977    }
1978    let final_path = dir.join(ANCHOR_FILE_NAME);
1979    let (temp_path, mut file) = create_unique_temp_file(dir, &final_path)?;
1980    file.write_all(&encode_anchor(&intent.anchor)?)
1981        .and_then(|_| file.sync_all())
1982        .map_err(|err| Error::Io(err.to_string()))?;
1983    drop(file);
1984    fs::rename(temp_path, final_path).map_err(|err| Error::Io(err.to_string()))
1985}
1986
1987fn install_replacement(
1988    dir: &Path,
1989    replacement: Option<&TruncateReplacement>,
1990    operation: &str,
1991) -> Result<()> {
1992    let Some(replacement) = replacement else {
1993        return Ok(());
1994    };
1995    let temp_path = dir.join(&replacement.temp_name);
1996    let final_path = dir.join(&replacement.final_name);
1997    match (temp_path.exists(), final_path.exists()) {
1998        (true, false) => {
1999            fs::rename(temp_path, final_path).map_err(|err| Error::Io(err.to_string()))?;
2000        }
2001        (true, true) => {
2002            let temp = fs::read(&temp_path).map_err(|err| Error::Io(err.to_string()))?;
2003            let final_bytes = fs::read(&final_path).map_err(|err| Error::Io(err.to_string()))?;
2004            if temp != final_bytes {
2005                return Err(Error::Decode(format!(
2006                    "{operation} replacement files disagree"
2007                )));
2008            }
2009            fs::remove_file(temp_path).map_err(|err| Error::Io(err.to_string()))?;
2010        }
2011        (false, true) => {}
2012        (false, false) => {
2013            return Err(Error::Decode(format!(
2014                "{operation} replacement file is missing"
2015            )));
2016        }
2017    }
2018    Ok(())
2019}
2020
2021fn truncate_suffix_with_hook(
2022    inner: &mut FileLogStoreInner,
2023    from: LogIndex,
2024    hook: &mut impl FnMut(TruncatePhase) -> Result<()>,
2025) -> Result<()> {
2026    let Some(position) = inner
2027        .segments
2028        .iter()
2029        .position(|segment| segment.end() >= from)
2030    else {
2031        return Ok(());
2032    };
2033
2034    let old_segment_names = inner.segments[position..]
2035        .iter()
2036        .map(|segment| {
2037            segment_file_name(
2038                IndexRange::new(segment.start(), segment.end())
2039                    .expect("closed segment range is valid"),
2040            )
2041        })
2042        .collect::<Vec<_>>();
2043
2044    let replacement_entries = (inner.segments[position].start() < from).then(|| {
2045        inner.segments[position]
2046            .entries
2047            .iter()
2048            .take_while(|entry| entry.index < from)
2049            .cloned()
2050            .collect::<Vec<_>>()
2051    });
2052    let replacement = match replacement_entries {
2053        Some(entries) if !entries.is_empty() => {
2054            let first = entries.first().expect("replacement is non-empty");
2055            let last = entries.last().expect("replacement is non-empty");
2056            let final_name = segment_file_name(IndexRange::new(first.index, last.index)?);
2057            let final_path = inner.dir.join(&final_name);
2058            let (temp_path, mut file) = create_unique_temp_file(&inner.dir, &final_path)?;
2059            file.write_all(&encode_segment(&entries))
2060                .and_then(|_| file.sync_all())
2061                .map_err(|err| Error::Io(err.to_string()))?;
2062            drop(file);
2063            sync_directory(&inner.dir)?;
2064            hook(TruncatePhase::ReplacementPrepared)?;
2065            Some(TruncateReplacement {
2066                temp_name: file_name(&temp_path)?,
2067                final_name,
2068            })
2069        }
2070        _ => None,
2071    };
2072
2073    let intent = TruncateIntent {
2074        old_segment_names,
2075        replacement,
2076    };
2077    publish_truncate_intent(&inner.dir, &intent, hook)?;
2078    apply_truncate_intent(&inner.dir, &intent, hook)
2079}
2080
2081fn recover_truncate_intent(dir: &Path) -> Result<()> {
2082    let path = dir.join(TRUNCATE_INTENT_FILE_NAME);
2083    if !path.exists() {
2084        return Ok(());
2085    }
2086    let bytes = fs::read(&path).map_err(|err| Error::Io(err.to_string()))?;
2087    let intent = decode_truncate_intent(&bytes)?;
2088    apply_truncate_intent(dir, &intent, &mut |_| Ok(()))
2089}
2090
2091fn publish_truncate_intent(
2092    dir: &Path,
2093    intent: &TruncateIntent,
2094    hook: &mut impl FnMut(TruncatePhase) -> Result<()>,
2095) -> Result<()> {
2096    let final_path = dir.join(TRUNCATE_INTENT_FILE_NAME);
2097    if final_path.exists() {
2098        return Err(Error::Io("truncate intent already exists".into()));
2099    }
2100    let (temp_path, mut file) = create_unique_temp_file(dir, &final_path)?;
2101    let bytes = encode_truncate_intent(intent)?;
2102    if let Err(err) = file
2103        .write_all(&bytes)
2104        .and_then(|_| file.sync_all())
2105        .map_err(|err| Error::Io(err.to_string()))
2106    {
2107        drop(file);
2108        let _ = fs::remove_file(&temp_path);
2109        return Err(err);
2110    }
2111    drop(file);
2112    fs::rename(&temp_path, &final_path).map_err(|err| Error::Io(err.to_string()))?;
2113    hook(TruncatePhase::IntentRenamed)?;
2114    sync_directory(dir)?;
2115    hook(TruncatePhase::IntentDurable)
2116}
2117
2118fn apply_truncate_intent(
2119    dir: &Path,
2120    intent: &TruncateIntent,
2121    hook: &mut impl FnMut(TruncatePhase) -> Result<()>,
2122) -> Result<()> {
2123    validate_truncate_intent(intent)?;
2124    for (position, name) in intent.old_segment_names.iter().enumerate() {
2125        let path = dir.join(name);
2126        if path.exists() {
2127            fs::remove_file(&path).map_err(|err| Error::Io(err.to_string()))?;
2128        }
2129        hook(TruncatePhase::OldSegmentRemoved(position))?;
2130    }
2131
2132    if let Some(replacement) = &intent.replacement {
2133        let temp_path = dir.join(&replacement.temp_name);
2134        let final_path = dir.join(&replacement.final_name);
2135        match (temp_path.exists(), final_path.exists()) {
2136            (true, false) => {
2137                fs::rename(&temp_path, &final_path).map_err(|err| Error::Io(err.to_string()))?;
2138            }
2139            (true, true) => {
2140                let temp = fs::read(&temp_path).map_err(|err| Error::Io(err.to_string()))?;
2141                let final_bytes =
2142                    fs::read(&final_path).map_err(|err| Error::Io(err.to_string()))?;
2143                if temp != final_bytes {
2144                    return Err(Error::Decode("truncate replacement files disagree".into()));
2145                }
2146                fs::remove_file(&temp_path).map_err(|err| Error::Io(err.to_string()))?;
2147            }
2148            (false, true) => {}
2149            (false, false) => {
2150                return Err(Error::Decode("truncate replacement file is missing".into()));
2151            }
2152        }
2153    }
2154    hook(TruncatePhase::ReplacementInstalled)?;
2155    sync_directory(dir)?;
2156    hook(TruncatePhase::AppliedDirectorySynced)?;
2157
2158    let intent_path = dir.join(TRUNCATE_INTENT_FILE_NAME);
2159    if intent_path.exists() {
2160        fs::remove_file(&intent_path).map_err(|err| Error::Io(err.to_string()))?;
2161    }
2162    hook(TruncatePhase::IntentRemoved)?;
2163    sync_directory(dir)?;
2164    hook(TruncatePhase::CompleteDirectorySynced)
2165}
2166
2167fn encode_truncate_intent(intent: &TruncateIntent) -> Result<Vec<u8>> {
2168    validate_truncate_intent(intent)?;
2169    let mut out = Vec::new();
2170    out.extend_from_slice(&TRUNCATE_INTENT_MAGIC);
2171    put_u16(&mut out, TRUNCATE_INTENT_VERSION);
2172    put_u16(
2173        &mut out,
2174        if intent.replacement.is_some() {
2175            TRUNCATE_INTENT_REPLACEMENT
2176        } else {
2177            0
2178        },
2179    );
2180    let count = u32::try_from(intent.old_segment_names.len())
2181        .map_err(|_| Error::Decode("too many truncate segments".into()))?;
2182    put_u32(&mut out, count);
2183    for name in &intent.old_segment_names {
2184        put_intent_name(&mut out, name)?;
2185    }
2186    if let Some(replacement) = &intent.replacement {
2187        put_intent_name(&mut out, &replacement.temp_name)?;
2188        put_intent_name(&mut out, &replacement.final_name)?;
2189    }
2190    let crc = crc32c(&out);
2191    put_u32(&mut out, crc);
2192    Ok(out)
2193}
2194
2195fn decode_truncate_intent(bytes: &[u8]) -> Result<TruncateIntent> {
2196    if bytes.len() < 16 || bytes.get(..4) != Some(TRUNCATE_INTENT_MAGIC.as_slice()) {
2197        return Err(Error::Decode("invalid truncate intent magic".into()));
2198    }
2199    let crc_offset = bytes.len() - 4;
2200    if crc32c(&bytes[..crc_offset]) != read_u32(bytes, crc_offset)? {
2201        return Err(Error::Decode("truncate intent crc mismatch".into()));
2202    }
2203    if read_u16(bytes, 4)? != TRUNCATE_INTENT_VERSION {
2204        return Err(Error::Decode("unsupported truncate intent version".into()));
2205    }
2206    let flags = read_u16(bytes, 6)?;
2207    if flags & !TRUNCATE_INTENT_REPLACEMENT != 0 {
2208        return Err(Error::Decode("invalid truncate intent flags".into()));
2209    }
2210    let count = read_u32(bytes, 8)? as usize;
2211    let mut cursor = 12;
2212    let mut old_segment_names = Vec::with_capacity(count);
2213    for _ in 0..count {
2214        old_segment_names.push(read_intent_name(bytes, &mut cursor, crc_offset)?);
2215    }
2216    let replacement = if flags & TRUNCATE_INTENT_REPLACEMENT != 0 {
2217        Some(TruncateReplacement {
2218            temp_name: read_intent_name(bytes, &mut cursor, crc_offset)?,
2219            final_name: read_intent_name(bytes, &mut cursor, crc_offset)?,
2220        })
2221    } else {
2222        None
2223    };
2224    if cursor != crc_offset {
2225        return Err(Error::Decode("trailing truncate intent bytes".into()));
2226    }
2227    let intent = TruncateIntent {
2228        old_segment_names,
2229        replacement,
2230    };
2231    validate_truncate_intent(&intent)?;
2232    Ok(intent)
2233}
2234
2235fn encode_compact_intent(intent: &CompactIntent) -> Result<Vec<u8>> {
2236    validate_compact_intent(intent)?;
2237    let mut flags = 0;
2238    if intent.previous_anchor.is_some() {
2239        flags |= COMPACT_INTENT_PREVIOUS_ANCHOR;
2240    }
2241    if intent.replacement.is_some() {
2242        flags |= COMPACT_INTENT_REPLACEMENT;
2243    }
2244    let mut out = Vec::new();
2245    out.extend_from_slice(&COMPACT_INTENT_MAGIC);
2246    put_u16(&mut out, COMPACT_INTENT_VERSION);
2247    put_u16(&mut out, flags);
2248    put_u32(
2249        &mut out,
2250        u32::try_from(intent.old_segment_names.len())
2251            .map_err(|_| Error::Decode("too many compact segments".into()))?,
2252    );
2253    put_blob(&mut out, &encode_anchor(&intent.anchor)?, "compact anchor")?;
2254    if let Some(previous) = &intent.previous_anchor {
2255        put_blob(
2256            &mut out,
2257            &encode_anchor(previous)?,
2258            "compact previous anchor",
2259        )?;
2260    }
2261    for name in &intent.old_segment_names {
2262        put_intent_name(&mut out, name)?;
2263    }
2264    if let Some(replacement) = &intent.replacement {
2265        put_intent_name(&mut out, &replacement.temp_name)?;
2266        put_intent_name(&mut out, &replacement.final_name)?;
2267    }
2268    let crc = crc32c(&out);
2269    put_u32(&mut out, crc);
2270    Ok(out)
2271}
2272
2273fn decode_compact_intent(bytes: &[u8]) -> Result<CompactIntent> {
2274    if bytes.len() < 20 || bytes.get(..4) != Some(COMPACT_INTENT_MAGIC.as_slice()) {
2275        return Err(Error::Decode("invalid compact intent magic".into()));
2276    }
2277    let crc_offset = bytes.len() - 4;
2278    if crc32c(&bytes[..crc_offset]) != read_u32(bytes, crc_offset)? {
2279        return Err(Error::Decode("compact intent crc mismatch".into()));
2280    }
2281    if read_u16(bytes, 4)? != COMPACT_INTENT_VERSION {
2282        return Err(Error::Decode("unsupported compact intent version".into()));
2283    }
2284    let flags = read_u16(bytes, 6)?;
2285    if flags & !(COMPACT_INTENT_PREVIOUS_ANCHOR | COMPACT_INTENT_REPLACEMENT) != 0 {
2286        return Err(Error::Decode("invalid compact intent flags".into()));
2287    }
2288    let count = read_u32(bytes, 8)? as usize;
2289    let mut cursor = 12;
2290    let anchor = decode_anchor(read_blob(bytes, &mut cursor, crc_offset, "compact anchor")?)?;
2291    let previous_anchor = if flags & COMPACT_INTENT_PREVIOUS_ANCHOR != 0 {
2292        Some(decode_anchor(read_blob(
2293            bytes,
2294            &mut cursor,
2295            crc_offset,
2296            "compact previous anchor",
2297        )?)?)
2298    } else {
2299        None
2300    };
2301    let mut old_segment_names = Vec::with_capacity(count);
2302    for _ in 0..count {
2303        old_segment_names.push(read_intent_name(bytes, &mut cursor, crc_offset)?);
2304    }
2305    let replacement = if flags & COMPACT_INTENT_REPLACEMENT != 0 {
2306        Some(TruncateReplacement {
2307            temp_name: read_intent_name(bytes, &mut cursor, crc_offset)?,
2308            final_name: read_intent_name(bytes, &mut cursor, crc_offset)?,
2309        })
2310    } else {
2311        None
2312    };
2313    if cursor != crc_offset {
2314        return Err(Error::Decode("trailing compact intent bytes".into()));
2315    }
2316    let intent = CompactIntent {
2317        previous_anchor,
2318        anchor,
2319        old_segment_names,
2320        replacement,
2321    };
2322    validate_compact_intent(&intent)?;
2323    Ok(intent)
2324}
2325
2326fn validate_compact_intent(intent: &CompactIntent) -> Result<()> {
2327    validate_anchor(&intent.anchor)?;
2328    if intent.old_segment_names.is_empty() {
2329        return Err(Error::Decode("compact intent has no old segments".into()));
2330    }
2331    if let Some(previous) = &intent.previous_anchor {
2332        validate_anchor(previous)?;
2333        if previous.cluster_id() != intent.anchor.cluster_id()
2334            || previous.epoch() != intent.anchor.epoch()
2335            || previous.recovery_generation() != intent.anchor.recovery_generation()
2336            || previous.compacted().index() >= intent.anchor.compacted().index()
2337        {
2338            return Err(Error::Decode(
2339                "compact intent anchor regression or conflict".into(),
2340            ));
2341        }
2342    }
2343    for (position, name) in intent.old_segment_names.iter().enumerate() {
2344        validate_closed_segment_name(name)?;
2345        if intent.old_segment_names[..position].contains(name) {
2346            return Err(Error::Decode("duplicate compact segment name".into()));
2347        }
2348    }
2349    if let Some(replacement) = &intent.replacement {
2350        validate_temp_name(&replacement.temp_name)?;
2351        validate_closed_segment_name(&replacement.final_name)?;
2352        if intent.old_segment_names.contains(&replacement.final_name) {
2353            return Err(Error::Decode(
2354                "compact replacement overlaps old segment name".into(),
2355            ));
2356        }
2357        let replacement_range = parse_closed_segment_name(&replacement.final_name)?;
2358        if replacement_range.start()
2359            != intent
2360                .anchor
2361                .compacted()
2362                .index()
2363                .checked_add(1)
2364                .ok_or_else(|| Error::Decode("compact anchor index overflow".into()))?
2365        {
2366            return Err(Error::Decode(
2367                "compact replacement does not start after anchor".into(),
2368            ));
2369        }
2370    }
2371    Ok(())
2372}
2373
2374fn validate_truncate_intent(intent: &TruncateIntent) -> Result<()> {
2375    if intent.old_segment_names.is_empty() {
2376        return Err(Error::Decode("truncate intent has no old segments".into()));
2377    }
2378    for (position, name) in intent.old_segment_names.iter().enumerate() {
2379        validate_closed_segment_name(name)?;
2380        if intent.old_segment_names[..position].contains(name) {
2381            return Err(Error::Decode("duplicate truncate segment name".into()));
2382        }
2383    }
2384    if let Some(replacement) = &intent.replacement {
2385        validate_temp_name(&replacement.temp_name)?;
2386        validate_closed_segment_name(&replacement.final_name)?;
2387        if intent.old_segment_names.contains(&replacement.final_name) {
2388            return Err(Error::Decode(
2389                "truncate replacement overlaps old segment name".into(),
2390            ));
2391        }
2392    }
2393    Ok(())
2394}
2395
2396fn validate_closed_segment_name(name: &str) -> Result<()> {
2397    validate_relative_name(name)?;
2398    let range = parse_closed_segment_name(name)?;
2399    if segment_file_name(range) != name {
2400        return Err(Error::Decode("non-canonical truncate segment name".into()));
2401    }
2402    Ok(())
2403}
2404
2405fn validate_temp_name(name: &str) -> Result<()> {
2406    validate_relative_name(name)?;
2407    if !name.starts_with('.') || !name.ends_with(".tmp") {
2408        return Err(Error::Decode("invalid truncate temp name".into()));
2409    }
2410    Ok(())
2411}
2412
2413fn validate_relative_name(name: &str) -> Result<()> {
2414    let mut components = Path::new(name).components();
2415    if name.is_empty()
2416        || !matches!(components.next(), Some(std::path::Component::Normal(_)))
2417        || components.next().is_some()
2418    {
2419        return Err(Error::Decode("unsafe truncate intent path".into()));
2420    }
2421    Ok(())
2422}
2423
2424fn put_intent_name(out: &mut Vec<u8>, name: &str) -> Result<()> {
2425    let len = u16::try_from(name.len())
2426        .map_err(|_| Error::Decode("truncate intent name is too long".into()))?;
2427    put_u16(out, len);
2428    out.extend_from_slice(name.as_bytes());
2429    Ok(())
2430}
2431
2432fn put_string(out: &mut Vec<u8>, value: &str, field: &str) -> Result<()> {
2433    let len =
2434        u16::try_from(value.len()).map_err(|_| Error::Decode(format!("{field} is too long")))?;
2435    put_u16(out, len);
2436    out.extend_from_slice(value.as_bytes());
2437    Ok(())
2438}
2439
2440fn read_string(bytes: &[u8], cursor: &mut usize, end: usize, field: &str) -> Result<String> {
2441    let value = read_intent_name(bytes, cursor, end)?;
2442    if value.is_empty() {
2443        return Err(Error::Decode(format!("{field} is empty")));
2444    }
2445    Ok(value)
2446}
2447
2448fn put_blob(out: &mut Vec<u8>, bytes: &[u8], field: &str) -> Result<()> {
2449    let len =
2450        u32::try_from(bytes.len()).map_err(|_| Error::Decode(format!("{field} is too large")))?;
2451    put_u32(out, len);
2452    out.extend_from_slice(bytes);
2453    Ok(())
2454}
2455
2456fn read_blob<'a>(bytes: &'a [u8], cursor: &mut usize, end: usize, field: &str) -> Result<&'a [u8]> {
2457    let len = read_u32(bytes, *cursor)? as usize;
2458    *cursor = cursor
2459        .checked_add(4)
2460        .ok_or_else(|| Error::Decode(format!("{field} cursor overflow")))?;
2461    let value_end = cursor
2462        .checked_add(len)
2463        .ok_or_else(|| Error::Decode(format!("{field} length overflow")))?;
2464    if value_end > end {
2465        return Err(Error::Decode(format!("short {field}")));
2466    }
2467    let value = &bytes[*cursor..value_end];
2468    *cursor = value_end;
2469    Ok(value)
2470}
2471
2472fn read_intent_name(bytes: &[u8], cursor: &mut usize, end: usize) -> Result<String> {
2473    let len = read_u16(bytes, *cursor)? as usize;
2474    *cursor = cursor
2475        .checked_add(2)
2476        .ok_or_else(|| Error::Decode("truncate intent cursor overflow".into()))?;
2477    let name_end = cursor
2478        .checked_add(len)
2479        .ok_or_else(|| Error::Decode("truncate intent name overflow".into()))?;
2480    if name_end > end {
2481        return Err(Error::Decode("short truncate intent name".into()));
2482    }
2483    let name = std::str::from_utf8(&bytes[*cursor..name_end])
2484        .map_err(|err| Error::Decode(err.to_string()))?
2485        .to_string();
2486    *cursor = name_end;
2487    Ok(name)
2488}
2489
2490fn file_name(path: &Path) -> Result<String> {
2491    path.file_name()
2492        .and_then(|name| name.to_str())
2493        .map(str::to_owned)
2494        .ok_or_else(|| Error::Decode("qlog temp filename is not UTF-8".into()))
2495}
2496
2497fn publish_closed_segment(dir: &Path, entries: &[LogEntry]) -> Result<PathBuf> {
2498    let first = entries
2499        .first()
2500        .ok_or_else(|| Error::Decode("cannot write empty segment".into()))?;
2501    let last = entries
2502        .last()
2503        .ok_or_else(|| Error::Decode("cannot write empty segment".into()))?;
2504    let range = IndexRange::new(first.index, last.index)?;
2505    let final_path = dir.join(segment_file_name(range));
2506    if final_path.exists() {
2507        return Err(Error::Io(format!(
2508            "qlog segment already exists: {}",
2509            final_path.display()
2510        )));
2511    }
2512
2513    let bytes = encode_segment(entries);
2514    let (temp_path, mut file) = create_unique_temp_file(dir, &final_path)?;
2515    if let Err(err) = file
2516        .write_all(&bytes)
2517        .and_then(|_| file.sync_all())
2518        .map_err(|err| Error::Io(err.to_string()))
2519    {
2520        drop(file);
2521        let _ = fs::remove_file(&temp_path);
2522        return Err(err);
2523    }
2524    drop(file);
2525    if let Err(err) = fs::rename(&temp_path, &final_path) {
2526        let _ = fs::remove_file(&temp_path);
2527        return Err(Error::Io(err.to_string()));
2528    }
2529    sync_directory(dir)?;
2530    Ok(final_path)
2531}
2532
2533fn create_unique_temp_file(dir: &Path, final_path: &Path) -> Result<(PathBuf, fs::File)> {
2534    let final_name = final_path
2535        .file_name()
2536        .and_then(|name| name.to_str())
2537        .expect("generated qlog filename is UTF-8");
2538    loop {
2539        let id = NEXT_TEMP_FILE_ID.fetch_add(1, Ordering::Relaxed);
2540        let path = dir.join(format!(".{final_name}.{}.{}.tmp", std::process::id(), id));
2541        match fs::OpenOptions::new()
2542            .write(true)
2543            .create_new(true)
2544            .open(&path)
2545        {
2546            Ok(file) => return Ok((path, file)),
2547            Err(err) if err.kind() == std::io::ErrorKind::AlreadyExists => continue,
2548            Err(err) => return Err(Error::Io(err.to_string())),
2549        }
2550    }
2551}
2552
2553fn sync_directory(dir: &Path) -> Result<()> {
2554    fs::File::open(dir)
2555        .and_then(|directory| directory.sync_all())
2556        .map_err(|err| Error::Io(err.to_string()))
2557}
2558
2559#[cfg(test)]
2560mod tests {
2561    use super::*;
2562    const INJECTED_CRASH: &str = "injected crash";
2563
2564    #[test]
2565    fn legacy_anchor_binary_version_is_rejected() {
2566        let entry = chain(&[b"one"]).pop().unwrap();
2567        let anchor = recovery_anchor(&entry);
2568        let mut bytes = encode_anchor(&anchor).unwrap();
2569        bytes[4..6].copy_from_slice(&3_u16.to_be_bytes());
2570        let crc_offset = bytes.len() - 4;
2571        let crc = crc32c(&bytes[..crc_offset]);
2572        bytes[crc_offset..].copy_from_slice(&crc.to_be_bytes());
2573        let dir = tempfile::tempdir().unwrap();
2574        let path = dir.path().join(ANCHOR_FILE_NAME);
2575        fs::write(&path, &bytes).unwrap();
2576
2577        assert!(matches!(
2578            FileLogStore::open(dir.path(), "cluster-a", 1, 1),
2579            Err(Error::Decode(_))
2580        ));
2581        assert_eq!(fs::read(path).unwrap(), bytes);
2582    }
2583
2584    #[test]
2585    fn compact_crash_before_durable_intent_preserves_genesis_log() {
2586        let dir = tempfile::tempdir().unwrap();
2587        let entries = chain(&[b"one", b"two", b"three", b"four", b"five", b"six"]);
2588        let store = segmented_store(dir.path(), &entries);
2589        let anchor = recovery_anchor(&entries[2]);
2590
2591        inject_compact_crash(&store, &anchor, CompactPhase::ReplacementPrepared);
2592        drop(store);
2593
2594        assert_reopens_with(dir.path(), &entries);
2595        assert!(read_anchor(dir.path()).unwrap().is_none());
2596        assert!(!dir.path().join(COMPACT_INTENT_FILE_NAME).exists());
2597    }
2598
2599    #[test]
2600    fn compact_rolls_forward_from_every_committed_phase() {
2601        let phases = [
2602            CompactPhase::IntentRenamed,
2603            CompactPhase::IntentDurable,
2604            CompactPhase::AnchorInstalled,
2605            CompactPhase::AnchorDurable,
2606            CompactPhase::OldSegmentRemoved(0),
2607            CompactPhase::OldSegmentRemoved(1),
2608            CompactPhase::ReplacementInstalled,
2609            CompactPhase::AppliedDirectorySynced,
2610            CompactPhase::IntentRemoved,
2611            CompactPhase::CompleteDirectorySynced,
2612        ];
2613
2614        for crash_phase in phases {
2615            let dir = tempfile::tempdir().unwrap();
2616            let entries = chain(&[b"one", b"two", b"three", b"four", b"five", b"six"]);
2617            let store = segmented_store(dir.path(), &entries);
2618            let anchor = recovery_anchor(&entries[2]);
2619
2620            inject_compact_crash(&store, &anchor, crash_phase);
2621            drop(store);
2622
2623            let reopened = FileLogStore::open(dir.path(), "cluster-a", 1, 1).unwrap();
2624            assert_eq!(read_anchor(dir.path()).unwrap(), Some(anchor.clone()));
2625            assert_eq!(
2626                reopened.read_range(IndexRange::new(1, 6).unwrap()).unwrap(),
2627                entries[3..]
2628            );
2629            assert_eq!(reopened.last_index().unwrap(), Some(6));
2630            assert_no_overlapping_closed_segments(dir.path());
2631            assert!(!dir.path().join(COMPACT_INTENT_FILE_NAME).exists());
2632        }
2633    }
2634
2635    #[test]
2636    fn corrupted_compact_intent_is_fatal_without_deleting_prefix() {
2637        let dir = tempfile::tempdir().unwrap();
2638        let entries = chain(&[b"one", b"two", b"three", b"four"]);
2639        let store = segmented_store(dir.path(), &entries);
2640        let anchor = recovery_anchor(&entries[2]);
2641        inject_compact_crash(&store, &anchor, CompactPhase::IntentDurable);
2642        drop(store);
2643        let files_before = closed_segment_bytes(dir.path());
2644        let intent_path = dir.path().join(COMPACT_INTENT_FILE_NAME);
2645        let mut bytes = fs::read(&intent_path).unwrap();
2646        bytes[4] ^= 1;
2647        fs::write(&intent_path, bytes).unwrap();
2648
2649        assert!(FileLogStore::open(dir.path(), "cluster-a", 1, 1).is_err());
2650        assert_eq!(closed_segment_bytes(dir.path()), files_before);
2651        assert!(read_anchor(dir.path()).unwrap().is_none());
2652    }
2653
2654    #[test]
2655    fn reopen_preserves_original_log_when_crash_precedes_durable_intent() {
2656        let dir = tempfile::tempdir().unwrap();
2657        let entries = chain(&[b"one", b"two", b"three", b"four", b"five", b"six"]);
2658        let store = segmented_store(dir.path(), &entries);
2659
2660        inject_truncate_crash(&store, 4, TruncatePhase::ReplacementPrepared);
2661        drop(store);
2662
2663        assert_reopens_with(dir.path(), &entries);
2664        assert_no_overlapping_closed_segments(dir.path());
2665        assert!(!dir.path().join(TRUNCATE_INTENT_FILE_NAME).exists());
2666    }
2667
2668    #[test]
2669    fn reopen_recovers_exact_prefix_from_every_durable_intent_phase() {
2670        let phases = [
2671            TruncatePhase::IntentRenamed,
2672            TruncatePhase::IntentDurable,
2673            TruncatePhase::OldSegmentRemoved(0),
2674            TruncatePhase::OldSegmentRemoved(1),
2675            TruncatePhase::ReplacementInstalled,
2676            TruncatePhase::AppliedDirectorySynced,
2677            TruncatePhase::IntentRemoved,
2678            TruncatePhase::CompleteDirectorySynced,
2679        ];
2680
2681        for crash_phase in phases {
2682            let dir = tempfile::tempdir().unwrap();
2683            let entries = chain(&[b"one", b"two", b"three", b"four", b"five", b"six"]);
2684            let store = segmented_store(dir.path(), &entries);
2685
2686            inject_truncate_crash(&store, 4, crash_phase);
2687            drop(store);
2688
2689            assert_reopens_with(dir.path(), &entries[..3]);
2690            assert_no_overlapping_closed_segments(dir.path());
2691            assert!(!dir.path().join(TRUNCATE_INTENT_FILE_NAME).exists());
2692        }
2693    }
2694
2695    #[test]
2696    fn corrupted_truncate_intent_is_fatal_without_removing_segments() {
2697        let dir = tempfile::tempdir().unwrap();
2698        let entries = chain(&[b"one", b"two", b"three", b"four"]);
2699        let store = segmented_store(dir.path(), &entries);
2700        inject_truncate_crash(&store, 3, TruncatePhase::IntentDurable);
2701        drop(store);
2702        let files_before = closed_segment_bytes(dir.path());
2703        let intent_path = dir.path().join(TRUNCATE_INTENT_FILE_NAME);
2704        let mut bytes = fs::read(&intent_path).unwrap();
2705        bytes[4] ^= 1;
2706        fs::write(&intent_path, bytes).unwrap();
2707
2708        assert!(FileLogStore::open(dir.path(), "cluster-a", 1, 1).is_err());
2709        assert_eq!(closed_segment_bytes(dir.path()), files_before);
2710    }
2711
2712    #[test]
2713    fn unsafe_truncate_intent_name_is_fatal_without_path_traversal() {
2714        let root = tempfile::tempdir().unwrap();
2715        let dir = root.path().join("log");
2716        fs::create_dir(&dir).unwrap();
2717        let victim = root.path().join("victim.qlog");
2718        fs::write(&victim, b"keep").unwrap();
2719        let bytes = encode_unchecked_test_intent("../victim.qlog");
2720        fs::write(dir.join(TRUNCATE_INTENT_FILE_NAME), bytes).unwrap();
2721
2722        assert!(FileLogStore::open(&dir, "cluster-a", 1, 1).is_err());
2723        assert_eq!(fs::read(victim).unwrap(), b"keep");
2724    }
2725
2726    fn inject_truncate_crash(store: &FileLogStore, from: LogIndex, crash_phase: TruncatePhase) {
2727        let mut inner = store.lock().unwrap();
2728        let err = truncate_suffix_with_hook(&mut inner, from, &mut |phase| {
2729            if phase == crash_phase {
2730                Err(Error::Io(INJECTED_CRASH.into()))
2731            } else {
2732                Ok(())
2733            }
2734        })
2735        .unwrap_err();
2736        assert_eq!(err, Error::Io(INJECTED_CRASH.into()));
2737    }
2738
2739    fn inject_compact_crash(
2740        store: &FileLogStore,
2741        anchor: &RecoveryAnchor,
2742        crash_phase: CompactPhase,
2743    ) {
2744        let mut inner = store.lock().unwrap();
2745        let err = compact_prefix_with_hook(&mut inner, anchor, &mut |phase| {
2746            if phase == crash_phase {
2747                Err(Error::Io(INJECTED_CRASH.into()))
2748            } else {
2749                Ok(())
2750            }
2751        })
2752        .unwrap_err();
2753        assert_eq!(err, Error::Io(INJECTED_CRASH.into()));
2754    }
2755
2756    fn segmented_store(dir: &Path, entries: &[LogEntry]) -> FileLogStore {
2757        fs::create_dir_all(dir).unwrap();
2758        publish_closed_segment(dir, &entries[..2]).unwrap();
2759        publish_closed_segment(dir, &entries[2..4]).unwrap();
2760        if entries.len() > 4 {
2761            publish_closed_segment(dir, &entries[4..]).unwrap();
2762        }
2763        FileLogStore::open(dir, "cluster-a", 1, 1).unwrap()
2764    }
2765
2766    fn assert_reopens_with(dir: &Path, expected: &[LogEntry]) {
2767        let reopened = FileLogStore::open(dir, "cluster-a", 1, 1).unwrap();
2768        assert_eq!(
2769            reopened.read_range(IndexRange::new(1, 6).unwrap()).unwrap(),
2770            expected
2771        );
2772        assert_eq!(
2773            reopened.last_index().unwrap(),
2774            expected.last().map(|entry| entry.index)
2775        );
2776    }
2777
2778    fn assert_no_overlapping_closed_segments(dir: &Path) {
2779        let mut ranges = fs::read_dir(dir)
2780            .unwrap()
2781            .filter_map(std::result::Result::ok)
2782            .filter_map(|entry| entry.file_name().into_string().ok())
2783            .filter(|name| name.ends_with(".qlog") && !name.ends_with("-open.qlog"))
2784            .map(|name| parse_closed_segment_name(&name).unwrap())
2785            .collect::<Vec<_>>();
2786        ranges.sort_by_key(IndexRange::start);
2787        assert!(ranges
2788            .windows(2)
2789            .all(|pair| pair[0].end() < pair[1].start()));
2790    }
2791
2792    fn closed_segment_bytes(dir: &Path) -> Vec<(String, Vec<u8>)> {
2793        let mut files = fs::read_dir(dir)
2794            .unwrap()
2795            .filter_map(std::result::Result::ok)
2796            .filter_map(|entry| {
2797                entry
2798                    .file_name()
2799                    .into_string()
2800                    .ok()
2801                    .map(|name| (name, entry.path()))
2802            })
2803            .filter(|(name, _)| name.ends_with(".qlog") && !name.ends_with("-open.qlog"))
2804            .map(|(name, path)| (name, fs::read(path).unwrap()))
2805            .collect::<Vec<_>>();
2806        files.sort_by(|left, right| left.0.cmp(&right.0));
2807        files
2808    }
2809
2810    fn encode_unchecked_test_intent(old_name: &str) -> Vec<u8> {
2811        let mut out = Vec::new();
2812        out.extend_from_slice(&TRUNCATE_INTENT_MAGIC);
2813        put_u16(&mut out, TRUNCATE_INTENT_VERSION);
2814        put_u16(&mut out, 0);
2815        put_u32(&mut out, 1);
2816        put_u16(&mut out, old_name.len() as u16);
2817        out.extend_from_slice(old_name.as_bytes());
2818        let crc = crc32c(&out);
2819        put_u32(&mut out, crc);
2820        out
2821    }
2822
2823    fn chain(payloads: &[&[u8]]) -> Vec<LogEntry> {
2824        let mut entries = Vec::new();
2825        let mut prev_hash = LogHash::ZERO;
2826        for (position, payload) in payloads.iter().enumerate() {
2827            let index = position as u64 + 1;
2828            let hash = LogEntry::calculate_hash(
2829                "cluster-a",
2830                index,
2831                1,
2832                1,
2833                EntryType::Command,
2834                prev_hash,
2835                payload,
2836            );
2837            entries.push(LogEntry {
2838                cluster_id: "cluster-a".into(),
2839                epoch: 1,
2840                config_id: 1,
2841                index,
2842                entry_type: EntryType::Command,
2843                payload: payload.to_vec(),
2844                prev_hash,
2845                hash,
2846            });
2847            prev_hash = hash;
2848        }
2849        entries
2850    }
2851
2852    fn recovery_anchor(entry: &LogEntry) -> RecoveryAnchor {
2853        RecoveryAnchor::new(
2854            entry.cluster_id.clone(),
2855            entry.epoch,
2856            ConfigurationState::active(entry.config_id, LogHash::ZERO),
2857            1,
2858            LogAnchor::new(entry.index, entry.hash),
2859            SnapshotIdentity::new(
2860                format!("snapshot-{:015}", entry.index),
2861                LogHash::digest(&[b"snapshot", &entry.index.to_be_bytes()]),
2862                4096,
2863                LogHash::from_bytes([6; 32]),
2864            ),
2865        )
2866    }
2867}