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 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 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}