Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
250 changes: 250 additions & 0 deletions crates/fff-core/src/dbs/env_pool.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,250 @@
use heed::{Env, EnvOpenOptions};
use std::collections::HashMap;
use std::fs;
use std::ops::Deref;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, LazyLock, Mutex, MutexGuard, PoisonError, Weak};
use std::thread;
use std::time::Duration;

use crate::error::{Error, Result};
use crate::lmdb::DbHealth;

pub(crate) struct EnvSpec {
pub label: &'static str,
pub map_size: usize,
pub max_dbs: u32,
pub size_cap_bytes: u64,
}

pub(crate) struct PooledEnv {
env: Env,
key: PathBuf,
/// lmdb's env spec label
label: &'static str,
map_size: usize,
max_dbs: u32,
health: DbHealth,
gc_started: AtomicBool,
dbi_lock: Mutex<()>,
}

impl Drop for PooledEnv {
fn drop(&mut self) {
let mut pool = POOL.lock().unwrap_or_else(PoisonError::into_inner);
// Only remove a dead entry: begin_exclusive_destroy may have removed ours.
if pool.get(&self.key).is_some_and(|w| w.strong_count() == 0) {
pool.remove(&self.key);
}
// heed closes the env right after this body; a concurrent reopen of the
// same path rides out that gap via env_closing_event in get_or_open.
}
}

// Cloneable handle to a process-shared LMDB env, derefs to `heed::Env`.
#[derive(Clone)]
pub(crate) struct SharedEnv(Arc<PooledEnv>);

impl Deref for SharedEnv {
type Target = Env;
fn deref(&self) -> &Env {
&self.0.env
}
}

impl std::fmt::Debug for SharedEnv {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_tuple("SharedEnv").field(&self.0.env).finish()
}
}

impl SharedEnv {
pub(crate) fn get_or_open(db_path: &Path, spec: &EnvSpec) -> Result<Self> {
fs::create_dir_all(db_path).map_err(Error::CreateDir)?;
let path = fs::canonicalize(db_path).map_err(|e| Error::EnvOpen {
db: spec.label,
source: heed::Error::Io(e),
})?;

let mut close_waits = 0u32;
let mut transient_retries = 0u32;

loop {
let mut open_failed = false;

{
let mut pool = POOL.lock().unwrap_or_else(PoisonError::into_inner);
if let Some(existing) = pool.get(&path).and_then(Weak::upgrade) {
drop(pool);
if existing.label != spec.label
|| existing.map_size != spec.map_size
|| existing.max_dbs != spec.max_dbs
{
return Err(Error::EnvSpecMismatch {
path,
open_as: existing.label,
requested_as: spec.label,
});
}
return Ok(Self(existing));
}

erase_if_oversized(&path, spec);
let result = unsafe {
let mut opts = EnvOpenOptions::new();
opts.map_size(spec.map_size);
if spec.max_dbs > 0 {
opts.max_dbs(spec.max_dbs);
}
opts.open(&path)
};

match result {
Ok(env) => {
let entry = Arc::new(PooledEnv {
env,
key: path.clone(),
label: spec.label,
map_size: spec.map_size,
max_dbs: spec.max_dbs,
health: DbHealth::new(),
gc_started: AtomicBool::new(false),
dbi_lock: Mutex::new(()),
});
pool.insert(path.clone(), Arc::downgrade(&entry));
drop(pool);
let shared = Self(entry);

match shared.clear_stale_readers() {
Ok(cleared_count) if cleared_count > 0 => {
tracing::info!(
cleared_count,
db = spec.label,
"reclaimed stale LMDB reader slots at open"
);
}
Ok(_) => {}
Err(e) => {
tracing::debug!("clear_stale_readers at open failed: {e}")
}
}

return Ok(shared);
}
Err(heed::Error::EnvAlreadyOpened) => open_failed = true,
// special handling cause we know this happens randomly
Err(e)
if is_transient_env_open_error(&e)
&& transient_retries < MAX_TRANSIENT_RETRIES =>
{
transient_retries += 1;
tracing::debug!(
path = %path.display(),
transient_retries,
error = ?e,
"transient LMDB env open error, retrying"
);
}
Err(e) => {
return Err(Error::EnvOpen {
db: spec.label,
source: e,
});
}
}
}

if open_failed {
close_waits += 1;
if close_waits > MAX_CLOSE_WAITS {
return Err(Error::EnvOpen {
db: spec.label,
source: heed::Error::EnvAlreadyOpened,
});
}

match heed::env_closing_event(&path) {
Some(event) => {
event.wait_timeout(CLOSE_WAIT);
}
None => thread::sleep(Duration::from_millis(2)),
}
} else {
thread::sleep(TRANSIENT_RETRY_SLEEP);
}
}
}

pub(crate) fn health(&self) -> &DbHealth {
&self.0.health
}

// First caller wins: GC runs once per opened env, not once per tracker.
pub(crate) fn try_start_gc(&self) -> bool {
!self.0.gc_started.swap(true, Ordering::AcqRel)
}

// LMDB forbids mdb_dbi_open from concurrent txns in the same process.
pub(crate) fn lock_dbi_open(&self) -> MutexGuard<'_, ()> {
self.0
.dbi_lock
.lock()
.unwrap_or_else(PoisonError::into_inner)
}

pub(crate) fn destroy(&self) -> Result<Option<heed::EnvClosingEvent>> {
let mut pool = POOL.lock().unwrap_or_else(PoisonError::into_inner);
let holders = Arc::strong_count(&self.0);

if holders > 1 {
return Err(Error::DbInUse {
db: self.0.label,
path: self.0.key.clone(),
holders: holders - 1,
});
}

pool.remove(&self.0.key);
Ok(heed::env_closing_event(&self.0.key))
}
}

static POOL: LazyLock<Mutex<HashMap<PathBuf, Weak<PooledEnv>>>> = LazyLock::new(Mutex::default);

const CLOSE_WAIT: Duration = Duration::from_millis(100);
const MAX_CLOSE_WAITS: u32 = 100;
const TRANSIENT_RETRY_SLEEP: Duration = Duration::from_millis(50);
const MAX_TRANSIENT_RETRIES: u32 = 8;

// Concurrent mdb_env_open calls on the same path can race on macOS
// this is for some reason fixable by simple retry of the open
fn is_transient_env_open_error(err: &heed::Error) -> bool {
match err {
heed::Error::Io(io) => matches!(
io.kind(),
std::io::ErrorKind::InvalidInput | std::io::ErrorKind::NotFound
),
_ => false,
}
}

fn erase_if_oversized(db_path: &Path, spec: &EnvSpec) {
let data = db_path.join("data.mdb");
let Ok(meta) = fs::metadata(&data) else {
return;
};

if meta.len() <= spec.size_cap_bytes {
return;
}

tracing::error!(
path = %db_path.display(),
size = meta.len(),
cap = spec.size_cap_bytes,
"LMDB db exceeds size cap, erasing"
);
let _ = fs::remove_file(&data);
let _ = fs::remove_file(db_path.join("lock.mdb"));
}
Comment thread
dmtrKovalenko marked this conversation as resolved.
13 changes: 7 additions & 6 deletions crates/fff-core/src/dbs/frecency.rs
Original file line number Diff line number Diff line change
@@ -1,10 +1,11 @@
use super::db_healthcheck::DbHealthChecker;
use super::lmdb::{DbHealth, LmdbStore, is_map_full};
use super::env_pool::SharedEnv;
use crate::error::{Error, Result};
use crate::file_picker::FFFMode;
use crate::git::is_modified_status;
use crate::lmdb::{DbHealth, LmdbStore, is_map_full};
use heed::Database;
use heed::types::{Bytes, SerdeBincode};
use heed::{Database, Env};
use std::time::{SystemTime, UNIX_EPOCH};
use std::{collections::VecDeque, path::Path};

Expand All @@ -19,7 +20,7 @@ const AI_MAX_HISTORY_DAYS: f64 = 7.0; // Only consider accesses within 7 days

#[derive(Debug)]
pub struct FrecencyTracker {
env: Env,
env: SharedEnv,
db: Database<Bytes, SerdeBincode<VecDeque<u64>>>,
health: DbHealth,
}
Expand Down Expand Up @@ -77,15 +78,15 @@ impl LmdbStore for FrecencyTracker {
// MAP_SIZE so we don't hit MDB_MAP_FULL before the open-time erase fires.
const SIZE_CAP_BYTES: u64 = 12 * 1024 * 1024;

fn env(&self) -> &Env {
fn shared_env(&self) -> &SharedEnv {
&self.env
}

fn health(&self) -> &DbHealth {
&self.health
}

fn purge_stale_data(env: &Env) -> Result<()> {
fn purge_stale_data(env: &SharedEnv) -> Result<()> {
let (deleted, pruned) = Self::purge_stale_entries(env)?;
if deleted > 0 || pruned > 0 {
tracing::info!(deleted, pruned, "Frecency GC purged entries");
Expand Down Expand Up @@ -121,7 +122,7 @@ impl FrecencyTracker {
/// Removes entries where all timestamps are older than MAX_HISTORY_DAYS,
/// and prunes stale timestamps from entries that still have recent ones.
/// Returns (deleted_count, pruned_count).
fn purge_stale_entries(env: &Env) -> Result<(usize, usize)> {
fn purge_stale_entries(env: &SharedEnv) -> Result<(usize, usize)> {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
Expand Down
Loading
Loading