Skip to main content

kernel/
readers.rs

1//! 2n: the reader table -- how the writer learns the oldest generation any
2//! live snapshot reader still needs, so page recycling never pulls a page
3//! out from under one.
4//!
5//! LMDB's shape, flock-based so it cannot go stale: each snapshot reader
6//! creates `readers/r-<pid>-<n>` holding its pinned generation and KEEPS an
7//! exclusive flock on it for the snapshot's lifetime. Close or crash, the
8//! kernel drops the lock with the fd. The writer, when it wants the oldest,
9//! walks the directory: a file whose lock it can grab is a dead reader's
10//! leavings (unlinked on the spot); a file whose lock is busy is a live
11//! reader. Its generation is stored in two identical framed copies, each with
12//! magic, version and CRC. Any malformed copy or disagreement means generation
13//! zero: ambiguity leaks space but cannot recycle a live snapshot's pages.
14//!
15//! Failure posture (the 2n risk ladder): an unreadable directory or file
16//! means "assume a reader at generation 0" -- recycling stops, the file
17//! grows like it always did, nothing can be wrongly reused.
18
19use std::io::{Read, Seek, Write};
20use std::path::{Path, PathBuf};
21
22use crate::Result;
23
24const SLOT_MAGIC: [u8; 8] = *b"SEKREAD\0";
25const SLOT_VERSION: u16 = 1;
26const FRAME_LEN: usize = 24;
27const SLOT_LEN: usize = FRAME_LEN * 2;
28
29fn encode_frame(generation: u64) -> [u8; FRAME_LEN] {
30    let mut frame = [0u8; FRAME_LEN];
31    frame[0..8].copy_from_slice(&SLOT_MAGIC);
32    frame[8..10].copy_from_slice(&SLOT_VERSION.to_le_bytes());
33    frame[12..20].copy_from_slice(&generation.to_le_bytes());
34    let crc = crc32c::crc32c(&frame[..20]);
35    frame[20..24].copy_from_slice(&crc.to_le_bytes());
36    frame
37}
38
39fn decode_slot(bytes: &[u8]) -> Option<u64> {
40    if bytes.len() != SLOT_LEN || bytes[..FRAME_LEN] != bytes[FRAME_LEN..] {
41        return None;
42    }
43    let frame = &bytes[..FRAME_LEN];
44    if frame[..8] != SLOT_MAGIC
45        || u16::from_le_bytes(frame[8..10].try_into().ok()?) != SLOT_VERSION
46        || frame[10..12] != [0, 0]
47        || u32::from_le_bytes(frame[20..24].try_into().ok()?) != crc32c::crc32c(&frame[..20])
48    {
49        return None;
50    }
51    Some(u64::from_le_bytes(frame[12..20].try_into().ok()?))
52}
53
54fn readers_dir(db: &Path) -> PathBuf { db.join("readers") }
55
56/// Held by a snapshot reader for its lifetime; the registration disappears
57/// (unlink + lock release) on drop, and the LOCK disappears even on crash.
58pub struct ReaderSlot {
59    path: PathBuf,
60    persistent: bool,
61    _file: std::fs::File, // holds the flock
62}
63
64impl ReaderSlot {
65    /// Register a reader pinned at `gen`. Failure to register is returned as
66    /// an error: an unregistered reader is exactly the unsafe state.
67    pub fn register(db: &Path, gen: u64) -> Result<ReaderSlot> {
68        Self::create(db, gen, true)
69    }
70
71    /// Install the conservative live lock before snapshot metadata is read.
72    /// The bytes need not be directory-durable: while this process lives,
73    /// completed writes are visible to the writer; if it crashes, the lock
74    /// disappears and the stale file is swept. The final generation update
75    /// below supplies the one existing fsync, avoiding another snapshot-open
76    /// barrier solely to close the race.
77    pub(crate) fn reserve(db: &Path) -> Result<ReaderSlot> {
78        Self::create(db, 0, false)
79    }
80
81    fn create(db: &Path, gen: u64, sync: bool) -> Result<ReaderSlot> {
82        if let Some(limits) = crate::limits::read(db)? {
83            return Self::create_fixed(db, gen, limits.readers);
84        }
85        Self::create_with_after_file(db, gen, sync, || {})
86    }
87
88    fn create_with_after_file(db: &Path, gen: u64, sync: bool, after_file: impl FnOnce()) -> Result<ReaderSlot> {
89        let dir = readers_dir(db);
90        std::fs::create_dir_all(&dir)?;
91        static N: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
92        // Publish only an initialized, already-locked inode. A visible unlocked
93        // slot can be swept by a writer before its creator obtains the lock.
94        let mut candidate = tempfile::Builder::new().prefix(".reader-pending-").tempfile_in(db)?;
95        after_file();
96        if !crate::io::try_lock_exclusive(candidate.as_file())? {
97            return Err(std::io::Error::new(std::io::ErrorKind::WouldBlock,
98                "reader slot file unexpectedly locked").into());
99        }
100        let frame = encode_frame(gen);
101        candidate.write_all(&frame)?;
102        candidate.write_all(&frame)?;
103        if sync { candidate.as_file().sync_all()?; }
104        loop {
105            let n = N.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
106            let path = dir.join(format!("r-{}-{}", std::process::id(), n));
107            match candidate.persist_noclobber(&path) {
108                Ok(file) => return Ok(ReaderSlot { path, _file: file, persistent: false }),
109                Err(e) if e.error.kind() == std::io::ErrorKind::AlreadyExists => candidate = e.file,
110                Err(e) => return Err(e.error.into()),
111            }
112        }
113    }
114
115    fn create_fixed(db: &Path, generation: u64, count: u32) -> Result<Self> {
116        let dir = readers_dir(db);
117        std::fs::create_dir_all(&dir)?;
118        // Never unlink fixed slots: all contenders must lock the same inode.
119        // A writer ignores unlocked slots; a locked, incomplete frame pins 0.
120        for n in 0..count {
121            let path = dir.join(format!("fixed-{n}"));
122            let file = std::fs::OpenOptions::new().read(true).write(true).create(true).truncate(false).open(&path)?;
123            if !crate::io::try_lock_exclusive(&file)? { continue; }
124            let mut slot = Self { path, _file: file, persistent: true };
125            slot.set_generation(generation)?;
126            return Ok(slot);
127        }
128        Err(crate::Error::ResourceLimit("snapshot reader slots full"))
129    }
130
131    /// Replace a conservative generation-zero reservation with the generation
132    /// selected from metadata. While the two framed copies are being changed,
133    /// a writer sees either 0, the final generation, or a mismatch that also
134    /// decodes as 0. Every intermediate state therefore stops recycling.
135    pub fn set_generation(&mut self, generation: u64) -> Result<()> {
136        self._file.rewind()?;
137        let frame = encode_frame(generation);
138        self._file.write_all(&frame)?;
139        self._file.write_all(&frame)?;
140        self._file.set_len(SLOT_LEN as u64)?;
141        self._file.sync_all()?;
142        Ok(())
143    }
144}
145
146impl Drop for ReaderSlot {
147    fn drop(&mut self) {
148        if !self.persistent { let _ = std::fs::remove_file(&self.path); }
149        // Released by unlock, not by the close that follows: a child process
150        // may hold a copy of this descriptor for an instant (`io::Locked`).
151        let _ = crate::io::unlock(&self._file);
152    }
153}
154
155/// The oldest generation any live reader pins, or `u64::MAX` when none.
156/// Dead readers' files are swept as they are met. Any ambiguity -- a file
157/// that cannot be opened or parsed while its lock is busy -- reports
158/// generation 0: recycling halts rather than guesses.
159pub fn oldest_live_reader(db: &Path) -> u64 {
160    live_generations(db).map_or(0, |g| g.first().copied().unwrap_or(u64::MAX))
161}
162
163/// A complete sorted set is needed to distinguish versions actually visible
164/// to readers from intermediate versions. Any ambiguity disables promotion.
165pub(crate) fn live_generations(db: &Path) -> Option<Vec<u64>> {
166    let dir = readers_dir(db);
167    let entries = match std::fs::read_dir(&dir) {
168        Ok(e) => e,
169        Err(ref e) if e.kind() == std::io::ErrorKind::NotFound => return Some(Vec::new()),
170        Err(_) => return None, // unreadable table: assume the oldest possible reader
171    };
172    let mut generations = Vec::new();
173    for entry in entries {
174        let Ok(entry) = entry else { return None };
175        let path = entry.path();
176        let f = match std::fs::File::open(&path) {
177            Ok(f) => f,
178            Err(e) if e.kind() == std::io::ErrorKind::NotFound => continue,
179            Err(_) => return None,
180        };
181        match crate::io::try_lock_exclusive(&f) {
182            Ok(true) => {
183                // lock acquired: the registering process is gone -- stale file
184                if !entry.file_name().to_string_lossy().starts_with("fixed-") {
185                    let _ = std::fs::remove_file(&path);
186                }
187                let _ = crate::io::unlock(&f);
188                continue;
189            }
190            Ok(false) => {
191                // busy lock: live reader -- read its pinned generation
192                let mut f2 = f;
193                let mut buf = [0u8; SLOT_LEN];
194                if f2.read_exact(&mut buf).is_err() { return None; }
195                let mut extra = [0u8; 1];
196                match f2.read(&mut extra) {
197                    Ok(0) => {}
198                    Ok(_) | Err(_) => return None,
199                }
200                let Some(generation) = decode_slot(&buf) else { return None };
201                if generation == 0 { return None; }
202                if generations.len() >= 4096 { return None; }
203                generations.push(generation);
204            }
205            Err(_) => return None, // ambiguous lock state: halt recycling
206        }
207    }
208    generations.sort_unstable();
209    generations.dedup();
210    Some(generations)
211}
212
213#[cfg(test)]
214mod tests {
215    use super::*;
216
217    #[test]
218    fn a_writer_sweep_during_creation_cannot_hide_a_live_reader() {
219        let d = tempfile::TempDir::new().unwrap();
220        let _reader = ReaderSlot::create_with_after_file(d.path(), 41, true, || {
221            assert_eq!(oldest_live_reader(d.path()), u64::MAX);
222        }).unwrap();
223        assert_eq!(oldest_live_reader(d.path()), 41,
224            "registration returned success but writer cannot see the reader");
225    }
226
227    #[test]
228    fn a_dropped_slot_disappears_and_a_live_one_pins() {
229        let d = tempfile::TempDir::new().unwrap();
230        assert_eq!(oldest_live_reader(d.path()), u64::MAX, "empty table: no pin");
231        let s1 = ReaderSlot::register(d.path(), 41).unwrap();
232        let _s2 = ReaderSlot::register(d.path(), 44).unwrap();
233        assert_eq!(oldest_live_reader(d.path()), 41);
234        drop(s1);
235        assert_eq!(oldest_live_reader(d.path()), 44);
236    }
237
238    #[test]
239    fn a_crashed_readers_file_is_swept_not_trusted() {
240        // simulate the crash leavings: a registration file with NO live lock
241        let d = tempfile::TempDir::new().unwrap();
242        let dir = readers_dir(d.path());
243        std::fs::create_dir_all(&dir).unwrap();
244        std::fs::write(dir.join("r-99999-0"), 7u64.to_le_bytes()).unwrap();
245        assert_eq!(oldest_live_reader(d.path()), u64::MAX,
246                   "an unlocked file is a dead reader, not a pin");
247        assert!(!dir.join("r-99999-0").exists(), "and it is swept");
248    }
249
250    #[test]
251    fn a_flipped_byte_in_a_live_slot_pins_generation_zero() {
252        let d = tempfile::TempDir::new().unwrap();
253        let _slot = ReaderSlot::register(d.path(), 41).unwrap();
254        let dir = readers_dir(d.path());
255        let path = std::fs::read_dir(&dir).unwrap().next().unwrap().unwrap().path();
256        let mut bytes = std::fs::read(&path).unwrap();
257        let last = bytes.len() - 1;
258        bytes[last] ^= 0x80;
259        std::fs::write(&path, bytes).unwrap();
260
261        assert_eq!(
262            oldest_live_reader(d.path()),
263            0,
264            "an unverifiable live reader must stop recycling, never advance its generation",
265        );
266    }
267}