use std::cell::RefCell;
use std::fs::{self, File, OpenOptions};
use std::marker::PhantomData;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Condvar, Mutex, OnceLock};
use anyhow::{Context, Result, bail};
pub const PERSIST_LOCK_FILE_NAME: &str = ".persist.lock";
pub const LOCK_POLL_INTERVAL: std::time::Duration = std::time::Duration::from_millis(20);
#[derive(Debug, Clone)]
pub struct IndexRemovedDuringPersist {
graph_dir: PathBuf,
}
impl IndexRemovedDuringPersist {
#[must_use]
pub fn new(graph_dir: &Path) -> Self {
Self {
graph_dir: graph_dir.to_path_buf(),
}
}
}
impl std::fmt::Display for IndexRemovedDuringPersist {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"the index directory {} was removed during the persist; no index was written",
self.graph_dir.display()
)
}
}
impl std::error::Error for IndexRemovedDuringPersist {}
pub const LOCK_IDENTITY_RETRIES: usize = 8;
#[derive(Debug, Clone)]
pub struct UnstableLockFile {
lock_path: PathBuf,
}
impl std::fmt::Display for UnstableLockFile {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"the persist lock {} has an unstable file identity: after {LOCK_IDENTITY_RETRIES} \
attempts the locked file was never the one the path names (a filesystem whose \
inode numbers are not stable?)",
self.lock_path.display()
)
}
}
impl std::error::Error for UnstableLockFile {}
enum TryTake {
Taken(IndexWriteLock),
Busy,
NoIndex,
}
#[derive(Debug)]
pub enum LockWait {
Held(IndexWriteLock),
Cancelled,
IndexRemoved,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
struct LockFileIdentity(Option<(u64, u64)>);
impl LockFileIdentity {
fn of(file: &File) -> Self {
Self(file.metadata().ok().and_then(|metadata| file_id(&metadata)))
}
fn matches(&self, path: &Path) -> bool {
match fs::metadata(path) {
Ok(metadata) => match (&self.0, file_id(&metadata)) {
(Some(ours), Some(now)) => *ours == now,
_ => true,
},
Err(_) => false,
}
}
}
#[cfg(unix)]
fn file_id(metadata: &fs::Metadata) -> Option<(u64, u64)> {
use std::os::unix::fs::MetadataExt;
Some((metadata.dev(), metadata.ino()))
}
#[cfg(not(unix))]
fn file_id(_metadata: &fs::Metadata) -> Option<(u64, u64)> {
None
}
pub(crate) const ROLLBACK_MARK: &str = ".rollback-";
pub(crate) const MARKER_BEGUN: &str = ".txn-begun";
pub(crate) const MARKER_COMMITTED: &str = ".txn-committed";
#[derive(Debug)]
struct HoldOwner {
file: File,
}
impl Drop for HoldOwner {
fn drop(&mut self) {
let _ = self.file.unlock();
}
}
struct HeldEntry {
dir: PathBuf,
identity: LockFileIdentity,
owner: std::rc::Weak<HoldOwner>,
}
thread_local! {
static HELD: RefCell<Vec<HeldEntry>> = const { RefCell::new(Vec::new()) };
}
#[derive(Debug)]
pub struct IndexWriteLock {
owner: std::rc::Rc<HoldOwner>,
graph_dir: PathBuf,
identity: LockFileIdentity,
saw_committed_index: bool,
_thread_bound: PhantomData<*const ()>,
}
impl IndexWriteLock {
pub fn acquire(graph_dir: &Path) -> Result<Self> {
if let Some(held) = Self::reenter(graph_dir) {
return Ok(held);
}
for _ in 0..LOCK_IDENTITY_RETRIES {
fs::create_dir_all(graph_dir)
.with_context(|| format!("Failed to create {}", graph_dir.display()))?;
let file = open_lock_file(graph_dir)?;
file.lock().with_context(|| {
format!(
"Failed to take the persist lock {}",
graph_dir.join(PERSIST_LOCK_FILE_NAME).display()
)
})?;
if let Some(identity) = Self::still_named(&file, graph_dir) {
return Ok(Self::held(file, graph_dir, identity).recovered());
}
}
Err(UnstableLockFile {
lock_path: graph_dir.join(PERSIST_LOCK_FILE_NAME),
}
.into())
}
pub fn acquire_if_present(graph_dir: &Path) -> Result<Option<Self>> {
if Self::current_hold(graph_dir).is_some() || graph_dir.is_dir() {
return Self::acquire(graph_dir).map(Some);
}
Ok(None)
}
pub fn acquire_unless_cancelled(
graph_dir: &Path,
cancelled: &dyn Fn() -> bool,
contended: &mut dyn FnMut(),
) -> Result<LockWait> {
Self::wait_unless_cancelled(graph_dir, true, cancelled, contended)
}
pub fn acquire_existing_unless_cancelled(
graph_dir: &Path,
cancelled: &dyn Fn() -> bool,
contended: &mut dyn FnMut(),
) -> Result<LockWait> {
Self::wait_unless_cancelled(graph_dir, false, cancelled, contended)
}
fn wait_unless_cancelled(
graph_dir: &Path,
create: bool,
cancelled: &dyn Fn() -> bool,
contended: &mut dyn FnMut(),
) -> Result<LockWait> {
if Self::held_by_this_thread(graph_dir) {
if cancelled() {
return Ok(LockWait::Cancelled);
}
return Ok(Self::reenter(graph_dir).map_or(LockWait::IndexRemoved, LockWait::Held));
}
let file = if create {
fs::create_dir_all(graph_dir)
.with_context(|| format!("Failed to create {}", graph_dir.display()))?;
open_lock_file(graph_dir)?
} else {
match open_lock_file_in_index(graph_dir)? {
Some(file) => file,
None => return Ok(LockWait::IndexRemoved),
}
};
let identity = LockFileIdentity::of(&file);
let lock_path = graph_dir.join(PERSIST_LOCK_FILE_NAME);
let opened_generation = commit_generation(&file);
let opened_committed = holds_committed_index(graph_dir);
let arrived = |file: File| -> LockWait {
if !identity.matches(&lock_path) {
return LockWait::IndexRemoved;
}
let committed_while_waiting = commit_generation(&file) > opened_generation;
let mut held = Self::held(file, graph_dir, identity.clone()).recovered();
held.saw_committed_index |= opened_committed || committed_while_waiting;
if cancelled() {
drop(held);
return LockWait::Cancelled;
}
LockWait::Held(held)
};
match file.try_lock() {
Ok(()) => return Ok(arrived(file)),
Err(fs::TryLockError::WouldBlock) => drop(file),
Err(fs::TryLockError::Error(e)) => {
return Err(e).with_context(|| {
format!("Failed to take the persist lock {}", lock_path.display())
});
}
}
contended();
let key = SharedLockKey::of(&lock_path, &identity);
loop {
let Some(shared) = SharedLockWait::attach(&key, &lock_path, &identity)? else {
return Ok(LockWait::IndexRemoved);
};
let mut state = shared.lock_state();
loop {
if let Some(answer) = state.answer.take() {
state.interested -= 1;
drop(state);
return match answer {
Ok(file) => Ok(arrived(file)),
Err(e) => Err(e).with_context(|| {
format!("Failed to take the persist lock {}", lock_path.display())
}),
};
}
if state.finished {
state.interested -= 1;
break;
}
state = shared
.answered
.wait_timeout(state, LOCK_POLL_INTERVAL)
.unwrap_or_else(std::sync::PoisonError::into_inner)
.0;
let gave_up = if cancelled() {
Some(LockWait::Cancelled)
} else if !identity.matches(&lock_path) {
Some(LockWait::IndexRemoved)
} else {
None
};
if let Some(outcome) = gave_up {
state.interested -= 1;
if state.interested == 0 {
drop(state.answer.take());
}
return Ok(outcome);
}
}
}
}
pub fn try_acquire(graph_dir: &Path) -> Result<Option<Self>> {
Ok(match Self::try_take(graph_dir)? {
TryTake::Taken(held) => Some(held),
TryTake::Busy | TryTake::NoIndex => None,
})
}
fn try_take(graph_dir: &Path) -> Result<TryTake> {
if Self::held_by_this_thread(graph_dir) {
return Ok(TryTake::Busy);
}
for _ in 0..LOCK_IDENTITY_RETRIES {
let Some(file) = open_lock_file_in_index(graph_dir)? else {
return Ok(TryTake::NoIndex);
};
match file.try_lock() {
Ok(()) => {
if let Some(identity) = Self::still_named(&file, graph_dir) {
return Ok(TryTake::Taken(
Self::held(file, graph_dir, identity).recovered(),
));
}
}
Err(fs::TryLockError::WouldBlock) => return Ok(TryTake::Busy),
Err(fs::TryLockError::Error(e)) => {
return Err(e).with_context(|| {
format!(
"Failed to try the persist lock {}",
graph_dir.join(PERSIST_LOCK_FILE_NAME).display()
)
});
}
}
}
log::warn!(
"{}",
UnstableLockFile {
lock_path: graph_dir.join(PERSIST_LOCK_FILE_NAME),
}
);
Ok(TryTake::NoIndex)
}
#[must_use]
pub fn held_by_this_thread(graph_dir: &Path) -> bool {
Self::current_hold(graph_dir).is_some()
}
#[must_use]
pub fn held_stale_by_this_thread(graph_dir: &Path) -> bool {
let lock_path = graph_dir.join(PERSIST_LOCK_FILE_NAME);
HELD.with(|held| {
held.borrow()
.iter()
.any(|entry| entry.dir == graph_dir && !entry.identity.matches(&lock_path))
})
}
#[must_use]
pub fn is_current(&self) -> bool {
self.identity
.matches(&self.graph_dir.join(PERSIST_LOCK_FILE_NAME))
}
#[must_use]
pub fn held_any_by_this_thread(graph_dir: &Path) -> bool {
HELD.with(|held| held.borrow().iter().any(|entry| entry.dir == graph_dir))
}
#[must_use]
pub fn reenter(graph_dir: &Path) -> Option<Self> {
Self::current_hold(graph_dir).map(|(identity, owner)| Self {
owner,
graph_dir: graph_dir.to_path_buf(),
identity,
saw_committed_index: holds_committed_index(graph_dir),
_thread_bound: PhantomData,
})
}
fn current_hold(graph_dir: &Path) -> Option<(LockFileIdentity, std::rc::Rc<HoldOwner>)> {
let lock_path = graph_dir.join(PERSIST_LOCK_FILE_NAME);
HELD.with(|held| {
held.borrow().iter().find_map(|entry| {
((entry.identity.0.is_some() || entry.dir == graph_dir)
&& entry.identity.matches(&lock_path))
.then(|| entry.owner.upgrade())
.flatten()
.map(|owner| (entry.identity.clone(), owner))
})
})
}
fn still_named(file: &File, graph_dir: &Path) -> Option<LockFileIdentity> {
#[cfg(test)]
if identity_seam::unstable(graph_dir) {
return None;
}
let identity = LockFileIdentity::of(file);
identity
.matches(&graph_dir.join(PERSIST_LOCK_FILE_NAME))
.then_some(identity)
}
fn recovered(mut self) -> Self {
if let Err(e) = recover_interrupted_persist(&self.graph_dir) {
log::error!(
"could not put back the index an interrupted persist left in {}: {e:#}",
self.graph_dir.display()
);
}
self.saw_committed_index = holds_committed_index(&self.graph_dir);
self
}
#[must_use]
pub fn saw_committed_index(&self) -> bool {
self.saw_committed_index
}
pub fn record_commit(&self) -> std::io::Result<()> {
let len = self.owner.file.metadata()?.len();
self.owner.file.set_len(len.saturating_add(1))
}
fn held(file: File, graph_dir: &Path, identity: LockFileIdentity) -> Self {
let owner = std::rc::Rc::new(HoldOwner { file });
HELD.with(|held| {
held.borrow_mut().push(HeldEntry {
dir: graph_dir.to_path_buf(),
identity: identity.clone(),
owner: std::rc::Rc::downgrade(&owner),
});
});
Self {
owner,
graph_dir: graph_dir.to_path_buf(),
identity,
saw_committed_index: false,
_thread_bound: PhantomData,
}
}
}
fn commit_generation(file: &File) -> u64 {
file.metadata().map_or(0, |metadata| metadata.len())
}
impl Drop for IndexWriteLock {
fn drop(&mut self) {
if std::rc::Rc::strong_count(&self.owner) == 1 {
let owner = std::rc::Rc::downgrade(&self.owner);
HELD.with(|held| {
held.borrow_mut()
.retain(|entry| !std::rc::Weak::ptr_eq(&entry.owner, &owner));
});
}
}
}
#[derive(Default)]
struct SharedLockWait {
state: Mutex<SharedLockState>,
answered: Condvar,
}
#[derive(Default)]
struct SharedLockState {
interested: usize,
answer: Option<std::io::Result<File>>,
finished: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
struct SharedLockKey {
path: Option<PathBuf>,
identity: LockFileIdentity,
}
impl SharedLockKey {
fn of(lock_path: &Path, identity: &LockFileIdentity) -> Self {
Self {
path: identity.0.is_none().then(|| lock_path.to_path_buf()),
identity: identity.clone(),
}
}
}
fn shared_lock_waits()
-> &'static Mutex<std::collections::HashMap<SharedLockKey, Arc<SharedLockWait>>> {
static WAITS: OnceLock<Mutex<std::collections::HashMap<SharedLockKey, Arc<SharedLockWait>>>> =
OnceLock::new();
WAITS.get_or_init(Default::default)
}
#[cfg(test)]
pub(crate) mod helper_seam {
use std::path::{Path, PathBuf};
use std::sync::Mutex;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Fault {
SpawnFails,
Panics,
}
static ARMED: Mutex<Vec<(PathBuf, Fault)>> = Mutex::new(Vec::new());
pub(crate) fn arm(lock_path: &Path, fault: Fault) {
ARMED
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push((lock_path.to_path_buf(), fault));
}
pub(crate) fn take(lock_path: &Path) -> Option<Fault> {
let mut armed = ARMED
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let at = armed.iter().position(|(path, _)| path == lock_path)?;
Some(armed.remove(at).1)
}
}
impl SharedLockWait {
fn lock_state(&self) -> std::sync::MutexGuard<'_, SharedLockState> {
self.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
fn attach(
key: &SharedLockKey,
lock_path: &Path,
identity: &LockFileIdentity,
) -> Result<Option<Arc<Self>>> {
let mut waits = shared_lock_waits()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(shared) = waits.get(key) {
let mut state = shared.lock_state();
if !state.finished {
state.interested += 1;
drop(state);
return Ok(Some(Arc::clone(shared)));
}
}
if !identity.matches(lock_path) {
return Ok(None);
}
let file = match OpenOptions::new().read(true).write(true).open(lock_path) {
Ok(file) => file,
Err(_) if !identity.matches(lock_path) => return Ok(None),
Err(e) => {
return Err(e).with_context(|| {
format!("Failed to open the persist lock {}", lock_path.display())
});
}
};
if LockFileIdentity::of(&file) != *identity {
return Ok(None);
}
let shared = Arc::new(Self::default());
shared.lock_state().interested = 1;
let (helper_shared, helper_key) = (Arc::clone(&shared), key.clone());
#[cfg(test)]
let fault = helper_seam::take(lock_path);
#[cfg(test)]
if fault == Some(helper_seam::Fault::SpawnFails) {
return Err(anyhow::anyhow!("seam: the helper thread could not start"))
.context("Failed to start the persist lock waiter");
}
std::thread::Builder::new()
.name("sqry-persist-lock-wait".to_string())
.spawn(move || {
let mut helper = HelperExit {
shared: helper_shared,
key: helper_key,
answer: None,
};
#[cfg(test)]
if fault == Some(helper_seam::Fault::Panics) {
panic!("seam: the persist lock helper panicked");
}
helper.answer = Some(file.lock().map(|()| file));
})
.context("Failed to start the persist lock waiter")?;
waits.insert(key.clone(), Arc::clone(&shared));
Ok(Some(shared))
}
}
struct HelperExit {
shared: Arc<SharedLockWait>,
key: SharedLockKey,
answer: Option<std::io::Result<File>>,
}
impl Drop for HelperExit {
fn drop(&mut self) {
let mut waits = shared_lock_waits()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if waits
.get(&self.key)
.is_some_and(|current| Arc::ptr_eq(current, &self.shared))
{
waits.remove(&self.key);
}
let mut state = self.shared.lock_state();
state.finished = true;
if state.interested > 0 {
state.answer = self.answer.take();
}
self.shared.answered.notify_all();
}
}
#[must_use]
pub fn holds_index_content(graph_dir: &Path) -> bool {
graph_dir.join("manifest.json").exists()
|| graph_dir.join("snapshot.sqry").exists()
|| has_rollback_names(graph_dir)
}
const COMMITTED_READ_ATTEMPTS: usize = 16;
#[must_use]
pub fn holds_committed_index(graph_dir: &Path) -> bool {
let manifest_path = graph_dir.join(super::MANIFEST_FILE_NAME);
let snapshot_path = graph_dir.join(super::SNAPSHOT_FILE_NAME);
let pair = || (presence(&manifest_path), presence(&snapshot_path));
let mut answer = false;
for _ in 0..COMMITTED_READ_ATTEMPTS {
let before = pair();
#[cfg(any(test, feature = "test-support"))]
committed_read_seam::run(graph_dir, committed_read_seam::Point::First);
let listing = RollbackListing::of(graph_dir);
#[cfg(any(test, feature = "test-support"))]
committed_read_seam::run(graph_dir, committed_read_seam::Point::Second);
let after = pair();
#[cfg(any(test, feature = "test-support"))]
committed_read_seam::run(graph_dir, committed_read_seam::Point::Third);
answer = listing.leaves_an_index(after.0.is_some(), after.1.is_some());
if pair_ids(&before) == pair_ids(&after) {
return answer;
}
}
answer
}
struct Seen {
id: Option<(u64, u64)>,
_open: Option<File>,
}
fn presence(path: &Path) -> Option<Seen> {
if let Some(file) = open_regular(path) {
return Some(Seen {
id: file.metadata().ok().and_then(|metadata| file_id(&metadata)),
_open: Some(file),
});
}
fs::metadata(path).ok().map(|metadata| Seen {
id: file_id(&metadata),
_open: None,
})
}
fn open_regular(path: &Path) -> Option<File> {
let mut options = OpenOptions::new();
options.read(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt;
options.custom_flags(libc::O_NONBLOCK);
}
let file = options.open(path).ok()?;
file.metadata().ok()?.is_file().then_some(file)
}
fn pair_ids(pair: &(Option<Seen>, Option<Seen>)) -> [Option<Option<(u64, u64)>>; 2] {
[
pair.0.as_ref().map(|seen| seen.id),
pair.1.as_ref().map(|seen| seen.id),
]
}
fn open_lock_file_in_index(graph_dir: &Path) -> Result<Option<File>> {
let path = graph_dir.join(PERSIST_LOCK_FILE_NAME);
let file = match OpenOptions::new().read(true).write(true).open(&path) {
Ok(file) => file,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
return create_lock_file_in_index(graph_dir, &path);
}
Err(e) => {
return Err(e)
.with_context(|| format!("Failed to open the persist lock {}", path.display()));
}
};
Ok(holds_index_content(graph_dir).then_some(file))
}
fn create_lock_file_in_index(graph_dir: &Path, path: &Path) -> Result<Option<File>> {
if !graph_dir.is_dir() || !holds_index_content(graph_dir) {
return Ok(None);
}
#[cfg(test)]
lock_create_seam::before_create(graph_dir);
let file = match OpenOptions::new()
.read(true)
.write(true)
.create_new(true)
.open(path)
{
Ok(file) => file,
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
#[cfg(test)]
lock_create_seam::after_already_exists(graph_dir);
return match OpenOptions::new().read(true).write(true).open(path) {
Ok(file) => Ok(holds_index_content(graph_dir).then_some(file)),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
Err(e) => Err(e)
.with_context(|| format!("Failed to open the persist lock {}", path.display())),
};
}
Err(_) if !graph_dir.is_dir() => return Ok(None),
Err(e) => {
return Err(e)
.with_context(|| format!("Failed to create the persist lock {}", path.display()));
}
};
#[cfg(test)]
lock_create_seam::after_create(graph_dir);
Ok(holds_index_content(graph_dir).then_some(file))
}
#[cfg(any(test, feature = "test-support"))]
#[doc(hidden)]
pub mod committed_read_seam {
use std::path::{Path, PathBuf};
use std::sync::Mutex;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Point {
First,
Second,
Third,
}
type Action = Box<dyn FnOnce() + Send>;
static ARMED: Mutex<Vec<(PathBuf, Point, Action)>> = Mutex::new(Vec::new());
pub fn arm(graph_dir: &Path, point: Point, action: Action) {
ARMED
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push((graph_dir.to_path_buf(), point, action));
}
pub(crate) fn run(graph_dir: &Path, point: Point) {
let action = {
let mut armed = ARMED
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
armed
.iter()
.position(|(dir, at, _)| dir == graph_dir && *at == point)
.map(|at| armed.remove(at).2)
};
if let Some(action) = action {
action();
}
}
}
#[cfg(test)]
pub(crate) mod identity_seam {
use std::path::{Path, PathBuf};
use std::sync::Mutex;
static UNSTABLE: Mutex<Vec<PathBuf>> = Mutex::new(Vec::new());
pub(crate) fn arm(graph_dir: &Path) {
UNSTABLE
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(graph_dir.to_path_buf());
}
pub(crate) fn unstable(graph_dir: &Path) -> bool {
UNSTABLE
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.iter()
.any(|dir| dir == graph_dir)
}
}
#[cfg(test)]
pub(crate) mod lock_create_seam {
use std::path::{Path, PathBuf};
use std::sync::Mutex;
type Action = Box<dyn FnOnce() + Send>;
static ARMED: Mutex<Vec<(PathBuf, Action)>> = Mutex::new(Vec::new());
static ARMED_BEFORE: Mutex<Vec<(PathBuf, Action)>> = Mutex::new(Vec::new());
static ARMED_AFTER_EXISTS: Mutex<Vec<(PathBuf, Action)>> = Mutex::new(Vec::new());
pub(crate) fn arm(graph_dir: &Path, action: Action) {
ARMED
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push((graph_dir.to_path_buf(), action));
}
pub(crate) fn arm_before(graph_dir: &Path, action: Action) {
ARMED_BEFORE
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push((graph_dir.to_path_buf(), action));
}
fn run(armed: &Mutex<Vec<(PathBuf, Action)>>, graph_dir: &Path) {
let action = {
let mut armed = armed
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
armed
.iter()
.position(|(dir, _)| dir == graph_dir)
.map(|at| armed.remove(at).1)
};
if let Some(action) = action {
action();
}
}
pub(crate) fn after_create(graph_dir: &Path) {
run(&ARMED, graph_dir);
}
pub(crate) fn before_create(graph_dir: &Path) {
run(&ARMED_BEFORE, graph_dir);
}
pub(crate) fn arm_after_already_exists(graph_dir: &Path, action: Action) {
ARMED_AFTER_EXISTS
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push((graph_dir.to_path_buf(), action));
}
pub(crate) fn after_already_exists(graph_dir: &Path) {
run(&ARMED_AFTER_EXISTS, graph_dir);
}
}
fn open_lock_file(graph_dir: &Path) -> Result<File> {
let path = graph_dir.join(PERSIST_LOCK_FILE_NAME);
OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&path)
.with_context(|| format!("Failed to open the persist lock {}", path.display()))
}
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct RecoveryOutcome {
pub restored: usize,
pub discarded: usize,
}
pub fn recover_interrupted_persist(graph_dir: &Path) -> Result<RecoveryOutcome> {
let mut outcome = RecoveryOutcome::default();
let listing = RollbackListing::of(graph_dir);
let manifest_path = graph_dir.join(super::MANIFEST_FILE_NAME);
let snapshot_path = graph_dir.join(super::SNAPSHOT_FILE_NAME);
for set in listing.sets.into_values() {
let (snapshot, manifest_back) = match set.fate(manifest_path.exists()) {
SetFate::Discard => {
set.discard();
outcome.discarded += 1;
continue;
}
SetFate::Restore {
snapshot,
manifest_back,
} => (snapshot, manifest_back),
};
match (snapshot, &set.snapshot_aside) {
(SnapshotRestore::PutBack, Some(aside)) => {
fs::rename(aside, &snapshot_path).with_context(|| {
format!(
"Failed to put the previous snapshot back from {}",
aside.display()
)
})?;
}
(SnapshotRestore::RemoveNew, _) => match fs::remove_file(&snapshot_path) {
Ok(()) => {}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => {
return Err(e).with_context(|| {
format!(
"Failed to remove the snapshot {} an interrupted persist wrote",
snapshot_path.display()
)
});
}
},
_ => {}
}
if manifest_back && let Some(aside) = &set.manifest_aside {
fs::rename(aside, &manifest_path).with_context(|| {
format!(
"Failed to put the previous manifest back from {}",
aside.display()
)
})?;
}
log::warn!(
"put back the index pair an interrupted persist set aside in {}",
graph_dir.display()
);
set.discard();
outcome.restored += 1;
}
Ok(outcome)
}
#[derive(Debug, Default)]
struct RollbackListing {
sets: std::collections::BTreeMap<(u64, u64, u64), RollbackSet>,
}
impl RollbackListing {
fn of(graph_dir: &Path) -> Self {
let mut listing = Self::default();
let Ok(entries) = fs::read_dir(graph_dir) else {
return listing;
};
for entry in entries.flatten() {
let name = entry.file_name().to_string_lossy().into_owned();
let Some(at) = name.find(ROLLBACK_MARK) else {
continue;
};
let Some(order) = parse_tag(&name[at + ROLLBACK_MARK.len()..]) else {
continue;
};
let set = listing.sets.entry(order).or_default();
let path = entry.path();
match &name[..at] {
".manifest.json" => set.manifest_aside = Some(path),
".snapshot.sqry" => set.snapshot_aside = Some(path),
MARKER_BEGUN => {
set.begun_without_snapshot = marker_says_no_snapshot(&path);
set.begun = Some(path);
}
MARKER_COMMITTED => set.committed = Some(path),
_ => set.other.push(path),
}
}
listing
}
fn leaves_an_index(&self, mut manifest: bool, mut snapshot: bool) -> bool {
for set in self.sets.values() {
if let SetFate::Restore {
snapshot: restore,
manifest_back,
} = set.fate(manifest)
{
match restore {
SnapshotRestore::PutBack => snapshot = true,
SnapshotRestore::RemoveNew => snapshot = false,
SnapshotRestore::Keep => {}
}
manifest |= manifest_back;
}
}
manifest || snapshot
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum SetFate {
Discard,
Restore {
snapshot: SnapshotRestore,
manifest_back: bool,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum SnapshotRestore {
PutBack,
RemoveNew,
Keep,
}
pub(crate) fn recover_for_reader(graph_dir: &Path, manifest_path: &Path) {
if manifest_path.exists()
|| IndexWriteLock::held_by_this_thread(graph_dir)
|| !has_rollback_names(graph_dir)
{
return;
}
let held = match IndexWriteLock::try_take(graph_dir) {
Ok(TryTake::Taken(_held)) => Ok(()),
Ok(TryTake::Busy) => {
IndexWriteLock::acquire_existing_unless_cancelled(graph_dir, &|| false, &mut || {})
.map(drop)
}
Ok(TryTake::NoIndex) => Ok(()),
Err(e) => Err(e),
};
if let Err(e) = held {
log::warn!(
"could not check {} for an interrupted persist: {e:#}",
graph_dir.display()
);
}
}
fn has_rollback_names(graph_dir: &Path) -> bool {
fs::read_dir(graph_dir).is_ok_and(|entries| {
entries
.flatten()
.any(|entry| entry.file_name().to_string_lossy().contains(ROLLBACK_MARK))
})
}
fn parse_tag(tag: &str) -> Option<(u64, u64, u64)> {
let mut parts = tag.split('-');
let pid = parts.next()?.parse().ok()?;
let secs = parts.next()?.parse().ok()?;
let seq = parts.next()?.parse().ok()?;
if parts.next().is_some() {
return None;
}
Some((secs, pid, seq))
}
fn marker_says_no_snapshot(marker: &Path) -> bool {
use std::io::Read;
let mut text = String::new();
open_regular(marker).is_some_and(|mut file| file.read_to_string(&mut text).is_ok())
&& text.contains("snapshot_existed=0")
}
#[derive(Debug, Default)]
struct RollbackSet {
manifest_aside: Option<PathBuf>,
snapshot_aside: Option<PathBuf>,
begun: Option<PathBuf>,
begun_without_snapshot: bool,
committed: Option<PathBuf>,
other: Vec<PathBuf>,
}
impl RollbackSet {
fn fate(&self, manifest_present: bool) -> SetFate {
if self.committed.is_some() || manifest_present {
return SetFate::Discard;
}
let snapshot = if self.snapshot_aside.is_some() {
SnapshotRestore::PutBack
} else if self.begun_without_snapshot {
SnapshotRestore::RemoveNew
} else {
SnapshotRestore::Keep
};
SetFate::Restore {
snapshot,
manifest_back: self.manifest_aside.is_some(),
}
}
fn discard(self) {
for path in [self.manifest_aside, self.snapshot_aside]
.into_iter()
.flatten()
.chain(self.other)
.chain([self.begun, self.committed].into_iter().flatten())
{
match fs::remove_file(&path) {
Ok(()) => {}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => log::warn!("could not remove {}: {e}", path.display()),
}
}
}
}
#[derive(Debug, Default)]
pub struct PersistGate {
state: Mutex<GateState>,
idle: Condvar,
}
#[derive(Debug, Default)]
struct GateState {
closed: bool,
in_flight: usize,
}
#[derive(Debug)]
pub struct PersistTicket<'a> {
gate: &'a PersistGate,
}
impl PersistGate {
#[must_use]
pub fn new() -> Self {
Self::default()
}
pub fn global() -> &'static Self {
static GATE: OnceLock<PersistGate> = OnceLock::new();
GATE.get_or_init(PersistGate::new)
}
pub fn enter(&self) -> Result<PersistTicket<'_>> {
let mut state = self.lock();
if state.closed {
bail!("the process is shutting down; the persist was not started");
}
state.in_flight += 1;
Ok(PersistTicket { gate: self })
}
pub fn close_and_wait(&self) {
let mut state = self.lock();
state.closed = true;
while state.in_flight > 0 {
state = self
.idle
.wait(state)
.unwrap_or_else(std::sync::PoisonError::into_inner);
}
}
#[must_use]
pub fn is_closed(&self) -> bool {
self.lock().closed
}
#[must_use]
pub fn in_flight(&self) -> usize {
self.lock().in_flight
}
fn lock(&self) -> std::sync::MutexGuard<'_, GateState> {
self.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
}
impl Drop for PersistTicket<'_> {
fn drop(&mut self) {
let mut state = self.gate.lock();
state.in_flight -= 1;
if state.in_flight == 0 {
self.gate.idle.notify_all();
}
}
}
pub fn close_persists_and_wait() {
PersistGate::global().close_and_wait();
}
#[doc(hidden)]
pub fn set_mid_persist_hook(hook: fn(&Path)) {
let _ = MID_PERSIST_HOOK.set(hook);
}
static MID_PERSIST_HOOK: OnceLock<fn(&Path)> = OnceLock::new();
pub(crate) fn run_mid_persist_hook(graph_dir: &Path) {
if let Some(hook) = MID_PERSIST_HOOK.get() {
hook(graph_dir);
}
}
#[cfg(all(test, target_os = "linux"))]
pub(crate) mod lock_waiters {
use std::os::unix::fs::MetadataExt;
use std::path::Path;
use std::time::{Duration, Instant};
pub(crate) fn someone_waits(graph_dir: &Path) -> bool {
let Ok(metadata) = std::fs::metadata(graph_dir.join(super::PERSIST_LOCK_FILE_NAME)) else {
return false;
};
let dev = metadata.dev();
let wanted = (
((dev >> 8) & 0xfff) | ((dev >> 32) & 0xffff_f000),
(dev & 0xff) | ((dev >> 12) & 0xffff_ff00),
metadata.ino(),
);
let locks = std::fs::read_to_string("/proc/locks").expect("read /proc/locks");
locks.lines().any(|line| {
let fields: Vec<&str> = line.split_whitespace().collect();
fields
.iter()
.position(|field| *field == "->")
.and_then(|at| fields.get(at + 5))
.and_then(|id| {
let mut parts = id.split(':');
let major = u64::from_str_radix(parts.next()?, 16).ok()?;
let minor = u64::from_str_radix(parts.next()?, 16).ok()?;
let ino = parts.next()?.parse::<u64>().ok()?;
parts.next().is_none().then_some((major, minor, ino))
})
.is_some_and(|waiter| waiter == wanted)
})
}
pub(crate) fn wait_until_someone_waits(graph_dir: &Path, gave_up: impl Fn() -> bool) -> bool {
let deadline = Instant::now() + Duration::from_secs(60);
loop {
if someone_waits(graph_dir) {
return true;
}
if gave_up() {
return false;
}
assert!(
Instant::now() < deadline,
"no thread waited for the persist lock of {}",
graph_dir.display()
);
std::thread::sleep(Duration::from_millis(5));
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
#[test]
fn the_lock_excludes_another_thread_and_is_reentrant_on_its_own() {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
let held = IndexWriteLock::acquire(&graph_dir).unwrap();
fs::write(graph_dir.join("manifest.json"), b"{}").unwrap();
let again = IndexWriteLock::acquire(&graph_dir).unwrap();
drop(again);
assert!(IndexWriteLock::held_by_this_thread(&graph_dir));
assert!(
!another_thread_takes(&graph_dir),
"another thread must not get a held lock"
);
drop(held);
assert!(!IndexWriteLock::held_by_this_thread(&graph_dir));
assert!(
another_thread_takes(&graph_dir),
"a released lock is free for another thread"
);
}
#[test]
#[cfg(target_os = "linux")]
fn two_threads_exclude_each_other_across_the_guard_lifecycle() {
use super::lock_waiters::wait_until_someone_waits;
use std::sync::mpsc;
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
let inside = Arc::new(AtomicBool::new(false));
let first = IndexWriteLock::acquire(&graph_dir).unwrap();
fs::write(graph_dir.join("manifest.json"), b"{}").unwrap();
inside.store(true, Ordering::SeqCst);
let again = IndexWriteLock::acquire(&graph_dir).unwrap();
let (to_main, from_other) = mpsc::channel::<&'static str>();
let (to_other, from_main) = mpsc::channel::<()>();
let other = {
let (graph_dir, inside) = (graph_dir.clone(), Arc::clone(&inside));
std::thread::spawn(move || {
to_main.send("waiting").unwrap();
let held = IndexWriteLock::acquire(&graph_dir).unwrap();
assert!(
!inside.swap(true, Ordering::SeqCst),
"the second thread got the lock while the first still held it"
);
to_main.send("holding").unwrap();
from_main.recv().unwrap();
inside.store(false, Ordering::SeqCst);
drop(held);
})
};
assert_eq!(from_other.recv().unwrap(), "waiting");
assert!(wait_until_someone_waits(&graph_dir, || other.is_finished()));
drop(again);
assert!(
wait_until_someone_waits(&graph_dir, || other.is_finished()),
"dropping a re-entrant hold must not release the lock"
);
assert!(
from_other.try_recv().is_err(),
"the second thread must wait while the first holds"
);
inside.store(false, Ordering::SeqCst);
drop(first);
assert!(!IndexWriteLock::held_by_this_thread(&graph_dir));
assert_eq!(from_other.recv().unwrap(), "holding");
assert!(
IndexWriteLock::try_acquire(&graph_dir).unwrap().is_none(),
"the first thread must not re-enter a lock the second now holds"
);
to_other.send(()).unwrap();
other.join().unwrap();
let back = IndexWriteLock::try_acquire(&graph_dir).unwrap();
assert!(back.is_some(), "a released lock is free again");
assert!(!inside.load(Ordering::SeqCst));
}
#[test]
#[cfg(target_os = "linux")]
fn a_reader_waits_out_a_live_set_aside_window() {
use crate::graph::unified::persistence::GraphStorage;
use std::sync::mpsc;
let dir = tempfile::tempdir().unwrap();
let storage = GraphStorage::new(dir.path());
fs::create_dir_all(storage.graph_dir()).unwrap();
fs::write(storage.manifest_path(), b"{\"old\": true}").unwrap();
let graph_dir = storage.graph_dir().to_path_buf();
let manifest = storage.manifest_path().to_path_buf();
let (to_main, from_writer) = mpsc::channel::<()>();
let (release, released) = mpsc::channel::<()>();
let writer = std::thread::spawn(move || {
let _held = IndexWriteLock::acquire(&graph_dir).unwrap();
let tag = format!("{ROLLBACK_MARK}{}-1-1", std::process::id());
fs::rename(&manifest, graph_dir.join(format!(".manifest.json{tag}"))).unwrap();
fs::write(
graph_dir.join(format!("{MARKER_BEGUN}{tag}")),
"snapshot_existed=1",
)
.unwrap();
to_main.send(()).unwrap();
released.recv().unwrap();
fs::write(&manifest, b"{\"new\": true}").unwrap();
for entry in fs::read_dir(&graph_dir).unwrap().flatten() {
if entry.file_name().to_string_lossy().contains(ROLLBACK_MARK) {
fs::remove_file(entry.path()).unwrap();
}
}
});
from_writer.recv().unwrap();
let reader = {
let root = dir.path().to_path_buf();
std::thread::spawn(move || GraphStorage::new(&root).exists())
};
super::lock_waiters::wait_until_someone_waits(storage.graph_dir(), || reader.is_finished());
release.send(()).unwrap();
assert!(
reader.join().unwrap(),
"a live transaction's set-aside window was read as no index"
);
writer.join().unwrap();
assert_eq!(
fs::read(storage.manifest_path()).unwrap(),
b"{\"new\": true}",
"the committed manifest stays; the reader recovered nothing over it"
);
}
#[test]
fn a_cancellable_wait_gives_up_on_a_cancel_and_takes_a_free_lock() {
use std::sync::mpsc;
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
let (held_tx, held_rx) = mpsc::channel::<()>();
let (release_tx, release_rx) = mpsc::channel::<()>();
let holder = {
let graph_dir = graph_dir.clone();
std::thread::spawn(move || {
let _held = IndexWriteLock::acquire(&graph_dir).unwrap();
held_tx.send(()).unwrap();
release_rx.recv().unwrap();
})
};
held_rx.recv().unwrap();
let cancelled = AtomicBool::new(false);
let mut contended = 0;
let gave_up = IndexWriteLock::acquire_unless_cancelled(
&graph_dir,
&|| cancelled.load(Ordering::SeqCst),
&mut || {
contended += 1;
cancelled.store(true, Ordering::SeqCst);
},
)
.unwrap();
assert!(
matches!(gave_up, LockWait::Cancelled),
"a cancelled wait holds nothing: {gave_up:?}"
);
assert_eq!(contended, 1, "the contended callback runs once");
assert!(!IndexWriteLock::held_by_this_thread(&graph_dir));
release_tx.send(()).unwrap();
holder.join().unwrap();
let taken =
IndexWriteLock::acquire_unless_cancelled(&graph_dir, &|| false, &mut || {}).unwrap();
assert!(matches!(taken, LockWait::Held(_)), "a free lock is taken");
assert!(IndexWriteLock::held_by_this_thread(&graph_dir));
}
#[test]
fn a_wait_ends_when_the_index_directory_is_removed() {
use std::sync::mpsc;
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
let (held_tx, held_rx) = mpsc::channel::<()>();
let (release_tx, release_rx) = mpsc::channel::<()>();
let holder = {
let graph_dir = graph_dir.clone();
std::thread::spawn(move || {
let _held = IndexWriteLock::acquire(&graph_dir).unwrap();
held_tx.send(()).unwrap();
release_rx.recv().unwrap();
})
};
held_rx.recv().unwrap();
let (contended_tx, contended_rx) = mpsc::channel::<()>();
let (done_tx, done_rx) = mpsc::channel::<String>();
{
let graph_dir = graph_dir.clone();
std::thread::spawn(move || {
let outcome =
IndexWriteLock::acquire_unless_cancelled(&graph_dir, &|| false, &mut || {
let _ = contended_tx.send(());
});
let _ = done_tx.send(format!("{outcome:?}"));
});
}
contended_rx.recv().unwrap();
fs::remove_dir_all(dir.path().join(".sqry")).unwrap();
let outcome = done_rx.recv_timeout(Duration::from_secs(3));
release_tx.send(()).unwrap();
holder.join().unwrap();
assert_eq!(
outcome.as_deref(),
Ok("Ok(IndexRemoved)"),
"the wait must end once the index is removed, while the holder still holds"
);
assert!(!graph_dir.exists(), "the wait created nothing again");
}
#[test]
fn a_cancellable_wait_is_not_starved_by_blocking_waiters() {
use std::sync::atomic::AtomicUsize;
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
let stop = Arc::new(AtomicBool::new(false));
let holds = Arc::new(AtomicUsize::new(0));
let blockers: Vec<_> = (0..2)
.map(|_| {
let (graph_dir, stop, holds) =
(graph_dir.clone(), Arc::clone(&stop), Arc::clone(&holds));
std::thread::spawn(move || {
while !stop.load(Ordering::SeqCst) {
let _held = IndexWriteLock::acquire(&graph_dir).unwrap();
holds.fetch_add(1, Ordering::SeqCst);
std::thread::sleep(Duration::from_millis(50));
}
})
})
.collect();
while holds.load(Ordering::SeqCst) < 2 {
std::thread::yield_now();
}
let deadline = std::time::Instant::now() + Duration::from_secs(5);
let outcome = IndexWriteLock::acquire_unless_cancelled(
&graph_dir,
&|| std::time::Instant::now() > deadline,
&mut || {},
)
.unwrap();
let won = matches!(outcome, LockWait::Held(_));
drop(outcome);
stop.store(true, Ordering::SeqCst);
for blocker in blockers {
blocker.join().unwrap();
}
assert!(
won,
"the cancellable wait never won a handoff in 5 s ({} blocking holds)",
holds.load(Ordering::SeqCst)
);
}
#[test]
#[cfg(target_os = "linux")]
fn a_stale_hold_is_not_re_entrant_once_the_lock_file_is_replaced() {
use super::lock_waiters::wait_until_someone_waits;
use std::sync::mpsc;
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
let (a_held_tx, a_held_rx) = mpsc::channel::<()>();
let (go_tx, go_rx) = mpsc::channel::<()>();
let (a_report_tx, a_report_rx) = mpsc::channel::<(bool, bool)>();
let a_returned = Arc::new(AtomicBool::new(false));
let a = {
let (graph_dir, a_returned) = (graph_dir.clone(), Arc::clone(&a_returned));
std::thread::spawn(move || {
let first = IndexWriteLock::acquire(&graph_dir).unwrap();
a_held_tx.send(()).unwrap();
go_rx.recv().unwrap();
let stale_reported_current = first.is_current();
let again = IndexWriteLock::acquire(&graph_dir).unwrap();
a_returned.store(true, Ordering::SeqCst);
a_report_tx
.send((stale_reported_current, again.is_current()))
.unwrap();
drop(again);
drop(first);
})
};
a_held_rx.recv().unwrap();
fs::remove_dir_all(dir.path().join(".sqry")).unwrap();
let (b_held_tx, b_held_rx) = mpsc::channel::<()>();
let (b_release_tx, b_release_rx) = mpsc::channel::<()>();
let b = {
let graph_dir = graph_dir.clone();
std::thread::spawn(move || {
let _held = IndexWriteLock::acquire(&graph_dir).unwrap();
b_held_tx.send(()).unwrap();
b_release_rx.recv().unwrap();
})
};
b_held_rx.recv().unwrap();
go_tx.send(()).unwrap();
let a_waited = wait_until_someone_waits(&graph_dir, || a_returned.load(Ordering::SeqCst));
b_release_tx.send(()).unwrap();
b.join().unwrap();
let (stale_reported_current, again_current) = a_report_rx.recv().unwrap();
a.join().unwrap();
assert!(
a_waited,
"the thread with a stale hold re-entered while another thread held the current lock"
);
assert!(
!stale_reported_current,
"a hold on a removed lock file is not current"
);
assert!(
again_current,
"the hold taken after the wait is on the current lock file"
);
}
#[test]
#[cfg(target_os = "linux")]
fn an_acquire_that_wins_an_unlinked_lock_file_retries_on_the_current_one() {
use super::lock_waiters::wait_until_someone_waits;
use std::sync::mpsc;
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
let (x_held_tx, x_held_rx) = mpsc::channel::<()>();
let (x_release_tx, x_release_rx) = mpsc::channel::<()>();
let x = {
let graph_dir = graph_dir.clone();
std::thread::spawn(move || {
let _held = IndexWriteLock::acquire(&graph_dir).unwrap();
x_held_tx.send(()).unwrap();
x_release_rx.recv().unwrap();
})
};
x_held_rx.recv().unwrap();
let a_done = Arc::new(AtomicBool::new(false));
let (a_held_tx, a_held_rx) = mpsc::channel::<()>();
let (a_release_tx, a_release_rx) = mpsc::channel::<()>();
let a = {
let (graph_dir, a_done) = (graph_dir.clone(), Arc::clone(&a_done));
std::thread::spawn(move || {
let held = IndexWriteLock::acquire(&graph_dir).unwrap();
a_done.store(true, Ordering::SeqCst);
a_held_tx.send(()).unwrap();
a_release_rx.recv().unwrap();
drop(held);
})
};
assert!(wait_until_someone_waits(&graph_dir, || a_done.load(Ordering::SeqCst)));
fs::remove_dir_all(dir.path().join(".sqry")).unwrap();
x_release_tx.send(()).unwrap();
x.join().unwrap();
a_held_rx.recv().unwrap();
let c_done = Arc::new(AtomicBool::new(false));
let c = {
let (graph_dir, c_done) = (graph_dir.clone(), Arc::clone(&c_done));
std::thread::spawn(move || {
let _held = IndexWriteLock::acquire(&graph_dir).unwrap();
c_done.store(true, Ordering::SeqCst);
})
};
let c_waited = wait_until_someone_waits(&graph_dir, || c_done.load(Ordering::SeqCst));
a_release_tx.send(()).unwrap();
a.join().unwrap();
c.join().unwrap();
assert!(
c_waited,
"a third acquirer took the recreated lock file while the second held the removed one"
);
}
fn hold_on_a_thread(
graph_dir: &Path,
) -> (std::sync::mpsc::Sender<()>, std::thread::JoinHandle<()>) {
use std::sync::mpsc;
let (held_tx, held_rx) = mpsc::channel::<()>();
let (release_tx, release_rx) = mpsc::channel::<()>();
let graph_dir = graph_dir.to_path_buf();
let holder = std::thread::spawn(move || {
let _held = IndexWriteLock::acquire(&graph_dir).unwrap();
held_tx.send(()).unwrap();
let _ = release_rx.recv();
});
held_rx.recv().unwrap();
(release_tx, holder)
}
#[test]
fn a_helper_that_fails_to_start_leaves_no_dead_wait() {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
let (release, holder) = hold_on_a_thread(&graph_dir);
helper_seam::arm(
&graph_dir.join(PERSIST_LOCK_FILE_NAME),
helper_seam::Fault::SpawnFails,
);
let first = IndexWriteLock::acquire_unless_cancelled(&graph_dir, &|| false, &mut || {});
assert!(first.is_err(), "the failed start is reported: {first:?}");
let deadline = std::time::Instant::now() + Duration::from_secs(5);
let mut release = Some(release);
let second = IndexWriteLock::acquire_unless_cancelled(
&graph_dir,
&|| std::time::Instant::now() > deadline,
&mut || {
drop(release.take());
},
)
.unwrap();
holder.join().unwrap();
assert!(
matches!(second, LockWait::Held(_)),
"a later wait attached to the dead entry: {second:?}"
);
}
#[test]
fn a_helper_that_panics_leaves_no_dead_wait() {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
let (release, holder) = hold_on_a_thread(&graph_dir);
helper_seam::arm(
&graph_dir.join(PERSIST_LOCK_FILE_NAME),
helper_seam::Fault::Panics,
);
let deadline = std::time::Instant::now() + Duration::from_secs(5);
let mut release = Some(release);
let outcome = IndexWriteLock::acquire_unless_cancelled(
&graph_dir,
&|| std::time::Instant::now() > deadline,
&mut || drop(release.take()),
)
.unwrap();
holder.join().unwrap();
assert!(
matches!(outcome, LockWait::Held(_)),
"the wait stayed attached to a panicked helper: {outcome:?}"
);
}
#[test]
fn an_existing_index_wait_takes_no_lock_in_an_emptied_directory() {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
fs::write(graph_dir.join("manifest.json"), b"{}").unwrap();
let emptied = graph_dir.clone();
lock_create_seam::arm(
&graph_dir,
Box::new(move || fs::remove_file(emptied.join("manifest.json")).unwrap()),
);
let outcome =
IndexWriteLock::acquire_existing_unless_cancelled(&graph_dir, &|| false, &mut || {})
.unwrap();
let left: Vec<_> = fs::read_dir(&graph_dir)
.unwrap()
.flatten()
.map(|entry| entry.file_name())
.collect();
assert!(
matches!(outcome, LockWait::IndexRemoved),
"a lock file created into a directory being emptied was taken: {outcome:?}"
);
assert_eq!(
left,
vec![std::ffi::OsString::from(PERSIST_LOCK_FILE_NAME)],
"only the created lock file is left"
);
assert!(
another_thread_takes_the_file(&graph_dir),
"the left lock file is still locked"
);
}
fn another_thread_takes_the_file(graph_dir: &Path) -> bool {
let path = graph_dir.join(PERSIST_LOCK_FILE_NAME);
std::thread::spawn(move || {
OpenOptions::new()
.read(true)
.write(true)
.open(&path)
.is_ok_and(|file| file.try_lock().is_ok())
})
.join()
.unwrap()
}
#[test]
fn try_acquire_takes_no_lock_in_an_emptied_directory() {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
assert!(IndexWriteLock::try_acquire(&graph_dir).unwrap().is_none());
assert!(
!graph_dir.join(PERSIST_LOCK_FILE_NAME).exists(),
"a lock file was created into a directory with no index"
);
fs::write(graph_dir.join("snapshot.sqry"), b"x").unwrap();
let emptied = graph_dir.clone();
lock_create_seam::arm(
&graph_dir,
Box::new(move || fs::remove_file(emptied.join("snapshot.sqry")).unwrap()),
);
assert!(IndexWriteLock::try_acquire(&graph_dir).unwrap().is_none());
assert!(
another_thread_takes_the_file(&graph_dir),
"the created lock file is gone or still locked"
);
}
#[test]
fn a_committed_index_is_what_recovery_would_leave() {
let begun = |tag: &str, snapshot: u8, manifest: u8| {
(
format!("{MARKER_BEGUN}{ROLLBACK_MARK}{tag}"),
format!("begun snapshot_existed={snapshot} manifest_existed={manifest}\n"),
)
};
let file = |name: &str, body: &str| (name.to_string(), body.to_string());
type Case<'a> = (&'a str, Vec<(String, String)>, bool);
let cases: Vec<Case<'_>> = vec![
("empty directory", vec![], false),
(
"lock file only",
vec![file(PERSIST_LOCK_FILE_NAME, "")],
false,
),
(
"manifest and snapshot",
vec![file("manifest.json", "{}"), file("snapshot.sqry", "s")],
true,
),
("snapshot only", vec![file("snapshot.sqry", "s")], true),
("manifest only", vec![file("manifest.json", "{}")], true),
(
"a killed first persist",
vec![begun("1-1-1", 0, 0), file("snapshot.sqry", "partial")],
false,
),
(
"a first persist killed before its snapshot",
vec![begun("1-1-1", 0, 0)],
false,
),
(
"a first persist killed after its manifest",
vec![
begun("1-1-1", 0, 0),
file("snapshot.sqry", "new"),
file("manifest.json", "{}"),
],
true,
),
(
"a killed update of a full index",
vec![
begun("1-1-1", 1, 1),
file(&format!(".manifest.json{ROLLBACK_MARK}1-1-1"), "{}"),
file(&format!(".snapshot.sqry{ROLLBACK_MARK}1-1-1"), "old"),
file("snapshot.sqry", "partial"),
],
true,
),
(
"a killed update of a snapshot-only index",
vec![
begun("1-1-1", 1, 0),
file(&format!(".snapshot.sqry{ROLLBACK_MARK}1-1-1"), "old"),
file("snapshot.sqry", "partial"),
],
true,
),
(
"an update killed before its snapshot was kept aside",
vec![begun("1-1-1", 1, 0), file("snapshot.sqry", "old")],
true,
),
(
"an older killed first persist and a newer killed update",
vec![
begun("1-1-1", 0, 0),
begun("2-2-2", 1, 1),
file(&format!(".manifest.json{ROLLBACK_MARK}2-2-2"), "{}"),
file(&format!(".snapshot.sqry{ROLLBACK_MARK}2-2-2"), "old"),
file("snapshot.sqry", "partial"),
],
true,
),
(
"a committed set left behind",
vec![
file(&format!("{MARKER_COMMITTED}{ROLLBACK_MARK}1-1-1"), ""),
file("snapshot.sqry", "new"),
file("manifest.json", "{}"),
],
true,
),
];
assert_eq!(cases.len(), 13, "cases");
for (case, files, expected) in cases {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
for (name, body) in &files {
fs::write(graph_dir.join(name), body).unwrap();
}
assert_eq!(holds_committed_index(&graph_dir), expected, "{case}");
recover_interrupted_persist(&graph_dir).unwrap();
let left = graph_dir.join("manifest.json").exists()
|| graph_dir.join("snapshot.sqry").exists();
assert_eq!(left, expected, "{case}: what recovery left");
assert!(!has_rollback_names(&graph_dir), "{case}: sets left");
}
let dir = tempfile::tempdir().unwrap();
assert!(
!holds_committed_index(&dir.path().join("missing")),
"a missing directory"
);
}
#[test]
fn a_first_persist_begun_and_rolled_back_during_the_read_is_no_index() {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
let marker = graph_dir.join(format!("{MARKER_BEGUN}{ROLLBACK_MARK}1-1-1"));
let snapshot = graph_dir.join("snapshot.sqry");
{
let (marker, snapshot) = (marker.clone(), snapshot.clone());
committed_read_seam::arm(
&graph_dir,
committed_read_seam::Point::First,
Box::new(move || {
fs::write(&marker, b"begun snapshot_existed=0 manifest_existed=0\n").unwrap();
fs::write(&snapshot, b"new").unwrap();
}),
);
}
committed_read_seam::arm(
&graph_dir,
committed_read_seam::Point::Second,
Box::new(move || {
fs::remove_file(&snapshot).unwrap();
fs::remove_file(&marker).unwrap();
}),
);
let answer = holds_committed_index(&graph_dir);
let left: Vec<_> = fs::read_dir(&graph_dir)
.unwrap()
.flatten()
.map(|entry| entry.file_name())
.collect();
assert!(left.is_empty(), "the writer left {left:?}");
assert!(!answer, "a phantom committed index");
}
#[test]
#[cfg(unix)]
fn a_fifo_in_the_index_does_not_block_the_read_or_recovery() {
let marker_name = format!("{MARKER_BEGUN}{ROLLBACK_MARK}1-1-1");
let cases: [(&str, &str, Vec<&str>, bool); 4] = [
("snapshot", "snapshot.sqry", vec![], true),
("manifest", "manifest.json", vec![], true),
("marker", marker_name.as_str(), vec!["snapshot.sqry"], true),
(
"marker, recovery",
marker_name.as_str(),
vec!["snapshot.sqry"],
true,
),
];
for (case, fifo, files, expected) in cases {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
for name in files {
fs::write(graph_dir.join(name), b"s").unwrap();
}
let fifo = graph_dir.join(fifo);
let made = std::process::Command::new("mkfifo")
.arg(&fifo)
.status()
.expect("run mkfifo");
assert!(made.success(), "{case}: mkfifo failed");
let (tx, rx) = std::sync::mpsc::channel();
{
let graph_dir = graph_dir.clone();
let recover = case.ends_with("recovery");
std::thread::spawn(move || {
let answer = if recover {
recover_interrupted_persist(&graph_dir).is_ok()
} else {
holds_committed_index(&graph_dir)
};
let _ = tx.send(answer);
});
}
let answer = rx.recv_timeout(Duration::from_secs(5));
if answer.is_err() {
let _ = OpenOptions::new().write(true).open(&fifo);
}
assert_eq!(answer, Ok(expected), "{case}");
}
}
#[test]
fn a_first_persist_rolled_back_after_the_last_look_is_no_index() {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
begin_a_first_persist(&graph_dir, "1-1-1");
let rolled_back = graph_dir.clone();
committed_read_seam::arm(
&graph_dir,
committed_read_seam::Point::Third,
Box::new(move || {
fs::remove_file(rolled_back.join("snapshot.sqry")).unwrap();
fs::remove_file(rolled_back.join(format!("{MARKER_BEGUN}{ROLLBACK_MARK}1-1-1")))
.unwrap();
}),
);
assert!(
!holds_committed_index(&graph_dir),
"a phantom committed index"
);
}
fn begin_a_first_persist(graph_dir: &Path, tag: &str) {
fs::write(
graph_dir.join(format!("{MARKER_BEGUN}{ROLLBACK_MARK}{tag}")),
b"begun snapshot_existed=0 manifest_existed=0\n",
)
.unwrap();
let partial = graph_dir.join(format!("snapshot.sqry.tmp-{tag}"));
fs::write(&partial, tag.as_bytes()).unwrap();
fs::rename(&partial, graph_dir.join("snapshot.sqry")).unwrap();
}
#[test]
fn a_first_persist_begun_after_the_listing_is_no_index() {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
let begun = graph_dir.clone();
committed_read_seam::arm(
&graph_dir,
committed_read_seam::Point::Second,
Box::new(move || begin_a_first_persist(&begun, "2-2-2")),
);
assert!(
!holds_committed_index(&graph_dir),
"a phantom committed index"
);
}
#[test]
#[cfg(unix)]
fn a_snapshot_replaced_during_the_read_is_seen_as_changed() {
use std::os::unix::fs::MetadataExt;
use std::sync::Mutex;
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
begin_a_first_persist(&graph_dir, "1-1-1");
let snapshot = graph_dir.join("snapshot.sqry");
let ino_a = fs::metadata(&snapshot).unwrap().ino();
let rolled_back = graph_dir.clone();
committed_read_seam::arm(
&graph_dir,
committed_read_seam::Point::First,
Box::new(move || {
fs::remove_file(rolled_back.join("snapshot.sqry")).unwrap();
fs::remove_file(rolled_back.join(format!("{MARKER_BEGUN}{ROLLBACK_MARK}1-1-1")))
.unwrap();
}),
);
let seen: Arc<Mutex<Option<(bool, u64)>>> = Arc::default();
{
let seen = Arc::clone(&seen);
let begun = graph_dir.clone();
let deleted = format!("{} (deleted)", snapshot.display());
committed_read_seam::arm(
&graph_dir,
committed_read_seam::Point::Second,
Box::new(move || {
let held = held_open(&deleted);
begin_a_first_persist(&begun, "2-2-2");
let ino_b = fs::metadata(begun.join("snapshot.sqry")).unwrap().ino();
*seen.lock().unwrap() = Some((held, ino_b));
}),
);
}
let answer = holds_committed_index(&graph_dir);
let (held, ino_b) = seen.lock().unwrap().expect("the seam ran");
assert!(held, "A's unlinked snapshot was not held open when B began");
assert_ne!(ino_b, ino_a, "B's snapshot took A's inode");
assert!(!answer, "a phantom committed index");
}
#[cfg(unix)]
fn held_open(target: &str) -> bool {
if !cfg!(target_os = "linux") {
return true;
}
let entries = fs::read_dir("/proc/self/fd").expect("read /proc/self/fd");
entries.flatten().any(|entry| {
fs::read_link(entry.path()).is_ok_and(|link| link.to_string_lossy() == target)
})
}
#[test]
#[cfg(target_os = "linux")]
fn a_wait_reports_an_index_it_saw() {
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Case {
AtOpen,
CommittedRecorded,
CommittedUnrecorded,
}
for case in [
Case::AtOpen,
Case::CommittedRecorded,
Case::CommittedUnrecorded,
] {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
let held = IndexWriteLock::acquire(&graph_dir).unwrap();
assert!(!held.saw_committed_index(), "{case:?}: no index yet");
if case == Case::AtOpen {
fs::write(graph_dir.join("manifest.json"), b"{}").unwrap();
}
let waiter = {
let graph_dir = graph_dir.clone();
std::thread::spawn(move || {
match IndexWriteLock::acquire_unless_cancelled(
&graph_dir,
&|| false,
&mut || {},
)
.unwrap()
{
LockWait::Held(held) => held.saw_committed_index(),
other => panic!("the wait answered {other:?}"),
}
})
};
assert!(
lock_waiters::wait_until_someone_waits(&graph_dir, || waiter.is_finished()),
"{case:?}: the wait was not contended"
);
if case != Case::AtOpen {
fs::write(graph_dir.join("manifest.json"), b"{}").unwrap();
}
if case == Case::CommittedRecorded {
held.record_commit().unwrap();
}
fs::remove_file(graph_dir.join("manifest.json")).unwrap();
drop(held);
assert_eq!(
waiter.join().unwrap(),
case != Case::CommittedUnrecorded,
"{case:?}"
);
}
}
#[test]
fn a_hold_reports_the_committed_index_it_found() {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
fs::write(
graph_dir.join(format!("{MARKER_BEGUN}{ROLLBACK_MARK}1-1-1")),
b"begun snapshot_existed=0 manifest_existed=0\n",
)
.unwrap();
fs::write(graph_dir.join("snapshot.sqry"), b"partial").unwrap();
let held = IndexWriteLock::acquire(&graph_dir).unwrap();
assert!(!held.saw_committed_index(), "a first persist rolled back");
assert!(!graph_dir.join("snapshot.sqry").exists());
fs::write(graph_dir.join("snapshot.sqry"), b"s").unwrap();
let again = IndexWriteLock::reenter(&graph_dir).expect("re-entry");
assert!(again.saw_committed_index(), "a snapshot in place");
drop((again, held));
}
fn existing_index_attempt(graph_dir: &Path, via_wait: bool) -> &'static str {
if via_wait {
match IndexWriteLock::acquire_existing_unless_cancelled(
graph_dir,
&|| false,
&mut || {},
)
.unwrap()
{
LockWait::Held(_) => "Held",
LockWait::Cancelled => "Cancelled",
LockWait::IndexRemoved => "IndexRemoved",
}
} else {
match IndexWriteLock::try_take(graph_dir).unwrap() {
TryTake::Taken(_) => "Taken",
TryTake::Busy => "Busy",
TryTake::NoIndex => "NoIndex",
}
}
}
fn assert_left_untouched(case: &str, other: &File, lock_path: &Path, graph_dir: &Path) {
assert!(
LockFileIdentity::of(other).matches(lock_path),
"{case}: the lock file another opener owns was removed"
);
assert!(
!IndexWriteLock::held_any_by_this_thread(graph_dir),
"{case}: a hold was kept"
);
assert!(
other.try_lock().is_ok(),
"{case}: the lock file another opener owns is still locked"
);
other.unlock().unwrap();
}
#[test]
#[cfg(unix)]
fn a_lock_file_another_opener_created_in_the_window_is_left_and_not_taken() {
for via_wait in [true, false] {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
fs::write(graph_dir.join("manifest.json"), b"{}").unwrap();
let lock_path = graph_dir.join(PERSIST_LOCK_FILE_NAME);
let (other_tx, other_rx) = std::sync::mpsc::channel::<File>();
{
let (lock_path, graph_dir) = (lock_path.clone(), graph_dir.clone());
lock_create_seam::arm_before(
&graph_dir.clone(),
Box::new(move || {
let other = OpenOptions::new()
.read(true)
.write(true)
.create_new(true)
.open(&lock_path)
.unwrap();
fs::remove_file(graph_dir.join("manifest.json")).unwrap();
other_tx.send(other).unwrap();
}),
);
}
let case = format!("via_wait={via_wait}");
let answer = existing_index_attempt(&graph_dir, via_wait);
let other = other_rx
.try_recv()
.unwrap_or_else(|_| panic!("{case}: the seam never ran"));
assert_left_untouched(&case, &other, &lock_path, &graph_dir);
let expected = if via_wait { "IndexRemoved" } else { "NoIndex" };
assert_eq!(
answer, expected,
"{case}: an index with no content was locked"
);
}
}
#[test]
#[cfg(unix)]
fn a_lock_file_another_opener_created_with_content_kept_is_taken() {
for via_wait in [true, false] {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
fs::write(graph_dir.join("manifest.json"), b"{}").unwrap();
let lock_path = graph_dir.join(PERSIST_LOCK_FILE_NAME);
let (other_tx, other_rx) = std::sync::mpsc::channel::<File>();
{
let lock_path = lock_path.clone();
lock_create_seam::arm_before(
&graph_dir,
Box::new(move || {
let other = OpenOptions::new()
.read(true)
.write(true)
.create_new(true)
.open(&lock_path)
.unwrap();
other_tx.send(other).unwrap();
}),
);
}
let case = format!("via_wait={via_wait}");
let taken = if via_wait {
match IndexWriteLock::acquire_existing_unless_cancelled(
&graph_dir,
&|| false,
&mut || {},
)
.unwrap()
{
LockWait::Held(held) => Ok(held),
other => Err(format!("{other:?}")),
}
} else {
match IndexWriteLock::try_take(&graph_dir).unwrap() {
TryTake::Taken(held) => Ok(held),
TryTake::Busy => Err("Busy".to_string()),
TryTake::NoIndex => Err("NoIndex".to_string()),
}
};
let other = other_rx
.try_recv()
.unwrap_or_else(|_| panic!("{case}: the seam never ran"));
let held = taken.unwrap_or_else(|answer| {
panic!("{case}: the other opener's file over an index was not taken: {answer}")
});
assert_eq!(
held.identity,
LockFileIdentity::of(&other),
"{case}: the lock was taken on another file than the other opener's"
);
assert!(
LockFileIdentity::of(&other).matches(&lock_path),
"{case}: the other opener's lock file was removed"
);
}
}
#[test]
fn a_lock_file_gone_again_after_another_opener_created_it_is_no_index() {
for via_wait in [true, false] {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
fs::write(graph_dir.join("manifest.json"), b"{}").unwrap();
let lock_path = graph_dir.join(PERSIST_LOCK_FILE_NAME);
{
let lock_path = lock_path.clone();
lock_create_seam::arm_before(
&graph_dir,
Box::new(move || {
OpenOptions::new()
.read(true)
.write(true)
.create_new(true)
.open(&lock_path)
.unwrap();
}),
);
}
{
let lock_path = lock_path.clone();
lock_create_seam::arm_after_already_exists(
&graph_dir,
Box::new(move || fs::remove_file(&lock_path).unwrap()),
);
}
let case = format!("via_wait={via_wait}");
let answer = if via_wait {
IndexWriteLock::acquire_existing_unless_cancelled(&graph_dir, &|| false, &mut || {})
.map(|outcome| format!("{outcome:?}"))
} else {
IndexWriteLock::try_take(&graph_dir).map(|taken| {
match taken {
TryTake::Taken(_) => "Taken",
TryTake::Busy => "Busy",
TryTake::NoIndex => "NoIndex",
}
.to_string()
})
}
.unwrap_or_else(|e| format!("Err({e:#})"));
let expected = if via_wait { "IndexRemoved" } else { "NoIndex" };
assert_eq!(answer, expected, "{case}");
let left: Vec<_> = fs::read_dir(&graph_dir)
.unwrap()
.flatten()
.map(|entry| entry.file_name())
.collect();
assert_eq!(
left,
vec![std::ffi::OsString::from("manifest.json")],
"{case}: something was created"
);
}
}
#[test]
fn a_lock_file_present_without_index_content_is_left_and_not_taken() {
for via_wait in [true, false] {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
let lock_path = graph_dir.join(PERSIST_LOCK_FILE_NAME);
let other = OpenOptions::new()
.read(true)
.write(true)
.create_new(true)
.open(&lock_path)
.unwrap();
let case = format!("via_wait={via_wait}");
let answer = existing_index_attempt(&graph_dir, via_wait);
assert_left_untouched(&case, &other, &lock_path, &graph_dir);
let expected = if via_wait { "IndexRemoved" } else { "NoIndex" };
assert_eq!(
answer, expected,
"{case}: an index with no content was locked"
);
}
}
#[test]
fn a_reader_recovery_creates_nothing_when_the_index_is_removed() {
use crate::graph::unified::persistence::GraphStorage;
let dir = tempfile::tempdir().unwrap();
let storage = GraphStorage::new(dir.path());
let graph_dir = storage.graph_dir().to_path_buf();
fs::create_dir_all(&graph_dir).unwrap();
fs::write(
graph_dir.join(format!(".manifest.json{ROLLBACK_MARK}1-1-1")),
b"{}",
)
.unwrap();
let sqry_dir = dir.path().join(".sqry");
let removed = sqry_dir.clone();
lock_create_seam::arm(
&graph_dir,
Box::new(move || fs::remove_dir_all(&removed).unwrap()),
);
let exists = storage.exists();
let entries: Vec<_> = fs::read_dir(&graph_dir)
.map(|entries| entries.flatten().map(|entry| entry.file_name()).collect())
.unwrap_or_default();
assert!(!exists, "no index is reported");
assert!(
!sqry_dir.exists(),
"the reader recreated the removed index directory: {entries:?}"
);
}
#[test]
#[cfg(target_os = "linux")]
fn a_waiting_reader_creates_nothing_when_the_index_is_removed() {
use super::lock_waiters::wait_until_someone_waits;
use crate::graph::unified::persistence::GraphStorage;
use std::sync::mpsc;
let dir = tempfile::tempdir().unwrap();
let root = dir.path().to_path_buf();
let graph_dir = GraphStorage::new(&root).graph_dir().to_path_buf();
let (held_tx, held_rx) = mpsc::channel::<()>();
let (release_tx, release_rx) = mpsc::channel::<()>();
let holder = {
let graph_dir = graph_dir.clone();
std::thread::spawn(move || {
let _held = IndexWriteLock::acquire(&graph_dir).unwrap();
fs::write(
graph_dir.join(format!(".manifest.json{ROLLBACK_MARK}1-1-1")),
b"{}",
)
.unwrap();
held_tx.send(()).unwrap();
let _ = release_rx.recv();
})
};
held_rx.recv().unwrap();
let reader = {
let root = root.clone();
std::thread::spawn(move || GraphStorage::new(&root).exists())
};
assert!(
wait_until_someone_waits(&graph_dir, || reader.is_finished()),
"the reader did not wait for the live transaction"
);
fs::remove_dir_all(root.join(".sqry")).unwrap();
release_tx.send(()).unwrap();
holder.join().unwrap();
let deadline = std::time::Instant::now() + Duration::from_secs(10);
while !reader.is_finished() {
assert!(
std::time::Instant::now() < deadline,
"the reader kept waiting after the index was removed"
);
std::thread::sleep(Duration::from_millis(5));
}
let exists = reader.join().unwrap();
let entries: Vec<_> = fs::read_dir(&graph_dir)
.map(|entries| entries.flatten().map(|entry| entry.file_name()).collect())
.unwrap_or_default();
assert!(!exists, "no index is reported");
assert!(
!root.join(".sqry").exists(),
"the waiting reader recreated the removed index directory: {entries:?}"
);
}
#[test]
#[cfg(unix)]
fn a_hold_is_found_under_another_spelling_of_the_directory() {
use crate::graph::unified::persistence::GraphStorage;
use std::sync::mpsc;
let dir = tempfile::tempdir().unwrap();
let real = dir.path().join("real");
fs::create_dir_all(&real).unwrap();
let link = dir.path().join("link");
std::os::unix::fs::symlink(&real, &link).unwrap();
let (done_tx, done_rx) = mpsc::channel::<(bool, bool)>();
std::thread::spawn(move || {
let real_graph = GraphStorage::new(&real).graph_dir().to_path_buf();
let link_graph = GraphStorage::new(&link).graph_dir().to_path_buf();
let held = IndexWriteLock::acquire(&real_graph).unwrap();
fs::write(
real_graph.join(format!(".manifest.json{ROLLBACK_MARK}1-1-1")),
b"{}",
)
.unwrap();
let reentered = IndexWriteLock::reenter(&link_graph).is_some();
let _ = GraphStorage::new(&link).exists();
drop(held);
let _ = done_tx.send((reentered, true));
});
let outcome = done_rx.recv_timeout(Duration::from_secs(10));
assert_eq!(
outcome,
Ok((true, true)),
"a hold under one spelling was not this thread's hold under another \
(Err(Timeout): the reader waited on its own lock)"
);
}
fn another_thread_takes(graph_dir: &Path) -> bool {
let graph_dir = graph_dir.to_path_buf();
std::thread::spawn(
move || match IndexWriteLock::try_take(&graph_dir).unwrap() {
TryTake::Taken(_) => true,
TryTake::Busy => false,
TryTake::NoIndex => panic!("the probe found no index in {}", graph_dir.display()),
},
)
.join()
.unwrap()
}
#[test]
fn the_lock_is_held_until_the_last_guard_drops() {
for (via_reenter, original_first) in
[(false, true), (false, false), (true, true), (true, false)]
{
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
let original = IndexWriteLock::acquire(&graph_dir).unwrap();
fs::write(graph_dir.join("manifest.json"), b"{}").unwrap();
let reentrant = if via_reenter {
IndexWriteLock::reenter(&graph_dir).expect("a current hold to re-enter")
} else {
IndexWriteLock::acquire(&graph_dir).unwrap()
};
let case = format!("via_reenter={via_reenter} original_first={original_first}");
let survivor = if original_first {
drop(original);
reentrant
} else {
drop(reentrant);
original
};
assert!(
IndexWriteLock::held_by_this_thread(&graph_dir),
"{case}: the thread lost its hold record with a guard alive"
);
assert!(survivor.is_current(), "{case}");
assert!(
!another_thread_takes(&graph_dir),
"{case}: another thread took the lock while a guard was alive"
);
drop(survivor);
assert!(!IndexWriteLock::held_by_this_thread(&graph_dir), "{case}");
assert!(
another_thread_takes(&graph_dir),
"{case}: the lock was not released when the last guard dropped"
);
}
}
#[test]
fn the_identity_retries_are_bounded() {
use std::sync::mpsc;
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
fs::write(graph_dir.join("manifest.json"), b"{}").unwrap();
identity_seam::arm(&graph_dir);
let (tx, rx) = mpsc::channel::<(String, bool)>();
{
let graph_dir = graph_dir.clone();
std::thread::spawn(move || {
let acquired = match IndexWriteLock::acquire(&graph_dir) {
Ok(_) => "Ok".to_string(),
Err(err) => {
let typed = err.chain().any(|cause| cause.is::<UnstableLockFile>());
format!("Err(typed={typed}): {err:#}")
}
};
let tried = IndexWriteLock::try_acquire(&graph_dir).unwrap().is_some();
let _ = tx.send((acquired, tried));
});
}
let (acquired, tried) = rx
.recv_timeout(Duration::from_secs(10))
.expect("the identity retries never ended");
assert!(
acquired.starts_with("Err(typed=true)") && acquired.contains("unstable"),
"{acquired}"
);
assert!(acquired.contains(".persist.lock"), "{acquired}");
assert!(
!tried,
"try_acquire took a lock whose identity never matched"
);
}
#[test]
fn cancellation_is_checked_with_and_without_a_prior_hold() {
for existing in [false, true] {
for prior_hold in [false, true] {
let dir = tempfile::tempdir().unwrap();
let graph_dir = dir.path().join(".sqry/graph");
fs::create_dir_all(&graph_dir).unwrap();
fs::write(graph_dir.join("manifest.json"), b"{}").unwrap();
let original = prior_hold.then(|| IndexWriteLock::acquire(&graph_dir).unwrap());
let calls = std::cell::Cell::new(0usize);
let cancelled = || {
calls.set(calls.get() + 1);
true
};
let outcome = if existing {
IndexWriteLock::acquire_existing_unless_cancelled(
&graph_dir,
&cancelled,
&mut || {},
)
} else {
IndexWriteLock::acquire_unless_cancelled(&graph_dir, &cancelled, &mut || {})
}
.unwrap();
let case = format!("existing={existing} prior_hold={prior_hold}");
assert!(
matches!(outcome, LockWait::Cancelled),
"{case}: {outcome:?} after {} cancel checks",
calls.get()
);
assert!(
calls.get() >= 1,
"{case}: the cancellation was never checked"
);
drop(outcome);
if let Some(original) = original {
assert!(
IndexWriteLock::held_by_this_thread(&graph_dir),
"{case}: the cancelled attempt dropped the original hold's record"
);
assert!(original.is_current(), "{case}");
assert!(
!another_thread_takes(&graph_dir),
"{case}: the cancelled attempt released the original hold"
);
drop(original);
}
assert!(
another_thread_takes(&graph_dir),
"{case}: the lock was not free at the end"
);
}
}
}
#[test]
fn a_re_entrant_wait_over_a_removed_lock_file_creates_nothing() {
for existing in [false, true] {
for whole_dir in [false, true] {
let dir = tempfile::tempdir().unwrap();
let sqry_dir = dir.path().join(".sqry");
let graph_dir = sqry_dir.join("graph");
fs::create_dir_all(&graph_dir).unwrap();
fs::write(graph_dir.join("manifest.json"), b"{}").unwrap();
let held = IndexWriteLock::acquire(&graph_dir).unwrap();
let lock_file = graph_dir.join(PERSIST_LOCK_FILE_NAME);
let remove = || {
if whole_dir {
let _ = fs::remove_dir_all(&sqry_dir);
} else {
let _ = fs::remove_file(&lock_file);
}
false
};
let outcome = if existing {
IndexWriteLock::acquire_existing_unless_cancelled(
&graph_dir,
&remove,
&mut || {},
)
} else {
IndexWriteLock::acquire_unless_cancelled(&graph_dir, &remove, &mut || {})
}
.unwrap();
let case = format!("existing={existing} whole_dir={whole_dir}");
assert!(
matches!(outcome, LockWait::IndexRemoved),
"{case}: {outcome:?}"
);
drop(outcome);
assert!(!lock_file.exists(), "{case}: the lock file was recreated");
if whole_dir {
assert!(!sqry_dir.exists(), "{case}: the directory was recreated");
}
drop(held);
}
}
}
#[test]
fn the_gate_waits_for_a_persist_in_flight_and_refuses_new_ones() {
let gate = Arc::new(PersistGate::new());
let finished = Arc::new(AtomicBool::new(false));
let entered = Arc::new(std::sync::Barrier::new(2));
let worker = {
let (gate, finished, entered) = (
Arc::clone(&gate),
Arc::clone(&finished),
Arc::clone(&entered),
);
std::thread::spawn(move || {
let _ticket = gate.enter().unwrap();
entered.wait();
std::thread::sleep(Duration::from_millis(300));
finished.store(true, Ordering::Release);
})
};
entered.wait();
gate.close_and_wait();
assert!(
finished.load(Ordering::Acquire),
"close_and_wait returned before the persist in flight finished"
);
assert!(gate.enter().is_err(), "a closed gate refuses a new persist");
worker.join().unwrap();
}
}