fix: switch PrsPathState.mu to std::sync::Mutex, fixing purgatory expiry leak

tokio::sync::Mutex was used for a purely-synchronous critical section
(git init, delete_ref, list_refs, remove_dir_all — no .await inside the
guard). This caused two problems:

1. The purgatory sweep used try_lock() because blocking_lock() panics
   inside a tokio runtime. When the lock was contended, the sweep fell
   through to self.pr_events.remove() unconditionally, dropping the
   purgatory entry and leaking the dangling refs/nostr/<event-id> ref
   until the next process restart.

2. The tokio::sync::Mutex was a latent hazard: one careless blocking_lock
   call inside the runtime would deadlock.

Fix: switch PrsPathState.mu to std::sync::Mutex<()>. All critical
sections are sync-only so no ergonomic loss. The purgatory sweep can now
call lock() directly, and pr_events.remove() only runs after the
filesystem cleanup succeeds. The async call sites in receive.rs and
pr_event.rs drop the .await (the guard is dropped before any await
point). ensure_repo_initialised is converted from async to sync since it
only uses std::process::Command and std::fs internally.
This commit is contained in:
DanConwayDev
2026-05-16 10:51:10 +01:00
parent 661bf584d7
commit b6e9d11d88
3 changed files with 45 additions and 48 deletions
+11 -5
View File
@@ -57,8 +57,8 @@ use hyper::body::Bytes;
use hyper::{Response, StatusCode};
use nostr_relay_builder::LocalRelay;
use nostr_sdk::prelude::*;
use std::sync::Mutex;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::sync::Mutex;
use tracing::{debug, error, info, warn};
use crate::git::authorization::{
@@ -98,6 +98,12 @@ use crate::sync::rejected_index::RejectedEventsIndex;
/// `in_flight.load() == 0` *and* `list_refs` returns empty while they
/// hold `mu` — the same mutex that gates `in_flight` updates — so a
/// repo can never be deleted while a push is mid-receive.
///
/// `mu` is a `std::sync::Mutex` (not `tokio::sync::Mutex`) because every
/// critical section is purely synchronous (git I/O, no `.await`). This
/// lets the purgatory sweep call `lock()` directly from sync context
/// instead of `try_lock()`, eliminating the leak-on-contention bug where
/// a purgatory entry was dropped even when the lock was busy.
pub struct PrsPathState {
pub mu: Mutex<()>,
pub in_flight: AtomicUsize,
@@ -240,8 +246,8 @@ pub async fn handle_prs_receive_pack(
let state = path_state(&repo_init_locks, &repo_path);
{
let _g = state.mu.lock().await;
if let Err(e) = ensure_repo_initialised(&repo_path).await {
let _g = state.mu.lock().expect("prs path mutex poisoned");
if let Err(e) = ensure_repo_initialised(&repo_path) {
error!(
"/prs/ receive-pack: failed to initialise repo at {}: {}",
repo_path.display(),
@@ -307,7 +313,7 @@ pub async fn handle_prs_receive_pack(
// leak it. Decrements `in_flight`; if no other push is in flight
// and the repo has zero refs left, removes the bare directory.
{
let _g = state.mu.lock().await;
let _g = state.mu.lock().expect("prs path mutex poisoned");
state.in_flight.fetch_sub(1, Ordering::Relaxed);
if state.in_flight.load(Ordering::Relaxed) == 0 {
if let Ok(refs) = list_refs(&repo_path) {
@@ -400,7 +406,7 @@ fn invalid_ref_reason(ref_name: &str) -> Option<String> {
///
/// The caller must hold the per-path mutex from [`PrsPathState`] for
/// `repo_path` before invoking this function.
async fn ensure_repo_initialised(repo_path: &Path) -> Result<(), GitError> {
fn ensure_repo_initialised(repo_path: &Path) -> Result<(), GitError> {
if repo_path.exists() {
return Ok(());
}
+1 -1
View File
@@ -129,7 +129,7 @@ impl PrEventPolicy {
&self.ctx.repo_init_locks,
&prs_repo,
);
let _guard = state.mu.lock().await;
let _guard = state.mu.lock().expect("prs path mutex poisoned");
if let Err(e) = crate::git::delete_ref(&prs_repo, &ref_name) {
tracing::warn!(
+33 -42
View File
@@ -1297,14 +1297,15 @@ impl Purgatory {
// For GRASP-06 `/prs/`-scoped placeholders, best-effort
// filesystem cleanup of the dangling refs/nostr/<event-id>
// ref and (if it leaves the bare repo with zero refs *and
// no push is in flight*) the repo directory itself. The
// cleanup is gated on a `try_lock` of the per-path mutex:
// if a push is currently holding the lock the cleanup is
// skipped this cycle and the dangling ref is simply left
// on disk. With the lock held, the `in_flight` check is
// what distinguishes "no push active" from "push mid-receive
// but the mutex is briefly free between init and end-of-push"
// — only the former is safe to follow with a `rm -rf`.
// no push is in flight*) the repo directory itself.
//
// We hold the per-path mutex for the duration of the
// filesystem operations. Because `mu` is a `std::sync::Mutex`
// (all critical sections are sync-only), we can call `lock()`
// directly here without a runtime hazard. The purgatory entry
// is only removed from `pr_events` after the filesystem work
// succeeds; if the lock is poisoned we skip cleanup this cycle
// and let the entry expire on the next sweep.
if let (Some(scope), Some(ctx)) = (prs_scope.as_ref(), self.prs_cleanup_ctx.get()) {
let repo_path = crate::grasp06::paths::prs_repo_path(
&ctx.git_data_path,
@@ -1314,44 +1315,34 @@ impl Purgatory {
if repo_path.exists() {
let state =
crate::grasp06::receive::path_state(&ctx.repo_init_locks, &repo_path);
match state.mu.try_lock() {
Ok(_guard) => {
let ref_name = format!("refs/nostr/{}", event_id_str);
if let Err(e) = crate::git::delete_ref(&repo_path, &ref_name) {
tracing::warn!(
repo = %repo_path.display(),
ref_name = %ref_name,
error = %e,
"Failed to delete dangling /prs/ ref during purgatory expiry",
);
} else if state.in_flight.load(std::sync::atomic::Ordering::Relaxed)
== 0
&& matches!(
crate::git::list_refs(&repo_path),
Ok(refs) if refs.is_empty()
)
{
if let Err(e) = std::fs::remove_dir_all(&repo_path) {
tracing::warn!(
repo = %repo_path.display(),
error = %e,
"Failed to remove zero-ref /prs/ repo during purgatory expiry",
);
} else {
tracing::debug!(
repo = %repo_path.display(),
"Removed zero-ref /prs/ repo during purgatory expiry",
);
}
}
}
Err(_) => {
let _guard = state.mu.lock().expect("prs path mutex poisoned");
let ref_name = format!("refs/nostr/{}", event_id_str);
if let Err(e) = crate::git::delete_ref(&repo_path, &ref_name) {
tracing::warn!(
repo = %repo_path.display(),
ref_name = %ref_name,
error = %e,
"Failed to delete dangling /prs/ ref during purgatory expiry",
);
} else if state.in_flight.load(std::sync::atomic::Ordering::Relaxed) == 0
&& matches!(
crate::git::list_refs(&repo_path),
Ok(refs) if refs.is_empty()
)
{
if let Err(e) = std::fs::remove_dir_all(&repo_path) {
tracing::warn!(
repo = %repo_path.display(),
error = %e,
"Failed to remove zero-ref /prs/ repo during purgatory expiry",
);
} else {
tracing::debug!(
repo = %repo_path.display(),
"/prs/ repo lock contended during purgatory expiry — skipping filesystem cleanup this cycle",
"Removed zero-ref /prs/ repo during purgatory expiry",
);
}
};
}
}
}