Coverage Report

Created: 2026-07-31 00:13

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/home/runner/work/feoxdb/feoxdb/src/core/store/recovery.rs
Line
Count
Source
1
use std::ops::Bound;
2
use std::sync::atomic::Ordering;
3
use std::sync::Arc;
4
5
use crate::constants::*;
6
use crate::core::record::{Record, TreeSlot};
7
use crate::error::{FeoxError, Result};
8
use crate::storage::format::{get_format_ref, retirement_marker_token, RecordFormat};
9
use crate::storage::io::DiskIO;
10
use crate::storage::metadata::Metadata;
11
use crate::storage::seq_token::{crc32c, header_range, SEQ_TOKEN_MIN_VERSION};
12
13
use super::FeoxStore;
14
15
const RECOVERY_SCAN_BLOCKS: u64 = 256;
16
const RECOVERY_EXPIRED_BATCH: usize = 256;
17
18
struct RecoveryScanner<'a> {
19
    disk_io: &'a DiskIO,
20
    total_sectors: u64,
21
    buffer_start: u64,
22
    buffer: Vec<u8>,
23
}
24
25
impl<'a> RecoveryScanner<'a> {
26
98
    fn new(disk_io: &'a DiskIO, total_sectors: u64) -> Self {
27
98
        Self {
28
98
            disk_io,
29
98
            total_sectors,
30
98
            buffer_start: total_sectors,
31
98
            buffer: Vec::new(),
32
98
        }
33
98
    }
34
35
3.24M
    fn block(&mut self, sector: u64) -> Result<&[u8]> {
36
3.24M
        self.fill_at(sector)
?0
;
37
3.24M
        let offset = usize::try_from(sector - self.buffer_start)
38
3.24M
            .map_err(|_| FeoxError::InvalidDevice)
?0
39
            * FEOX_BLOCK_SIZE;
40
3.24M
        Ok(&self.buffer[offset..offset + FEOX_BLOCK_SIZE])
41
3.24M
    }
42
43
25
    fn visit_blocks(
44
25
        &mut self,
45
25
        start: u64,
46
25
        blocks: u64,
47
25
        mut visit: impl FnMut(u64, &[u8]) -> bool,
48
25
    ) -> Result<bool> {
49
25
        let end = start
50
25
            .checked_add(blocks)
51
25
            .filter(|end| *end <= self.total_sectors)
52
25
            .ok_or(FeoxError::CorruptedRecord)
?0
;
53
25
        let mut sector = start;
54
55
50
        while sector < end {
56
27
            self.fill_at(sector)
?0
;
57
27
            let buffered_blocks = (self.buffer.len() / FEOX_BLOCK_SIZE) as u64;
58
27
            let buffer_end = self.buffer_start + buffered_blocks;
59
27
            let chunk_end = end.min(buffer_end);
60
27
            let offset = usize::try_from(sector - self.buffer_start)
61
27
                .map_err(|_| FeoxError::InvalidDevice)
?0
62
                * FEOX_BLOCK_SIZE;
63
27
            let len = usize::try_from(chunk_end - sector).map_err(|_| FeoxError::InvalidDevice)
?0
64
                * FEOX_BLOCK_SIZE;
65
27
            if !visit(sector, &self.buffer[offset..offset + len]) {
66
2
                return Ok(false);
67
25
            }
68
25
            sector = chunk_end;
69
        }
70
71
23
        Ok(true)
72
25
    }
73
74
3.24M
    fn fill_at(&mut self, sector: u64) -> Result<()> {
75
3.24M
        let buffered_blocks = (self.buffer.len() / FEOX_BLOCK_SIZE) as u64;
76
3.24M
        if sector >= self.buffer_start && 
sector3.24M
< self.buffer_start + buffered_blocks {
77
3.23M
            return Ok(());
78
12.7k
        }
79
12.7k
        if sector >= self.total_sectors {
80
0
            return Err(FeoxError::CorruptedRecord);
81
12.7k
        }
82
83
12.7k
        let blocks = (self.total_sectors - sector).min(RECOVERY_SCAN_BLOCKS);
84
12.7k
        self.buffer = self.disk_io.read_sectors_sync(sector, blocks)
?0
;
85
12.7k
        self.buffer_start = sector;
86
12.7k
        Ok(())
87
3.24M
    }
88
}
89
90
impl FeoxStore {
91
149
    pub(super) fn load_indexes(&mut self) -> Result<()> {
92
149
        if self.memory_only {
93
0
            return Ok(());
94
149
        }
95
96
        // A fresh device has nothing to scan, but its signature must be on disk
97
        // before the first record is written: without it a process that crashed
98
        // before its first flush_all recovers zero records and re-marks the whole
99
        // device free over live data.
100
149
        if self.fresh_device {
101
48
            if let Some(ref disk_io) = self.disk_io {
102
48
                let mut metadata = self._metadata.write();
103
48
                disk_io.read().initialize_store_metadata(&mut metadata)
?0
;
104
0
            }
105
48
            return Ok(());
106
101
        }
107
108
101
        if let Some(ref disk_io) = self.disk_io {
109
101
            let metadata_data = disk_io.read().read_metadata()
?0
;
110
101
            if metadata_data.len() < FEOX_SIGNATURE_SIZE
111
101
                || &metadata_data[..FEOX_SIGNATURE_SIZE] != FEOX_SIGNATURE
112
            {
113
3
                return Err(FeoxError::InvalidMetadata);
114
98
            }
115
116
98
            let metadata =
117
98
                Metadata::from_bytes(&metadata_data).ok_or(FeoxError::InvalidMetadata)
?0
;
118
98
            self.format_version = metadata.version;
119
98
            *self._metadata.write() = metadata;
120
98
            self.scan_and_rebuild_indexes()
?10
;
121
0
        }
122
123
88
        Ok(())
124
149
    }
125
126
98
    pub(super) fn scan_and_rebuild_indexes(&mut self) -> Result<()> {
127
98
        if self.memory_only || self.device_size == 0 {
128
0
            return Ok(());
129
98
        }
130
131
98
        let disk_io = self.disk_io.as_ref().ok_or(FeoxError::NoDevice)
?0
;
132
133
        // Get the appropriate format handler
134
98
        let metadata_version = self.format_version;
135
98
        let format = get_format_ref(metadata_version);
136
137
98
        let total_sectors = self.device_size / FEOX_BLOCK_SIZE as u64;
138
98
        let disk = disk_io.read();
139
98
        let mut allocation_journal = disk.read_allocation_journal(total_sectors)
?0
;
140
98
        if self.read_only {
141
26
            allocation_journal.sort_unstable_by_key(|extent| extent.0);
142
72
        } else if !allocation_journal.is_empty() {
143
5
            disk.replay_allocation_journal(&allocation_journal)
?0
;
144
67
        }
145
146
98
        let mut scanner = RecoveryScanner::new(&disk, total_sectors);
147
98
        let mut sector = FEOX_DATA_START_BLOCK;
148
98
        let mut journal_index = 0;
149
98
        let mut last_end = FEOX_DATA_START_BLOCK;
150
98
        let mut retired_extents = Vec::new();
151
98
        let recovery_time = self.enable_ttl.then(|| 
self9
.
get_timestamp_pub9
());
152
98
        self.stats.disk_usage.store(0, Ordering::Relaxed);
153
154
3.24M
        'scan: while sector < total_sectors {
155
3.24M
            if self.read_only {
156
8.43k
                while let Some(&(
start1
,
sectors1
)) = allocation_journal.get(journal_index) {
157
1
                    let end = start
158
1
                        .checked_add(sectors as u64)
159
1
                        .ok_or(FeoxError::CorruptedRecord)
?0
;
160
1
                    if sector < start {
161
0
                        break;
162
1
                    }
163
1
                    journal_index += 1;
164
1
                    if sector < end {
165
1
                        sector = end;
166
1
                        continue 'scan;
167
0
                    }
168
                }
169
3.24M
            }
170
171
3.24M
            let data = scanner.block(sector)
?0
;
172
173
3.24M
            if data.len() < SECTOR_HEADER_SIZE {
174
0
                sector += 1;
175
0
                continue;
176
3.24M
            }
177
178
            // Check for deletion marker first
179
3.24M
            if data.len() >= DELETION_MARKER.len() && &data[..8] == DELETION_MARKER {
180
34
                if metadata_version < SEQ_TOKEN_MIN_VERSION
181
32.7k
                    && 
data[8..].iter()13
.
all13
(|byte| *byte == 0)
182
                {
183
8
                    if !self.allow_ambiguous_legacy_recovery {
184
4
                        return Err(FeoxError::AmbiguousLegacyTombstone);
185
4
                    }
186
4
                    self.ambiguous_legacy_markers += 1;
187
4
                    sector += 1;
188
4
                    continue;
189
26
                }
190
191
26
                if data.len() < DELETION_MARKER_SIZE {
192
0
                    return Err(FeoxError::CorruptedRecord);
193
26
                }
194
195
26
                let expected = retirement_marker_token(sector, data);
196
26
                let found = u16::from_le_bytes([data[16], data[17]]);
197
26
                if expected != found {
198
2
                    return Err(FeoxError::CorruptedRecord);
199
24
                }
200
201
24
                let extent = u64::from_le_bytes(data[8..16].try_into().unwrap());
202
24
                let Some(extent_end) = sector.checked_add(extent) else {
203
0
                    return Err(FeoxError::CorruptedRecord);
204
                };
205
24
                if extent == 0 || 
extent_end > total_sectors23
{
206
1
                    return Err(FeoxError::CorruptedRecord);
207
23
                }
208
23
                let mut needs_repair = data[18] != RETIREMENT_COMPLETE;
209
23
                if !needs_repair && 
extent > 121
{
210
                    needs_repair =
211
16
                        !
scanner14
.
visit_blocks14
(
sector + 114
,
extent - 114
, |chunk_sector, tails| {
212
16
                            tails
213
16
                                .chunks_exact(FEOX_BLOCK_SIZE)
214
16
                                .enumerate()
215
599
                                .
all16
(|(index, tail)| {
216
599
                                    let tail_sector = chunk_sector + index as u64;
217
599
                                    is_complete_retirement_block(
218
599
                                        tail,
219
599
                                        tail_sector,
220
599
                                        extent - (tail_sector - sector),
221
                                    )
222
599
                                })
223
16
                        })
?0
;
224
9
                }
225
23
                if needs_repair && 
!self.read_only4
{
226
3
                    retired_extents.push((sector, extent as usize));
227
20
                }
228
23
                sector = extent_end;
229
23
                continue;
230
3.24M
            }
231
232
3.24M
            let marker = u16::from_le_bytes([data[0], data[1]]);
233
234
3.24M
            if marker != SECTOR_MARKER {
235
3.24M
                sector += 1;
236
3.24M
                continue;
237
584
            }
238
239
584
            if data.len() < SECTOR_HEADER_SIZE + 2 {
240
0
                if metadata_version >= SEQ_TOKEN_MIN_VERSION {
241
0
                    return Err(FeoxError::CorruptedRecord);
242
0
                }
243
0
                sector += 1;
244
0
                continue;
245
584
            }
246
247
584
            if header_range(format, data).is_none() {
248
1
                if metadata_version >= SEQ_TOKEN_MIN_VERSION {
249
1
                    return Err(FeoxError::CorruptedRecord);
250
0
                }
251
0
                sector += 1;
252
0
                continue;
253
583
            }
254
255
583
            let seq_num = u16::from_le_bytes([data[2], data[3]]);
256
257
583
            if (metadata_version < SEQ_TOKEN_MIN_VERSION && 
seq_num != 0332
)
258
583
                || (metadata_version >= SEQ_TOKEN_MIN_VERSION && 
seq_num == 0251
)
259
            {
260
0
                return Err(FeoxError::CorruptedRecord);
261
583
            }
262
263
            // Parse the record using format trait
264
583
            let (key, value_len, timestamp, ttl_expiry) = match format.parse_record(data) {
265
583
                Some(parsed) => parsed,
266
                None => {
267
0
                    if metadata_version >= SEQ_TOKEN_MIN_VERSION {
268
0
                        return Err(FeoxError::CorruptedRecord);
269
0
                    }
270
0
                    sector += 1;
271
0
                    continue;
272
                }
273
            };
274
275
583
            if key.len() > MAX_KEY_SIZE || value_len == 0 || value_len > MAX_VALUE_SIZE {
276
0
                if metadata_version >= SEQ_TOKEN_MIN_VERSION {
277
0
                    return Err(FeoxError::CorruptedRecord);
278
0
                }
279
0
                sector += 1;
280
0
                continue;
281
583
            }
282
283
            // Calculate total size using format trait
284
583
            let total_size = format.total_size(key.len(), value_len);
285
583
            let sectors_needed = total_size.div_ceil(FEOX_BLOCK_SIZE);
286
583
            let head_crc =
287
583
                (metadata_version >= SEQ_TOKEN_MIN_VERSION).then(|| 
record_crc_head251
(
sector251
,
data251
));
288
289
            // A record whose extent leaves the device is rejected, never propagated:
290
            // one such header used to make every subsequent open fail.
291
583
            let extent_in_bounds = sector
292
583
                .checked_add(sectors_needed as u64)
293
583
                .is_some_and(|end| end <= total_sectors);
294
583
            if sectors_needed == 0 || !extent_in_bounds {
295
0
                if metadata_version >= SEQ_TOKEN_MIN_VERSION {
296
0
                    return Err(FeoxError::CorruptedRecord);
297
0
                }
298
0
                sector += 1;
299
0
                continue;
300
583
            }
301
583
            let extent_end = sector + sectors_needed as u64;
302
583
            if self.read_only && 
journal_overlaps24
(
&allocation_journal24
,
journal_index24
,
extent_end24
) {
303
0
                return Err(FeoxError::CorruptedRecord);
304
583
            }
305
306
583
            if metadata_version >= SEQ_TOKEN_MIN_VERSION {
307
251
                let mut crc = head_crc.unwrap();
308
251
                if sectors_needed > 1 {
309
11
                    scanner.visit_blocks(sector + 1, sectors_needed as u64 - 1, |_, tail| {
310
11
                        crc = crc32c(crc, tail);
311
11
                        true
312
11
                    })
?0
;
313
240
                }
314
251
                let actual = record_token(crc);
315
251
                if seq_num != actual {
316
2
                    return Err(FeoxError::CorruptedRecord);
317
249
                }
318
332
            }
319
320
581
            self.version_clock.observe(&key, timestamp);
321
581
            let mut record = Record::new(key.clone(), Vec::new(), timestamp);
322
581
            record.sector.store(sector, Ordering::Release);
323
581
            record.value_len = value_len;
324
581
            record.ttl_expiry.store(ttl_expiry, Ordering::Release);
325
581
            record.clear_value();
326
327
581
            let record_arc = Arc::new(record);
328
581
            let key_len = key.len();
329
330
581
            let existing = self.hash_table.read(&key, |_, record| 
Arc::clone7
(
record7
));
331
581
            if existing
332
581
                .as_ref()
333
581
                .is_some_and(|record| 
record.timestamp7
>
timestamp7
)
334
            {
335
2
                if !self.read_only {
336
1
                    retired_extents.push((sector, sectors_needed));
337
1
                }
338
2
                sector += sectors_needed as u64;
339
2
                continue;
340
579
            }
341
342
579
            if let Some(
existing5
) = existing {
343
5
                let existing_sectors = format
344
5
                    .total_size(existing.key.len(), existing.value_len)
345
5
                    .div_ceil(FEOX_BLOCK_SIZE);
346
5
                let existing_sector = existing.sector.load(Ordering::Acquire);
347
5
                self.free_space
348
5
                    .write()
349
5
                    .release_sectors(existing_sector, existing_sectors as u64)
?0
;
350
5
                self.stats.memory_usage.fetch_sub(
351
5
                    self.calculate_record_size(existing.key.len(), existing.value_len),
352
5
                    Ordering::Relaxed,
353
                );
354
5
                self.stats.disk_usage.fetch_sub(
355
5
                    (existing_sectors * FEOX_BLOCK_SIZE) as u64,
356
5
                    Ordering::Relaxed,
357
                );
358
5
                self.note_ttl_transition(existing.ttl_expiry.load(Ordering::Acquire), 0);
359
5
                if !self.read_only {
360
3
                    retired_extents.push((existing_sector, existing_sectors));
361
3
                
}2
362
574
            } else {
363
574
                self.stats.record_count.fetch_add(1, Ordering::Relaxed);
364
574
            }
365
366
579
            if sector > last_end {
367
14
                self.free_space
368
14
                    .write()
369
14
                    .release_sectors(last_end, sector - last_end)
?0
;
370
565
            }
371
579
            last_end = sector + sectors_needed as u64;
372
373
579
            self.hash_table.upsert(key.clone(), Arc::clone(&record_arc));
374
579
            self.tree
375
579
                .insert(key, TreeSlot::new(Arc::clone(&record_arc)));
376
377
579
            let record_size = self.calculate_record_size(key_len, value_len);
378
579
            self.stats
379
579
                .memory_usage
380
579
                .fetch_add(record_size, Ordering::Relaxed);
381
579
            self.note_ttl_transition(0, ttl_expiry);
382
383
            // Track disk usage
384
579
            self.stats
385
579
                .disk_usage
386
579
                .fetch_add((sectors_needed * FEOX_BLOCK_SIZE) as u64, Ordering::Relaxed);
387
388
579
            sector += sectors_needed as u64;
389
        }
390
391
88
        if let Some(
now9
) = recovery_time {
392
9
            self.remove_expired_recovery_winners(now, format, &mut retired_extents)
?0
;
393
79
        }
394
395
88
        if !self.read_only {
396
64
            disk.retire_extents(&retired_extents)
?0
;
397
24
        }
398
399
88
        if last_end < total_sectors {
400
82
            self.free_space
401
82
                .write()
402
82
                .release_sectors(last_end, total_sectors - last_end)
?0
;
403
6
        }
404
405
88
        Ok(())
406
98
    }
407
408
9
    fn remove_expired_recovery_winners(
409
9
        &self,
410
9
        now: u64,
411
9
        format: &dyn RecordFormat,
412
9
        retired_extents: &mut Vec<(u64, usize)>,
413
9
    ) -> Result<()> {
414
9
        let mut after = None;
415
416
        loop {
417
18
            let (last, expired) = {
418
18
                let guard = &crossbeam_epoch::pin();
419
18
                let mut cursor = match after.as_deref() {
420
9
                    Some(key) => self.tree.lower_bound(Bound::Excluded(key)),
421
9
                    None => self.tree.front(),
422
                };
423
18
                let mut last = None;
424
18
                let mut expired = Vec::new();
425
18
                let mut visited = 0;
426
427
325
                while visited < RECOVERY_EXPIRED_BATCH {
428
324
                    let Some(
entry307
) = cursor else {
429
17
                        break;
430
                    };
431
307
                    let record = Arc::clone(entry.value().load(guard));
432
307
                    last = Some(entry.key().clone());
433
307
                    cursor = entry.next();
434
307
                    let expiry = record.ttl_expiry.load(Ordering::Acquire);
435
307
                    if expiry > 0 && 
now > expiry305
{
436
303
                        expired.push((record.key.clone(), record));
437
303
                    
}4
438
307
                    visited += 1;
439
                }
440
441
18
                (last, expired)
442
            };
443
444
18
            let Some(
last9
) = last else {
445
9
                break;
446
            };
447
9
            after = Some(last);
448
449
303
            for (key, record) in 
expired9
{
450
303
                let removed = match self.hash_table.entry(key.clone()) {
451
303
                    scc::hash_map::Entry::Occupied(entry) if Arc::ptr_eq(entry.get(), &record) => {
452
303
                        self.tree.remove(&key);
453
303
                        let _ = entry.remove();
454
303
                        true
455
                    }
456
0
                    _ => false,
457
                };
458
303
                if !removed {
459
0
                    continue;
460
303
                }
461
462
303
                let sectors = format
463
303
                    .total_size(record.key.len(), record.value_len)
464
303
                    .div_ceil(FEOX_BLOCK_SIZE);
465
303
                let sector = record.sector.load(Ordering::Acquire);
466
303
                record.refcount.store(0, Ordering::Release);
467
303
                self.free_space
468
303
                    .write()
469
303
                    .release_sectors(sector, sectors as u64)
?0
;
470
303
                self.stats.record_count.fetch_sub(1, Ordering::Relaxed);
471
303
                self.stats.memory_usage.fetch_sub(
472
303
                    self.calculate_record_size(record.key.len(), record.value_len),
473
303
                    Ordering::Relaxed,
474
                );
475
303
                self.stats
476
303
                    .disk_usage
477
303
                    .fetch_sub((sectors * FEOX_BLOCK_SIZE) as u64, Ordering::Relaxed);
478
303
                self.note_ttl_transition(record.ttl_expiry.load(Ordering::Acquire), 0);
479
303
                if !self.read_only {
480
303
                    retired_extents.push((sector, sectors));
481
303
                
}0
482
            }
483
        }
484
485
9
        Ok(())
486
9
    }
487
}
488
489
24
fn journal_overlaps(journal: &[(u64, usize)], index: usize, extent_end: u64) -> bool {
490
24
    journal
491
24
        .get(index)
492
24
        .is_some_and(|(start, _)| 
*start0
<
extent_end0
)
493
24
}
494
495
251
fn record_crc_head(sector: u64, head: &[u8]) -> u32 {
496
251
    let mut crc = crc32c(0, &sector.to_le_bytes());
497
251
    crc = crc32c(crc, &head[..2]);
498
251
    crc = crc32c(crc, &[0, 0]);
499
251
    crc32c(crc, &head[SECTOR_HEADER_SIZE..])
500
251
}
501
502
251
fn record_token(crc: u32) -> u16 {
503
251
    match ((crc >> 16) ^ (crc & 0xFFFF)) as u16 {
504
0
        0 => 1,
505
251
        token => token,
506
    }
507
251
}
508
509
599
fn is_complete_retirement_block(data: &[u8], sector: u64, remaining: u64) -> bool {
510
599
    data.len() >= DELETION_MARKER_SIZE
511
599
        && &data[..8] == DELETION_MARKER
512
597
        && u64::from_le_bytes(data[8..16].try_into().unwrap()) == remaining
513
597
        && data[18] == RETIREMENT_COMPLETE
514
597
        && u16::from_le_bytes([data[16], data[17]]) == retirement_marker_token(sector, data)
515
599
}