1use std::fs::{File, OpenOptions};
2use std::io::{Read, Write};
3use std::path::PathBuf;
4
5#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
6pub enum FileId {
7 Wal,
8 Snapshot,
9 SnapshotBak,
13 Roles,
16 CommitTimes,
22}
23
24impl FileId {
25 fn name(self) -> &'static str {
26 match self {
27 FileId::Wal => "wal.bin",
28 FileId::Snapshot => "snapshot.bin",
29 FileId::SnapshotBak => "snapshot.bin.bak",
30 FileId::Roles => "roles.json",
31 FileId::CommitTimes => "commit_times.bin",
32 }
33 }
34}
35
36pub trait Fs {
37 fn append(&mut self, file: FileId, data: &[u8]) -> std::io::Result<()>;
38 fn sync(&mut self, file: FileId) -> std::io::Result<()>;
39 fn read(&self, file: FileId) -> std::io::Result<Vec<u8>>;
40 fn write_atomic(&mut self, file: FileId, data: &[u8]) -> std::io::Result<()>;
41 fn snapshot_path(&self) -> Option<std::path::PathBuf> {
47 None
48 }
49
50 fn wal_path(&self) -> Option<std::path::PathBuf> {
55 None
56 }
57 fn read_prefix(&self, file: FileId, n: usize) -> std::io::Result<Vec<u8>> {
66 let mut bytes = self.read(file)?;
67 bytes.truncate(n);
68 Ok(bytes)
69 }
70
71 fn try_lock_exclusive(&self) -> std::io::Result<bool> {
89 Ok(true)
90 }
91
92 fn unlock(&self) -> std::io::Result<()> {
96 Ok(())
97 }
98
99 fn wal_len(&self) -> std::io::Result<u64> {
108 Ok(self.read(FileId::Wal)?.len() as u64)
109 }
110
111 fn read_range(&self, file: FileId, from: u64) -> std::io::Result<Vec<u8>> {
120 let bytes = self.read(file)?;
121 let from = from.min(bytes.len() as u64) as usize;
122 Ok(bytes[from..].to_vec())
123 }
124
125 fn snapshot_ident(&self) -> std::io::Result<Option<(u64, u64)>> {
136 let len = self.read(FileId::Snapshot)?.len() as u64;
137 Ok(if len == 0 { None } else { Some((len, 0)) })
138 }
139
140 fn list_archives(&self) -> std::io::Result<Vec<u64>> {
149 Ok(vec![])
150 }
151
152 fn read_archive(&self, _n: u64) -> std::io::Result<Vec<u8>> {
158 Ok(vec![])
159 }
160
161 fn archive_wal(&mut self, _n: u64) -> std::io::Result<()> {
171 Err(std::io::Error::other(
172 "archive_wal not supported by this Fs implementation",
173 ))
174 }
175
176 fn delete_archive(&mut self, _n: u64) -> std::io::Result<()> {
182 Ok(())
183 }
184
185 fn read_horizon_floor(&self) -> std::io::Result<u64> {
190 Ok(0)
191 }
192
193 fn write_horizon_floor(&mut self, _floor: u64) -> std::io::Result<()> {
198 Ok(())
199 }
200
201 fn has_genesis_marker(&self) -> bool {
210 false
211 }
212
213 fn write_genesis_marker(&mut self) -> std::io::Result<()> {
220 Ok(())
221 }
222
223 fn delete_genesis_marker(&mut self) -> std::io::Result<()> {
231 Ok(())
232 }
233}
234
235pub trait FsIntrospect {
236 fn total_appended(&self) -> usize;
237 fn sync_count(&self) -> usize {
238 0
239 }
240}
241
242pub const LOCK_FILE: &str = "LOCK";
248
249#[derive(Debug, Default)]
257struct LockState {
258 file: Option<File>,
259 held: bool,
260}
261
262#[derive(Debug)]
263pub struct RealFs {
264 dir: PathBuf,
265 lock: std::sync::Mutex<LockState>,
270}
271
272impl RealFs {
273 pub fn new(dir: &std::path::Path) -> std::io::Result<Self> {
274 std::fs::create_dir_all(dir)?;
275 Ok(Self {
276 dir: dir.to_path_buf(),
277 lock: std::sync::Mutex::new(LockState::default()),
278 })
279 }
280
281 pub fn dir(&self) -> &std::path::Path {
283 &self.dir
284 }
285
286 fn path(&self, file: FileId) -> PathBuf {
287 self.dir.join(file.name())
288 }
289}
290
291impl Fs for RealFs {
292 fn append(&mut self, file: FileId, data: &[u8]) -> std::io::Result<()> {
293 let mut f = OpenOptions::new()
294 .create(true)
295 .append(true)
296 .open(self.path(file))?;
297 f.write_all(data)
298 }
299
300 fn sync(&mut self, file: FileId) -> std::io::Result<()> {
301 let f = File::open(self.path(file))?;
302 full_sync(&f)
303 }
304
305 fn read(&self, file: FileId) -> std::io::Result<Vec<u8>> {
306 match File::open(self.path(file)) {
307 Ok(mut f) => {
308 let mut buf = Vec::new();
309 f.read_to_end(&mut buf)?;
310 Ok(buf)
311 }
312 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(Vec::new()),
313 Err(e) => Err(e),
314 }
315 }
316
317 fn write_atomic(&mut self, file: FileId, data: &[u8]) -> std::io::Result<()> {
318 let tmp = self.dir.join(format!("{}.tmp", file.name()));
319 {
320 let mut f = File::create(&tmp)?;
321 f.write_all(data)?;
322 full_sync(&f)?;
323 }
324 std::fs::rename(&tmp, self.path(file))?;
325 sync_dir(&self.dir)
326 }
327
328 fn snapshot_path(&self) -> Option<std::path::PathBuf> {
329 Some(self.path(FileId::Snapshot))
330 }
331
332 fn wal_path(&self) -> Option<std::path::PathBuf> {
333 Some(self.path(FileId::Wal))
334 }
335
336 fn read_prefix(&self, file: FileId, n: usize) -> std::io::Result<Vec<u8>> {
337 use std::io::Read as _;
338 match File::open(self.path(file)) {
339 Ok(mut f) => {
340 let mut buf = vec![0u8; n];
341 let read = f.read(&mut buf)?;
342 buf.truncate(read);
343 Ok(buf)
344 }
345 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(Vec::new()),
346 Err(e) => Err(e),
347 }
348 }
349
350 fn try_lock_exclusive(&self) -> std::io::Result<bool> {
351 let mut state = self.lock.lock().unwrap_or_else(|e| e.into_inner());
352 if state.held {
353 return Ok(true);
354 }
355 if state.file.is_none() {
356 state.file = Some(
357 OpenOptions::new()
358 .create(true)
359 .read(true)
360 .write(true)
361 .truncate(false)
362 .open(self.dir.join(LOCK_FILE))?,
363 );
364 }
365 let f = state.file.as_ref().expect("lock file just opened");
366 match f.try_lock() {
367 Ok(()) => {
368 state.held = true;
369 Ok(true)
370 }
371 Err(std::fs::TryLockError::WouldBlock) => Ok(false),
372 Err(std::fs::TryLockError::Error(e)) => Err(e),
373 }
374 }
375
376 fn unlock(&self) -> std::io::Result<()> {
377 let mut state = self.lock.lock().unwrap_or_else(|e| e.into_inner());
378 if !state.held {
379 return Ok(());
380 }
381 state.held = false;
384 match state.file.as_ref() {
385 Some(f) => f.unlock(),
386 None => Ok(()),
387 }
388 }
389
390 fn wal_len(&self) -> std::io::Result<u64> {
391 match std::fs::metadata(self.path(FileId::Wal)) {
392 Ok(m) => Ok(m.len()),
393 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(0),
394 Err(e) => Err(e),
395 }
396 }
397
398 fn read_range(&self, file: FileId, from: u64) -> std::io::Result<Vec<u8>> {
399 use std::io::{Read as _, Seek as _, SeekFrom};
400 let mut f = match File::open(self.path(file)) {
401 Ok(f) => f,
402 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
403 Err(e) => return Err(e),
404 };
405 let len = f.metadata()?.len();
406 if from >= len {
407 return Ok(Vec::new());
408 }
409 f.seek(SeekFrom::Start(from))?;
410 let mut buf = Vec::with_capacity((len - from) as usize);
411 f.read_to_end(&mut buf)?;
412 Ok(buf)
413 }
414
415 fn snapshot_ident(&self) -> std::io::Result<Option<(u64, u64)>> {
416 let m = match std::fs::metadata(self.path(FileId::Snapshot)) {
417 Ok(m) => m,
418 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
419 Err(e) => return Err(e),
420 };
421 let mtime_nanos = m
424 .modified()
425 .ok()
426 .and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
427 .map(|d| d.as_nanos() as u64)
428 .unwrap_or(0);
429 Ok(Some((m.len(), mtime_nanos)))
430 }
431
432 fn list_archives(&self) -> std::io::Result<Vec<u64>> {
433 let mut ns = Vec::new();
434 for entry in std::fs::read_dir(&self.dir)? {
435 let entry = entry?;
436 let name = entry.file_name();
437 let s = name.to_string_lossy();
438 if let Some(mid) = s
439 .strip_prefix("wal.")
440 .and_then(|r| r.strip_suffix(".archive"))
441 {
442 if let Ok(n) = mid.parse::<u64>() {
443 ns.push(n);
444 }
445 }
446 }
447 ns.sort_unstable();
448 Ok(ns)
449 }
450
451 fn read_archive(&self, n: u64) -> std::io::Result<Vec<u8>> {
452 let path = self.dir.join(format!("wal.{n}.archive"));
453 match std::fs::read(&path) {
454 Ok(b) => Ok(b),
455 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(vec![]),
456 Err(e) => Err(e),
457 }
458 }
459
460 fn archive_wal(&mut self, n: u64) -> std::io::Result<()> {
461 let wal_path = self.path(FileId::Wal);
462 let archive_path = self.dir.join(format!("wal.{n}.archive"));
463 std::fs::rename(&wal_path, &archive_path)?;
464 sync_dir(&self.dir)
465 }
466
467 fn delete_archive(&mut self, n: u64) -> std::io::Result<()> {
468 let path = self.dir.join(format!("wal.{n}.archive"));
469 match std::fs::remove_file(&path) {
470 Ok(()) => sync_dir(&self.dir),
471 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
472 Err(e) => Err(e),
473 }
474 }
475
476 fn read_horizon_floor(&self) -> std::io::Result<u64> {
477 let path = self.dir.join("wal.floor");
478 match std::fs::read(&path) {
479 Ok(b) if b.len() >= 8 => Ok(u64::from_le_bytes([
480 b[0], b[1], b[2], b[3], b[4], b[5], b[6], b[7],
481 ])),
482 Ok(_) => Ok(0),
483 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(0),
484 Err(e) => Err(e),
485 }
486 }
487
488 fn write_horizon_floor(&mut self, floor: u64) -> std::io::Result<()> {
489 let tmp = self.dir.join("wal.floor.tmp");
490 {
491 let mut f = File::create(&tmp)?;
492 f.write_all(&floor.to_le_bytes())?;
493 full_sync(&f)?;
494 }
495 std::fs::rename(&tmp, self.dir.join("wal.floor"))?;
496 sync_dir(&self.dir)
497 }
498
499 fn has_genesis_marker(&self) -> bool {
500 self.dir.join("wal.genesis").exists()
501 }
502
503 fn write_genesis_marker(&mut self) -> std::io::Result<()> {
504 let path = self.dir.join("wal.genesis");
505 {
506 let mut f = File::create(&path)?;
507 f.write_all(b"")?;
508 full_sync(&f)?;
509 }
510 sync_dir(&self.dir)
511 }
512
513 fn delete_genesis_marker(&mut self) -> std::io::Result<()> {
514 match std::fs::remove_file(self.dir.join("wal.genesis")) {
515 Ok(()) => sync_dir(&self.dir),
516 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
517 Err(e) => Err(e),
518 }
519 }
520}
521
522fn full_sync(file: &File) -> std::io::Result<()> {
523 #[cfg(target_os = "macos")]
524 {
525 use std::os::unix::io::AsRawFd;
526 let fd = file.as_raw_fd();
527 let rc = unsafe { libc::fcntl(fd, libc::F_FULLFSYNC) };
528 if rc == -1 {
529 return Err(std::io::Error::last_os_error());
530 }
531 Ok(())
532 }
533 #[cfg(not(target_os = "macos"))]
534 {
535 file.sync_all()
536 }
537}
538
539pub fn sync_wal_at(dir: &std::path::Path) -> std::io::Result<()> {
547 let path = dir.join(FileId::Wal.name());
548 let f = match std::fs::File::open(&path) {
549 Ok(f) => f,
550 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(()),
551 Err(e) => return Err(e),
552 };
553 full_sync(&f)
554}
555
556pub fn truncate_wal_at(dir: &std::path::Path, len: u64) -> std::io::Result<()> {
566 let path = dir.join(FileId::Wal.name());
567 let f = match OpenOptions::new().write(true).open(&path) {
568 Ok(f) => f,
569 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(()),
570 Err(e) => return Err(e),
571 };
572 f.set_len(len)?;
573 f.sync_all() }
575
576fn sync_dir(dir: &std::path::Path) -> std::io::Result<()> {
577 let d = File::open(dir)?;
578 d.sync_all()
579}
580
581#[cfg(test)]
582mod tests {
583 use super::*;
584
585 fn tmp() -> std::path::PathBuf {
586 let d = std::env::temp_dir().join(format!("graphdb-fs-{}", std::process::id()));
587 let _ = std::fs::remove_dir_all(&d);
588 d
589 }
590
591 #[test]
592 fn append_read_and_atomic_write() {
593 let mut fs = RealFs::new(&tmp()).unwrap();
594 assert_eq!(fs.read(FileId::Wal).unwrap(), Vec::<u8>::new()); fs.append(FileId::Wal, b"ab").unwrap();
596 fs.append(FileId::Wal, b"cd").unwrap();
597 fs.sync(FileId::Wal).unwrap();
598 assert_eq!(fs.read(FileId::Wal).unwrap(), b"abcd");
599 fs.write_atomic(FileId::Snapshot, b"snap1").unwrap();
600 fs.write_atomic(FileId::Snapshot, b"snap2").unwrap(); assert_eq!(fs.read(FileId::Snapshot).unwrap(), b"snap2");
602 fs.write_atomic(FileId::Wal, b"").unwrap(); assert_eq!(fs.read(FileId::Wal).unwrap(), Vec::<u8>::new());
604 }
605
606 #[test]
607 fn write_atomic_replaces_and_still_readable() {
608 let d = std::env::temp_dir().join(format!("graphdb-fs-atomic-{}", std::process::id()));
612 let _ = std::fs::remove_dir_all(&d);
613 let mut fs = RealFs::new(&d).unwrap();
614 fs.write_atomic(FileId::Snapshot, b"snap1").unwrap();
615 fs.write_atomic(FileId::Snapshot, b"snap2").unwrap();
616 assert_eq!(fs.read(FileId::Snapshot).unwrap(), b"snap2");
617 }
618}