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