1use 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
56pub struct ReaderSlot {
59 path: PathBuf,
60 persistent: bool,
61 _file: std::fs::File, }
63
64impl ReaderSlot {
65 pub fn register(db: &Path, gen: u64) -> Result<ReaderSlot> {
68 Self::create(db, gen, true)
69 }
70
71 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 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 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 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 let _ = crate::io::unlock(&self._file);
152 }
153}
154
155pub fn oldest_live_reader(db: &Path) -> u64 {
160 live_generations(db).map_or(0, |g| g.first().copied().unwrap_or(u64::MAX))
161}
162
163pub(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, };
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 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 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, }
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 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}