mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
fix(nostr): use advisory lifecycle locks
This commit is contained in:
Generated
+11
@@ -660,6 +660,16 @@ dependencies = [
|
||||
"percent-encoding",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "fs2"
|
||||
version = "0.4.3"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "9564fc758e15025b46aa6643b1b77d047d1a56a1aea6e01002ac0c7026876213"
|
||||
dependencies = [
|
||||
"libc",
|
||||
"winapi",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "fs_extra"
|
||||
version = "1.3.0"
|
||||
@@ -1396,6 +1406,7 @@ dependencies = [
|
||||
"dashmap",
|
||||
"dotenvy",
|
||||
"flate2",
|
||||
"fs2",
|
||||
"futures-util",
|
||||
"grasp-audit",
|
||||
"http-body-util",
|
||||
|
||||
@@ -30,6 +30,7 @@ futures-util = "0.3"
|
||||
base64 = "0.22"
|
||||
flate2 = "1.0"
|
||||
tar = "0.4"
|
||||
fs2 = "0.4"
|
||||
|
||||
# Metrics
|
||||
prometheus = "0.14"
|
||||
|
||||
@@ -5,6 +5,7 @@
|
||||
//! lock storage is owned here so repository runtime coordination is not coupled
|
||||
//! to holding/deletion storage.
|
||||
|
||||
use std::fs::File;
|
||||
use std::hash::{Hash, Hasher};
|
||||
use std::io::ErrorKind;
|
||||
use std::path::{Path, PathBuf};
|
||||
@@ -12,6 +13,7 @@ use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use dashmap::DashMap;
|
||||
use fs2::FileExt;
|
||||
use tokio::sync::{OwnedRwLockReadGuard, OwnedRwLockWriteGuard, RwLock as AsyncRwLock};
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, std::hash::Hash)]
|
||||
@@ -26,19 +28,7 @@ pub struct LifecycleReadGuard {
|
||||
|
||||
pub struct LifecycleWriteGuard {
|
||||
_memory_guard: OwnedRwLockWriteGuard<()>,
|
||||
lock_path: Option<PathBuf>,
|
||||
}
|
||||
|
||||
impl Drop for LifecycleWriteGuard {
|
||||
fn drop(&mut self) {
|
||||
if let Some(path) = &self.lock_path {
|
||||
if let Err(e) = std::fs::remove_file(path) {
|
||||
if e.kind() != ErrorKind::NotFound {
|
||||
tracing::warn!(path = %path.display(), error = %e, "Failed to remove lifecycle lock file");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
_lock_file: Option<File>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
@@ -54,7 +44,9 @@ pub struct RepositoryLifecycleLocks {
|
||||
/// short-lived CLI write operations such as `holding-eject`, which run in a
|
||||
/// separate process and cannot share this in-memory map. Read-side Git Smart
|
||||
/// HTTP operations share the in-memory read lock only, while write-side
|
||||
/// operations also take the lock file.
|
||||
/// operations also take an advisory lock on an open lock file. The lock file
|
||||
/// is intentionally left in place; only the advisory lock matters, and the
|
||||
/// OS releases it automatically if the process dies.
|
||||
locks: Arc<DashMap<RecoveryScope, Arc<AsyncRwLock<()>>>>,
|
||||
lock_dir: Option<PathBuf>,
|
||||
}
|
||||
@@ -161,16 +153,16 @@ impl RepositoryLifecycle {
|
||||
.value()
|
||||
.clone();
|
||||
let memory_guard = lock.write_owned().await;
|
||||
let lock_path = self.acquire_scope_file_lock(&scope).await;
|
||||
let lock_file = self.acquire_scope_file_lock(&scope).await;
|
||||
guards.push(LifecycleWriteGuard {
|
||||
_memory_guard: memory_guard,
|
||||
lock_path,
|
||||
_lock_file: lock_file,
|
||||
});
|
||||
}
|
||||
guards
|
||||
}
|
||||
|
||||
async fn acquire_scope_file_lock(&self, scope: &RecoveryScope) -> Option<PathBuf> {
|
||||
async fn acquire_scope_file_lock(&self, scope: &RecoveryScope) -> Option<File> {
|
||||
let lock_dir = self.locks.lock_dir.as_ref()?;
|
||||
if let Err(e) = std::fs::create_dir_all(lock_dir) {
|
||||
tracing::warn!(path = %lock_dir.display(), error = %e, "Failed to create lifecycle lock directory");
|
||||
@@ -182,13 +174,29 @@ impl RepositoryLifecycle {
|
||||
let lock_path = lock_dir.join(format!("{:016x}.lock", hasher.finish()));
|
||||
|
||||
loop {
|
||||
match std::fs::OpenOptions::new()
|
||||
let file = match std::fs::OpenOptions::new()
|
||||
.read(true)
|
||||
.write(true)
|
||||
.create_new(true)
|
||||
.create(true)
|
||||
.truncate(false)
|
||||
.open(&lock_path)
|
||||
{
|
||||
Ok(_) => return Some(lock_path),
|
||||
Err(e) if e.kind() == ErrorKind::AlreadyExists => {
|
||||
Ok(file) => file,
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
owner = %scope.owner_pubkey_hex,
|
||||
identifier = %scope.identifier,
|
||||
path = %lock_path.display(),
|
||||
error = %e,
|
||||
"Failed to open lifecycle lock file; continuing with in-process lock only"
|
||||
);
|
||||
return None;
|
||||
}
|
||||
};
|
||||
|
||||
match file.try_lock_exclusive() {
|
||||
Ok(()) => return Some(file),
|
||||
Err(e) if e.kind() == ErrorKind::WouldBlock => {
|
||||
tokio::time::sleep(Duration::from_millis(25)).await;
|
||||
}
|
||||
Err(e) => {
|
||||
@@ -197,7 +205,7 @@ impl RepositoryLifecycle {
|
||||
identifier = %scope.identifier,
|
||||
path = %lock_path.display(),
|
||||
error = %e,
|
||||
"Failed to acquire lifecycle lock file; continuing with in-process lock only"
|
||||
"Failed to lock lifecycle lock file; continuing with in-process lock only"
|
||||
);
|
||||
return None;
|
||||
}
|
||||
@@ -339,4 +347,37 @@ mod tests {
|
||||
.expect("second reader task should join");
|
||||
drop(second_read_guard);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn lifecycle_write_lock_ignores_stale_lock_file() {
|
||||
let git_data_dir = tempfile::tempdir().expect("temp git data dir");
|
||||
let lifecycle = RepositoryLifecycle::for_git_data_path(git_data_dir.path());
|
||||
let owner_hex = Keys::generate().public_key().to_hex();
|
||||
let identifier = "stale-lock-file";
|
||||
let scope = RecoveryScope {
|
||||
owner_pubkey_hex: owner_hex.clone(),
|
||||
identifier: identifier.to_string(),
|
||||
};
|
||||
|
||||
let lock_dir = git_data_dir.path().join(".lifecycle-locks");
|
||||
std::fs::create_dir_all(&lock_dir).expect("create lock dir");
|
||||
let mut hasher = std::collections::hash_map::DefaultHasher::new();
|
||||
scope.hash(&mut hasher);
|
||||
let lock_path = lock_dir.join(format!("{:016x}.lock", hasher.finish()));
|
||||
std::fs::write(&lock_path, b"stale process died before cleanup")
|
||||
.expect("create stale lock file");
|
||||
|
||||
let write_guard = tokio::time::timeout(
|
||||
Duration::from_millis(100),
|
||||
lifecycle.write_repository(&owner_hex, identifier),
|
||||
)
|
||||
.await
|
||||
.expect("stale lock file should not block advisory lock acquisition");
|
||||
drop(write_guard);
|
||||
|
||||
assert!(
|
||||
lock_path.exists(),
|
||||
"advisory lock files remain in place so file existence cannot become a stale lock"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user