Skip to main content

feoxdb/core/store/
recovery.rs

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