Packages

Fast, safe, and ergonomic LMDB storage for Elixir, powered by Rust

Current section

Files

Jump to
lean_lmdb native lean_lmdb_nif src lifecycle.rs
Raw

native/lean_lmdb_nif/src/lifecycle.rs

use crate::{InternalResult, atoms};
use heed::types::Bytes;
use heed::{Database, Env, EnvFlags, EnvOpenOptions, WithoutTls};
use rustler::{Resource, ResourceArc, Term};
use std::collections::HashMap;
use std::fs;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, LazyLock, Mutex, MutexGuard, Weak};
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum LifecycleError {
InvalidPath,
PathNotFound,
InvalidEnvironmentOptions,
InvalidEnvironment,
IncompatibleEnvironmentOptions,
EnvironmentOpenFailed,
ResourceIdExhausted,
InvalidDatabaseName,
ReadOnly,
TransactionFailed,
DatabaseCreateFailed,
DatabaseOpenFailed,
DatabaseNotFound,
}
impl LifecycleError {
pub(crate) fn atom(self) -> rustler::Atom {
match self {
Self::InvalidPath => atoms::invalid_path(),
Self::PathNotFound => atoms::path_not_found(),
Self::InvalidEnvironmentOptions => atoms::invalid_environment_options(),
Self::InvalidEnvironment => atoms::invalid_environment(),
Self::IncompatibleEnvironmentOptions => atoms::incompatible_environment_options(),
Self::EnvironmentOpenFailed => atoms::environment_open_failed(),
Self::ResourceIdExhausted => atoms::resource_id_exhausted(),
Self::InvalidDatabaseName => atoms::invalid_database_name(),
Self::ReadOnly => atoms::read_only(),
Self::TransactionFailed => atoms::transaction_failed(),
Self::DatabaseCreateFailed => atoms::database_create_failed(),
Self::DatabaseOpenFailed => atoms::database_open_failed(),
Self::DatabaseNotFound => atoms::database_not_found(),
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum Durability {
Sync,
NoMetaSync,
NoSync,
}
impl Durability {
pub(crate) fn from_code(code: u8) -> InternalResult<Self> {
match code {
0 => Ok(Self::Sync),
1 => Ok(Self::NoMetaSync),
2 => Ok(Self::NoSync),
_ => Err(LifecycleError::InvalidEnvironmentOptions),
}
}
pub(crate) fn flag(self) -> EnvFlags {
match self {
Self::Sync => EnvFlags::empty(),
Self::NoMetaSync => EnvFlags::NO_META_SYNC,
Self::NoSync => EnvFlags::NO_SYNC,
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) struct EnvironmentOptions {
pub(crate) map_size: usize,
pub(crate) max_dbs: u32,
pub(crate) max_readers: u32,
pub(crate) read_only: bool,
pub(crate) durability: Durability,
pub(crate) fixed_map: bool,
pub(crate) read_ahead: bool,
}
pub(crate) type BinaryDatabase = Database<Bytes, Bytes>;
pub(crate) struct EnvironmentState {
pub(crate) id: u64,
pub(crate) path: PathBuf,
pub(crate) options: EnvironmentOptions,
pub(crate) env: Env<WithoutTls>,
pub(crate) databases: Mutex<HashMap<String, BinaryDatabase>>,
}
pub(crate) struct EnvironmentResource {
pub(crate) state: Arc<EnvironmentState>,
}
#[rustler::resource_impl]
impl Resource for EnvironmentResource {}
pub(crate) struct DatabaseResource {
pub(crate) state: Arc<EnvironmentState>,
pub(crate) database: BinaryDatabase,
pub(crate) name: String,
}
#[rustler::resource_impl]
impl Resource for DatabaseResource {}
pub(crate) static REGISTRY: LazyLock<Mutex<HashMap<PathBuf, Weak<EnvironmentState>>>> =
LazyLock::new(|| Mutex::new(HashMap::new()));
pub(crate) static NEXT_RESOURCE_ID: AtomicU64 = AtomicU64::new(1);
pub(crate) fn registry() -> MutexGuard<'static, HashMap<PathBuf, Weak<EnvironmentState>>> {
// A panic in unrelated native code must not permanently disable lifecycle calls.
match REGISTRY.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
}
}
pub(crate) const MIN_MAP_SIZE: u64 = 1024 * 1024;
pub(crate) const MAX_MAP_SIZE: u64 = 1024 * 1024 * 1024 * 1024;
#[allow(clippy::too_many_arguments)]
pub(crate) fn validate_options(
map_size: u64,
max_dbs: u32,
max_readers: u32,
read_only: bool,
create: bool,
durability_code: u8,
fixed_map: bool,
read_ahead: bool,
) -> InternalResult<EnvironmentOptions> {
let durability = Durability::from_code(durability_code)?;
if !(MIN_MAP_SIZE..=MAX_MAP_SIZE).contains(&map_size)
|| !map_size.is_multiple_of(page_size::get() as u64)
|| !(1..=1024).contains(&max_dbs)
|| !(1..=4096).contains(&max_readers)
|| (read_only && create)
|| (read_only && durability != Durability::Sync)
{
return Err(LifecycleError::InvalidEnvironmentOptions);
}
let map_size =
usize::try_from(map_size).map_err(|_| LifecycleError::InvalidEnvironmentOptions)?;
Ok(EnvironmentOptions {
map_size,
max_dbs,
max_readers,
read_only,
durability,
fixed_map,
read_ahead,
})
}
pub(crate) fn canonical_environment_path(path: &str, create: bool) -> InternalResult<PathBuf> {
if path.is_empty() {
return Err(LifecycleError::InvalidPath);
}
let requested = Path::new(path);
if create {
fs::create_dir_all(requested).map_err(|_| LifecycleError::InvalidPath)?;
} else if !requested.is_dir() {
return Err(LifecycleError::PathNotFound);
}
let canonical = fs::canonicalize(requested).map_err(|_| LifecycleError::InvalidPath)?;
if canonical.to_str().is_none() {
return Err(LifecycleError::InvalidPath);
}
Ok(canonical)
}
pub(crate) fn next_resource_id() -> InternalResult<u64> {
NEXT_RESOURCE_ID
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |id| {
(id != u64::MAX).then_some(id + 1)
})
.map_err(|_| LifecycleError::ResourceIdExhausted)
}
pub(crate) fn open_shared_environment(
path: PathBuf,
options: EnvironmentOptions,
create: bool,
) -> InternalResult<Arc<EnvironmentState>> {
// Keep the lock across the native open. This is intentionally conservative: it
// makes lookup-and-create atomic and prevents duplicate opens for all aliases.
let mut entries = registry();
entries.retain(|_, state| state.strong_count() != 0);
// LMDB may create files during open, so enforce create=false ourselves.
if !create && !path.join("data.mdb").is_file() {
return Err(LifecycleError::PathNotFound);
}
if let Some(state) = entries.get(&path).and_then(Weak::upgrade) {
return if state.options == options {
Ok(state)
} else {
Err(LifecycleError::IncompatibleEnvironmentOptions)
};
}
// Weak::upgrade can fail while the final Arc is still running its destructor.
// heed keeps a closing event until its own process-wide entry is gone; waiting
// here prevents a transient EnvAlreadyOpened during immediate drop-and-reopen.
if let Some(closing) = heed::env_closing_event(&path) {
closing.wait();
}
let mut builder = EnvOpenOptions::new().read_txn_without_tls();
builder
.map_size(options.map_size)
.max_dbs(options.max_dbs)
.max_readers(options.max_readers);
let mut flags = options.durability.flag();
if options.read_only {
flags |= EnvFlags::READ_ONLY;
}
if options.fixed_map {
flags |= EnvFlags::FIXED_MAP;
}
if !options.read_ahead {
flags |= EnvFlags::NO_READ_AHEAD;
}
if !flags.is_empty() {
// SAFETY: only LMDB's documented synchronization modes, READ_ONLY,
// experimental FIXED_MAP, and the NO_READ_AHEAD performance hint are
// exposed. NO_LOCK and WRITE_MAP remain unavailable; the process
// registry still prevents duplicate Env values.
unsafe {
builder.flags(flags);
}
}
// SAFETY: transactions never escape a NIF call, locking is never disabled,
// and this registry prevents duplicate Env values in this process.
let env = unsafe { builder.open(&path) }.map_err(|_| LifecycleError::EnvironmentOpenFailed)?;
let state = Arc::new(EnvironmentState {
id: next_resource_id()?,
path: path.clone(),
options,
env,
databases: Mutex::new(HashMap::new()),
});
entries.insert(path, Arc::downgrade(&state));
Ok(state)
}
pub(crate) fn validate_database_name(name: &[u8]) -> InternalResult<&str> {
if name.is_empty() || name.len() > 511 || name.contains(&0) {
return Err(LifecycleError::InvalidDatabaseName);
}
std::str::from_utf8(name).map_err(|_| LifecycleError::InvalidDatabaseName)
}
pub(crate) fn database_cache(
state: &EnvironmentState,
) -> MutexGuard<'_, HashMap<String, BinaryDatabase>> {
match state.databases.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
}
}
pub(crate) fn create_named_database(
state: &EnvironmentState,
name: &str,
) -> InternalResult<BinaryDatabase> {
if state.options.read_only {
return Err(LifecycleError::ReadOnly);
}
// LMDB requires DBI opens to be serialized within a process. The cache lock
// covers lookup, DBI creation, and commit, then CRUD can use the immutable
// cached handle without taking this lock.
let mut databases = database_cache(state);
if let Some(database) = databases.get(name) {
return Ok(*database);
}
let mut transaction = state
.env
.write_txn()
.map_err(|_| LifecycleError::TransactionFailed)?;
let database = state
.env
.create_database::<Bytes, Bytes>(&mut transaction, Some(name))
.map_err(|_| LifecycleError::DatabaseCreateFailed)?;
transaction
.commit()
.map_err(|_| LifecycleError::TransactionFailed)?;
databases.insert(name.to_owned(), database);
Ok(database)
}
pub(crate) fn open_named_database(
state: &EnvironmentState,
name: &str,
) -> InternalResult<BinaryDatabase> {
let mut databases = database_cache(state);
if let Some(database) = databases.get(name) {
return Ok(*database);
}
let transaction = state
.env
.read_txn()
.map_err(|_| LifecycleError::TransactionFailed)?;
let database = state
.env
.open_database::<Bytes, Bytes>(&transaction, Some(name))
.map_err(|_| LifecycleError::DatabaseOpenFailed)?
.ok_or(LifecycleError::DatabaseNotFound)?;
// heed requires committing read transactions which opened a DBI, especially
// when multiple OS processes use the environment.
transaction
.commit()
.map_err(|_| LifecycleError::TransactionFailed)?;
databases.insert(name.to_owned(), database);
Ok(database)
}
pub(crate) fn decode_environment(
term: Term<'_>,
) -> InternalResult<ResourceArc<EnvironmentResource>> {
term.decode()
.map_err(|_| LifecycleError::InvalidEnvironment)
}
pub(crate) fn decode_database(term: Term<'_>) -> InternalResult<ResourceArc<DatabaseResource>> {
term.decode()
.map_err(|_| LifecycleError::InvalidEnvironment)
}