use std::fs::{File, OpenOptions};
use std::io::Read;
use std::path::{Path, PathBuf};
use std::sync::atomic::{fence, AtomicBool, AtomicU64, Ordering};
use std::sync::Arc;
use crossbeam_epoch::Guard;
use crossbeam_utils::CachePadded;
use memmap2::{Mmap, MmapOptions};
use parking_lot::{Mutex, RwLock};
use crate::error::from_fsys;
use crate::storage::arc_cell::{ArcCell, MmapView};
use crate::storage::flush::FlushPolicy;
use crate::storage::format;
use crate::storage::meta::{self, MetaHeader};
use crate::{Error, Result};
pub(crate) const FSYS_FRAME_MAGIC: [u8; 4] = [0x46, 0x53, 0x59, 0x01];
pub(crate) const FSYS_FRAME_OVERHEAD: u64 = 12;
pub(crate) const FSYS_PRE_PAYLOAD_BYTES: u64 = 8;
const FSYS_POST_PAYLOAD_BYTES: u64 = 4;
pub(crate) const FSYS_MAX_PAYLOAD: u64 = (1 << 28) - 1;
const SERIALIZE_APPENDS: bool = cfg!(windows);
pub(crate) struct Store {
path: PathBuf,
journal: RwLock<Option<fsys::JournalHandle>>,
fs: fsys::Handle,
read_file: Mutex<File>,
mmap: ArcCell<Mmap>,
mmap_len: CachePadded<AtomicU64>,
swap_seq: CachePadded<AtomicU64>,
policy: FlushPolicy,
meta: RwLock<MetaHeader>,
meta_deferred: AtomicBool,
created: bool,
append_order: Mutex<()>,
}
impl std::fmt::Debug for Store {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Store")
.field("path", &self.path)
.field("policy", &self.policy)
.field("tail", &self.tail())
.finish()
}
}
impl Store {
pub(crate) fn open_with_policy(
path: PathBuf,
flags: u32,
policy: FlushPolicy,
iouring_sqpoll_idle_ms: Option<u32>,
) -> Result<Self> {
let mut fs_builder = fsys::builder().tune_for(fsys::Workload::Database);
if let Some(idle_ms) = iouring_sqpoll_idle_ms {
fs_builder = fs_builder.sqpoll(idle_ms);
}
let fs = fs_builder.build().map_err(from_fsys)?;
let created = prepare_data_file(&path)?;
let (meta, meta_deferred) = match meta::read(&path)? {
Some(existing) => (existing, false),
None => (MetaHeader::fresh(flags), true),
};
let read_file = OpenOptions::new()
.read(true)
.open(&path)
.map_err(Error::Io)?;
let initial_mmap = unsafe { Mmap::map(&read_file)? };
let mmap_len = initial_mmap.len() as u64;
Ok(Self {
path,
journal: RwLock::new(None),
fs,
read_file: Mutex::new(read_file),
mmap: ArcCell::new(Arc::new(initial_mmap)),
mmap_len: CachePadded::new(AtomicU64::new(mmap_len)),
swap_seq: CachePadded::new(AtomicU64::new(0)),
policy,
meta: RwLock::new(meta),
meta_deferred: AtomicBool::new(meta_deferred),
created,
append_order: Mutex::new(()),
})
}
pub(crate) fn open_journal(&mut self) -> Result<()> {
let empty = empty_mapping(self.read_file.get_mut())?;
self.mmap.set_mut(Arc::new(empty));
self.mmap_len.store(0, Ordering::Release);
let journal = self
.fs
.journal_with(&self.path, journal_options())
.map_err(from_fsys)?;
*self.journal.write() = Some(journal);
self.remap_exact()
}
pub(crate) fn finish_open(&self) -> Result<()> {
let new_entries = self.created || self.meta_deferred.load(Ordering::Acquire);
if self.meta_deferred.load(Ordering::Acquire) {
let header = *self.meta.read();
meta::write_with(&self.fs, &self.path, &header)?;
self.meta_deferred.store(false, Ordering::Release);
}
if new_entries {
sync_dir(&self.path)?;
}
Ok(())
}
pub(crate) fn path(&self) -> &Path {
&self.path
}
pub(crate) fn header(&self) -> Result<MetaHeader> {
Ok(*self.meta.read())
}
pub(crate) fn tail(&self) -> u64 {
self.journal.read().as_ref().map_or_else(
|| self.mmap_len.load(Ordering::Acquire),
|journal| journal.next_lsn().as_u64(),
)
}
pub(crate) fn fs(&self) -> &fsys::Handle {
&self.fs
}
#[inline]
pub(crate) fn mapped<'a>(&'a self, end_offset: u64, guard: &'a Guard) -> Result<MmapView<'a>> {
let view = self.mmap.load(guard);
if view.bytes().len() as u64 >= end_offset {
return Ok(view);
}
self.refresh_mmap()?;
Ok(self.mmap.load(guard))
}
pub(crate) fn payload<'a>(
&'a self,
payload_start: u64,
guard: &'a Guard,
) -> Result<Option<(&'a [u8], MmapView<'a>)>> {
let Ok(start) = usize::try_from(payload_start) else {
return Ok(None);
};
let view = self.mapped(payload_start.saturating_add(1), guard)?;
if let Ok(payload) = format::payload_at(view.bytes(), start) {
return Ok(Some((payload, view)));
}
let view = self.mapped(self.tail(), guard)?;
Ok(format::payload_at(view.bytes(), start)
.ok()
.map(|payload| (payload, view)))
}
pub(crate) fn pinned_mapping(&self) -> Result<Arc<Mmap>> {
let guard = crossbeam_epoch::pin();
Ok(self.mapped(self.tail(), &guard)?.to_arc())
}
pub(crate) fn read_begin(&self) -> u64 {
loop {
let seq = self.swap_seq.load(Ordering::Acquire);
if seq & 1 == 0 {
return seq;
}
std::thread::yield_now();
}
}
pub(crate) fn read_validate(&self, seq: u64) -> bool {
fence(Ordering::Acquire);
self.swap_seq.load(Ordering::Relaxed) == seq
}
pub(crate) fn begin_swap(&self) {
let _previous = self.swap_seq.fetch_add(1, Ordering::AcqRel);
fence(Ordering::Release);
}
pub(crate) fn end_swap(&self) {
let _previous = self.swap_seq.fetch_add(1, Ordering::Release);
}
pub(crate) fn install_file(&self, journal: fsys::JournalHandle, read_file: File, mmap: Mmap) {
let old_journal = {
let mut journal_guard = self.journal.write();
journal_guard.replace(journal)
};
{
let mut file_guard = self.read_file.lock();
let new_len = mmap.len() as u64;
self.mmap.store(Arc::new(mmap));
*file_guard = read_file;
self.mmap_len.store(new_len, Ordering::Release);
}
drop(old_journal);
}
pub(crate) fn append(&self, payload: &[u8]) -> Result<u64> {
let payload_len = payload.len() as u64;
let guard = self.journal.read();
let journal = guard
.as_ref()
.ok_or(Error::InvalidConfig(JOURNAL_NOT_OPEN))?;
let end_lsn = {
let _ordered = SERIALIZE_APPENDS.then(|| self.append_order.lock());
journal.append(payload).map_err(from_fsys)?.as_u64()
};
let payload_start = end_lsn - FSYS_POST_PAYLOAD_BYTES - payload_len;
if matches!(self.policy, FlushPolicy::WriteThrough) {
journal
.sync_through(fsys::Lsn::new(end_lsn))
.map_err(from_fsys)?;
}
Ok(payload_start)
}
pub(crate) fn append_with<F>(&self, fill_payload: F) -> Result<u64>
where
F: FnOnce(&mut Vec<u8>) -> Result<()>,
{
let mut buf = Vec::with_capacity(64);
fill_payload(&mut buf)?;
self.append(&buf)
}
pub(crate) fn append_batch<'a, I>(&self, payloads: I) -> Result<Vec<u64>>
where
I: IntoIterator<Item = &'a [u8]>,
{
let payloads: Vec<&[u8]> = payloads.into_iter().collect();
if payloads.is_empty() {
return Ok(Vec::new());
}
let guard = self.journal.read();
let journal = guard
.as_ref()
.ok_or(Error::InvalidConfig(JOURNAL_NOT_OPEN))?;
let end_lsn = {
let _ordered = SERIALIZE_APPENDS.then(|| self.append_order.lock());
journal.append_batch(&payloads).map_err(from_fsys)?.as_u64()
};
let starts = batch_payload_starts(end_lsn, &payloads);
if matches!(self.policy, FlushPolicy::WriteThrough) {
journal
.sync_through(fsys::Lsn::new(end_lsn))
.map_err(from_fsys)?;
}
Ok(starts)
}
pub(crate) fn flush(&self) -> Result<()> {
let guard = self.journal.read();
let Some(journal) = guard.as_ref() else {
return Ok(());
};
let target = journal.next_lsn();
journal.sync_through(target).map_err(from_fsys)
}
pub(crate) fn persist_meta(&self) -> Result<()> {
let header = *self.meta.read();
meta::write_with(&self.fs, &self.path, &header)
}
#[cfg(feature = "encrypt")]
pub(crate) fn set_encryption_metadata(
&self,
salt: [u8; meta::META_SALT_LEN],
verify: [u8; meta::META_VERIFY_LEN],
) -> Result<()> {
{
let mut guard = self.meta.write();
guard.encryption_salt = salt;
guard.encryption_verify = verify;
guard.flags |= meta::FLAG_ENCRYPTED;
}
self.meta_deferred.store(true, Ordering::Release);
if self.journal.read().is_some() {
self.persist_meta()?;
self.meta_deferred.store(false, Ordering::Release);
}
Ok(())
}
fn refresh_mmap(&self) -> Result<()> {
let file_guard = self.read_file.lock();
if (file_guard.metadata()?.len()) <= self.mmap_len.load(Ordering::Acquire) {
return Ok(());
}
let new_mmap = unsafe { Mmap::map(&*file_guard)? };
let new_len = new_mmap.len() as u64;
if new_len > self.mmap_len.load(Ordering::Acquire) {
self.mmap.store(Arc::new(new_mmap));
self.mmap_len.store(new_len, Ordering::Release);
}
Ok(())
}
fn remap_exact(&mut self) -> Result<()> {
let new_mmap = unsafe { Mmap::map(&*self.read_file.get_mut())? };
let new_len = new_mmap.len() as u64;
self.mmap.set_mut(Arc::new(new_mmap));
self.mmap_len.store(new_len, Ordering::Release);
Ok(())
}
pub(crate) fn open_reader(&self) -> Result<fsys::JournalReader> {
fsys::JournalReader::open(&self.path).map_err(from_fsys)
}
}
impl Drop for Store {
fn drop(&mut self) {
let _ignored = self.flush();
}
}
const JOURNAL_NOT_OPEN: &str = "journal is not open yet";
pub(crate) fn journal_options() -> fsys::JournalOptions {
fsys::JournalOptions::new().write_lifetime_hint(Some(fsys::WriteLifetimeHint::Long))
}
pub(crate) fn batch_payload_starts(end_lsn: u64, payloads: &[&[u8]]) -> Vec<u64> {
let total_frame_size: u64 = payloads
.iter()
.map(|p| FSYS_FRAME_OVERHEAD + p.len() as u64)
.sum();
let mut cursor = end_lsn - total_frame_size;
let mut starts = Vec::with_capacity(payloads.len());
for payload in payloads {
starts.push(cursor + FSYS_PRE_PAYLOAD_BYTES);
cursor += FSYS_FRAME_OVERHEAD + payload.len() as u64;
}
starts
}
fn prepare_data_file(path: &Path) -> Result<bool> {
match std::fs::metadata(path) {
Ok(meta) if meta.is_dir() => Err(Error::InvalidConfig(
"database path names a directory, not a file",
)),
Ok(meta) if meta.len() == 0 => Ok(false),
Ok(meta) => {
let mut head = [0_u8; 4];
let want = usize::try_from(meta.len().min(4)).unwrap_or(4);
let mut file = File::open(path)?;
file.read_exact(&mut head[..want])?;
let head = &head[..want];
if head[0] == 0 || head == &FSYS_FRAME_MAGIC[..want] {
Ok(false)
} else {
Err(Error::MagicMismatch)
}
}
Err(err) if err.kind() == std::io::ErrorKind::NotFound => {
drop(crate::private_fs::create_new_private_file(path)?);
Ok(true)
}
Err(err) => Err(Error::Io(err)),
}
}
fn empty_mapping(file: &File) -> Result<Mmap> {
let mmap = unsafe { MmapOptions::new().len(0).map(file)? };
Ok(mmap)
}
pub(crate) fn sync_dir(path: &Path) -> Result<()> {
let dir = match path.parent() {
Some(parent) if !parent.as_os_str().is_empty() => parent,
_ => Path::new("."),
};
sync_dir_handle(dir)
}
#[cfg(unix)]
fn sync_dir_handle(dir: &Path) -> Result<()> {
File::open(dir)?.sync_all()?;
Ok(())
}
#[cfg(windows)]
fn sync_dir_handle(dir: &Path) -> Result<()> {
use std::os::windows::fs::OpenOptionsExt;
const FILE_WRITE_DATA: u32 = 0x0000_0002;
const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
const UNSUPPORTED: [i32; 3] = [1, 50, 87];
let handle = OpenOptions::new()
.access_mode(FILE_WRITE_DATA)
.custom_flags(FILE_FLAG_BACKUP_SEMANTICS)
.open(dir)?;
match handle.sync_all() {
Err(err) if err.raw_os_error().is_some_and(|c| UNSUPPORTED.contains(&c)) => Ok(()),
other => other.map_err(Error::Io),
}
}
#[cfg(not(any(unix, windows)))]
fn sync_dir_handle(_dir: &Path) -> Result<()> {
Ok(())
}
pub(crate) fn remove_if_exists(path: &Path) -> Result<()> {
match std::fs::remove_file(path) {
Ok(()) => Ok(()),
Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(err) => Err(Error::Io(err)),
}
}
#[cfg(test)]
mod tests {
use super::*;
fn tmp_dir(label: &str) -> PathBuf {
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0_u128, |d| d.as_nanos());
let mut p = std::env::temp_dir();
p.push(format!("emdb-store-{label}-{}-{nanos}", std::process::id()));
std::fs::create_dir_all(&p).expect("mkdir");
p
}
#[test]
fn test_sync_dir_existing_directory_succeeds() {
let dir = tmp_dir("syncdir");
let file = dir.join("x");
std::fs::write(&file, b"x").expect("write");
sync_dir(&file).expect("sync dir");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn test_sync_dir_missing_directory_errors() {
let dir = tmp_dir("syncdir-missing");
let missing = dir.join("nope").join("x");
assert!(sync_dir(&missing).is_err());
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn test_prepare_data_file_foreign_bytes_rejected_untouched() {
let dir = tmp_dir("foreign");
let path = dir.join("notes.txt");
std::fs::write(&path, b"hello, not a journal").expect("write");
assert!(matches!(
prepare_data_file(&path),
Err(Error::MagicMismatch)
));
assert_eq!(std::fs::read(&path).expect("read"), b"hello, not a journal");
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn test_prepare_data_file_short_magic_prefix_accepted() {
let dir = tmp_dir("short");
let path = dir.join("db");
std::fs::write(&path, &FSYS_FRAME_MAGIC[..2]).expect("write");
assert!(!prepare_data_file(&path).expect("accepted"));
std::fs::write(&path, [0_u8; 3]).expect("write");
assert!(!prepare_data_file(&path).expect("accepted"));
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn test_prepare_data_file_missing_is_created() {
let dir = tmp_dir("create");
let path = dir.join("db");
assert!(prepare_data_file(&path).expect("created"));
assert_eq!(std::fs::metadata(&path).expect("meta").len(), 0);
assert!(!prepare_data_file(&path).expect("exists"));
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn test_batch_payload_starts_layout() {
let a: &[u8] = b"abc";
let b: &[u8] = b"";
let c: &[u8] = b"zz";
let starts = batch_payload_starts(141, &[a, b, c]);
assert_eq!(starts, vec![108, 123, 135]);
}
#[test]
fn test_remove_if_exists_missing_is_ok() {
let dir = tmp_dir("rm");
remove_if_exists(&dir.join("absent")).expect("ok");
let _ = std::fs::remove_dir_all(&dir);
}
}