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